Project an awaited reviewer slot as a gap on the shared review substrate

Task acceptance returns at the dispatch barrier just like plan review, and
the shared projection described a reviewer that had simply not answered yet
as transport_status provider_transport_error with parse_status malformed and
the custody prose as its reason. The panel reason listed "slot:Pending
dispatch..." plus "slot:degraded" as if an empty body had been parsed. Those
words reach the Reviews card, the Logs page, headless and exported results
and the reflection and task-summary prompts.

One decision point, in the per-actor projection: a row is awaiting only when
its operation_state is pending_dispatch AND it carries no answer (no parsed
object, blank text). Such a row projects transport_status and parse_status
"awaiting" and a host sentence that promises nothing; a stored "awaiting" on
a row that no longer waits is recomputed. A row that carries an answer is
judged by its answer exactly as before.

aggregate_review_actors reads that projected word: an awaited slot keeps its
place in the fail-closed clause, is no longer listed as an error or as a
parse-degraded answer, and the panel gets one sentence, "awaiting k of n
reviewer slot(s): ids". compact_review_projection folds the panel words over
the collected slots only, so a real failure beside a wait keeps its own word,
and derives the panel words instead of trusting run-level values recorded for
a slot that was only awaited.

Aggregate, quorum and every per-actor gate field are identical to the base
for every input (232-case base-versus-candidate harness, 0 mismatches):
2 PASS + 1 awaited stays DEGRADED on a fail-closed surface and PASS on task
acceptance, FAIL + awaited stays FAIL, an empty roster still reports
quorum_not_met. in_flight timeouts, custody_lost, settled failures,
not_dispatched refusals and reservation-shaped rows project and aggregate
byte-identically. No stored actor field, gate, aggregate value or projection
key changes; aggregate_review_actors grows to 154 lines (advisory signal).

Co-authored-by: Ouroboros <311266734+ouroboros-agent@users.noreply.github.com>
This commit is contained in:
Ouroboros 2026-09-20 19:13:04 +03:00 • committed by Ouroboros
parent e20f4b5757
commit 527596e5fd
4 changed files with 579 additions and 27 deletions

View file

@ -40,7 +40,7 @@ The packet is delivery-conditional; the FULL packet is not. Every triad row reac
Paid identity binds that semantic subject together with substantive nonempty obligation dispositions (`acceptance_paid_identity`); forensic source hashes and ingress counters alone do not buy a panel. A resubmit with the same paid identity reuses its recorded verdict for free: a clean replay can authorize acceptance, a non-clean replay keeps its verdict and the `identical_acceptance_refused` outcome — no repeated payment, and no cosmetic edit needed when real criteria or evidence change.
The configured slots are independent actors with adaptive quorum (`config.adaptive_quorum`: 2-of-N for N≥3, both for N=2, a single reviewer as loud `single_reviewer_no_diversity`; a fewer-responded shortfall stays a loud infra quorum failure); each receives one substantive interaction on its bound route — at most two physical sends for a packet row, one bounded episode for a native row, one delegated session for a session row — and a retrieving verdict is equally authoritative. Transport status, parse status, semantic verdict, criterion support, route, quorum contribution and binding hashes stay distinct, so an unavailable or malformed response cannot masquerade as a negative judgment; a panel that refuses before any transport projects `not_dispatched` on every row and on the panel — a transport state distinct from `success`, `timeout` and `provider_transport_error`, never a verdict. `PASS`/`FAIL`/`DEGRADED` are reviewer verdicts; the host-owned completion decision is separately `accepted`, `revision_requested` or `finalized_unaccepted`, written only by `loop_acceptance._set_acceptance_decision`. A clean quorum supplies critic approval; an informed Advisory author finish supplies separate current-author authority. Neither changes the original verdict. Material FAIL and typed unavailable outcomes reach Main before it chooses how to respond, including after the last paid panel; a no-quorum outcome with real minority findings retains that partial feedback. Blocking without fresh approval remains unaccepted, and a terminal technical failure keeps `reason=review_degraded` when no permitted author completion follows.
The configured slots are independent actors with adaptive quorum (`config.adaptive_quorum`: 2-of-N for N≥3, both for N=2, a single reviewer as loud `single_reviewer_no_diversity`; a fewer-responded shortfall stays a loud infra quorum failure); each receives one substantive interaction on its bound route — at most two physical sends for a packet row, one bounded episode for a native row, one delegated session for a session row — and a retrieving verdict is equally authoritative. Transport status, parse status, semantic verdict, criterion support, route, quorum contribution and binding hashes stay distinct, so an unavailable or malformed response cannot masquerade as a negative judgment; a panel that refuses before any transport projects `not_dispatched` on every row and on the panel, and a slot released at the dispatch barrier projects `awaiting` for transport and parse until it settles — states distinct from `success`, `timeout` and `provider_transport_error`, never a failure or a verdict. `PASS`/`FAIL`/`DEGRADED` are reviewer verdicts; the host-owned completion decision is separately `accepted`, `revision_requested` or `finalized_unaccepted`, written only by `loop_acceptance._set_acceptance_decision`. A clean quorum supplies critic approval; an informed Advisory author finish supplies separate current-author authority. Neither changes the original verdict. Material FAIL and typed unavailable outcomes reach Main before it chooses how to respond, including after the last paid panel; a no-quorum outcome with real minority findings retains that partial feedback. Blocking without fresh approval remains unaccepted, and a terminal technical failure keeps `reason=review_degraded` when no permitted author completion follows.
A clean criterion is evidence-resolved, not merely well argued: reviewer `evidence_refs` must be exact members of the packet's enumerable reference vocabulary, and a claim id resolves only through `acceptance_support_refs` linked to a passing host receipt for that claim. Agent-supplied, declared-intent, unattested and non-resolving sections never certify success; an OPEN plan wave binds nothing — its claims are disclosed as `acceptance_claims_source='none_open_plan_wave'` beside a non-binding `plan_claims_exhibit` inside `DECLARED_INTENT_SECTIONS`, so citing it never resolves and the task is distinguishable from one that never had claims. An unresolved reference keeps the actor's record for audit but removes its clean contribution (`criteria_refs_unresolved`). This total, fail-closed resolver is why the task cannot certify itself by echoing its expected outcome.

View file

@ -4,13 +4,14 @@ from __future__ import annotations
from typing import Any, Callable, Dict, List
from ouroboros.review_projection import AWAITING_PROJECTION, awaiting_panel_reason
from ouroboros.triad_review import parse_review_findings
def contract_valid_actors(result: Any) -> List[Dict[str, Any]]:
"""Actors with a DELIBERATE, CONTRACT-VALID reviewer object: parsed dict,
recognizable verdict, parse_status not "malformed" — so a contract-DEMOTED or
garbage response never votes (commit triad #1).
recognizable verdict, parse_status neither "malformed" nor "awaiting" — so a
contract-DEMOTED, garbage or unanswered row never votes (commit triad #1).
Owner ratification 2026-08-30: the acceptance-dialogue reducer now counts
votes over ``_contributing_actors`` (a slot whose verdict did not reach the
@ -26,7 +27,7 @@ def contract_valid_actors(result: Any) -> List[Dict[str, Any]]:
for actor in (getattr(result, "actors", None) or []):
row = actor if isinstance(actor, dict) else asdict(actor)
parsed = row.get("parsed")
if str(row.get("parse_status") or "") == "malformed":
if str(row.get("parse_status") or "") in {"malformed", AWAITING_PROJECTION}:
continue
if isinstance(parsed, dict) and str(
parsed.get("verdict") or parsed.get("status") or ""
@ -74,6 +75,7 @@ def aggregate_review_actors(
# slot must NOT poison a clean quorum PASS.
actor_errors: List[str] = []
parse_degraded: List[str] = []
awaiting: List[str] = []
fail_count = 0
pass_count = 0
classify_tier = bool(
@ -86,10 +88,6 @@ def aggregate_review_actors(
or str((request.policy or {}).get("hardness") or "") == advisory_hardness
)
for actor in actors:
if actor.status in {"error", "not_dispatched"}:
actor_errors.append(f"{actor.slot_id}:{actor.error}")
elif actor.status != "ok":
actor_errors.append(f"{actor.slot_id}:{actor.status}")
parsed, findings, signal = parse_review_findings(actor.raw_text)
actor.parsed = parsed
actor.signal = signal
@ -104,6 +102,16 @@ def aggregate_review_actors(
"coverage", "reason",
):
setattr(actor, key, truth[key])
# A slot released at the dispatch barrier is a gap: it holds a participation fault's
# fail-closed place without being reported as one. The projection decides it, from the
# row as a mapping (the typed predicate never matches a dataclass).
held = actor.status in {"error", "not_dispatched"} and truth["transport_status"] == AWAITING_PROJECTION
if held:
awaiting.append(actor.slot_id)
elif actor.status in {"error", "not_dispatched"}:
actor_errors.append(f"{actor.slot_id}:{actor.error}")
elif actor.status != "ok":
actor_errors.append(f"{actor.slot_id}:{actor.status}")
all_findings.extend(
{**item, "slot_id": actor.slot_id, "model": actor.model}
for item in findings
@ -160,7 +168,7 @@ def aggregate_review_actors(
)
elif signal == "PASS":
pass_count += 1
elif signal == "DEGRADED":
elif signal == "DEGRADED" and not held:
parse_degraded.append(f"{actor.slot_id}:degraded")
min_successful = max(
@ -171,16 +179,18 @@ def aggregate_review_actors(
if fail_count >= 1:
aggregate = "FAIL"
elif pass_count >= min_successful and not (
fail_closed_on_errors and actor_errors and request.surface != "task_acceptance"
fail_closed_on_errors and (actor_errors or awaiting) and request.surface != "task_acceptance"
):
aggregate = "PASS"
else:
aggregate = "DEGRADED"
if not degraded_reasons:
if not degraded_reasons and not awaiting:
degraded_reasons.append(
f"quorum_not_met: pass_count={pass_count} < min_successful={min_successful}"
)
if awaiting:
degraded_reasons.insert(0, awaiting_panel_reason(awaiting, len(actors), aggregate))
participating_ids = {
actor.slot_id
for actor in actors

View file

@ -18,6 +18,13 @@ import logging
from dataclasses import asdict
from typing import Any, Dict, List, TYPE_CHECKING
from ouroboros.review_records import review_slot_awaiting
# A slot released at the dispatch barrier has no transport or parse event:
# its projection is a gap, never a transport failure or a malformed answer.
AWAITING_PROJECTION = "awaiting"
_AWAITING_REASON = "No answer recorded: the host returned at the dispatch barrier before this reviewer answered."
if TYPE_CHECKING: # annotation-only names; lazy under future annotations, never imported at runtime
from ouroboros.review_records import ReviewActorRecord, ReviewRequest
@ -65,11 +72,22 @@ def _public_review_reason(value: Any) -> str:
return str(_sub().redact_projection(text).value)
def awaiting_panel_reason(slot_ids: List[str], configured: int, aggregate: str) -> str:
"""The one host sentence for the slots of a panel released at the dispatch barrier."""
return (f"awaiting {len(slot_ids)} of {configured} reviewer slot(s): {', '.join(slot_ids)}"
+ (" — no verdict" if aggregate == "DEGRADED" else ""))
def _review_actor_projection(actor: Any, surface: str) -> Dict[str, Any]:
row = actor if isinstance(actor, dict) else asdict(actor)
parsed = row.get("parsed") if isinstance(row.get("parsed"), (dict, list)) else None
# The one decision point: a slot is awaited only while it carries no answer. A row that
# carries one is judged by its answer, whatever its custody state says.
awaiting = review_slot_awaiting(row) and parsed is None and not str(row.get("raw_text") or "").strip()
usage = row.get("usage") if isinstance(row.get("usage"), dict) else {}
explicit_parse = str(row.get("parse_status") or "")
if explicit_parse == AWAITING_PROJECTION:
explicit_parse = "" # derived state: true only while the row is awaiting, recomputed below
semantic = str(row.get("semantic_verdict") or "").upper()
if not semantic and isinstance(parsed, dict):
semantic = str(parsed.get("verdict") or parsed.get("status") or "").upper()
@ -82,7 +100,9 @@ def _review_actor_projection(actor: Any, surface: str) -> Dict[str, Any]:
)
error = str(row.get("error") or "")
transport = str(row.get("transport_status") or "")
if not transport:
if awaiting:
transport = AWAITING_PROJECTION # the typed state outranks a word stored by an earlier projection
elif not transport or transport == AWAITING_PROJECTION:
not_dispatched = (
str(row.get("status") or "") == "not_dispatched"
or str(row.get("operation_state") or "") == "not_dispatched"
@ -121,6 +141,8 @@ def _review_actor_projection(actor: Any, surface: str) -> Dict[str, Any]:
if reason:
break
reason = reason or error or ("Reviewer response was malformed or absent." if not valid else "")
if awaiting:
reason = _AWAITING_REASON
model = str(usage.get("resolved_model") or row.get("model") or "")
provider = str(usage.get("provider") or row.get("provider") or "")
if not provider:
@ -145,7 +167,7 @@ def _review_actor_projection(actor: Any, surface: str) -> Dict[str, Any]:
"slot_id": str(row.get("slot_id") or ""), "model": model, "provider": provider,
"actor_role": str(row.get("actor_role") or f"{surface} reviewer"),
"transport_status": transport,
"parse_status": explicit_parse or ("valid" if valid else "malformed"),
"parse_status": AWAITING_PROJECTION if awaiting else (explicit_parse or ("valid" if valid else "malformed")),
"semantic_verdict": semantic if valid else "",
"outcome_tier": outcome_tier if valid else "",
"dialogue_status": dialogue_vote if valid else "",
@ -263,7 +285,9 @@ def compact_review_projection(review_runs: Any) -> Dict[str, Any]:
policy = request.get("policy") if isinstance(request.get("policy"), dict) else {}
min_successful = max(1, int(policy.get("min_successful_slots") or 1))
contributing = sum(1 for actor in actors if actor["quorum_contribution"])
transport_statuses = [actor["transport_status"] for actor in actors]
awaited = [actor["slot_id"] for actor in actors if actor["transport_status"] == AWAITING_PROJECTION]
collected = [actor for actor in actors if actor["transport_status"] != AWAITING_PROJECTION]
transport_statuses = [actor["transport_status"] for actor in collected]
transport = (
"success" if transport_statuses and all(s == "success" for s in transport_statuses)
else ("partial" if "success" in transport_statuses else (
@ -272,16 +296,35 @@ def compact_review_projection(review_runs: Any) -> Dict[str, Any]:
else "provider_transport_error")
))
)
parse = "valid" if collected and all(a["parse_status"] == "valid" for a in collected) else "malformed"
reasons = raw_run.get("degraded_reasons") if isinstance(raw_run.get("degraded_reasons"), list) else []
reasons = [str(item) for item in reasons]
aggregate = str(raw_run.get("aggregate_signal") or "UNKNOWN").upper()
if awaited:
# The panel words follow the typed rows, so a run-level word recorded for an
# awaited slot cannot outlive it. Awaited slots speak for the panel only when
# every collected slot is clean; a real failure beside a wait keeps its own word.
transport = AWAITING_PROJECTION if (not collected or transport == "success") else transport
parse = AWAITING_PROJECTION if (not collected or parse == "valid") else parse
note = awaiting_panel_reason(awaited, len(actors), aggregate)
reason = "; ".join([note] + [
item for item in reasons if item != note and item.split(":", 1)[0] not in awaited
])
else:
transport = str(raw_run.get("transport_status") or transport)
parse = str(raw_run.get("parse_status") or parse)
# v6.74.0 (A6): the fallback reason is the structured panel_reason
# reducer — it names the real blocker (tier + finding / degraded
# causes) instead of an opaque aggregate label. An explicitly
# recorded reason still wins.
reason = str(raw_run.get("reason") or "; ".join(reasons) or _sub().panel_reason(raw_run))
panel: Dict[str, Any] = {
"panel_id": str(raw_run.get("panel_id") or f"panel_{index + 1}"),
"surface": surface,
"authority": str(raw_run.get("authority") or "unspecified"),
"aggregate_signal": str(raw_run.get("aggregate_signal") or "UNKNOWN").upper(),
"transport_status": str(raw_run.get("transport_status") or transport),
"parse_status": str(raw_run.get("parse_status") or (
"valid" if actors and all(a["parse_status"] == "valid" for a in actors) else "malformed"
)),
"aggregate_signal": aggregate,
"transport_status": transport,
"parse_status": parse,
"coverage": {
"actors_configured": len(actors),
"transport_success": sum(1 for actor in actors if actor["transport_status"] == "success"),
@ -289,14 +332,7 @@ def compact_review_projection(review_runs: Any) -> Dict[str, Any]:
"quorum_contributing": contributing,
},
"quorum": {"required": min_successful, "contributed": contributing, "configured": len(actors)},
# v6.74.0 (A6): the fallback reason is the structured panel_reason
# reducer — it names the real blocker (tier + finding / degraded
# causes) instead of an opaque aggregate label. An explicitly
# recorded reason still wins.
"reason": _public_review_reason(
str(raw_run.get("reason") or "; ".join(str(item) for item in reasons)
or _sub().panel_reason(raw_run)),
),
"reason": _public_review_reason(reason),
"enforcement_impact": _review_enforcement_impact(raw_run),
"actors": actors,
"superseded": bool(raw_run.get("superseded_by_revision")),

View file

@ -0,0 +1,506 @@
"""A reviewer slot released at the dispatch barrier is a gap on the shared review
substrate: never a transport failure, a malformed answer or a verdict.
Two halves, both directions each. The WORDS about an awaited slot change (the
per-actor projection, the panel fold, the aggregate's reasons, the prompt row).
The ARITHMETIC does not: aggregate, quorum and the per-actor gate fields equal the
literal values pinned below for every input, and the loud rows (expired window, lost
custody, settled failure, typed refusal, reservation shape) keep their words. A slot
is awaited only while it carries no answer: a row that carries one is judged by its
answer exactly as before, words included, whatever its custody state says.
"""
from __future__ import annotations
import dataclasses
import json
import threading
import time
from types import SimpleNamespace
import pytest
from ouroboros.review_actor_aggregation import aggregate_review_actors, contract_valid_actors
from ouroboros.review_evidence import _acceptance_panel_prompt_row
from ouroboros.review_projection import (
AWAITING_PROJECTION,
_review_actor_projection,
compact_review_projection,
)
from ouroboros.review_records import (
HARDNESS_ADVISORY_VISIBLE,
ReviewActorRecord,
review_outcome_received,
review_slot_awaiting,
)
from ouroboros.review_verdict import _criteria_shape_valid
PASS_TEXT = json.dumps({
"verdict": "PASS", "outcome_tier": "solved", "summary": "done",
"criteria_used": [{"criterion": "works", "status": "supported", "evidence_refs": ["tool:1"]}],
"findings": [],
})
FAIL_TEXT = json.dumps({
"verdict": "FAIL", "outcome_tier": "best_effort", "summary": "broken", "completion_coach": "fix x",
"criteria_used": [{"criterion": "works", "status": "rejected"}],
"findings": [{"severity": "critical", "item": "x is broken", "recommendation": "fix x"}],
})
DEGRADED_TEXT = json.dumps({"verdict": "DEGRADED", "summary": "cannot judge", "findings": []})
PENDING_ERROR = "Pending dispatch; the physical review operation is in flight (window 21600s)"
TIMEOUT_ERROR = "Timeout after 1800s; physical review operation remains in flight"
CUSTODY_ERROR = "Exact review custody is unavailable; refusing a second paid dispatch"
REFUSAL_ERROR = "preflight_oversize: assembled acceptance prompt exceeds this slot's cap"
RUN_FAILED_ERROR = "delegated review session ended failed: harness_unavailable"
AWAITING_REASON = "No answer recorded: the host returned at the dispatch barrier before this reviewer answered."
def row(shape: str, slot: str) -> ReviewActorRecord:
"""One actor per constructor shape of the review substrate and custody."""
def make(**fields):
return ReviewActorRecord(slot_id=slot, model="m/" + slot, **fields)
op = "op-" + slot
pending = dict(status="error", operation_id=op, operation_state="pending_dispatch",
late_result_pending=True, error=PENDING_ERROR)
shapes = {
"pass": lambda: make(status="ok", raw_text=PASS_TEXT, operation_id=op),
"fail": lambda: make(status="ok", raw_text=FAIL_TEXT, operation_id=op),
"degraded": lambda: make(status="ok", raw_text=DEGRADED_TEXT, operation_id=op),
"pending": lambda: make(**pending),
"timeout": lambda: make(status="error", operation_id=op, operation_state="in_flight",
late_result_pending=True, error=TIMEOUT_ERROR),
"custody_lost": lambda: make(status="error", operation_state="custody_lost",
late_result_pending=True, error=CUSTODY_ERROR),
"not_dispatched": lambda: make(status="not_dispatched", operation_id=op,
operation_state="not_dispatched", error=REFUSAL_ERROR),
"run_failed": lambda: make(status="error", operation_id=op, operation_state="settled",
transport_status="provider_transport_error",
failure_code="run_failed", error=RUN_FAILED_ERROR),
"empty": lambda: make(status="empty", raw_text="", operation_id=op),
"malformed": lambda: make(status="ok", raw_text="I think it is fine.", operation_id=op),
# What a frozen commit reservation would look like as a substrate actor.
"reservation": lambda: make(status="error", error="", operation_id=op,
operation_state="in_flight", late_result_pending=True),
"late_settled_pass": lambda: make(status="ok", raw_text=PASS_TEXT, operation_id=op,
operation_state="late_settled"),
# No constructor mints the shapes below: a pending-dispatch row that carries an answer.
"held_raw_pass": lambda: make(**pending, raw_text=PASS_TEXT),
"held_raw_fail": lambda: make(**pending, raw_text=FAIL_TEXT),
"held_raw_prose": lambda: make(**pending, raw_text="I think it is fine."),
"ok_pending": lambda: make(status="ok", raw_text=PASS_TEXT, operation_id=op,
operation_state="pending_dispatch", late_result_pending=True),
}
return shapes[shape]()
POLICIES = {
"acceptance": ("task_acceptance", {"min_successful_slots": 2, "fail_closed_on_errors": True,
"classify_outcome_tier": True}),
"acceptance_flat": ("task_acceptance", {"min_successful_slots": 2, "fail_closed_on_errors": True}),
"fail_closed": ("plan_review", {"min_successful_slots": 2, "fail_closed_on_errors": True}),
"open": ("skill_review", {"min_successful_slots": 2}),
"quorum3": ("plan_review", {"min_successful_slots": 3, "fail_closed_on_errors": True}),
"quorum0": ("task_acceptance", {"min_successful_slots": 0}),
"quorum_negative": ("skill_review", {"min_successful_slots": -2}),
"quorum9": ("task_acceptance", {"min_successful_slots": 9}),
}
# The gate-visible row per shape and the aggregate per (case, policy), captured by running
# the sprint base commit on these fixtures: (signal, semantic_verdict, enforcement_impact,
# quorum_contribution). An awaited slot changes words only, so every value here still holds.
BASE_ROW = {
**{shape: ("PASS", "PASS", "supports_pass", True) for shape in (
"pass", "late_settled_pass", "held_raw_pass", "ok_pending")},
"fail": ("FAIL", "FAIL", "veto", True),
"held_raw_fail": ("FAIL", "FAIL", "veto", True),
"degraded": ("DEGRADED", "DEGRADED", "abstains", False),
**{shape: ("DEGRADED", "", "abstains", False) for shape in (
"pending", "timeout", "custody_lost", "not_dispatched", "run_failed", "empty",
"malformed", "reservation", "held_raw_prose")},
}
P, F, D = "PASS", "FAIL", "DEGRADED"
BASE_AGGREGATE = { # columns follow POLICIES
("pass", "pass", "pass"): (P, P, P, P, P, P, P, D),
("pass", "pending", "pending"): (D, D, D, D, D, P, P, D),
("pass", "pass", "pending"): (P, P, D, P, D, P, P, D),
("fail", "pending", "pending"): (F, F, F, F, F, F, F, F),
("pending", "pending", "pending"): (D, D, D, D, D, D, D, D),
("pending", "not_dispatched", "not_dispatched"): (D, D, D, D, D, D, D, D),
("run_failed", "pending", "pending"): (D, D, D, D, D, D, D, D),
("degraded", "pending", "pending"): (D, D, D, D, D, D, D, D),
("timeout", "pending", "pending"): (D, D, D, D, D, D, D, D),
("custody_lost", "pending", "pending"): (D, D, D, D, D, D, D, D),
("reservation", "reservation", "reservation"): (D, D, D, D, D, D, D, D),
("pass", "custody_lost", "custody_lost"): (D, D, D, D, D, P, P, D),
("pass", "pass", "custody_lost"): (P, P, D, P, D, P, P, D),
("timeout", "timeout", "timeout"): (D, D, D, D, D, D, D, D),
("pass", "pass", "timeout"): (P, P, D, P, D, P, P, D),
("not_dispatched", "not_dispatched", "not_dispatched"): (D, D, D, D, D, D, D, D),
("pass", "run_failed", "run_failed"): (D, D, D, D, D, P, P, D),
("pass", "pass", "empty"): (P, P, D, P, D, P, P, D),
("pass", "pass", "malformed"): (P, P, P, P, D, P, P, D),
("late_settled_pass", "late_settled_pass", "late_settled_pass"): (P, P, P, P, P, P, P, D),
("late_settled_pass", "late_settled_pass", "pending"): (P, P, D, P, D, P, P, D),
(): (D, D, D, D, D, D, D, D),
# A pending-dispatch row that carries an answer votes exactly as the base votes it.
("pass", "held_raw_pass", "pending"): (P, P, D, P, D, P, P, D),
("pass", "held_raw_pass", "pass"): (P, P, D, P, D, P, P, D),
("pass", "held_raw_pass", "run_failed"): (P, P, D, P, D, P, P, D),
("pass", "pass", "held_raw_fail"): (F, F, F, F, F, F, F, F),
("pass", "pass", "held_raw_prose"): (P, P, D, P, D, P, P, D),
("pass", "ok_pending", "pending"): (P, P, D, P, D, P, P, D),
("ok_pending", "pass", "pass"): (P, P, P, P, P, P, P, D),
}
INPUT_FIELDS = ("status", "operation_state", "late_result_pending", "error", "failure_code",
"operation_id", "raw_text")
def aggregate(shapes, policy="acceptance"):
surface, settings = POLICIES[policy]
actors = [row(shape, f"s{index + 1}") for index, shape in enumerate(shapes)]
slots = [SimpleNamespace(slot_id=actor.slot_id, role_hint="") for actor in actors]
result = aggregate_review_actors(
request=SimpleNamespace(surface=surface, policy=dict(settings)), slots=slots, actors=actors,
slots_by_id={slot.slot_id: slot for slot in slots},
actor_projection=_review_actor_projection, criteria_shape_valid=_criteria_shape_valid,
advisory_hardness=HARDNESS_ADVISORY_VISIBLE,
)
return result, actors
def panel_of(shapes, policy="acceptance", **run_fields):
surface, settings = POLICIES[policy]
result, actors = aggregate(shapes, policy)
run = {"request": {"surface": surface, "policy": dict(settings)}, "authority": "host_root",
"actors": [dataclasses.asdict(actor) for actor in actors], **result, **run_fields}
return compact_review_projection([run])["panels"][0]
# ── the per-actor projection ────────────────────────────────────────────────
def test_an_awaited_row_projects_a_gap_and_keeps_its_custody_identity():
actor = row("pending", "s1")
truth = _review_actor_projection(actor, "task_acceptance")
assert truth["transport_status"] == truth["parse_status"] == AWAITING_PROJECTION == "awaiting"
assert truth["semantic_verdict"] == "" and truth["quorum_contribution"] is False
assert truth["reason"] == AWAITING_REASON
assert "Pending dispatch" not in truth["reason"] and "window" not in truth["reason"]
assert "yet" not in truth["reason"]
assert (truth["operation_state"], truth["late_result_pending"], truth["operation_id"]) == (
"pending_dispatch", True, "op-s1")
# The stored row is the floor: the projection never rewrites it.
assert (actor.status, actor.error, actor.failure_code) == ("error", PENDING_ERROR, "")
@pytest.mark.parametrize("shape, transport, parse, reason", [
("timeout", "timeout", "malformed", TIMEOUT_ERROR),
("custody_lost", "provider_transport_error", "malformed", CUSTODY_ERROR),
("run_failed", "provider_transport_error", "malformed", RUN_FAILED_ERROR),
("not_dispatched", "not_dispatched", "malformed", REFUSAL_ERROR),
("reservation", "provider_transport_error", "malformed", "Reviewer response was malformed or absent."),
("empty", "success", "malformed", "Reviewer response was malformed or absent."),
])
def test_a_row_that_is_not_awaited_keeps_its_failure_words(shape, transport, parse, reason):
truth = _review_actor_projection(row(shape, "s1"), "task_acceptance")
assert (truth["transport_status"], truth["parse_status"], truth["reason"]) == (transport, parse, reason)
assert AWAITING_PROJECTION not in (truth["transport_status"], truth["parse_status"])
def test_the_typed_state_outranks_words_stored_by_an_earlier_projection():
stored = {"slot_id": "t", "model": "m", "status": "error", "error": PENDING_ERROR,
"transport_status": "provider_transport_error", "parse_status": "malformed",
"reason": PENDING_ERROR, "operation_state": "pending_dispatch", "late_result_pending": True}
truth = _review_actor_projection(stored, "task_acceptance")
assert (truth["transport_status"], truth["parse_status"], truth["reason"]) == (
"awaiting", "awaiting", AWAITING_REASON)
# The derived word never outlives the wait: a settled answer is projected afresh.
settled = {"slot_id": "t", "model": "m", "status": "ok", "transport_status": "awaiting",
"parse_status": "awaiting", "operation_state": "late_settled",
"raw_text": PASS_TEXT, "parsed": json.loads(PASS_TEXT)}
truth = _review_actor_projection(settled, "task_acceptance")
assert (truth["transport_status"], truth["parse_status"], truth["semantic_verdict"]) == (
"success", "valid", "PASS")
# ... and a wait that ended in lost custody gets its alarm back.
lost = {**stored, "transport_status": "awaiting", "parse_status": "awaiting", "reason": "",
"error": CUSTODY_ERROR, "operation_state": "custody_lost"}
truth = _review_actor_projection(lost, "task_acceptance")
assert (truth["transport_status"], truth["parse_status"], truth["reason"]) == (
"provider_transport_error", "malformed", CUSTODY_ERROR)
# ── arithmetic: identical to the base for every input ────────────────────────
@pytest.mark.parametrize("policy", list(POLICIES))
@pytest.mark.parametrize("shapes", list(BASE_AGGREGATE), ids=lambda shapes: "+".join(shapes) or "empty")
def test_aggregate_quorum_and_gate_fields_equal_the_base(shapes, policy):
before = [dataclasses.asdict(row(shape, f"s{index + 1}")) for index, shape in enumerate(shapes)]
result, actors = aggregate(shapes, policy)
expected = BASE_AGGREGATE[shapes][list(POLICIES).index(policy)]
assert result["aggregate_signal"] == expected
assert result["degraded"] is (expected == "DEGRADED")
for shape, actor, stored in zip(shapes, actors, before):
assert (actor.signal, actor.semantic_verdict, actor.enforcement_impact,
actor.quorum_contribution) == BASE_ROW[shape], shape
assert {key: getattr(actor, key) for key in INPUT_FIELDS} == {key: stored[key] for key in INPUT_FIELDS}
panel = panel_of(shapes, policy)
assert panel["aggregate_signal"] == expected
assert panel["quorum"] == {
"required": max(1, int(POLICIES[policy][1]["min_successful_slots"] or 1)),
"contributed": sum(BASE_ROW[shape][3] for shape in shapes), "configured": len(shapes)}
def test_an_awaited_slot_holds_the_fail_closed_place_without_becoming_a_pass():
# Both directions of the hold: a fail-closed surface stays DEGRADED, the advisory
# acceptance surface keeps its quorum PASS, and a reviewer FAIL stays a FAIL.
assert aggregate(("pass", "pass", "pending"), "fail_closed")[0]["aggregate_signal"] == "DEGRADED"
assert aggregate(("pass", "pass", "pending"), "acceptance")[0]["aggregate_signal"] == "PASS"
assert aggregate(("pass", "pass", "pass"), "fail_closed")[0]["aggregate_signal"] == "PASS"
assert aggregate(("fail", "pending", "pending"), "fail_closed")[0]["aggregate_signal"] == "FAIL"
assert aggregate(("pass", "pass", "pending"), "quorum3")[0]["aggregate_signal"] == "DEGRADED"
# ── the aggregate's reasons ─────────────────────────────────────────────────
def test_awaited_slots_are_reported_once_and_never_as_a_fault():
assert aggregate(("pass", "pending", "pending"))[0]["degraded_reasons"] == [
"awaiting 2 of 3 reviewer slot(s): s2, s3 — no verdict"]
assert aggregate(("pending", "pending", "pending"))[0]["degraded_reasons"] == [
"awaiting 3 of 3 reviewer slot(s): s1, s2, s3 — no verdict"]
# A quorum that is already met, or a veto, is not "no verdict".
assert aggregate(("pass", "pass", "pending"))[0]["degraded_reasons"] == [
"awaiting 1 of 3 reviewer slot(s): s3"]
assert aggregate(("pass", "pass", "pending"), "fail_closed")[0]["degraded_reasons"] == [
"awaiting 1 of 3 reviewer slot(s): s3 — no verdict"]
assert aggregate(("fail", "pending", "pending"))[0]["degraded_reasons"] == [
"awaiting 2 of 3 reviewer slot(s): s2, s3"]
@pytest.mark.parametrize("shapes, real", [
(("run_failed", "pending", "pending"), [f"s1:{RUN_FAILED_ERROR}", "s1:degraded"]),
(("timeout", "pending", "pending"), [f"s1:{TIMEOUT_ERROR}", "s1:degraded"]),
(("custody_lost", "pending", "pending"), [f"s1:{CUSTODY_ERROR}", "s1:degraded"]),
(("degraded", "pending", "pending"), ["s1:degraded"]),
])
def test_a_real_failure_beside_a_wait_stays_named(shapes, real):
assert aggregate(shapes)[0]["degraded_reasons"] == [
"awaiting 2 of 3 reviewer slot(s): s2, s3 — no verdict", *real]
def test_a_refusal_beside_a_wait_keeps_both_of_its_entries():
assert aggregate(("pending", "not_dispatched", "not_dispatched"))[0]["degraded_reasons"] == [
"awaiting 1 of 3 reviewer slot(s): s1 — no verdict",
f"s2:{REFUSAL_ERROR}", f"s3:{REFUSAL_ERROR}", "s2:degraded", "s3:degraded"]
@pytest.mark.parametrize("shapes, reasons", [
(("pass", "run_failed", "run_failed"),
[f"s2:{RUN_FAILED_ERROR}", f"s3:{RUN_FAILED_ERROR}", "s2:degraded", "s3:degraded"]),
(("not_dispatched",) * 3,
[f"s1:{REFUSAL_ERROR}", f"s2:{REFUSAL_ERROR}", f"s3:{REFUSAL_ERROR}",
"s1:degraded", "s2:degraded", "s3:degraded"]),
(("timeout",) * 3,
[f"s1:{TIMEOUT_ERROR}", f"s2:{TIMEOUT_ERROR}", f"s3:{TIMEOUT_ERROR}",
"s1:degraded", "s2:degraded", "s3:degraded"]),
(("pass", "custody_lost", "custody_lost"),
[f"s2:{CUSTODY_ERROR}", f"s3:{CUSTODY_ERROR}", "s2:degraded", "s3:degraded"]),
(("reservation",) * 3, ["s1:", "s2:", "s3:", "s1:degraded", "s2:degraded", "s3:degraded"]),
(("pass", "pass", "empty"), ["s3:empty", "s3:degraded"]),
(("pass", "pass", "malformed"), ["s3:degraded"]),
((), ["quorum_not_met: pass_count=0 < min_successful=2"]),
(("pass", "pass", "pass"), []),
])
def test_reasons_without_an_awaited_slot_are_the_base_lists(shapes, reasons):
assert aggregate(shapes)[0]["degraded_reasons"] == reasons
# ── one decision point: awaited only while the row carries no answer ─────────
def test_the_hold_is_decided_from_the_row_as_a_mapping(monkeypatch):
actor = row("pending", "s1")
assert review_slot_awaiting(actor) is False # the typed predicate never matches a dataclass
assert review_slot_awaiting(dataclasses.asdict(actor)) is True
import ouroboros.review_projection as projection
seen = []
monkeypatch.setattr(projection, "review_slot_awaiting",
lambda value: seen.append(type(value)) or review_slot_awaiting(value))
result, actors = aggregate(("pass", "pending", "run_failed"))
assert seen == [dict, dict, dict] # the projection asks about every row, always as a mapping
assert result["degraded_reasons"] == [
"awaiting 1 of 3 reviewer slot(s): s2 — no verdict", f"s3:{RUN_FAILED_ERROR}", "s3:degraded"]
# The aggregation and the stamped row agree with the projection, never with a second reading.
assert [actor.transport_status for actor in actors] == ["success", "awaiting", "provider_transport_error"]
@pytest.mark.parametrize("shape, words, vote", [
("held_raw_pass", ("provider_transport_error", "valid", "PASS", "solved", "done"), True),
("held_raw_fail", ("provider_transport_error", "valid", "FAIL", "best_effort", "broken"), True),
("held_raw_prose", ("provider_transport_error", "malformed", "", "", PENDING_ERROR), False),
("ok_pending", ("success", "valid", "PASS", "solved", "done"), True),
])
def test_a_pending_row_that_carries_an_answer_is_judged_by_its_answer(shape, words, vote):
result, actors = aggregate(("pass", "pass", shape))
actor, stored = actors[-1], dataclasses.asdict(actors[-1])
truth = _review_actor_projection(stored, "task_acceptance")
assert (truth["transport_status"], truth["parse_status"], truth["semantic_verdict"],
truth["outcome_tier"], truth["reason"]) == words
assert actor.quorum_contribution is vote
assert ("s3" in {voter["slot_id"] for voter in contract_valid_actors(SimpleNamespace(actors=[stored]))}) is vote
# Its reasons are the base's: a participation fault by name, never an awaited slot.
assert result["degraded_reasons"] == {
"held_raw_pass": [f"s3:{PENDING_ERROR}"], "held_raw_fail": [f"s3:{PENDING_ERROR}"],
"held_raw_prose": [f"s3:{PENDING_ERROR}", "s3:degraded"], "ok_pending": [],
}[shape]
panel = panel_of(("pass", "pass", shape))
assert "awaiting" not in json.dumps(panel)
assert (panel["transport_status"], panel["coverage"]["transport_success"]) == (
("success", 3) if shape == "ok_pending" else ("partial", 2))
def test_a_stored_parsed_answer_on_a_pending_row_projects_and_votes_as_before():
stale = {"slot_id": "s1", "model": "m", "status": "error", "error": PENDING_ERROR,
"operation_state": "pending_dispatch", "late_result_pending": True,
"parsed": {**json.loads(PASS_TEXT), "dialogue_status": "continue_actionable"}}
truth = _review_actor_projection(stale, "task_acceptance")
assert {key: truth[key] for key in ("transport_status", "parse_status", "semantic_verdict", "outcome_tier",
"dialogue_status", "coverage", "reason", "findings")} == {
"transport_status": "provider_transport_error", "parse_status": "valid", "semantic_verdict": "PASS",
"outcome_tier": "solved", "dialogue_status": "continue_actionable",
"coverage": {"criteria_total": 1, "findings": 0}, "reason": "done", "findings": []}
assert [voter["slot_id"] for voter in contract_valid_actors(SimpleNamespace(actors=[stale]))] == ["s1"]
# Without the answer the same row is a gap, and a gap never votes.
silent = {key: value for key, value in stale.items() if key != "parsed"}
truth = _review_actor_projection(silent, "task_acceptance")
assert (truth["transport_status"], truth["parse_status"], truth["semantic_verdict"]) == ("awaiting", "awaiting", "")
assert contract_valid_actors(SimpleNamespace(actors=[{**silent, **truth, "parsed": json.loads(PASS_TEXT)}])) == []
# ── the panel ───────────────────────────────────────────────────────────────
@pytest.mark.parametrize("shapes, transport, parse", [
(("pending", "pending", "pending"), "awaiting", "awaiting"),
(("pass", "pass", "pending"), "awaiting", "awaiting"),
(("fail", "pending", "pending"), "awaiting", "awaiting"),
(("degraded", "pending", "pending"), "awaiting", "awaiting"),
# A real failure beside a wait keeps its own word.
(("run_failed", "pending", "pending"), "provider_transport_error", "malformed"),
(("custody_lost", "pending", "pending"), "provider_transport_error", "malformed"),
(("timeout", "pending", "pending"), "timeout", "malformed"),
(("pending", "not_dispatched", "not_dispatched"), "not_dispatched", "malformed"),
(("pass", "run_failed", "pending"), "partial", "malformed"),
# No awaited slot: the base fold, literally.
(("pass", "pass", "pass"), "success", "valid"),
(("pass", "pass", "timeout"), "partial", "malformed"),
(("pass", "pass", "malformed"), "success", "malformed"),
(("timeout",) * 3, "timeout", "malformed"),
(("not_dispatched",) * 3, "not_dispatched", "malformed"),
(("reservation",) * 3, "provider_transport_error", "malformed"),
((), "provider_transport_error", "malformed"),
])
def test_panel_words_follow_the_collected_slots(shapes, transport, parse):
panel = panel_of(shapes)
assert (panel["transport_status"], panel["parse_status"]) == (transport, parse)
answered = sum(shape in {"pass", "fail", "degraded", "malformed"} for shape in shapes)
assert panel["coverage"]["transport_success"] == answered # counts stay over every slot
def test_a_settled_panel_projects_the_base_reason_and_an_explicit_one_still_wins():
failed = panel_of(("pass", "run_failed", "run_failed"))
assert failed["reason"] == (f"s2:{RUN_FAILED_ERROR}; s3:{RUN_FAILED_ERROR}; s2:degraded; s3:degraded")
explicit = panel_of(("pass", "run_failed", "run_failed"), reason="recorded by the host",
transport_status="timeout", parse_status="valid")
assert (explicit["reason"], explicit["transport_status"], explicit["parse_status"]) == (
"recorded by the host", "timeout", "valid")
def test_a_run_recorded_with_failure_words_about_a_wait_reprojects_as_awaiting():
noise = [f"s1:{RUN_FAILED_ERROR}", f"s2:{PENDING_ERROR}", f"s3:{PENDING_ERROR}",
"s1:degraded", "s2:degraded", "s3:degraded"]
panel = panel_of(("run_failed", "pending", "pending"), degraded_reasons=noise,
reason="; ".join(noise), transport_status="provider_transport_error",
parse_status="malformed")
assert panel["reason"] == ("awaiting 2 of 3 reviewer slot(s): s2, s3 — no verdict; "
f"s1:{RUN_FAILED_ERROR}; s1:degraded")
waiting = panel_of(("pending", "pending", "pending"), reason=f"s1:{PENDING_ERROR}",
degraded_reasons=[f"s1:{PENDING_ERROR}", "s1:degraded"],
transport_status="provider_transport_error", parse_status="malformed")
assert (waiting["transport_status"], waiting["parse_status"]) == ("awaiting", "awaiting")
assert waiting["reason"] == "awaiting 3 of 3 reviewer slot(s): s1, s2, s3 — no verdict"
prompt_row = _acceptance_panel_prompt_row(waiting)
assert (prompt_row["transport_status"], prompt_row["parse_status"], prompt_row["aggregate_signal"]) == (
"awaiting", "awaiting", "DEGRADED")
assert prompt_row["reason"] == waiting["reason"] and "Pending dispatch" not in json.dumps(prompt_row)
def test_a_superseded_awaited_panel_keeps_its_words_and_its_flag():
panel = panel_of(("pass", "pending", "pending"), superseded_by_revision=True)
assert panel["superseded"] is True
assert (panel["transport_status"], panel["aggregate_signal"]) == ("awaiting", "DEGRADED")
# ── end to end through the substrate ────────────────────────────────────────
def test_a_released_acceptance_run_reads_as_awaiting_and_then_as_its_answers(tmp_path, monkeypatch):
import ouroboros.review_custody as custody
from ouroboros.loop_acceptance_review import acceptance_run_pending
from ouroboros.review_dispatch import collect_task_acceptance_run
from ouroboros.review_substrate import ReviewRequest, ReviewSlot, run_review_request
release, settled, entered = threading.Event(), threading.Semaphore(0), threading.Semaphore(0)
original_settle = custody._settle_review_attempt
def settle(*args, **kwargs):
try:
return original_settle(*args, **kwargs)
finally:
settled.release()
monkeypatch.setattr(custody, "_settle_review_attempt", settle)
class HeldModel:
def chat(self, **kwargs):
entered.release()
assert release.wait(10), "fixture did not release its model"
return {"content": PASS_TEXT}, {"prompt_tokens": 5, "completion_tokens": 2}
ctx = SimpleNamespace(task_id="awaiting-root", task_attempt=1, drive_root=tmp_path,
budget_drive_root=tmp_path, task_metadata={}, pending_events=[], event_queue=None)
request = ReviewRequest(surface="task_acceptance", task_id=ctx.task_id, goal="goal", subject="result",
evidence={"requirement": "exact"}, retry_key="awaiting-subject",
policy={"min_successful_slots": 2}, drain_deadline=time.monotonic())
slots = [ReviewSlot(slot_id=f"s{index}", model=f"model/{index}", effort="high", timeout_sec=20)
for index in (1, 2, 3)]
try:
first = run_review_request(request, slots=slots, drive_root=tmp_path, usage_ctx=ctx, llm=HeldModel())
assert all(entered.acquire(timeout=5) for _ in slots)
# FLOOR: the stored rows and every gate reading of them are untouched.
for actor in first.actors:
assert actor["status"] == "error" and actor["error"].startswith("Pending dispatch;")
assert (actor["operation_state"], actor["late_result_pending"], actor["failure_code"]) == (
"pending_dispatch", True, "")
assert (actor["signal"], actor["quorum_contribution"], actor["enforcement_impact"]) == (
"DEGRADED", False, "abstains")
assert (first.aggregate_signal, first.degraded) == ("DEGRADED", True)
assert acceptance_run_pending(first) and not review_outcome_received(first.actors)
# The words about those rows.
assert {(actor["transport_status"], actor["parse_status"]) for actor in first.actors} == {
("awaiting", "awaiting")}
assert first.degraded_reasons == ["awaiting 3 of 3 reviewer slot(s): s1, s2, s3 — no verdict"]
frozen = json.loads(json.dumps(dataclasses.asdict(first)))
panel = compact_review_projection([{**frozen, "authority": "host_root"}])["panels"][0]
assert (panel["transport_status"], panel["parse_status"], panel["aggregate_signal"]) == (
"awaiting", "awaiting", "DEGRADED")
assert panel["reason"] == first.degraded_reasons[0]
assert panel["coverage"]["transport_success"] == 0
release.set()
assert all(settled.acquire(timeout=5) for _ in slots)
result = collect_task_acceptance_run(frozen, drive_root=tmp_path, usage_ctx=ctx)
assert not acceptance_run_pending(result)
assert (result.aggregate_signal, result.degraded_reasons) == ("PASS", [])
assert [actor["operation_id"] for actor in result.actors] == [
actor["operation_id"] for actor in first.actors]
done = compact_review_projection([{**dataclasses.asdict(result), "authority": "host_root"}])["panels"][0]
assert (done["transport_status"], done["parse_status"]) == ("success", "valid")
assert "awaiting" not in json.dumps(done)
finally:
release.set()