diff --git a/docs/PERSISTENCE.md b/docs/PERSISTENCE.md index b88afcdd1..542f8241a 100644 --- a/docs/PERSISTENCE.md +++ b/docs/PERSISTENCE.md @@ -106,7 +106,7 @@ scanned data-relative path to be covered by a row here (count-anchored both ways | `state/code_intel//inventory.json` | `ouroboros/code_intelligence.py` | own `schema_version: 2` (older/malformed rebuilt silently) | per-repo rewrite in place; stale roots age-pruned at startup by `inventory.json` mtime past GC retention (pure cache) | pure derived cache; one full re-index | | `state/extension_reconcile/` (+`failed/`) | `ouroboros/extension_reconcile_queue.py` | none — not needed (one-shot markers) | consumed by server loop; after 5 attempts moved to `failed/`, where markers age-prune past GC retention (the failure fact stays durable in events.jsonl) | pending worker→server reconciles lost; re-toggle heals | | `state/workspace_executor_processes/` | `ouroboros/workspace_executor.py` | own `schema_version: 1` + owner tag | unlink on stop; stale rows filtered at read (pid/cmd-sha) | service processes survive unreaped | -| `state/acceptance_fence_acks/..json` | `supervisor/events_worker_reports.py` (writer, one file per pooled-worker request), `ouroboros/agent.py` (the requester reads and unlinks only its own `req`; direct turns fence in-process and write none) | none — transport ack | inline GC on write (255 newest / 3600 s); a late ack nobody waits for ages out here | the waiting worker re-sends once, then records a typed `supervisor_ack_unavailable` unknown — never read as a refusal | +| `state/acceptance_fence_acks/..json` | `supervisor/events_worker_reports.py` (writer, one file per pooled-worker request), `ouroboros/agent.py` (the requester reads and unlinks only its own `req`; direct turns fence in-process and write none) | none — transport ack | inline GC on write (255 newest / 3600 s); a late ack nobody waits for ages out here | the waiting worker re-sends a transition once (a read never), then records a typed `supervisor_ack_unavailable` unknown — never read as a refusal | | `state/headless_tasks//` (child data drives) | `ouroboros/headless.py` | child `state/state.json`: `schema_version: 1` | GC-retention prune at startup (terminal + age; skips artifacts-not-terminal / refs-unpromoted) | in-flight child drives and unpromoted child refs lost | | `state/cx/` (managed Claudexor runtime) | `ouroboros/claudexor_runtime.py` | meta files versioned (`_NODE_META_SCHEMA_VERSION: 2`, pin 1) | deliberate keep-all (rollback selects older pins); staging/displaced temporaries reaped | re-downloaded on demand; nothing durable lost | | `state/betterleaks/` (managed runtime + cache) | `ouroboros/betterleaks_runtime.py` | own manifests (`schema_version: 1`) | keep-all archive cache — accepted (managed runtime) | re-downloaded on demand | diff --git a/docs/architecture/10-key-invariants.md b/docs/architecture/10-key-invariants.md index e00f7398f..fcd601622 100644 --- a/docs/architecture/10-key-invariants.md +++ b/docs/architecture/10-key-invariants.md @@ -26,7 +26,7 @@ This chapter is the short list of properties the rest of the book must not contr 22. **A host-owed round never parks.** A turn parks behind a review panel only when the panel is the sole thing it waits for; a turn in which the host has just spoken to the model never parks (`loop._finalize_loop_candidate`). 23. **Every call has a bound, and a recorder speaks only for what it collected.** A reviewer's tool call runs under the loop's per-tool timeout narrowed by the inherited dispatch deadline; a call that outlives it is abandoned — its late value sources no receipt and no coverage. A deadline recorder reconciles the turn's own panel at $0 before it writes a terminal reason. Owners: `review_native_episode.py`, `loop_tool_execution.py`, `acceptance_settlement.py`. 24. **The thread that answers workers runs only queue-bounded work.** Work that scales with history, the daemon or the network runs off-thread, reads its candidates before it reads liveness (one in-memory live source under `_queue_lock`), stops mutating when its loop generation ends and reaches the daemon attach-only once a stop is in flight. Named residuals on the loop thread: the 300-s zombie reconcile and the usage-ledger lock in the heartbeat handler. Owner: `ouroboros/server_maintenance.py`. -25. **An answer that has not arrived is a gap — never a refusal, a failure, a verdict or an owner message.** A direct turn applies its acceptance fence in-process (admission lock, then `_queue_lock` — never the reverse); a pooled request is idempotent by token, acknowledged per request (`..json`) and re-sent once; only `sealed` is a seal, an absent row is never read as one, and a fence that did not answer buys no model round: the panel runs as advice on `admission_fence_available=false`, delivery seals again, and a blocking install accepts a reviewer-approved answer with the typed note `admission_close_unconfirmed`. Owners: `ouroboros/agent.py`, `supervisor/queue_transitions.py`, `ouroboros/loop_delivery.py`. +25. **An answer that has not arrived is a gap — never a refusal, a failure, a verdict or an owner message.** A direct turn applies its acceptance fence in-process (admission lock, then `_queue_lock` — never the reverse); a pooled request is idempotent by token, acknowledged per request (`..json`), and a transition is re-sent once while a read never is; only `sealed` is a seal, an absent row is never read as one, and a fence that did not answer buys no model round: the panel runs as advice on `admission_fence_available=false`, delivery seals again, and a blocking install accepts a reviewer-approved answer with the typed note `admission_close_unconfirmed`. Owners: `ouroboros/agent.py`, `supervisor/queue_transitions.py`, `ouroboros/loop_delivery.py`. ### 10.1 Continuity data-flow map diff --git a/ouroboros/acceptance_settlement.py b/ouroboros/acceptance_settlement.py index 6723ff915..c556369f4 100644 --- a/ouroboros/acceptance_settlement.py +++ b/ouroboros/acceptance_settlement.py @@ -331,6 +331,7 @@ def forced_rail_panel_verdict(tools_ctx: Any, llm_trace: Dict[str, Any], rail_re """ from ouroboros.loop_acceptance_review import acceptance_run_pending from ouroboros.loop_delivery import delivery_subject_hash + from ouroboros.loop_messages import owner_source_sha256 from ouroboros.outcomes import ACCEPTANCE_ACCEPTED from ouroboros.review_dispatch import reconcile_pending_acceptance_runs from ouroboros.review_verdict import task_acceptance_is_clean @@ -347,6 +348,9 @@ def forced_rail_panel_verdict(tools_ctx: Any, llm_trace: Dict[str, Any], rail_re log.debug("a forced rail could not collect its own acceptance panel", exc_info=True) if acceptance_run_pending(run): return {"reason": "review_degraded", "review_pending": True} + reviewed_source = str(run.get("owner_source_sha256") or "") + if run.get("superseded_by_revision") or (reviewed_source and reviewed_source != str(owner_source_sha256(tools_ctx) or "")): + return {"reason": rail_reason} # that panel judged an earlier revision or older owner premises: this answer was never reviewed reviewed = (run.get("request") or {}).get("subject", "") if (not task_acceptance_is_clean(SimpleNamespace(**run)) or run.get("subject_hash") != delivery_subject_hash(tools_ctx, llm_trace, reviewed)): diff --git a/ouroboros/loop_acceptance_review.py b/ouroboros/loop_acceptance_review.py index b4e23e81f..232fd3d85 100644 --- a/ouroboros/loop_acceptance_review.py +++ b/ouroboros/loop_acceptance_review.py @@ -857,7 +857,7 @@ def _apply_task_acceptance_result( "dissent_noted": bool(dissent), }) ctx.tools._ctx._task_acceptance_improvement_passes = ctx.passes_done + 1 - if not _loop()._end_task_acceptance_fence(ctx.tools._ctx, outcome="revision"): + if _loop()._end_task_acceptance_fence(ctx.tools._ctx, outcome="revision").status == "refused": # a gap is not a refusal ctx.tools._ctx._task_acceptance_reviewed = True _loop()._set_acceptance_decision(ctx.llm_trace, { "status": ACCEPTANCE_FINALIZED_UNACCEPTED, diff --git a/ouroboros/loop_delivery.py b/ouroboros/loop_delivery.py index b439bb012..fe2d78b10 100644 --- a/ouroboros/loop_delivery.py +++ b/ouroboros/loop_delivery.py @@ -1128,7 +1128,12 @@ def _seal_admission_before_delivery(tools: ToolRegistry, limit_ctx: Any, llm_tra opened, _token = _loop()._begin_task_acceptance_fence(tool_ctx, limit_ctx.task_id) answer = opened and _loop()._end_task_acceptance_fence(tool_ctx, outcome="terminal") own_seal = (opened.status, opened.reason) == ("refused", "sealed") - if getattr(tool_ctx, "_task_acceptance_fence_generation_mismatch", False) or not (answer or own_seal or answer.status == "unknown"): + from ouroboros.loop_messages import _pending_owner_input_kinds + + # The wait for a silent supervisor is long enough for the owner to write: their mail is durable + # before any generation moves, so the local mailbox is read once more before a gap delivers. + gap = not answer and not own_seal and answer.status == "unknown" and not _pending_owner_input_kinds(tool_ctx) + if getattr(tool_ctx, "_task_acceptance_fence_generation_mismatch", False) or not (answer or own_seal or gap): _loop()._supersede_task_acceptance_for_owner_followup(tool_ctx, llm_trace) admission_lock = getattr(tool_ctx, "owner_message_admission_lock", None) admission_agent = getattr(tool_ctx, "owner_message_admission_agent", None) diff --git a/supervisor/queue_snapshot.py b/supervisor/queue_snapshot.py index b7fef9105..f50fae19c 100644 --- a/supervisor/queue_snapshot.py +++ b/supervisor/queue_snapshot.py @@ -423,10 +423,14 @@ def restore_pending_from_snapshot( # they are named; both lists describe work the stop caught, and one fence # call gives them the one cancel-intent path custody settles. direct_roots = dict(_queue().PRIOR_DIRECT_ROOTS) + # An in-process supervisor revival re-runs queue init while direct turns of THIS process are + # alive: the roster then names live work, not what a stop caught. + from supervisor.active_activity import get_direct_activity_registry + live_direct = {str(row.get("activity_id") or "") for row in get_direct_activity_registry().snapshot()} running_rows = snap.get("running") fenced_running = _fence_snapshot_running_rows( (running_rows if isinstance(running_rows, list) else []) - + [{"id": task_id} for task_id in direct_roots.get("task_ids") or []], + + [{"id": task_id} for task_id in direct_roots.get("task_ids") or [] if task_id not in live_direct], restored_ids={str(task.get("id") or "") for task in snapshot_pending}, ) if terminalized is not None: diff --git a/tests/test_admission_close_unconfirmed.py b/tests/test_admission_close_unconfirmed.py index 4fd5dadab..753fd15c6 100644 --- a/tests/test_admission_close_unconfirmed.py +++ b/tests/test_admission_close_unconfirmed.py @@ -278,6 +278,24 @@ def test_the_final_seal_reads_the_queues_typed_answer(tmp_path, monkeypatch, enf assert llm_trace["review_decision"]["admission_released"] is False +def test_owner_mail_that_arrived_during_the_silent_wait_is_not_delivered_over(tmp_path, monkeypatch): + """The wait for a silent supervisor is long enough for the owner to write, and their mail is + durable before any generation moves: a gap delivers only after the local mailbox was read once + more. Quiet direction: the same gap with an empty mailbox delivers with the note (table above).""" + from ouroboros.loop_delivery import _seal_admission_before_delivery + from ouroboros.owner_mailbox import write_owner_message + + monkeypatch.setenv("OUROBOROS_REVIEW_ENFORCEMENT", "blocking") + tools, limit_ctx = _seal_context(tmp_path, False) + tools._ctx.begin_acceptance_fence = _gap + assert write_owner_message(tmp_path, "Use the blue variant instead", "root-1", msg_id="owner-blue") + llm_trace = {"review_decision": {"eligibility": "eligible"}, "review_runs": [], + "acceptance_decision": {"status": ACCEPTANCE_ACCEPTED, "reason": "clean_pass", "source": "task_acceptance_review"}} + assert _seal_admission_before_delivery(tools, limit_ctx, llm_trace) is False + assert llm_trace["acceptance_decision"]["reason"] == "owner_followup" + assert llm_trace["acceptance_decision"]["status"] == "revision_requested" + + def test_a_locally_seen_owner_change_survives_a_silent_end(tmp_path): """`refused(generation_mismatch)` is a real owner follow-up even when the transport went silent: the local comparison decides the flag, not the missing ack. The quiet direction: diff --git a/tests/test_forced_acceptance_panel_collection.py b/tests/test_forced_acceptance_panel_collection.py index a2bf245be..d9aacd6c2 100644 --- a/tests/test_forced_acceptance_panel_collection.py +++ b/tests/test_forced_acceptance_panel_collection.py @@ -110,6 +110,21 @@ def test_settled_clean_pass_on_a_different_subject_does_not_accept(recorder): assert "review_pending" not in decision +@pytest.mark.parametrize("stale", [{"superseded_by_revision": True}, {"owner_source_sha256": "an-older-owner-corpus"}]) +def test_a_clean_pass_that_no_longer_speaks_for_this_turn_does_not_accept(recorder, stale): + """A panel that judged an earlier revision, or older owner premises, is not this answer's + review: the rail keeps its own "never reviewed" reason. The quiet direction is the + same-subject test above — an un-superseded PASS on the current owner source accepts.""" + ctx, record = recorder + trace = _trace(ctx) + subject = delivery_subject_hash(ctx, trace, ANSWER) + trace["review_runs"] = [{**_run(subject, actors=[_clean_actor()]), **stale}] + record(trace) + decision = trace["acceptance_decision"] + assert decision["status"] == "finalized_unaccepted" + assert decision["reason"] == BYPASS + + def test_pending_panel_stays_unaccepted_with_review_pending_and_keeps_its_rows(recorder): ctx, record = recorder trace = _trace(ctx) diff --git a/tests/test_restart_child_orphans.py b/tests/test_restart_child_orphans.py index 0f561b89e..96d7af8be 100644 --- a/tests/test_restart_child_orphans.py +++ b/tests/test_restart_child_orphans.py @@ -256,6 +256,28 @@ def test_a_direct_root_killed_by_the_window_is_cancelled_not_orphaned(qenv, monk assert load_task_result(qenv.drive, "direct-child")["status"] == STATUS_CANCELLED +def test_a_supervisor_revival_never_fences_a_direct_turn_alive_in_this_process(qenv, monkeypatch): + """An in-process supervisor revival re-runs queue init while direct turns of THIS process + are still running: the roster then names live work, not what a stop caught. The quiet + direction is the window-close test above — a turn nobody runs any more IS fenced.""" + from supervisor.active_activity import get_direct_activity_registry + + _boot_pool(qenv, monkeypatch) + write_task_result(qenv.drive, "live-direct-turn", STATUS_RUNNING, chat_id=1) + _restart_snapshot(qenv, monkeypatch, running=[]) + _direct_roots_fragment(qenv.drive, [{"task_id": "live-direct-turn", "chat_id": 1}]) + registry = get_direct_activity_registry() + registry.register("live-direct-turn", 1) + try: + qenv.q.init(qenv.drive) + fenced: list = [] + assert qenv.q.restore_pending_from_snapshot(terminalized=fenced) == 0 + finally: + registry.unregister("live-direct-turn") + assert fenced == [] + assert ci.active_intent(qenv.drive, "live-direct-turn") is None + + @pytest.mark.parametrize("fragment", ["incomplete", "missing"]) def test_a_direct_root_the_roster_does_not_name_is_never_fabricated(qenv, monkeypatch, fragment): """A turn the roster does not name gets no invented row: an ``incomplete`` diff --git a/tests/test_v678_acceptance_state.py b/tests/test_v678_acceptance_state.py index e42945796..2349f690c 100644 --- a/tests/test_v678_acceptance_state.py +++ b/tests/test_v678_acceptance_state.py @@ -74,7 +74,8 @@ def _apply_ctx(tmp_path, *, prior_trace=None, passes_done=0, budget_profile=None task_metadata={}, task_contract={"budget_profile": budget_profile} if budget_profile else {}, is_direct_chat=False, - end_acceptance_fence=(lambda **_k: {"ok": True}) if fence_ok else (lambda **_k: {"ok": False}), + end_acceptance_fence=(lambda **_k: (_ for _ in ()).throw(TimeoutError("no ack"))) if fence_ok == "silent" + else (lambda **_k: {"ok": True}) if fence_ok else (lambda **_k: {"ok": False}), _task_acceptance_fence_token="tok", ) trace = dict(prior_trace or {}) @@ -187,6 +188,12 @@ _DIALOGUE_TERMINAL = dict( "fence_reopen_failed", _ACTIONABLE_FAIL, {"fence_ok": False}, (ACCEPTANCE_FINALIZED_UNACCEPTED, "fence_reopen_failed"), ), + ( + # A supervisor that never answered the reopen is a gap, not a refusal: the improvement + # round still happens (the queued end reaches the queue first; the next begin re-adopts). + "fence_reopen_unanswered", _ACTIONABLE_FAIL, {"fence_ok": "silent"}, + (ACCEPTANCE_REVISION_REQUESTED, "improvement_capsule"), + ), ("dialogue_terminal", _DIALOGUE_TERMINAL, {}, (ACCEPTANCE_REVISION_REQUESTED, "improvement_capsule")), ("review_degraded", _NO_QUORUM, {}, (ACCEPTANCE_REVISION_REQUESTED, "improvement_capsule")), (