From 0672b6891e7e5a8962e08cf2a766f35996dba5ea Mon Sep 17 00:00:00 2001 From: Ouroboros <311266734+ouroboros-agent@users.noreply.github.com> Date: Sat, 26 Sep 2026 05:36:10 +0300 Subject: [PATCH] fix: preserve post-work interruption and explicit stop truth --- docs/architecture/06-agent-core.md | 4 +- docs/architecture/11-frozen-contracts-v1.md | 2 +- .../inventories/FROZEN_CONTRACTS_INVENTORY.md | 2 +- ouroboros/agent_task_pipeline.py | 23 ++- ouroboros/consolidator.py | 8 +- ouroboros/improvement_backlog.py | 77 ++++---- ouroboros/loop_acceptance.py | 4 + ouroboros/loop_acceptance_review.py | 5 + ouroboros/post_task_synthesis.py | 55 +++--- ouroboros/project_dialogue.py | 10 +- ouroboros/reflection.py | 8 + ouroboros/review_records.py | 18 ++ ouroboros/size_ratchet_manifest.py | 2 - ouroboros/transport_custody.py | 9 + skills/telegram/lib/telegram_quiz.py | 18 +- tests/test_acceptance_author_stop.py | 104 +++++++++++ tests/test_post_task_model_wait.py | 172 ++++++++++++++++++ tests/test_telegram_quiz_state.py | 49 +++++ web/modules/chat_decision.js | 9 +- web/modules/log_events.js | 19 +- web/tests/fixtures/outcome_phase_parity.json | 14 ++ 21 files changed, 500 insertions(+), 112 deletions(-) diff --git a/docs/architecture/06-agent-core.md b/docs/architecture/06-agent-core.md index bc317daac..cd520aae8 100644 --- a/docs/architecture/06-agent-core.md +++ b/docs/architecture/06-agent-core.md @@ -54,7 +54,7 @@ A clean criterion is evidence-resolved, not only well argued: reviewer `evidence Actionable findings enter the durable obligation dialogue with stable identity: fix, rebut with an evidence-bearing disposition, or ask the reviewers to declare the issue unreachable or a stable disagreement. A re-raise must name an existing obligation id or is disclosed as new; a valid rebuttal retires the row, an invalid one reopens it with both positions preserved, and the reviewer's `reviewer_rebuttal_response` rides into the next panel's catalog so it can tell "already answered" from "never answered". Each panel receives the bounded `acceptance_dialogue_history` OUTSIDE the hashed evidence material — so reading the history cannot mint a fresh paid binding. -Every received review outcome permits ordinary author inspection, correction or rebuttal, including the last finite panel. `loop_acceptance.merge_agent_acceptance_stance` binds `author_action=finish` to the outcome exposed in Main's request and its tool/owner-directive/evidence state; `stop` stands as recorded; queueing feedback or predeclaring is insufficient. `loop_acceptance_review._finish_advisory_author` permits informed Advisory completion on unchanged or revised work after criticism or disclosed review unavailability, even when paid capacity remains. The current author hash stays separate from the critic hash/verdict. Under Blocking, corrections may be retained and the attempt stopped, but disputed advancement still requires fresh reviewer approval; a later independently initiated task can resume under its own admission, never an automatic root or budget reset. `stop` means unfinished work and grants no approval. The earlier-revision PASS rule for final task-message delivery above remains distinct; it grants no new commit, plan or skill authority. Typed `dialogue_status` votes retain the reviewer's position: invalid votes abstain, material-free continue is `continue_without_findings`, and zero valid votes is `inconclusive`. A terminal opinion may close another critique of that case, but cannot suppress feedback or choose the author's stop. Auto and Required share these enforcement rules once eligible; explicit task-local response limits and real task rails remain. Cyber retains its own BIBLE P0/P3 authority without rewriting verdict, cost or custody. +Every received review outcome permits ordinary author inspection, correction or rebuttal, including the last finite panel. `loop_acceptance.merge_agent_acceptance_stance` binds `author_action=finish` to the outcome exposed in Main's request and its tool/owner-directive/evidence state; `stop` (unfinished work, no approval) stands, bound to no subject, until the author's next decision or owner input; queueing feedback or predeclaring is insufficient. `loop_acceptance_review._finish_advisory_author` permits informed Advisory completion on unchanged or revised work after criticism or disclosed review unavailability, even when paid capacity remains. The current author hash stays separate from the critic hash/verdict. Under Blocking, corrections may be retained and the attempt stopped, but disputed advancement still requires fresh reviewer approval; a later independently initiated task can resume under its own admission, never an automatic root or budget reset. The earlier-revision PASS rule for final task-message delivery above remains distinct; it grants no new commit, plan or skill authority. Typed `dialogue_status` votes retain the reviewer's position: invalid votes abstain, material-free continue is `continue_without_findings`, and zero valid votes is `inconclusive`. A terminal opinion may close another critique of that case, but cannot suppress feedback or choose the author's stop. Auto and Required share these enforcement rules once eligible; explicit task-local response limits and real task rails remain. Cyber retains its own BIBLE P0/P3 authority without rewriting verdict, cost or custody. Pacing predicts no review duration: a task has three host-owned rails — a deadline, a paid-cycle cap and a wallet — and a panel starts iff the cycle cap has room, the wallet can buy one work-order send per paid row, and MORE than the configured floor (`OUROBOROS_ACCEPTANCE_REVIEW_EST_SEC`, never below 200 s) remains above the finalization reserve; that floor applies only to a new critic. Author reaction uses ordinary remaining task time above the existing finalization reserve, budget, cancellation and round limits, without a reviewer-sized floor or adaptive multiplier. The two refusals remain `review_skipped_deadline_reserve` and `improvement_window_inside_reserve` for their respective owners. `task_pacing` owns both predicates (`review_launch_allowed`, `improvement_pass_allowed`); the launch rule is evaluated ONCE per panel, at loop admission before a panel is built — the paid dispatch claim (`review_dispatch.task_acceptance_paid_dispatch_stamp`) checks cancellation and the paid-cycle wallet only, and no other surface evaluates time. That claim is minted in the locked `task_results.task_acceptance_review_accounting` at first physical reviewer dispatch, and a claim without a recoverable terminal host run is UNKNOWN, never permission to re-dispatch — the double-spend fence. Once launched, a review is clamped to the owner deadline and the task ceiling with the per-send money fence: a panel the deadline cuts is a typed DEGRADED outcome, never a free skip. Disclosed residual: a panel whose evidence build consumed the margin after admission still dispatches and may be cut, and packet repair/retry sends and native rounds mean the total is NOT bounded to one floor wave. Panel durations are `task_acceptance_review_timing` telemetry (`delivery`, `deliveries`, `native_rounds`, `native_rows`) that no gate reads back; the structured review axis is mirrored as top-level `review_status`. A revision row names the pass it starts and the causes the wave recorded, never the aggregate word as its own explanation — a DEGRADED wave that still fed an improvement capsule is not the no-quorum terminal that shares the word. @@ -537,7 +537,7 @@ Only roots synthesize; `root_phase_checkpoint` makes paid synthesis at-most-once Synthesis receives a sealed final package from the durable result — the submitted final text, its artifact manifest and completion_observations. Full redacted action observations live in the canonical artifact store (`task.budget_drive_root or drive_root`), in the write-once `source_handles/context_checkpoints` store with verified `task_source` refs, before compact publication and outside deliverables and inferred readiness; their native reader `get_task_result(include_completion_source=true)` returns complete length/hash first, then explicit `source_start_char`/`source_end_char` ranges (`artifacts.text_source_range_projection`, the shared work-order range contract), with bytes, kind, path containment and SHA checked before any excerpt. Packet-only reflection receives per-send-tool counts, each family's latest recorded return, and task-related skill readiness with coverage; full-source references are for later readers, not evidence the synthesizer has read. Positive observed facts correct error-trace impressions, while tool success does not prove owner receipt, empty material does not prove absence, and skill readiness does not attribute an owner's action to the task. Before context cleanup, `agent_task_pipeline.emit_task_results` also freezes `review_evidence.task_inputs` through `post_task_synthesis.capture_task_inputs`: `run_origin`, the existing task-local owner corpus, intact question/answer provenance and the canonical split-root verification-receipt union. Reflection receives the same complete redacted content through `reflection.task_inputs_prompt_section`, separate from bounded trace/review excerpts. A zero return code is positive evidence; an unrelated later pass cannot resolve another check's failure. Peer proposals stay attributed, and unavailable input is not evidence that approval or verification never existed. Recovery uses these stored observations and inputs, not a later conversation. A free `host_task_facts` row precedes paid stages (or follows result persistence when Stop skips them): no model call or narrative; its metrics, routing and cost serve history. `files_rescued` (TZ-2 C2) is a stat-only file count of the artifact stores: positive, zero or unknown if unreadable, `hash_computed: false`; the stop receipt repeats it so no salvageable text never means no files. -Pooled workers retain their slot until post-task work settles; early final-answer delivery is independent of that timing. Native work stays on its registered actor; `TaskModelWait` remains reachable through `POST_TASK_SYNTHESIS_INFLIGHT`, and detached work owns a separate live wait. Mailbox cleanup waits for the checkpoint. The solve-phase absolute ceiling does not cut a settled root's running post-work, but Stop, calendar deadline, monetary admission, per-call and idle rails still bind. The stage owner reads returned memory errors: budget or unresolved provider attempts skip later paid stages and degrade the checkpoint; an ordinary failure stays local — later stages still run, completed reflection actions still apply — but the stage that lost work is unfinished, so the checkpoint reads `degraded` with no skipped list, never `completed` (TZ-2 C3). A running checkpoint after restart degrades without replaying a paid request. +Pooled workers retain their slot until post-task work settles; early final-answer delivery is independent of that timing. Native work stays on its registered actor; `TaskModelWait` remains reachable through `POST_TASK_SYNTHESIS_INFLIGHT`, and detached work owns a separate live wait. Mailbox cleanup waits for the checkpoint. The solve-phase absolute ceiling does not cut a settled root's running post-work, but Stop, calendar deadline, monetary admission, per-call and idle rails still bind. The stage owner reads returned and raised failures alike: budget, a control or an unresolved attempt on any provider's chain (`transport_custody.outcome_unknown_on_chain`) skips later paid stages, degrading the checkpoint; an ordinary failure stays local (later stages run, completed reflection actions apply), but its stage lost work: `degraded` with no skipped list, never `completed` (TZ-2 C3). A running checkpoint after restart degrades without replaying a paid request. #### Project registry and lease diff --git a/docs/architecture/11-frozen-contracts-v1.md b/docs/architecture/11-frozen-contracts-v1.md index 2f1c19085..2b2f333a6 100644 --- a/docs/architecture/11-frozen-contracts-v1.md +++ b/docs/architecture/11-frozen-contracts-v1.md @@ -28,7 +28,7 @@ This chapter owns the ABI promise: which typed shapes and their parsing, normali | `ChatOutbound.initiator` — additive origin label of a self-initiated turn (`"consciousness"` on every frame and chat/progress/summary row of a consciousness wake-up and of the roots it starts; absent on an owner's turn); stamped by the turn's own event queue and the agent's frame meta, persisted by `log_chat`/the authored summary row/the task result, replayed by history on each row | `ouroboros/gateway/contracts.py`, `supervisor/log_addressing.py`, `ouroboros/subagent_messages.py`, `supervisor/message_bus.py`, `ouroboros/gateway/history.py`, `web/modules/api_types.js` | `tests/test_consciousness_initiator_label.py`, `tests/test_consciousness_wake_lane.py`, `web/tests/consciousness_label.test.js` | | `project_thread` stamp on all seven outbound frame types, stamped at the message-bus broadcast choke; a stamped frame is never adopted by Main (`chat_activity.mainThreadAccepts`) | `supervisor/message_bus.py`, `ouroboros/projects_registry.py` | `tests/test_message_bus.py`, `web/tests/chat_thread_routing.test.js` | | Media/link envelopes — media `task_id`/`size_bytes`/`download_url`; `LinkAction {label,url}` with at most twelve absolute HTTP(S) actions; `links` in `WS_MESSAGE_TYPES`; `chat.links` host topic | `ouroboros/gateway/contracts.py`, `ouroboros/tools/core.py`, `ouroboros/event_bus.py` | `tests/test_contracts.py` | -| Owner quiz ABI — `QuizOption {label, detail?, recommended?}` (the asker marks its recommendation on that option; the web card badges it, Telegram stars its button, the durable block keeps `recommended_index`), `QuizOutbound` (quiz_id, question, options (0–6; empty is an open question answered in the owner's words), stake, `assumption` (required for optional clarification), additive `wait_for_answer` for a live pooled or ordinary-conversation root that must wait, lifecycle state open/answered/expired_terminal/superseded, optional host-written `host_facts` sentence, also on history rows and Main pointers), separate `QuizStateOutbound` discriminator, `chat.quiz` host topic (+ optional event-only `project_name`); the producer is the one escalation verb `escalate(question, options, stake, assumption, wait_for_answer=False, max_wait_minutes=None)` (the bound applies to a required wait only, never past the task's own deadline: the wait resumes with a system notice and the card stays open; named on an optional question it takes the omitted path and the asker's receipt says so) — a ROOT asks the owner, a SUBAGENT delivers a typed frame to its nearest LIVE ancestor, which answers via `forward_to_worker` or escalates verbatim, so the owner sees only what no ancestor answered; answers arrive through the ONE ingress `POST /api/decisions` (family ids `quiz:{task_id}:{quiz_id}`, `routing:{client_message_id}:{routing_token}`; `interaction:` reserved), request-id idempotent, first answer wins, validated against the STORED options; `option_index` is optional for the quiz family alone — a comment-only answer writes NO `answered_index`, because a stored 0 would replay as "chose the first option"; injected as the typed `KIND_QUIZ_ANSWER` mailbox control and broadcast as `quiz_state` (carrying the recorded `comment` when the owner answered in their own words, so the live card shows `Owner's answer:` exactly as replay does); expiry is structural only (the task-done seam flips open quizzes to `expired_terminal`, and the SAME reconcile closes the paired `owner_wait` so a terminal task never projects `quiz=expired_terminal` beside `owner_wait=waiting`; `owner_wait.set_owner_wait`'s refusal to continue waiting on a terminal result is preserved, not caught); a LATE answer to such an expired card is nevertheless ACCEPTED at the same ingress (the projection records `answered_after_terminal`) and, because no mailbox will ever be drained, is delivered as the owner's OWN message into the card's chat through the named ingress `supervisor.message_bus.accept_local_message`, idempotent on `client_message_id = quiz_late_answer::`, its provenance in the message's own `late_answer` metadata rather than a substituted `client_surface`; the 2xx says `forwarded` so no surface claims a delivery that did not happen, and 409 is left for what is genuinely settled (an already answered card, a non-root addressee). History replay merges the projection state | `ouroboros/gateway/contracts.py`, `ouroboros/gateway/task_decision.py`, `ouroboros/owner_quiz.py`, `ouroboros/tools/core.py` | `tests/test_gateway_parity.py`, `tests/test_quiz_display.py`, `tests/test_quiz_answer.py`, `web/tests/chat_decision.test.js` | +| Owner quiz ABI — `QuizOption {label, detail?, recommended?}` (the asker marks its recommendation on that option; the web card badges it, Telegram stars its button, the durable block keeps `recommended_index`), `QuizOutbound` (quiz_id, question, options (0–6; empty is an open question answered in the owner's words), stake, `assumption` (required for optional clarification), additive `wait_for_answer` for a live pooled or ordinary-conversation root that must wait, lifecycle state open/answered/expired_terminal/superseded, optional host-written `host_facts` sentence, also on history rows and Main pointers), separate `QuizStateOutbound` discriminator, `chat.quiz` host topic (+ optional event-only `project_name`); the producer is the one escalation verb `escalate(question, options, stake, assumption, wait_for_answer=False, max_wait_minutes=None)` (the bound applies to a required wait only, never past the task's own deadline: the wait resumes with a system notice and the card stays open; named on an optional question it takes the omitted path and the asker's receipt says so) — a ROOT asks the owner, a SUBAGENT delivers a typed frame to its nearest LIVE ancestor, which answers via `forward_to_worker` or escalates verbatim, so the owner sees only what no ancestor answered; answers arrive through the ONE ingress `POST /api/decisions` (family ids `quiz:{task_id}:{quiz_id}`, `routing:{client_message_id}:{routing_token}`; `interaction:` reserved), request-id idempotent, first answer wins, validated against the STORED options; `option_index` is optional for the quiz family alone — a comment-only answer writes NO `answered_index`, because a stored 0 would replay as "chose the first option"; injected as the typed `KIND_QUIZ_ANSWER` mailbox control and broadcast as `quiz_state` (carrying the recorded `comment` when the owner answered in their own words, so the live card shows `Owner's answer:` exactly as replay does); expiry is structural only (the task-done seam flips open quizzes to `expired_terminal`, and the SAME reconcile closes the paired `owner_wait` so a terminal task never projects `quiz=expired_terminal` beside `owner_wait=waiting`; `owner_wait.set_owner_wait`'s refusal to continue waiting on a terminal result is preserved, not caught); a LATE answer to such an expired card is ACCEPTED at the same ingress (the projection records `answered_after_terminal`) and, because no mailbox will ever be drained, is delivered as the owner's OWN message into the card's chat through the named ingress `supervisor.message_bus.accept_local_message`, idempotent on `client_message_id = quiz_late_answer::`, its provenance in the message's own `late_answer` metadata rather than a substituted `client_surface`; the 2xx says `forwarded` so no surface claims a delivery that did not happen, and 409 is left for what is settled (an already answered card, a non-root addressee). History replay merges the projection state | `ouroboros/gateway/contracts.py`, `ouroboros/gateway/task_decision.py`, `ouroboros/owner_quiz.py`, `ouroboros/tools/core.py` | `tests/test_gateway_parity.py`, `tests/test_quiz_display.py`, `tests/test_quiz_answer.py`, `web/tests/chat_decision.test.js` | | Managed update ABI — preflight, `UpdateMergePlan`, pinned apply, process-local `update_progress`, `update_progress_changed` invalidation and boot-only `update_status_ready` | `ouroboros/gateway/contracts.py` | `tests/test_update_apply_routing.py` | | `ChatOutbound.review_projection` — bounded actor findings via `utils.truncate_review_artifact`, at most `MAX_PROJECTED_ACTOR_FINDINGS` rows (`review_execution_projection.py`) | `ouroboros/gateway/contracts.py` | `tests/test_review_substrate_v2.py`, `web/tests/review_truth.test.js` | | Skill preflight statuses — `preflight_failed` is fresh-only; a stale failure surfaces as `preflight_failed_stale`; absence means the caller could not know | `ouroboros/skill_review_status.py` | `tests/test_skill_preflight_repair.py`, `web/tests/skill_preflight_repair.test.js` | diff --git a/docs/inventories/FROZEN_CONTRACTS_INVENTORY.md b/docs/inventories/FROZEN_CONTRACTS_INVENTORY.md index 8f2c9f71e..5e9dc0750 100644 --- a/docs/inventories/FROZEN_CONTRACTS_INVENTORY.md +++ b/docs/inventories/FROZEN_CONTRACTS_INVENTORY.md @@ -2,7 +2,7 @@ Machine extraction of `docs/ARCHITECTURE.md` §11.1 (the frozen-ABI SSOT), regenerated by `python scripts/regenerate_inventories.py`. Do not edit — edit the owning chapter named in the Source line and regenerate; `tests/test_generated_inventories.py` pins byte-identity and the resolution invariants (a §11.1 row whose owner or anchor file disappeared from the tree = red). -Source: `docs/architecture/11-frozen-contracts-v1.md`, physical LF lines 7-40; UTF-8 SHA-256 `9c0c9c16e9d2ae8445b9f731a88b24f4258d221a67309cade694c8945091de33`. +Source: `docs/architecture/11-frozen-contracts-v1.md`, physical LF lines 7-40; UTF-8 SHA-256 `2aa8ff73bf8a7723250c834d62a91bd1f305583b210cea4ab324839d21250042`. - table rows: **29** - browser-envelope prose owners: diff --git a/ouroboros/agent_task_pipeline.py b/ouroboros/agent_task_pipeline.py index 9ce544c9f..448324df4 100644 --- a/ouroboros/agent_task_pipeline.py +++ b/ouroboros/agent_task_pipeline.py @@ -42,7 +42,6 @@ from ouroboros.subagents import envelope_from_task, substrate_result_fields from ouroboros.subagent_messages import initiator_meta, subagent_message_meta from ouroboros.utils import utc_now_iso, append_jsonl, truncate_review_artifact as _truncate_with_notice from ouroboros.utils import in_worker_process -from ouroboros.llm_claudexor import propagate_model_error from ouroboros.post_task_checkpoint import ( POST_TASK_SYNTHESIS_INFLIGHT as _POST_TASK_SYNTHESIS_INFLIGHT, POST_TASK_SYNTHESIS_LOCK as _POST_TASK_SYNTHESIS_LOCK, @@ -201,8 +200,15 @@ def _run_post_task_processing_async( if is_presence_task(task_snapshot): return "" # Project facts stay scoped; generic process lessons remain global. - _update_improvement_backlog(env, reflection_entry) failure = "" + try: + _update_improvement_backlog(env, reflection_entry) + except Exception as error: + propagate_paid_interruption(error) + # A lost append or grooming pass is this stage's own typed + # failure (degraded, nothing skipped); the chooser still runs. + failure = "backlog_update_failed" + log.warning("Improvement backlog update failed", exc_info=True) try: from ouroboros.post_task_evolution import maybe_promote @@ -258,13 +264,14 @@ def _run_post_task_processing_async( stage_errors = True log.warning("Post-task stage %s failed for %s: %s", name, stage_task_id, stage_reason) except Exception as error: - if isinstance(error, BudgetExceeded): + # The adapters' own classifier: a control, the wallet and an + # unresolved attempt on any provider's chain stop later paid work. + try: + propagate_paid_interruption(error) + except BudgetExceeded: interrupted = "budget_exhausted" - else: - try: - propagate_model_error(error) - except Exception as control: - interrupted = post_task_interruption(control) + except Exception as control: + interrupted = post_task_interruption(control) if interrupted: skipped = [stage for stage, _run in stages[index + 1:]] log.warning("Post-task paid stage %s interrupted for %s: %s", diff --git a/ouroboros/consolidator.py b/ouroboros/consolidator.py index 46faf192e..9e9f4d982 100644 --- a/ouroboros/consolidator.py +++ b/ouroboros/consolidator.py @@ -938,12 +938,10 @@ def _call_consolidation_llm( break except Exception as error: from ouroboros.llm_claudexor import propagate_model_error + from ouroboros.loop_llm_call import classify_llm_exception + from ouroboros.transport_custody import outcome_unknown_on_chain from ouroboros.usage_accounting import BudgetExceeded propagate_model_error(error) - from ouroboros.loop_llm_call import classify_llm_exception - from ouroboros.transport_custody import _capture_on_chain - - capture = _capture_on_chain(error) if getattr(error, "route", None): # A refusal belongs to the actual account, which can differ from # catalog discovery. Rebind its facts without masking the refusal @@ -952,7 +950,7 @@ def _call_consolidation_llm( preflight = isinstance(error, SummarizerContextOverflow) or not invoked kind = ("budget_exhausted" if isinstance(error, BudgetExceeded) else "context_overflow" if isinstance(error, SummarizerContextOverflow) - else "provider_outcome_unknown" if getattr(capture, "state", "") in {"dispatched", "unresolved"} + else "provider_outcome_unknown" if outcome_unknown_on_chain(error) else classify_llm_exception(error).kind) message = str(error) usage = dict(getattr(error, "usage", None) or {}) diff --git a/ouroboros/improvement_backlog.py b/ouroboros/improvement_backlog.py index d45207bdf..5a559be7c 100644 --- a/ouroboros/improvement_backlog.py +++ b/ouroboros/improvement_backlog.py @@ -533,7 +533,8 @@ def groom_backlog(drive_root: Any, *, cap: int = _GROOM_CAP) -> int: merge near-dupes, mark resolved, re-rank, cap to <=cap. Hand-added items and fingerprinted blocks with unmodelled bytes pass through UNCHANGED. Runs on a size trigger and re-serializes through the locked parser-safe writer. Returns - the number written, or 0 on a bad/empty/oversized reply or concurrent change.""" + the number written, or 0 on a bad/empty/oversized reply or concurrent change; + a failed model call raises to the caller (never the 0 of a pass not needed).""" import json as _json path = backlog_path(drive_root) @@ -558,46 +559,42 @@ def groom_backlog(drive_root: Any, *, cap: int = _GROOM_CAP) -> int: if not fp_items: return 0 + from ouroboros.config import get_light_model + from ouroboros.llm import LLMClient + from ouroboros.llm_observability import chat_observed + + # One destructive grooming call sees every complete stored record, + # including evidence/context/custom fields and immutable manual items. + # If this full prompt cannot be served, the raised failure preserves the + # file; there is no smaller second call or prefix-authorized rewrite. The + # caller's stage classifies it: an interruption stops later paid post-work. + complete = [dict(it) for it in items] + prompt = _GROOM_PROMPT.format(cap=cap, items_json=_json.dumps(complete, ensure_ascii=False)) + resp, usage = chat_observed( + LLMClient(), + drive_root=pathlib.Path(drive_root), + task_id="backlog_groom", + call_type="backlog_groom", + model_role="light", + messages=[{"role": "user", "content": prompt}], + model=get_light_model(), + reasoning_effort="low", + max_tokens=8192, + ) + if usage: + try: + from supervisor.state import update_budget_from_usage + + update_budget_from_usage(usage) + except Exception: + pass + content = (resp.get("content") or "").strip() + start, end = content.find("["), content.rfind("]") try: - from ouroboros.config import get_light_model - from ouroboros.llm import LLMClient - from ouroboros.llm_observability import chat_observed - - # One destructive grooming call sees every complete stored record, - # including evidence/context/custom fields and immutable manual items. - # If this full prompt cannot be served, the exception path preserves the - # file; there is no smaller second call or prefix-authorized rewrite. - complete = [dict(it) for it in items] - prompt = _GROOM_PROMPT.format(cap=cap, items_json=_json.dumps(complete, ensure_ascii=False)) - client = LLMClient() - resp, usage = chat_observed( - client, - drive_root=pathlib.Path(drive_root), - task_id="backlog_groom", - call_type="backlog_groom", - model_role="light", - messages=[{"role": "user", "content": prompt}], - model=get_light_model(), - reasoning_effort="low", - max_tokens=8192, - ) - if usage: - try: - from supervisor.state import update_budget_from_usage - - update_budget_from_usage(usage) - except Exception: - pass - content = (resp.get("content") or "").strip() - start, end = content.find("["), content.rfind("]") - if start < 0 or end <= start: - return 0 - kept_raw = _json.loads(content[start:end + 1]) - if not isinstance(kept_raw, list): - return 0 - except Exception as exc: - from ouroboros.post_task_synthesis import propagate_paid_interruption - propagate_paid_interruption(exc) # budget, unknown outcome and controls stop later paid post-work + kept_raw = _json.loads(content[start:end + 1]) if 0 <= start < end else None + except ValueError: + return 0 # a bad reply, like an empty one: the file is preserved + if not isinstance(kept_raw, list): return 0 # Anti-wipe: every kept item MUST map to an existing fingerprinted item — the diff --git a/ouroboros/loop_acceptance.py b/ouroboros/loop_acceptance.py index 93443fb95..052131898 100644 --- a/ouroboros/loop_acceptance.py +++ b/ouroboros/loop_acceptance.py @@ -738,6 +738,10 @@ def merge_agent_acceptance_stance(trace: Dict[str, Any], decision: dict, ctx: An # fresh read, so an unverifiable stance is never honoured as ready. "evidence_fingerprint": observed_delivery_evidence(ctx, trace), } + from ouroboros.review_records import recorded_author_stop + + if merged.get("agent_finish_intent") and recorded_author_stop(previous): + ctx._task_acceptance_reviewed = False # the author's next decision reopens an honoured stop trace["acceptance_decision"] = merged diff --git a/ouroboros/loop_acceptance_review.py b/ouroboros/loop_acceptance_review.py index 728a8a87c..461f1e2f3 100644 --- a/ouroboros/loop_acceptance_review.py +++ b/ouroboros/loop_acceptance_review.py @@ -672,6 +672,11 @@ def _finish_advisory_author(ctx: _TaskAcceptanceContext) -> bool: _loop()._supersede_task_acceptance_for_owner_followup(ctx.tools._ctx, ctx.llm_trace) return True ctx.tools._ctx._task_acceptance_reviewed = True + if action == "stop": + # A stop binds no subject, so an earlier panel's cannot reopen review on the next + # delivery pass; only the author's next decision (merge_agent_acceptance_stance) + # or owner input does. + ctx.tools._ctx._task_acceptance_reviewed_subject = "" ctx.tools._ctx._task_acceptance_pending = "" _loop()._mark_root_acceptance_checkpoint( ctx.tools._ctx, ctx.llm_trace, status=author["reviewer_signal"].lower(), pass_index=ctx.passes_done, diff --git a/ouroboros/post_task_synthesis.py b/ouroboros/post_task_synthesis.py index e78ade877..794afa9a7 100644 --- a/ouroboros/post_task_synthesis.py +++ b/ouroboros/post_task_synthesis.py @@ -273,26 +273,20 @@ def _update_improvement_backlog( env: Any, reflection_entry: Dict[str, Any] | None, ) -> int: - """Persist LLM-nominated follow-up improvements into the durable backlog.""" - try: - from ouroboros.improvement_backlog import append_backlog_items + """Persist LLM-nominated follow-up improvements into the durable backlog. - candidates = list((reflection_entry or {}).get("backlog_candidates") or []) - if not candidates: - return 0 - added = append_backlog_items(env.drive_root, candidates) - try: - from ouroboros.improvement_backlog import groom_backlog + Returns the number appended; 0 is a genuine no-op (nothing nominated), never a + swallowed failure. An append or grooming failure raises to the promotion stage, + which isolates an ordinary one and stops later paid work on an interruption. + """ + from ouroboros.improvement_backlog import append_backlog_items, groom_backlog - groom_backlog(env.drive_root) # size-triggered; no-op while small - except Exception as error: - propagate_paid_interruption(error) - log.debug("Backlog grooming failed", exc_info=True) - return added - except Exception as error: - propagate_paid_interruption(error) - log.debug("Improvement backlog update failed", exc_info=True) + candidates = list((reflection_entry or {}).get("backlog_candidates") or []) + if not candidates: return 0 + added = append_backlog_items(env.drive_root, candidates) + groom_backlog(env.drive_root) # size-triggered; no-op while small + return added def _apply_reflection_memory_actions( @@ -524,8 +518,12 @@ def _record_task_facts(env: Any, task: Dict[str, Any], usage: Dict[str, Any], stored_result = _atp().load_task_result(result_root, task_id) or {} review_projection = _compact_review_projection(llm_trace) # TZ-2 C2: how many files the task rescued into its store(s) — positive, zero or - # unknown — by stat alone; the fact discloses that no hash was computed. - files_rescued = rescued_files_fact(task_id, artifact_store_roots(canonical_root, task_id, child_root=result_root)) + # unknown — by stat alone; the fact discloses that no hash was computed. A split + # non-Project root synthesizes on the canonical drive (parent env and task): its + # actor store is then the row's recorded ``child_drive_root``, else this drive. + canonical = result_root.resolve(strict=False) == canonical_root.resolve(strict=False) + files_rescued = rescued_files_fact(task_id, artifact_store_roots( + canonical_root, task_id, task=task, child_root=None if canonical else result_root)) append_canonical_task_summary(canonical_root, { "ts": utc_now_iso(), "direction": "system", "type": "task_summary", "summary_kind": "host_task_facts", "summary_id": f"task-facts:{task_id}", @@ -555,15 +553,18 @@ POST_TASK_INTERRUPT_KINDS = frozenset({"budget_exhausted", "provider_outcome_unk def propagate_paid_interruption(error: BaseException) -> None: """Re-raise what must stop later paid post-work; return for an ordinary failure. - ``propagate_model_error`` carries the control and unknown-provider facts; the - wallet's ``BudgetExceeded`` is the third (TZ-2 C3). A stage adapter that only - logged it let the coordinator run the next paid stage and write ``completed``. - Anything else returns, so the caller isolates the failure to its own stage. + ``propagate_model_error`` carries the control and typed unknown-provider facts; + the wallet's ``BudgetExceeded`` and an unresolved attempt on ANY provider's + exception chain — the consolidator's own classifier, never a broadened global + one — are the others (TZ-2 C3). A stage adapter that only logged them let the + coordinator run the next paid stage and write ``completed``. Anything else + returns, so the caller isolates the failure to its own stage. """ propagate_model_error(error) + from ouroboros.transport_custody import outcome_unknown_on_chain from ouroboros.usage_accounting import BudgetExceeded - if isinstance(error, BudgetExceeded): + if isinstance(error, BudgetExceeded) or outcome_unknown_on_chain(error): raise error @@ -579,9 +580,9 @@ def post_task_interruption(control: BaseException) -> str: if reason: return reason code = str(getattr(control, "code", "") or "") - capture = getattr(control, "physical_attempt_capture", None) - if (not code or code == "model_outcome_unknown" - or getattr(capture, "state", None) in {"dispatched", "unresolved"}): + from ouroboros.transport_custody import outcome_unknown_on_chain + + if not code or code == "model_outcome_unknown" or outcome_unknown_on_chain(control): return "provider_outcome_unknown" return code diff --git a/ouroboros/project_dialogue.py b/ouroboros/project_dialogue.py index 6a4273020..2bb64d2bf 100644 --- a/ouroboros/project_dialogue.py +++ b/ouroboros/project_dialogue.py @@ -20,6 +20,7 @@ from typing import Any, Callable, Dict, Iterable, List, Optional from ouroboros.acceptance_preparation import incident_cause_clauses from ouroboros.platform_layer import acquire_exclusive_file_lock, release_exclusive_file_lock +from ouroboros.review_records import recorded_author_stop from ouroboros.task_finalization import TERMINAL_ORIGIN_HOST_SALVAGE from ouroboros.utils import append_jsonl, iter_jsonl_objects, jsonl_append_lock_path, replace_atomic, strip_markdown, utc_now_iso @@ -1307,9 +1308,10 @@ def _completion_verdict(result: Dict[str, Any], event: Dict[str, Any]) -> str: if raw_reason == REASON_OWNER_REQUESTED_FINALIZATION: clause = "" # an owner-requested stop is a success and carries its own marker elif (status and (status != ACCEPTANCE_ACCEPTED or cause in TASK_CAUSE_PHRASES) - and (phase in {"done", "warn"} or (cause == "author_stop" and phase == "error"))): + and (phase in {"done", "warn"} or (recorded_author_stop(decision) and phase == "error"))): # An explicit author stop is the fact that ended the task (its objective is - # blocked, so the card is red); the typed sentence speaks over the delivery step. + # blocked, so the card is red); the decision's typed reason — its TRUE cause, + # not always ``author_stop`` — speaks over the delivery step. clause = TASK_CAUSE_PHRASES.get(cause, cause) elif phase == "cancelled" and isinstance(origin, dict) and origin: # The recorded cause and the relation the record PROVES (#1061). @@ -1341,8 +1343,8 @@ def _author_stop_rationale(decision: Dict[str, Any]) -> str: """The agent's own reason for an explicit stop, beside the typed sentence (TZ-2 C4). Only the AUTHOR's recorded rationale reaches the row: the reviewer rationale - stays in the card. The twin of ``authorStopRationale``.""" - author = decision.get("author_disposition") if decision.get("reason") == "author_stop" else None + stays in the card; a finish carries none. The twin of ``authorStopRationale``.""" + author = decision.get("author_disposition") if recorded_author_stop(decision) else None return " ".join(strip_markdown(str(author.get("rationale") or "")).split()) if isinstance(author, dict) else "" diff --git a/ouroboros/reflection.py b/ouroboros/reflection.py index f20c6db20..8202cd204 100644 --- a/ouroboros/reflection.py +++ b/ouroboros/reflection.py @@ -660,6 +660,14 @@ def generate_reflection( backlog_candidates = [] memory_actions = [] reflection_route = "unknown" + # The placeholder is a stage that lost its work, never a clean one: the + # post-task coordinator reads this typed row and degrades the checkpoint + # while later stages still run; an interruption row keeps precedence there. + from ouroboros.utils import sanitize_tool_result_for_log + + memory_operation_errors = [*memory_operation_errors, { + "kind": "reflection_failed", "label": "Task reflection", + "message": sanitize_tool_result_for_log(str(e)) or type(e).__name__}] return { "ts": utc_now_iso(), diff --git a/ouroboros/review_records.py b/ouroboros/review_records.py index 931d015df..0de0e35f0 100644 --- a/ouroboros/review_records.py +++ b/ouroboros/review_records.py @@ -111,6 +111,24 @@ def validate_author_disposition( return normalized +def recorded_author_stop(decision: Any) -> bool: + """Whether a task acceptance decision records the author's explicit stop (TZ-2 C4). + + The typed ``author_stop`` reason, or the structured stop the producer records + under its TRUE terminal cause when the review rounds ran out: ``author_action`` + and the disposition's ``action`` are both ``stop``. A finish is never a stop. + The twin of ``log_events.explicitAuthorStop``. + """ + from ouroboros.outcomes import REASON_REVIEW_CYCLES_EXHAUSTED + + if not isinstance(decision, dict): + return False + author = decision.get("author_disposition") + return decision.get("reason") == "author_stop" or ( + decision.get("reason") == REASON_REVIEW_CYCLES_EXHAUSTED and decision.get("author_action") == "stop" + and isinstance(author, dict) and author.get("action") == "stop") + + def build_author_disposition_from_mapping( value: Any, *, subject_hash: str, reviewer_signal: str = "", enforcement: str = "", ) -> Dict[str, Any]: diff --git a/ouroboros/size_ratchet_manifest.py b/ouroboros/size_ratchet_manifest.py index 00b45999c..d0a00e979 100644 --- a/ouroboros/size_ratchet_manifest.py +++ b/ouroboros/size_ratchet_manifest.py @@ -114,8 +114,6 @@ BAND_PATHS = { "ouroboros/claudexor_daemon.py": "Installation daemon lifecycle owns marker and authenticated endpoint stop authority, confirmed self-started handles, and duplicate-start refusal; process signal and ledger mechanics remain in process_custody. No new lifecycle store or scheduler.", "ouroboros/claudexor_runtime.py": "Exact byte verification and delivery now have a shared owner for engine and skill resources; this module retains engine pin, archive installation and platform-specific contracts.", "ouroboros/cli.py": "The existing command-line transport keeps task-event negotiation, bounded replay deduplication and result rendering together; the additive cursor does not introduce a second CLI or task engine.", - "ouroboros/consolidator.py": "Shrunk from the 1501-1600 giant band after per-room consolidation moved room draft/correction into room_consolidation.py; the block/era orchestration, chunk atomicity and route-fit machinery still share this owner.", - "ouroboros/context.py": "Entered the band from the 1501-1600 zone (1590 lines) by the v7 D03 extraction of the runtime-section fact builders into ouroboros/context_runtime_facts.py; shrink-only residue of the split, not new growth.", "ouroboros/context_compaction.py": "Existing compaction owns propagation of typed model outcomes; unchanged semantic compaction policy.", "ouroboros/delegate_custody.py": "D07 DEL1 split brought the custody monolith DOWN from the 1600 hard cap into the band (1600->1305); reconcile family extracted to delegate_custody_reconcile.py, shrink-only direction", "ouroboros/delegate_recovery.py": "Native owner waits and delegated sessions share the existing planned-restart transaction and repeated-cleanup preservation owner.", diff --git a/ouroboros/transport_custody.py b/ouroboros/transport_custody.py index 43c10c54a..9d1148830 100644 --- a/ouroboros/transport_custody.py +++ b/ouroboros/transport_custody.py @@ -145,6 +145,15 @@ def _capture_on_chain(error: BaseException) -> Any: return capture +def outcome_unknown_on_chain(error: BaseException) -> bool: + """Whether ``error``'s chain carries a dispatched attempt without a terminal provider fact. + + Provider-independent: a generic API exception, or a wrapper whose explicit + cause carries the capture, reads exactly like the typed Claudexor error. + """ + return getattr(_capture_on_chain(error), "state", None) in {"dispatched", "unresolved"} + + def _requests_protocol_death(exc: BaseException) -> Any: """The innermost typed death inside a requests body-disconnect wrapper, or None. diff --git a/skills/telegram/lib/telegram_quiz.py b/skills/telegram/lib/telegram_quiz.py index d554d48ab..ceeee8699 100644 --- a/skills/telegram/lib/telegram_quiz.py +++ b/skills/telegram/lib/telegram_quiz.py @@ -424,13 +424,17 @@ async def apply_retained_fact(api, token: str, lang: str, *, client_factory) -> async def _mark_answered(api, client, record: Dict[str, Any], answer: str, lang: str) -> None: token = mint_token(str(record.get("task_id") or ""), str(record.get("quiz_id") or "")) - remember_state(api, token, record, "answered") - message_id = int(record.get("message_id") or 0) - if not message_id: - return - await client.edit_message_text_with_inline_keyboard( - int(record.get("chat_id") or 0), message_id, _answered_text(record, answer, lang), [], parse_mode="", - ) + # The card's lifecycle lock (``_CARD_LOCKS``): an expiry edit already on the wire + # lands first, then this answer re-reads the stored card and settles it last. + async with card_lock(token): + record = quiz_for_token(api, token) or record + remember_state(api, token, record, "answered") + message_id = int(record.get("message_id") or 0) + if not message_id: + return + await client.edit_message_text_with_inline_keyboard( + int(record.get("chat_id") or 0), message_id, _answered_text(record, answer, lang), [], parse_mode="", + ) async def answer_from_callback( diff --git a/tests/test_acceptance_author_stop.py b/tests/test_acceptance_author_stop.py index 778e2e49b..df572582a 100644 --- a/tests/test_acceptance_author_stop.py +++ b/tests/test_acceptance_author_stop.py @@ -270,3 +270,107 @@ def test_the_row_reason_slot_carries_the_stop_rationale(): # No rationale recorded (a malformed or stale disposition): the typed sentence alone. record["outcome_axes"]["review"]["acceptance_decision"]["author_disposition"] = "" assert _completion_verdict(record, {}) == TASK_CAUSE_PHRASES["author_stop"] + + +def test_a_stop_after_an_earlier_panel_keeps_its_finality_through_the_keep_round(tmp_path, monkeypatch): + """(c) FAIL panel → action-only stop → keep: the stop honoured after an EARLIER + panel used to lose its finality on the next delivery pass — the stale reviewed + subject of that panel cleared the reviewed latch, the finish intent had been + consumed, the keep/replace control re-armed, a SECOND panel ran and Cyber's + advisory ``author_finish`` replaced the stop. One panel, final reason + ``author_stop``, the agent's rationale on the record, and never Done.""" + from ouroboros.outcomes import derive_loop_outcome + + run = _run_stop_loop(tmp_path, monkeypatch, [ + "The export endpoint ships.", # reviewed: FAIL, capsule fed back + _stop_tool_call(), # action-only stop after the feedback + json.dumps({"delivery_control": "keep"}), # the host's control round: keep the stop + STOP_TEXT, json.dumps({"delivery_control": "keep"}), STOP_TEXT, + json.dumps({"delivery_control": "keep"}), STOP_TEXT, + ], stop_services=False) + assert run.panels == ["The export endpoint ships."], "an honoured stop must never buy a second panel" + decision = run.trace["acceptance_decision"] + assert decision["reason"] == "author_stop" and decision["author_action"] == "stop" + assert decision["author_disposition"]["action"] == "stop" + assert decision["author_disposition"]["rationale"] == STOP_RATIONALE + assert run.result == STOP_TEXT + axes = derive_loop_outcome(run.result, run.usage, run.trace)["outcome_axes"] + assert axes["objective"]["status"] == "fail" and axes["objective"]["reason"] == "author_stop" + + +def _honoured_stop_after_a_panel(tmp_path, monkeypatch): + """A FAIL panel, its feedback exposed, then an action-only stop the host honours.""" + from ouroboros.acceptance_settlement import expose_acceptance_feedback + from ouroboros.loop_acceptance import merge_agent_acceptance_stance + + fx = _finality_pass(tmp_path, monkeypatch) + assert fx.run("initial answer") is True + expose_acceptance_feedback(fx.trace, fx.messages, "author-root") + fx.trace["tool_calls"].append({"tool": "task_acceptance_review", "args": {}}) + merge_agent_acceptance_stance(fx.trace, {"explicit_finish": True, "author_action": "stop", + "rationale": STOP_RATIONALE}, fx.ctx) + assert fx.run(STOP_TEXT) is False + assert fx.trace["acceptance_decision"]["reason"] == "author_stop" + return fx + + +def test_an_honoured_stop_holds_until_the_authors_next_decision_replaces_it(tmp_path, monkeypatch): + """(d) The earlier panel's reviewed subject no longer reopens review on a later + delivery pass (the stop binds no subject), over the same or a changed text; the + author's next explicit decision is still heard — a finish replaces the stop + without another panel.""" + from ouroboros.loop_acceptance import merge_agent_acceptance_stance + + fx = _honoured_stop_after_a_panel(tmp_path, monkeypatch) + for text in (STOP_TEXT, "A restated unfinished answer."): + assert fx.run(text) is False + assert fx.trace["acceptance_decision"]["reason"] == "author_stop" + fx.trace["tool_calls"].append({"tool": "task_acceptance_review", "args": {}}) + merge_agent_acceptance_stance(fx.trace, {"disposition": "partial", "explicit_finish": True, + "author_action": "finish", + "rationale": "Fixed the empty-input case after all."}, fx.ctx) + assert fx.run("revised answer") is False + decision = fx.trace["acceptance_decision"] + assert decision["reason"] == "author_finish" and decision["author_disposition"]["action"] == "finish" + assert fx.calls == ["initial answer"], "neither the held stop nor the finish bought a panel" + + +def test_owner_input_after_an_honoured_stop_takes_the_ordinary_review_path(tmp_path, monkeypatch): + """(e) New owner input supersedes the stop exactly as it supersedes any terminal + acceptance: Main answers it and that answer is reviewed normally.""" + import ouroboros.loop as loop_mod + + fx = _honoured_stop_after_a_panel(tmp_path, monkeypatch) + loop_mod._supersede_task_acceptance_for_owner_followup(fx.ctx, fx.trace) + assert fx.trace["acceptance_decision"]["reason"] == "owner_followup" + assert fx.run("The answer to the owner's follow-up.") is True + assert fx.calls == ["initial answer", "The answer to the owner's follow-up."] + assert fx.trace["acceptance_decision"]["reason"] != "author_stop" + + +def test_a_stop_recorded_under_exhausted_rounds_keeps_its_cause_and_its_rationale(): + """A2: when the review rounds ran out, ``_finish_advisory_author`` records the + stop under its TRUE cause, ``review_cycles_exhausted``, with ``author_action`` + and the disposition's action ``stop``. The row keeps that cause's sentence and + still carries the author's rationale; a finish, or a bare inherited stop action + without the author's stop disposition, carries none.""" + from ouroboros.project_dialogue import TASK_CAUSE_PHRASES, _author_stop_rationale, _completion_verdict + from ouroboros.review_records import recorded_author_stop + + author = {"action": "stop", "rationale": "Not ready: the export endpoint still fails.", "source": "author"} + structured = {"status": "finalized_unaccepted", "reason": "review_cycles_exhausted", + "author_action": "stop", "author_disposition": author} + finish = {"status": "finalized_unaccepted", "reason": "review_cycles_exhausted", "author_action": "finish", + "author_disposition": {"action": "finish", "rationale": "Shipping as is.", "source": "author"}} + bare = {"status": "finalized_unaccepted", "reason": "review_cycles_exhausted", "author_action": "stop"} + assert recorded_author_stop(structured) and recorded_author_stop({"reason": "author_stop"}) + assert _author_stop_rationale(structured) == "Not ready: the export endpoint still fails." + for other in (finish, bare, {**structured, "reason": "capsule_spent"}, {}, None): + assert not recorded_author_stop(other) + assert _author_stop_rationale(finish) == _author_stop_rationale(bare) == "" + record = {"status": "completed", "reason_code": "final_message", "outcome_axes": { + "execution": {"status": "ok"}, + "objective": {"status": "fail", "source": "task_acceptance_review", "reason": "review_cycles_exhausted"}, + "review": {"status": "skipped", "acceptance_decision": structured}}} + assert _completion_verdict(record, {}) == ( + TASK_CAUSE_PHRASES["review_cycles_exhausted"][:-1] + " · Not ready: the export endpoint still fails.") diff --git a/tests/test_post_task_model_wait.py b/tests/test_post_task_model_wait.py index 4ce528aa9..b9efc2eba 100644 --- a/tests/test_post_task_model_wait.py +++ b/tests/test_post_task_model_wait.py @@ -683,3 +683,175 @@ def test_ordinary_promotion_failure_degrades_the_checkpoint_but_keeps_the_global assert checkpoint["post_task_synthesis"] == "degraded" assert not checkpoint.get("post_task_stop_reason") assert f.stages[-2:] == ["backlog", "promotion-model"] and len(callbacks) == 1 + + +def _generic_unknown(carrier): + """A generic API exception (not the typed Claudexor error) whose physical attempt + was dispatched with no terminal provider fact — directly or on ``__cause__``.""" + inner = RuntimeError("connection reset mid-request") + inner.physical_attempt_capture = SimpleNamespace(state="unresolved") + if carrier == "direct": + return inner + try: + raise RuntimeError("request failed") from inner + except RuntimeError as wrapped: + return wrapped + + +def _real_promotion_chooser(monkeypatch, f): + from ouroboros import post_task_evolution as promotion + + monkeypatch.setattr("ouroboros.post_task_evolution.maybe_promote", f.real_promote) + monkeypatch.setattr(config, "get_post_task_evolution_enabled", lambda: True) + monkeypatch.setattr(config, "get_runtime_mode", lambda: "advanced") + monkeypatch.setattr(config, "get_post_task_evolution_cadence", lambda: "llm") + monkeypatch.setattr(promotion, "_eligible", lambda *_: True) + monkeypatch.setattr(promotion, "_is_canonical_run", lambda *_: True) + monkeypatch.setattr(promotion, "_closed_objectives_digest", lambda *_: "") + + +@pytest.mark.parametrize("carrier", ["direct", "cause"]) +@pytest.mark.parametrize("seam", ["groom", "chooser"]) +def test_generic_unknown_outcome_in_a_real_adapter_stops_every_later_paid_call(phase, monkeypatch, seam, carrier): + """F1: only the typed Claudexor error and ``BudgetExceeded`` were interruptions, + so a generic API exception carrying an unresolved attempt was swallowed by + ``groom_backlog`` and the promotion stage then bought the chooser and the global + callback. The consolidator's chain classifier now reads it provider-independently + in the REAL grooming and chooser adapters: nothing paid runs after it.""" + import functools + from ouroboros import improvement_backlog, llm_observability, post_task_synthesis + + f = phase + error = _generic_unknown(carrier) + + def dispatch(*_args, call_type="", **_kwargs): + f.stages.append(call_type) + raise error + + monkeypatch.setattr(llm_observability, "chat_observed", dispatch) + monkeypatch.setattr(pipeline, "_update_improvement_backlog", post_task_synthesis._update_improvement_backlog) + if seam == "groom": + monkeypatch.setattr(improvement_backlog, "groom_backlog", + functools.partial(improvement_backlog.groom_backlog, cap=0)) + monkeypatch.setattr("ouroboros.post_task_evolution.maybe_promote", lambda *a: f.stages.append("chooser")) + else: + _real_promotion_chooser(monkeypatch, f) + monkeypatch.setattr(pipeline, "_run_reflection", lambda *a, **k: f.stages.append("reflection") or { + "backlog_candidates": [{"summary": "one auto item", "category": "process", "source": "execution_reflection"}], + "memory_actions": [{"type": "knowledge_write"}]}) + applied, callbacks = [], [] + monkeypatch.setattr(pipeline, "_apply_reflection_memory_actions", lambda *a, **k: applied.append(1)) + pipeline._run_post_task_processing_async( + f.env, f.task, {"rounds": 3}, {}, {}, f.root / "logs", event_queue=f.events, + on_reflection=lambda *args: callbacks.append(args)) + assert f.done.wait(5) + checkpoint = load_task_result(f.root, f.task["id"])["root_phase_checkpoint"] + assert checkpoint["post_task_synthesis"] == "degraded" + assert checkpoint["post_task_stop_reason"] == "provider_outcome_unknown:skipped=" + paid = "backlog_groom" if seam == "groom" else "post_task_evolution_decision" + assert f.stages == ["facts", "chat", "scratch", "reflection", paid], "no paid call after an unknown outcome" + assert callbacks == [] and not f.engine.creates + assert applied == [1], "the completed reflection's free actions are kept, applied once" + + +@pytest.mark.parametrize("outcome", ["failure", "no_op"]) +def test_ordinary_grooming_failure_degrades_the_checkpoint_and_keeps_the_chooser(phase, monkeypatch, outcome): + """F3: ``groom_backlog`` turned a confirmed ordinary provider failure into 0, the + backlog adapter returned added-or-0 and the promotion stage ignored it, so a stage + that lost its grooming wrote ``completed``. Through the REAL adapters the failure + reaches the stage: degraded, nothing skipped, the chooser and the global callback + still run. A genuine no-op (the backlog is below the grooming trigger) completes.""" + import functools + from ouroboros import improvement_backlog, llm_observability, post_task_synthesis + + f = phase + + def dispatch(*_args, **_kwargs): + f.stages.append("groom-model") + raise RuntimeError("grooming provider failed") + + monkeypatch.setattr(llm_observability, "chat_observed", dispatch) + if outcome == "failure": + monkeypatch.setattr(improvement_backlog, "groom_backlog", + functools.partial(improvement_backlog.groom_backlog, cap=0)) + monkeypatch.setattr(pipeline, "_update_improvement_backlog", post_task_synthesis._update_improvement_backlog) + monkeypatch.setattr(pipeline, "_run_reflection", lambda *a, **k: f.stages.append("reflection") or { + "backlog_candidates": [{"summary": "one auto item", "category": "process", "source": "execution_reflection"}]}) + monkeypatch.setattr("ouroboros.post_task_evolution.maybe_promote", lambda *a: f.stages.append("chooser")) + callbacks = [] + pipeline._run_post_task_processing_async( + f.env, f.task, {"rounds": 3}, {}, {}, f.root / "logs", event_queue=f.events, + on_reflection=lambda *args: callbacks.append(args)) + assert f.done.wait(5) + checkpoint = load_task_result(f.root, f.task["id"])["root_phase_checkpoint"] + assert checkpoint["post_task_synthesis"] == ("degraded" if outcome == "failure" else "completed") + assert not checkpoint.get("post_task_stop_reason") + assert f.stages == (["facts", "chat", "scratch", "reflection"] + + (["groom-model"] if outcome == "failure" else []) + ["chooser"]) + assert len(callbacks) == 1 + assert improvement_backlog.load_backlog_items(f.root), "the appended candidate survives a failed grooming" + + +def test_reflection_preparation_failure_degrades_the_checkpoint_with_a_typed_row(phase, monkeypatch): + """F3: ``generate_reflection``'s own catch returned a placeholder WITHOUT + ``memory_operation_errors``, so the coordinator read a clean reflection and wrote + ``completed`` over a stage that lost its work. Through the REAL adapters (no stub + of ``generate_reflection``) the placeholder carries a typed row: degraded, nothing + skipped, no reflection call bought, and the promotion stage still runs.""" + from ouroboros import consolidator, post_task_synthesis, reflection + + f = phase + monkeypatch.setattr(pipeline, "_run_reflection", post_task_synthesis._run_reflection) + monkeypatch.setattr(reflection, "should_generate_reflection", lambda *a, **k: True) + + def unwritable(*_args, **_kwargs): + f.stages.append("retain") + raise RuntimeError("retention store unwritable") + + monkeypatch.setattr(consolidator, "retain_memory_source", unwritable) + entries = [] + monkeypatch.setattr(reflection, "append_reflection_routed", lambda _env, _task, entry: entries.append(entry)) + launch(f) + assert f.done.wait(5) + checkpoint = load_task_result(f.root, f.task["id"])["root_phase_checkpoint"] + assert checkpoint["post_task_synthesis"] == "degraded" + assert not checkpoint.get("post_task_stop_reason") + assert f.stages == ["facts", "chat", "scratch", "retain", "backlog"] + assert not f.engine.creates, "a failed preparation buys no reflection call" + [entry] = entries + assert entry["reflection"].startswith("(reflection generation failed") + assert [row["kind"] for row in entry["memory_operation_errors"]] == ["reflection_failed"] + assert "retention store unwritable" in entry["memory_operation_errors"][0]["message"] + + +def test_split_root_facts_row_counts_the_actor_store_when_synthesis_runs_canonically(phase, monkeypatch): + """F2 (TZ-2 C2): a split NON-Project root synthesizes with the parent env and task, + so ``env.drive_root`` IS the canonical drive; passing it as the child store folded + two identical canonical stores into one, and a file present only in the actor's + child store (before copy-back) was never walked — a confirmed-looking zero. Through + the real dispatch the actor store is the row's recorded ``child_drive_root``.""" + import json + from ouroboros import post_task_synthesis + from ouroboros.headless import task_artifacts_dir + + f = phase + child = f.root.parent / "child" + child.mkdir() + (task_artifacts_dir(child, f.task["id"]) / "report.md").write_text("r", encoding="utf-8") + monkeypatch.setattr(pipeline, "_record_task_facts", post_task_synthesis._record_task_facts) + monkeypatch.setattr(pipeline, "_run_reflection", lambda *a, **k: None) + child_env = SimpleNamespace(drive_root=child, repo_dir=f.root.parent, drive_path=lambda rel: child / rel) + child_task = {**f.task, "budget_drive_root": str(f.root), "drive_root": str(child)} + parent_task = {**child_task, "drive_root": str(f.root), "child_drive_root": str(child)} + pipeline._dispatch_root_post_task( + child_env, child_task, "Already answered", None, [], {"rounds": 3}, {}, {}, child / "logs", + budget_drive_root=str(f.root), split_drive=True, project_scoped=False, project_task=False, + parent_env=f.env, parent_task=parent_task) + assert f.done.wait(5) + rows = [json.loads(line) for line in (f.root / "logs" / "chat.jsonl").read_text(encoding="utf-8").splitlines()] + [row] = [r for r in rows if r.get("summary_id") == f"task-facts:{f.task['id']}"] + fact = row["files_rescued"] + assert (fact["count"], fact["state"], fact["hash_computed"]) == (1, "positive", False) + assert [s["store"] for s in fact["stores"]] == [ + str(task_artifacts_dir(f.root, f.task["id"], create=False)), + str(task_artifacts_dir(child, f.task["id"], create=False))] diff --git a/tests/test_telegram_quiz_state.py b/tests/test_telegram_quiz_state.py index 113945f5f..611f6c25a 100644 --- a/tests/test_telegram_quiz_state.py +++ b/tests/test_telegram_quiz_state.py @@ -332,3 +332,52 @@ def test_concurrent_edits_are_serialized_so_an_expiry_never_lands_over_an_answer assert [bool(edit[3]) for edit in sent] == [True, False] # expiry first, then the answer assert sent[-1][2].endswith("\nAnswered: 2. postgres") assert card.record()["state"] == "answered" + + +@pytest.mark.parametrize("path", ["callback", "reply"]) +def test_a_telegram_answer_waits_for_an_in_flight_lifecycle_edit_and_lands_last(card, monkeypatch, path): + """Finding A1: the owner's own Telegram answer settled the card OUTSIDE the card + lock, so an expiry edit already past its state read (``follow_lifecycle`` holds + the lock while its edit is on the wire) landed after the answered edit and put + the buttons back while the stored state said ``answered``. The answer takes the + same lock: the toast never waits, the expiry lands first, the answered edit last.""" + card.send(wait_for_answer=True) + plugin, api = card.plugin, card.api + gate = asyncio.Event() + sent = [] # every edit in the order Telegram would receive it + + async def slow_edit(self, chat_id, message_id, text, keyboard, parse_mode="HTML"): + if keyboard: # the expiry edit (buttons restored) is the slow one + await gate.wait() + sent.append((chat_id, message_id, text, keyboard)) + return True + + monkeypatch.setattr(Client, "edit_message_text_with_inline_keyboard", slow_edit) + client = Client("token") + + async def post(_api, _path, body): + return 200, {"ok": True, "state": "answered", "answered_index": 1} + + async def answer(): + if path == "callback": + await plugin.telegram_quiz.answer_from_callback( + api, client, f"qz:{card.token}:1", cb_id="cb", update_id=7, lang="en", post=post) + else: + await plugin.telegram_quiz.answer_from_reply( + api, client, card.record(), "2. postgres", chat_id=42, update_id=7, lang="en", post=post) + + async def scenario(): + expiry = asyncio.ensure_future(plugin._make_quiz_state(api)(_fact(state="expired_terminal"))) + await asyncio.sleep(0) + assert card.record()["state"] == "expired_terminal" and sent == [] # stored; its edit is in flight + answered = asyncio.ensure_future(answer()) + await asyncio.sleep(0) + assert client.toasts or client.sent # the outcome is told at once, never behind the lock + gate.set() + await asyncio.gather(expiry, answered) + + asyncio.run(scenario()) + assert [bool(edit[3]) for edit in sent] == [True, False] # expiry first, then the answer + assert sent[-1][2].endswith("\nAnswered: 2. postgres") + assert card.record()["state"] == "answered" + assert card.api.logs == [] diff --git a/web/modules/chat_decision.js b/web/modules/chat_decision.js index 26071f77c..1af28fec0 100644 --- a/web/modules/chat_decision.js +++ b/web/modules/chat_decision.js @@ -369,11 +369,10 @@ export function createChatDecision({ const corrupt = normalized.some( (option) => !option || typeof option !== 'object' || !String(option.label || '').trim()); const optionsKnown = Array.isArray(src.options) && !corrupt && normalized.length <= MAX_QUIZ_OPTIONS; - const options = optionsKnown ? normalized : []; return { quizId: String(src.quiz_id || ''), question: String((nested ? msg.text : src.question) || ''), - options, + options: optionsKnown ? normalized : [], optionsKnown, stake: String(src.stake || ''), assumption: String(src.assumption || ''), @@ -414,7 +413,6 @@ export function createChatDecision({ button.append(badge); } - async function submitAnswer(card, quiz, index, comment, settle = setCardState) { if (card.dataset.pending === '1') return; card.dataset.pending = '1'; @@ -561,10 +559,9 @@ export function createChatDecision({ } const buttons = card.querySelectorAll('.chat-quiz-option'); buttons.forEach((btn, i) => { - const disabled = !answerable; const chosen = state === 'answered' && answeredIndex !== null && i === answeredIndex; - if (btn.disabled !== disabled) { - btn.disabled = disabled; + if (btn.disabled !== !answerable) { + btn.disabled = !answerable; changed = true; } if (btn.classList.contains('chosen') !== chosen) { diff --git a/web/modules/log_events.js b/web/modules/log_events.js index ed690fb4c..2ad293315 100644 --- a/web/modules/log_events.js +++ b/web/modules/log_events.js @@ -533,6 +533,13 @@ function joinCauseClauses(clauses) { .join(' · '); } +// The explicit author stop (TZ-2 C4): the typed author_stop, or the structured stop recorded under its TRUE +// cause when the review rounds ran out. The twin of review_records.recorded_author_stop. +const explicitAuthorStop = (d) => d?.reason === 'author_stop' || (d?.reason === 'review_cycles_exhausted' + && d.author_action === 'stop' && d.author_disposition?.action === 'stop'); +// Its AUTHOR's reason beside the typed sentence (never a finish's or reviewer prose). Twin: project_dialogue._author_stop_rationale. +const authorStopRationale = (d) => (explicitAuthorStop(d) ? String(d.author_disposition?.rationale || '').split(/\s+/).filter(Boolean).join(' ') : ''); + // Standing limitations of the delivered answer, from facts already on the // record: deferred children, and a plan review still open at delivery (the // result's terminal_plan_review_open flag), worded by its class when one is @@ -540,11 +547,6 @@ function joinCauseClauses(clauses) { // The open-review classes that state a standing limitation (never the merely awaited case). const PLAN_REVIEW_OPEN_CLASSES = new Set(['plan_review_unanswered', 'plan_review_none_answered', 'plan_review_answered_open']); -// The AUTHOR's own reason for an explicit stop, beside the typed sentence (TZ-2 C4); the reviewer -// rationale stays in the card. The twin of project_dialogue._author_stop_rationale. -const authorStopRationale = (d) => (d?.reason === 'author_stop' && typeof d.author_disposition === 'object' - ? String(d.author_disposition?.rationale || '').split(/\s+/).filter(Boolean).join(' ') : ''); - function terminalLimitations(record, reason, held = false) { const deferred = Number(record?.outcome_axes?.objective?.deferred_count || 0) > 0; const planKey = reason === 'plan_review_advisory' ? reason : 'terminal_plan_review_open'; @@ -576,11 +578,10 @@ export function taskReasonDetail(evt) { if (taskStoppedWithSummary(evt)) { // An owner-requested stop is a success and carries its own marker instead. clause = ''; - } else if (((severity !== 'error' && severity !== 'cancelled') || (decisionCause === 'author_stop' && severity === 'error')) + } else if (((severity !== 'error' && severity !== 'cancelled') || (explicitAuthorStop(decision) && severity === 'error')) && decision?.status && (decision.status !== 'accepted' || Object.hasOwn(TASK_CAUSE_PHRASES, decisionCause))) { - // A REVIEW-caused warning, or the explicit author stop that ended a red card, is explained - // by the host's acceptance decision in its own typed reason (an accepted decision only when - // it has a sentence); the stored reviewer rationale stays in the card, the result and Logs. + // A REVIEW-caused warning, or the explicit author stop that ended a red card, speaks through the + // decision's TRUE typed reason (an accepted one only with a sentence); reviewer prose stays in the card. clause = taskReasonPhrase(decisionCause); } else if (severity === "cancelled" && origin && typeof origin === "object" && !Array.isArray(origin) && Object.keys(origin).length) { diff --git a/web/tests/fixtures/outcome_phase_parity.json b/web/tests/fixtures/outcome_phase_parity.json index 173d31843..19b08cb13 100644 --- a/web/tests/fixtures/outcome_phase_parity.json +++ b/web/tests/fixtures/outcome_phase_parity.json @@ -622,6 +622,20 @@ "headline": "Failed", "acceptance_clause": "Ouroboros stopped with unfinished work; no review approval was granted · Not ready: the export endpoint still fails on empty input." }, + { + "name": "an author stop recorded after the review rounds ran out keeps its true terminal cause and still carries the agent's own reason", + "record": {"status": "completed", "reason_code": "final_message", "outcome_axes": {"execution": {"status": "ok"}, "objective": {"status": "fail", "source": "task_acceptance_review", "reason": "review_cycles_exhausted", "outcome_tier": "blocked_with_evidence"}, "review": {"status": "skipped", "acceptance_decision": {"status": "finalized_unaccepted", "reason": "review_cycles_exhausted", "author_action": "stop", "enforcement": "advisory", "review_capacity": {"reason": "review_cycles_exhausted"}, "rationale": "reviewer prose that never reaches the row", "author_disposition": {"disposition": "", "action": "stop", "rationale": "Not ready: the export endpoint still fails on empty input.", "subject_hash": "s1", "reviewer_signal": "", "enforcement": "advisory", "recorded_at": "2026-09-25T00:00:00+00:00", "source": "author"}}}}}, + "phase": "error", + "headline": "Failed", + "acceptance_clause": "The task used up its review rounds before the answer was signed off · Not ready: the export endpoint still fails on empty input." + }, + { + "name": "a finish beside the exhausted review rounds is no stop: the author's finish rationale never reaches the row", + "record": {"status": "completed", "reason_code": "final_message", "outcome_axes": {"execution": {"status": "ok"}, "review": {"status": "degraded", "acceptance_decision": {"status": "finalized_unaccepted", "reason": "review_cycles_exhausted", "author_action": "finish", "enforcement": "advisory", "author_disposition": {"disposition": "partial", "action": "finish", "rationale": "Shipping as is.", "subject_hash": "s1", "reviewer_signal": "", "enforcement": "advisory", "recorded_at": "2026-09-25T00:00:00+00:00", "source": "author"}}}}}, + "phase": "warn", + "headline": "Done with warnings", + "acceptance_clause": "The task used up its review rounds before the answer was signed off." + }, { "name": "a blocking plan review whose reviewers were too few holds the work, and the card says so", "record": {"status": "completed", "reason_code": "final_message", "outcome_axes": {"execution": {"status": "ok"}, "objective": {"status": "fail", "source": "plan_review_quorum_unreachable", "reason": "plan_review_quorum_unreachable", "outcome_tier": "blocked_with_evidence"}, "review": {"status": "skipped"}}},