fix: preserve live actors, follow-up addresses and exception evidence

Version-neutral contribution candidate. Exact final review and deleted-project behavior remain outstanding.

Co-authored-by: Ouroboros <311266734+ouroboros-agent@users.noreply.github.com>
This commit is contained in:
Anton 2026-09-10 20:23:35 +03:00
parent eec71d53a0
commit e70b6058b4
21 changed files with 687 additions and 23 deletions

View file

@ -82,7 +82,7 @@ server.py (Starlette+uvicorn) ← HTTP + WebSocket on configurable host:port (de
├── cli.py ← Source/headless CLI over gateway tasks, logs, settings, skills, marketplace, local-model, and MCP wrappers
├── packaged_cli.py ← Packaged desktop CLI bridge: resolves bundle roots, bootstraps the launcher-managed repo, delegates to cli.py
├── packaged_cli_install.py ← Packaged CLI installer planning/execution for user-local command shims
├── agent.py ← Task orchestrator; the dispatch-note pair lives in `subagent_dispatch_notes.py` (same-name re-exports)
├── agent.py ← Task orchestrator; the dispatch-note pair lives in `subagent_dispatch_notes.py` (same-name re-exports). The outer catch around `run_llm_loop` owns the TERMINAL PROJECTION but not the evidence: the loop attaches its own accumulated usage/trace to the raised exception (`loop.py`, same in-memory objects, original exception re-raised unchanged), and this handler projects those instead of the untouched pre-loop defaults — so a lifecycle failure after real rounds can no longer publish a `0 calls` trace. The loop's tally rides `loop_outcome.usage` (ABI-3's honest loop plane); the top-level `total_rounds`/`prompt_tokens`/`completion_tokens` stay the LEDGER's answer from `reconstruct_task_cost`. An internal `task_exception` is `failure.kind = "runtime"`, never a fabricated provider failure
├── agent_startup_checks.py ← Worker-boot verification: dirty repo, version sync, budget, memory files, health checks
├── agent_task_pipeline.py ← Task execution pipeline orchestration; freezes one shared non-final subtree-cost snapshot for summary/reflection before the terminal checkpoint records final spend; hands the summary and reflection prompts the commit/advisory review lens PLUS the task's own acceptance-panel projection, and an absence statement names the lens it describes; calls the swarm-efficiency rollup owned by task_finalization.py at pipeline end
├── agent_dispatch.py, post_task_synthesis.py ← The agent's delegated-child dispatch seam, and the post-task synthesis workers the pipeline runs after a terminal result
@ -122,7 +122,7 @@ server.py (Starlette+uvicorn) ← HTTP + WebSocket on configurable host:port (de
├── model_send_seal.py ← The runtime invariant `model-visible ⟺ logged` for `model_send`: a reconstruction mismatch is a typed durable fact, and the call is NOT blocked — dispatch authority stays with the pre-existing in-memory identity re-check
├── cancel_intents.py ← Durable cancel-intent projection: compact locked `state/cancel_intents.json` of ACTIVE intents (request id, claim owner/pid + claim GENERATION fencing every mutation, `scope` recording single-vs-cascade so a watchdog replay re-runs the right shape) + forensic `cancel_intent` ledger rows; the ONE ingress `request_cancel` for the agent tool, HTTP single/cascade, and boot migration of legacy latch files — intent never rides the canonical task status; reads are strict and fail closed per §10 (typed `CancelIntentProjectionCorrupt`; enforcement degradation is owner-visible, never a silent "no intent"); a quarantined malformed row discloses once per row content — a ~20 s watchdog must not append the same disclosure forever, and a restart re-announcing once is honest; owns `claim_is_abandoned` and the `allow_settled_target` live-ownership exception (§10, cancellation custody)
├── owner_hurry.py ← Owner "hurry": a typed TASK-LOCAL acceleration latch, never a chat message; the durable `owner_hurry` projection is written by `update_json_locked` touching only its own keys — never `write_task_result`, whose status-regression guard could drop concurrent terminal fields — keyed by the real attempt identity `task["_attempt"]`; while latched, the next acceptance panel is skipped with zero reviewer calls (`acceptance_skip_applied`), remaining improvement passes overlay to 0 through `effective_budget_profile` (the immutable task_contract is never rewritten), and force-plan becomes task-locally advisory; the effect DIES WITH THE ATTEMPT (`retry_reset` on every same-id requeue producer), a never-applied request is marked `not_applied_before_terminal`, and the non-chat `owner_hurry` events are hidden from chat by `log_events.js`
├── owner_quiz.py ← Owner-quiz lifecycle projection: worker-side `record_asked`, request-id-idempotent first-answer-wins `record_answered` (option index validated against the STORED labels), structural-only `reconcile_terminal` (open → expired_terminal at task done; no host TTL), `quiz_states` replay; same locked-writer idiom as owner_hurry, touching only the `owner_quiz` key
├── owner_quiz.py ← Owner-quiz lifecycle projection: worker-side `record_asked`, request-id-idempotent first-answer-wins `record_answered` (option index validated against the STORED labels), structural-only `reconcile_terminal` (open → expired_terminal at task done; no host TTL) which also closes the PAIRED `owner_wait` under the same task-result authority — repairable on a second pass after a partial pair write, so a terminal task never replays as both expired and still waiting — `quiz_states` replay; same locked-writer idiom as owner_hurry, touching only the `owner_quiz` and paired `owner_wait` keys
├── owner_wait.py ← Native owner-answer continuation: completed-tool source checkpoint, original-process sleep, and same-task recovery only through an acknowledged planned-restart handoff; active capacity belongs to supervisor/worker_owner_wait.py
├── routing_wait.py ← Root-parameterized SSOT of the durable routing-receipt waits (`wait_for_promotion_admission`, `wait_for_routing_annotation`); tools/control.py keeps thin wrappers, so the gateway picker dispatcher confirms clicks through the SAME receipts the LLM routing tools poll
├── outcomes.py ← Typed task-outcome and acceptance-decision authority keeping the lifecycle/execution/objective/review/artifact/verification/child-absorption axes separate; policy denials, cosmetic exits, and ignored outcomes never masquerade as genuine tool failures; receipt reconciliation lives in `_outcome_receipts.py`, trace classification in `_outcome_tool_errors.py`; a verification ledger above the inline threshold rides as a stub whose `summary` is re-projected from the refreshed artifact file at finalization, and the stub is never a source for entries or outcome axes (the ledger embeds the task contract, which its entry count excludes, so a stub is the normal shape for a swarm root)
@ -307,7 +307,7 @@ server.py (Starlette+uvicorn) ← HTTP + WebSocket on configurable host:port (de
├── task_continuation.py ← Durable review continuation state
├── task_results.py ← Durable task results `task_results/<id>.json`; locked `task_acceptance_review_accounting` — the claim is minted at first physical reviewer dispatch, and a claim without a recoverable terminal host run is UNKNOWN, never permission to re-dispatch (double-spend fence); its read-only root review-capacity projection for configured-session wakes is WALLET and cancellation only (`root_task_id`, `cap_cycles`, `claimed_cycles`, `remaining_cycles`, `binding_seen`, `dedupe`, `state`, `reason` — no time axis: the launch rule `task_pacing.review_launch_allowed` is evaluated once per panel at loop admission (owner R55; the paid claim inside the dispatch stamp checks cancellation and the wallet only), and a descendant reads its own window from the coordination `time` fact)
├── task_result_schema.py ← Task-result schema admission: the `_schema_version` stamp, the classifier, and the quarantine an unstamped, future, malformed or retired-key row lands in
├── task_status.py ← Effective-status SSOT, lineage, bounded waits; worker-side `task_has_live_queue_ownership` (§10, cancellation custody)
├── task_status.py ← Effective-status SSOT, lineage, bounded waits; worker-side `task_has_live_queue_ownership` (§10, cancellation custody); the DESTRUCTIVE orphan predicate keeps the same fail-open-toward-liveness polarity — an in-process direct actor (`supervisor.active_activity` registry, deliberately absent from PENDING/RUNNING) and a missing/invalid/stale queue snapshot can never prove a task dead
├── git_shell_policy.py ← Structural git argv classifiers for the shell guards
├── protected_artifacts.py ← Execute-only black-box policy for protected artifacts
├── shell_parse.py ← One shell normalization shared by guard and execution (`recover_stringified_argv`, `normalize_check_argv`, `shell_tokens`, `shell_segments`, `canonical_command_text`): quoted shell punctuation is data, not syntax — over-splitting on a quoted `&&` is the fail-safe direction; `split_redirections` is the ONE redirect grammar read by `writer_target_rows`, whose per-segment `(argv, targets, inline_code, unprovable)` facts include bounded shell-`-c` recursion, stdin-program heredocs only when no inline/script operand exists, Python AST targets/independent UNKNOWN, sequential `cd`/`pushd`/`env -C` cwd changes, and non-concrete find/xargs placeholders; uncertainty widens only its own row/body mentions, while the separate mention lane keeps unmangled Windows drive/UNC spellings
@ -448,7 +448,7 @@ server.py (Starlette+uvicorn) ← HTTP + WebSocket on configurable host:port (de
│ ├── project_journal.py ← journal_write/read, workpad_read/write, journal_tail_digest (over-limit rejected); owns `mirror_tree_coordination_to_journal`, the durable-journal mirror of tree coordination
│ ├── presence.py ← configure_presence, initiate_presence, typed completion/cancel
│ ├── task_tree.py ← tree_note/tree_read (storage SSOT: task_tree_ledger.py)
│ ├── followup.py ← One deferred follow-up into `state/scheduled_tasks.json`; exactly one trigger (`once` ISO or 5-field cron+tz); cap 2 pending, over-limit refused
│ ├── followup.py ← One deferred follow-up into `state/scheduled_tasks.json`; exactly one trigger (`once` ISO or 5-field cron+tz); cap 2 pending, over-limit refused; preserves the SOURCE address (top-level `project_id` so `resolve_project_id` sees it, plus the originating `chat_id`) instead of defaulting to the global owner chat — an unscoped task's follow-up still takes the existing `owner_chat_id` default
│ ├── join_ledger.py ← Child-result absorption: validates lineage + exact hashes; dispositions integrated/irrelevant/deferred; `CHILD_RESULT_STALE`; keeps peek_task/discard_child_result
│ ├── delegate.py ← Delegation facade verbs `delegate_start` (with `retry_of`), `delegate_wait`, `delegate_cancel`, `delegate_answer`; the host pre-start rides the same `delegate_start(prompt="")` wrapper and the shared `subagent_runtime.exact_start`; supervision/recovery/custody/transport live in the named leaf modules (§6)
│ ├── delegate_integration.py ← Delegated-patch integration: `_mutation_authority`, `_provision_snapshot` (registered before the start intent), retry-binding validation, `_capture_terminal_patch`; the skill-payload cluster (`_payload_mutation_authority`, `_rebind_payload_reference`, `_write_payload_patch_artifacts` — git diff --binary; reserved paths refuse the WHOLE apply as `blocked_reserved_paths` with the candidate preserved) and `integrate_payload_patch` (CAS, index-free git apply in NO-REPOSITORY mode — `GIT_CEILING_DIRECTORIES` pinned at the resolved PARENT of the payload, because git still searches the ceiling entry itself, so an ancestor Git worktree above the runtime data root cannot make git treat the payload as a subdirectory prefix and silently skip every hunk at rc=0 while `--numstat` prints nothing; the same env bounds the `--numstat -z` touched-path reader; a post-apply live hash equal to the BASELINE with a non-empty touched set is refused typed as `INTEGRATE_APPLY_NO_OP` with the apply intent resolved and nothing disposed; the busy-check refuses a second delegation only while a run whose OWNER TASK IS STILL LIVE, or whose terminality cannot be proven, has open custody on the same payload, QUEUES the extension reconcile request via `request_extension_reconcile`)
@ -1001,7 +1001,7 @@ A headless task is ADDRESSED when it is admitted, not when it is displayed (`log
`queue_snapshot.json` is an atomic recovery and diagnostic projection, not a second scheduler. It carries pending and running rows, acceptance and root-budget fences, resident/active/parked worker counts, assignable capacity, and any pool-disabled reason. Startup restores a recent snapshot into an otherwise-empty pending queue and never resurrects ordinary RUNNING work. A selected native owner-wait handoff remains eligible beyond snapshot age only through its current waiting source and acknowledged planned-restart transaction. Terminal tasks stay terminal, a task with an active durable cancel intent (or a legacy cancel-requested latch file) is left for cancellation custody, descendants below an accepted or sealed root finalize as cancelled, and malformed durable fence evidence fails closed. Snapshot capture copies the live containers under the queue lock because concurrent HTTP mutation can otherwise crash the supervisor mid-iteration.
A required owner wait keeps its task RUNNING and retains the same worker, command queue, live browser and services. The completed-tool checkpoint and task-result `owner_wait` projection reach durable storage before the snapshot confirms the park. The worker lends only its `active_capacity`; the ordinary assignment tick replenishes the configured active capacity through the existing spawn/readiness owner. Addressed input requests a wake, which waits for active capacity and retires only an idle replacement before granting the original worker. Resumption preserves the task attempt and start timestamp and marks its continuation authority consumed before dispatch; no completion, new admission or replay of completed tools occurs. Waiting spares only the idle rail. Stop, deadline and absolute ceiling keep their authority, and cancellation, timeout or crash retires an inactive worker or replaces it when active capacity is missing, without exceeding the configured active limit. A process crash or automatic timeout retry preserves the current-attempt checkpoint as evidence and does not blindly replay the task; cold continuation is limited to confirmed planned restart, which cannot preserve an OS browser session. Every shutdown cleanup recognizes the same restart transaction. Cold preparation restores the original CostCeiling before constructing either context projection and retains hard clocks, refreshes genuine assignment progress, and rebinds ContextFit to the saved model. Cold continuation completes the saved round's pending budget decision before advancing to another model round or applying queued route overrides. TaskModelWait retains its per-role choices (including Auto), auto-continue choices and completed quota union; both native execution and supervisor timing retain the original execution basis and clock revision. Calendar deadlines and time spent waiting for the owner are unchanged. Successful grant consumes the queue task's resume locator; source evidence remains available through a later crash or timeout.
A required owner wait keeps its task RUNNING and retains the same worker, command queue, live browser and services. The completed-tool checkpoint and task-result `owner_wait` projection reach durable storage before the snapshot confirms the park. The worker lends only its `active_capacity`; the ordinary assignment tick replenishes the configured active capacity through the existing spawn/readiness owner. Addressed input requests a wake, which waits for active capacity and retires only an idle replacement before granting the original worker. Resumption preserves the task attempt and start timestamp and marks its continuation authority consumed before dispatch; no completion, new admission or replay of completed tools occurs. Waiting spares only the idle rail. Stop, deadline and absolute ceiling keep their authority, and cancellation, timeout or crash retires an inactive worker or replaces it when active capacity is missing, without exceeding the configured active limit. A process crash or automatic timeout retry preserves the current-attempt checkpoint as evidence and does not blindly replay the task; cold continuation is limited to confirmed planned restart, which cannot preserve an OS browser session. Every shutdown cleanup recognizes the same restart transaction. An owner-requested MANAGED UPDATE is such a planned restart: its writer fence prepares the owner-wait/planned-restart handoff before stopping the pool and passes those exact task ids as `preserve_running_task_ids`, so a parked wait is requeued instead of interrupted, and a preparation failure blocks the update (repo-writer admission re-opens) rather than terminalizing the wait. Immediately before the re-exec request the already-created transaction is armed through the existing one-shot environment handoff so DIRECT-server mode can acknowledge the same transaction on its same-PID successor (launcher mode still acknowledges by observing exit code 42), and a failed restart callback disarms that token. A manual ROLLBACK returns the tree to an older runtime, so it deliberately parks nothing and keeps the ordinary interrupt semantics. Cold preparation restores the original CostCeiling before constructing either context projection and retains hard clocks, refreshes genuine assignment progress, and rebinds ContextFit to the saved model. Cold continuation completes the saved round's pending budget decision before advancing to another model round or applying queued route overrides. TaskModelWait retains its per-role choices (including Auto), auto-continue choices and completed quota union; both native execution and supervisor timing retain the original execution basis and clock revision. Calendar deadlines and time spent waiting for the owner are unchanged. Successful grant consumes the queue task's resume locator; source evidence remains available through a later crash or timeout.
Ordinary Main/Project roots retain their in-process actor through the same completed-tool owner-wait boundary. `owner_wait.direct_owner_wait` uses that actor's existing mailbox and TaskModelWait controls; no pooled slot is held, lent or synthesized. Pending owner events are flushed before sleeping, and the loop handles controls/deadlines before its post-tool budget decision after wake. Its saved source is evidence, not an automatic cold-restart grant for an ordinary conversation. An addressed text or existing quiz answer resumes the same live stack and browser.
@ -2094,7 +2094,7 @@ and failed source writes disclose `applied_source_status="unavailable"`.
| `StateResponse.active_direct_turns`/`active_chat_activities` (phases queued/working/finalizing/budget_paused — one predicate `budget_pause_fact` decides budget pause); `TypingOutbound` activity fields; `ChatOutbound.task_phase`/`task_terminal_status` | `ouroboros/gateway/contracts.py`, `supervisor/active_activity.py`, `web/modules/chat_activity.js` | `tests/test_gateway_parity.py` + the activity test files |
| `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?}`, `QuizOutbound` (quiz_id, question, options, 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), separate `QuizStateOutbound` discriminator, `chat.quiz` host topic; the producer is the one escalation verb `escalate(question, options, stake, assumption, wait_for_answer=False)` — 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 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?}`, `QuizOutbound` (quiz_id, question, options, 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), separate `QuizStateOutbound` discriminator, `chat.quiz` host topic; the producer is the one escalation verb `escalate(question, options, stake, assumption, wait_for_answer=False)` — 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) and 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, `update_status_ready` WS notice | `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` |

View file

@ -65,7 +65,7 @@ Rows may import columns (`[graph].allowed`). `·` = forbidden direction.
## Hidden coupling (classified out of the strict graph)
- lazy-only cross-domain pairs: **95**
- lazy-only cross-domain pairs: **96**
- D01->D08
- D01->D10
- D01->D11
@ -145,6 +145,7 @@ Rows may import columns (`[graph].allowed`). `·` = forbidden direction.
- D17->D05
- D17->D06
- D17->D07
- D17->D08
- D17->D09
- D17->D12
- D18->D01

View file

@ -933,21 +933,43 @@ class OuroborosAgent:
"traceback": truncate_for_log(tb, 2000),
})
text = f"⚠️ Error during processing: {type(e).__name__}: {e}"
usage = {
"execution_status": "infra_failed",
"reason_code": "task_exception",
}
captured_usage = getattr(e, "_ouroboros_loop_usage", None)
captured_trace = getattr(e, "_ouroboros_loop_trace", None)
if isinstance(captured_usage, dict):
usage = dict(captured_usage)
else:
usage = {}
if isinstance(captured_trace, dict):
llm_trace = captured_trace
usage.update(
execution_status="infra_failed",
reason_code="task_exception",
)
try:
from ouroboros.outcomes import collect_trace_refs, derive_loop_outcome
from ouroboros.agent_task_pipeline import build_trace_summary
from ouroboros.task_results import STATUS_FAILED, write_task_result
# CW3: an ephemeral decision turn leaves no durable task_result even on error.
if not bool(task.get("_ephemeral_turn")):
# The loop's own tally rides ``loop_outcome.usage``
# (ABI-3's honest loop plane); the top-level
# total_rounds/prompt_tokens/completion_tokens stay
# the LEDGER's answer, written by the finalization
# pipeline from reconstruct_task_cost.
loop_outcome = derive_loop_outcome(text, usage, llm_trace)
trace_refs = loop_outcome.get("trace_refs") or collect_trace_refs(usage, llm_trace)
write_task_result(
self.env.drive_root,
str(task.get("id") or ""),
STATUS_FAILED,
result=text,
reason_code="task_exception",
outcome_axes=infra_failed_axes("task_exception", review_trigger="agent_exception"),
outcome_axes=loop_outcome.get("outcome_axes") or infra_failed_axes(
"task_exception", review_trigger="agent_exception",
),
loop_outcome=loop_outcome,
trace_summary=build_trace_summary(llm_trace),
trace_refs=trace_refs,
)
except Exception:
pass

