mirror of
https://github.com/razzant/ouroboros.git
synced 2026-10-03 04:07:04 +00:00
Store exact delegated request bodies outside event rows
Retain complete raw request envelopes before durable start references, replay legacy inline and CAS-backed requests, and remove only the equal task-contract event mirror. Focused request/custody tests and the real S11 mock flow passed; phase-wide checks and external review pending.
This commit is contained in:
parent
40a93e091a
commit
598895ca72
9 changed files with 276 additions and 12 deletions
|
|
@ -124,7 +124,7 @@ scanned data-relative path to be covered by a row here (count-anchored both ways
|
|||
|---|---|---|---|---|
|
||||
| `logs/chat.jsonl` | `supervisor/message_bus.py` (+presence, project summaries) | `direction` + optional `type`; no version — accepted (projection replayed by chain-aware readers) | rotated 800 KB → `archive/chat_*.jsonl`; archive chain WARN at 100 MB | newest generation lost; consolidation cursor reports gap (recoverable) |
|
||||
| `logs/progress.jsonl` | `supervisor/message_bus.py` (+plan review) | `type: send_message`, `is_progress` | rotated 800 KB → `archive/progress_*.jsonl`; 8 MB WARN = rotation broken | current segment lost; readers archive-chain-aware |
|
||||
| `logs/events.jsonl` | ~60 modules via `append_jsonl` (+`delegate_custody.emit`) | universal `type` discriminator — accepted (per-type payloads owned by emitters) | rotated 800 KB → `archive/events_*.jsonl`; custody readers (replay, fault tail-scan, `complete_custody_rows`, settled-terminal chain cursor, legacy-usage import, swarm rollup, worker-boot verify) are chain-aware; 8 MB live WARN = rotation broken; 100 MB chain WARN = replay degradation | delegated-run custody destroyed (chain incl. archive segments): open runs invisible/unreapable; lineage, citations, legacy-usage source lost |
|
||||
| `logs/events.jsonl` | ~60 modules via `append_jsonl` (+`delegate_custody.emit`) | universal `type` discriminator — accepted (per-type payloads owned by emitters) | rotated 800 KB → `archive/events_*.jsonl`; custody readers (replay, fault tail-scan, `complete_custody_rows`, settled-terminal chain cursor, legacy-usage import, swarm rollup, worker-boot verify) are chain-aware; delegated start rows reference their complete raw replay envelope in `observability/blobs/**` (blob before event), with inline legacy rows still readable; `task_received` omits only an equal metadata mirror of its top-level task contract; 8 MB live WARN = rotation broken; 100 MB chain WARN = replay degradation | delegated-run custody destroyed (chain incl. archive segments): open runs invisible/unreapable; lineage, citations, legacy-usage source lost |
|
||||
| `logs/tools.jsonl` | `ouroboros/loop_tool_execution.py` (+budget-drive mirror) | `type: tool_call`, untruncated args | rotated 800 KB → `archive/tools_*.jsonl`; tail readers (api_logs_tail, task_events) archive-backfill; 8 MB WARN = rotation broken | untruncated tool record + `result_ref` pointers lost |
|
||||
| `logs/supervisor.jsonl` | supervisor family, `process_custody`, gateway control, server shutdown | `type` (+secondary `event_type`) | rotated 800 KB → `archive/supervisor_*.jsonl` + 8 MB tripwire; tail readers (`memory.read_jsonl_tail`, api_logs_tail) archive-backfill | reap receipts, rescue disclosures, shutdown causes lost |
|
||||
| `logs/task_reflections.jsonl` | `ouroboros/reflection.py` (+ project-scoped copy under `projects/<id>/logs/`) | full rows unversioned; pointer rows `type: project_reflection_pointer` | rotated 800 KB → `archive/task_reflections_*.jsonl` + 8 MB tripwire; tail-20 read archive-backfills; project-scoped copies follow project retention (never age-pruned) | inter-task memory-carry signal lost |
|
||||
|
|
|
|||
|
|
@ -268,7 +268,7 @@ An undisclosed spend contributes `0.0` to `accounted_usd` — inventing a conser
|
|||
|
||||
**Transport.** `gateways/claudexor.py` is pure transport (descriptor read, `/v2` handshake, the `config.CLAUDEXOR_MIN_VERSION` 3.2.0 floor). The daemon bearer token grants the entire `/v2` surface, so it never leaves this module, and the HTTP client runs `trust_env=False` so an ambient proxy cannot intercept the loopback control plane. Production starts obtain a handshaken owned gateway from `claudexor_daemon.ensure_owned_gateway` — exact reviewed engine/Node pins (`claudexor_runtime.py`; the reviewed pin IS the next-spawn selection — no mutable `current` pointer, no background updater — and `OUROBOROS_CLAUDEXOR_BIN` is the explicit operator override), never a PATH install. Keeping lifecycle above transport keeps account status side-effect-free and harness mechanisms out of Ouroboros (`claudexor_daemon.py`/`claudexor_runtime.py` docstrings; stop path, spawn latch and typed start failures: §9). Native child processes are contained by env token (`process_containment.py`, `OURO_PROC_CONTAINER_*`), because a surviving descendant can become invisible to parent-child traversal once its controller exits; an alive-or-undeterminable member is an honest hard-block answer, never a kill guarantee.
|
||||
|
||||
**Custody is durable, because the run is not ours to kill.** A delegated run lives inside the daemon and survives our worker, so the AUTHORITY is the durable `delegate_run_*` custody rows (`delegate_run_started` and friends) on the canonical/budget root (`ouroboros/delegate_custody.py`); one compact incident projection, `<drive_root>/logs/containment_faults.jsonl`, exists because the event log grows without bound and a tail-bounded scan can bury an unresolved fault. A lookup answers OWNED, FOREIGN or UNKNOWN — collapsing UNKNOWN into "not yours" made a restarted owner indistinguishable from an intruder. An ABSENT custody log is a positively established clean state; an EXISTING-but-unreadable one audits as `delegated_run_state_unknown:custody_log_unreadable`, never cleanly reconciled. Every INTENDED start mints a fresh per-intention invocation UUID as the wire `Idempotency-Key`; the content hash is only the LOOKUP identity — a content-stable wire key would hand a deliberate re-run the finished old run — and reuse happens only by explicit token (`pending_invocation_id`/`retry_of`, replaying the STORED canonical body under the SAME key). `reconcile_orphaned_runs` visits every open run whose owning task left the live set and settles the terminal ones, but it CANCELS only behind a deliberate verdict: a durable owner result that is readable, truly terminal and finished by the task's own decision. A custody row carries its owner's kind, and task association confers no lifecycle authority: a run a review surface registered (`RunCustody.review_owned`, durable `source` under the review substrate) belongs to its panel, bounded by the slot window and its own `maxSeconds`, so the sweep spares it unless the owner task was itself cancelled — a `left_live` row names the panel, and a task that consciously finished under a running acceptance panel keeps its reviewer alive; a pending review invocation is likewise retained, never re-posted by the generic recovery, because the review substrate owns its rejoin. The owner-cancel kill boundary STATES the verdict it is about to write, because it audits custody before that write, so an owner cancellation stops the paid run at the boundary rather than at the next sweep. A provider or transport death, a worker crash, a missing or unreadable result all SPARE the run, left live with a durable `left_live` reconciliation row: an undignified nanny death must not kill a healthy paid run, and "unknown" is exactly that case. `maxSeconds` is the damage limitation for a spared orphan, never custody. SETTLED is published before registration retirement, after the ledger row lands; settlement and registration retirement are separate durable duties, a failed retirement stays replayable on `project_owned` for the later sweep, and writing `settled` over a suppressed ledger failure would turn a lock timeout into a permanent leak. A start whose row did not land reports `started_uncustodied`: no supervision, no replacement, until the original run is proven absent or terminal.
|
||||
**Custody is durable, because the run is not ours to kill.** A delegated run lives inside the daemon and survives our worker, so the AUTHORITY is the durable `delegate_run_*` custody rows (`delegate_run_started` and friends) on the canonical/budget root (`ouroboros/delegate_custody.py`); one compact incident projection, `<drive_root>/logs/containment_faults.jsonl`, exists because the event log grows without bound and a tail-bounded scan can bury an unresolved fault. A lookup answers OWNED, FOREIGN or UNKNOWN — collapsing UNKNOWN into "not yours" made a restarted owner indistinguishable from an intruder. An ABSENT custody log is a positively established clean state; an EXISTING-but-unreadable one audits as `delegated_run_state_unknown:custody_log_unreadable`, never cleanly reconciled. Every INTENDED start mints a fresh per-intention invocation UUID as the wire `Idempotency-Key`; the content hash is only the LOOKUP identity — a content-stable wire key would hand a deliberate re-run the finished old run — and reuse happens only by explicit token (`pending_invocation_id`/`retry_of`, replaying the STORED canonical body under the SAME key). The complete replay envelope lives unredacted in the existing private observability CAS, written before `delegate_run_start_requested`; the event carries `request_ref` and `prompt_chars`, keeping a large reviewer packet out of every custody scan. `delegate_pending.request_body` resolves legacy inline first, then verifies the CAS reference for both invocation readers. JSON values and the canonical request digest stay unchanged, including the thread fields needed to reproduce its wire projection; a missing/corrupt blob leaves the request unknown for the existing typed recovery refusals. Pending scans resolve only surviving invocations, while legacy event/archive bytes remain untouched. `reconcile_orphaned_runs` visits every open run whose owning task left the live set and settles the terminal ones, but it CANCELS only behind a deliberate verdict: a durable owner result that is readable, truly terminal and finished by the task's own decision. A custody row carries its owner's kind, and task association confers no lifecycle authority: a run a review surface registered (`RunCustody.review_owned`, durable `source` under the review substrate) belongs to its panel, bounded by the slot window and its own `maxSeconds`, so the sweep spares it unless the owner task was itself cancelled — a `left_live` row names the panel, and a task that consciously finished under a running acceptance panel keeps its reviewer alive; a pending review invocation is likewise retained, never re-posted by the generic recovery, because the review substrate owns its rejoin. The owner-cancel kill boundary STATES the verdict it is about to write, because it audits custody before that write, so an owner cancellation stops the paid run at the boundary rather than at the next sweep. A provider or transport death, a worker crash, a missing or unreadable result all SPARE the run, left live with a durable `left_live` reconciliation row: an undignified nanny death must not kill a healthy paid run, and "unknown" is exactly that case. `maxSeconds` is the damage limitation for a spared orphan, never custody. SETTLED is published before registration retirement, after the ledger row lands; settlement and registration retirement are separate durable duties, a failed retirement stays replayable on `project_owned` for the later sweep, and writing `settled` over a suppressed ledger failure would turn a lock timeout into a permanent leak. A start whose row did not land reports `started_uncustodied`: no supervision, no replacement, until the original run is proven absent or terminal.
|
||||
|
||||
**No terminal or cancel claim without a verified receipt.** `delegate_cancel` returns `confirmed` (read back terminal), `requested`, `failed` or `containment_fault_run_may_still_be_live`; the last two hold a durable CRITICAL containment fault until a receipt or settlement clears it — an overpowered run that may still be alive is an incident, not a reassuring string — and the state read decides, so a refused control is never a verdict about the RUN. One `daemon_says_absent` predicate decides everywhere that a 404 is the daemon ANSWERING that the resource is gone (scoped to the daemon that answered), never a failure to find out; custody closes such a run `delegate_run_closed_absent` (unreachable, not settled), inventing no terminal detail, usage or spend. Results are delivered, not severed: `delegate_wait` stages the whole terminal detail atomically under `task_drive/delegated_runs/<run>.json` with a typed `output_delivery` block, and cut fields are renamed `*_preview` so a partial read of head-truncated JSON fails loudly instead of looking like an answer.
|
||||
|
||||
|
|
|
|||
|
|
@ -847,7 +847,9 @@ def hot_store_growth_notes(env: Any) -> list:
|
|||
"WARNING: HOT STORE GROWTH — the events chain (logs/events.jsonl + "
|
||||
f"archive/events_*.jsonl) totals {events_chain_size / 1_000_000:.1f} MB "
|
||||
f"(threshold {EVENTS_ARCHIVE_SCAN_WARN_BYTES // 1_000_000} MB). Custody "
|
||||
"replay scans this chain on ownership questions. Investigate chain "
|
||||
"replay scans this chain on ownership questions. Legacy segments retain "
|
||||
"inline delegated request bodies; new start rows reference the observability "
|
||||
"store, without shrinking existing history. Investigate chain "
|
||||
"indexing/compaction; archives are durable history and are never deleted."
|
||||
)
|
||||
from ouroboros.context_budget import RETAINED_EXECUTION_DRIVES_WARN_COUNT
|
||||
|
|
|
|||
|
|
@ -711,6 +711,8 @@ def invocation_record(drive_root: Any, invocation_id: str, *,
|
|||
and isolation facts are likewise replayed rather than re-derived.
|
||||
``rows`` reuses a caller's single event snapshot, as the other replay views do.
|
||||
"""
|
||||
from ouroboros.delegate_pending import request_body
|
||||
|
||||
target = str(invocation_id or "").strip()
|
||||
if not target:
|
||||
return None
|
||||
|
|
@ -726,7 +728,7 @@ def invocation_record(drive_root: Any, invocation_id: str, *,
|
|||
"surface": str(row.get("surface") or ""),
|
||||
"slot_id": str(row.get("slot_id") or ""),
|
||||
"operation_id": str(row.get("operation_id") or ""),
|
||||
"request": row.get("request") if isinstance(row.get("request"), dict) else None,
|
||||
"request": request_body(drive_root, row),
|
||||
"route": str(row.get("route") or ""),
|
||||
"project_id": str(row.get("project_id") or ""),
|
||||
"project_owned": bool(row.get("project_owned")),
|
||||
|
|
@ -774,7 +776,22 @@ def record_start_requested(drive_root: Any, **payload: Any) -> bool:
|
|||
Returns whether the row LANDED; the caller must not POST when it did not —
|
||||
a run whose request row never reached disk is live, mutating and unfindable
|
||||
if the worker dies before ``record_started``.
|
||||
|
||||
The full replay envelope goes to raw CAS before its event reference. Use
|
||||
``write_blob``, never a redacted ``persist_call`` projection: request values
|
||||
must retain the engine's canonical JSON digest on an idempotent retry.
|
||||
"""
|
||||
body = payload.get("request")
|
||||
if isinstance(body, dict) and body:
|
||||
from ouroboros.observability import write_blob
|
||||
|
||||
try:
|
||||
ref = write_blob(pathlib.Path(drive_root), body, kind="json")
|
||||
except Exception:
|
||||
log.warning("delegate custody request body could not be stored", exc_info=True)
|
||||
return False
|
||||
payload = {key: value for key, value in payload.items() if key != "request"}
|
||||
payload.update(request_ref=ref, prompt_chars=len(str(body.get("prompt") or "")))
|
||||
return emit(drive_root, START_REQUESTED, payload)
|
||||
|
||||
|
||||
|
|
|
|||
|
|
@ -2,6 +2,7 @@
|
|||
|
||||
from __future__ import annotations
|
||||
|
||||
import pathlib
|
||||
from typing import Any, Dict, List, Optional
|
||||
|
||||
|
||||
|
|
@ -28,6 +29,7 @@ def pending_invocations(
|
|||
"slot_id": str(row.get("slot_id") or ""),
|
||||
"operation_id": str(row.get("operation_id") or ""),
|
||||
"request": row.get("request") if isinstance(row.get("request"), dict) else None,
|
||||
"request_ref": row.get("request_ref"),
|
||||
"route": str(row.get("route") or ""),
|
||||
"project_id": str(row.get("project_id") or ""),
|
||||
"project_owned": bool(row.get("project_owned")),
|
||||
|
|
@ -65,11 +67,38 @@ def pending_invocations(
|
|||
and state.get(invocation_id) != "started"
|
||||
):
|
||||
state[invocation_id] = "failed_definite"
|
||||
return [
|
||||
record for invocation_id, record in found.items()
|
||||
if state.get(invocation_id, "pending") == "pending"
|
||||
and isinstance(record["request"], dict) and record["request"]
|
||||
]
|
||||
pending = []
|
||||
for invocation_id, record in found.items():
|
||||
if state.get(invocation_id, "pending") != "pending":
|
||||
continue
|
||||
# Resolve only survivors, not every historical start on each sweep.
|
||||
body = request_body(drive_root, record)
|
||||
record.pop("request_ref")
|
||||
if body:
|
||||
record["request"] = body
|
||||
pending.append(record)
|
||||
return pending
|
||||
|
||||
|
||||
__all__ = ["pending_invocations"]
|
||||
def request_body(drive_root: Any, row: Dict[str, Any]) -> Optional[Dict[str, Any]]:
|
||||
"""Resolve the canonical replay envelope, legacy inline first, then raw CAS.
|
||||
|
||||
Unreadable references leave the request unknown for the caller's existing
|
||||
refusal paths; never rebuild a paid invocation from current settings.
|
||||
"""
|
||||
inline = row.get("request")
|
||||
if isinstance(inline, dict) and inline:
|
||||
return inline
|
||||
ref = row.get("request_ref")
|
||||
if not isinstance(ref, dict) or not ref:
|
||||
return None
|
||||
from ouroboros.observability import read_blob_ref
|
||||
|
||||
try:
|
||||
body = read_blob_ref(pathlib.Path(drive_root), ref, expected_kind="json")
|
||||
except Exception:
|
||||
return None
|
||||
return body if isinstance(body, dict) and body else None
|
||||
|
||||
|
||||
__all__ = ["pending_invocations", "request_body"]
|
||||
|
|
|
|||
|
|
@ -1179,6 +1179,11 @@ def sanitize_task_for_event(
|
|||
metadata["origin_message_text"] = truncate_for_log(nested, threshold)
|
||||
sanitized["metadata"] = metadata
|
||||
|
||||
# The event needs one contract; the live task keeps both mirrors.
|
||||
contract = sanitized.get("task_contract")
|
||||
if isinstance(metadata, dict) and contract and metadata.get("task_contract") == contract:
|
||||
sanitized["metadata"] = {key: value for key, value in metadata.items() if key != "task_contract"}
|
||||
|
||||
text = task.get("text")
|
||||
if not isinstance(text, str):
|
||||
return sanitized
|
||||
|
|
|
|||
|
|
@ -407,6 +407,9 @@ class TestHotStoreGrowthInvariant:
|
|||
assert "HOT STORE GROWTH" in result
|
||||
assert "events chain" in result
|
||||
assert "never deleted" in result
|
||||
assert "Legacy segments retain inline delegated request bodies" in result
|
||||
assert "without shrinking existing history" in result
|
||||
assert segment.stat().st_size == EVENTS_ARCHIVE_SCAN_WARN_BYTES + 1
|
||||
|
||||
def test_isolated_benchmark_sentinel_suppresses_warnings(self, tmp_path):
|
||||
from supervisor.state import ISOLATED_BENCHMARK_SENTINEL
|
||||
|
|
|
|||
204
tests/test_delegate_request_storage.py
Normal file
204
tests/test_delegate_request_storage.py
Normal file
|
|
@ -0,0 +1,204 @@
|
|||
"""Canonical delegation requests stay replayable outside the event row."""
|
||||
|
||||
import copy
|
||||
import hashlib
|
||||
import json
|
||||
from pathlib import Path
|
||||
|
||||
import pytest
|
||||
|
||||
from ouroboros import delegate_custody as custody
|
||||
from ouroboros import observability
|
||||
from ouroboros.delegate_pending import request_body
|
||||
from tests._delegated_transport_shared import (
|
||||
_LiveRunStub, _nanny_ctx,
|
||||
_owned_gateway_uses_each_test_transport, # noqa: F401 -- autouse fixture
|
||||
)
|
||||
from tests._review_session_route_shared import fake_route as _fake_route
|
||||
|
||||
fake_route = _fake_route
|
||||
|
||||
|
||||
def _start_rows(root):
|
||||
return [row for row in custody._iter_rows(custody.event_log_path(root))
|
||||
if row["type"] == custody.START_REQUESTED]
|
||||
|
||||
|
||||
def _digest(body):
|
||||
# The engine digests parsed JSON with stable keys, not wire whitespace.
|
||||
raw = json.dumps(body, ensure_ascii=False, sort_keys=True, separators=(",", ":"))
|
||||
return hashlib.sha256(raw.encode("utf-8")).hexdigest()
|
||||
|
||||
|
||||
@pytest.mark.parametrize("surface", ["delegate", "review"])
|
||||
@pytest.mark.parametrize("failure", [None, "blob", "event"])
|
||||
def test_canonical_body_and_event_precede_dispatch(tmp_path, monkeypatch, request, surface, failure):
|
||||
from ouroboros.gateways import claudexor as gateway
|
||||
from ouroboros.review_execution import ReviewRouteUnavailable
|
||||
from ouroboros.tools import delegate
|
||||
from tests._review_session_route_shared import _run_session_directly
|
||||
|
||||
prompt = "Inspect these exact bytes: Привет\n" + "x" * 50_000
|
||||
order, bodies = [], []
|
||||
real_write, real_emit = observability.write_blob, custody.emit
|
||||
|
||||
def write(root, body, **kwargs):
|
||||
if isinstance(body, dict) and body.get("prompt") == prompt:
|
||||
assert order == []
|
||||
if failure == "blob":
|
||||
raise OSError("fixture refuses body storage")
|
||||
bodies.append(copy.deepcopy(body))
|
||||
ref = real_write(root, body, **kwargs)
|
||||
order.append("blob")
|
||||
return ref
|
||||
return real_write(root, body, **kwargs)
|
||||
|
||||
def emit(root, kind, payload):
|
||||
if kind == custody.START_REQUESTED:
|
||||
assert order == ["blob"]
|
||||
assert "request" not in payload
|
||||
assert request_body(root, payload) == bodies[0]
|
||||
if failure == "event":
|
||||
return False
|
||||
order.append("event")
|
||||
return real_emit(root, kind, payload)
|
||||
|
||||
if surface == "review":
|
||||
stub_type = request.getfixturevalue("fake_route")
|
||||
else:
|
||||
stub = _LiveRunStub()
|
||||
stub_type = type(stub)
|
||||
monkeypatch.setenv("OUROBOROS_SUBAGENT_HARNESS", "some-route=weak-model:low")
|
||||
monkeypatch.setattr(gateway, "ClaudexorGateway", lambda *a, **k: stub)
|
||||
original_start = stub_type.start_run
|
||||
|
||||
def dispatch(self, body, *, idempotency_key=""):
|
||||
assert order == ["blob", "event"]
|
||||
stored = custody.invocation_record(tmp_path, idempotency_key)
|
||||
assert stored["request"] == body == bodies[0]
|
||||
assert _digest(stored["request"]) == _digest(body)
|
||||
order.append("dispatch")
|
||||
return original_start(self, body, idempotency_key=idempotency_key)
|
||||
|
||||
monkeypatch.setattr(observability, "write_blob", write)
|
||||
monkeypatch.setattr(custody, "emit", emit)
|
||||
monkeypatch.setattr(stub_type, "start_run", dispatch)
|
||||
if surface == "review":
|
||||
if failure:
|
||||
with pytest.raises(ReviewRouteUnavailable) as caught:
|
||||
_run_session_directly(tmp_path, prompt=prompt)
|
||||
assert caught.value.code == "start_request_row_unwritable"
|
||||
else:
|
||||
_run_session_directly(tmp_path, prompt=prompt)
|
||||
else:
|
||||
result = json.loads(delegate._delegate_start(_nanny_ctx(tmp_path), prompt).text)
|
||||
assert result["status"] == ("refused" if failure else "started")
|
||||
if failure:
|
||||
assert result["reason"] == "start_request_row_unwritable"
|
||||
if failure:
|
||||
assert "dispatch" not in order
|
||||
assert _start_rows(tmp_path) == []
|
||||
else:
|
||||
row = _start_rows(tmp_path)[0]
|
||||
assert row["prompt_chars"] == len(prompt)
|
||||
assert len(json.dumps(row).encode("utf-8")) < 2_000
|
||||
assert order == ["blob", "event", "dispatch"]
|
||||
|
||||
|
||||
def test_raw_review_envelope_and_thread_projection_round_trip(tmp_path):
|
||||
from ouroboros.review_thread_continuity import start_review_thread_turn
|
||||
|
||||
body = {"prompt": "Привет\n" + "x" * 80_000,
|
||||
"instructions": "Authorization: Bearer EXAMPLE-NOT-A-REAL-CREDENTIAL",
|
||||
"authPreference": "subscription", "mode": "ask", "access": "readonly",
|
||||
"scope": {"kind": "project", "root": str(tmp_path)},
|
||||
"harnesses": ["codex"], "primaryHarness": "codex", "maxSeconds": 300,
|
||||
"credentialProfileId": None, "_use_thread": True, "_thread_id": "thread-1",
|
||||
"execution": {"delegated": False, "isolation": "in_place"},
|
||||
"outputSchema": {"type": "object", "properties": {"passed": {"type": "boolean"}}}}
|
||||
original = copy.deepcopy(body)
|
||||
assert custody.record_start_requested(tmp_path, invocation_id="review", task_id="task", request=body)
|
||||
row = _start_rows(tmp_path)[0]
|
||||
assert row["request_ref"]["kind"] == "json"
|
||||
restored = custody.invocation_record(tmp_path, "review")["request"]
|
||||
assert restored == body == original
|
||||
assert _digest(restored) == _digest(body)
|
||||
projections = []
|
||||
|
||||
class Gateway:
|
||||
def start_thread_turn(self, thread_id, request, *, idempotency_key):
|
||||
projections.append((thread_id, request, idempotency_key))
|
||||
return {"runId": "run"}
|
||||
|
||||
for envelope in (body, restored):
|
||||
start_review_thread_turn(Gateway(), envelope["_thread_id"], envelope, idempotency_key="review")
|
||||
assert projections[0] == projections[1]
|
||||
assert not {"_thread_id", "_use_thread", "scope", "execution"} & projections[0][1].keys()
|
||||
|
||||
|
||||
def test_legacy_inline_wins_even_beside_an_unreadable_ref(tmp_path, monkeypatch):
|
||||
inline = {"prompt": "legacy request"}
|
||||
assert custody.emit(tmp_path, custody.START_REQUESTED, {
|
||||
"invocation_id": "legacy", "request": inline, "request_ref": {"path": "missing"}})
|
||||
def refuse_read(*args, **kwargs):
|
||||
raise AssertionError("legacy inline must not read a blob")
|
||||
monkeypatch.setattr(observability, "read_blob_ref", refuse_read)
|
||||
assert custody.invocation_record(tmp_path, "legacy")["request"] == inline
|
||||
assert custody.pending_invocations(tmp_path)[0]["request"] == inline
|
||||
|
||||
|
||||
@pytest.mark.parametrize("damage", ["missing", "invalid_size", "corrupt"])
|
||||
def test_unreadable_ref_keeps_request_unknown_and_retry_refuses(tmp_path, damage):
|
||||
from ouroboros.tools.delegate_integration import _validated_invocation
|
||||
|
||||
assert custody.record_start_requested(tmp_path, invocation_id="lost", task_id="task",
|
||||
request={"prompt": "exact original"})
|
||||
row = _start_rows(tmp_path)[0]
|
||||
if damage == "missing":
|
||||
Path(row["request_ref"]["path"]).unlink()
|
||||
elif damage == "corrupt":
|
||||
Path(row["request_ref"]["path"]).write_bytes(b"not gzip")
|
||||
else:
|
||||
row["request_ref"]["size"] = "invalid"
|
||||
custody.event_log_path(tmp_path).write_text(json.dumps(row) + "\n", encoding="utf-8")
|
||||
assert custody.invocation_record(tmp_path, "lost")["request"] is None
|
||||
assert custody.pending_invocations(tmp_path) == []
|
||||
record, refusal = _validated_invocation(tmp_path, "lost", "task", "exact original")
|
||||
assert record is None
|
||||
assert json.loads(refusal.text)["reason"] == "invocation_request_unrecorded"
|
||||
|
||||
|
||||
def test_pending_scan_only_reads_blobs_for_surviving_invocations(tmp_path, monkeypatch):
|
||||
for identity in ("started", "refused", "pending"):
|
||||
assert custody.record_start_requested(tmp_path, invocation_id=identity,
|
||||
request={"prompt": identity})
|
||||
assert custody.emit(tmp_path, custody.STARTED, {"invocation_id": "started", "run_id": "run"})
|
||||
assert custody.emit(tmp_path, custody.START_FAILED, {"invocation_id": "refused", "definite": True})
|
||||
reads, original = [], observability.read_blob_ref
|
||||
def read(*args, **kwargs):
|
||||
body = original(*args, **kwargs)
|
||||
reads.append(body["prompt"])
|
||||
return body
|
||||
monkeypatch.setattr(observability, "read_blob_ref", read)
|
||||
pending = custody.pending_invocations(tmp_path)
|
||||
assert [row["invocation_id"] for row in pending] == ["pending"]
|
||||
assert reads == ["pending"]
|
||||
assert "request_ref" not in pending[0]
|
||||
|
||||
|
||||
@pytest.mark.parametrize("different", [False, True])
|
||||
def test_task_event_omits_only_an_equal_contract_mirror_without_mutation(tmp_path, different):
|
||||
from ouroboros.contracts.task_contract import attach_task_contract
|
||||
from ouroboros.utils import sanitize_task_for_event
|
||||
|
||||
task = attach_task_contract({"id": "task", "text": "Inspect the files", "metadata": {"source": "web"}})
|
||||
if different:
|
||||
task["metadata"]["task_contract"] = {"other": "contract"}
|
||||
original = copy.deepcopy(task)
|
||||
projected = sanitize_task_for_event(task, tmp_path / "logs")
|
||||
assert task == original
|
||||
assert projected["task_contract"] == original["task_contract"]
|
||||
assert projected["metadata"]["source"] == "web"
|
||||
assert ("task_contract" in projected["metadata"]) is different
|
||||
if different:
|
||||
assert projected["metadata"]["task_contract"] == original["metadata"]["task_contract"]
|
||||
|
|
@ -309,17 +309,21 @@ def test_a_retry_testifies_about_the_stored_invocation_not_the_current_config(
|
|||
assert retried["root"] == str(root_a)
|
||||
|
||||
# ... and so do the durable rows, attempt and custody alike.
|
||||
from ouroboros.delegate_pending import request_body
|
||||
|
||||
rows = [json.loads(line) for line
|
||||
in (drive / "logs" / "events.jsonl").read_text().splitlines() if line.strip()]
|
||||
in (drive / "logs" / "events.jsonl").read_text(encoding="utf-8").splitlines() if line.strip()]
|
||||
attempts = [r for r in rows if r.get("type") == dc.START_REQUESTED
|
||||
and r.get("invocation_id") == token]
|
||||
started = [r for r in rows if r.get("type") == dc.STARTED
|
||||
and r.get("run_id") == retried["run_id"]][-1]
|
||||
original = attempts[0]
|
||||
assert request_body(drive, original) == bodies[0]
|
||||
for row in attempts[1:]:
|
||||
for fact in ("route", "project_id", "project_owned", "idempotency_key",
|
||||
"max_seconds", "request"):
|
||||
"max_seconds", "request_ref", "prompt_chars"):
|
||||
assert row[fact] == original[fact], f"retry attempt re-derived {fact}"
|
||||
assert request_body(drive, row) == bodies[0], "retry must retain the original canonical body"
|
||||
for fact, expected in (("route", "route-a"), ("model", "model-old"),
|
||||
("effort", "low"), ("root", str(root_a)),
|
||||
("project_id", prj_a), ("project_owned", True),
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue