mirror of
https://github.com/razzant/ouroboros.git
synced 2026-10-02 19:58:46 +00:00
Retain pending invocation custody when request blobs are unavailable
This commit is contained in:
parent
3ed6a1b815
commit
afbf9b1f85
7 changed files with 140 additions and 16 deletions
|
|
@ -168,7 +168,7 @@ scanned data-relative path to be covered by a row here (count-anchored both ways
|
|||
| `projects/<id>/knowledge_history.jsonl`, `projects/<id>/knowledge_journal.jsonl` | `ouroboros/knowledge.py` through the shared knowledge write lock | append rows retain source topic, revision and operation facts | follows the owning project shelf; retained with the project until explicit deletion | history/provenance lost while authored project notes remain |
|
||||
| `archive/**` (rotated segments, `rescue/`, `usage_import/`, `managed_repo/`) | rotation + `supervisor/git_ops_rescue.py`, `usage_legacy_import.py`, `launcher_bootstrap.py` | segments inherit source shape; usage_import carries sha256 sidecar | UNBOUNDED BY DESIGN — durable history, never GC'd (P1) — accepted | memory horizon truncated; rescue copies of uncommitted work destroyed |
|
||||
| `archive/usage_ledger/segment_*.jsonl` | `ouroboros/usage_compaction.py` (exact pre-compaction ledger bytes, written + fsync'd BEFORE the live swap) | each segment is a whole valid ledger generation; hash-pinned by the live `usage_baseline` header (`source_sha256`), chained recursively through each segment's own leading header | UNBOUNDED BY DESIGN — the folded monetary history, never GC'd (P1); read via `archived_attempt_ids` (tamper-evident, per-attempt joins for the model-send reverse sweep) | folded per-attempt monetary history unrecoverable; live aggregates (baseline block) survive, but seal/attempt joins for folded ids break — the mirror of the ledger row: deleting the ARCHIVE alone under a stamped ledger is a typed chain break on every history question, deleting the LEDGER alone raises `generation newer` for surviving newer, non-prefix segments; once fresh compactions reach those generations, old unreferenced segments are skipped and their attempt IDs are absent ([history readers](USAGE_COMPACTION.md#10-history-readers-model-send-reconciliation-and-audits)); reset both together |
|
||||
| `observability/{calls,blobs,salvaged}/**` | `ouroboros/observability.py` (private 0700/0600, CAS gzip); model-send records beside the call manifests: `model_send_seal` block + write-once `<attempt>.model_send_violation.json` typed facts (`ouroboros/model_send_seal.py`) | call manifests `schema_version: 1` + custody/redaction honesty markers; blob refs sha-verified on read; `model_send_seal.seal_version: 1` with the `canonical_json_v1` basis string | preserved indefinitely BY CONTRACT (the startup census counts, never deletes); the inert `OUROBOROS_OBSERVABILITY_RETENTION_DAYS` knob is RETIRED (key in `RETIRED_SETTING_KEYS`) | every recorded `result_ref`/`manifest_ref` dangles (strict readers raise); salvaged outputs unrecoverable; a lost seal on a seam-dispatched attempt surfaces as a typed `unlogged_attempt` fact at the next startup sweep |
|
||||
| `observability/{calls,blobs,salvaged}/**` | `ouroboros/observability.py` (private 0700/0600, CAS gzip); model-send records beside the call manifests: `model_send_seal` block + write-once `<attempt>.model_send_violation.json` typed facts (`ouroboros/model_send_seal.py`) | call manifests `schema_version: 1` + custody/redaction honesty markers; blob refs sha-verified on read; `model_send_seal.seal_version: 1` with the `canonical_json_v1` basis string | preserved indefinitely BY CONTRACT (the startup census counts, never deletes); the inert `OUROBOROS_OBSERVABILITY_RETENTION_DAYS` knob is RETIRED (key in `RETIRED_SETTING_KEYS`) | every recorded `result_ref`/`manifest_ref` dangles (strict readers raise); pending delegated request bodies become unknown, while custody identity still blocks replacement and protects snapshots; replay refuses without the recorded body; salvaged outputs unrecoverable; a lost seal on a seam-dispatched attempt surfaces as a typed `unlogged_attempt` fact at the next startup sweep |
|
||||
| `claudexor/**` | EXTERNAL writer — the claudexord daemon (Ouroboros only mkdirs, appends `daemon.log`, writes `ouroboros-owned.json` marker) | marker unversioned | daemon-owned; grows unbounded under our root — disclosed external plane | owner harness logins/profiles lost (fresh device-auth required) |
|
||||
| `playwright-browsers/` | `ouroboros/tools/browser.py` (vendor install) | none — vendor tree | no GC — accepted (vendor cache) | re-downloaded on next browser use |
|
||||
|
||||
|
|
|
|||
|
|
@ -270,7 +270,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). 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.
|
||||
**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 while preserving pending custody identity for start blockers, terminal audits and snapshot retention; recovery retains that invocation without POSTing. Known pending review tokens take the same missing-request path, while absent or definitely refused skill-review records keep their existing fresh-retry behavior. 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.
|
||||
|
||||
|
|
|
|||
|
|
@ -53,8 +53,9 @@ def pending_invocations(drive_root: Any,
|
|||
The launched-never-collected class one step EARLIER than ``open_runs``: a
|
||||
worker death between the accepted POST and ``record_started`` leaves only the
|
||||
``START_REQUESTED`` row. Facts come from the FIRST request row (the minting,
|
||||
same rule as ``invocation_record``); a record whose canonical body never
|
||||
landed is excluded (nothing byte-identical can be replayed). ``rows`` shares
|
||||
same rule as ``invocation_record``). Legacy rows without a body stay excluded;
|
||||
an unreadable request reference retains identity with ``request=None``.
|
||||
Reconciliation cannot replay it without that body. ``rows`` shares
|
||||
one pre-read snapshot with ``replay`` (atomic payload busy claim)."""
|
||||
from ouroboros.delegate_pending import pending_invocations as replay_pending
|
||||
|
||||
|
|
@ -233,8 +234,14 @@ def _recover_pending_invocation(drive_root: Any, gateway: Any,
|
|||
"reason": "review_panel_owns_invocation"}
|
||||
_custody().emit(drive_root, _custody().RECONCILED, result)
|
||||
return result
|
||||
body = record.get("request")
|
||||
if not isinstance(body, dict) or not body:
|
||||
result = {"invocation_id": invocation_id, "task_id": task_id,
|
||||
"action": "invocation_retained", "reason": "invocation_request_unrecorded"}
|
||||
_custody().emit(drive_root, _custody().RECONCILED, result)
|
||||
return result
|
||||
try:
|
||||
handle = gateway.start_run(dict(record["request"]), idempotency_key=invocation_id)
|
||||
handle = gateway.start_run(dict(body), idempotency_key=invocation_id)
|
||||
except ClaudexorUnavailable as exc:
|
||||
status = int(getattr(exc, "status_code", 0) or 0)
|
||||
if 400 <= status < 500:
|
||||
|
|
@ -259,7 +266,6 @@ def _recover_pending_invocation(drive_root: Any, gateway: Any,
|
|||
"action": "recovery_pending"}
|
||||
_custody().emit(drive_root, _custody().RECONCILED, result)
|
||||
return result
|
||||
body = record["request"]
|
||||
execution = body.get("execution") if isinstance(body.get("execution"), dict) else {}
|
||||
scope = body.get("scope") if isinstance(body.get("scope"), dict) else {}
|
||||
custody = _custody().RunCustody(
|
||||
|
|
|
|||
|
|
@ -73,8 +73,9 @@ def pending_invocations(
|
|||
continue
|
||||
# Resolve only survivors, not every historical start on each sweep.
|
||||
body = request_body(drive_root, record)
|
||||
record.pop("request_ref")
|
||||
if body:
|
||||
ref = record.pop("request_ref")
|
||||
# An unreadable stored body does not discharge the pending start.
|
||||
if body or ref is not None:
|
||||
record["request"] = body
|
||||
pending.append(record)
|
||||
return pending
|
||||
|
|
|
|||
|
|
@ -778,10 +778,10 @@ def run_delegated_review_session(
|
|||
run_id, started_custody = owned_started_review_custody(
|
||||
custody, custody_drive, record, task_id)
|
||||
run_request, invocation_id = record.get("request"), retry_token
|
||||
elif (record is not None and record["state"] == "pending"
|
||||
and isinstance(record.get("request"), dict) and record["request"]):
|
||||
run_request, invocation_id = record["request"], retry_token
|
||||
recovering = bool(run_id) or run_request is not None
|
||||
elif record is not None and record["state"] == "pending":
|
||||
run_request, invocation_id = record.get("request"), retry_token
|
||||
# Known pending custody survives body loss; recovery reports the missing request.
|
||||
recovering = bool(run_id or invocation_id)
|
||||
if retry_token and not recovering and surface != "skill_review":
|
||||
raise ReviewRouteUnavailable(
|
||||
"delegated retry token has no durable invocation; refusing a second paid run",
|
||||
|
|
|
|||
|
|
@ -210,9 +210,9 @@ def _validated_invocation(drive: Any, retry_token: str, task_id: str,
|
|||
if not isinstance(body, dict) or not body:
|
||||
return None, _fail("delegate_start", "invocation_request_unrecorded",
|
||||
"That invocation's durable row carries no canonical request "
|
||||
"body, so it cannot be replayed byte-identically. Start a "
|
||||
"new run with a plain "
|
||||
"delegate_start(subagent_id=..., prompt=...).",
|
||||
"body, so it cannot be replayed byte-identically. Its outcome "
|
||||
"remains unknown; restore the recorded request or reconcile "
|
||||
"the original invocation before starting a replacement.",
|
||||
retry_of=retry_token)
|
||||
if str(body.get("prompt") or "") != text:
|
||||
return None, _fail("delegate_start", "retry_prompt_mismatch",
|
||||
|
|
|
|||
|
|
@ -162,12 +162,129 @@ def test_unreadable_ref_keeps_request_unknown_and_retry_refuses(tmp_path, damage
|
|||
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) == []
|
||||
pending = custody.pending_invocations(tmp_path)
|
||||
assert [row["invocation_id"] for row in pending] == ["lost"]
|
||||
assert pending[0]["request"] is None and "request_ref" not in pending[0]
|
||||
record, refusal = _validated_invocation(tmp_path, "lost", "task", "exact original")
|
||||
assert record is None
|
||||
assert json.loads(refusal.text)["reason"] == "invocation_request_unrecorded"
|
||||
|
||||
|
||||
@pytest.mark.parametrize("inline_request", [None, {}, [], "invalid"])
|
||||
def test_malformed_legacy_inline_without_a_reference_keeps_existing_behavior(tmp_path, inline_request):
|
||||
assert custody.emit(tmp_path, custody.START_REQUESTED, {
|
||||
"invocation_id": "legacy-empty", "task_id": "owner", "request": inline_request})
|
||||
assert custody.pending_invocations(tmp_path) == []
|
||||
|
||||
|
||||
@pytest.mark.parametrize("damage", [None, "missing", "corrupt"])
|
||||
def test_lost_ack_keeps_the_original_start_claim_when_its_body_is_unreadable(tmp_path, monkeypatch, damage):
|
||||
from ouroboros.delegate_recovery import unsettled_start_ids
|
||||
from ouroboros.delegate_terminal import _audit_task_custody
|
||||
from ouroboros.gateways import claudexor as gateway
|
||||
from ouroboros.tools import delegate
|
||||
|
||||
keys = []
|
||||
|
||||
class LostAcknowledgment(_LiveRunStub):
|
||||
def start_run(self, body, *, idempotency_key=""):
|
||||
keys.append(idempotency_key)
|
||||
if len(keys) == 1:
|
||||
raise gateway.ClaudexorUnavailable("daemon_unreachable", "fixture accepted start then lost reply")
|
||||
return {"runId": "original-run"}
|
||||
|
||||
stub = LostAcknowledgment()
|
||||
monkeypatch.setenv("OUROBOROS_SUBAGENT_HARNESS", "some-route=weak-model:low")
|
||||
monkeypatch.setattr(gateway, "ClaudexorGateway", lambda *args, **kwargs: stub)
|
||||
ctx = _nanny_ctx(tmp_path, task_id="owner")
|
||||
first = json.loads(delegate._delegate_start(ctx, "Inspect fixture files.").text)
|
||||
token = first["pending_invocation_id"]
|
||||
assert keys == [token]
|
||||
blob = Path(_start_rows(tmp_path)[0]["request_ref"]["path"])
|
||||
if damage == "missing":
|
||||
blob.unlink()
|
||||
elif damage == "corrupt":
|
||||
blob.write_bytes(b"fixture corrupt gzip")
|
||||
assert unsettled_start_ids(tmp_path, "owner")["pending_invocation_ids"] == [token]
|
||||
audit = {"task_id": "owner", "audit_status": "ok", "unreconciled": []}
|
||||
_audit_task_custody(tmp_path, "owner", audit, emit_evidence=False)
|
||||
assert audit["pending_invocation_ids"] == [token]
|
||||
assert audit["unreconciled"] == [f"invocation:{token}"]
|
||||
fresh = json.loads(delegate._delegate_start(ctx, "Inspect fixture files.").text)
|
||||
assert fresh["reason"] == "replacement_requires_settlement" and keys == [token]
|
||||
retry = json.loads(delegate._delegate_start(ctx, "Inspect fixture files.", retry_of=token).text)
|
||||
if damage:
|
||||
assert retry["reason"] == "invocation_request_unrecorded" and keys == [token]
|
||||
assert "Start a new run" not in retry["detail"]
|
||||
else:
|
||||
assert retry["status"] == "started" and keys == [token, token]
|
||||
|
||||
|
||||
@pytest.mark.parametrize("damage", ["missing", "corrupt"])
|
||||
def test_missing_pending_body_keeps_snapshot_and_orphan_recovery_without_a_post(tmp_path, damage):
|
||||
from types import SimpleNamespace
|
||||
from ouroboros.delegate_custody_reconcile import _recover_pending_invocation
|
||||
|
||||
assert custody.record_start_requested(
|
||||
tmp_path, invocation_id="pending", task_id="owner", snapshot_id="snapshot",
|
||||
request={"prompt": "Preserve this pending invocation"})
|
||||
blob = Path(_start_rows(tmp_path)[0]["request_ref"]["path"])
|
||||
if damage == "missing":
|
||||
blob.unlink()
|
||||
else:
|
||||
blob.write_bytes(b"fixture corrupt gzip")
|
||||
pending = custody.pending_invocations(tmp_path)
|
||||
assert len(pending) == 1 and pending[0]["request"] is None
|
||||
assert custody.open_snapshot_ids(tmp_path) == {"snapshot"}
|
||||
gateway = SimpleNamespace(start_run=lambda *args, **kwargs: pytest.fail("unknown body must not POST"))
|
||||
result = _recover_pending_invocation(tmp_path, gateway, pending[0])
|
||||
assert result["action"] == "invocation_retained"
|
||||
assert result["reason"] == "invocation_request_unrecorded"
|
||||
assert custody.invocation_record(tmp_path, "pending")["state"] == "pending"
|
||||
assert custody.open_snapshot_ids(tmp_path) == {"snapshot"}
|
||||
review_result = _recover_pending_invocation(tmp_path, gateway, {**pending[0], "source": "review_substrate"})
|
||||
assert review_result["reason"] == "review_panel_owns_invocation"
|
||||
assert custody.emit(tmp_path, custody.START_FAILED, {"invocation_id": "pending", "definite": True})
|
||||
assert custody.pending_invocations(tmp_path) == [] and custody.open_snapshot_ids(tmp_path) == set()
|
||||
|
||||
|
||||
@pytest.mark.parametrize("prior", ["absent", "refused", "intact", "missing", "corrupt"])
|
||||
def test_skill_review_retries_distinguish_known_pending_from_absent_or_refused(tmp_path, fake_route, prior):
|
||||
from ouroboros.gateways.claudexor import ClaudexorUnavailable
|
||||
from ouroboros.review_execution import ReviewRouteUnavailable
|
||||
from tests._review_session_route_shared import _run_session_directly
|
||||
|
||||
state = {"pending_invocation_id": "absent-token"} if prior == "absent" else {}
|
||||
if prior != "absent":
|
||||
fake_route.start_error = ClaudexorUnavailable(
|
||||
"fixture_refused" if prior == "refused" else "daemon_unreachable",
|
||||
"fixture initial transport outcome", status_code=400 if prior == "refused" else 0)
|
||||
with pytest.raises(ClaudexorUnavailable):
|
||||
_run_session_directly(tmp_path, surface="skill_review", slot_id="skill-slot", retry_state=state)
|
||||
row = _start_rows(tmp_path)[0]
|
||||
token = row["invocation_id"]
|
||||
state["pending_invocation_id"] = token
|
||||
if prior == "missing":
|
||||
Path(row["request_ref"]["path"]).unlink()
|
||||
elif prior == "corrupt":
|
||||
Path(row["request_ref"]["path"]).write_bytes(b"fixture corrupt gzip")
|
||||
else:
|
||||
token = state["pending_invocation_id"]
|
||||
if prior in {"missing", "corrupt"}:
|
||||
with pytest.raises(ReviewRouteUnavailable) as caught:
|
||||
_run_session_directly(tmp_path, surface="skill_review", slot_id="skill-slot", retry_state=state)
|
||||
assert caught.value.code == "review_recovery_request_missing"
|
||||
assert state["pending_invocation_id"] == token
|
||||
assert [key for gateway in fake_route.instances for key in gateway.start_keys] == [token]
|
||||
else:
|
||||
_run_session_directly(tmp_path, surface="skill_review", slot_id="skill-slot", retry_state=state)
|
||||
keys = [key for gateway in fake_route.instances for key in gateway.start_keys]
|
||||
if prior == "intact":
|
||||
assert keys == [token, token]
|
||||
else:
|
||||
assert keys[-1] != token and len(keys) == (1 if prior == "absent" else 2)
|
||||
|
||||
|
||||
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,
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue