diff --git a/tests/system_e2e/harness.py b/tests/system_e2e/harness.py index 2c3628432..6fd9d2c67 100644 --- a/tests/system_e2e/harness.py +++ b/tests/system_e2e/harness.py @@ -135,9 +135,18 @@ SCENARIOS = { # FOUR tests: executor starvation alone, a held Host authentication alone, # both together, and the owner's Panic under the combined load. "S33": ("Presence waits keep owner controls answering: 12 events (2 held at the model, slot and same-conversation waits) on a 12-thread default executor and/or one held Host authentication; while held, /api/state, a v1 receipt and the owner's Stop (durable cancel intent, cancelled terminal) answer inside 10s/15s windows (health alone never passes); a disconnected turn and its retry are ONE model call and a later replay answers the same projection; Panic under the combined load ends the whole tree", LANE_MOCK), + # Plan review's ANSWER CHANNEL under blocking enforcement at the shipped cycle cap, + # over the asynchronous barrier route with three DISTINCT keyless reviewer models + # (the stub answers per seat by the wire model id). + "S34": ("plan review addressed answer, BLOCKING at the shipped cap: t1 objects below quorum -> $0 reject -> the identical envelope with the answer re-asks t1 ALONE (t2/t3 kept at $0 as replayed rows) over the barrier route -> t1 retires -> GREEN closed, two paid cycles, the task completes under blocking", LANE_MOCK), + "S35": ("plan review no-need path, BLOCKING: t1 asks the author (need_evidence), t2 leaves a note; a $0 accept closes the wave GREEN with no second panel (three reviewer calls, one paid cycle) and the task completes under blocking", LANE_MOCK), } MOCK_SLUG = "openai-compatible::mock-model" +# The three DISTINCT reviewer slugs of the per-seat scenarios: seat t rides +# ``-t``, so ``default_slot_binder`` (the wire ``model`` field) names +# the seat that made the call and a ReviewScript step can answer per seat. +DISTINCT_MOCK_MODEL_IDS = tuple(f"mock-model-t{i}" for i in (1, 2, 3)) # --------------------------------------------------------------------------- # Prompt markers the stub classifies review-organ calls by (roast F22). @@ -638,13 +647,21 @@ class ScriptedStubModel(LoopbackModelServer): def __init__(self, script=None, *, final_answer: str = "Final answer: scripted scenario complete.", latency_sec: float = 0.0, review_script: "ReviewScript | None" = None, - gate: "ModelGate | None" = None) -> None: + gate: "ModelGate | None" = None, model_ids=None) -> None: super().__init__(latency_sec=latency_sec, gate=gate) self.script = list(script or []) self.final_answer = final_answer self.review_script = review_script + # As on ReplayModel: extra wire ids to advertise on /models, so the + # capability-evidence window probe confirms a window for each distinct slot route. + self._explicit_model_ids = [str(m) for m in model_ids] if model_ids else None self._script_index = 0 + def _model_ids(self) -> list[str]: + if self._explicit_model_ids is not None: + return sorted(set(self._explicit_model_ids) | {"mock-model"}) + return super()._model_ids() + def _next_step(self, _body) -> dict | None: if self._script_index >= len(self.script): return None @@ -1087,7 +1104,7 @@ class KeylessIsolatedServer(IsolatedServer): self.candidate.release() -def keyless_reviewer_slots(*, advisory: bool = False) -> str: +def keyless_reviewer_slots(*, advisory: bool = False, distinct_models: bool = False) -> str: """The structured ``OUROBOROS_REVIEWER_SLOTS`` value pinning every reviewer row to the loopback stub. @@ -1102,10 +1119,15 @@ def keyless_reviewer_slots(*, advisory: bool = False) -> str: the stub (wave 3a): the advisory pre-review then runs the bounded NATIVE inspection episode against the loopback model instead of being unavailable keyless (which the commit gate compensates with an audited bypass). + + ``distinct_models=True`` pins seat ``t`` to its own slug + (``DISTINCT_MOCK_MODEL_IDS``) so a per-seat ReviewScript can tell the seats + apart on the wire; the stub must advertise those ids (``model_ids``). """ row = {"kind": "api_chat", "target_id": MOCK_SLUG} payload = { - "triad": [{"slot_id": f"t{i}", "route": dict(row)} for i in (1, 2, 3)], + "triad": [{"slot_id": f"t{i}", "route": {**row, **({"target_id": f"{MOCK_SLUG}-t{i}"} if distinct_models else {})}} + for i in (1, 2, 3)], "scope": [{"slot_id": "s1", "route": dict(row)}], } if advisory: diff --git a/tests/system_e2e/test_plan_review_addressed_answer.py b/tests/system_e2e/test_plan_review_addressed_answer.py new file mode 100644 index 000000000..53fd04e09 --- /dev/null +++ b/tests/system_e2e/test_plan_review_addressed_answer.py @@ -0,0 +1,339 @@ +"""S34-S35 — plan review's ANSWER CHANNEL on a real isolated server, keyless. + +Both scenarios run under BLOCKING enforcement at the shipped cycle cap (2), over the +asynchronous barrier route (a fresh dispatch returns at the dispatch barrier with an +open custody-pending wave; the settled-wave mailbox frame is the cue; the identical +envelope resubmitted after it is the $0 collection), with three DISTINCT keyless +reviewer models so the stub answers PER SEAT by the wire model id. + +* S34 — the addressed answer: cycle 1 ends REVIEW_REQUIRED (t1 blocking below + quorum, t2/t3 clean); the author records a $0 reject; the IDENTICAL envelope sent + with that answer asks t1 ALONE (one paid cycle, same fingerprint) while t2 and t3 + keep their recorded answers at $0 as replayed rows; t1 retires its finding; the + collection closes the wave GREEN; two cycles paid; the task completes under + blocking. The durable chain is exact: the barrier snapshot of cycle 2 names the + cycle-1 wave it answered, and t1's second request carries the prior cycle's + rationale. +* S35 — the no-need path: t1 asks the author (need_evidence) and t2 leaves a note; + a $0 accept of the question closes the wave GREEN with no second panel (three + reviewer calls in total, one paid cycle) and the task completes under blocking. + +Every assertion reads durable artifacts (the stored task row, the immutable wave +artifacts, the tool log) and the stub's own call ledger — never an HTTP 200 alone. +""" + +from __future__ import annotations + +import json +import re + +import pytest + +from ouroboros.tools.review_synthesis import PLAN_REVIEW_CONTROL_PREFIX +from tests.system_e2e.harness import ( + DISTINCT_MOCK_MODEL_IDS, + LANE_MOCK, + ArtifactOracle, + ReviewScript, + body_text, + default_slot_binder, + keyless_reviewer_slots, + keyless_settings, + require_lane, + start_server, + submit_running, + wait_durable_result, +) +from tests.system_e2e.test_system_scenarios_w3a import ( + _ROOT_TASK_ID_RE, + _WAIT_ROUNDS_MAX, + _WAIT_WINDOW_SEC, + _Again, + _HoldingStubModel, + _tool_rows, +) + +CLEAN = "[]\nNO_FINDINGS" +GOAL = "Write the addressed-answer smoke note." +PLAN = "Draft the note with its marker line, verify the marker, then finish." +SPEC = { + "in_scope": ["addressed-answer smoke note"], + "acceptance_claims": ["The note carries the S34_MARKER line."], + # Required on every submitted spec; the smoke note changes no repository file. + "affected_paths": [], +} +T1_OBJECTION = "The spec has no invariant pinning the exact bytes of the S34_MARKER line." +REJECT_RATIONALE = "claim_1 already pins the exact marker line; an invariant would duplicate it." +T1_BLOCKING = json.dumps([{ + "id": "f1", "class": "blocking", "breaks": "claim_1", + "summary": T1_OBJECTION, + "recommendation": "Add an invariant naming the exact marker bytes.", +}]) +T1_QUESTION = json.dumps([{ + "id": "q1", "class": "need_evidence", "breaks": "claim_1", + "summary": "Which exact marker line will the note carry?", + "recommendation": "", +}]) +T2_NOTE = json.dumps([{ + "id": "n1", "class": "note", "breaks": "", + "summary": "Keep the note to one paragraph.", + "recommendation": "", +}]) + +# The settled-wave frame (plan_review_collect.announce_released_settlement) carries +# the fingerprint AND the counts; the addressed cycle re-uses the fingerprint, so the +# frames are told apart by their full text, never by the fingerprint alone. +_FRAME_RE = re.compile(r"Plan review wave ([0-9a-f?]+): (\d+) of (\d+) reviewer slot\(s\) settled") +_FINGERPRINT_RE = re.compile(r"\*\*Plan fingerprint:\*\* `([0-9a-f]{64})`") + + +def _seat(body: dict) -> str: + """The seat a plan-review call came from: the ``-t`` tail of the wire model id.""" + return default_slot_binder(body).rsplit("-", 1)[-1] + + +def _per_seat(answers: dict): + """A ReviewScript step that answers BY SEAT. An unexpected seat gets a loud + unparseable text, so the durable wave (not a hang) names the defect.""" + def step(body: dict) -> str: + return answers.get(_seat(body), f"E2E_SCRIPT_ERROR: unexpected reviewer seat {default_slot_binder(body)!r}") + return step + + +def _envelope(**extra) -> dict: + """The ONE envelope of the scenario; ``extra`` rides beside it (review_disposition).""" + return {"tool": "plan_task", "arguments": {"goal": GOAL, "plan": PLAN, "spec": SPEC, **extra}} + + +def _after_frames(count: int, then): + """Hold the script until ``count`` DISTINCT settled-wave frames are visible in the + transcript, then emit ``then`` (a step or a callable step). Until then the step + waits on this task's own id (``wait_task`` returns early on + ``owner_mailbox_pending``), the route the plan-review contract text names.""" + waits = {"rounds": 0} + + def step(body: dict) -> dict: + text = body_text(body) + frames = list(dict.fromkeys(_FRAME_RE.findall(text))) + if len(frames) >= count: + return then(body) if callable(then) else then + waits["rounds"] += 1 + task_ids = _ROOT_TASK_ID_RE.findall(text) + if not task_ids or waits["rounds"] > _WAIT_ROUNDS_MAX: + return {"final": ( + f"E2E_SCRIPT_ERROR: fewer than {count} settled-wave frame(s) after " + f"{waits['rounds']} round(s); frames seen: {frames}; own task id visible: {bool(task_ids)}")} + return _Again({"tool": "wait_task", "arguments": { + "task_id": task_ids[-1], "timeout_sec": _WAIT_WINDOW_SEC}}) + + return step + + +def _answer_step(items: list, *, with_envelope: bool): + """Answer the recorded wave named by the LAST ``Plan fingerprint:`` line the host + printed: alone (a $0 disposition) or beside the identical envelope (the addressed + re-ask).""" + def step(body: dict) -> dict: + found = _FINGERPRINT_RE.findall(body_text(body)) + if not found: + return {"final": "E2E_SCRIPT_ERROR: no plan fingerprint visible in the transcript"} + disposition = {"review_fingerprint": found[-1], "items": list(items)} + if with_envelope: + return _envelope(review_disposition=disposition) + return {"tool": "plan_task", "arguments": {"review_disposition": disposition}} + + return step + + +def _control(text: str) -> dict: + lines = [line for line in str(text).splitlines() if line.startswith(PLAN_REVIEW_CONTROL_PREFIX)] + assert len(lines) == 1, text[-600:] + return json.loads(lines[0][len(PLAN_REVIEW_CONTROL_PREFIX):]) + + +def _settings(stub) -> dict: + return keyless_settings( + stub, OUROBOROS_RUNTIME_MODE="advanced", OUROBOROS_REVIEW_ENFORCEMENT="blocking", + OUROBOROS_REVIEWER_SLOTS=keyless_reviewer_slots(distinct_models=True), + ) + + +def _wave_artifacts(oracle: ArtifactOracle) -> list: + paths = sorted((oracle.data_root / "task_results" / "artifacts").rglob("plan-review-wave-*.json")) + return [json.loads(path.read_text(encoding="utf-8")) for path in paths] + + +def _plan_results(oracle: ArtifactOracle, task_id: str) -> list: + """The FULL text of every ``plan_task`` result of the task, in call order, read + through each tool row's persisted call trace (the direct tools.jsonl row keeps a + bounded preview; the exact result is the digest-verified observability blob).""" + from ouroboros.observability import read_call_payload + + task_drive = oracle.task_drive(task_id) + out = [] + for row in _tool_rows(task_drive, "plan_task"): + call_id = str((row.get("result_ref") or {}).get("call_id") or "") + assert call_id, row + _manifest, payload, _ref = read_call_payload(task_drive.data_root, task_id=task_id, call_id=call_id) + out.append(str(payload.get("result") or "")) + return out + + +# =========================================================================== +# S34 — the addressed answer under blocking, barrier route +# =========================================================================== + +@pytest.mark.integration +@pytest.mark.serial +def test_s34_addressed_answer_reasks_the_objector_alone_and_closes_green(e2e_clone, tmp_path_factory): + require_lane(LANE_MOCK) + from ouroboros.tools.plan_review_artifacts import read_wave + + root = tmp_path_factory.mktemp("s34") + reject = {"finding_id": "t1:f1", "decision": "reject", "rationale": REJECT_RATIONALE} + review_script = ReviewScript({"plan_review": [ + *([_per_seat({"t1": T1_BLOCKING, "t2": CLEAN, "t3": CLEAN})] * 3), # cycle 1: three seats + _per_seat({"t1": CLEAN}), # cycle 2: t1 alone + ]}) + stub = _HoldingStubModel( + [_envelope(), # dispatch: the barrier + _after_frames(1, _envelope()), # $0 collect: REVIEW_REQUIRED + _answer_step([reject], with_envelope=False), # $0 reject, wave stays open + _answer_step([reject], with_envelope=True), # addressed re-ask: t1 alone + _after_frames(2, _envelope())], # $0 collect: GREEN closed + review_script=review_script, model_ids=DISTINCT_MOCK_MODEL_IDS, + ) + with stub: + server = start_server(e2e_clone, root, _settings(stub)) + try: + task_id = submit_running( + server, "Plan the note through plan_task, answer the reviewer, then finish.") + result = server.wait_task(task_id, timeout=600) + assert result.get("status") == "completed", result + oracle = ArtifactOracle(server.data_root) + stored = wait_durable_result(oracle, task_id) + assert stored.get("status") == "completed", stored.get("status") + + # The durable chronicle: two PAID cycles on ONE fingerprint, GREEN and closed. + state = stored.get("plan_review_state") + assert isinstance(state, dict), sorted(stored) + assert int(state.get("cycles_paid") or 0) == 2, state + waves = [w for w in (state.get("waves") or []) if isinstance(w, dict)] + last = waves[-1] + fingerprint = str(last.get("request_fingerprint") or "") + assert len(fingerprint) == 64 and all(w.get("request_fingerprint") == fingerprint for w in waves), waves + assert last.get("aggregate") == "GREEN" and last.get("closed") is True, last + assert int(last.get("cycle_index") or 0) == 2 and last.get("paid") is True, last + assert last.get("findings") == [] and last.get("dispositions") == [], last + actors = {str(a.get("slot_id")): a for a in last.get("actors") or []} + assert sorted(actors) == ["t1", "t2", "t3"], actors + for sid in ("t2", "t3"): # kept at $0: the replayed cycle-1 rows, never a send + kept = actors[sid] + assert kept.get("operation_state") == "not_dispatched" and kept.get("cost") == 0.0, kept + assert kept.get("replayed_from", {}).get("cycle_index") == 1, kept + assert kept["replayed_from"].get("request_fingerprint") == fingerprint, kept + assert "replayed_from" not in actors["t1"] and actors["t1"].get("ok") is True, actors["t1"] + assert str(actors["t1"].get("model") or "").endswith("mock-model-t1"), actors["t1"] + + # The stub saw 3 + 1 plan-review calls: three distinct seats, then t1 alone, + # whose second request carries the prior cycle's rationale. + plan_bodies = [body for kind, body in stub.calls if kind == "plan_review"] + assert len(plan_bodies) == 4, stub.kinds() + assert sorted(default_slot_binder(b) for b in plan_bodies[:3]) == sorted(DISTINCT_MOCK_MODEL_IDS) + assert default_slot_binder(plan_bodies[3]) == "mock-model-t1", plan_bodies[3].get("model") + second_request = body_text(plan_bodies[3]) + assert "PRIOR CYCLES" in second_request and REJECT_RATIONALE in second_request, second_request[-3000:] + assert T1_OBJECTION in second_request, "t1's own cycle-1 answer continues its transcript" + + # The exact artifact chain: the cycle-2 BARRIER snapshot names the seats it asked + # again and the cycle-1 wave it answered; that reference reads back as the + # answered wave (cycle 1, same fingerprint, the reject recorded); the collected + # cycle-2 wave carries the same predecessor reference and is closed GREEN. + barrier = [p for p in _wave_artifacts(oracle) + if p.get("custody_pending") and int(p.get("cycle_index") or 0) == 2] + assert len(barrier) == 1, [(p.get("cycle_index"), p.get("custody_pending")) for p in _wave_artifacts(oracle)] + addressed = barrier[0].get("addressed") + assert addressed and addressed["slots"] == ["t1"] and addressed["finding_ids"] == ["t1:f1"], addressed + assert addressed["kept"] == ["t2", "t3"], addressed + answered = read_wave(oracle.data_root, task_id, addressed["wave_artifact"]) + assert int(answered.get("cycle_index") or 0) == 1 and answered.get("request_fingerprint") == fingerprint + assert answered.get("closed") is False and answered.get("aggregate") == "REVIEW_REQUIRED", answered.get("aggregate") + assert [(d["finding_id"], d["decision"]) for d in answered.get("dispositions") or []] == [("t1:f1", "reject")] + assert last.get("previous_wave_artifact") == addressed["wave_artifact"], last.get("previous_wave_artifact") + collected = read_wave(oracle.data_root, task_id, last["wave_artifact"]) + assert collected.get("closed") is True and collected.get("aggregate") == "GREEN" + assert int(collected.get("cycle_index") or 0) == 2 + + # FIVE plan_task calls bought exactly one re-ask: dispatch, collect, the $0 + # reject, the addressed re-ask (barrier), collect. + results = _plan_results(oracle, task_id) + assert len(results) == 5, [text[:120] for text in results] + assert [(_control(t)["outcome"], _control(t)["closed"]) for t in results] == [ + ("DEGRADED", False), ("REVIEW_REQUIRED", False), ("REVIEW_REQUIRED", False), + ("DEGRADED", False), ("GREEN", True)], results + assert "REVIEW CUSTODY PENDING" in results[0] and "REVIEW CUSTODY PENDING" in results[3] + assert "**Addressed answer:** asked again t1 on t1:f1; kept at $0: t2, t3." in results[3], results[3] + assert "kept its cycle-1 answer at $0" in results[4], results[4] + review_script.assert_consumed() + assert stub.script_consumed() + finally: + server.stop() + + +# =========================================================================== +# S35 — the no-need path under blocking +# =========================================================================== + +@pytest.mark.integration +@pytest.mark.serial +def test_s35_no_need_path_closes_at_zero_cost_under_blocking(e2e_clone, tmp_path_factory): + require_lane(LANE_MOCK) + root = tmp_path_factory.mktemp("s35") + accept = {"finding_id": "t1:q1", "decision": "accept", + "rationale": "The note carries the literal line S34_MARKER as its first line."} + review_script = ReviewScript({"plan_review": [ + _per_seat({"t1": T1_QUESTION, "t2": T2_NOTE, "t3": CLEAN})] * 3}) + stub = _HoldingStubModel( + [_envelope(), # dispatch: the barrier + _after_frames(1, _envelope()), # $0 collect: REVIEW_REQUIRED + _answer_step([accept], with_envelope=False)], # $0 accept: GREEN closed + review_script=review_script, model_ids=DISTINCT_MOCK_MODEL_IDS, + ) + with stub: + server = start_server(e2e_clone, root, _settings(stub)) + try: + task_id = submit_running( + server, "Plan the note through plan_task, answer the reviewer's question, then finish.") + result = server.wait_task(task_id, timeout=600) + assert result.get("status") == "completed", result + oracle = ArtifactOracle(server.data_root) + stored = wait_durable_result(oracle, task_id) + assert stored.get("status") == "completed", stored.get("status") + + state = stored.get("plan_review_state") + assert isinstance(state, dict), sorted(stored) + assert int(state.get("cycles_paid") or 0) == 1, state + waves = [w for w in (state.get("waves") or []) if isinstance(w, dict)] + last = waves[-1] + assert last.get("aggregate") == "GREEN" and last.get("closed") is True, last + assert int(last.get("cycle_index") or 0) == 1 and "addressed" not in last, last + assert [(f["finding_id"], f["class"]) for f in last.get("findings") or []] == [ + ("t1:q1", "need_evidence"), ("t2:n1", "note")], last.get("findings") + assert [(d["finding_id"], d["decision"]) for d in last.get("dispositions") or []] == [("t1:q1", "accept")] + assert state.get("need_evidence_seen") == [], state.get("need_evidence_seen") # no locator asked + actors = {str(a.get("slot_id")): a for a in last.get("actors") or []} + assert sorted(actors) == ["t1", "t2", "t3"] and not any("replayed_from" in a for a in actors.values()) + + # Exactly one panel of three distinct seats; the accept sent nothing. + plan_bodies = [body for kind, body in stub.calls if kind == "plan_review"] + assert len(plan_bodies) == 3, stub.kinds() + assert sorted(default_slot_binder(b) for b in plan_bodies) == sorted(DISTINCT_MOCK_MODEL_IDS) + results = _plan_results(oracle, task_id) + assert len(results) == 3, [text[:120] for text in results] + assert [(_control(t)["outcome"], _control(t)["closed"]) for t in results] == [ + ("DEGRADED", False), ("REVIEW_REQUIRED", False), ("GREEN", True)], results + assert "(cached exact review — no reviewer was called)" not in results[1], results[1] + review_script.assert_consumed() + assert stub.script_consumed() + finally: + server.stop() diff --git a/tests/test_plan_review_answer_route.py b/tests/test_plan_review_answer_route.py new file mode 100644 index 000000000..3d35b9c8c --- /dev/null +++ b/tests/test_plan_review_answer_route.py @@ -0,0 +1,252 @@ +"""The ROUTE of the addressed answer, end to end and in process: the real engine, the real +``task_results``/artifact store on ``tmp_path``, the fake review substrate of +``tests.test_plan_review_engine`` and the production finalization gate +(``owner_hurry.force_plan_decision``), all under BLOCKING enforcement at the shipped +cycle cap of 2 — the configuration every install starts from. + +Five routes: the objector retires after the addressed re-ask (GREEN, closed, the gate +releases, the closed wave's claims bind acceptance); two objectors at quorum retire; +the objector never retires and the cap is spent (the typed cap state, the honest exits); +the objector fails to answer the re-ask (its finding is carried, never GREEN); and the +no-need path (a question and a note close at $0 with no second transport call). +""" +from __future__ import annotations + +import json + +from ouroboros.contracts.task_contract import effective_acceptance_claims +from ouroboros.outcomes import derive_loop_outcome +from ouroboros.owner_hurry import force_plan_decision +from ouroboros.review_cycles import review_max_cycles +from ouroboros.task_results import closed_plan_review_wave +from ouroboros.tools.plan_review_artifacts import read_wave +from tests.test_plan_review_answer_channel import _actors, _answer, _blocking, _ids, _item, _question, _reject +from tests.test_plan_review_engine import ( # noqa: F401 + CLEAN, DECK_SPEC, _call, _control, _finding, _state, harness, +) + +REJECT_RATIONALE = "the budget line is already approved" + + +def _gate(ctx) -> dict: + return force_plan_decision(ctx, {}, enforcement="blocking") + + +def _wave(h) -> dict: + return _state(h)["waves"][-1] + + +def _events(h, event_type: str) -> list: + return [e["data"] for e in list(h.events.queue) + if e.get("type") == "log_event" and e.get("data", {}).get("type") == event_type] + + +def _addressed(ctx, fp: str, *fids: str) -> str: + """The identical envelope carrying the reject of every named finding.""" + return _call(ctx, review_disposition={"review_fingerprint": fp, "items": [ + _item(fid, "reject", REJECT_RATIONALE) for fid in fids]}) + + +def _bound_claims(state: dict) -> tuple[list, str]: + return effective_acceptance_claims({}, closed_plan_review_wave(state)) + + +def _objection(h, answers: dict) -> tuple: + """Cycle 1 at the shipped cap: the configured seats answer ``answers``; the wave is + open, one cycle is paid, and the blocking gate holds.""" + assert review_max_cycles() == 2, "the route runs at the shipped default cap" + sub = h.install(answers) + ctx = h.make_ctx() + out = _call(ctx) + assert not _control(out)["closed"], out + assert len(sub.calls) == 1 and _state(h)["cycles_paid"] == 1 + gate = _gate(ctx) + assert gate["allow"] is False and gate["status"] == "open", gate + assert _bound_claims(_state(h)) == ([], "") # no closed authority: nothing binds acceptance + return ctx, _wave(h)["request_fingerprint"], sub + + +# ----------------------------------------------------------------- the happy route + +def test_e2e_reject_addressed_retire_green_under_blocking(harness): # noqa: F811 + """REVIEW_REQUIRED (s1 blocking below quorum) → the gate holds → a $0 reject → the + gate still holds → the identical envelope with the reject asks s1 alone (s2, s3 kept + at $0) → s1 retires → GREEN closed, two cycles paid, the gate releases and the closed + wave's claims bind acceptance.""" + ctx, fp, sub1 = _objection(harness, {"s1": _blocking(), "s2": CLEAN, "s3": CLEAN}) + assert _control(_call(ctx)) == {"outcome": "REVIEW_REQUIRED", "closed": False} # replays free + assert len(sub1.calls) == 1 + rejected = _reject(ctx, fp) + assert _control(rejected) == {"outcome": "REVIEW_REQUIRED", "closed": False} + assert len(sub1.calls) == 1 and _state(harness)["cycles_paid"] == 1 # the answer costs nothing + assert _ids(_wave(harness)) == [("s1:f1", "reject")] + assert _gate(ctx)["allow"] is False, "a reject alone never closes a blocking finding under blocking" + + sub2 = harness.install({"s1": CLEAN, "s2": CLEAN, "s3": CLEAN}) + result = _addressed(ctx, fp, "s1:f1") + assert _control(result) == {"outcome": "GREEN", "closed": True} + assert len(sub2.calls) == 1 and [s.slot_id for s in sub2.calls[0]["slots"]] == ["s1"] + assert "**Addressed answer:** asked again s1 on s1:f1; kept at $0: s2, s3." in result + + state = _state(harness) + wave = state["waves"][-1] + assert state["cycles_paid"] == 2 + assert wave["request_fingerprint"] == fp and wave["cycle_index"] == 2 and wave["paid"] + assert wave["aggregate"] == "GREEN" and wave["closed"] is True and wave["findings"] == [] + assert wave["addressed"] == {"slots": ["s1"], "finding_ids": ["s1:f1"], "kept": ["s2", "s3"], + "wave_artifact": wave["previous_wave_artifact"]} + actors = _actors(harness) + for sid in ("s2", "s3"): + assert actors[sid]["operation_state"] == "not_dispatched" and actors[sid]["cost"] == 0.0 + assert actors[sid]["replayed_from"] == {"request_fingerprint": fp, "cycle_index": 1, + "operation_id": f"op-{sid}"} + assert "replayed_from" not in actors["s1"] and actors["s1"]["operation_state"] == "settled" + # The exact chain: the addressed wave names the cycle-1 wave it answered, byte-exact. + prior = read_wave(harness.drive, ctx.task_id, wave["addressed"]["wave_artifact"]) + assert prior["cycle_index"] == 1 and prior["request_fingerprint"] == fp + assert _ids(prior) == [("s1:f1", "reject")] and not prior["closed"] + + gate = _gate(ctx) + assert gate["allow"] is True and gate["status"] == "closed" and gate["closed"] is True + claims, source = _bound_claims(state) + assert source == "plan_review" + assert [c["claim"] for c in claims] == DECK_SPEC["acceptance_claims"] and claims[0]["id"] == "claim_1" + + +def test_e2e_revise_plan_two_objectors_retire(harness): # noqa: F811 + """Two objectors at quorum → REVISE_PLAN, which no disposition closes; both rejects + ride the identical envelope, both seats are re-asked (s3 kept), both retire → GREEN.""" + ctx, fp, sub1 = _objection(harness, {"s1": _blocking(), "s2": _blocking(), "s3": CLEAN}) + assert _wave(harness)["aggregate"] == "REVISE_PLAN" + rejected = _answer(ctx, fp, _item("s1:f1", "reject", REJECT_RATIONALE), _item("s2:f1", "reject", REJECT_RATIONALE)) + assert _control(rejected) == {"outcome": "REVISE_PLAN", "closed": False} + assert len(sub1.calls) == 1 and _gate(ctx)["allow"] is False + + sub2 = harness.install({"s1": CLEAN, "s2": CLEAN, "s3": CLEAN}) + result = _addressed(ctx, fp, "s1:f1", "s2:f1") + assert _control(result) == {"outcome": "GREEN", "closed": True} + assert len(sub2.calls) == 1 and [s.slot_id for s in sub2.calls[0]["slots"]] == ["s1", "s2"] + state = _state(harness) + wave = state["waves"][-1] + assert state["cycles_paid"] == 2 and wave["closed"] and wave["aggregate"] == "GREEN" + assert wave["addressed"]["slots"] == ["s1", "s2"] and wave["addressed"]["kept"] == ["s3"] + actors = _actors(harness) + assert actors["s3"]["replayed_from"]["cycle_index"] == 1 and actors["s3"]["cost"] == 0.0 + assert all("replayed_from" not in actors[sid] for sid in ("s1", "s2")) + assert _gate(ctx)["allow"] is True and _bound_claims(state)[1] == "plan_review" + + +# ---------------------------------------------------------- the honest failure routes + +def test_e2e_unresolved_objector_spends_the_cap_and_exits_honestly(harness): # noqa: F811 + """The re-ask at the final permitted cycle: s1 re-emits its finding → the wave stays + open, both cycles are spent, the typed cap state lands at once (no further envelope + needed), finalization is RELEASED for an honest blocked exit while no closed authority + exists, and any further re-ask is refused typed at $0 with the answers kept.""" + ctx, fp, _sub1 = _objection(harness, {"s1": _blocking(), "s2": CLEAN, "s3": CLEAN}) + _reject(ctx, fp) + sub2 = harness.install({"s1": _blocking(), "s2": CLEAN, "s3": CLEAN}) + result = _addressed(ctx, fp, "s1:f1") + assert _control(result) == {"outcome": "REVIEW_REQUIRED", "closed": False} + assert len(sub2.calls) == 1 and [s.slot_id for s in sub2.calls[0]["slots"]] == ["s1"] + + state = _state(harness) + wave = state["waves"][-1] + assert state["cycles_paid"] == 2 and wave["paid"] and not wave["closed"] + assert [f["finding_id"] for f in wave["findings"] if f["class"] == "blocking"] == ["s1:f1"] + assert wave["dispositions"] == [], "the re-asked seat's finding starts undispositioned" + assert wave["cycles_exhausted"] is True + assert state["current_attempt"] == {"fingerprint": fp, "status": "cycles_exhausted", + "reason": "2/2 paid plan-review cycles spent"} + exhausted = _events(harness, "review_cycles_exhausted") + assert exhausted and exhausted[-1]["surface"] == "plan_review" and exhausted[-1]["cycles_paid"] == 2 + gate = _gate(ctx) + assert gate["status"] == "cycles_exhausted" and gate["allow"] is True and gate["closed"] is False + assert gate["review_capacity_reason"] == "review_cycles_exhausted" and gate["outcome"] == "REVIEW_REQUIRED" + assert closed_plan_review_wave(state) is None and _bound_claims(state) == ([], "") + objective = derive_loop_outcome("done", {}, {"force_plan_decision": gate, "tool_calls": []})["outcome_axes"]["objective"] + assert (objective["status"], objective["outcome_tier"], objective["reason"]) == ( + "fail", "blocked_with_evidence", "review_cycles_exhausted") + + # A further re-ask at the spent cap: refused typed, nothing sent, the answer recorded. + _reject(ctx, fp) + sub3 = harness.install({"s1": CLEAN, "s2": CLEAN, "s3": CLEAN}) + refused = _addressed(ctx, fp, "s1:f1") + assert refused.startswith("⚠️ PLAN_REVIEW_CYCLES_EXHAUSTED") and "blocked_with_evidence" in refused + assert _control(refused) == {"outcome": "REVIEW_REQUIRED", "closed": False} + assert not sub3.calls and _state(harness)["cycles_paid"] == 2 + assert _ids(_wave(harness)) == [("s1:f1", "reject")] + # The plain identical envelope replays the open wave free; the cap state stands. + replay = _call(ctx) + assert _control(replay) == {"outcome": "REVIEW_REQUIRED", "closed": False} and "cached exact review" in replay + assert not sub3.calls and _state(harness)["cycles_paid"] == 2 + assert _gate(ctx)["status"] == "cycles_exhausted" and _wave(harness)["cycles_exhausted"] is True + + +def test_e2e_failed_objector_never_closes(harness): # noqa: F811 + """s1 answers garbage to the re-ask: its recorded finding is CARRIED (disclosed), the + wave stays REVIEW_REQUIRED and never turns GREEN; absence never retires a finding. + The gate before the re-ask holds; the final permitted cycle ending open lands the cap + state, so finalization is released for a blocked exit with no closed authority.""" + ctx, fp, _sub1 = _objection(harness, {"s1": _blocking(), "s2": CLEAN, "s3": CLEAN}) + _reject(ctx, fp) + assert _gate(ctx)["allow"] is False + sub2 = harness.install({"s1": "garbage, not an array", "s2": CLEAN, "s3": CLEAN}) + result = _addressed(ctx, fp, "s1:f1") + assert _control(result) == {"outcome": "REVIEW_REQUIRED", "closed": False} + assert len(sub2.calls) == 1 and [s.slot_id for s in sub2.calls[0]["slots"]] == ["s1"] + assert "did not answer; its earlier finding is still listed" in result + + state = _state(harness) + wave = state["waves"][-1] + assert state["cycles_paid"] == 2 and wave["paid"] and not wave["closed"] + assert wave["aggregate"] == "REVIEW_REQUIRED" + assert [f["finding_id"] for f in wave["findings"]] == ["s1:f1"] + actors = _actors(harness) + assert not actors["s1"]["ok"] and "findings_carried_absent_answer:1" in actors["s1"]["disclosures"] + assert actors["s2"]["replayed_from"]["cycle_index"] == 1 and actors["s3"]["cost"] == 0.0 + assert wave["addressed"]["slots"] == ["s1"] and wave["addressed"]["kept"] == ["s2", "s3"] + assert closed_plan_review_wave(state) is None and _bound_claims(state) == ([], "") + gate = _gate(ctx) + assert gate["closed"] is False and gate["outcome"] == "REVIEW_REQUIRED" + assert gate["status"] == "cycles_exhausted" and gate["allow"] is True # the cap rail, not a verdict + # Nothing turns it GREEN afterwards: the answers stay recorded, no seat is asked again. + sub3 = harness.install({"s1": CLEAN, "s2": CLEAN, "s3": CLEAN}) + later = _addressed(ctx, fp, "s1:f1") + assert "PLAN_REVIEW_CYCLES_EXHAUSTED" in later and not sub3.calls + assert _wave(harness)["aggregate"] == "REVIEW_REQUIRED" and not _wave(harness)["closed"] + + +# ---------------------------------------------------------------- the no-need path + +def test_e2e_no_need_path_closes_at_zero_cost(harness): # noqa: F811 + """s1 asks the author (need_evidence on claim_1, no locator) and s2 leaves a note: + REVIEW_REQUIRED holds the gate; answering the note alone changes nothing (notes never + hold); the $0 accept of the question closes GREEN with no second transport call, one + cycle paid, and the gate releases with the claims bound.""" + note = json.dumps([_finding("n1", "note", summary="a thought", rec="")]) + ctx, fp, sub = _objection(harness, {"s1": json.dumps([_question("q1")]), "s2": note, "s3": CLEAN}) + wave = _wave(harness) + assert wave["aggregate"] == "REVIEW_REQUIRED" + assert [(f["finding_id"], f["class"]) for f in wave["findings"]] == [("s1:q1", "need_evidence"), ("s2:n1", "note")] + assert _state(harness)["need_evidence_seen"] == [] # no locator: the envelope stays the same + + noted = _answer(ctx, fp, _item("s2:n1", "accept", "noted")) + assert _control(noted) == {"outcome": "REVIEW_REQUIRED", "closed": False} + assert _gate(ctx)["allow"] is False and len(sub.calls) == 1 + + closed = _answer(ctx, fp, _item("s1:q1", "accept", "the author's word")) + assert _control(closed) == {"outcome": "GREEN", "closed": True} + assert len(sub.calls) == 1, "closing by answer sends nothing" + state = _state(harness) + wave = state["waves"][-1] + assert state["cycles_paid"] == 1 and wave["closed"] and wave["aggregate"] == "GREEN" + assert wave["cycle_index"] == 1 and "addressed" not in wave + assert _ids(wave) == [("s2:n1", "accept"), ("s1:q1", "accept")] + assert [f["finding_id"] for f in wave["findings"]] == ["s1:q1", "s2:n1"] # the record keeps both + gate = _gate(ctx) + assert gate["allow"] is True and gate["status"] == "closed" + claims, source = _bound_claims(state) + assert source == "plan_review" and [c["claim"] for c in claims] == DECK_SPEC["acceptance_claims"] + # Two-sided: an identical envelope after closure replays the verdict free, sends nothing. + assert _control(_call(ctx)) == {"outcome": "GREEN", "closed": True} and len(sub.calls) == 1