diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index 8d9d029b7..faaaad7f7 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -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). 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.py ← Task orchestrator; the dispatch-note pair lives in `subagent_dispatch_notes.py` (same-name re-exports). The outer catch around `run_llm_loop` delegates terminal projection to `_task_exception_terminal`: `_LoopExitContext.attach_exception_evidence` retains the loop's own accumulated objects on the original exception, while missing captures are explicitly unknown. A failed cold-source read never supplies fabricated zero counters or unverified checkpoint bytes; the trace summary, loop usage, task metrics and post-task summary preserve that absence. 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) 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_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, including an answered quiz whose worker resume was not yet granted, while preserving the recorded answer — `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) @@ -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; 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 + │ ├── 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. The existing scheduler consumes a one-shot on the typed tombstoned-Project refusal, retaining its failed task and last error; deleting Projects and transient refusals remain retryable │ ├── 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. 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. +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. At the common direct re-exec seam, a valid managed update in the same restartable phases allowed by `_safe_restart_serialized` arms the already-prepared transaction through the existing one-shot environment handoff. This includes assisted updates whose waiting tasks are already PENDING; the same-PID successor acknowledges it (launcher mode observes exit code 42). No-resume flags suppress this token, and an aborted update's leftover recovery record alone cannot authorize a later manual Restart or rollback. A failed restart callback retains the existing disarm behavior. 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. diff --git a/docs/DEVELOPMENT.md b/docs/DEVELOPMENT.md index 2b6cd4594..7671dede5 100644 --- a/docs/DEVELOPMENT.md +++ b/docs/DEVELOPMENT.md @@ -1894,6 +1894,11 @@ owner, owed terminal delivery, cascade postconditions — lives in ARCHITECTURE plus execution root, task, attempt and seen ids. Read/parse/stat failure or torn data is not proof; never cache it. Check the in-memory incoming queue every tick. A changed source re-enters the full revocation-aware reader; no TTL or ACK in peek. +- Terminal quiz reconciliation closes the paired wait even if the answer arrived + before worker capacity was granted; keep the answer and source unchanged. A + failed loop without captured evidence reports unknown counts through the existing + summary/outcome/metrics producers. Never infer zero work or read an unverified + checkpoint to fill the gap (`tests/test_autonomy_review_fixes.py`). - Cancellation observations use `task_status.observe_cancellation_target` before the existing intent write. They name the resolved physical target, separate task-result update/start facts from queue freshness, and optionally include diff --git a/ouroboros/agent.py b/ouroboros/agent.py index 40968d6b6..45460af00 100644 --- a/ouroboros/agent.py +++ b/ouroboros/agent.py @@ -87,6 +87,46 @@ def _authority_source_terminal(refusal: Dict[str, Any]): return text, usage, {"reasoning_notes": ["authority_source_unavailable"], "tool_calls": []} +def _task_exception_terminal(env: Any, task: Dict[str, Any], exc: Exception, drive_logs: pathlib.Path): + """Project captured loop evidence, or its explicit absence, without recovery. + + A failed cold-source read is never permission to read unverified checkpoint + bytes as usage. The loop tally stays in loop_outcome; top-level money and + counters remain the existing ledger reconstruction's answer. + """ + captured_usage = getattr(exc, "_ouroboros_loop_usage", None) + captured_trace = getattr(exc, "_ouroboros_loop_trace", None) + usage = dict(captured_usage) if isinstance(captured_usage, dict) else {"loop_evidence_unavailable": True} + llm_trace = captured_trace if isinstance(captured_trace, dict) else { + "reasoning_notes": [], "tool_calls": [], "loop_evidence_unavailable": True, + } + usage.update(execution_status="infra_failed", reason_code="task_exception") + text = f"⚠️ Error during processing: {type(exc).__name__}: {exc}" + append_jsonl(drive_logs / "events.jsonl", { + "ts": utc_now_iso(), "type": "task_error", "task_id": task.get("id"), + "error": repr(exc), "traceback": truncate_for_log(traceback.format_exc(), 2000), + }) + 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 + + # Ephemeral decision turns leave no durable task result, including errors. + if not bool(task.get("_ephemeral_turn")): + loop_outcome = derive_loop_outcome(text, usage, llm_trace) + write_task_result( + env.drive_root, str(task.get("id") or ""), STATUS_FAILED, + result=text, reason_code="task_exception", loop_outcome=loop_outcome, + outcome_axes=loop_outcome.get("outcome_axes") or infra_failed_axes( + "task_exception", review_trigger="agent_exception"), + trace_summary=build_trace_summary(llm_trace), + trace_refs=loop_outcome.get("trace_refs") or collect_trace_refs(usage, llm_trace), + ) + except Exception: + log.debug("Failed to persist task exception projection", exc_info=True) + return text, usage, llm_trace + + def _sync_task_project_scope(task: Dict[str, Any], ctx: Any) -> None: project_id = str(getattr(ctx, "project_id", "") or "").strip() if project_id and not str(task.get("project_id") or "").strip(): @@ -926,53 +966,7 @@ class OuroborosAgent: # Empty events leave its queue slot/project owned until # supervisor cancellation kills and settles this task. return [] - tb = traceback.format_exc() - append_jsonl(drive_logs / "events.jsonl", { - "ts": utc_now_iso(), "type": "task_error", - "task_id": task.get("id"), "error": repr(e), - "traceback": truncate_for_log(tb, 2000), - }) - text = f"⚠️ Error during processing: {type(e).__name__}: {e}" - 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=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 + text, usage, llm_trace = _task_exception_terminal(self.env, task, e, drive_logs) try: from ouroboros.task_continuation import capture_review_continuation_from_state capture_review_continuation_from_state( diff --git a/ouroboros/agent_task_pipeline.py b/ouroboros/agent_task_pipeline.py index 939848b95..5dac71150 100644 --- a/ouroboros/agent_task_pipeline.py +++ b/ouroboros/agent_task_pipeline.py @@ -514,6 +514,8 @@ def emit_task_results( n_tool_calls = len(llm_trace.get("tool_calls", [])) n_tool_errors = sum(1 for tc in llm_trace.get("tool_calls", []) if isinstance(tc, dict) and tc.get("is_error")) + if llm_trace.get("loop_evidence_unavailable"): + n_tool_calls = n_tool_errors = None try: from supervisor.state import reconstruct_task_cost diff --git a/ouroboros/gateway/control.py b/ouroboros/gateway/control.py index 1de3c7977..180ff3bef 100644 --- a/ouroboros/gateway/control.py +++ b/ouroboros/gateway/control.py @@ -497,17 +497,6 @@ 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: diff --git a/ouroboros/loop.py b/ouroboros/loop.py index 9b98b64fc..3926456ff 100644 --- a/ouroboros/loop.py +++ b/ouroboros/loop.py @@ -642,15 +642,7 @@ def run_llm_loop( 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 + exit_ctx.attach_exception_evidence(exc) raise finally: # No stale active latch behind an in-process exit (a crash skips this frame, keeping the latch for recovery). diff --git a/ouroboros/loop_budget.py b/ouroboros/loop_budget.py index b40566156..abe435829 100644 --- a/ouroboros/loop_budget.py +++ b/ouroboros/loop_budget.py @@ -352,6 +352,19 @@ class _LoopExitContext: accumulated_usage: Dict[str, Any] llm_trace: Dict[str, Any] + def attach_exception_evidence(self, exc: Exception) -> None: + """The caller owns terminal projection; this loop owns its evidence. + + Keep the same in-memory objects on the original exception so a lifecycle + failure cannot erase a completed multi-round trace. A failed attachment + leaves the caller's explicit unknown projection, never a new exception. + """ + try: + setattr(exc, "_ouroboros_loop_usage", self.accumulated_usage) + setattr(exc, "_ouroboros_loop_trace", self.llm_trace) + except Exception: + log.debug("Loop exception evidence could not be attached", exc_info=True) + def _handle_budget_exceeded( exc: BudgetExceeded, diff --git a/ouroboros/outcomes.py b/ouroboros/outcomes.py index 60472f326..23f66c1e9 100644 --- a/ouroboros/outcomes.py +++ b/ouroboros/outcomes.py @@ -977,9 +977,9 @@ def _loop_usage_snapshot(usage: Dict[str, Any], resource_limit: Dict[str, Any]) round(float(usage["cost"]), 6) if usage.get("cost") is not None else None ), - "prompt_tokens": int(usage.get("prompt_tokens") or 0), - "completion_tokens": int(usage.get("completion_tokens") or 0), - "total_rounds": int(usage.get("rounds") or 0), + "prompt_tokens": None if usage.get("loop_evidence_unavailable") else int(usage.get("prompt_tokens") or 0), + "completion_tokens": None if usage.get("loop_evidence_unavailable") else int(usage.get("completion_tokens") or 0), + "total_rounds": None if usage.get("loop_evidence_unavailable") else int(usage.get("rounds") or 0), **({"resource_limit": resource_limit} if resource_limit else {}), } diff --git a/ouroboros/owner_quiz.py b/ouroboros/owner_quiz.py index 8033e7882..23b4b3b93 100644 --- a/ouroboros/owner_quiz.py +++ b/ouroboros/owner_quiz.py @@ -21,7 +21,7 @@ like the hurry projection): Structural expiry only (owner decision 30=A): a quiz dies with its author — ``reconcile_terminal`` runs on the task-done seam; there is no host TTL. -The writer mutates ONLY the ``owner_quiz`` key via ``update_json_locked`` +The writers mutate ``owner_quiz`` and its paired terminal ``owner_wait`` via ``update_json_locked`` (never ``write_task_result`` — its status-regression guard can drop the write), so concurrent terminal writers merge around it. """ @@ -224,10 +224,11 @@ def reconcile_terminal(drive_root: Any, task_id: str) -> List[str]: block.update({"state": STATE_EXPIRED_TERMINAL, "reconciled_at": stamp}) expired.append(str(key)) terminal_quizzes.append(str(key)) - elif state == STATE_EXPIRED_TERMINAL: + elif state in (STATE_EXPIRED_TERMINAL, STATE_ANSWERED): # 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. + # idempotent; an accepted answer can also await worker capacity + # when the task ends. Neither case rewrites the quiz's answer. terminal_quizzes.append(str(key)) return True if expired else _KEEP diff --git a/ouroboros/post_task_synthesis.py b/ouroboros/post_task_synthesis.py index 8fed971fa..99fa06b49 100644 --- a/ouroboros/post_task_synthesis.py +++ b/ouroboros/post_task_synthesis.py @@ -43,6 +43,8 @@ def _atp(): def build_trace_summary(llm_trace: dict) -> str: """Return a compact human-readable summary of tool calls and agent notes.""" + if llm_trace.get("loop_evidence_unavailable"): + return "## Tool trace (call count unknown)\nThe failed loop supplied no verified execution trace." tool_calls = llm_trace.get("tool_calls", []) or [] notes = llm_trace.get("reasoning_notes", []) or [] @@ -289,8 +291,9 @@ def _run_task_summary(env, llm, task, usage, llm_trace, drive_logs, review_evide task_id = str(task.get("id") or "unknown") canonical_root = pathlib.Path(task.get("budget_drive_root") or drive_logs.parent) summary_id = f"task-narrative:{task_id}" - n_tool_calls = len(llm_trace.get("tool_calls", []) or []) - rounds = int(usage.get("rounds") or 0) + n_tool_calls = None if llm_trace.get("loop_evidence_unavailable") else len(llm_trace.get("tool_calls", []) or []) + rounds = None if usage.get("loop_evidence_unavailable") else int(usage.get("rounds") or 0) + round_text = "round count unknown" if rounds is None else f"{rounds}r" cost_text = _synthesis_cost_text(usage) outcome_axes = normalize_outcome_axes(usage) reason_code = str(usage.get("reason_code") or "") @@ -316,11 +319,11 @@ def _run_task_summary(env, llm, task, usage, llm_trace, drive_logs, review_evide canonical_root, result_root, row, status=str(stored_result.get("status") or ""), ) # Skip LLM summary for trivial tasks. - if n_tool_calls == 0 and rounds <= 1: + if n_tool_calls in (None, 0) and (rounds is None or rounds <= 1): goal = _truncate_with_notice(task.get("text", ""), 200) summary_text = ( f"Task {task_id} ({task.get('type', 'user')}): " - f"{goal}. {rounds}r, {cost_text}." + project_thread_note_for_task(task) + f"{goal}. {round_text}, {cost_text}." + project_thread_note_for_task(task) ) _append_summary(summary_text) return @@ -335,7 +338,7 @@ def _run_task_summary(env, llm, task, usage, llm_trace, drive_logs, review_evide review_section = "(review evidence unavailable)" prompt = _TASK_SUMMARY_PROMPT.format( task_id=task_id, goal=goal or "(no goal text)", - task_type=task.get("type", "user"), rounds=rounds, + task_type=task.get("type", "user"), rounds="unknown" if rounds is None else rounds, cost_text=cost_text, usage_snapshot=_synthesis_usage_snapshot_text(usage), sealed_final=sealed_final_prompt_section(sealed_final), @@ -364,7 +367,7 @@ def _run_task_summary(env, llm, task, usage, llm_trace, drive_logs, review_evide log.warning("Task summary LLM call failed, using fallback", exc_info=True) summary_text = ( f"Task {task_id} ({task.get('type', 'user')}): " - f"{_truncate_with_notice(goal, 200)}. {rounds}r, {cost_text}." + f"{_truncate_with_notice(goal, 200)}. {round_text}, {cost_text}." ) if summary_text: summary_text += project_thread_note_for_task(task) diff --git a/ouroboros/reflection.py b/ouroboros/reflection.py index 028791944..2e99b8663 100644 --- a/ouroboros/reflection.py +++ b/ouroboros/reflection.py @@ -503,13 +503,13 @@ def generate_reflection( "task_id": task.get("id", ""), "task_type": str(task.get("type", "")), "goal": goal, - "rounds": int(usage_dict.get("rounds", 0)), + "rounds": None if usage_dict.get("loop_evidence_unavailable") else int(usage_dict.get("rounds", 0)), "cost_usd": ( round(float(usage_dict["cost"]), 4) if usage_dict.get("cost") is not None else None ), - "error_count": error_count, + "error_count": None if llm_trace.get("loop_evidence_unavailable") else error_count, "key_markers": markers, "review_evidence": review_evidence or {}, "reflection": reflection_text, diff --git a/ouroboros/server_restart.py b/ouroboros/server_restart.py index 8c3868971..5032149ab 100644 --- a/ouroboros/server_restart.py +++ b/ouroboros/server_restart.py @@ -19,6 +19,8 @@ from typing import Any from ouroboros.server_process import DATA_DIR, _owner_restart_requested, _restart_requested, log +_RESTARTABLE_UPDATE_PHASES = frozenset({"pending_boot_smoke", "applying_replace"}) + def _owned_live_task_ids(ctx: Any) -> list: """Every id this generation's cancel intent can address: pooled tasks, @@ -225,8 +227,7 @@ def _safe_restart_serialized(safe_restart_fn, *, reason: str, unsynced_policy: s "An update intent marker with no update transaction could not be removed; " "restart was deferred rather than applying an orphaned update." ) - allowed_phases = {"pending_boot_smoke", "applying_replace"} - if status == "valid" and str(tx.get("phase") or "") not in allowed_phases: + if status == "valid" and str(tx.get("phase") or "") not in _RESTARTABLE_UPDATE_PHASES: return False, "Managed update merge is still being resolved; restart was deferred." return safe_restart_fn(reason=reason, unsynced_policy=unsynced_policy) finally: diff --git a/ouroboros/size_ratchet_manifest.py b/ouroboros/size_ratchet_manifest.py index ef5e41356..2a4eb1668 100644 --- a/ouroboros/size_ratchet_manifest.py +++ b/ouroboros/size_ratchet_manifest.py @@ -214,6 +214,7 @@ BAND_PATHS = { "tests/test_terminal_durability_v664.py": "Entered the band from 974 lines: terminal durability coverage now pins retry-admission failure custody so an unpersisted terminal row cannot publish task_done or lose the retry marker.", "tests/test_timeout_policy.py": "Adaptive timeout and custody regression suite covers raw-deadline admission, explicit finalization reserve, transport bounds, and late-result reconciliation.", "tests/test_tool_result.py": "F3.1 typed-organ pins carried with the D02 organ (D04 entry 9): the closed code table, the one legacy-text adapter, the publish/sidecar seam and the meta-boundary contracts pin one organ in one suite; sibling suites (meta_boundaries, t46, classification differential) already hold the spill-over families.", + "tests/test_transport_death_retry.py": "Physical transport-repeat and exceptional-loop evidence tests share the same scripted provider and real tool-execution fixture; retain this coherent contract suite below the module cap.", "tests/test_tree_cost_ceiling.py": "Budget-rail coverage entered the band with the cache-split ownership, candidate-predicate, soft-landing, probe-confirmed stop and one-row-per-delegated-run regressions; one focused suite for the tree cost ceiling.", "tests/test_ui_smoke_project_continuity.py": "Playwright smoke of the Project continuity contracts (panel/Main re-homing, lifecycle rows, the Main-root project pointer): each test drives one end-to-end owner flow across both surfaces, so the cross-surface assertions cannot be split into smaller files without losing what they prove.", "tests/test_update_letter.py": "Entered the band from 969 lines: the update-letter contract gained the HTTP-200 body-error verdicts (a body overflow earns the same single Low retry, other body errors are typed provider failures) and the attempt ids of calls that raised (update-letter sprint 2026-09-04).", diff --git a/server.py b/server.py index 3bdae82aa..492807344 100644 --- a/server.py +++ b/server.py @@ -168,6 +168,21 @@ def _has_active_evolution_transaction() -> bool: def _restart_current_process(host: str, port: int) -> None: + # Every direct restart reaches this seam, including an assisted update whose + # native waits were already moved to PENDING before its resolver ran. + try: + from ouroboros.delegate_recovery import PLANNED_RESTART_TRANSACTION_ENV, arm_active_planned_restart_transaction + from ouroboros.server_restart import _RESTARTABLE_UPDATE_PHASES + from supervisor.update_merge import read_update_tx_strict + + if any((DATA_DIR / "state" / name).exists() for name in ("owner_restart_no_resume.flag", "panic_stop.flag")): + os.environ.pop(PLANNED_RESTART_TRANSACTION_ENV, None) + else: + status, tx = read_update_tx_strict() + if status == "valid" and tx.get("phase") in _RESTARTABLE_UPDATE_PHASES: + arm_active_planned_restart_transaction(DATA_DIR) + except Exception: + log.warning("Direct restart transaction could not be armed; continuation remains unconfirmed", exc_info=True) _restart_current_process_impl( host, port, repo_dir=REPO_DIR, log=log, owner_initiated=_owner_restart_requested.is_set(), diff --git a/supervisor/events_worker_reports.py b/supervisor/events_worker_reports.py index fb60bd08d..6078784e3 100644 --- a/supervisor/events_worker_reports.py +++ b/supervisor/events_worker_reports.py @@ -111,8 +111,8 @@ def _handle_task_metrics(evt: Dict[str, Any], ctx: Any) -> None: "task_id": str(evt.get("task_id") or ""), "task_type": str(evt.get("task_type") or ""), "duration_sec": round(float(evt.get("duration_sec") or 0.0), 3), - "tool_calls": int(evt.get("tool_calls") or 0), - "tool_errors": int(evt.get("tool_errors") or 0), + "tool_calls": None if evt.get("tool_calls") is None else int(evt["tool_calls"]), + "tool_errors": None if evt.get("tool_errors") is None else int(evt["tool_errors"]), "outcome_axes": normalize_outcome_axes(evt), "reason_code": str(evt.get("reason_code") or ""), } diff --git a/supervisor/queue_schedules.py b/supervisor/queue_schedules.py index fbdad1349..ed060f228 100644 --- a/supervisor/queue_schedules.py +++ b/supervisor/queue_schedules.py @@ -373,9 +373,12 @@ def check_scheduled_tasks() -> None: record["last_task_id"] = task["id"] record_scheduled_admission(task, admitted, record) if trigger_type == "once": - if not (isinstance(admitted, dict) and admitted.get("_admission_blocked")): - # Consumed ONLY when admission succeeded (durable receipt, never re-fired); a - # refused admission left the record enabled with last_error → next tick retries. + refused = isinstance(admitted, dict) and admitted.get("_admission_blocked") + permanent = (refused == "project_routing_fence" + and admitted.get("_project_lifecycle") == "tombstoned") + if not refused or permanent: + # A consumed receipt includes a permanent target refusal; + # keep its failed task and last_error. Transient refusals retry. record["enabled"] = False record["completed_at"] = now.isoformat() record["next_run_at"] = "" diff --git a/tests/test_autonomy_review_fixes.py b/tests/test_autonomy_review_fixes.py new file mode 100644 index 000000000..ae0718794 --- /dev/null +++ b/tests/test_autonomy_review_fixes.py @@ -0,0 +1,234 @@ +"""Review regressions at the existing restart, wait and schedule owners.""" + +import json +import os +import threading +from types import SimpleNamespace + +import pytest + +from ouroboros import delegate_recovery, owner_quiz, owner_wait +from ouroboros.task_results import load_task_result, write_task_result +from tests.test_owner_wait_restart import restart_case as _restart_case, restore_stale_snapshot + +restart_case = _restart_case + + +@pytest.mark.parametrize("terminal", ["cancelled", "failed"]) +def test_answered_quiz_closes_wait_before_resume_capacity_is_granted(tmp_path, terminal): + write_task_result(tmp_path, "t", "running") + owner_quiz.record_asked(tmp_path, "t", quiz_id="q", question="Proceed?", options=["Yes"]) + wait = {"quiz_id": "q", "wait_id": "w", "state": "waiting", "source_ref": {"path": "preserved"}} + owner_wait.set_owner_wait(tmp_path, "t", wait) + answer = owner_quiz.record_answered(tmp_path, "t", quiz_id="q", option_index=0, + request_id="answer-1", comment="Keep this exact answer") + write_task_result(tmp_path, "t", terminal, result="terminal work") + + assert owner_quiz.reconcile_terminal(tmp_path, "t") == [] + row = load_task_result(tmp_path, "t") + assert row["owner_quiz"]["q"] == answer["block"] + assert row["owner_wait"]["state"] == "expired_terminal" + assert row["owner_wait"]["source_ref"] == wait["source_ref"] + assert row["result"] == "terminal work" + assert owner_quiz.reconcile_terminal(tmp_path, "t") == [] + assert load_task_result(tmp_path, "t") == row + + +@pytest.mark.parametrize("lifecycle", ["tombstoned", "deleting"]) +def test_project_followup_consumes_only_a_permanent_refusal(tmp_path, monkeypatch, lifecycle): + from ouroboros.projects_registry import create_project, begin_project_deletion, complete_project_deletion + from ouroboros.tools.followup import _handle_schedule_followup + from ouroboros.tools.registry import ToolContext + from supervisor import queue, queue_schedules + + root = tmp_path / "data" + monkeypatch.setattr(queue, "DRIVE_ROOT", root) + for name, value in {"PENDING": [], "RUNNING": {}, "ADMISSION_RESERVATIONS": {}, + "BUDGET_ROOT_FENCES": {}, "ACCEPTANCE_FENCES": {}, + "QUEUE_SEQ_COUNTER_REF": {"value": 0}}.items(): + monkeypatch.setattr(queue, name, value) + monkeypatch.setattr(queue, "persist_queue_snapshot", lambda **_: None) + monkeypatch.setattr(queue_schedules, "resync_skill_schedules", lambda *_: {}) + project = create_project(root, "project-a", name="Project A") + ctx = ToolContext(repo_dir=tmp_path, drive_root=root, task_id="source-task", + project_id=project["id"], current_chat_id=project["chat_id"]) + assert _handle_schedule_followup(ctx, run_at="2000-01-01T00:00:00Z", + objective="Continue in the same project").startswith("FOLLOWUP_SCHEDULED") + begin_project_deletion(root, project["id"]) + if lifecycle == "tombstoned": + complete_project_deletion(root, project["id"]) + queue.check_scheduled_tasks() + first = queue.list_scheduled_tasks(root)["tasks"][0] + failed = load_task_result(root, first["last_task_id"]) + assert failed["status"] == "failed" and failed["reason_code"] == "project_routing_fence" + queue.check_scheduled_tasks() + second = queue.list_scheduled_tasks(root)["tasks"][0] + assert queue.PENDING == [] + if lifecycle == "tombstoned": + assert second["last_task_id"] == first["last_task_id"] + assert second["enabled"] is False and second["completed_at"] + assert len(list((root / "task_results").glob("*.json"))) == 1 + assert second["failure_count"] == 1 and "project_routing_fence" in second["last_error"] + else: + assert second["last_task_id"] != first["last_task_id"] + assert second["enabled"] is True and not second.get("completed_at") + + +class ReachedExec(BaseException): + """Stop before OS effects while retaining the real exec environment.""" + + +def direct_exec_environment(monkeypatch, root): + import server + import ouroboros.server_control as control + + captured = {} + monkeypatch.setattr(server, "DATA_DIR", root) + monkeypatch.setattr(server, "_owner_restart_requested", threading.Event()) + monkeypatch.setenv("OUROBOROS_SERVER_HOST", "127.0.0.1") + def execvpe(_file, _argv, env): + captured.update(env) + raise ReachedExec() + monkeypatch.setattr(control.os, "execvpe", execvpe) + with pytest.raises(ReachedExec): + server._restart_current_process("127.0.0.1", 8765) + return captured + + +def test_assisted_update_carries_the_already_parked_wait_to_direct_successor(restart_case, monkeypatch): + import server + from ouroboros.gateway import control + from supervisor import git_ops, update_merge, workers, queue + + case = restart_case + monkeypatch.setattr(git_ops, "DRIVE_ROOT", case.root) + monkeypatch.setattr(workers, "close_repo_writer_admission", lambda *_: None) + monkeypatch.setattr(workers, "drain_repo_writers", lambda: []) + real_kill = workers.kill_workers + def kill(**kwargs): + return real_kill(**{**kwargs, "reconcile_delegate_custody": False, "archive_service_logs": False}) + monkeypatch.setattr(workers, "kill_workers", kill) + assert control._quiesce_repo_writers("assisted") == [] + assert not workers.RUNNING and not case.process.is_alive() + successor = workers.PENDING[0] + transaction_id = successor["_owner_wait_resume"]["restart_transaction_id"] + monkeypatch.setattr(update_merge, "active_update_tx", lambda: {"phase": "pending_boot_smoke"}) + monkeypatch.setattr(update_merge, "read_update_tx_strict", lambda: ("valid", {"phase": "pending_boot_smoke"})) + monkeypatch.setattr(server, "_safe_restart_serialized", lambda *_a, **_kw: (True, "ok")) + monkeypatch.setattr(server, "_request_restart_exit", lambda: None) + monkeypatch.setattr(server, "_planned_delegate_restart_transaction_id", "") + state = {} + ctx = SimpleNamespace(DRIVE_ROOT=case.root, REPO_DIR=case.root, RUNNING=workers.RUNNING, + load_state=lambda: state.copy(), save_state=lambda s: state.update(s), safe_restart=lambda **_: (True, "ok"), + kill_workers=kill, persist_queue_snapshot=queue.persist_queue_snapshot) + server._perform_supervisor_restart(ctx) + assert server._planned_delegate_restart_transaction_id == "" + env = direct_exec_environment(monkeypatch, case.root) + assert env[delegate_recovery.PLANNED_RESTART_TRANSACTION_ENV] == transaction_id + monkeypatch.setenv(delegate_recovery.PLANNED_RESTART_TRANSACTION_ENV, env[delegate_recovery.PLANNED_RESTART_TRANSACTION_ENV]) + assert restore_stale_snapshot(case) == 1 + assert workers.PENDING[0]["_owner_wait_resume"]["wait_id"] == case.wait["wait_id"] + assert delegate_recovery._read_restart_transaction(case.root, transaction_id)["status"] == "normal_exit_acknowledged" + + +@pytest.mark.parametrize("intent", ["manual_restart", "manual_rollback"]) +def test_aborted_update_does_not_lend_old_handoff_authority_to_manual_restart(tmp_path, monkeypatch, intent): + from supervisor import update_merge + + transaction = {"transaction_id": "aborted", "status": "prepared", "supervisor_pid": os.getpid()} + delegate_recovery._write_restart_transaction(tmp_path, transaction) + delegate_recovery._active_restart_transaction_path(tmp_path).write_text(json.dumps({"transaction_id": "aborted"})) + monkeypatch.delenv(delegate_recovery.PLANNED_RESTART_TRANSACTION_ENV, raising=False) + monkeypatch.setattr(update_merge, "active_update_tx", lambda: {}) + monkeypatch.setattr(update_merge, "read_update_tx_strict", lambda: ("absent", {})) + if intent == "manual_restart": + (tmp_path / "state" / "owner_restart_no_resume.flag").write_text("owner restart") + env = direct_exec_environment(monkeypatch, tmp_path) + assert delegate_recovery.PLANNED_RESTART_TRANSACTION_ENV not in env + assert delegate_recovery._read_restart_transaction(tmp_path, "aborted")["status"] == "prepared" + + +def test_owner_no_resume_flag_dominates_a_restartable_update(tmp_path, monkeypatch): + from supervisor import update_merge + + delegate_recovery._write_restart_transaction(tmp_path, + {"transaction_id": "prior", "status": "prepared", "supervisor_pid": os.getpid()}) + delegate_recovery._active_restart_transaction_path(tmp_path).write_text(json.dumps({"transaction_id": "prior"})) + (tmp_path / "state" / "owner_restart_no_resume.flag").write_text("owner restart") + monkeypatch.setenv(delegate_recovery.PLANNED_RESTART_TRANSACTION_ENV, "prior") + monkeypatch.setattr(update_merge, "read_update_tx_strict", lambda: ("valid", {"phase": "pending_boot_smoke"})) + assert delegate_recovery.PLANNED_RESTART_TRANSACTION_ENV not in direct_exec_environment(monkeypatch, tmp_path) + + +def test_pre_loop_checkpoint_failure_reports_unknown_counts_without_reading_corrupt_source(restart_case, monkeypatch): + from ouroboros import agent as agent_module, agent_task_pipeline, loop + from ouroboros.agent import Env, OuroborosAgent + + case = restart_case + case.task["_owner_wait_resume"] = {**case.wait, "restart_transaction_id": "prepared-test"} + monkeypatch.setattr(OuroborosAgent, "_log_worker_boot_once", lambda _: None) + monkeypatch.setattr(agent_module, "build_llm_messages", lambda **_: ([], {})) + monkeypatch.setattr(agent_task_pipeline, "_run_post_task_processing_async", lambda *_a, **_kw: None) + def unreadable(_ctx): + raise ValueError("owner-wait source checksum mismatch") + monkeypatch.setattr(loop, "load_owner_wait", unreadable) + monkeypatch.setattr("ouroboros.llm.LLMClient.chat", lambda *_a, **_kw: pytest.fail("no generation during failed recovery")) + agent = OuroborosAgent(Env(repo_dir=case.root, drive_root=case.root)) + events = agent._handle_task_scoped(case.task) + stored = load_task_result(case.root, case.task_id) + assert stored["status"] == "failed" and stored["reason_code"] == "task_exception" + assert "unknown" in stored["trace_summary"] and "0 calls" not in stored["trace_summary"] + assert stored["loop_outcome"]["usage"]["total_rounds"] is None + assert stored["loop_outcome"]["usage"]["prompt_tokens"] is None + metrics = next(event for event in events if event["type"] == "task_metrics") + assert metrics["tool_calls"] is None and metrics["tool_errors"] is None + from supervisor.events_worker_reports import _handle_task_metrics + published = [] + metrics_ctx = SimpleNamespace(DRIVE_ROOT=case.root, RUNNING={}, + append_jsonl=lambda _path, row: published.append(row), + bridge=SimpleNamespace(push_log=lambda row: None)) + _handle_task_metrics(metrics, metrics_ctx) + assert published[0]["tool_calls"] is None and published[0]["tool_errors"] is None + assert stored["owner_wait"]["source_ref"] == case.wait["source_ref"] + + +def test_unknown_exception_summary_does_not_invent_zero_rounds_or_buy_a_model_call(tmp_path, monkeypatch): + from ouroboros.post_task_synthesis import _run_task_summary + + rows = [] + monkeypatch.setattr("ouroboros.project_dialogue.append_authored_task_summary", + lambda _root, _result_root, row, **_: rows.append(row)) + monkeypatch.setattr("ouroboros.llm_observability.chat_observed", lambda *_a, **_kw: pytest.fail("no new paid summary")) + _run_task_summary(SimpleNamespace(drive_root=tmp_path), None, {"id": "unknown", "text": "Recover work"}, + {"loop_evidence_unavailable": True}, + {"loop_evidence_unavailable": True, "tool_calls": []}, tmp_path / "logs") + assert rows[0]["tool_calls"] is None and rows[0]["rounds"] is None + assert "round count unknown" in rows[0]["text"] + + +def test_failed_exception_attachment_keeps_the_original_error_and_unknown_projection(tmp_path): + from ouroboros.agent import _task_exception_terminal + from ouroboros.loop_budget import _LoopExitContext + + class FixedError(RuntimeError): + def __setattr__(self, _key, _value): + raise TypeError("attributes unavailable") + + error = FixedError("original failure") + _LoopExitContext(None, None, "t", None, tmp_path, {"rounds": 7}, {"tool_calls": [1]}).attach_exception_evidence(error) + text, usage, trace = _task_exception_terminal(SimpleNamespace(drive_root=tmp_path), {"id": "t"}, error, tmp_path) + assert "FixedError: original failure" in text + assert usage["loop_evidence_unavailable"] is trace["loop_evidence_unavailable"] is True + + +def test_unknown_exception_evidence_stays_unknown_in_actual_reflection(tmp_path, monkeypatch): + from ouroboros import post_task_synthesis, reflection, llm_observability + captured = [] + monkeypatch.setattr(reflection, 'append_reflection_routed', lambda env, task, row: captured.append(row)) + monkeypatch.setattr(llm_observability, 'chat_observed', lambda *a, **kw: ({'content': 'Reflection over the disclosed unknown trace.'}, {})) + task = {'id': 'cold-source-failure', 'type': 'task', 'text': 'Continue saved workspace work', 'workspace_root': str(tmp_path / 'workspace'), 'drive_root': str(tmp_path)} + usage = {'loop_evidence_unavailable': True, 'execution_status': 'infra_failed', 'reason_code': 'task_exception'} + trace = {'loop_evidence_unavailable': True, 'tool_calls': [], 'reasoning_notes': []} + post_task_synthesis._run_reflection(SimpleNamespace(drive_root=tmp_path), None, task, usage, trace, {}) + assert len(captured) == 1, 'The real reflection path must execute and publish one row.' + assert captured[0]['rounds'] is None and captured[0]['error_count'] is None, captured[0] diff --git a/tests/test_cancel_protocol_inventory_s6.py b/tests/test_cancel_protocol_inventory_s6.py index df7909efc..13de6b111 100644 --- a/tests/test_cancel_protocol_inventory_s6.py +++ b/tests/test_cancel_protocol_inventory_s6.py @@ -52,7 +52,7 @@ _TERMINAL_TOKENS = ( # constant from the sticky set; "dynamic" is a variable or expression that can # carry one, which counts because the reducer, not the caller, decides. TERMINAL_WRITERS = { - ('ouroboros/agent.py::OuroborosAgent._handle_task_scoped', 'STATUS_FAILED'): 'terminal', + ('ouroboros/agent.py::_task_exception_terminal', 'STATUS_FAILED'): 'terminal', ('ouroboros/agent_task_pipeline.py::_store_task_result', 'status'): 'dynamic', ('ouroboros/delegate_terminal.py::record_terminal_reconciliation', 'str(existing.get("status") or STATUS_RUNNING)'): 'dynamic', # F6 upstream sync: the cursor/backfill refresh rewrites a stale stored diff --git a/tests/test_restart_reconnect.py b/tests/test_restart_reconnect.py index bba77c176..499c518af 100644 --- a/tests/test_restart_reconnect.py +++ b/tests/test_restart_reconnect.py @@ -349,9 +349,12 @@ def test_owner_restart_copy_is_explicit_about_stopped_task(): assert "stable_skip_flag.unlink(missing_ok=True)" in source # Checkout gate first (a refusal leaves the server intact), then the durable # no-resume intent, then the owned-work stop, then the owner's stop notice. - notice = source.index("Stopping active task. New settings apply to the next message.") - assert (source.index("_safe_restart_serialized(") < source.index("owner_restart_no_resume.flag") - < source.index("_stop_owned_work(ctx)") < notice) + owner_restart = source.split('elif lowered.startswith("/restart"):', 1)[1].split( + 'elif lowered == "/review"', 1 + )[0] + notice = owner_restart.index("Stopping active task. New settings apply to the next message.") + assert (owner_restart.index("_safe_restart_serialized(") < owner_restart.index("owner_restart_no_resume.flag") + < owner_restart.index("_stop_owned_work(ctx)") < notice) stop = _read("ouroboros/server_restart.py").split("def _stop_owned_work", 1)[1] assert (stop.index("request_cancel(") < stop.index("ctx.kill_workers(") < stop.index("reconcile_orphaned_runs(") < stop.index("stop_outcome()")) diff --git a/tests/test_server_shutdown.py b/tests/test_server_shutdown.py index c36ef5585..1c82aa591 100644 --- a/tests/test_server_shutdown.py +++ b/tests/test_server_shutdown.py @@ -193,6 +193,8 @@ def test_managed_update_restart_arms_prepared_transaction_for_direct_reexec( import ouroboros.delegate_recovery as delegate_recovery import ouroboros.gateway.control as control import supervisor.git_ops as git_ops + import supervisor.update_merge as update_merge + import server prepared = { "transaction_id": "tx-owner-wait", @@ -208,6 +210,9 @@ def test_managed_update_restart_arms_prepared_transaction_for_direct_reexec( json.dumps(prepared), encoding="utf-8", ) monkeypatch.setattr(git_ops, "DRIVE_ROOT", tmp_path) + monkeypatch.setattr(server, "DATA_DIR", tmp_path) + monkeypatch.setattr(server, "_restart_current_process_impl", lambda *_a, **_kw: None) + monkeypatch.setattr(update_merge, "read_update_tx_strict", lambda: ("valid", {"phase": "pending_boot_smoke"})) monkeypatch.delenv(delegate_recovery.PLANNED_RESTART_TRANSACTION_ENV, raising=False) request = type("Request", (), {"app": type("App", (), {"state": type("State", (), { @@ -216,6 +221,8 @@ def test_managed_update_restart_arms_prepared_transaction_for_direct_reexec( response = control._restart_response(request, strategy="auto_merge", plan={}) assert response.status_code == 200 + assert delegate_recovery.PLANNED_RESTART_TRANSACTION_ENV not in os.environ + server._restart_current_process("127.0.0.1", 8765) assert os.environ[delegate_recovery.PLANNED_RESTART_TRANSACTION_ENV] == "tx-owner-wait"