View file

@ -102,6 +102,32 @@ def _active_restart_transaction_path(drive_root: Any) -> pathlib.Path:
return pathlib.Path(drive_root) / "state" / "delegate_recovery_transactions" / "active.json"
def arm_active_planned_restart_transaction(drive_root: Any) -> str:
"""Pass a prepared restart transaction to a direct re-exec successor.
Launcher-managed exits acknowledge the same durable transaction by waiting
for exit code 42. A direct server re-exec has no launcher, so it carries
the already-created transaction id through the existing one-shot
environment handoff consumed by ``_ack_direct_exec_successor``.
"""
try:
active = json.loads(_active_restart_transaction_path(drive_root).read_text(encoding="utf-8"))
except (OSError, ValueError):
return ""
if not isinstance(active, dict):
return ""
transaction_id = str(active.get("transaction_id") or "")
row = _read_restart_transaction(drive_root, transaction_id) if transaction_id else {}
if (
not transaction_id
or row.get("status") != "prepared"
or int(row.get("supervisor_pid") or 0) != os.getpid()
):
return ""
os.environ[PLANNED_RESTART_TRANSACTION_ENV] = transaction_id
return transaction_id
def _read_restart_transaction(drive_root: Any, transaction_id: str) -> dict[str, Any]:
try:
data = json.loads(
@ -994,6 +1020,7 @@ __all__ = [
"CAUSE_WORKER_CRASH",
"NO_RESUME_CAUSES",
"PLANNED_RESTART_TRANSACTION_ENV",
"arm_active_planned_restart_transaction",
"acknowledge_observed_restart_exit",
"adopt_handoff",
"authority_fingerprint_from_context",

View file

@ -872,6 +872,7 @@ lazy_only = [
"D17->D05",
"D17->D06",
"D17->D07",
"D17->D08",
"D17->D09",
"D17->D12",
"D18->D01",

View file

@ -405,9 +405,33 @@ def _quiesce_repo_writers(reason: str) -> list[str]:
if blocked:
open_repo_writer_admission()
return [f"active:{label}" for label in blocked]
preserve_running_task_ids: set[str] = set()
if reason != "manual_rollback":
try:
import uuid
from ouroboros.delegate_recovery import prepare_planned_restart_handoffs
from ouroboros.owner_wait import prepare_owner_wait_handoffs
from supervisor import workers as worker_state
restart_transaction_id = uuid.uuid4().hex
owner_wait_ids = prepare_owner_wait_handoffs(
DRIVE_ROOT, worker_state.RUNNING, restart_transaction_id,
)
preserve_running_task_ids = prepare_planned_restart_handoffs(
DRIVE_ROOT,
worker_state.RUNNING,
restart_transaction_id=restart_transaction_id,
additional_task_ids=owner_wait_ids,
)
except Exception as exc:
open_repo_writer_admission()
log.warning("Managed update owner-wait handoff preparation failed", exc_info=True)
return [f"owner_wait_handoff:{type(exc).__name__}: {exc}"]
survivors = kill_workers_for_update(
result_reason="Task interrupted by an owner-requested managed update.",
terminal_status="interrupted",
preserve_running_task_ids=preserve_running_task_ids,
)
if survivors:
return survivors
@ -473,6 +497,17 @@ def _rollback_fenced_update(reason: str, error: str, **extra: Any) -> JSONRespon
def _restart_response(request: Request, *, strategy: str, plan: dict) -> JSONResponse:
# Managed update prepared owner-wait/restart custody before stopping the
# pool. Arm the existing one-shot transaction token immediately before
# asking the server to re-exec, so direct-server mode can acknowledge the
# same transaction on the successor (launcher mode uses exit code 42).
try:
from supervisor.git_ops import DRIVE_ROOT
from ouroboros.delegate_recovery import arm_active_planned_restart_transaction
arm_active_planned_restart_transaction(DRIVE_ROOT)
except Exception:
log.warning("managed update restart transaction could not be armed", exc_info=True)
try:
restarting = _request_restart(request)
except Exception as exc:
@ -482,6 +517,13 @@ def _restart_response(request: Request, *, strategy: str, plan: dict) -> JSONRes
else:
restart_error = "restart callback is unavailable" if not restarting else ""
if not restarting:
try:
from ouroboros.delegate_recovery import PLANNED_RESTART_TRANSACTION_ENV
import os
os.environ.pop(PLANNED_RESTART_TRANSACTION_ENV, None)
except Exception:
log.warning("failed to disarm managed update restart transaction", exc_info=True)
return JSONResponse(
{
"status": "restart_required",

View file

@ -641,6 +641,17 @@ def run_llm_loop(
_delegate_hold_close(tools, drive_logs=drive_logs, task_id=task_id, detail="budget")
return _handle_budget_exceeded(
exc, exit_ctx, limit_ctx=limit_ctx, episode=transport_wait)
except Exception as exc:
# The caller still owns the terminal projection, but the loop owns the
# accumulated evidence. Keep the same in-memory objects on the raised
# exception so an unexpected lifecycle failure cannot turn a completed
# multi-round trace into an empty ``0 calls`` result.
try:
setattr(exc, "_ouroboros_loop_usage", accumulated_usage)
setattr(exc, "_ouroboros_loop_trace", llm_trace)
except Exception:
pass
raise
finally:
# No stale active latch behind an in-process exit (a crash skips this frame, keeping the latch for recovery).
_delegate_hold_close(tools, drive_logs=drive_logs, task_id=task_id, detail="loop_exit")

View file

@ -1053,7 +1053,12 @@ def derive_loop_outcome(final_text: str, usage: Dict[str, Any], llm_trace: Dict[
if usage_status == RESULT_INFRA_FAILED:
execution_status = EXECUTION_INFRA_FAILED
reason_code = usage_reason or REASON_PROVIDER_FAILURE
failure = {"kind": "provider", "reason_code": reason_code}
# An internal lifecycle error is a RUNTIME failure — the same kind the
# host-fallback prefix table below already assigns to this exact
# terminal text; calling it a provider failure made the two paths of
# this one function contradict each other.
failure_kind = "runtime" if reason_code == REASON_TASK_EXCEPTION else "provider"
failure = {"kind": failure_kind, "reason_code": reason_code}
# The overflow salvage keeps `llm_api_error`; a waited-out outage or the unknown
# no-resend fence may leave the same sticky kind behind under its own reason code,
# and the published projection must not contradict the terminal that chose it.

View file

@ -215,15 +215,52 @@ def reconcile_terminal(drive_root: Any, task_id: str) -> List[str]:
answered block."""
stamp = utc_now_iso()
expired: List[str] = []
terminal_quizzes: List[str] = []
def _mutator(quizzes: Dict[str, Dict[str, Any]]) -> Any:
for key, block in quizzes.items():
if str(block.get("state") or STATE_OPEN) == STATE_OPEN:
state = str(block.get("state") or STATE_OPEN)
if state == STATE_OPEN:
block.update({"state": STATE_EXPIRED_TERMINAL, "reconciled_at": stamp})
expired.append(str(key))
terminal_quizzes.append(str(key))
elif state == STATE_EXPIRED_TERMINAL:
# A previous call may have committed quiz expiry before the
# paired task-result repair failed. Keep the second pass
# idempotent so it can close the wait without a new quiz event.
terminal_quizzes.append(str(key))
return True if expired else _KEEP
_mutate_projection(drive_root, task_id, _mutator)
if terminal_quizzes:
from ouroboros.task_results import (
require_writable_task_result_schema,
stamp_task_result_schema,
)
def _close_owner_wait(current: Dict[str, Any]) -> Optional[Dict[str, Any]]:
require_writable_task_result_schema(current)
wait = current.get("owner_wait")
if not isinstance(wait, dict) or str(wait.get("state") or "") != "waiting":
return None
quiz_id = str(wait.get("quiz_id") or "")
if quiz_id not in terminal_quizzes:
return None
updated = dict(current)
updated["owner_wait"] = {
**wait,
"state": STATE_EXPIRED_TERMINAL,
"reconciled_at": stamp,
}
return stamp_task_result_schema(updated)
# The quiz and its waiting continuation share the same task-result
# authority. Close the paired wait after the quiz projection so a
# terminal task cannot replay as both expired and still waiting.
update_json_locked(
_quiz_result_path(drive_root, task_id),
_close_owner_wait,
)
return expired

View file

@ -389,7 +389,21 @@ def _is_stale_orphan_running_task(
task_id: str,
result: Dict[str, Any],
events_index: Optional[_EventsTailIndex] = None,
queue_snapshot: Optional[Dict[str, Any]] = None,
) -> bool:
# Direct-chat actors are deliberately absent from PENDING/RUNNING. The
# process-local registry is the authoritative owner for that execution;
# queue snapshots and pooled worker_boot rows cannot prove it dead.
try:
from supervisor.active_activity import get_direct_activity_registry
if get_direct_activity_registry().get(str(task_id or "")) is not None:
return False
except Exception:
# This helper is also imported by worker-side readers where the server's
# direct registry is not available. Absence of that optional observation
# is not itself evidence of liveness, so retain the existing pooled path.
pass
status = str(result.get("status") or "").lower()
# ``interrupted`` is the transient pre-requeue marker (A.11): a record still
# carrying it with no queued retry after a worker restart is the same orphan
@ -415,6 +429,22 @@ def _is_stale_orphan_running_task(
pass
if heartbeat and time.time() - heartbeat < _ORPHAN_RUNNING_GRACE_SECONDS:
return False
# A stale/missing snapshot cannot prove that a pooled owner is gone (the
# GR7-1a polarity ``task_has_live_queue_ownership`` already uses). The
# destructive reconciler must keep the old row until a fresh snapshot or a
# separate positive recovery fact exists. A batch caller passes the
# snapshot it already read, like ``events_index`` above.
try:
snapshot = (
queue_snapshot if isinstance(queue_snapshot, dict)
else _load_queue_snapshot(pathlib.Path(drive_root))
)
if snapshot.get("_snapshot_missing") or snapshot.get("_snapshot_invalid"):
return False
if _snapshot_is_stale(snapshot):
return False
except Exception:
return False
if events_index is None:
events_index = _EventsTailIndex(pathlib.Path(drive_root))
latest_task_event = max(heartbeat, events_index.latest_event_ts(task_id))
@ -737,7 +767,8 @@ def effective_task_result(
parent_status = str(merged.get("status") or "").lower()
if parent_status not in FINAL_STATUSES:
queue_status, queue_task = _queue_task_status(_load_queue_snapshot(pathlib.Path(drive_root)), task_id)
queue_snapshot = _load_queue_snapshot(pathlib.Path(drive_root))
queue_status, queue_task = _queue_task_status(queue_snapshot, task_id)
if queue_status and queue_status != "unknown":
merged["status"] = _merge_queue_status(parent_status, queue_status)
for key in (
@ -769,7 +800,10 @@ def effective_task_result(
bundle,
"task ended before artifact finalization",
)
elif _is_stale_orphan_running_task(pathlib.Path(drive_root), task_id, merged, _events_index):
elif _is_stale_orphan_running_task(
pathlib.Path(drive_root), task_id, merged, _events_index,
queue_snapshot=queue_snapshot,
):
orphan_reason = (
"interrupted_retry_lost"
if parent_status == STATUS_INTERRUPTED

View file

@ -189,6 +189,14 @@ def _handle_schedule_followup(ctx: ToolContext, **params) -> str:
)
metadata_src = getattr(ctx, "task_metadata", None)
root_task_id = metadata_src.get("root_task_id") if isinstance(metadata_src, dict) else None
project_id = str(
getattr(ctx, "project_id", "")
or (metadata_src.get("project_id") if isinstance(metadata_src, dict) else "")
or ""
).strip()
source_chat_id = getattr(ctx, "current_chat_id", None)
if source_chat_id in (None, "") and isinstance(metadata_src, dict):
source_chat_id = metadata_src.get("chat_id")
record = {
"id": f"followup-{task_id}-{uuid.uuid4().hex[:6]}",
"name": f"Follow-up of task {task_id}",
@ -202,6 +210,7 @@ def _handle_schedule_followup(ctx: ToolContext, **params) -> str:
"text": objective,
"description": objective,
**({"context": context} if context else {}),
**({"project_id": project_id} if project_id else {}),
"metadata": {
"source": FOLLOWUP_SOURCE,
"origin_task_id": task_id,
@ -209,6 +218,7 @@ def _handle_schedule_followup(ctx: ToolContext, **params) -> str:
# to task_id, never become the literal string "None".
"origin_root_task_id": str(root_task_id or "") or task_id,
},
**({"chat_id": source_chat_id} if source_chat_id not in (None, "") else {}),
},
}
presence = metadata_src.get("presence") if isinstance(metadata_src, dict) else None

View file

@ -255,7 +255,7 @@ def _task_from_schedule(record: Dict[str, Any]) -> Dict[str, Any]:
"delegation_role": "root",
"metadata": metadata,
}
for key in ("attachments", "context", "expected_output", "constraints", "deadline_at"):
for key in ("attachments", "context", "expected_output", "constraints", "deadline_at", "project_id"):
if key in template:
task[key] = template[key]
allowed_resources = normalize_allowed_resources(template.get("allowed_resources") or metadata.get("allowed_resources") or {})

View file

@ -525,7 +525,10 @@ def reap_orphaned_workers() -> int:
@_serialized_worker_lifecycle
def kill_workers_for_update(*, result_reason: str, terminal_status: str = "interrupted") -> List[str]:
def kill_workers_for_update(
*, result_reason: str, terminal_status: str = "interrupted",
preserve_running_task_ids: Optional[set[str]] = None,
) -> List[str]:
"""Stop the current pool and return anything whose death could not be proven."""
with _queue_lock:
fenced = list(_pool().WORKERS.values())
@ -536,6 +539,7 @@ def kill_workers_for_update(*, result_reason: str, terminal_status: str = "inter
terminal_status=terminal_status,
disable_reason="managed_update",
preserve_pending=True,
preserve_running_task_ids=set(preserve_running_task_ids or ()),
)
if kill_ok is False:
teardown_error = "teardown:queue_snapshot_persist_failed"

View file

@ -89,6 +89,28 @@ def test_observation_names_resolved_retry_and_discloses_a_later_target_change(tm
assert widened["observation"]["matches_cancel_target"] is False
def test_agent_cancel_of_a_foreign_task_records_origin_without_a_parent_decision(tmp_path):
"""A cancel outside the caller's own lineage must not discharge the REAL
parent's disposition duty (D#7). ``requested_by`` is the parent-decision
trigger, so the initiator fact rides ``source`` + ``request_origin``."""
write_task_result(tmp_path, "stranger", "running",
parent_task_id="someone-else", root_task_id="someone-else",
delegation_role="subagent")
_queue(tmp_path, "stranger")
response = _cancel_task(_caller(tmp_path), "stranger", "not mine but stop it")
assert response.startswith("Cancel requested:")
intent = cancel_intents.active_intent(tmp_path, "stranger")
assert intent["requested_by"] == ""
assert intent["source"] == "agent_tool"
assert intent["observation"]["request_origin"] == {
"kind": "agent_task", "task_id": "parent"}
from supervisor.cancel_publication import _intent_outcome_fields
fields = _intent_outcome_fields(intent)
assert "parent_decision" not in fields
assert fields["cancel_observation"] == intent["observation"]
def test_completed_child_keeps_its_result(tmp_path):
write_task_result(tmp_path, "child", "completed", parent_task_id="parent", root_task_id="parent", delegation_role="subagent", result="finished work")
atomic_write_json(tmp_path / "state" / "queue_snapshot.json", {"ts": utc_now_iso(), "running": [], "pending": []})
@ -121,3 +143,12 @@ def test_http_observation_records_transport_without_inventing_owner_identity(tmp
assert intent["observation"]["request_origin"] == {
"kind": "http_client", "source": "http_cascade" if cascade else "http_single"}
assert "delegated_execution" not in intent["observation"]
# The durable origin fact rides ``source`` + ``observation.request_origin``.
# ``requested_by`` is NOT a display label: any non-empty value is the
# PARENT-DECISION trigger, so naming the owner here would stamp
# ``parent_decision=cancelled`` on a child whose real parent decided
# nothing, and the parent's forced-finalization would treat that child
# result as already dispositioned (D#7).
from supervisor.cancel_publication import _intent_outcome_fields
assert "parent_decision" not in _intent_outcome_fields(intent)
assert intent["source"] == ("http_cascade" if cascade else "http_single")

View file

@ -304,11 +304,12 @@ def test_loop_outcome_distinguishes_success_empty_and_provider_failure():
runtime_error = derive_loop_outcome(
"⚠️ Error during processing: RuntimeError: boom",
{"rounds": 1},
{"rounds": 1, "execution_status": "infra_failed", "reason_code": "task_exception"},
{"tool_calls": []},
)
assert runtime_error["outcome_axes"]["execution"]["status"] == EXECUTION_INFRA_FAILED
assert runtime_error["reason_code"] == "task_exception"
assert runtime_error["failure"]["kind"] == "runtime"
deep_unavailable = derive_loop_outcome(
"❌ Deep self-review unavailable: no key",
@ -1312,3 +1313,78 @@ def test_refreshing_an_omitted_ledger_stub_is_an_identity():
assert refresh_verification_ledger_artifacts(
dict(stub), {"status": status, "artifacts": [], "errors": []},
) == stub
def test_task_exception_publishes_the_loop_s_accumulated_evidence(monkeypatch, tmp_path):
"""The outer agent catch owns the terminal projection, not the evidence.
A lifecycle failure after real rounds must publish the loop's accumulated
trace and usage — never a ``0 calls`` projection built from the untouched
pre-loop defaults — and must not call an internal error a provider failure.
"""
from ouroboros import agent as agent_module
from ouroboros import agent_task_pipeline
from ouroboros.agent import Env, OuroborosAgent
from ouroboros.task_results import STATUS_FAILED, load_task_result
repo, drive = tmp_path / "repo", tmp_path / "drive"
repo.mkdir()
drive.mkdir()
monkeypatch.setattr(OuroborosAgent, "_log_worker_boot_once", lambda self: None)
monkeypatch.setattr(agent_module, "build_llm_messages", lambda **_kwargs: ([], {}))
# The durable result is written before post-task cognition starts; the
# reflection thread is not this seam's subject.
monkeypatch.setattr(
agent_task_pipeline, "_run_post_task_processing_async", lambda *_a, **_kw: None)
def die_after_real_work(**_kwargs):
# Exactly what ``run_llm_loop`` attaches on an unexpected exit: the SAME
# in-memory accumulators the loop was filling.
exc = RuntimeError("owner wait refused a terminal continuation")
exc._ouroboros_loop_usage = {
"rounds": 7, "prompt_tokens": 4321, "completion_tokens": 210,
"execution_id": "exec_lifecycle_failure",
}
exc._ouroboros_loop_trace = {
"reasoning_notes": ["planned the edit"],
"tool_calls": [{
"tool": "write_file", "tool_call_id": "call-1", "result": "ok",
"trace_ref": {"call_id": "tool_write_file_1"},
}],
}
raise exc
monkeypatch.setattr(agent_module, "run_llm_loop", die_after_real_work)
agent = OuroborosAgent(Env(repo_dir=repo, drive_root=drive))
events = agent._handle_task_scoped({
"id": "lifecycle-fail", "type": "task", "chat_id": 1, "text": "do it",
"drive_root": str(drive), "budget_drive_root": str(drive),
})
stored = load_task_result(drive, "lifecycle-fail")
assert stored["status"] == STATUS_FAILED
assert stored["reason_code"] == "task_exception"
# The published trace counts the call that really happened.
assert stored["trace_summary"].startswith("## Tool trace (1 calls")
assert stored["trace_refs"]["execution_id"] == "exec_lifecycle_failure"
assert [ref["call_id"] for ref in stored["trace_refs"]["tool_call_refs"]] == [
"tool_write_file_1"]
# The loop's own tally rides the honest loop plane; an internal lifecycle
# error is a runtime failure, not a provider one.
assert stored["loop_outcome"]["usage"]["total_rounds"] == 7
assert stored["loop_outcome"]["usage"]["prompt_tokens"] == 4321
assert stored["loop_outcome"]["usage"]["completion_tokens"] == 210
execution = stored["outcome_axes"]["execution"]
assert execution["status"] == EXECUTION_INFRA_FAILED
assert execution["failure"] == {"kind": "runtime", "reason_code": "task_exception"}
# The original exception stays the evidence of what failed.
assert "RuntimeError: owner wait refused a terminal continuation" in stored["result"]
error_events = [
json.loads(line)
for line in (drive / "logs" / "events.jsonl").read_text(encoding="utf-8").splitlines()
if line.strip()
]
task_error = next(row for row in error_events if row.get("type") == "task_error")
assert "owner wait refused a terminal continuation" in task_error["error"]
assert "die_after_real_work" in task_error["traceback"]
assert any(event.get("type") == "task_done" for event in events)

View file

@ -84,6 +84,47 @@ def test_structural_expiry_flips_open_only(tmp_path):
assert reconcile_terminal(tmp_path, "t1") == []
def test_structural_expiry_closes_the_paired_owner_wait(tmp_path):
record_asked(tmp_path, "t1", quiz_id="q1", question="?", options=["A"])
path = _result_path(tmp_path, "t1")
row = json.loads(path.read_text())
row["owner_wait"] = {"quiz_id": "q1", "wait_id": "w1", "state": "waiting"}
path.write_text(json.dumps(row))
assert reconcile_terminal(tmp_path, "t1") == ["q1"]
updated = json.loads(path.read_text())
assert updated["owner_quiz"]["q1"]["state"] == STATE_EXPIRED_TERMINAL
assert updated["owner_wait"]["state"] == STATE_EXPIRED_TERMINAL
assert updated["owner_wait"]["reconciled_at"]
def test_structural_expiry_repairs_wait_after_partial_pair_write_failure(tmp_path, monkeypatch):
import ouroboros.owner_quiz as owner_quiz
from ouroboros.utils import update_json_locked as real_update_json_locked
record_asked(tmp_path, "t1", quiz_id="q1", question="?", options=["A"])
path = _result_path(tmp_path, "t1")
row = json.loads(path.read_text())
row["owner_wait"] = {"quiz_id": "q1", "wait_id": "w1", "state": "waiting"}
path.write_text(json.dumps(row))
calls = {"count": 0}
def fail_pair_write(*args, **kwargs):
calls["count"] += 1
if calls["count"] == 2:
raise OSError("simulated paired projection failure")
return real_update_json_locked(*args, **kwargs)
monkeypatch.setattr(owner_quiz, "update_json_locked", fail_pair_write)
with pytest.raises(OSError):
reconcile_terminal(tmp_path, "t1")
assert json.loads(path.read_text())["owner_quiz"]["q1"]["state"] == STATE_EXPIRED_TERMINAL
monkeypatch.setattr(owner_quiz, "update_json_locked", real_update_json_locked)
assert reconcile_terminal(tmp_path, "t1") == []
assert json.loads(path.read_text())["owner_wait"]["state"] == STATE_EXPIRED_TERMINAL
def test_projection_survives_concurrent_result_fields(tmp_path):
record_asked(tmp_path, "t1", quiz_id="q1", question="?", options=["A"] * 2)
path = _result_path(tmp_path, "t1")

View file

@ -320,6 +320,43 @@ def test_schedule_followup_root_id_falls_back_to_task_id_never_the_string_none(t
assert record["task"]["metadata"]["origin_root_task_id"] == "root-3"
def test_schedule_followup_preserves_source_project_and_chat(tmp_path):
ctx = _ctx(tmp_path, task_id="project-task")
ctx.project_id = "memory-atlas"
ctx.current_chat_id = 233966548
assert _followup(ctx).startswith("FOLLOWUP_SCHEDULED")
from supervisor.queue import list_scheduled_tasks
record = list_scheduled_tasks(pathlib.Path(tmp_path / "data").resolve())["tasks"][0]
assert record["task"]["chat_id"] == 233966548
assert record["task"]["project_id"] == "memory-atlas"
from supervisor.queue_schedules import _task_from_schedule
from ouroboros.project_facts import resolve_project_id
queued = _task_from_schedule(record)
assert queued["project_id"] == "memory-atlas"
assert resolve_project_id(queued) == "memory-atlas"
assert queued["chat_id"] == 233966548
def test_schedule_followup_of_an_unscoped_task_invents_no_project_address(tmp_path):
"""Preserving a source address must not become a new addressing policy: an
unscoped task's follow-up keeps the existing owner-chat default."""
from supervisor.queue import list_scheduled_tasks
from supervisor.queue_schedules import _task_from_schedule
from ouroboros.project_facts import resolve_project_id
assert _followup(_ctx(tmp_path, task_id="plain-task")).startswith("FOLLOWUP_SCHEDULED")
record = list_scheduled_tasks(pathlib.Path(tmp_path / "data").resolve())["tasks"][0]
assert "project_id" not in record["task"]
assert "chat_id" not in record["task"]
queued = _task_from_schedule(record)
assert resolve_project_id(queued) == ""
assert queued["chat_id"] == 0 # the existing owner_chat_id default, unchanged
# ------------------------------------------------- gateway + digest + queue GC

View file

@ -1,6 +1,10 @@
import json
import os
import threading
from types import SimpleNamespace
import pytest
def _stop_restart_watcher(server):
"""Stop the restart watcher ``server.main()`` starts (an unnamed daemon polling
@ -90,6 +94,147 @@ def test_pre_transaction_update_quiesce_also_preserves_pending(monkeypatch):
assert server._managed_update_pending_kwargs() == {"preserve_pending": True}
def test_managed_update_quiesce_passes_owner_wait_handoffs(monkeypatch, tmp_path):
import ouroboros.gateway.control as control
import ouroboros.owner_wait as owner_wait
import ouroboros.delegate_recovery as delegate_recovery
import supervisor.git_ops as git_ops
import supervisor.workers as workers
captured = {}
monkeypatch.setattr(git_ops, "DRIVE_ROOT", tmp_path)
monkeypatch.setattr(workers, "close_repo_writer_admission", lambda _reason: None)
monkeypatch.setattr(workers, "drain_repo_writers", lambda: [])
monkeypatch.setattr(owner_wait, "prepare_owner_wait_handoffs", lambda *_args: {"waiting"})
monkeypatch.setattr(
delegate_recovery,
"prepare_planned_restart_handoffs",
lambda *_args, **kwargs: set(kwargs["additional_task_ids"]),
)
def fake_kill(**kwargs):
captured.update(kwargs)
return ["worker:still-running"]
monkeypatch.setattr(workers, "kill_workers_for_update", fake_kill)
blockers = control._quiesce_repo_writers("smart")
assert blockers == ["worker:still-running"]
assert captured["preserve_running_task_ids"] == {"waiting"}
def test_manual_rollback_quiesce_keeps_cancellation_semantics(monkeypatch, tmp_path):
"""A rollback returns the tree to an OLDER runtime, so no owner wait may be
handed to it: the pool stop keeps its ordinary interrupt/cancel semantics."""
import ouroboros.gateway.control as control
import ouroboros.delegate_recovery as delegate_recovery
import ouroboros.owner_wait as owner_wait
import supervisor.git_ops as git_ops
import supervisor.workers as workers
captured = {}
monkeypatch.setattr(git_ops, "DRIVE_ROOT", tmp_path)
monkeypatch.setattr(workers, "close_repo_writer_admission", lambda _reason: None)
monkeypatch.setattr(workers, "drain_repo_writers", lambda: [])
monkeypatch.setattr(
owner_wait, "prepare_owner_wait_handoffs",
lambda *_args: pytest.fail("a rollback must not park an owner wait"),
)
monkeypatch.setattr(
delegate_recovery, "prepare_planned_restart_handoffs",
lambda *_args, **_kwargs: pytest.fail("a rollback must not prepare a handoff"),
)
monkeypatch.setattr(workers, "kill_workers_for_update", lambda **kwargs: captured.update(kwargs) or [])
monkeypatch.setattr(
"ouroboros.tools.services.kill_all_services", lambda *_a, **_kw: [])
monkeypatch.setattr(
"ouroboros.process_custody.quiesce_custodied_services", lambda *_a: (True, []))
assert control._quiesce_repo_writers("manual_rollback") == []
assert captured["preserve_running_task_ids"] == set()
assert captured["terminal_status"] == "interrupted"
def test_failed_owner_wait_handoff_blocks_the_update_instead_of_killing_the_pool(
monkeypatch, tmp_path,
):
"""Custody preparation is the gate: if it cannot park the wait, the pool is
never stopped and repo-writer admission re-opens."""
import ouroboros.gateway.control as control
import ouroboros.owner_wait as owner_wait
import supervisor.git_ops as git_ops
import supervisor.workers as workers
reopened = []
monkeypatch.setattr(git_ops, "DRIVE_ROOT", tmp_path)
monkeypatch.setattr(workers, "close_repo_writer_admission", lambda _reason: None)
monkeypatch.setattr(workers, "drain_repo_writers", lambda: [])
monkeypatch.setattr(workers, "open_repo_writer_admission", lambda: reopened.append(True))
def refuse(*_args):
raise OSError("owner wait source bytes are unreadable")
monkeypatch.setattr(owner_wait, "prepare_owner_wait_handoffs", refuse)
monkeypatch.setattr(
workers, "kill_workers_for_update",
lambda **_kwargs: pytest.fail("the pool must not stop after a failed handoff"),
)
blockers = control._quiesce_repo_writers("smart")
assert blockers == ["owner_wait_handoff:OSError: owner wait source bytes are unreadable"]
assert reopened == [True]
def test_managed_update_restart_arms_prepared_transaction_for_direct_reexec(
monkeypatch, tmp_path,
):
import ouroboros.delegate_recovery as delegate_recovery
import ouroboros.gateway.control as control
import supervisor.git_ops as git_ops
prepared = {
"transaction_id": "tx-owner-wait",
"status": "prepared",
"supervisor_pid": os.getpid(),
}
tx_dir = tmp_path / "state" / "delegate_recovery_transactions"
tx_dir.mkdir(parents=True)
(tx_dir / "active.json").write_text(
'{"transaction_id":"tx-owner-wait"}', encoding="utf-8",
)
(tx_dir / "tx-owner-wait.json").write_text(
json.dumps(prepared), encoding="utf-8",
)
monkeypatch.setattr(git_ops, "DRIVE_ROOT", tmp_path)
monkeypatch.delenv(delegate_recovery.PLANNED_RESTART_TRANSACTION_ENV, raising=False)
request = type("Request", (), {"app": type("App", (), {"state": type("State", (), {
"request_restart": lambda _self, owner=False: None,
})()})()})()
response = control._restart_response(request, strategy="auto_merge", plan={})
assert response.status_code == 200
assert os.environ[delegate_recovery.PLANNED_RESTART_TRANSACTION_ENV] == "tx-owner-wait"
def test_failed_managed_update_restart_disarms_transaction_token(monkeypatch, tmp_path):
import ouroboros.delegate_recovery as delegate_recovery
import ouroboros.gateway.control as control
monkeypatch.setenv(delegate_recovery.PLANNED_RESTART_TRANSACTION_ENV, "stale-tx")
request = type("Request", (), {"app": type("App", (), {"state": type("State", (), {
"request_restart": None,
})()})()})()
response = control._restart_response(request, strategy="auto_merge", plan={})
assert response.status_code == 200
assert json.loads(response.body)["status"] == "restart_required"
assert delegate_recovery.PLANNED_RESTART_TRANSACTION_ENV not in os.environ
def test_ordinary_restart_disarms_orphan_update_intent(monkeypatch):
import server
import supervisor.git_ops as git_ops

View file

@ -1607,7 +1607,10 @@ def test_effective_status_repairs_stale_running_infra_failure_when_queue_empty(t
},
)
(tmp_path / "state").mkdir(exist_ok=True)
(tmp_path / "state" / "queue_snapshot.json").write_text('{"pending": [], "running": []}', encoding="utf-8")
(tmp_path / "state" / "queue_snapshot.json").write_text(
'{"ts": "2027-01-15T08:00:00+00:00", "pending": [], "running": []}',
encoding="utf-8",
)
effective = load_effective_task_result(tmp_path, "providerfail")
@ -1660,7 +1663,10 @@ def test_effective_status_repairs_orphan_running_after_worker_restart(tmp_path,
},
)
(tmp_path / "state").mkdir(exist_ok=True)
(tmp_path / "state" / "queue_snapshot.json").write_text('{"pending": [], "running": []}', encoding="utf-8")
(tmp_path / "state" / "queue_snapshot.json").write_text(
'{"ts": "2027-01-15T08:00:00+00:00", "pending": [], "running": []}',
encoding="utf-8",
)
events = tmp_path / "logs" / "events.jsonl"
append_jsonl(events, {"ts": "2026-05-28T00:00:01+00:00", "type": "llm_round", "task_id": "cc4db6fa"})
append_jsonl(events, {"ts": "2026-05-28T00:00:02+00:00", "type": "worker_boot"})
@ -1696,7 +1702,10 @@ def test_reconcile_durably_finalizes_orphaned_running_task(tmp_path, monkeypatch
result="Task is running.", ts="2026-05-28T00:00:00+00:00",
)
(tmp_path / "state").mkdir(exist_ok=True)
(tmp_path / "state" / "queue_snapshot.json").write_text('{"pending": [], "running": []}', encoding="utf-8")
(tmp_path / "state" / "queue_snapshot.json").write_text(
'{"ts": "2027-01-15T08:00:00+00:00", "pending": [], "running": []}',
encoding="utf-8",
)
events = tmp_path / "logs" / "events.jsonl"
append_jsonl(events, {"ts": "2026-05-28T00:00:01+00:00", "type": "llm_round", "task_id": "orphan1"})
append_jsonl(events, {"ts": "2026-05-28T00:00:02+00:00", "type": "worker_boot"})
@ -3407,6 +3416,80 @@ def test_handle_text_response_keeps_full_reasoning_note():
assert updated["reasoning_notes"] == [content]
def _orphan_shaped_running_task(tmp_path, task_id, *, snapshot_ts):
"""The exact fixture the pooled orphan reconciler accepts as proof of death:
an aged ``running`` row, an empty queue snapshot and a LATER worker boot."""
from ouroboros.task_results import STATUS_RUNNING, write_task_result
from ouroboros.utils import append_jsonl
write_task_result(
tmp_path, task_id, STATUS_RUNNING,
result="Task is running.", ts="2026-05-28T00:00:00+00:00",
)
(tmp_path / "state").mkdir(exist_ok=True)
(tmp_path / "state" / "queue_snapshot.json").write_text(
json.dumps({"ts": snapshot_ts, "pending": [], "running": []}), encoding="utf-8",
)
events = tmp_path / "logs" / "events.jsonl"
append_jsonl(events, {"ts": "2026-05-28T00:00:01+00:00", "type": "llm_round", "task_id": task_id})
append_jsonl(events, {"ts": "2026-05-28T00:00:02+00:00", "type": "worker_boot"})
def test_orphan_reconcile_never_terminalizes_a_live_direct_activity(tmp_path, monkeypatch):
"""A direct-chat actor is deliberately absent from PENDING/RUNNING, so a
FOREIGN ``worker_boot`` plus a fresh empty snapshot cannot prove it dead."""
from ouroboros.task_results import STATUS_RUNNING, load_task_result
from ouroboros.task_status import (
load_effective_task_result, reconcile_orphaned_running_tasks,
)
from supervisor import active_activity
from supervisor.active_activity import (
DirectActivityRegistry, get_direct_activity_registry,
)
monkeypatch.setattr(time, "time", lambda: 1_800_000_000.0)
# Reliable fixture isolation for a process-global (DEVELOPMENT.md, parallel
# pass): monkeypatch reverses exactly this singleton, so the real
# ``get_direct_activity_registry`` seam is still the one under test.
monkeypatch.setattr(active_activity, "_DIRECT_ACTIVITY_REGISTRY", DirectActivityRegistry())
_orphan_shaped_running_task(tmp_path, "direct-live", snapshot_ts="2027-01-15T08:00:00+00:00")
registry = get_direct_activity_registry()
registry.register("direct-live", chat_id=1)
assert reconcile_orphaned_running_tasks(tmp_path) == 0
assert load_effective_task_result(tmp_path, "direct-live")["status"] == STATUS_RUNNING
assert load_task_result(tmp_path, "direct-live")["status"] == STATUS_RUNNING
# Non-vacuous: the SAME evidence is an ordinary orphan once the actor is gone,
# so the guard — not the fixture — is what kept the live row alive.
registry.unregister("direct-live")
assert reconcile_orphaned_running_tasks(tmp_path) == 1
healed = load_task_result(tmp_path, "direct-live")
assert healed["reason_code"] == "orphaned_running_after_worker_restart"
def test_orphan_reconcile_does_not_use_a_stale_snapshot_as_death_proof(tmp_path, monkeypatch):
"""GR7-1a polarity for the DESTRUCTIVE reconciler: an out-of-date snapshot
cannot prove a pooled owner is gone, but a fresh one still can."""
from ouroboros.task_results import STATUS_RUNNING, load_task_result
from ouroboros.task_status import reconcile_orphaned_running_tasks
monkeypatch.setattr(time, "time", lambda: 1_800_000_000.0)
_orphan_shaped_running_task(tmp_path, "pooled-live", snapshot_ts="2026-05-28T00:00:03+00:00")
assert reconcile_orphaned_running_tasks(tmp_path) == 0
assert load_task_result(tmp_path, "pooled-live")["status"] == STATUS_RUNNING
(tmp_path / "state" / "queue_snapshot.json").write_text(
json.dumps({"ts": "2027-01-15T08:00:00+00:00", "pending": [], "running": []}),
encoding="utf-8",
)
assert reconcile_orphaned_running_tasks(tmp_path) == 1
assert load_task_result(tmp_path, "pooled-live")["reason_code"] == (
"orphaned_running_after_worker_restart"
)
def test_request_restart_latches_reason_until_task_end(tmp_path, monkeypatch):
from ouroboros.tools import control as control_module
from ouroboros.tools import control_runtime

View file

@ -17,6 +17,7 @@ from types import SimpleNamespace
from ouroboros.gateway.tasks import api_task_get, api_tasks_list
from ouroboros.task_results import write_task_result
from ouroboros.utils import utc_now_iso
def _request(data, **params):
@ -41,9 +42,12 @@ def _write_raw(data, task_id, **fields):
def _seed_queue_snapshot(data):
# A FRESH ``ts``: an empty snapshot only proves a running row is orphaned
# while it is still current — an undated or out-of-date snapshot fails open
# toward liveness (GR7-1a), and the orphan projection would never run.
(data / "state").mkdir(parents=True, exist_ok=True)
(data / "state" / "queue_snapshot.json").write_text(
'{"pending": [], "running": []}', encoding="utf-8"
json.dumps({"ts": utc_now_iso(), "pending": [], "running": []}), encoding="utf-8"
)

View file

@ -536,6 +536,59 @@ def test_primary_round_dispatch_recovers_after_two_deaths(tmp_path, monkeypatch,
assert TRANSPORT_DEATHS_KEY not in usage
def _tool_round(name, args, call_id):
return (
{"role": "assistant", "content": "", "tool_calls": [{
"id": call_id, "type": "function",
"function": {"name": name, "arguments": json.dumps(args)},
}]},
{"prompt_tokens": 11, "completion_tokens": 2},
)
def test_unexpected_loop_error_carries_accumulated_evidence_to_owner_projection(
tmp_path, monkeypatch,
):
"""The outer agent catch must not turn a failed multi-round loop into 0 calls.
Round 1 runs the REAL tool executor, so ``llm_trace`` holds a recorded call
with its durable trace ref; round 2's executor then dies. The loop owns that
accumulated evidence, so it must reach the raiser instead of being replaced
by an empty projection.
"""
real_handle_tool_calls = loop_mod.handle_tool_calls
executed = {"count": 0}
def explode_after_the_first_batch(*args, **kwargs):
executed["count"] += 1
if executed["count"] == 1:
return real_handle_tool_calls(*args, **kwargs)
raise RuntimeError("tool executor crashed after the provider response")
monkeypatch.setattr(loop_mod, "handle_tool_calls", explode_after_the_first_batch)
llm = _ScriptedLLM(
_tool_round("write_file", {"root": "task_drive", "path": "a.txt", "content": "a"}, "call-1"),
_tool_round("write_file", {"root": "task_drive", "path": "b.txt", "content": "b"}, "call-2"),
)
with pytest.raises(RuntimeError) as caught:
run_llm_loop(**_loop_kwargs(tmp_path, llm, []))
# The ORIGINAL exception propagates: same object, same traceback, no
# wrapper that would hide where the lifecycle actually failed.
assert str(caught.value) == "tool executor crashed after the provider response"
assert [entry.name for entry in caught.traceback][-1] == "explode_after_the_first_batch"
usage = getattr(caught.value, "_ouroboros_loop_usage")
trace = getattr(caught.value, "_ouroboros_loop_trace")
assert usage["rounds"] == 2
assert usage["prompt_tokens"] == 22
assert usage["completion_tokens"] == 4
assert TRANSPORT_DEATHS_KEY not in usage
# Durable evidence for the completed batch survives the failure.
assert [call["tool_call_id"] for call in trace["tool_calls"]] == ["call-1"]
assert trace["tool_calls"][0]["trace_ref"]["call_id"]
@pytest.mark.parametrize("turn_flag", [None, "is_direct_chat", "is_ephemeral_turn"])
def test_counter_survives_the_wait_episodes_free_redial_of_the_same_round(tmp_path, monkeypatch, no_sleep, turn_flag):
"""death → released ConnectError → wait episode → free redial → death →