mirror of
https://github.com/razzant/ouroboros.git
synced 2026-10-02 19:58:46 +00:00
321 lines
16 KiB
Python
321 lines
16 KiB
Python
"""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
|
|
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,
|
|
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"]
|