fix(#1196): close Astra exact-diff findings (Fable run-94419e602eaf): consumption hold, fence-selected admission, hold revalidation, paused-seconds carrier, snapshot ts gating, control routing, continuation bootstrap

This commit is contained in:
Anton 2026-09-24 15:19:36 +03:00
parent 5ffed43ae3
commit b52f4f1aa2
13 changed files with 478 additions and 40 deletions

View file

@ -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

View file

@ -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))

View file

@ -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)

View file

@ -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")

View file

@ -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"

View file

@ -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",

View file

@ -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"):

View file

@ -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,

View file

@ -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 "")

View file

@ -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)

View file

@ -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

View file

@ -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

View file

@ -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]