diff --git a/docs/architecture/06-agent-core.md b/docs/architecture/06-agent-core.md index c6cf246fb..59c81f3ac 100644 --- a/docs/architecture/06-agent-core.md +++ b/docs/architecture/06-agent-core.md @@ -625,23 +625,23 @@ Affordability probes reserve nothing; full-cap ledger admission decides. Unknown #### Exact budget pause and Resume -`budget_pause.py` pauses pooled/direct work under the SAME ID on global exhaustion, graceful ceiling, either last-fit stop, soft landing or refused dispatch, without a paid final. Missing task id/root/continuation owner retains terminal behavior with `exact_pause_unavailable`. Increases never wake tasks. +`budget_pause.py` pauses pooled/direct work under the SAME ID on global exhaustion, graceful ceiling, either last-fit stop, soft landing or refused dispatch, without a paid final. A missing task id/root/continuation owner keeps terminal behavior (`exact_pause_unavailable`). Increases never wake tasks. | Owner | Ordered contract and reason | |---|---| -| `request_pause`, `usage_accounting.reserve_attempt` | Close `DispatchFenced`; review POST/replay and Light extraction also refuse (`budget_pausing_no_send`, `budget_pausing_no_extraction`). Persist `pausing` before waiting (create a direct RUNNING stub if absent), so crash custody cannot replay ambiguous work; `has_budget_pause_checkpoint` fails closed on unreadable evidence. | -| `local_producer_observation` | Hold the worker nonterminal until sent reviews settle through `review_custody` AND timed-out tool futures finish their late settlement callbacks (`register_tool_future`/`hold_tool_settlement`); `done()` alone is insufficient, unobservable is unknown. Failed writes/producers publish `budget_pause_hold` on checkpoint change and retry, never claim a pause or buy a final; a publication failing after quiescence retries the SAME prepared snapshot (custody stops requested, source stored once), discarded if a producer goes live again. Only Stop/Panic/deadline/cancel ends the hold as `abandoned`, through the loop's model-wait control rails (a no-call terminal, never a task exception); controls read the canonical budget root on split drives. | -| `observe_external_runs` | Re-read custody; unreadable is `custody_read=failed`, not empty. Pre-terminal subscription coverage is unproven, so unsettled runs request `cancel_and_verify`, retaining requested/confirmed/unknown outcomes. Unknown stop permits no second writer. | +| `request_pause`, `usage_accounting.reserve_attempt` | Close `DispatchFenced`; review POST/replay and Light extraction also refuse (`budget_pausing_no_send`, `budget_pausing_no_extraction`). Persist `pausing` before waiting (direct RUNNING stub if absent) so crash custody cannot replay ambiguous work; `has_budget_pause_checkpoint` fails closed on unreadable evidence. | +| `local_producer_observation` | Hold the worker nonterminal until sent reviews settle through `review_custody` AND timed-out tool futures finish their late settlement callbacks (`register_tool_future`/`hold_tool_settlement`); `done()` alone is insufficient, unobservable is unknown. Failed writes/producers publish `budget_pause_hold` on change and retry, never claim a pause or buy a final; a publication failing after quiescence retries the SAME prepared snapshot, discarded if a producer goes live again. Only Stop/Panic/deadline/cancel ends the hold as `abandoned`, via the loop's model-wait control rails (a no-call terminal, never a task exception); controls read the canonical budget root on split drives. | +| `observe_external_runs` | Re-read custody; unreadable is `custody_read=failed`, not empty. Pre-terminal subscription coverage is unproven, so unsettled runs request `cancel_and_verify`, retaining typed outcomes. Unknown stop permits no second writer. | | `owner_wait.continuation_state` | Retains cognition, opaque acceptance/preparation, usage/rounds/clocks. `resume_point` retains the budget tail and unanswered call IDs; Resume inserts execution-UNKNOWN host rows, never replays tools. Source/row precedes `BudgetPauseRequested`; no task_done/result/Main final. | -| `events_budget.install_exact_budget_pause` | Validate event/source; park RUNNING or `parkable_direct_task` in PENDING as `_budget_pause`, persist snapshot, then mark `paused`. Root scope fences the tree. Park, projection and grant share the queue lock; pause-id/state CAS turns late publication into `budget_pause_park_superseded`. Direct records retain `_is_direct_chat`, omit inline image bytes, and release the local fence once the actor has unwound (inline park, or a failed park's ordinary path; the durable row owns the hold). Resumed direct work uses a pooled worker; `log_addressing.address_task_event` preserves the lane. | -| `worker_health._complete_exact_budget_pause_after_death`, `queue_snapshot._park_pausing_running_rows` | Complete a saved pause after death/restart. Revoke an UNCONSUMED grant against its current identity and re-park; a CONSUMED grant takes terminal crash custody, never ordinary retry. Failed revocation retains source/grant in a nonterminal hold, explicitly memory-only if persistence fails. | -| `budget_pause_restore_refusal`, `events_budget.budget_hold_fact` | Pauses restore at any snapshot age without wake. Missing/unreadable/mismatched/terminal records or sources, root acceptance or malformed snapshot fences retain `_budget_pause` under `_budget_pause_hold`, never drop/cancel. Regrant revalidates. | +| `events_budget.install_exact_budget_pause` | Validate event/source; park RUNNING or `parkable_direct_task` in PENDING as `_budget_pause`, persist snapshot, then mark `paused`. Root scope fences the tree. Park, projection and grant share the queue lock; pause-id/state CAS turns late publication into `budget_pause_park_superseded`. Direct records retain `_is_direct_chat`, omit inline image bytes and release the local fence once the actor unwinds (the durable row owns the hold). Resumed direct work uses a pooled worker; `log_addressing.address_task_event` keeps the lane. | +| `worker_health._complete_exact_budget_pause_after_death`, `queue_snapshot._park_pausing_running_rows` | Complete a saved pause after death/restart. Revoke an UNCONSUMED grant against its current identity and re-park; a CONSUMED grant takes terminal crash custody, never ordinary retry. Failed revocation retains source/grant in a nonterminal hold (memory-only if persistence fails). | +| `budget_pause_restore_refusal`, `events_budget.budget_hold_fact` | Pauses restore at any snapshot age or stamp validity without wake. Missing/unreadable/mismatched/terminal records or sources, root acceptance or malformed snapshot fences retain `_budget_pause` under `_budget_pause_hold`, never drop/cancel. Regrant revalidates. | -**Owner grant.** `POST /api/tasks/{id}/resume` → `queue_transitions.resume_budget_paused_task` → `budget_resume.grant_exact_budget_resume`: require authoritative positive global/root headroom, readable source, clear cancel/restart/Panic/deadline/lifetime rails and an unpaused ancestor root; re-read failed custody. Root headroom is ONE fresh successful ledger observation (`usage_accounting.refresh_root_accounting(strict=True)`): money never answers from the display cache — a stale snapshot after a failed read, an unreadable ledger or a degraded tree refuses typed; display readers keep their bounded-stale fallback. Pause ID/generation and `resume_generation` bind a single-use grant. `revoke_exact_budget_resume` re-parks after lost money/restart, preserves newer identities and retries unwritten revocation before regrant. `model_wait.execution_elapsed_seconds` subtracts quota-union and cumulative `paused_duration_sec`/`budget_paused_sec` from original `started_at` across continuations, including owner waits; None stays unlimited. Money/round/review wallets never reset. +**Owner grant.** `POST /api/tasks/{id}/resume` → `queue_transitions.resume_budget_paused_task` → `budget_resume.grant_exact_budget_resume`: require authoritative positive global/root headroom, readable source, clear cancel/restart/Panic/deadline/lifetime rails and an unpaused ancestor root; re-read failed custody. Root headroom is ONE fresh strict ledger read (`usage_accounting.refresh_root_accounting(strict=True)`): a cached snapshot after a failed read, an unreadable ledger or a degraded tree refuses typed; display readers keep their bounded-stale fallback. Pause ID/generation and `resume_generation` bind a single-use grant. `revoke_exact_budget_resume` re-parks after lost money/restart, keeps newer identities and retries unwritten revocation before regrant. `model_wait.execution_elapsed_seconds` subtracts quota-union and cumulative `paused_duration_sec`/`budget_paused_sec` from original `started_at` across continuations, including owner waits; None stays unlimited. Money/round/review wallets never reset. -**Descendant selection.** `events_budget.hold_root_resume_descendants` makes children eligible, not runnable: exact rows remain paused; zero-dispatch/fence-only siblings get `root_fence_lifted_pending_selection`, rebound to the current root grant on each Resume. `resume_child_task` (`budget_resume_child`, recorded `selected_by`) explicitly selects within lineage: root may select stored descendants, intermediate parent only direct children; owner selection uses the same path. With a live latch, `budget_fence_selected` binds only that member to its fence generation, leaving siblings fenced; legacy root Resume selects only root. New root pauses invalidate pending child grants; assignment rechecks root grant and fence. Other roots/cancelled/completed/stopped members never revive. +**Descendant selection.** `events_budget.hold_root_resume_descendants` makes children eligible, not runnable: exact rows remain paused; zero-dispatch/fence-only siblings get `root_fence_lifted_pending_selection`, rebound to the current root grant on each Resume. `resume_child_task` (`budget_resume_child`, recorded `selected_by`) selects within lineage: root selects stored descendants, an intermediate parent only direct children; owner selection shares the path. With a live latch, `budget_fence_selected` binds one member to its fence generation, leaving siblings fenced; legacy root Resume selects only root; `reserve_attempt` admits the member selected against that exact fence. New root pauses invalidate pending child grants; assignment rechecks root grant and fence for exact grants and zero-dispatch selections (stale → unselected hold). Other roots/cancelled/completed/stopped members never revive. -**Consumption.** `resume_paused_loop` rejects spent/revoked/foreign grants and discloses drift/custody before effects. Graceful Resume refreshes ledger-authorized headroom through the same strict tree read and admits ONE fitting reservation (`_second_reservation_fits`, `budget_resume_last_fit_admitted`), avoiding the early-margin re-pause. Hard exhaustion needs owner increase; full-cap admission still binds. Browser/services/hidden state and crash recovery are not restored. +**Consumption.** `resume_paused_loop` rejects spent/revoked/foreign grants and discloses drift/custody before effects; a consumption publication that fails HOLDS (`resume_grant_consumption_unwritable`; identities retained), a grant revoked underneath re-parks; never FAILED. Graceful Resume refreshes ledger-authorized headroom via the same strict tree read and admits ONE fitting reservation (`_second_reservation_fits`, `budget_resume_last_fit_admitted`), avoiding the early-margin re-pause. Hard exhaustion needs owner increase; full-cap admission binds. Browser/services/hidden state and crash recovery are not restored. #### Monetary authority and projections diff --git a/ouroboros/budget_pause.py b/ouroboros/budget_pause.py index 56cc0b74a..ee80b85b6 100644 --- a/ouroboros/budget_pause.py +++ b/ouroboros/budget_pause.py @@ -614,6 +614,9 @@ _HOLD_POLL_SEC = 0.5 HOLD_PRODUCERS_UNSETTLED = "local_producers_unsettled" HOLD_PAUSE_RECORD_UNWRITABLE = "pause_record_unwritable" HOLD_CHECKPOINT_UNWRITABLE = "continuation_source_unwritable" +# A Resume whose grant consumption could not be published: the grant and the +# checkpoint stay exactly as the durable row carries them; the task HOLDS. +HOLD_GRANT_CONSUMPTION_UNWRITABLE = "resume_grant_consumption_unwritable" # A ``pausing`` row the task's own controls ended before it became a pause. STATE_ABANDONED = "abandoned" @@ -1117,14 +1120,55 @@ def resume_paused_loop(tools: Any, state: Dict[str, Any], messages: list, trace: root = pathlib.Path(ctx.budget_drive_root or ctx.drive_root) grant = dict(row.get("grant") or {}) grant["consumed_at"] = time.time() - # Compare-and-set on the exact pause, state and grant this worker was - # handed: a grant revoked or superseded between the dispatch and this - # write refuses here instead of consuming a grant that is no longer live. - set_budget_pause(root, ctx.task_id, {**row, "state": STATE_RESUMED, "grant": grant, - "resumed_at": grant["consumed_at"]}, - expected_pause_id=str(row.get("pause_id") or ""), - expected_state=STATE_RESUME_GRANTED, - expected_grant_id=str(grant.get("grant_id") or "")) + consumed = {**row, "state": STATE_RESUMED, "grant": grant, "resumed_at": grant["consumed_at"]} + published = "" + while True: + # Compare-and-set on the exact pause, state and grant this worker was + # handed: a grant revoked or superseded between the dispatch and this + # write refuses here instead of consuming a grant that is no longer live. + try: + set_budget_pause(root, ctx.task_id, consumed, + expected_pause_id=str(row.get("pause_id") or ""), + expected_state=STATE_RESUME_GRANTED, + expected_grant_id=str(grant.get("grant_id") or "")) + break + except Exception as exc: + # A consumption that cannot be published is never a paid terminal + # (#1196): the worker HOLDS, typed and nonterminal, retrying the SAME + # publication so the grant and checkpoint identities stay exactly as + # the durable row carries them. A row the durable authority already + # returned to a live pause (a revocation landed first) re-parks under + # it: there is no grant to consume and nothing to run on. + try: + current = budget_pause_row(root, ctx.task_id) + except Exception: + current = {} + if (current.get("pause_id") == row.get("pause_id") + and current.get("state") in {STATE_PAUSING, STATE_PAUSED} + and int(current.get("task_attempt") or 0) == int(ctx.task_attempt or 1) + and current.get("source_ref")): + usage.pop("budget_pause_hold", None) + raise BudgetPauseRequested(current) from exc + hold = _hold_row(HOLD_GRANT_CONSUMPTION_UNWRITABLE, exc=exc) + usage["budget_pause_hold"] = dict(hold) + if str(hold.get("error") or "") != published: + published = str(hold.get("error") or "") + log.warning("Budget resume for %s is HELD (nonterminal): grant consumption unpublished %s", + ctx.task_id, published) + _publish_hold(ctx, {"pause_id": row.get("pause_id"), "rail": row.get("rail"), + "state": STATE_RESUME_GRANTED, **hold}) + control = _hold_control_reason(ctx) + if control: + usage["exact_pause_unavailable"] = "hold_ended_by_control" + usage["budget_pause_hold"] = {**hold, "ended_by": control} + _publish_hold(ctx, {"pause_id": row.get("pause_id"), "rail": row.get("rail"), + "state": "hold_ended", "hold_reason": hold["hold_reason"], + "ended_by": control}) + from ouroboros.model_wait import ModelWaitInterrupted + + raise ModelWaitInterrupted(control) from exc + time.sleep(_HOLD_POLL_SEC) + usage.pop("budget_pause_hold", None) ctx._budget_paused_sec = float(grant.get("paused_duration_sec") or row.get("paused_duration_sec") or 0.0) ctx.budget_pause_resume = None setattr(ctx, "_budget_pause_generation", int(row.get("pause_generation") or 0)) diff --git a/ouroboros/loop.py b/ouroboros/loop.py index c28957a3b..6e7869d74 100644 --- a/ouroboros/loop.py +++ b/ouroboros/loop.py @@ -478,9 +478,7 @@ def run_llm_loop( continuation = saved or saved_pause tool_schemas = continuation["tool_schemas"] if continuation else initial_tool_schemas(tools, context_mode=active_context_mode) - tool_schemas, _enabled_extra_tools = _setup_dynamic_tools( - tools, tool_schemas, messages, context_mode=active_context_mode - ) + tool_schemas, _enabled_extra_tools = _setup_dynamic_tools(tools, tool_schemas, messages, context_mode=active_context_mode) ctx.event_queue, ctx.task_id, ctx.messages = event_queue, task_id, messages stateful_executor = StatefulToolExecutor() exit_ctx = _LoopExitContext( @@ -729,8 +727,10 @@ def run_llm_loop( pending_tool_budget, pending_tool_calls = True, tool_calls except BudgetExceeded as exc: _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) + try: + return _handle_budget_exceeded(exc, exit_ctx, limit_ctx=limit_ctx, episode=transport_wait) + except ModelWaitInterrupted as interrupted: # a refused-dispatch HOLD ended by control: the sibling clause below never sees it + return _loop_exit_after_exception(interrupted, limit_ctx, exit_ctx, llm_trace, transport_wait) except Exception as exc: # A budget-pause HOLD ended by control rejoins the model-wait rails; else re-raise with evidence. return _loop_exit_after_exception(exc, limit_ctx, exit_ctx, llm_trace, transport_wait) diff --git a/ouroboros/owner_wait.py b/ouroboros/owner_wait.py index 175bc7c2a..f036b5965 100644 --- a/ouroboros/owner_wait.py +++ b/ouroboros/owner_wait.py @@ -393,6 +393,8 @@ def restore_continuation_state(tools: Any, state: dict, messages: list, trace: d NOT restored: they died with the previous process and stay invalidated.""" from ouroboros.loop_delivery import DeliveryCandidate + from ouroboros.model_wait import budget_paused_seconds + ctx = tools._ctx messages[:] = state["messages"] trace.update(state["trace"]) @@ -400,6 +402,11 @@ def restore_continuation_state(tools: Any, state: dict, messages: list, trace: d seen.update(state["seen"]) ctx._loop_mailbox_seen_ids = seen ctx._owner_directives = state["owner_directives"] + # The cumulative budget-paused carrier rides EVERY same-ID continuation + # (#1196, F5): a cold owner-wait restore of a task that had been budget + # paused keeps it, so a later pause row and the delegate clock start from + # the same cumulative value; a budget grant overrides it with its own. + ctx._budget_paused_sec = budget_paused_seconds(state.get("model_wait") or {}) for key, value in {**state["route"], **state["delivery"], **state["acceptance"]}.items(): setattr(ctx, key, value) candidate = state.get("delivery_candidate") diff --git a/ouroboros/subagent_bootstrap.py b/ouroboros/subagent_bootstrap.py index 0aac96bcd..81622092b 100644 --- a/ouroboros/subagent_bootstrap.py +++ b/ouroboros/subagent_bootstrap.py @@ -163,6 +163,41 @@ def bootstrap_before_context(ctx: Any, task: Mapping[str, Any], dispatch: Any) - return _with_coordination_context(ctx, recovery) actor_bootstrap = getattr(ctx, "_configured_actor_bootstrap", {}) actor_bootstrap = actor_bootstrap if isinstance(actor_bootstrap, dict) else {} + if isinstance(task.get("_budget_pause_resume"), dict): + # Same-ID budget continuation (#1196): the paused attempt's own episode + # already decided this leaf's physical start and the restored transcript + # carries its receipts. The host hydrates the durable custody facts and + # mints NO second invocation over a settled or disposed leaf — a + # replacement start is the model's explicit decision after admission. + from ouroboros import delegate_custody as custody + from ouroboros.delegate_evidence import task_execution_evidence + + try: + evidence = task_execution_evidence( + custody.custody_root(ctx), str(getattr(ctx, "task_id", "") or task.get("id") or "")) + except Exception: + evidence = {"evidence_read_failed": True} + evidence = evidence if isinstance(evidence, dict) else {"evidence_read_failed": True} + started = int(evidence.get("delegated_runs_started") or 0) + if started: + _mark_physical_activity(ctx) + elif evidence.get("evidence_read_failed"): + # Unreadable custody may hide a prior run: UNKNOWN fences a new start. + actor_bootstrap.update({ + "zero_run_evidence_status": "unknown", + "zero_run_evidence_gaps": ["custody_evidence_unreadable"], + "exact_start_pending": False, + }) + try: + payload = json.loads(actor_ready) + except (TypeError, ValueError): + payload = {} + payload["status"] = "configured_session_budget_continuation" + payload["continuation"] = { + "delegated_runs_started": started, "physical_start": "not_repeated", + "custody_read": "failed" if evidence.get("evidence_read_failed") else "ok", + } + return _with_coordination_context(ctx, json.dumps(payload, ensure_ascii=False, indent=2)) fenced = ( bool(actor_bootstrap.get("zero_run_receipt_recorded")) or str(actor_bootstrap.get("zero_run_evidence_status") or "") == "unknown" diff --git a/ouroboros/usage_accounting.py b/ouroboros/usage_accounting.py index 73d94b220..8067db6f9 100644 --- a/ouroboros/usage_accounting.py +++ b/ouroboros/usage_accounting.py @@ -825,8 +825,23 @@ def reserve_attempt(request: AttemptRequest) -> AttemptReservation: for row in rows: if not isinstance(row, dict): raise UsageAccountingError(f"invalid root budget fence row: {snapshot_path}") - if (str(row.get("root_task_id") or "") == root_task_id - and str(row.get("status") or "") in {"active", "paused"}): + if (str(row.get("root_task_id") or "") != root_task_id + or str(row.get("status") or "") not in {"active", "paused"}): + continue + # ONE member explicitly selected against THIS fence generation is + # admitted (owner Q9, #1196): the queue recorded that selection on + # the row itself, and the latch still refuses every unselected member. + fence_id = str(row.get("fence_id") or "") + selected = False + for bucket in ("running", "pending"): + for entry in (snapshot.get(bucket) or []) if isinstance(snapshot, dict) else []: + member = entry.get("task") if isinstance(entry, dict) else None + if not isinstance(member, dict) or str(member.get("id") or "") != scope.task_id: + continue + hold = member.get("_budget_pause_hold") + selected = bool(fence_id and isinstance(hold, dict) and hold.get("selected") + and str(hold.get("fence_id") or "") == fence_id) + if not selected: raise BudgetExceeded( f"root model dispatch paused pending explicit resume for {scope.root_task_id}", limit_scope="root", diff --git a/supervisor/events_budget.py b/supervisor/events_budget.py index 85bdf37ee..f3e4e348f 100644 --- a/supervisor/events_budget.py +++ b/supervisor/events_budget.py @@ -596,20 +596,29 @@ def budget_fence_selected(task: Any, fence: Any) -> bool: def budget_resume_dispatch_allowed(q: Any, task: Dict[str, Any]) -> bool: - """An exact child selection belongs to the CURRENT root grant and fence only.""" + """A child selection belongs to the CURRENT root grant and fence only. + + Both carriers are revalidated at dispatch: an exact grant handoff and a + zero-dispatch hold selection alike name the root grant they were selected + under, so a root that is pausing or paused again (with or without a new + fence) admits neither on its old selection (#1196, owner Q9). + """ handoff = task.get("_budget_pause_resume") - if not isinstance(handoff, dict): + hold = task.get(BUDGET_HOLD_KEY) if isinstance(task.get(BUDGET_HOLD_KEY), dict) else None + exact = isinstance(handoff, dict) + carrier = handoff if exact else (hold if hold is not None and hold.get("selected") else None) + if carrier is None: return True root_id = str(task.get("root_task_id") or task.get("id") or "") fence_id = str((q.BUDGET_ROOT_FENCES.get(root_id) or {}).get("fence_id") or "") - if fence_id != str(handoff.get("root_fence_id") or ""): + if fence_id != str(carrier.get("root_fence_id" if exact else "fence_id") or ""): return False if root_id == str(task.get("id") or ""): return True root_grant = live_root_resume_grant(q, root_id, pathlib.Path(task.get("budget_drive_root") or q.DRIVE_ROOT)) - return bool(handoff.get("selected_by") == "owner" and not handoff.get("root_grant_id") and not root_grant - or root_grant and handoff.get("root_grant_id") == root_grant["grant_id"] - and handoff.get("root_resume_generation") == root_grant["generation"]) + return bool(carrier.get("selected_by") == "owner" and not carrier.get("root_grant_id") and not root_grant + or root_grant and carrier.get("root_grant_id") == root_grant["grant_id"] + and int(carrier.get("root_resume_generation") or 0) == int(root_grant["generation"])) def hold_budget_row(task: Dict[str, Any], *, reason: str, detail: str = "", @@ -738,8 +747,8 @@ def select_held_budget_row(q: Any, task: Dict[str, Any], hold: Dict[str, Any], root_task_id = str(hold.get("root_task_id") or task.get("root_task_id") or task_id) if hold.get("reason") not in {HOLD_ROOT_FENCE_LIFTED, HOLD_ROOT_FENCE_MEMBER_SELECTION}: return {"ok": False, "error": str(hold.get("reason") or "budget_hold_unresolved")} - if selected_by: - root_grant = live_root_resume_grant(q, root_task_id, result_root) + root_grant = live_root_resume_grant(q, root_task_id, result_root) if task_id != root_task_id else {} + if selected_by and task_id != root_task_id: if not root_grant: return {"ok": False, "error": "root_resume_grant_missing", "root_task_id": root_task_id, "action": "resume_root_first"} @@ -755,6 +764,11 @@ def select_held_budget_row(q: Any, task: Dict[str, Any], hold: Dict[str, Any], "selected_by": str(selected_by or "owner")} if task_id == root_task_id: selection.update(root_grant_id=uuid.uuid4().hex, root_resume_generation=1) + elif root_grant: + # The selection names the root grant it was made under, so dispatch can + # recheck it exactly like an exact grant handoff (budget_resume_dispatch_allowed). + selection.update(root_grant_id=root_grant["grant_id"], + root_resume_generation=int(root_grant["generation"])) prior_pause = task.pop("_budget_pause", None) task[BUDGET_HOLD_KEY] = selection if not q.persist_queue_snapshot(reason="budget_hold_selected"): diff --git a/supervisor/queue_snapshot.py b/supervisor/queue_snapshot.py index 788acc72f..181f4fe5f 100644 --- a/supervisor/queue_snapshot.py +++ b/supervisor/queue_snapshot.py @@ -671,9 +671,18 @@ def restore_pending_from_snapshot( return 0 ts = str(snap.get("ts") or "") ts_unix = _queue().parse_iso_to_ts(ts) + # Timestamp validity and freshness gate ORDINARY rows only (#1196): a + # readable snapshot whose stamp is missing or malformed is treated as + # stale, so an identifiable exact pause or an acknowledged owner-wait + # handoff is still retained under its own durable authority instead of + # vanishing (a Resume would then answer task_not_pending over an intact checkpoint). if ts_unix is None: - return 0 - stale = (time.time() - ts_unix) > max_age_sec + _queue().append_jsonl( + _queue().DRIVE_ROOT / "logs" / "supervisor.jsonl", + {"ts": utc_now_iso(), "type": "queue_restore_snapshot_timestamp_invalid", + "snapshot_ts": ts[:64], "action": "treated_as_stale"}, + ) + stale = ts_unix is None or (time.time() - ts_unix) > max_age_sec from ouroboros.task_results import ( _TRULY_TERMINAL_STATUSES, STATUS_CANCEL_REQUESTED, STATUS_CANCELLED, load_task_result, write_task_result, diff --git a/supervisor/worker_assignment.py b/supervisor/worker_assignment.py index 38e5f3424..1075b56f4 100644 --- a/supervisor/worker_assignment.py +++ b/supervisor/worker_assignment.py @@ -13,6 +13,7 @@ import pathlib import time from typing import Any, Dict +from ouroboros.model_wait import budget_paused_seconds from supervisor.events_budget import budget_fence_selected, budget_hold_fact from supervisor.queue import _queue_lock @@ -378,9 +379,27 @@ def assign_tasks() -> None: from supervisor.events_budget import budget_resume_dispatch_allowed if not budget_resume_dispatch_allowed(queue, candidate): - from supervisor.budget_resume import revoke_exact_budget_resume + if isinstance(candidate.get("_budget_pause_resume"), dict): + from supervisor.budget_resume import revoke_exact_budget_resume - revoke_exact_budget_resume(candidate, "root_resume_generation_stale") + revoke_exact_budget_resume(candidate, "root_resume_generation_stale") + else: + # A zero-dispatch selection whose root grant or fence is + # no longer live returns to an UNSELECTED hold (#1196, Q9): + # the row keeps its hold identity, drops the dead grant + # binding, and the next selection records the live one. + from supervisor.events_budget import ( + BUDGET_HOLD_KEY, HOLD_ROOT_FENCE_LIFTED, hold_budget_row, + ) + + stale = candidate.get(BUDGET_HOLD_KEY) or {} + hold_budget_row( + candidate, reason=str(stale.get("reason") or HOLD_ROOT_FENCE_LIFTED), + detail="root_resume_generation_stale", + extra={**{key: stale[key] for key in ("root_task_id", "fence_id") + if key in stale}, + "stale_root_grant_id": str(stale.get("root_grant_id") or "")}, + result_root=pathlib.Path(candidate.get("budget_drive_root") or _pool().DRIVE_ROOT)) queue.persist_queue_snapshot(reason="stale_child_resume_held") continue if (root_task_id in queue.BUDGET_ROOT_FENCES @@ -443,9 +462,11 @@ def assign_tasks() -> None: **({"model_wait_quota_clock": dict(resume["model_wait_quota_clock"])} if resume.get("model_wait_quota_clock") else {}), # Separate paused-interval carrier (#1196): the original - # started_at is untouched; lifetime rails subtract this. - **({"budget_paused_sec": float(resume["paused_duration_sec"] or 0.0)} - if resume.get("paused_duration_sec") else {}), + # started_at is untouched; lifetime rails subtract this. ONE + # reader for either handoff: a budget grant names it + # ``paused_duration_sec``, an owner-wait restart ``budget_paused_sec``. + **({"budget_paused_sec": budget_paused_seconds(resume)} + if budget_paused_seconds(resume) > 0 else {}), "soft_sent": False, "attempt": int(task.get("_attempt") or 1), } task_type = str(task.get("type") or "") diff --git a/tests/test_budget_pause_exact_resume.py b/tests/test_budget_pause_exact_resume.py index c9b6f528f..0d2b3f8b2 100644 --- a/tests/test_budget_pause_exact_resume.py +++ b/tests/test_budget_pause_exact_resume.py @@ -627,3 +627,112 @@ def test_runbook_does_not_promise_a_managed_outage_window_the_runtime_has_no_rai source = _repo_file("devtools", "benchmarks", "continual_learning", "RUNBOOK.md") assert "6h operation window from episode entry" not in source assert "OUROBOROS_TASK_ABS_CEILING_SEC" in source and "idle reaper" in source + + +# --------------------------------------------------------------------------- consumption publication holds + +def _granted_loop(tmp_path, monkeypatch, task_id): + from ouroboros import budget_pause, owner_wait + + queue, state, _workers = _install_queue(tmp_path, monkeypatch) + monkeypatch.setattr(state, "budget_remaining", lambda _st, **_k: 5.0) + task, row = _parked(tmp_path, monkeypatch, task_id=task_id) + assert queue.resume_budget_paused_task(task_id)["ok"] is True + ctx, _limit = _loop_ctx(tmp_path, task_id) + ctx.budget_pause_resume = task["_budget_pause_resume"] + state_blob = budget_pause.load_budget_pause(ctx) + monkeypatch.setattr(owner_wait, "rebind_restored_route", lambda *_a, **_k: (None, "max")) + return budget_pause, ctx, state_blob, row + + +def test_failed_grant_consumption_publication_holds_then_consumes_never_terminalizes(tmp_path, monkeypatch): + """#1196 review finding 2: a raise from the state=resumed/consumed_at write used + to leave ``resume_paused_loop`` on the generic loop exception path and end in + FAILED (``_task_exception_terminal``). A publication that fails HOLDS — typed, + nonterminal, retrying the SAME compare-and-set — and consumes once it lands.""" + budget_pause, ctx, state_blob, _row = _granted_loop(tmp_path, monkeypatch, "consume-hold") + real = budget_pause.set_budget_pause + failures = {"left": 2} + + def flaky(*args, **kwargs): + if failures["left"] > 0: + failures["left"] -= 1 + raise OSError("disk full") + return real(*args, **kwargs) + + monkeypatch.setattr(budget_pause, "set_budget_pause", flaky) + published = [] + monkeypatch.setattr(budget_pause, "_publish_hold", lambda _ctx, hold: published.append(dict(hold))) + usage = {} + budget_pause.resume_paused_loop(SimpleNamespace(_ctx=ctx), state_blob, [], {}, usage, set(), + budget_remaining_usd=5.0) + consumed = budget_pause.budget_pause_row(tmp_path, "consume-hold") + assert consumed["state"] == budget_pause.STATE_RESUMED and consumed["grant"]["consumed_at"] + assert "budget_pause_hold" not in usage and "execution_status" not in usage + assert [hold["hold_reason"] for hold in published] == [budget_pause.HOLD_GRANT_CONSUMPTION_UNWRITABLE] + assert published[0]["state"] == budget_pause.STATE_RESUME_GRANTED and "OSError" in published[0]["error"] + + +def test_grant_consumption_hold_ended_by_control_retains_the_grant_and_checkpoint(tmp_path, monkeypatch): + from ouroboros.model_wait import ModelWaitInterrupted + from tests._budget_pause_exact_helpers import _controls + + budget_pause, ctx, state_blob, row = _granted_loop(tmp_path, monkeypatch, "consume-stop") + + def unwritable(*_args, **_kwargs): + raise OSError("read-only drive") + + monkeypatch.setattr(budget_pause, "set_budget_pause", unwritable) + monkeypatch.setattr(budget_pause, "_hold_control_reason", _controls("", "deadline")) + usage = {} + with pytest.raises(ModelWaitInterrupted) as raised: + budget_pause.resume_paused_loop(SimpleNamespace(_ctx=ctx), state_blob, [], {}, usage, set(), + budget_remaining_usd=5.0) + assert raised.value.control_reason == "deadline" + durable = budget_pause.budget_pause_row(tmp_path, "consume-stop") + grant = state_blob["_pause_row"]["grant"] + # Grant and checkpoint identities are exactly what the row carried: nothing consumed, nothing rewritten. + assert durable["state"] == budget_pause.STATE_RESUME_GRANTED + assert durable["grant"]["grant_id"] == grant["grant_id"] and not durable["grant"].get("consumed_at") + assert durable["source_ref"] == row["source_ref"] and durable["pause_id"] == row["pause_id"] + assert usage["budget_pause_hold"]["hold_reason"] == budget_pause.HOLD_GRANT_CONSUMPTION_UNWRITABLE + assert usage["budget_pause_hold"]["ended_by"] == "deadline" + assert usage["exact_pause_unavailable"] == "hold_ended_by_control" + + +def test_grant_consumption_over_a_revoked_grant_reparks_under_the_live_pause(tmp_path, monkeypatch): + budget_pause, ctx, state_blob, row = _granted_loop(tmp_path, monkeypatch, "consume-revoked") + live = budget_pause.budget_pause_row(tmp_path, "consume-revoked") + budget_pause.set_budget_pause(tmp_path, "consume-revoked", { + **live, "state": budget_pause.STATE_PAUSED, "grant": {**live["grant"], "revoked_at": time.time()}}) + usage = {} + with pytest.raises(budget_pause.BudgetPauseRequested) as raised: + budget_pause.resume_paused_loop(SimpleNamespace(_ctx=ctx), state_blob, [], {}, usage, set(), + budget_remaining_usd=5.0) + assert raised.value.pause["pause_id"] == row["pause_id"] + assert raised.value.pause["state"] == budget_pause.STATE_PAUSED and "budget_pause_hold" not in usage + + +@pytest.mark.parametrize("stamp", ["", "not-a-timestamp"]) +def test_paused_row_survives_a_snapshot_with_an_invalid_timestamp_while_ordinary_rows_do_not( + tmp_path, monkeypatch, stamp): + """#1196 review finding 6: a readable snapshot with valid paused rows but a + missing/invalid ``ts`` returned zero before ``_retain_snapshot_pending`` ran, + so no paused carrier was restored and Resume answered task_not_pending over + an intact durable checkpoint. Timestamp validity gates ORDINARY rows only.""" + queue, _state, workers = _install_queue(tmp_path, monkeypatch) + _parked(tmp_path, monkeypatch, task_id="ts-paused") + workers.PENDING.append({"id": "ts-ordinary", "type": "task", "chat_id": 0, "_attempt": 1, "text": "x"}) + queue.persist_queue_snapshot(reason="test") + workers.PENDING[:] = [] + snap = json.loads(queue.QUEUE_SNAPSHOT_PATH.read_text()) + if stamp: + snap["ts"] = stamp + else: + snap.pop("ts", None) + queue.QUEUE_SNAPSHOT_PATH.write_text(json.dumps(snap)) + assert queue.restore_pending_from_snapshot() == 1 + assert [task["id"] for task in workers.PENDING] == ["ts-paused"] + assert workers.PENDING[0]["_budget_pause"]["exact_continuation"] is True # retained, not woken + rows = [json.loads(line) for line in (tmp_path / "logs" / "supervisor.jsonl").read_text().splitlines()] + assert any(row["type"] == "queue_restore_snapshot_timestamp_invalid" for row in rows) diff --git a/tests/test_budget_pause_holds.py b/tests/test_budget_pause_holds.py index aca61b824..674e598de 100644 --- a/tests/test_budget_pause_holds.py +++ b/tests/test_budget_pause_holds.py @@ -558,3 +558,74 @@ def test_an_unwritten_grant_rollback_holds_the_row_instead_of_granting_forever(t assert second["released_hold"] == HOLD_REVOCATION_UNWRITTEN row = budget_pause.budget_pause_row(tmp_path, "f6-1") assert row["grant"]["grant_id"] == second["grant_id"] and row["resume_generation"] == 2 + + +def test_a_root_selected_against_its_retained_fence_reserves_while_siblings_stay_refused(tmp_path, monkeypatch): + """#1196 review finding 3: a legacy zero-dispatch root Resume records the + selection but keeps the root latch, and ``reserve_attempt`` used to refuse + EVERY reservation under that fence — Resume "succeeded", the worker's first + send hit the fence, and the root re-paused without a model call. Reservation + admission now honours the selection recorded against the exact fence; the + unselected siblings are still refused at the same gate (owner Q9).""" + from ouroboros import usage_accounting as accounting + from supervisor.events_budget import _set_root_budget_pause_locked + + queue, state, workers = _install_queue(tmp_path, monkeypatch) + monkeypatch.setattr(state, "budget_remaining", lambda _st, **_k: 5.0) + root = _fenced_member(workers, "root-r", "root-r") + root.pop("parent_task_id") + _fenced_member(workers, "sib-r", "root-r") + fence = _set_root_budget_pause_locked("root-r", {}) + assert queue.resume_budget_paused_task("root-r")["ok"] is True + sent = [] + _idle_worker(workers, sent) + workers.assign_tasks() + assert [task["id"] for task in sent] == ["root-r"] + assert queue.BUDGET_ROOT_FENCES["root-r"]["fence_id"] == fence["fence_id"] # latch retained + + def _reserve(task_id): + with accounting.usage_scope(accounting.UsageScope( + drive_root=tmp_path, task_id=task_id, root_task_id="root-r", global_limit_usd=100.0)): + return accounting.reserve_attempt(accounting.AttemptRequest( + model="fixture", provider="openai", reservation_usd=1.0)) + + assert _reserve("root-r").attempt_id # the selected row's first send is admitted + with pytest.raises(accounting.BudgetExceeded) as refused: + _reserve("sib-r") + assert refused.value.limit_scope == "root" + # A selection recorded against an OLDER fence generation is no key to a new latch. + queue.BUDGET_ROOT_FENCES["root-r"] = {**fence, "fence_id": "fence-next"} + queue.persist_queue_snapshot(reason="test") + with pytest.raises(accounting.BudgetExceeded): + _reserve("root-r") + + +def test_owner_wait_restart_assignment_and_the_cold_loop_keep_the_budget_paused_carrier(tmp_path, monkeypatch): + """#1196 review finding 5: after a budget Resume an owner-wait checkpoint + stores ``budget_paused_sec`` but assignment read only ``paused_duration_sec``, + so after a planned restart the RUNNING row charged the old pause as execution; + the cold owner-wait restore also left ``ctx._budget_paused_sec`` unset. One + shared reader (``model_wait.budget_paused_seconds``) serves either handoff.""" + from ouroboros import owner_wait + + queue, state, workers = _install_queue(tmp_path, monkeypatch) + monkeypatch.setattr(state, "budget_remaining", lambda _st, **_k: 5.0) + started = time.time() - 1000.0 + workers.PENDING.append({ + "id": "wait-6", "type": "task", "chat_id": 0, "_attempt": 1, + "_owner_wait_resume": {"wait_id": "w6", "restart_transaction_id": "tx-6", "started_at": started, + "budget_paused_sec": 600.0, "model_wait_quota_clock": {}}, + }) + sent = [] + _idle_worker(workers, sent) + workers.assign_tasks() + assert [task["id"] for task in sent] == ["wait-6"] + meta = workers.RUNNING["wait-6"] + assert meta["started_at"] == pytest.approx(started) and meta["budget_paused_sec"] == 600.0 + + ctx, _limit = _loop_ctx(tmp_path, "wait-6") + state_blob = {"messages": [], "trace": {}, "usage": {}, "seen": [], "owner_directives": [], + "route": {}, "delivery": {}, "acceptance": {}, "delivery_candidate": None, + "model_wait": {"budget_paused_sec": 600.0}} + owner_wait.restore_continuation_state(SimpleNamespace(_ctx=ctx), state_blob, [], {}, {}, set()) + assert ctx._budget_paused_sec == 600.0 diff --git a/tests/test_budget_pause_safety.py b/tests/test_budget_pause_safety.py index bbe3c007c..e5e3e7a7c 100644 --- a/tests/test_budget_pause_safety.py +++ b/tests/test_budget_pause_safety.py @@ -420,3 +420,58 @@ def test_an_interruption_the_model_call_rail_already_routed_is_not_routed_twice( _run(loop_mod, registry, tmp_path, "stop-once") assert routed == ["cancelled"] and raised.value.control_rails_seen is True assert isinstance(getattr(raised.value, "_ouroboros_loop_usage", None), dict) + + +def test_a_zero_dispatch_selection_is_rechecked_against_the_live_root_grant_at_dispatch(tmp_path, monkeypatch): + """#1196 review finding 4: a child carrying a SELECTED ``_budget_pause_hold`` + (no exact grant handoff) returned True from ``budget_resume_dispatch_allowed`` + immediately, so after a global-scope re-pause of its root — which raises no + new fence — it dispatched on the old selection. Both carriers are now bound + to the root grant they were selected under; a stale one returns to an + UNSELECTED hold and is re-bound and re-selectable under the next root Resume.""" + queue, state, workers = _install_queue(tmp_path, monkeypatch) + monkeypatch.setattr(state, "budget_remaining", lambda *_a, **_k: 5.0) + root, _ = _parked(tmp_path, monkeypatch, task_id="root", scope="root") + sibling = _fenced_member(workers, "sibling", "root") + first = queue.resume_budget_paused_task("root") + assert first["ok"] and queue.resume_budget_paused_task("sibling", selected_by="root")["ok"] + assert sibling["_budget_pause_hold"]["selected"] is True + assert sibling["_budget_pause_hold"]["root_grant_id"] == first["grant_id"] + workers.PENDING.remove(root) # the root ran on, then paused AGAIN under global scope: no new fence + _parked(tmp_path, monkeypatch, task_id="root", scope="global") + assert "root" not in queue.BUDGET_ROOT_FENCES + sent = [] + _idle_worker(workers, sent) + workers.assign_tasks() + assert sent == [] and sibling["_budget_pause_hold"]["selected"] is False + assert sibling["_budget_pause_hold"]["detail"] == "root_resume_generation_stale" + second = queue.resume_budget_paused_task("root") + assert second["ok"] and second["grant_id"] != first["grant_id"] + workers.assign_tasks() + assert [task["id"] for task in sent] == ["root"] # the root Resume made the sibling eligible only + assert queue.resume_budget_paused_task("sibling", selected_by="root")["ok"] + workers.WORKERS[0].busy_task_id = None + workers.assign_tasks() + assert [task["id"] for task in sent] == ["root", "sibling"] + assert sent[-1]["_budget_pause_hold"]["root_grant_id"] == second["grant_id"] + + +def test_a_hold_ended_by_control_inside_the_refused_dispatch_rail_is_a_truthful_terminal(tmp_path, monkeypatch): + """#1196 review finding 7: ``_handle_budget_exceeded`` runs INSIDE the loop's + ``except BudgetExceeded`` clause; a hold it entered through ``request_pause`` + that a deadline ended raised ``ModelWaitInterrupted`` past the sibling + ``except Exception`` and reached task_exception. It now rejoins the common + control rails: a no-call ``deadline_local`` terminal, no model call.""" + from ouroboros import budget_pause, usage_accounting + from ouroboros.model_wait import ModelWaitInterrupted + + loop_mod, registry = _bare_loop(tmp_path, monkeypatch) + monkeypatch.setattr(loop_mod, "_call_round_model", lambda _call: (_ for _ in ()).throw( + usage_accounting.BudgetExceeded("global model budget exhausted", limit_scope="global"))) + monkeypatch.setattr(usage_accounting, "usage_breakdown", lambda *_a, **_k: {"physical_calls": 3}) + monkeypatch.setattr(budget_pause, "request_pause", + lambda *_a, **_k: (_ for _ in ()).throw(ModelWaitInterrupted("deadline"))) + text, usage, trace = _run(loop_mod, registry, tmp_path, "hold-budget-rail") + assert usage["execution_status"] == "failed" and usage["reason_code"] == "deadline_local" + assert trace["forced_finalization"]["control_reason"] == "deadline" + assert isinstance(text, str) and text diff --git a/tests/test_configured_session_prestart.py b/tests/test_configured_session_prestart.py index b7b2ebf72..8c0abf71b 100644 --- a/tests/test_configured_session_prestart.py +++ b/tests/test_configured_session_prestart.py @@ -405,3 +405,61 @@ def test_precustody_refusals_leave_a_durable_start_blocked_row(tmp_path): assert "configured_work_order_unavailable" in out.text reasons = [row["reason"] for row in _rows(tmp_path)] assert reasons[-1] == "configured_work_order_unavailable" + + +def test_a_budget_continuation_never_pre_starts_the_leaf_again_and_hydrates_custody(monkeypatch, tmp_path): + """#1196 review finding 1: a configured ``agent_session`` nanny that paused + AFTER its leaf settled and its patch was disposed used to reach + ``_pre_start_leaf`` again on Resume (no crash handoff, no unsettled custody), + minting a second invocation with the original canonical assignment before + the grant was consumed or the transcript restored. A same-ID budget + continuation now bypasses the physical start, hydrates the durable custody + facts onto the actor bootstrap and leaves any replacement to the model.""" + import ouroboros.delegate_evidence as evidence_mod + import ouroboros.subagent_runtime as runtime + from ouroboros.subagent_bootstrap import bootstrap_before_context + + starts = [] + monkeypatch.setattr(runtime, "exact_start", lambda _ctx, prompt, _spec: ( + starts.append(prompt) or _fail("delegate_start", "start_probe", "probe"))) + monkeypatch.setattr(evidence_mod, "task_execution_evidence", lambda _root, _tid: { + "delegated_runs_started": 1, "delegated_runs_settled": 1, "delegated_runs_succeeded": 1}) + snapshot = _snapshot(_settings(_session_row()), "session-builder") + dispatch = SimpleNamespace( + executor="harness", blocked=False, + executor_resolution=SimpleNamespace(route=SimpleNamespace(route_id="codex")), + ) + + def _ctx(): + return SimpleNamespace(task_id="child-resumed", drive_root=tmp_path, + budget_drive_root=str(tmp_path), task_metadata={}) + + task = {"id": "child-resumed", "configured_subagent": snapshot, + "task_contract": {"objective": "Build"}, + "_budget_pause_resume": {"pause_id": "p1", "grant_id": "g1", "grant_generation": 1}} + ctx = _ctx() + receipt = json.loads(bootstrap_before_context(ctx, task, dispatch)) + assert starts == [] # no second physical invocation + assert receipt["status"] == "configured_session_budget_continuation" + assert receipt["continuation"] == { + "delegated_runs_started": 1, "physical_start": "not_repeated", "custody_read": "ok"} + bootstrap = ctx._configured_actor_bootstrap + assert bootstrap["physical_started"] is True and bootstrap["exact_start_pending"] is False + assert ctx._nanny_physical_activity_seed is True # nanny economics see the adopted run + assert not hasattr(ctx, "_configured_startup_refusal") # never an unrun $0 terminal + + # Unreadable custody may hide a prior run: still no start, typed UNKNOWN fence. + def _boom(_root, _tid): + raise OSError("custody log unreadable") + monkeypatch.setattr(evidence_mod, "task_execution_evidence", _boom) + ctx_unknown = _ctx() + receipt_unknown = json.loads(bootstrap_before_context(ctx_unknown, task, dispatch)) + assert starts == [] and receipt_unknown["continuation"]["custody_read"] == "failed" + unknown = ctx_unknown._configured_actor_bootstrap + assert unknown["zero_run_evidence_status"] == "unknown" and unknown["exact_start_pending"] is False + assert unknown["physical_started"] is False + + # Control: the same task WITHOUT the continuation handoff pre-starts the exact leaf. + task.pop("_budget_pause_resume") + bootstrap_before_context(_ctx(), task, dispatch) + assert len(starts) == 1 and "Build" in starts[0]