mirror of
https://github.com/razzant/ouroboros.git
synced 2026-10-03 04:07:04 +00:00
plan-review: merge answers by finding_id; answer and review in one call; drop the automatic delta
A reviewer's question or objection reached the author, but the author's answer could not travel back as a normal move: a second review_disposition call REPLACED the wave's whole answer list (8 of 9 answers were lost on one live wave), answering and re-asking in one call was refused as PLAN_REVIEW_DISPOSITION_MIXED_ENVELOPE (21 refusals in 7 tasks), items beside an author finish were refused, and the host bought a delta panel by itself whenever every blocking finding carried a valid reject (a path that ran 0 times in production). Now answers MERGE by finding_id across calls (plan_spec.merge_dispositions: a later answer supersedes only its own id; two entries for one id in ONE call stay contradictory and the closure table keeps the finding open), both in the engine's closure/exact artifact and in the durable writer under its lock. An envelope sent beside review_disposition items is validated FIRST through the read-only prepare seam (an invalid envelope records nothing, so the answered wave stays current), the answers are recorded merged, and the envelope is then reviewed with them in view (_apply_disposition(then_review=...)). An author finish or stop carries its items: they are validated against the critic wave and recorded before the author source; the kept guard (finish while reviewers run) still refuses and writes nothing. The automatic earned delta and plan_spec.blocking_fully_rejected are deleted: an identical envelope without items always replays free, and re-judgement is the mind's explicit move. Owner authority: DECISIONS "Communication" item 2 (communication is an option, not an obligation; reuse surfaces, no if-else); D4 with Q-v = A; "Blocking installs: no skip". _apply_disposition is split into _disposition_items and _record_disposition so the author path reuses them. task_results.py stays under its hard line cap (1596/1600). Tests: tests/test_plan_review_answer_channel.py (merge across calls, supersede-own-id, same-call duplicate, durable writer, author finish with items, refused finish writes nothing, unknown id beside a finish, invalid envelope records nothing, changed envelope with items reviews every slot with the rationale in PRIOR CYCLES); plan_spec units for merge_dispositions; the MIXED-envelope pins rewritten to the validate-first form; the earned-delta tests deleted or rewritten as "replay is free until the mind addresses". Co-authored-by: Ouroboros <311266734+ouroboros-agent@users.noreply.github.com>
This commit is contained in:
parent
91cb19e32a
commit
55b7113b75
13 changed files with 436 additions and 251 deletions
|
|
@ -82,9 +82,12 @@ Accept, reject or defer findings. Disposition-only
|
|||
`plan_task(review_disposition={review_fingerprint, items:[{finding_id, decision, rationale}]})`
|
||||
closes `need_evidence` at $0, one item per required finding; under advisory a
|
||||
reasoned reject also closes a below-quorum blocking finding.
|
||||
Duplicate, conflicting, unknown, stale, incomplete, mixed or vacuous calls
|
||||
Answers merge by `finding_id` across calls: a later answer supersedes only its
|
||||
own id, and two entries for one id in ONE call stay contradictory and open.
|
||||
Duplicate, conflicting, unknown, stale, incomplete or vacuous calls
|
||||
return typed argument errors before recording; `plan_review._handle_plan_task` ignores
|
||||
default-empty optional fields. Do not replay plans for dispositions.
|
||||
default-empty optional fields. An envelope sent with items records them first,
|
||||
then reviews.
|
||||
|
||||
Only explicit `review_disposition.author_action`, author disposition and critic
|
||||
fingerprint select corrected goal/plan/spec. Exact `current_attempt.author_subject`
|
||||
|
|
@ -121,7 +124,7 @@ post-consolidation reader or the acceptance directive ledger that task
|
|||
acceptance keeps (`review_evidence._accept_owner_directives`). JSONL records
|
||||
and chat line selectors split on physical LF only, never on valid Unicode
|
||||
inside a message. Each consumer redacts at its boundary and discloses missing
|
||||
source or ranges; a replay or earned paid retry of the same author request
|
||||
source or ranges; a replay or an addressed re-ask of the same author request
|
||||
keeps its recorded snapshot (complete and uncapped; the inline view is the
|
||||
conversation only, with a pointer naming exact omitted line ranges),
|
||||
disclosing later messages as unreviewed. Follow
|
||||
|
|
|
|||
|
|
@ -10,7 +10,8 @@ default would silently swallow ``"unlimited"``. ``review_max_cycles()`` returns
|
|||
``Optional[int]`` — ``None`` means unlimited.
|
||||
|
||||
Per-gate meaning of the ONE number — on every gate it counts PAID cycles, and
|
||||
identical material is never re-reviewed for pay:
|
||||
identical material is never re-reviewed for pay unless the mind sends it with
|
||||
answers (a plan re-ask addressed by ``review_disposition`` items is one paid cycle):
|
||||
|
||||
* plan review — paid reviewer-panel cycles per task (the engine consumes the
|
||||
getter; this module only exposes it);
|
||||
|
|
|
|||
|
|
@ -1486,8 +1486,8 @@ def record_plan_review_wave(
|
|||
waves.append(recorded)
|
||||
if not state.get("series_id"):
|
||||
state["series_id"] = fingerprint[:16]
|
||||
# C-07: replacement writes don't charge again; a fully-rejected wave's
|
||||
# earned delta advances cycle_index and charges its new physical panel.
|
||||
# C-07: replacement writes don't charge again; an addressed answer advances
|
||||
# cycle_index and charges its new physical panel.
|
||||
already_paid = any(
|
||||
w.get("paid") and int(w.get("cycle_index") or 0) >= int(wave.get("cycle_index") or 0)
|
||||
for w in previous
|
||||
|
|
@ -1541,8 +1541,9 @@ def record_plan_review_dispositions(
|
|||
recorded_at: str = "",
|
||||
author_disposition: Optional[Dict[str, Any]] = None,
|
||||
) -> Dict[str, Any]:
|
||||
"""Store the agent's dispositions on one FULL wave and its resulting closure.
|
||||
Only note-only closed waves accept annotations. Closure authority remains
|
||||
"""Record this call's answers MERGED into one FULL wave's answers by ``finding_id``
|
||||
(``plan_spec.merge_dispositions``: a later answer supersedes only its own id) and the
|
||||
resulting closure. Only note-only closed waves accept annotations. Closure authority remains
|
||||
``plan_spec.closure_after_disposition`` (``aggregate`` is the verdict that
|
||||
table says to record — GREEN when a REVIEW_REQUIRED open set emptied); this
|
||||
writer is rule-free. Other closed waves are immutable."""
|
||||
|
|
@ -1556,7 +1557,9 @@ def record_plan_review_dispositions(
|
|||
raise ValueError("PLAN_REVIEW_DISPOSITION_STALE: a newer attempt supersedes this wave")
|
||||
if wave.get("closed") and not plan_review_notes_are_annotatable(wave):
|
||||
raise ValueError("PLAN_REVIEW_DISPOSITION_IMMUTABLE: a closed wave cannot be changed")
|
||||
wave["dispositions"] = copy.deepcopy(list(dispositions))
|
||||
from ouroboros.tools.plan_spec import merge_dispositions
|
||||
|
||||
wave["dispositions"] = merge_dispositions(wave.get("dispositions"), dispositions)
|
||||
wave["disposition_recorded_at"] = recorded_at or utc_now_iso()
|
||||
if closure_notes is not None:
|
||||
wave["closure_notes"] = list(closure_notes)
|
||||
|
|
|
|||
|
|
@ -21,9 +21,10 @@ verdict); need_evidence closes by disposition at $0; under advisory enforcement
|
|||
reject with its rationale also closes a below-quorum blocking finding (per
|
||||
finding), under blocking it stays open until a changed spec is reviewed or the
|
||||
reviewer retires it; a REVIEW_REQUIRED whose open set empties is recorded GREEN.
|
||||
REVISE_PLAN never closes by disposition. A subsequent paid delta review may
|
||||
evaluate a changed spec or a justified rejection when another paid cycle is
|
||||
available. Under blocking enforcement an open wave HOLDS
|
||||
REVISE_PLAN never closes by disposition. Answers merge by ``finding_id`` across
|
||||
calls; an envelope sent with ``review_disposition`` items records them first and
|
||||
is then reviewed; no host path buys a panel the mind did not send. Under
|
||||
blocking enforcement an open wave HOLDS
|
||||
finalization (``owner_hurry.force_plan_decision``); at the cap the typed
|
||||
``plan_review_cycles_exhausted`` result + event leave the honest exits: owner
|
||||
unstick or a ``blocked_with_evidence`` terminal. Advisory proceeds open under the
|
||||
|
|
@ -372,30 +373,42 @@ def _handle_plan_task(ctx: ToolContext, **params) -> str:
|
|||
if raw_disposition is not None and not disposition_vacuous:
|
||||
if isinstance(raw_disposition, dict) and "author_action" in raw_disposition:
|
||||
return _apply_author_subject(ctx, raw_disposition, params if envelope_fields else None)
|
||||
if envelope_fields:
|
||||
return _typed_refusal(
|
||||
ctx, "TOOL_ARG_ERROR",
|
||||
"ERROR: PLAN_REVIEW_DISPOSITION_MIXED_ENVELOPE: disposition mode accepts "
|
||||
"review_disposition only; a changed plan needs a new review-mode call "
|
||||
"without review_disposition. No plan attempt was recorded. Non-empty plan fields: "
|
||||
+ _argument_values(params, envelope_fields),
|
||||
)
|
||||
if not isinstance(raw_disposition, dict):
|
||||
return _typed_refusal(
|
||||
ctx, "TOOL_ARG_ERROR",
|
||||
"ERROR: PLAN_REVIEW_DISPOSITION_INVALID: review_disposition must be an object",
|
||||
)
|
||||
if not envelope_fields:
|
||||
return _apply_disposition(ctx, raw_disposition)
|
||||
# Answers travelling with an envelope: the envelope is validated FIRST through the
|
||||
# read-only prepare (an invalid one records nothing — a superseding raw attempt would
|
||||
# make the answered wave stale and its answers undeliverable), the answers are
|
||||
# recorded merged, then the envelope is reviewed with them in view.
|
||||
request = _request_of(params)
|
||||
try:
|
||||
prepared = _prepare_plan_inputs(ctx, request, _planning_state_location(ctx)[0])
|
||||
except ValueError as exc:
|
||||
return _typed_refusal(ctx, "TOOL_ERROR", f"ERROR: PLAN_REVIEW_STATE_INVALID: {exc}")
|
||||
if prepared.get("error"):
|
||||
return _typed_refusal(ctx, str(prepared.get("code") or "TOOL_ARG_ERROR"), prepared["error"])
|
||||
return _apply_disposition(ctx, raw_disposition, then_review=lambda items: _review(ctx, request))
|
||||
if "review_disposition" in params and not envelope_fields:
|
||||
return _typed_refusal(
|
||||
ctx, "TOOL_ARG_ERROR",
|
||||
"ERROR: PLAN_REVIEW_DISPOSITION_EMPTY: submit goal, plan and spec for review "
|
||||
"mode, or a complete review_disposition as the only field. No plan attempt was recorded.",
|
||||
)
|
||||
request = _PlanRequest(
|
||||
return _review(ctx, _request_of(params))
|
||||
|
||||
|
||||
def _request_of(params: dict) -> "_PlanRequest":
|
||||
return _PlanRequest(
|
||||
goal=str(params.get("goal") or ""), plan=str(params.get("plan") or ""), spec=params.get("spec"),
|
||||
reviewer_effort=params.get("reviewer_effort"),
|
||||
)
|
||||
|
||||
|
||||
def _review(ctx: ToolContext, request: "_PlanRequest") -> str:
|
||||
try: # the ToolEntry envelope is the outer settlement bound (plan_review_collect.run_plan_coroutine)
|
||||
return _collect.run_plan_coroutine(_run_plan_review_async(ctx, request))
|
||||
except Exception as e:
|
||||
|
|
@ -676,20 +689,9 @@ async def _run_plan_review_async(ctx: ToolContext, request: _PlanRequest, *, col
|
|||
return _publish_rendered_wave(ctx, existing, cap=cap, cycles_paid=cycles_paid,
|
||||
enforcement=enforcement, cached=True, reminder=reminder, historical_feedback=historical)
|
||||
resume_in_flight = _plan_wave_has_in_flight(existing)
|
||||
# Identical requests replay free unless authority lapsed; fully rejected
|
||||
# blocking findings are the one earned-delta exception (4e133c8a).
|
||||
earned_delta = (
|
||||
str(existing.get("aggregate")) in {"REVISE_PLAN", "REVIEW_REQUIRED"}
|
||||
and not existing.get("closed")
|
||||
# D1: only VALID rejections earn the panel — raw items are persisted for
|
||||
# disclosure, but an invalid or contradictory one must not buy a cycle.
|
||||
and plan_spec.blocking_fully_rejected(
|
||||
existing.get("findings"), existing.get("dispositions"))
|
||||
)
|
||||
if earned_delta and not resume_in_flight:
|
||||
# the rejected wave IS the previous cycle for the delta panel
|
||||
previous_override = existing
|
||||
elif bool(existing.get("closed")) and not resume_in_flight:
|
||||
# Identical requests replay free unless authority lapsed: no host path buys a
|
||||
# panel the mind did not send (an addressed re-ask carries its items explicitly).
|
||||
if bool(existing.get("closed")) and not resume_in_flight:
|
||||
# A closed verdict is earned authority; later roster changes govern
|
||||
# future panels and do not retroactively void it (accepted 3a).
|
||||
return _publish_rendered_wave(ctx, existing, cap=cap, cycles_paid=cycles_paid,
|
||||
|
|
@ -1040,6 +1042,7 @@ def _apply_author_subject(ctx: ToolContext, disposition: dict, envelope: Optiona
|
|||
from ouroboros.observability import redact_projection
|
||||
from ouroboros.review_custody import review_retry_cancelled
|
||||
|
||||
items: List[dict] = []
|
||||
try:
|
||||
if set(disposition) - {"review_fingerprint", "items", "author_disposition", "author_action"}:
|
||||
raise ValueError("unknown disposition fields: " + _argument_values(disposition,
|
||||
|
|
@ -1060,7 +1063,11 @@ def _apply_author_subject(ctx: ToolContext, disposition: dict, envelope: Optiona
|
|||
if wave:
|
||||
wave = _authority_wave(root, task_id, wave)
|
||||
if disposition.get("items"):
|
||||
raise ValueError("submit per-finding dispositions separately before selecting the current author plan")
|
||||
if not wave:
|
||||
raise ValueError("items answer the findings of a recorded wave; this outcome has none")
|
||||
items, item_error = _disposition_items(wave, disposition.get("items"))
|
||||
if item_error:
|
||||
raise ValueError(item_error.removeprefix("ERROR: PLAN_REVIEW_DISPOSITION_INVALID: "))
|
||||
from ouroboros.review_records import review_outcome_received
|
||||
|
||||
if action == "finish" and review_enforcement_blocks("blocking") and not unavailable:
|
||||
|
|
@ -1079,6 +1086,8 @@ def _apply_author_subject(ctx: ToolContext, disposition: dict, envelope: Optiona
|
|||
else:
|
||||
raise ValueError("include goal, plan and spec to retain the exact current author plan")
|
||||
enforcement = get_review_enforcement()
|
||||
if items: # the same call's answers land on the critic wave first, merged by finding_id
|
||||
wave, _closure = _record_disposition(root, task_id, wave, items, fingerprint=critic_fp, enforcement=enforcement)
|
||||
author = build_author_disposition_from_mapping(disposition.get("author_disposition"),
|
||||
subject_hash=fingerprint, reviewer_signal=str((wave or {}).get("aggregate") or "unavailable"), enforcement=enforcement)
|
||||
author["action"] = action
|
||||
|
|
@ -1102,6 +1111,7 @@ def _apply_author_subject(ctx: ToolContext, disposition: dict, envelope: Optiona
|
|||
signal = str((wave or {}).get("aggregate") or "DEGRADED")
|
||||
historical = bool(wave) and fingerprint != critic_fp
|
||||
text = (f"Current author plan saved: {fingerprint}. Critic subject: {critic_fp}. "
|
||||
+ (f"{len(items)} answer(s) recorded on the critic wave, merged by finding_id. " if items else "")
|
||||
+ (f"Earlier plan {critic_fp} was {signal}; this revised plan has no verdict of its own. "
|
||||
if historical else "")
|
||||
+ "No reviewer called and no cycle consumed; original findings and custody remain unchanged. "
|
||||
|
|
@ -1114,7 +1124,66 @@ def _apply_author_subject(ctx: ToolContext, disposition: dict, envelope: Optiona
|
|||
"author_action": action, "author_disposition": author}, text)
|
||||
|
||||
|
||||
def _apply_disposition(ctx: ToolContext, disposition: dict) -> str:
|
||||
def _disposition_items(wave: dict, raw_items: Any) -> tuple[List[dict], str]:
|
||||
"""Bound and validate one call's answers against the wave's findings → ``(items, error)``."""
|
||||
if not isinstance(raw_items, list):
|
||||
return [], "ERROR: PLAN_REVIEW_DISPOSITION_INVALID: items must be an array"
|
||||
if len(raw_items) > 2 * len(wave.get("findings") or []) + 8: # bounded like the findings they answer
|
||||
return [], "ERROR: PLAN_REVIEW_DISPOSITION_INVALID: more items than findings could need"
|
||||
items: List[dict] = []
|
||||
for index, item in enumerate(raw_items):
|
||||
if not isinstance(item, dict):
|
||||
return [], f"ERROR: PLAN_REVIEW_DISPOSITION_INVALID: items[{index}] must be an object"
|
||||
items.append({
|
||||
"finding_id": str(item.get("finding_id") or "").strip()[:plan_spec.MAX_ID_CHARS * 2],
|
||||
"decision": str(item.get("decision") or "").strip().lower()[:40], # enum-like, bounded
|
||||
"rationale": plan_spec.bounded_text(item.get("rationale"), plan_spec.MAX_FINDING_TEXT_CHARS),
|
||||
})
|
||||
known = {str(f.get("finding_id") or "") for f in wave.get("findings") or []}
|
||||
unknown_ids = sorted({i["finding_id"] for i in items if i["finding_id"] not in known})
|
||||
if unknown_ids:
|
||||
return [], ("ERROR: PLAN_REVIEW_DISPOSITION_INVALID: unknown finding ids " + ", ".join(unknown_ids)
|
||||
+ "; valid ids: " + ", ".join(sorted(known)))
|
||||
return items, ""
|
||||
|
||||
|
||||
def _record_disposition(root: pathlib.Path, task_id: str, wave: dict, items: List[dict], *,
|
||||
fingerprint: str, enforcement: str, author_record: Optional[dict] = None) -> tuple[dict, dict]:
|
||||
"""Merge this call's answers into the wave's (``plan_spec.merge_dispositions``), close over
|
||||
the MERGED answers, write the exact artifact, then the hot record → ``(stored, closure)``.
|
||||
Store errors (OSError/TimeoutError/ValueError, incl. a stale or immutable wave) propagate."""
|
||||
merged = plan_spec.merge_dispositions(wave.get("dispositions"), items)
|
||||
closure = plan_spec.closure_after_disposition(
|
||||
str(wave.get("aggregate") or ""), wave.get("findings") or [], merged, enforcement,
|
||||
)
|
||||
disposition_recorded_at = utc_now_iso()
|
||||
prior_ref = wave.get("wave_artifact") if isinstance(wave.get("wave_artifact"), dict) else {}
|
||||
closure_notes = [*closure["notes"], *([] if prior_ref else [
|
||||
"exact_artifact_absent: v2 wave had no exact wave_artifact reference",
|
||||
])]
|
||||
exact = _read_plan_review_wave_artifact(root, task_id, prior_ref) if prior_ref else dict(wave)
|
||||
exact.update({
|
||||
"dispositions": merged, "closed": bool(closure["closed"]),
|
||||
"aggregate": closure["aggregate"], "closure_notes": closure_notes,
|
||||
"disposition_recorded_at": disposition_recorded_at,
|
||||
"supersedes_wave_artifact": prior_ref,
|
||||
})
|
||||
if author_record is not None:
|
||||
exact["author_disposition"] = author_record
|
||||
disposition_ref = _persist_plan_review_wave_artifact(root, task_id, exact)
|
||||
stored = record_plan_review_dispositions(
|
||||
root, task_id, fingerprint=fingerprint, dispositions=merged,
|
||||
closed=bool(closure["closed"]), aggregate=closure["aggregate"], closure_notes=closure_notes,
|
||||
wave_artifact=disposition_ref, recorded_at=disposition_recorded_at,
|
||||
author_disposition=author_record,
|
||||
)
|
||||
return stored, closure
|
||||
|
||||
|
||||
def _apply_disposition(ctx: ToolContext, disposition: dict, *, then_review=None) -> str:
|
||||
"""``then_review(items)`` — the review of an envelope sent beside the answers — replaces
|
||||
this call's own rendering at its two non-refusal exits, so nothing is published twice;
|
||||
refusals and the already-closed note stand alone (a closed, immutable wave takes no answers)."""
|
||||
def _bad(text: str) -> str: # every refusal below is an argument-shape refusal
|
||||
return _typed_refusal(ctx, "TOOL_ARG_ERROR", text)
|
||||
|
||||
|
|
@ -1157,26 +1226,15 @@ def _apply_disposition(ctx: ToolContext, disposition: dict) -> str:
|
|||
except PlanReviewSourceUnavailable as exc:
|
||||
return _plan_unavailable(ctx, str(exc), "plan_review_exact_artifact_unavailable")
|
||||
if not disposition.get("items") and not disposition.get("author_disposition"):
|
||||
return text
|
||||
return then_review([]) if then_review is not None else text
|
||||
cycles_paid = int(state.get("cycles_paid") or 0)
|
||||
if wave.get("closed") and not plan_review_notes_are_annotatable(wave):
|
||||
return _publish_rendered_wave(ctx, wave, cap=cap, cycles_paid=cycles_paid, enforcement=enforcement,
|
||||
cached=True,
|
||||
notes=["already_closed: this wave is closed; the disposition is not re-applied"])
|
||||
raw_items = disposition.get("items")
|
||||
if not isinstance(raw_items, list):
|
||||
return _bad("ERROR: PLAN_REVIEW_DISPOSITION_INVALID: items must be an array")
|
||||
if len(raw_items) > 2 * len(wave.get("findings") or []) + 8: # bounded like the findings they answer
|
||||
return _bad("ERROR: PLAN_REVIEW_DISPOSITION_INVALID: more items than findings could need")
|
||||
items: List[dict] = []
|
||||
for index, item in enumerate(raw_items):
|
||||
if not isinstance(item, dict):
|
||||
return _bad(f"ERROR: PLAN_REVIEW_DISPOSITION_INVALID: items[{index}] must be an object")
|
||||
items.append({
|
||||
"finding_id": str(item.get("finding_id") or "").strip()[:plan_spec.MAX_ID_CHARS * 2],
|
||||
"decision": str(item.get("decision") or "").strip().lower()[:40], # enum-like, bounded
|
||||
"rationale": plan_spec.bounded_text(item.get("rationale"), plan_spec.MAX_FINDING_TEXT_CHARS),
|
||||
})
|
||||
items, item_error = _disposition_items(wave, disposition.get("items"))
|
||||
if item_error:
|
||||
return _bad(item_error)
|
||||
author_record = None
|
||||
if disposition.get("author_disposition") is not None:
|
||||
from ouroboros.review_custody import review_retry_cancelled
|
||||
|
|
@ -1194,38 +1252,9 @@ def _apply_disposition(ctx: ToolContext, disposition: dict) -> str:
|
|||
author_record = build_author_disposition_from_mapping(disposition["author_disposition"], subject_hash=fingerprint, reviewer_signal=str(wave.get("aggregate") or ""), enforcement=enforcement)
|
||||
except ValueError as exc:
|
||||
return _bad("ERROR: PLAN_REVIEW_DISPOSITION_INVALID: " + str(exc))
|
||||
known = {str(f.get("finding_id") or "") for f in wave.get("findings") or []}
|
||||
unknown_ids = sorted({i["finding_id"] for i in items if i["finding_id"] not in known})
|
||||
if unknown_ids:
|
||||
return _bad(
|
||||
"ERROR: PLAN_REVIEW_DISPOSITION_INVALID: unknown finding ids " + ", ".join(unknown_ids)
|
||||
+ "; valid ids: " + ", ".join(sorted(known)),
|
||||
)
|
||||
closure = plan_spec.closure_after_disposition(
|
||||
str(wave.get("aggregate") or ""), wave.get("findings") or [], items, enforcement,
|
||||
)
|
||||
disposition_recorded_at = utc_now_iso()
|
||||
try:
|
||||
prior_ref = wave.get("wave_artifact") if isinstance(wave.get("wave_artifact"), dict) else {}
|
||||
closure_notes = [*closure["notes"], *([] if prior_ref else [
|
||||
"exact_artifact_absent: v2 wave had no exact wave_artifact reference",
|
||||
])]
|
||||
exact = _read_plan_review_wave_artifact(root, task_id, prior_ref) if prior_ref else dict(wave)
|
||||
exact.update({
|
||||
"dispositions": list(items), "closed": bool(closure["closed"]),
|
||||
"aggregate": closure["aggregate"], "closure_notes": closure_notes,
|
||||
"disposition_recorded_at": disposition_recorded_at,
|
||||
"supersedes_wave_artifact": prior_ref,
|
||||
})
|
||||
if author_record is not None:
|
||||
exact["author_disposition"] = author_record
|
||||
disposition_ref = _persist_plan_review_wave_artifact(root, task_id, exact)
|
||||
stored = record_plan_review_dispositions(
|
||||
root, task_id, fingerprint=fingerprint, dispositions=items,
|
||||
closed=bool(closure["closed"]), aggregate=closure["aggregate"], closure_notes=closure_notes,
|
||||
wave_artifact=disposition_ref, recorded_at=disposition_recorded_at,
|
||||
author_disposition=author_record,
|
||||
)
|
||||
stored, closure = _record_disposition(root, task_id, wave, items, fingerprint=fingerprint,
|
||||
enforcement=enforcement, author_record=author_record)
|
||||
except (OSError, TimeoutError, ValueError) as exc:
|
||||
return _typed_refusal(
|
||||
ctx, "TOOL_ERROR", "ERROR: PLAN_REVIEW_STATE_PERSIST_FAILED: " + str(exc))
|
||||
|
|
@ -1236,5 +1265,7 @@ def _apply_disposition(ctx: ToolContext, disposition: dict) -> str:
|
|||
"review closed." if closure["closed"] else
|
||||
f"{open_count} finding{'s'[:open_count != 1]} remain{'s'[:open_count == 1]} open." if open_count else
|
||||
"review stays open."))
|
||||
if then_review is not None:
|
||||
return then_review(items)
|
||||
return _publish_rendered_wave(ctx, stored, cap=cap, cycles_paid=cycles_paid,
|
||||
enforcement=enforcement, notes=list(closure["notes"]))
|
||||
|
|
|
|||
|
|
@ -803,36 +803,16 @@ def aggregate(slot_results: Iterable[Mapping[str, Any]], *, quorum: Optional[int
|
|||
return {"aggregate": verdict, "reasons": reasons, "counts": counts, "findings": flat}
|
||||
|
||||
|
||||
def blocking_fully_rejected(findings, dispositions) -> bool:
|
||||
"""Whether EVERY blocking finding carries exactly one VALID reject disposition.
|
||||
|
||||
The earned-delta admission (a fully-rejected REVISE_PLAN wave buys its promised
|
||||
delta cycle) must not trust raw items: an empty rationale, an unknown id, or a
|
||||
contradictory accept+reject pair is refused by the closure table and must not
|
||||
earn a paid panel (delta-review finding D1)."""
|
||||
blocking = [f for f in findings or [] if isinstance(f, Mapping) and f.get("class") == "blocking"]
|
||||
if not blocking:
|
||||
return False
|
||||
known = {str(f.get("finding_id") or f.get("id") or "") for f in findings or [] if isinstance(f, Mapping)}
|
||||
seen: dict[str, str] = {}
|
||||
invalid: set[str] = set()
|
||||
for item in dispositions or []:
|
||||
if not isinstance(item, Mapping):
|
||||
continue
|
||||
fid = str(item.get("finding_id") or "").strip()
|
||||
decision = str(item.get("decision") or "").strip().lower()
|
||||
ok = decision in DISPOSITION_DECISIONS and bool(str(item.get("rationale") or "").strip())
|
||||
if fid not in known:
|
||||
return False # D3: an unknown id makes the whole disposition invalid — no earned cycle
|
||||
if not ok or fid in seen:
|
||||
invalid.add(fid)
|
||||
continue
|
||||
seen[fid] = decision
|
||||
return all(
|
||||
(fid := str(f.get("finding_id") or f.get("id") or "")) not in invalid
|
||||
and seen.get(fid) == "reject"
|
||||
for f in blocking
|
||||
)
|
||||
def merge_dispositions(prior, new) -> list[dict]:
|
||||
"""One wave's answers after another call: every earlier answer survives unless this
|
||||
call answers the same ``finding_id``, which then supersedes ALL earlier entries for that
|
||||
id. Two entries for one id inside ONE call stay as written — ``closure_after_disposition``
|
||||
refuses that contradiction and the finding stays open. Non-mapping entries are dropped."""
|
||||
fresh = [dict(item) for item in new or [] if isinstance(item, Mapping)]
|
||||
answered = {str(item.get("finding_id") or "").strip() for item in fresh}
|
||||
kept = [dict(item) for item in prior or [] if isinstance(item, Mapping)
|
||||
and str(item.get("finding_id") or "").strip() not in answered]
|
||||
return kept + fresh
|
||||
|
||||
|
||||
def closure_after_disposition(
|
||||
|
|
|
|||
|
|
@ -673,7 +673,9 @@ class TestPlanReviewDispositionEnvelope(unittest.TestCase):
|
|||
run.assert_not_called()
|
||||
self.assertFalse((root / "task_results" / "parent.json").exists())
|
||||
|
||||
def test_mixed_disposition_and_plan_envelope_is_rejected_without_mutation(self):
|
||||
def test_answer_beside_an_envelope_naming_no_wave_is_refused_without_mutation(self):
|
||||
"""The valid envelope is prepared read-only, then the answers are bound to their wave:
|
||||
none holds this fingerprint, so the call is refused before any review or write."""
|
||||
import tempfile
|
||||
import ouroboros.tools.plan_review as pr
|
||||
from ouroboros.tools.registry import ToolContext
|
||||
|
|
@ -682,18 +684,14 @@ class TestPlanReviewDispositionEnvelope(unittest.TestCase):
|
|||
root = pathlib.Path(raw)
|
||||
ctx = ToolContext(repo_dir=root, drive_root=root)
|
||||
ctx.task_id = "parent"
|
||||
disposition = {
|
||||
"review_fingerprint": "f" * 64,
|
||||
"items": [{"finding_id": "slot_1:f1", "decision": "reject", "rationale": "one"}],
|
||||
}
|
||||
disposition = {"review_fingerprint": "f" * 64,
|
||||
"items": [{"finding_id": "slot_1:f1", "decision": "reject", "rationale": "one"}]}
|
||||
with patch.object(pr, "_record_raw_plan_request_with_reference") as record, patch.object(
|
||||
pr, "_run_plan_review_async",
|
||||
) as run:
|
||||
out = pr._handle_plan_task(
|
||||
ctx, plan="P changed", goal="G", spec={"in_scope": ["a"]},
|
||||
review_disposition=disposition,
|
||||
)
|
||||
self.assertIn("PLAN_REVIEW_DISPOSITION_MIXED_ENVELOPE", out)
|
||||
out = pr._handle_plan_task(ctx, plan="P changed", goal="G", spec={"in_scope": ["a"], "affected_paths": []},
|
||||
review_disposition=disposition)
|
||||
self.assertIn("PLAN_REVIEW_DISPOSITION_UNBINDABLE", out)
|
||||
record.assert_not_called()
|
||||
run.assert_not_called()
|
||||
self.assertFalse((root / "task_results" / "parent.json").exists())
|
||||
|
|
@ -723,34 +721,32 @@ class TestPlanReviewDispositionEnvelope(unittest.TestCase):
|
|||
apply_.assert_called_once_with(ctx, disposition)
|
||||
run.assert_not_called()
|
||||
|
||||
def test_meaningful_or_invalid_padding_beside_a_disposition_is_still_mixed(self):
|
||||
def test_meaningful_or_invalid_padding_beside_a_disposition_is_refused_by_the_envelope_form(self):
|
||||
"""Only schema-equivalent emptiness is ignored: a non-empty list, an unknown spec
|
||||
key or a wrong type is meaning (or an error) and keeps the typed refusal — a
|
||||
vacuity rule must never discard an invalid value to make a call pass."""
|
||||
key or a wrong type is an envelope, validated FIRST — its typed form refusal comes
|
||||
before the answers are touched; a vacuity rule never discards an invalid value."""
|
||||
import tempfile
|
||||
import ouroboros.tools.plan_review as pr
|
||||
from ouroboros.tools.registry import ToolContext
|
||||
|
||||
ctx = ToolContext(repo_dir=pathlib.Path("."), drive_root=pathlib.Path("."))
|
||||
with tempfile.TemporaryDirectory() as raw:
|
||||
root = pathlib.Path(raw)
|
||||
ctx = ToolContext(repo_dir=root, drive_root=root)
|
||||
ctx.task_id = "parent"
|
||||
disposition = {"review_fingerprint": "f" * 64, "items": []}
|
||||
for padding in (
|
||||
{"spec": {"in_scope": [""]}},
|
||||
{"spec": {"unknown": ""}},
|
||||
{"goal": []},
|
||||
{"plan": "P changed"},
|
||||
):
|
||||
with self.subTest(padding=padding):
|
||||
with patch.object(pr, "_apply_disposition") as apply_, patch.object(
|
||||
for padding in ({"spec": {"in_scope": [""]}}, {"spec": {"unknown": ""}}, {"goal": []}, {"plan": "P changed"}):
|
||||
with self.subTest(padding=padding), patch.object(pr, "_apply_disposition") as apply_, patch.object(
|
||||
pr, "_run_plan_review_async",
|
||||
) as run:
|
||||
out = pr._handle_plan_task(ctx, review_disposition=disposition, **padding)
|
||||
self.assertIn("PLAN_REVIEW_DISPOSITION_MIXED_ENVELOPE", out)
|
||||
self.assertTrue("PLAN_SPEC_INVALID" in out or "PLAN_RESOURCE_FORM_REQUIRED" in out, out)
|
||||
apply_.assert_not_called()
|
||||
run.assert_not_called()
|
||||
self.assertFalse((root / "task_results" / "parent.json").exists())
|
||||
|
||||
def test_disposition_with_an_empty_item_beside_a_plan_is_not_vacuous(self):
|
||||
"""``items=[{}]`` says something malformed, not nothing: beside a plan it is a
|
||||
mixed envelope (refused typed), never silently promoted into review mode."""
|
||||
"""``items=[{}]`` says something malformed, not nothing: beside a plan it is an answer
|
||||
with an envelope whose (malformed) form is refused typed, never plain review mode."""
|
||||
import ouroboros.tools.plan_review as pr
|
||||
from ouroboros.tools.registry import ToolContext
|
||||
|
||||
|
|
@ -761,7 +757,7 @@ class TestPlanReviewDispositionEnvelope(unittest.TestCase):
|
|||
ctx, plan="P", goal="G", spec={},
|
||||
review_disposition={"review_fingerprint": "", "items": [{}]},
|
||||
)
|
||||
self.assertIn("PLAN_REVIEW_DISPOSITION_MIXED_ENVELOPE", out)
|
||||
self.assertIn("PLAN_RESOURCE_FORM_REQUIRED", out)
|
||||
run.assert_not_called()
|
||||
|
||||
def test_state_lookup_failure_is_error_not_absence(self):
|
||||
|
|
@ -830,11 +826,16 @@ class TestPlanReviewDispositionEnvelope(unittest.TestCase):
|
|||
out = pr._handle_plan_task(ctx, plan="P", goal="G", spec={}, review_disposition=filler)
|
||||
self.assertEqual(out, "reviewed")
|
||||
run.assert_called_once()
|
||||
# A rationale is a statement; with it the disposition is real and still refused beside a plan.
|
||||
# A rationale is a statement; with it the disposition is real: beside a valid plan it is
|
||||
# bound to its wave first, and an empty fingerprint names none.
|
||||
import tempfile
|
||||
|
||||
spoken = {**filler, "author_disposition": {"disposition": "rejected", "rationale": "no"}}
|
||||
with patch.object(pr, "_run_plan_review_async") as run:
|
||||
out = pr._handle_plan_task(ctx, plan="P", goal="G", spec={}, review_disposition=spoken)
|
||||
self.assertIn("PLAN_REVIEW_DISPOSITION_MIXED_ENVELOPE", out)
|
||||
with tempfile.TemporaryDirectory() as raw, patch.object(pr, "_run_plan_review_async") as run:
|
||||
ctx = ToolContext(repo_dir=pathlib.Path(raw), drive_root=pathlib.Path(raw))
|
||||
ctx.task_id = "parent"
|
||||
out = pr._handle_plan_task(ctx, plan="P", goal="G", spec={"affected_paths": []}, review_disposition=spoken)
|
||||
self.assertIn("review_fingerprint is required", out)
|
||||
run.assert_not_called()
|
||||
|
||||
def test_duplicate_plan_calls_use_existing_sequential_tool_lane(self):
|
||||
|
|
|
|||
202
tests/test_plan_review_answer_channel.py
Normal file
202
tests/test_plan_review_answer_channel.py
Normal file
|
|
@ -0,0 +1,202 @@
|
|||
"""The plan-review ANSWER CHANNEL: answers merge by ``finding_id`` across calls (a later
|
||||
answer supersedes only its own id; a same-call duplicate stays contradictory), an author
|
||||
finish carries its answers, and an envelope sent beside answers is validated BEFORE anything
|
||||
is recorded — the answers land first, then the envelope is reviewed with them in view.
|
||||
|
||||
Driven through the real engine, the real ``task_results``/artifact store and the fake
|
||||
review substrate of ``tests.test_plan_review_engine``.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
|
||||
from ouroboros.task_results import (STATUS_RUNNING, load_plan_review_state, plan_review_wave,
|
||||
record_plan_review_dispositions, record_plan_review_wave, write_task_result)
|
||||
from ouroboros.tools import plan_review as pr
|
||||
from ouroboros.tools.plan_review_artifacts import authority_wave
|
||||
from tests.test_plan_review_engine import ( # noqa: F401
|
||||
CLEAN, DECK_SPEC, _call, _control, _finding, _state, _user_text, harness,
|
||||
)
|
||||
|
||||
|
||||
def _question(fid: str, breaks: str = "claim_1") -> dict:
|
||||
return _finding(fid, "need_evidence", breaks=breaks, summary="which one?", rec="")
|
||||
|
||||
|
||||
def _item(fid: str, decision: str = "accept", rationale: str = "the author's word") -> dict:
|
||||
return {"finding_id": fid, "decision": decision, "rationale": rationale}
|
||||
|
||||
|
||||
def _answer(ctx, fp: str, *items: dict) -> str:
|
||||
return pr._handle_plan_task(ctx, review_disposition={"review_fingerprint": fp, "items": list(items)})
|
||||
|
||||
|
||||
def _two_questions(h):
|
||||
"""A REVIEW_REQUIRED wave with two questions to the author (s1:q1, s2:q2); s3 is clean."""
|
||||
sub = h.install({"s1": json.dumps([_question("q1")]), "s2": json.dumps([_question("q2", "goal")]), "s3": CLEAN})
|
||||
ctx = h.make_ctx()
|
||||
out = _call(ctx)
|
||||
assert _control(out) == {"outcome": "REVIEW_REQUIRED", "closed": False}
|
||||
return ctx, _state(h)["waves"][-1]["request_fingerprint"], sub
|
||||
|
||||
|
||||
def _ids(wave: dict) -> list:
|
||||
return [(d["finding_id"], d["decision"]) for d in wave.get("dispositions") or []]
|
||||
|
||||
|
||||
# ------------------------------------------------------------------ item 1: merge by id
|
||||
|
||||
def test_answers_across_calls_merge_and_close_the_wave(harness): # noqa: F811
|
||||
"""Call A answers one question, call B the other: B's closure sees A's answer, and the
|
||||
hot record and the exact artifact both hold BOTH answers (before: B erased A)."""
|
||||
ctx, fp, _sub = _two_questions(harness)
|
||||
first = _answer(ctx, fp, _item("s1:q1"))
|
||||
assert _control(first) == {"outcome": "REVIEW_REQUIRED", "closed": False}
|
||||
second = _answer(ctx, fp, _item("s2:q2"))
|
||||
assert _control(second) == {"outcome": "GREEN", "closed": True}
|
||||
hot = _state(harness)["waves"][-1]
|
||||
assert _ids(hot) == [("s1:q1", "accept"), ("s2:q2", "accept")] and hot["closed"]
|
||||
exact = authority_wave(harness.drive, ctx.task_id, hot)
|
||||
assert _ids(exact) == [("s1:q1", "accept"), ("s2:q2", "accept")] and exact["closed"]
|
||||
|
||||
|
||||
def test_a_later_answer_supersedes_only_its_own_id(harness): # noqa: F811
|
||||
ctx, fp, _sub = _two_questions(harness)
|
||||
_answer(ctx, fp, _item("s1:q1"))
|
||||
_answer(ctx, fp, _item("s1:q1", "defer", "deferred to the owner"))
|
||||
assert _ids(_state(harness)["waves"][-1]) == [("s1:q1", "defer")]
|
||||
closed = _answer(ctx, fp, _item("s2:q2"))
|
||||
assert _control(closed) == {"outcome": "GREEN", "closed": True}
|
||||
assert _ids(_state(harness)["waves"][-1]) == [("s1:q1", "defer"), ("s2:q2", "accept")]
|
||||
|
||||
|
||||
def test_same_call_duplicate_stays_contradictory(harness): # noqa: F811
|
||||
"""The guard stays active: accept AND reject for one id in ONE call is refused as a
|
||||
contradiction and keeps the finding open; a later single answer then closes it."""
|
||||
ctx, fp, _sub = _two_questions(harness)
|
||||
twice = _answer(ctx, fp, _item("s1:q1", "accept"), _item("s1:q1", "reject", "no"), _item("s2:q2"))
|
||||
assert "duplicate_disposition:s1:q1" in twice
|
||||
assert _control(twice) == {"outcome": "REVIEW_REQUIRED", "closed": False}
|
||||
later = _answer(ctx, fp, _item("s1:q1"))
|
||||
assert _control(later) == {"outcome": "GREEN", "closed": True}
|
||||
assert _ids(_state(harness)["waves"][-1]) == [("s2:q2", "accept"), ("s1:q1", "accept")]
|
||||
|
||||
|
||||
def test_the_durable_writer_merges_by_finding_id(tmp_path):
|
||||
"""``record_plan_review_dispositions`` applied to prior ``[a]`` and new ``[b]`` stores
|
||||
``[a, b]``; with new ``[a']`` it stores ``[a']`` — never replace-on-write."""
|
||||
write_task_result(tmp_path, "t1", STATUS_RUNNING, result="running")
|
||||
fp = "c" * 64
|
||||
record_plan_review_wave(tmp_path, "t1", {
|
||||
"schema_version": 2, "cycle_index": 1, "request_fingerprint": fp,
|
||||
"spec": {"goal": "g", "acceptance_claims": []}, "spec_hash": "b" * 64,
|
||||
"findings": [{"finding_id": "s1:q1", "class": "need_evidence", "breaks": "goal"},
|
||||
{"finding_id": "s2:q2", "class": "need_evidence", "breaks": "goal"}],
|
||||
"aggregate": "REVIEW_REQUIRED", "closed": False, "dispositions": [], "paid": True,
|
||||
})
|
||||
a = _item("s1:q1")
|
||||
b = _item("s2:q2", "defer", "later")
|
||||
record_plan_review_dispositions(tmp_path, "t1", fingerprint=fp, dispositions=[a], closed=False)
|
||||
record_plan_review_dispositions(tmp_path, "t1", fingerprint=fp, dispositions=[b], closed=False)
|
||||
assert plan_review_wave(load_plan_review_state(tmp_path, "t1"), fp)["dispositions"] == [a, b]
|
||||
a2 = _item("s1:q1", "reject", "no")
|
||||
record_plan_review_dispositions(tmp_path, "t1", fingerprint=fp, dispositions=[a2], closed=False)
|
||||
assert plan_review_wave(load_plan_review_state(tmp_path, "t1"), fp)["dispositions"] == [b, a2]
|
||||
|
||||
|
||||
# ------------------------------------------------- item 2: answers travel with an envelope
|
||||
|
||||
def test_author_finish_with_items_records_answers_then_selects_the_plan(harness, monkeypatch): # noqa: F811
|
||||
"""Advisory: one call answers the critic's findings AND selects the corrected plan; the
|
||||
answers land on the critic wave first, merged by id, and no new panel runs."""
|
||||
harness.state["enforcement"] = "advisory"
|
||||
monkeypatch.setenv("OUROBOROS_REVIEW_ENFORCEMENT", "advisory")
|
||||
transport = harness.install({"s1": json.dumps([_finding("f1", "blocking", breaks="claim_1")]),
|
||||
"s2": CLEAN, "s3": CLEAN})
|
||||
ctx = harness.make_ctx()
|
||||
first = _call(ctx)
|
||||
assert _control(first) == {"outcome": "REVIEW_REQUIRED", "closed": False}
|
||||
critic_fp = _state(harness)["waves"][-1]["request_fingerprint"]
|
||||
spec = {**DECK_SPEC, "acceptance_claims": ["the corrected claim"]}
|
||||
result = _call(ctx, spec, plan="Corrected complete plan.", review_disposition={
|
||||
"review_fingerprint": critic_fp, "author_action": "finish",
|
||||
"items": [_item("s1:f1", "reject", "the budget line is already approved")],
|
||||
"author_disposition": {"disposition": "partial", "rationale": "Corrected the budget."}})
|
||||
assert "Current author plan saved" in result
|
||||
assert "1 answer(s) recorded on the critic wave, merged by finding_id." in result
|
||||
after = load_plan_review_state(harness.drive, ctx.task_id)
|
||||
assert len(transport.calls) == 1 and after["cycles_paid"] == 1
|
||||
critic = next(w for w in after["waves"] if w["request_fingerprint"] == critic_fp)
|
||||
assert _ids(critic) == [("s1:f1", "reject")]
|
||||
assert critic["closed"] and critic["aggregate"] == "GREEN" # advisory: a reasoned reject closes it
|
||||
assert after["current_attempt"]["author_subject"]["review_fingerprint"] == critic_fp
|
||||
|
||||
|
||||
def test_finish_while_reviewers_run_is_refused_and_records_no_answers(harness, monkeypatch): # noqa: F811
|
||||
"""The kept guard fires first: a finish while reviewers are still running is refused, and
|
||||
the items it carried are NOT written (nothing lands when the call refuses)."""
|
||||
from ouroboros import review_records
|
||||
|
||||
harness.install({"s1": json.dumps([_finding("f1", "blocking", breaks="claim_1")]), "s2": CLEAN, "s3": CLEAN})
|
||||
ctx = harness.make_ctx()
|
||||
_call(ctx)
|
||||
before = load_plan_review_state(harness.drive, ctx.task_id)
|
||||
fp = before["waves"][-1]["request_fingerprint"]
|
||||
monkeypatch.setattr(review_records, "review_outcome_received", lambda *_a, **_kw: False)
|
||||
refused = pr._handle_plan_task(ctx, review_disposition={
|
||||
"review_fingerprint": fp, "author_action": "finish",
|
||||
"items": [_item("s1:f1", "reject", "already approved")],
|
||||
"author_disposition": {"disposition": "accepted", "rationale": "Proceed."}})
|
||||
assert "PLAN_AUTHOR_SUBJECT_INVALID" in refused and "reviewers are still running" in refused
|
||||
assert load_plan_review_state(harness.drive, ctx.task_id) == before
|
||||
|
||||
|
||||
def test_author_finish_items_are_validated_against_the_critic_wave(harness): # noqa: F811
|
||||
"""An unknown finding id beside a finish is refused as a whole, before any write."""
|
||||
harness.install({"s1": json.dumps([_finding("f1", "blocking", breaks="claim_1")]), "s2": CLEAN, "s3": CLEAN})
|
||||
ctx = harness.make_ctx()
|
||||
_call(ctx)
|
||||
before = load_plan_review_state(harness.drive, ctx.task_id)
|
||||
fp = before["waves"][-1]["request_fingerprint"]
|
||||
refused = pr._handle_plan_task(ctx, review_disposition={
|
||||
"review_fingerprint": fp, "author_action": "stop", "items": [_item("s9:zz", "reject", "phantom")],
|
||||
"author_disposition": {"disposition": "partial", "rationale": "Stopping."}})
|
||||
assert "PLAN_AUTHOR_SUBJECT_INVALID" in refused and "unknown finding ids s9:zz" in refused
|
||||
assert load_plan_review_state(harness.drive, ctx.task_id) == before
|
||||
|
||||
|
||||
def test_invalid_envelope_beside_answers_records_nothing(harness): # noqa: F811
|
||||
"""Validate-first: a malformed envelope beside answers is refused by its form, and neither
|
||||
the answers nor a superseding raw attempt are written (the answered wave stays current)."""
|
||||
harness.install({"s1": json.dumps([_question("q1")]), "s2": CLEAN, "s3": CLEAN})
|
||||
ctx = harness.make_ctx()
|
||||
_call(ctx)
|
||||
before = load_plan_review_state(harness.drive, ctx.task_id)
|
||||
fp = before["waves"][-1]["request_fingerprint"]
|
||||
refused = pr._handle_plan_task(ctx, goal="Ship the deck", plan="Changed prose.",
|
||||
spec={"in_scope": ["one table"]}, # the legacy form: no affected_paths list
|
||||
review_disposition={"review_fingerprint": fp, "items": [_item("s1:q1")]})
|
||||
assert "PLAN_RESOURCE_FORM_REQUIRED" in refused
|
||||
assert load_plan_review_state(harness.drive, ctx.task_id) == before
|
||||
|
||||
|
||||
def test_changed_envelope_with_items_records_the_answers_then_reviews_every_slot(harness, monkeypatch): # noqa: F811
|
||||
"""A CHANGED envelope with items is an ordinary full wave: the answers are stored on the
|
||||
answered wave first, every slot reviews the new envelope, and the packet's PRIOR CYCLES
|
||||
section carries the rationale."""
|
||||
monkeypatch.setenv("OUROBOROS_REVIEW_MAX_CYCLES", "5")
|
||||
harness.install({"s1": json.dumps([_finding("f1", "blocking", breaks="claim_1")]), "s2": CLEAN, "s3": CLEAN})
|
||||
ctx = harness.make_ctx()
|
||||
_call(ctx)
|
||||
old_fp = _state(harness)["waves"][-1]["request_fingerprint"]
|
||||
sub2 = harness.install({"s1": CLEAN, "s2": CLEAN, "s3": CLEAN})
|
||||
result = _call(ctx, plan="A revised outline with the budget line fixed.", review_disposition={
|
||||
"review_fingerprint": old_fp, "items": [_item("s1:f1", "reject", "the budget line is already approved")]})
|
||||
assert _control(result) == {"outcome": "GREEN", "closed": True}
|
||||
assert [s.slot_id for s in sub2.calls[0]["slots"]] == ["s1", "s2", "s3"]
|
||||
packet = _user_text(sub2.calls[0]["request"].messages[1]["content"])
|
||||
assert "PRIOR CYCLES" in packet and "the budget line is already approved" in packet
|
||||
state = _state(harness)
|
||||
old = next(w for w in state["waves"] if w["request_fingerprint"] == old_fp)
|
||||
assert _ids(old) == [("s1:f1", "reject")]
|
||||
assert state["cycles_paid"] == 2 and state["waves"][-1]["request_fingerprint"] != old_fp
|
||||
|
|
@ -1212,9 +1212,11 @@ def test_session_reviewer_gets_redacted_evidence_inline_never_raw_locators(harne
|
|||
assert "plan notes" in task # the evidence text itself IS inline (redacted)
|
||||
|
||||
|
||||
def test_fully_rejected_revise_plan_wave_earns_the_promised_delta_cycle(harness, monkeypatch):
|
||||
"""Final-gate finding (scope, 4e133c8a): after a full reject-disposition of a REVISE_PLAN
|
||||
wave, re-calling the SAME envelope must buy the delta cycle, not replay forever."""
|
||||
def test_rejected_blocking_findings_replay_free_until_an_answer_is_addressed(harness, monkeypatch):
|
||||
"""No host path buys a panel the mind did not send: after a full reject-disposition of a
|
||||
REVISE_PLAN wave the identical envelope REPLAYS the recorded wave at $0 (the answers stay
|
||||
recorded on it); re-judgement is the mind's explicit move — the same envelope sent WITH
|
||||
review_disposition items (the addressed answer)."""
|
||||
monkeypatch.setenv("OUROBOROS_REVIEW_MAX_CYCLES", "5")
|
||||
blocking = json.dumps([_finding("f1", "blocking", breaks="claim_1")])
|
||||
sub = harness.install({"s1": blocking, "s2": blocking, "s3": CLEAN})
|
||||
|
|
@ -1227,64 +1229,32 @@ def test_fully_rejected_revise_plan_wave_earns_the_promised_delta_cycle(harness,
|
|||
{"finding_id": "s2:f1", "decision": "reject", "rationale": "the deadline is fine"},
|
||||
]})
|
||||
assert "revise_plan_not_closable_by_disposition" in rejected
|
||||
again = _call(ctx) # SAME envelope
|
||||
assert len(sub.calls) == 2, "the fully-rejected wave earns a paid delta cycle"
|
||||
assert "cycle 2" in again
|
||||
assert _state(harness)["cycles_paid"] == 2
|
||||
assert "PRIOR CYCLES" in _user_text(sub.calls[1]["request"].messages[1]["content"])
|
||||
|
||||
|
||||
def test_invalid_or_contradictory_rejections_do_not_earn_the_delta_cycle(harness, monkeypatch):
|
||||
"""Delta-review finding D1: only VALID rejections buy the promised delta panel."""
|
||||
monkeypatch.setenv("OUROBOROS_REVIEW_MAX_CYCLES", "5")
|
||||
blocking = json.dumps([_finding("f1", "blocking", breaks="claim_1")])
|
||||
sub = harness.install({"s1": blocking, "s2": blocking, "s3": CLEAN})
|
||||
ctx = harness.make_ctx()
|
||||
first = _call(ctx)
|
||||
assert "REVISE_PLAN" in first
|
||||
fp = _state(harness)["waves"][-1]["request_fingerprint"]
|
||||
# contradictory accept+reject for one finding, valid reject for the other
|
||||
pr._handle_plan_task(ctx, review_disposition={"review_fingerprint": fp, "items": [
|
||||
{"finding_id": "s1:f1", "decision": "accept", "rationale": "ok"},
|
||||
{"finding_id": "s1:f1", "decision": "reject", "rationale": "no"},
|
||||
{"finding_id": "s2:f1", "decision": "reject", "rationale": "the deadline is fine"},
|
||||
]})
|
||||
again = _call(ctx)
|
||||
assert len(sub.calls) == 1, "a contradictory rejection must replay, not buy a panel"
|
||||
again = _call(ctx) # SAME envelope, no items: a free replay, never a bought panel
|
||||
assert len(sub.calls) == 1 and "cached exact review" in again
|
||||
assert _state(harness)["cycles_paid"] == 1
|
||||
assert "PLAN_REVIEW_CYCLES_EXHAUSTED" not in again and "cycle 1" in again
|
||||
assert [d["finding_id"] for d in _state(harness)["waves"][-1]["dispositions"]] == ["s1:f1", "s2:f1"]
|
||||
|
||||
|
||||
def test_dispatched_degraded_delta_attempt_pays_and_becomes_current(harness, monkeypatch):
|
||||
"""B2 wave-record authority change (deliberate, supersedes delta-review D2 for
|
||||
DISPATCHED waves): the earned delta panel ran — garbage answers and all — so it
|
||||
pays its cycle and replaces the paid predecessor. The old free-retry
|
||||
preservation survives only for nothing-dispatched waves (next test)."""
|
||||
def test_dispatched_degraded_wave_pays_and_the_identical_envelope_redispatches(harness, monkeypatch):
|
||||
"""B2 wave-record authority: a panel that RAN and came back DEGRADED (garbage answers,
|
||||
not window-spent lanes) pays its cycle like any other paid wave, and — having no
|
||||
structural epoch — the identical envelope RE-DISPATCHES instead of replaying it."""
|
||||
monkeypatch.setenv("OUROBOROS_REVIEW_MAX_CYCLES", "5")
|
||||
blocking = json.dumps([_finding("f1", "blocking", breaks="claim_1")])
|
||||
harness.install({"s1": blocking, "s2": blocking, "s3": CLEAN})
|
||||
ctx = harness.make_ctx()
|
||||
_call(ctx)
|
||||
fp = _state(harness)["waves"][-1]["request_fingerprint"]
|
||||
pr._handle_plan_task(ctx, review_disposition={"review_fingerprint": fp, "items": [
|
||||
{"finding_id": "s1:f1", "decision": "reject", "rationale": "fine"},
|
||||
{"finding_id": "s2:f1", "decision": "reject", "rationale": "fine"},
|
||||
]})
|
||||
harness.install({"s1": "garbage not an array", "s2": "also garbage", "s3": "nope"})
|
||||
degraded = _call(ctx) # the earned delta panel comes back DEGRADED — but it RAN
|
||||
ctx = harness.make_ctx()
|
||||
degraded = _call(ctx)
|
||||
assert _control(degraded) == {"outcome": "DEGRADED", "closed": False}
|
||||
state = _state(harness)
|
||||
fp = state["waves"][-1]["request_fingerprint"]
|
||||
waves = [w for w in state["waves"] if w.get("request_fingerprint") == fp]
|
||||
assert len(waves) == 1 and waves[0]["aggregate"] == "DEGRADED"
|
||||
assert waves[0]["paid"] is True and not waves[0].get("degraded_retries")
|
||||
assert state["cycles_paid"] == 2, "the dispatched delta panel charged its cycle"
|
||||
# Fix 2: this DEGRADED wave has NO structural epoch (garbage answers, not
|
||||
# window-spent lanes), so the identical envelope RE-DISPATCHES, never replays.
|
||||
sub3 = harness.install({"s1": CLEAN, "s2": CLEAN, "s3": CLEAN})
|
||||
assert state["cycles_paid"] == 1, "the dispatched panel charged its cycle"
|
||||
sub2 = harness.install({"s1": CLEAN, "s2": CLEAN, "s3": CLEAN})
|
||||
fresh = _call(ctx)
|
||||
assert len(sub3.calls) == 1 and "cached exact review" not in fresh
|
||||
assert len(sub2.calls) == 1 and "cached exact review" not in fresh
|
||||
assert _control(fresh) == {"outcome": "GREEN", "closed": True}
|
||||
assert _state(harness)["cycles_paid"] == 3
|
||||
assert _state(harness)["cycles_paid"] == 2
|
||||
|
||||
|
||||
def test_nothing_dispatched_wave_stays_unpaid_and_preserves_the_paid_predecessor(tmp_path):
|
||||
|
|
@ -1356,20 +1326,6 @@ def test_engine_denies_the_runtime_data_plane_as_evidence(harness, tmp_path, mon
|
|||
assert any(str(data_root) in d for d in denied)
|
||||
|
||||
|
||||
def test_unknown_disposition_id_does_not_earn_the_delta_cycle():
|
||||
"""Delta-review D3: an unknown finding id in the disposition invalidates the earn."""
|
||||
from ouroboros.tools import plan_spec
|
||||
|
||||
findings = [{"finding_id": "s1:f1", "id": "f1", "class": "blocking", "breaks": "goal", "summary": "x"}]
|
||||
ok = plan_spec.blocking_fully_rejected(findings, [
|
||||
{"finding_id": "s1:f1", "decision": "reject", "rationale": "no"},
|
||||
{"finding_id": "s9:zz", "decision": "reject", "rationale": "phantom"},
|
||||
])
|
||||
assert ok is False
|
||||
assert plan_spec.blocking_fully_rejected(findings, [
|
||||
{"finding_id": "s1:f1", "decision": "reject", "rationale": "no"}]) is True
|
||||
|
||||
|
||||
def test_schema_conformant_clean_session_verdict_counts_as_clean():
|
||||
"""Delta-review D4 (rejected with proof): the substrate canonicalizes a schema-conformant
|
||||
`{"findings": []}` to a bare `[]`, and the repo's own clean-verdict rule already accepts a
|
||||
|
|
|
|||
|
|
@ -508,30 +508,6 @@ def test_worst_case_state_successor_receives_the_current_decision_core(tmp_path)
|
|||
assert decision_fact in preview
|
||||
|
||||
|
||||
def test_below_quorum_blocking_rejection_earns_the_promised_delta_cycle(harness, monkeypatch):
|
||||
"""R9-5: a REVIEW_REQUIRED wave carrying ONE below-quorum blocking finding cannot be closed
|
||||
by disposition (C-08); the closure table promises "the next paid delta cycle" for its
|
||||
rejection — so an identical envelope after a valid reject must RUN that paid delta panel,
|
||||
not replay the cached wave forever."""
|
||||
monkeypatch.setenv("OUROBOROS_REVIEW_MAX_CYCLES", "5")
|
||||
blocking = json.dumps([_finding("f1", "blocking", breaks="claim_1")])
|
||||
harness.install({"s1": blocking, "s2": CLEAN, "s3": CLEAN}) # 1 of 3 < quorum(2)
|
||||
ctx = harness.make_ctx()
|
||||
out = _call(ctx)
|
||||
assert _control(out) == {"outcome": "REVIEW_REQUIRED", "closed": False}
|
||||
fp = _state(harness)["waves"][-1]["request_fingerprint"]
|
||||
replay = _call(ctx) # nothing rejected yet: idempotent replay
|
||||
assert "cached" in replay.lower() and _state(harness)["cycles_paid"] == 1
|
||||
closed = pr._handle_plan_task(ctx, review_disposition={"review_fingerprint": fp, "items": [
|
||||
{"finding_id": "s1:f1", "decision": "reject", "rationale": "the visa is already granted"},
|
||||
]})
|
||||
assert _control(closed) == {"outcome": "REVIEW_REQUIRED", "closed": False}
|
||||
sub = harness.install({"s1": CLEAN, "s2": CLEAN, "s3": CLEAN})
|
||||
delta = _call(ctx) # same envelope: the rejection now buys the promised delta panel
|
||||
assert _control(delta) == {"outcome": "GREEN", "closed": True}
|
||||
assert len(sub.calls) == 1 and _state(harness)["cycles_paid"] == 2
|
||||
|
||||
|
||||
def test_disposition_inputs_are_bounded_at_entry(harness):
|
||||
"""R9-4: a disposition is bounded like the findings it answers — the rationale text and the
|
||||
item count — so a $0 closure can always be persisted."""
|
||||
|
|
|
|||
|
|
@ -1013,3 +1013,33 @@ def test_the_findings_contract_advertises_locator_forms_and_range_selectors() ->
|
|||
assert token in PLAN_FINDINGS_ARRAY_CONTRACT, token
|
||||
assert _PLAN_FINDING_ELEMENT_SCHEMA.startswith("{")
|
||||
assert _PLAN_FINDING_ELEMENT_SCHEMA.endswith("}")
|
||||
|
||||
|
||||
def test_merge_dispositions_keeps_earlier_answers_and_supersedes_only_the_answered_id() -> None:
|
||||
"""One wave's answers after another call: order is kept-then-fresh, a later answer
|
||||
replaces every earlier entry for ITS id only, and non-mapping entries are dropped."""
|
||||
from ouroboros.tools.plan_spec import merge_dispositions
|
||||
|
||||
a = {"finding_id": "s1:q1", "decision": "accept", "rationale": "yes"}
|
||||
b = {"finding_id": "s2:q2", "decision": "defer", "rationale": "later"}
|
||||
a2 = {"finding_id": "s1:q1", "decision": "reject", "rationale": "no"}
|
||||
assert merge_dispositions([], [a]) == [a]
|
||||
assert merge_dispositions([a], [b]) == [a, b]
|
||||
assert merge_dispositions([a, b], [a2]) == [b, a2]
|
||||
assert merge_dispositions([a, a2], [b]) == [a, a2, b] # an earlier contradiction is not rewritten
|
||||
assert merge_dispositions([a, "junk", None], [b, 3]) == [a, b]
|
||||
assert merge_dispositions(None, None) == []
|
||||
|
||||
|
||||
def test_merge_dispositions_keeps_a_within_call_duplicate_for_closure_to_refuse() -> None:
|
||||
"""Two entries for one id in ONE call stay as written: ``closure_after_disposition`` reads
|
||||
them as a contradiction (``duplicate_disposition``) and keeps the finding open."""
|
||||
from ouroboros.tools.plan_spec import closure_after_disposition, merge_dispositions
|
||||
|
||||
twice = [{"finding_id": "1:q1", "decision": "accept", "rationale": "yes"},
|
||||
{"finding_id": "1:q1", "decision": "reject", "rationale": "no"}]
|
||||
merged = merge_dispositions([], twice)
|
||||
assert merged == twice
|
||||
question = {"finding_id": "1:q1", "id": "q1", "class": "need_evidence", "breaks": "claim_1", "summary": "?"}
|
||||
closure = closure_after_disposition("REVIEW_REQUIRED", [question], merged, "blocking")
|
||||
assert closure["closed"] is False and "duplicate_disposition:1:q1" in closure["notes"]
|
||||
|
|
|
|||
|
|
@ -46,17 +46,13 @@ def test_a_question_holds_the_wave_until_a_free_disposition_under_both_enforceme
|
|||
assert closed["closed"] is True and closed["open_ids"] == []
|
||||
|
||||
|
||||
def test_a_question_only_wave_is_review_required_and_never_earns_a_paid_delta_cycle():
|
||||
def test_a_question_only_wave_is_review_required():
|
||||
rows = [{"slot": "s1", "model": "m", "ok": True, "findings": [
|
||||
{"id": "q1", "class": "need_evidence", "breaks": "claim_1", "locator": "", "summary": "?"}]},
|
||||
{"slot": "s2", "model": "m", "ok": True, "findings": []},
|
||||
{"slot": "s3", "model": "m", "ok": True, "findings": []}]
|
||||
agg = plan_spec.aggregate(rows)
|
||||
assert agg["aggregate"] == "REVIEW_REQUIRED" and agg["counts"]["need_evidence"] == 1
|
||||
# The earned-delta guard (no blocking finding -> False) is the only thing between a
|
||||
# question wave and a paid panel bought by rejecting the question; pinned here.
|
||||
assert plan_spec.blocking_fully_rejected(
|
||||
agg["findings"], [{"finding_id": "s1:q1", "decision": "reject", "rationale": "not needed"}]) is False
|
||||
|
||||
|
||||
def test_the_output_contract_states_both_need_evidence_forms():
|
||||
|
|
|
|||
|
|
@ -410,7 +410,7 @@ def test_plan_handler_wrapper_preserves_native_meta_for_all_projection_paths(
|
|||
],
|
||||
},
|
||||
},
|
||||
"PLAN_REVIEW_DISPOSITION_MIXED_ENVELOPE",
|
||||
"PLAN_RESOURCE_FORM_REQUIRED",
|
||||
),
|
||||
(
|
||||
{"review_disposition": {"review_fingerprint": "", "items": []}},
|
||||
|
|
@ -432,6 +432,7 @@ def test_plan_task_argument_refusals_are_typed_at_the_registry_boundary(
|
|||
from ouroboros.tools import plan_review
|
||||
|
||||
registry = ToolRegistry(repo_dir=tmp_path, drive_root=tmp_path)
|
||||
registry._ctx.task_id = "t46" # answers beside an envelope are validated against the task's state
|
||||
monkeypatch.setattr(safety, "check_safety", lambda *_args, **_kwargs: (True, ""))
|
||||
monkeypatch.setattr(
|
||||
plan_review,
|
||||
|
|
|
|||
|
|
@ -94,6 +94,8 @@ def test_advisory_author_can_select_current_plan_after_no_dispatch_outcome(harne
|
|||
|
||||
|
||||
def test_full_plan_needs_explicit_action_even_with_a_prior_wave(harness, monkeypatch): # noqa: F811
|
||||
"""A changed plan beside a stance (no author_action) is REVIEWED — the stance lands on the
|
||||
critic wave, the envelope goes to the panel — and is never SELECTED as the author plan."""
|
||||
h = harness
|
||||
h.state["enforcement"] = "advisory"
|
||||
monkeypatch.setenv("OUROBOROS_REVIEW_ENFORCEMENT", "advisory")
|
||||
|
|
@ -102,12 +104,17 @@ def test_full_plan_needs_explicit_action_even_with_a_prior_wave(harness, monkeyp
|
|||
for slot in ("s1", "s2", "s3")})
|
||||
_call(ctx)
|
||||
before = load_plan_review_state(h.drive, ctx.task_id)
|
||||
critic_fp = before["waves"][-1]["request_fingerprint"]
|
||||
result = _call(ctx, plan="Changed plan without a finish action.", review_disposition={
|
||||
"review_fingerprint": before["waves"][-1]["request_fingerprint"], "items": [],
|
||||
"review_fingerprint": critic_fp, "items": [],
|
||||
"author_disposition": {"disposition": "partial", "rationale": "This is a stance, not a finish choice."}})
|
||||
assert "PLAN_REVIEW_DISPOSITION_MIXED_ENVELOPE" in result
|
||||
assert load_plan_review_state(h.drive, ctx.task_id) == before
|
||||
assert len(transport.calls) == 1
|
||||
assert "Current author plan saved" not in result and "ERROR:" not in result
|
||||
after = load_plan_review_state(h.drive, ctx.task_id)
|
||||
assert len(transport.calls) == 2 and after["cycles_paid"] == 2
|
||||
assert after["current_attempt"]["fingerprint"] != critic_fp
|
||||
assert not after["current_attempt"].get("author_subject"), "a full plan is never selected without author_action"
|
||||
critic = next(w for w in after["waves"] if w["request_fingerprint"] == critic_fp)
|
||||
assert critic["author_disposition"]["rationale"] == "This is a stance, not a finish choice."
|
||||
|
||||
|
||||
@pytest.mark.parametrize("action", [None, "none"])
|
||||
|
|
@ -166,13 +173,12 @@ def test_neutral_action_records_finding_answers_but_no_author_finish(harness, mo
|
|||
|
||||
@pytest.mark.parametrize("field,value", [("goal", "Changed goal"), ("plan", "Changed plan"),
|
||||
("spec", {"in_scope": ["Changed scope"]})])
|
||||
def test_neutral_action_still_rejects_a_real_mixed_envelope(harness, field, value): # noqa: F811
|
||||
def test_neutral_action_beside_a_malformed_envelope_is_refused_by_its_form(harness, field, value): # noqa: F811
|
||||
ctx = harness.make_ctx()
|
||||
transport = harness.install({})
|
||||
result = pr._handle_plan_task(ctx, **{field: value}, reviewer_effort="low", review_disposition={
|
||||
"review_fingerprint": "f" * 64, "items": [], "author_action": "none"})
|
||||
assert "PLAN_REVIEW_DISPOSITION_MIXED_ENVELOPE" in result
|
||||
assert field + "=" in result and "Changed" in result
|
||||
assert "PLAN_SPEC_INVALID" in result or "PLAN_RESOURCE_FORM_REQUIRED" in result, result
|
||||
assert not transport.calls and not (harness.drive / "task_results" / (ctx.task_id + ".json")).exists()
|
||||
|
||||
|
||||
|
|
@ -244,8 +250,7 @@ def test_omitted_action_keeps_intentional_legacy_advisory_finish(harness, monkey
|
|||
assert not after["waves"][-1]["closed"] and len(transport.calls) == 1
|
||||
|
||||
|
||||
def test_mixed_envelope_error_discloses_bounded_field_previews(harness): # noqa: F811
|
||||
def test_envelope_form_refusal_beside_answers_stays_bounded(harness): # noqa: F811
|
||||
result = pr._handle_plan_task(harness.make_ctx(), plan="Long plan " * 10_000,
|
||||
review_disposition={"review_fingerprint": "f" * 64, "items": []})
|
||||
assert "PLAN_REVIEW_DISPOSITION_MIXED_ENVELOPE" in result and "plan=" in result
|
||||
assert "OMISSION NOTE" in result and len(result) < 1200
|
||||
assert "PLAN_SPEC_INVALID" in result and len(result) < 1200
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue