tests: end-to-end addressed answer and no-need path

The route proof of the answer channel (work order item 7, amendments A10/A11),
under blocking enforcement at the shipped cycle cap of 2 in every lane.

In process (tests/test_plan_review_answer_route.py): the real engine, the real
task_results/artifact store, the fake substrate and owner_hurry.force_plan_decision.
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 (typed cycles_exhausted state, the
finalization release of D27, the blocked_with_evidence exit, a further re-ask
refused at $0 with the answers kept); the objector answers garbage (finding
carried per D4 / Q-v, never GREEN, no closed authority); the no-need path closes
at $0 with no second transport call.

Keyless system lane (tests/system_e2e/test_plan_review_addressed_answer.py):
S34 drives the addressed answer over the asynchronous barrier route with three
DISTINCT reviewer models (t1 re-asked alone, t2/t3 kept as replayed rows, 3 + 1
reviewer calls, the cycle-2 barrier snapshot names the answered cycle-1 wave,
t1's second request carries the reject rationale and its own first answer);
S35 is the no-need path (one paid cycle, three reviewer calls). The harness
gains keyless_reviewer_slots(distinct_models=True), ScriptedStubModel(model_ids)
and the S34/S35 manifest rows.

Disclosed: after the $0 collection over the barrier the collected wave carries
previous_wave_artifact and the replayed_from rows but not the typed `addressed`
block (only the barrier snapshot does); S34 asserts the chain through that
snapshot and does not change production.

Co-authored-by: Ouroboros <311266734+ouroboros-agent@users.noreply.github.com>
This commit is contained in:
Ouroboros 2026-09-26 21:10:12 +03:00 • committed by Ouroboros
parent fc76e64116
commit afb9652b11
3 changed files with 616 additions and 3 deletions

View file

@ -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<i> rides
# ``<MOCK_SLUG>-t<i>``, 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<i>`` 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:

View file

@ -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<i>`` 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()

View file

@ -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