diff --git a/docs/architecture/05-supervisor-loop.md b/docs/architecture/05-supervisor-loop.md index 6898f7bc9..2824751c6 100644 --- a/docs/architecture/05-supervisor-loop.md +++ b/docs/architecture/05-supervisor-loop.md @@ -26,7 +26,7 @@ At actual execution start, `agent._persist_running_record` mirrors a split root' Pooled completion separates a finished file-save attempt from publishable terminal truth: the worker prepares its terminal task files (`headless.prepare_terminal_task_files`) before its buffered `task_done`, and `events_task_done` publishes only when `headless.terminal_task_files_ready` confirms CURRENT — on a split drive the child-bound copyback, with no workspace artifact finalization pending. Legacy or faulted completions go through `enqueue_terminal_file_recovery` (`worker_health` prepares and recovers, `task_reaper` runs the queues): unreadable or missing CURRENT keeps RUNNING ownership and retries on the health cadence — never a false Done, never replayed model work — and only CONFIRMED absence of a terminal source reaches the lifecycle-fault owner, which marks execution `infra_failed` while an early sticky completed status, the authored answer, review and cost survive. -A required owner wait keeps the task RUNNING, its worker and browser; a completed-tool checkpoint precedes parking. The worker lends active capacity; resumption needs a grant within the configured cap. Resumption never replays: same attempt, same start timestamp, continuation authority consumed before dispatch. Any mail wakes; only an owner answer answers. Stop, deadline and ceiling still bind. Ordinary Main/Project roots wait in-process instead (`owner_wait.direct_owner_wait`): no pooled slot is held, lent or synthesized, and the saved source is evidence, never a cold-restart grant — an addressed text or quiz answer resumes the same live stack and browser. +A required owner wait keeps the task RUNNING, its worker and browser after a completed-tool checkpoint. Capacity is lent; a capped grant resumes the same attempt/start without replay, consuming continuation authority before dispatch. Any mail wakes, only an owner answer answers; Stop, deadline and ceiling still bind. A settled Project result rejects fast mailbox delivery even while post-work holds RUNNING: no solve loop remains to read it before cleanup. Ordinary Main/Project waits stay in-process (`owner_wait.direct_owner_wait`), lend no pooled slot and resume the live stack/browser on addressed text or quiz; a saved source grants no cold restart. A crash or automatic timeout retry keeps the current-attempt checkpoint as evidence and never blindly replays the task; cold continuation exists only for a confirmed planned restart, which cannot preserve an OS browser session. An owner-requested MANAGED UPDATE is such a restart: its writer fence prepares the handoff before stopping the pool and passes the parked ids as `preserve_running_task_ids`, so a wait is requeued instead of interrupted, and a preparation failure blocks the update rather than terminalizing the wait; the re-exec seam arms the transaction only in the phases `_safe_restart_serialized` allows (launcher mode observes exit code 42), no-resume flags suppress it, and an aborted update's leftover record never authorizes a later manual Restart or rollback. A manual ROLLBACK returns to an older runtime, so it deliberately parks nothing. Cold continuation resumes the saved economics — original CostCeiling, hard clocks, the saved model's ContextFit, TaskModelWait choices, the pending budget decision of the saved round — and calendar deadlines and owner-wait time are unchanged. @@ -51,7 +51,7 @@ Beside stopping sits the owner "hurry" control: a typed task-local `kind=hurry` The event bus is process-lifetime rather than worker-generation-lifetime: respawns reuse one manager-backed queue shared by workers and direct chat, because a force-killed producer can corrupt a raw multiprocessing feeder frame and a queue rebuilt on pool rotation strands surviving producers on the old endpoint. Live-frame publication of persisted rows is exactly-once and process-symmetric: `ouroboros/utils.py::append_jsonl` streams only runtime `logs/*.jsonl` rows into the process log sink (never `chat.jsonl`, never state/memory/receipt stores), and each process suppresses the types whose live delivery has a dedicated owner (`WORKER_LOG_SINK_SUPPRESSED_TYPES`, the server superset `SERVER_LOG_SINK_SUPPRESSED_TYPES`). One persisted event produces exactly one live frame (`tests/test_log_forwarding.py`); an LLM call failure is one durable `llm_api_error` row and nothing else. -Heartbeat and progress are different evidence: a heartbeat proves a process or loop is alive; owner-visible progress and model-usage events prove the task advanced. Fresh descendant progress or queued descendants can keep an orchestrator alive, while an explicit deadline, absolute ceiling, cancellation and budget stop remain hard. After the typed finalization episode (§6), timeout handling freezes its decision under the queue lock, marks the worker `reaping`, and hands kill, join, salvage, retry and respawn to the single off-loop reaper; an orchestrator with live descendants is not blindly retried, because a retry would replay its plan and spawn a competing tree. No retry or new assignment may occupy a timed-out slot until the original process is provably dead: if kill and join cannot establish death, the reaper keeps a low-rank RUNNING result and the `reaping` slot, emits a visible `task_reaper_wedged` receipt and restart hint, and writes no terminal, `task_done`, retry or respawn — one slot is sacrificed rather than letting a still-running process race a replacement and overwrite its result; the next supervisor generation reconciles the record after old-generation process custody. Typed timeout codes (`queue_timeouts.TIMEOUT_TERMINAL_REASONS`) retain their `task_incident` identity; owner-facing grace, kill and salvage text use `project_dialogue.TASK_CAUSE_PHRASES`, with unknown codes still raw. +Heartbeat and progress are different evidence: a heartbeat proves a process or loop is alive; owner-visible progress and model-usage events prove the task advanced. Progress keeps an orchestrator alive, but deadline, Stop and budget remain hard. The absolute ceiling ends solve/owner-stop finalization, not settled post-work; Stop, deadline, money and per-call bounds still apply (§6). After the typed finalization episode (§6), timeout handling freezes its decision under the queue lock, marks the worker `reaping`, and hands kill, join, salvage, retry and respawn to the single off-loop reaper; an orchestrator with live descendants is not blindly retried, because a retry would replay its plan and spawn a competing tree. No retry or new assignment may occupy a timed-out slot until the original process is provably dead: if kill and join cannot establish death, the reaper keeps a low-rank RUNNING result and the `reaping` slot, emits a visible `task_reaper_wedged` receipt and restart hint, and writes no terminal, `task_done`, retry or respawn — one slot is sacrificed rather than letting a still-running process race a replacement and overwrite its result; the next supervisor generation reconciles the record after old-generation process custody. Typed timeout codes (`queue_timeouts.TIMEOUT_TERMINAL_REASONS`) retain their `task_incident` identity; owner-facing grace, kill and salvage text use `project_dialogue.TASK_CAUSE_PHRASES`, with unknown codes still raw. A spawned or respawned slot is not assignable until its child's PID-bound `worker_ready` row arrives (`supervisor/worker_pool_lifecycle.py`). A live child's own `worker_starting` row, emitted before extension loading and agent construction, permits one extension of `WORKER_READY_WINDOW_SEC` to `WORKER_READY_CEILING_SEC` (300 seconds from birth, both in `runtime_limits.py`); foreign or pre-spawn rows cannot extend another slot. `worker_ready_window_extended` records that decision. A silent child keeps the original window, and logging failure cannot block startup. After `WORKER_READY_MAX_ATTEMPTS` failed attempts, `Worker.readiness_exhausted` is final for that exact slot — late events cannot reopen it. Total exhaustion, distinguished from busy/booting/reaping capacity and from a live owner-wait stack, closes pooled ingress (owner `/review` included) without blocking direct chat/control or boot/update recovery; once RUNNING completion custody has settled, `disable_exhausted_worker_pool` fails unstarted PENDING work honestly with a Restart hint, and a new task cannot clear the latch. Readiness stays separate from liveness and task idle time; a watcher error releases only still-booting, non-exhausted slots to the crash detector (`worker_ready_released`). Linux workers use forkserver; macOS and Windows use spawn. diff --git a/docs/architecture/06-agent-core.md b/docs/architecture/06-agent-core.md index 5a45cff54..3354e60c6 100644 --- a/docs/architecture/06-agent-core.md +++ b/docs/architecture/06-agent-core.md @@ -535,7 +535,7 @@ The advisory rows a reflection or summary reads are ATTRIBUTED. Advisory runs ar Only roots synthesize; `root_phase_checkpoint` makes paid synthesis at-most-once across restart, while children contribute evidence. Durable-result persistence owes the answer as `final::` in `supervisor/terminal_delivery.py`'s bounded outbox (§5; normal/cancel/reap). `send_message` delivers immediately; the retained buffered copy shares its ID for durable dedupe. Replays use bounded backoff. Exhaustion/eviction preserves full text on disk, emits `terminal_delivery_exhausted` and a chat notice; external delivery remains at-least-once. Buffered `task_done` stays last to retain the slot/child drive during synthesis; a hung-synthesis reap need not lose the delivered answer. Project roots keep early answers in Project. Their canonical row and deferred Main mirror use `terminal_projection.settle_terminal_projection` via task-done/checkpoint/startup/maintenance; §3 "Main rows and host-stamped card rows" owns readiness, retirement and limits. -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: no model call or narrative; its metrics, routing and cost serve history. +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. 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. A typed stop, budget refusal or unknown paid outcome skips later paid stages and degrades the checkpoint; ordinary stage failures are isolated and already-produced reflection actions still apply. A running checkpoint after restart degrades without replaying a paid request. diff --git a/ouroboros/agent_task_pipeline.py b/ouroboros/agent_task_pipeline.py index 3856fdc52..3670b79c4 100644 --- a/ouroboros/agent_task_pipeline.py +++ b/ouroboros/agent_task_pipeline.py @@ -705,6 +705,15 @@ def emit_task_results( loop_outcome=loop_outcome, cost_fields=task_cost_fields, ) stored_result = load_task_result(env.drive_root, str(task.get("id") or "")) or {} + if _root_outbox and task.get("_skip_post_task_synthesis"): + # Stop before post-task dispatch forbids paid synthesis, not the free + # factual row. Record it after durable result write; the structural + # root predicate deliberately excludes stopped roots from recovery. + fact_usage = {**usage, "outcome_axes": outcome_axes, "reason_code": reason_code} + if _typed_routing_action: + fact_usage["typed_routing_action"] = _typed_routing_action + _record_task_facts(env, task, _pre_synthesis_usage_snapshot(env, task, fact_usage), + llm_trace, drive_logs) artifact_bundle = stored_result.get("artifact_bundle") if isinstance(stored_result.get("artifact_bundle"), dict) else {} review_projection = stored_result.get("review_projection") or {} pending_events.append({ diff --git a/ouroboros/server_owner_routing.py b/ouroboros/server_owner_routing.py index 9d7e6a879..89084aee2 100644 --- a/ouroboros/server_owner_routing.py +++ b/ouroboros/server_owner_routing.py @@ -224,6 +224,17 @@ def _route_project_chat_to_running_task( ) if live_meta is None and not still_pending: return "" + # A worker may remain RUNNING to finish paid post-work after its + # answer/result settled. Its solve loop no longer drains this + # mailbox; accepting an owner follow-up here would label it + # delivered, then terminal cleanup would erase it unread. + # Check the actor's own drive: split-root copyback can lag the + # already-settled result while the worker still owns this slot. + from ouroboros.task_results import load_task_result + from ouroboros.task_status import SETTLED_STATUSES + + if (load_task_result(task_drive, tid) or {}).get("status") in SETTLED_STATUSES: + return "" # Phase A: a task whose cancellation is PENDING must not accept a # new owner message — same refusal the steer_task route makes, # checked inside this admission transaction. Falling through to diff --git a/tests/test_agent_task_pipeline.py b/tests/test_agent_task_pipeline.py index 47fbea356..c54c176b0 100644 --- a/tests/test_agent_task_pipeline.py +++ b/tests/test_agent_task_pipeline.py @@ -421,7 +421,7 @@ def test_stopped_direct_turn_pays_no_post_task_synthesis(tmp_path, monkeypatch): ``root_phase_checkpoint`` is seeded for the boot reconciler to re-pay. A positive control (the same turn, not stopped) proves the recording model would have seen the paid reflection call. Neither turn buys a narrative: - the dispatched worker writes one free host facts row, the stopped turn none.""" + both keep free host facts even when Stop prevents the worker from starting.""" import ouroboros.llm as llm_mod from ouroboros.outcomes import REASON_OWNER_REQUESTED_FINALIZATION from supervisor.owner_stop import REASON_OWNER_STOPPED_DIRECT_TURN @@ -507,8 +507,8 @@ def test_stopped_direct_turn_pays_no_post_task_synthesis(tmp_path, monkeypatch): rows = [json.loads(line) for line in chat_log.read_text(encoding="utf-8").splitlines()] if chat_log.exists() else [] return [row.get("summary_kind") for row in rows if row.get("type") == "task_summary" and row.get("task_id") == task_id] - stopped_kinds = _summary_kinds("stopped1") # no post-task worker was dispatched at all - assert "host_task_facts" not in stopped_kinds and "authored_root_summary" not in stopped_kinds + stopped_kinds = _summary_kinds("stopped1") # no paid post-task worker was dispatched + assert stopped_kinds == ["host_task_facts"] # free facts still reach pruned-result readers _task, _events, control_calls = _turn("control1", stopped=False) assert len(control_calls) >= 1, control_calls diff --git a/tests/test_tz2_last_drain.py b/tests/test_tz2_last_drain.py new file mode 100644 index 000000000..70fba6e73 --- /dev/null +++ b/tests/test_tz2_last_drain.py @@ -0,0 +1,93 @@ +"""A finalizing worker cannot consume a new owner turn after its last loop drain.""" + +import pytest + +from tests.test_project_routing_v664 import _ctx + + +@pytest.mark.parametrize("split", [False, True]) +def test_settled_project_root_does_not_accept_unread_owner_mail(tmp_path, monkeypatch, split): + from ouroboros.owner_mailbox import drain_owner_entries, _mailbox_path + from ouroboros.projects_registry import create_project + from ouroboros.server_owner_routing import _route_project_chat_to_running_task + from ouroboros.task_results import write_task_result + from supervisor.terminal_delivery import cleanup_settled_owner_mailbox + + project = create_project(tmp_path, "last-drain") + chat_id = int(project["chat_id"]) + actor_root = tmp_path / "child-drive" if split else tmp_path + actor_root.mkdir(exist_ok=True) + task = {"id": "last-drain-root", "chat_id": chat_id, "root_task_id": "last-drain-root", + "delegation_role": "root", "drive_root": str(actor_root)} + # The answer has settled, but the pooled worker is still RUNNING for post-work. + # Split-root copyback need not have published the terminal row in the canonical root yet. + write_task_result(actor_root, task["id"], "completed", result="answer", + root_phase_checkpoint={"post_task_synthesis": "running"}) + ctx = _ctx(tmp_path, running={task["id"]: {"task": task}}) + monkeypatch.setattr("ouroboros.server_owner_routing._addressable_root_tasks", + lambda *_: [{"task_id": task["id"], "project_id": project["id"]}]) + + routed = _route_project_chat_to_running_task(ctx, chat_id, "late owner clarification", "late-1") + assert routed == "" # fall through to a fresh decision turn; no false delivered receipt + assert drain_owner_entries(tmp_path, task["id"]) == [] + if not split: + write_task_result(tmp_path, task["id"], "completed", + root_phase_checkpoint={"post_task_synthesis": "completed"}) + cleanup_settled_owner_mailbox(tmp_path, task["id"], task) + assert not _mailbox_path(actor_root, task["id"]).exists() + + +def test_running_project_root_still_receives_owner_mail(tmp_path, monkeypatch): + from ouroboros.owner_mailbox import drain_owner_entries + from ouroboros.projects_registry import create_project + from ouroboros.server_owner_routing import _route_project_chat_to_running_task + from ouroboros.task_results import write_task_result + + project = create_project(tmp_path, "live-loop") + chat_id = int(project["chat_id"]) + task = {"id": "live-loop-root", "chat_id": chat_id, "root_task_id": "live-loop-root", + "delegation_role": "root", "drive_root": str(tmp_path)} + write_task_result(tmp_path, task["id"], "running", result="") + ctx = _ctx(tmp_path, running={task["id"]: {"task": task}}) + monkeypatch.setattr("ouroboros.server_owner_routing._addressable_root_tasks", + lambda *_: [{"task_id": task["id"], "project_id": project["id"]}]) + assert _route_project_chat_to_running_task(ctx, chat_id, "continue", "live-1") == task["id"] + assert [row["text"] for row in drain_owner_entries(tmp_path, task["id"])] == ["continue"] + + +def test_settled_project_followup_enters_new_decision_turn(tmp_path, monkeypatch): + import server + from ouroboros.owner_mailbox import drain_owner_entries + from ouroboros.projects_registry import create_project + from ouroboros.task_results import write_task_result + + project = create_project(tmp_path, "new-turn") + chat_id = int(project["chat_id"]) + task = {"id": "finished-loop", "chat_id": chat_id, "root_task_id": "finished-loop", + "delegation_role": "root", "drive_root": str(tmp_path)} + write_task_result(tmp_path, task["id"], "completed", result="answered", + root_phase_checkpoint={"post_task_synthesis": "running"}) + calls = [] + ctx = _ctx(tmp_path, running={task["id"]: {"task": task}}, + direct=lambda *_a, **_k: calls.append("new_turn")) + monkeypatch.setattr("ouroboros.server_owner_routing._addressable_root_tasks", + lambda *_: [{"task_id": task["id"], "project_id": project["id"]}]) + monkeypatch.setattr("supervisor.message_bus.log_chat", lambda *_a, **_k: None) + + class Bridge: + def get_updates(self, offset=0, timeout=1): + return [{"update_id": 1, "message": {"chat": {"id": chat_id}, "from": {"id": 1}, + "text": "new question after final answer", "source": "web", + "client_message_id": "post-work-1"}}] + + def send_routing_ack(self, *args, **kwargs): + calls.append((args, kwargs)) + + def broadcast(self, _payload): + return None + + server._process_bridge_updates(Bridge(), 0, ctx) + assert "new_turn" in calls + assert not any(isinstance(call, tuple) and call[1].get("action") == "mailbox_delivery" + for call in calls) + assert drain_owner_entries(tmp_path, task["id"]) == []