mirror of
https://github.com/razzant/ouroboros.git
synced 2026-10-03 04:07:04 +00:00
Close four gaps the review panel found in the night-class change
- An unanswered fence reopen after a FAIL panel was read as a refusal and ended the task as fence_reopen_failed: only a typed refusal does that now; a gap lets the improvement round happen (the queued end reaches the queue first and the next begin re-adopts the row). - A gap at the final seal delivered without looking at the local mailbox again: the wait for a silent supervisor is long enough for the owner to write, and their mail is durable before any generation moves. It is read once more, and pending owner input takes the existing revision path instead of the note. - The forced-rail recorder accepted on a clean PASS without the two predicates its neighbour applies: a superseded panel or one that judged older owner premises is not this answer's review — the rail keeps its own reason. - An in-process supervisor revival re-runs queue init while direct turns of this very process are alive: the restore no longer fences a roster row that the live direct-activity registry still holds. Docs: a fence transition is re-sent once, a read never.
This commit is contained in:
parent
375dc07688
commit
4acbeac092
10 changed files with 81 additions and 6 deletions
|
|
@ -106,7 +106,7 @@ scanned data-relative path to be covered by a row here (count-anchored both ways
|
|||
| `state/code_intel/<root-sha>/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/<token>.<req>.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/<token>.<req>.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/<id>/` (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 |
|
||||
|
|
|
|||
|
|
@ -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 (`<token>.<req>.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 (`<token>.<req>.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
|
||||
|
||||
|
|
|
|||
|
|
@ -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)):
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
|
|
|||
|
|
@ -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:
|
||||
|
|
|
|||
|
|
@ -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:
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
|
|
|||
|
|
@ -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``
|
||||
|
|
|
|||
|
|
@ -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")),
|
||||
(
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue