Close nested coordination custody and review races

Co-authored-by: Ouroboros <311266734+ouroboros-agent@users.noreply.github.com>
This commit is contained in:
Ouroboros 2026-08-24 10:28:53 +03:00
parent 508531ba53
commit 81194970b2
27 changed files with 1900 additions and 351 deletions

View file

@ -147,7 +147,7 @@ server.py (Starlette+uvicorn) ← HTTP + WebSocket on configurable host:port (de
├── subagent_runtime.py ← Active configured-actor runtime: immutable task-start snapshots, exact `subagent_id` selection, API-child versus session-nanny dispatch, typed alternatives, exact-start binding, parent cognitive-route inheritance, and the bounded deterministic legacy-input seam.
├── subagent_work_order.py ← One complete external work-order compiler from the scheduled child's objective/context/output/constraints/acceptance/authority; the single 250,000-character total wire budget sends fitting orders complete, and over-budget orders become a full-SHA/source-selector partial lens rather than a prefix. The generic live manifest interaction capability is read here; no harness name is interpreted.
├── subagent_bootstrap.py ← Actor-first episode and recovery seam for configured session nannies: expose the immutable route/work-order facts before a new physical start, let the host LLM choose children/leaf/zero-run, inject a typed route-unavailable fact without silent fallback, and preserve exact startup/recovery receipts. An over-budget order stays pending until the existing source-range interaction is positively available; otherwise the configured bridge returns a typed source-channel refusal. Pending recovery replays the stored compact body and canonical full-order fingerprint.
├── delegate_supervision.py ← Event-only sleeping-nanny loop over the low-level wait: quiet windows renew without a model call; terminal/interaction/fault/addressed-message/control or one explicit reasoned checkpoint becomes a durable pending/acknowledged wake.
├── delegate_supervision.py ← Event-only sleeping-nanny loop over the low-level wait: quiet windows renew without a model call; terminal/interaction/fault/addressed-message/control or one explicit reasoned checkpoint becomes a durable pending/acknowledged wake. Startup and each newly minted meaningful wake carry one fresh coordination context (parent intent, time, tree spend, host-visible descendants and root review capacity); replay returns the stored snapshot unchanged.
├── delegate_start_instructions.py ← Stable harness host instructions plus the bounded, separately hashed actor-first coordination appendix; canonical work-order authority remains a separate immutable field.
├── delegate_recovery.py ← Narrow exact-leaf recovery for a proven non-signal worker crash and an explicit planned-self-restart handoff; validates task/config/worktree/authority/run/cursor bindings before adoption and vetoes every no-resume cause.
├── delegate_pending.py ← Durable pending-invocation replay helpers: preserve the original idempotency key and canonical start body when a start response/custody write is uncertain.
@ -222,7 +222,7 @@ server.py (Starlette+uvicorn) ← HTTP + WebSocket on configurable host:port (de
├── server_runtime.py ← Server startup/onboarding and WebSocket liveness helpers
├── server_web.py ← Static web file helpers (NoCacheStaticFiles, web dir resolver)
├── task_continuation.py ← Durable per-task review continuation state across restart/outage
├── task_results.py ← Durable task result/status files (task_results/<id>.json)
├── task_results.py ← Durable task result/status files (`task_results/<id>.json`), including the canonical root's strict locked `task_acceptance_review_accounting` exact-binding claims, their immutable digest, and the read-only root review-capacity projection used by configured-session wakes; the claim count is the one tree-wide paid acceptance-cycle authority and a claim without a recoverable terminal host run is unknown, never permission to re-dispatch.
├── task_status.py ← Effective task-status SSOT: child-drive result merge, lineage lookup, bounded waits
├── git_shell_policy.py ← Structural git argv classifiers for shell safety guards
├── protected_artifacts.py ← Task-contract protected artifact policy helpers for execute-only black-box references
@ -323,6 +323,7 @@ server.py (Starlette+uvicorn) ← HTTP + WebSocket on configurable host:port (de
│ ├── join_ledger.py ← Soft-join decision authority: validates direct lineage and exact current child-result hashes for tagged `tree_note(kind="decision")` dispositions (`integrated`, `irrelevant`, `deferred`) — single-child or batch `children` array form, each batch entry validated individually — appends the sole authoritative task-tree row, rejects stale hashes as `CHILD_RESULT_STALE`, and keeps `peek_task`, `discard_child_result`, constraint override, cancellation, and shared child-decision helpers. The hash covers status, full result, trace summary, artifact status, and stable artifact identities, not cost/timestamps/queue diagnostics/parent decisions; task-result fields are derived read projections only.
│ ├── delegate.py ← Nanny verb facade: exact `delegate_start(subagent_id=..., prompt=...)` / idempotent `retry_of`, event-only `delegate_wait` with optional one-shot reasoned checkpoint, custody-gated `delegate_cancel`, and typed `delegate_answer`. A configured session's host actor may start the snapshotted leaf after its first ordinary episode; direct starts, actor-first starts and recovery share `subagent_runtime.exact_start`. Host-derived authority, private snapshot integration, exact `skill_payload` resource selection, terminal output staging, verified settlement and no-widening rules remain unchanged. Supervision is `delegate_supervision.py`, recovery is `delegate_recovery.py`, custody is `delegate_custody.py`, and transport is `gateways/claudexor.py`.
│ ├── delegate_integration.py ← The C1 integration seam of the nanny verbs (extracted from tools/delegate.py for the module-size gate; delegate.py re-exports it so sibling code and tests keep one name): `_mutation_authority` derives the unified host authority record `{target_root, source, capture_mode}` for a mutating run (acting write_root vs B5 external-workspace root, typed refusals on any disagreement), `_provision_snapshot` registers + describes the private execution snapshot durably BEFORE any start intent, `_resolve_retry_invocation`/`_validated_invocation`/`_retry_binding_refusal` rebuild a retried start from its ONE durable invocation record and re-prove the C1 binding (pre-C1 mutating rows refused; moved workspace refused; GC-collected snapshot refused), and `_capture_terminal_patch` idempotently captures a terminal mutating run's diff from its snapshot for explicit `integrate_delegated_patch` disposition. It also owns the exact SKILL-PAYLOAD delegation cluster (restored D10 target class, owner option A 2026-08-14): `_payload_mutation_authority` grants `authority_source="skill_payload"` only on a fresh `ResolvedResourceBinding` for a top-level principal (busy-check refuses a second delegation while the same payload has open custody), the host-minted semantic `resource_ref` {root, source, skill_name, target, baseline payload hash} rides custody durably and is re-resolved by `_rebind_payload_reference` at retry and apply, `_write_payload_patch_artifacts` captures over the skill-loader inventory (`git diff --binary` transport so UTF-8-with-NUL survives; junk the loader excludes never enters; an added/modified non-UTF-8 file is a typed capture failure; reserved lifecycle/control paths are reported as `blocked_reserved_paths`, never silently filtered), and `integrate_payload_patch` applies the candidate into the live NON-Git payload — fresh binding must equal the recorded target, whole-payload content-hash CAS (already-applied content disposes idempotently), reserved paths refuse the WHOLE apply with the candidate preserved, index-free `git apply` with the live payload as cwd, no `.git`/staging created there, and any payload-mutating apply outcome — success or post-apply hash mismatch — queues `request_extension_reconcile` while enablement/grants stay untouched and the existing review goes stale by content hash
│ ├── delegate_start_claims.py ← One short pre-transport transaction for a fresh actor start: it serializes the final zero-run/custody recheck and `START_REQUESTED` append, nesting the existing exact-payload claim only when that resource is selected; transport and waiting stay outside every claim lock.
│ └── subagent_integration.py ← integrate_subagent_patch: parent's manifest-first integration of an acting subagent's workspace.patch. For self_worktree children it applies into ctx.active_repo_dir() (sha256-verified, 3-way --index, protected-path gated, top-only lineage check, genesis refused), stages but never commits. For external_workspace children it verifies the child wrote in the same active external workspace and records an audited verdict without re-applying the patch; (v6.58.0) a NON-workspace parent integrating a COOP child (write_root = a host-minted tree under the subagent-projects root) gets a read-only verification + a SUCCESSFUL `coop_already_in_tree` no-op verdict instead of a parent-missing error — the work is already in the shared tree, which `coop_checkpoint.checkpoint_commit_coop_roots` checkpoint-commits at root finalization. Also compare_subagent_patches: read-only best-of-N helper that shows several children's candidate patches side by side for LLM-first synthesis
├── process_containment.py ← `ProcessContainer` for the hermetic gate: env-token membership (`OURO_PROC_CONTAINER_*`, /proc environ on Linux, `ps -E` on macOS, kill-on-close Job Object on Windows), read from LIVE kernel state at reap time; an alive or undeterminable member is an honest hard-block answer, never a kill guarantee. Policy layer over platform_layer's OS primitives
└── platform_layer.py ← Cross-platform process/path/locking helpers, including the public descendant-enumeration seam reused by launcher cleanup and the Windows Job Object ABI with explicit argtypes/restype
@ -709,7 +710,7 @@ A chat instance has an explicit resource lifecycle. `ws.on()` registers listener
`app.js` normally keeps at most one live Project chat instance. Closing or switching stashes only its scroll intent and destroys the instance. The narrow exception is unsendable client state: staged `File` objects or an upload already in flight. Such an instance is hidden and marked pending instead of destroyed, is reused if the Project is reopened, and returns to the ordinary destroy policy after that work settles. Typed but unsent text survives separately in per-thread session storage. This prevents hidden Project rooms from accumulating listeners, repainting, or acknowledging unseen revisions without discarding data the server cannot reconstruct.
Chat and progress logs rotate as one timeline and archived segments remain durable. Interactive history uses bounded, archive-aware readers that expand only far enough to satisfy the requested thread's filtered quota; Project history cannot be satisfied by unrelated Main rows. Terminal annotation occurs after the emitted window is chosen, and display reads avoid materializing or rebasing artifacts. The task-event stream performs an archive-aware replay and then follows appended bytes, handling rotation and newly discovered children; its terminal event performs the one materializing task read needed to deliver artifact-bearing final truth. Children are discovered by a scandir name-diff over the main root's task-results directory. For a subagent, queue insertion and the first durable scheduled-result row form one queue-lock transition: if the result write fails, the still-pending row is removed before assignment and the admission becomes a typed rejection. A follow tick therefore decodes only result files it has not successfully read yet instead of re-projecting the whole store; the disclosed residual is that the per-tick directory scan itself stays proportional to the size of the task-results directory. The reason is bounded UI latency without treating log rotation as conversation loss or mutating artifact state merely to render status.
Chat and progress logs rotate as one timeline and archived segments remain durable. Interactive history uses bounded, archive-aware readers that expand only far enough to satisfy the requested thread's filtered quota; Project history cannot be satisfied by unrelated Main rows. Terminal annotation occurs after the emitted window is chosen, and display reads avoid materializing or rebasing artifacts. The task-event stream performs an archive-aware replay and then follows appended bytes, handling rotation and newly discovered children; its terminal event performs the one materializing task read needed to deliver artifact-bearing final truth. Children are discovered by a scandir name-diff over the main root's task-results directory. For a subagent, queue insertion and the first durable scheduled-result row form one queue-lock transition: if the result write fails, the still-pending row is removed before assignment and the admission becomes a typed rejection. A replay whose exact task id already has live or durable custody exits before write-surface provisioning and is rechecked under that lock, so it cannot append a second physical task, replace the accepted transition id, or publish duplicate progress. A follow tick therefore decodes only result files it has not successfully read yet instead of re-projecting the whole store; the disclosed residual is that the per-tick directory scan itself stays proportional to the size of the task-results directory. The reason is bounded UI latency without treating log rotation as conversation loss or mutating artifact state merely to render status.
The agent-facing `chat_history` reader uses the same live-plus-rotated timeline and may narrow it by exact provider, account, conversation, thread, actor, and inclusive date bounds before applying the established count/offset/text-search window. Presence provenance is therefore searchable as structured transport fact rather than only as flattened user prose.
@ -1180,7 +1181,7 @@ Advisory availability is evaluated from the current configured slot and route, n
`usage_ledger.py` owns the durable append-only physical-attempt ledger: cross-process locking, sequence and transition validation, append+fsync, replay, and loud tail quarantine. `usage_accounting.py` is the one-way policy layer above it: route pricing, reservations, settlement, scopes, budget fences, imports, projections, and admission. The substrate never imports policy. This is a structural boundary around the monetary authority: a pricing or budget-policy change cannot redefine valid ledger storage, and a locking or repair change cannot silently change what an attempt costs. Compatibility events, state mirrors, task fields, and UI projections may carry attempt ids and derived totals but never become a second charge source.
### Delegated subagents (Claudexor transport + the nanny)
Children coordinate through `tree_note` and `tree_read`; only the parent may use `override_delegation_constraint`. A `review_requested` note carries an exact evidence reference/hash, wakes the waited/direct parent, preserves distinct typed concerns, and starts no reviewer; the parent/root chooses an ordinary critic or root acceptance path and carries the hash into its existing accounting/deduplication seam. Both read-only and acting children can use the existing descendant-scoped `forward_to_worker`, `peek_task`, `cancel_task`, and `discard_child_result` controls for their own children; the durable ancestry check grants no unrelated-task reach. Workspace children retain scoped `knowledge_read` and `knowledge_list`, and recursive delegation never widens filesystem, budget, depth, deadline, commit or owner authority.
Children coordinate through `tree_note` and `tree_read`; only the parent may use `override_delegation_constraint`. A `review_requested` note carries an exact evidence reference/hash, wakes the waited/direct parent, preserves distinct typed concerns, and starts no reviewer or paid cycle. The parent/root may inspect it, hire an ordinary critic child, or let host-verified bytes enter the final root acceptance packet. Only the latter uses the acceptance wallet: immediately before reviewer transport the canonical root result atomically claims the complete candidate/evidence/fence binding under the effective cycle cap. A pre-existing claim without its terminal host run is typed unknown and cannot be re-dispatched. Both read-only and acting children can use the existing descendant-scoped `forward_to_worker`, `peek_task`, `cancel_task`, and `discard_child_result` controls for their own children; the durable ancestry check grants no unrelated-task reach. Workspace children retain scoped `knowledge_read` and `knowledge_list`, and recursive delegation never widens filesystem, budget, depth, deadline, commit or owner authority.
`OUROBOROS_SUBAGENTS` is the active task-actor SSOT. Its strict `{enabled, items}` value contains at most ten `ConfiguredSubagent` rows. Each row has a stable `subagent_id`, display name, owner-authored English `recommended_use`, one normalized route (`api_model` or `agent_session`), optional effort and an optional session credential pin. The description is selection context only: host code never parses, ranks or maps its words to task text. API models and session harnesses therefore occupy one LLM-selectable list without pretending they have the same topology.
@ -1424,6 +1425,11 @@ with a compact coverage=partial source-request lens. The external actor asks
for exact source character ranges, and its nanny answers through the existing
waiting_on_user/delegate_answer transport from the canonical-work-order projection
of `get_task_result` (the same renderer the host validates).
The manifest observation is a point-in-time preflight, not a lease: the route may
lose its interaction capability before the later start POST. The probe is not
delivery evidence; only durable verified source-range coverage authorizes completion,
so a raced run remains `cannot_verify` and patch application stays refused until
coverage is complete.
An unavailable or unverified channel returns a typed source-channel refusal before
POST; `cannot_verify` remains the distinct verdict for incomplete interaction
evidence after a run exists. Crash recovery replays the durable compact request body and the full-order
@ -1455,6 +1461,19 @@ Quiet renewal emits supervision telemetry and makes zero LLM calls. The existing
external-wait lease still protects legitimate host-side silence from the idle rail;
deadline, absolute ceiling, budget and cancellation remain outer bounds.
The actor-first startup receipt and every newly minted meaningful wake also carry
one host-rendered `coordination_context`: the complete parent-authored advisory
`delegation_budget.intent_note`, explicit-deadline time remaining, known/partial/unknown
root-tree settled and accounted spend, active host-visible descendants, and the
canonical root's remaining paid acceptance capacity. Vendor-internal descendants are
explicitly opaque. These are facts for LLM judgment, not thresholds or a scheduler;
replay returns the stored snapshot byte-for-byte, while the next acknowledged wake
recomputes it from the existing authorities. If the combined wake exceeds the tool
budget, the complete context stays in the existing exact wake source and the bounded
envelope carries a typed source-only projection. Active descendants are known only
from a fresh queue snapshot plus the targeted parent chains of those live rows; stale
queue state is unknown and unrelated historical corruption does not poison the subtree.
The nanny may deliberately request one future inspection by supplying both
`checkpoint_after_sec` and a free-text `checkpoint_reason`. It wakes once at that
time, or an earlier real event consumes the checkpoint. Every later inspection needs

View file

@ -1177,6 +1177,16 @@ Before every commit, verify the following:
invocation is proven absent or terminal and any captured physical result is
explicitly disposed; replay the original pending invocation/idempotency key after
worker loss.
A fresh physical start and `delegation_zero_run` are mutually exclusive actor
decisions. Rebuild all run/start/patch blockers from one custody-log snapshot and
hold the existing short per-task file-lock seam only across the final recheck plus
START_REQUESTED/zero-run append (`delegate_start_claims` owns the start side);
never hold it across transport or waiting.
Treat supervisor delivery of one `schedule_subagent` event as at-least-once:
an exact task id with live or durable custody is an idempotent no-op before
write-surface provisioning, and the same identity check runs again under the
queue lock immediately before enqueue. Never use semantic duplicate judgement
as the physical identity fence.
The complete external work-order wire budget is one total 250,000-character
limit, not a model-context claim and not a per-field prefix rule. A brief that
fits is sent byte-complete. A brief above that limit is never silently prefixed:
@ -1199,6 +1209,11 @@ Before every commit, verify the following:
source selector; timeout or another resolution remains incomplete. Until the union covers the whole brief, terminal delivery carries
`work_order_verification.status=cannot_verify`, and `integrate_delegated_patch`
may reject the captured result but may not apply it.
The manifest observation is a point-in-time preflight, not a lease: capability
may change before the later POST. Never call the probe delivery evidence or add
a second probe/lease to pretend the race vanished. Durable verified range coverage
remains the authority; a raced run stays `cannot_verify` and its patch stays
unapplied until coverage is complete.
- `delegate_wait` is an event-only model sleep. Renew bounded transport windows in
`delegate_supervision` with zero LLM calls; journal progress may stream to the
owner but is not a wake. Wake only for terminal/interaction/fault, an addressed
@ -1212,6 +1227,16 @@ Before every commit, verify the following:
and replay it rather than advancing the coordination cursor. On wake the nanny retains its full ordinary
tool surface and inherited parent cognitive route; no-co-building is a
prompt/review/receipt role contract, not a host allowlist.
Actor-first startup and every newly minted meaningful wake carry one fresh
`coordination_context`: full parent-authored advisory `intent_note`, explicit
deadline time remaining, known/partial/unknown tree spend, active host-visible
descendants and root acceptance capacity. Vendor-internal descendants stay opaque.
Persist this context inside the pending wake so replay is identical; recompute only
after acknowledgement on a later wake. When the combined wake spills, preserve the
complete context in the exact source and keep only a typed bounded projection in the
envelope. These facts inform the LLM and never become an automatic fan-out, hurry,
review or stop state machine. Treat active descendants as known only from a fresh
queue snapshot and targeted live-row ancestry; stale/missing lineage is unknown.
- Recovery is cause-specific. A proven non-signal worker crash and an explicit
planned-self-restart transaction may adopt the same exact run before orphan
cleanup/LLM/start. Owner restart, signals, panic, timeout/deadline, explicit
@ -1296,9 +1321,14 @@ Before every commit, verify the following:
separate typed concerns even when they reference the same bytes, but never starts
or waits for a reviewer. The full hash remains visible through `tree_read`/`peek_task`.
The parent/root decides whether to inspect, spawn an ordinary critic, or use the
root-owned acceptance path; if it launches a check, it carries the hash into the
existing accounting/deduplication path. The referenced bytes remain caller-authored
evidence until that parent/critic actually verifies them.
root-owned acceptance path. The beacon itself spends no cycle, and its hash remains
caller-authored until host bytes are actually read and verified. Immediately before
a real root acceptance transport, strictly claim the complete candidate/evidence/fence
binding in canonical `task_acceptance_review_accounting` under the root task-result
lock. Cap check and exact-binding dedupe are one mutation; missing/malformed authority,
lock failure, or a prior claim without a recoverable terminal run is typed unknown and
starts no reviewer. Ordinary critic children remain ordinary budgeted tasks, not a role-
parsed hidden review flow.
- Subagent changes must keep writes, commits, review mutation, runtime control,
tool expansion, skills lifecycle, and shell blocked — except bounded task-tree
coordination via `tree_note`/`tree_read`, parent-only

View file

@ -23,6 +23,7 @@ import json
import logging
import pathlib
import uuid
from contextlib import contextmanager
from dataclasses import dataclass, field
from typing import Any, Callable, Dict, Iterator, List, Optional, Tuple
@ -264,6 +265,36 @@ def custody_log_unreadable(drive_root: Any) -> bool:
return True
return False
@contextmanager
def actor_decision_lock(drive_root: Any, task_id: str) -> Iterator[None]:
"""Serialize one actor's mutually exclusive zero-run/start decisions.
The custody log and the verification-receipt store are separate append-only
authorities. This short critical section joins only their decision edge:
re-read current authority, then append either a zero-run receipt or a fresh
START_REQUESTED row. Transport, waiting and settlement stay outside it.
"""
from ouroboros.platform_layer import (
acquire_exclusive_file_lock,
release_exclusive_file_lock,
)
identity = str(task_id or "").strip()
if not identity:
raise ValueError("actor decision lock requires a task id")
digest = hashlib.sha256(identity.encode("utf-8")).hexdigest()[:24]
lock_path = pathlib.Path(drive_root) / "state" / "delegate_actor_claims" / f"{digest}.lock"
lock_path.parent.mkdir(parents=True, exist_ok=True)
fd = acquire_exclusive_file_lock(lock_path, timeout_sec=20.0, stale_sec=120.0)
if fd is None:
raise TimeoutError("actor decision lock is unavailable")
try:
yield
finally:
release_exclusive_file_lock(lock_path, fd)
def _iter_rows(path: pathlib.Path, tail_bytes: Optional[int] = None) -> Iterator[Dict[str, Any]]:
try:
with path.open("rb") as handle:
@ -1462,41 +1493,9 @@ def _retire_recovered_registration(gateway: Any, record: Dict[str, Any]) -> bool
return False
def _capture_stranded_patch(drive_root: Any, run: RunCustody) -> Dict[str, Any]:
"""Capture a reconciled mutating run's diff into the ordinary patch artifact.
The reconcile path is the ONLY terminal observer a dead-owner run gets, so
without this the child's work stayed in the snapshot with no captured patch
and no apply/reject material — stranded, invisible, and one binding loss away
from GC. Called ONLY where a terminal receipt PROVES the run is over (C1-R2):
a run closed absent/unreadable has unknowable state, and freezing a patch
there would put a "captured" receipt over work the child might still be
writing — those runs are captured lazily at disposition instead. Reuses the
one existing capture primitive (idempotent, durable ``PATCH_CAPTURED`` row);
capture ONLY — the apply/reject decision belongs to a live owner and is NEVER
taken by a sweep. Fail-soft: a capture error is disclosed in the reconcile
row, and the snapshot persists either way because the run has no recorded
disposition.
"""
if not (run.execution_root and run.settled and not run.patch_disposed):
return {}
try:
from ouroboros.tools.delegate_integration import capture_terminal_patch_for_drive
block = capture_terminal_patch_for_drive(drive_root, run) or {}
except Exception:
log.warning("Reconcile patch capture failed for %s", run.run_id, exc_info=True)
return {"patch_capture": "failed", "patch_disposition": "pending"}
return {"patch_capture": str(block.get("status") or ""),
"patch_artifact": block.get("patch_artifact"),
# The typed disposition-pending disclosure: this rides the durable
# RECONCILED row, and the health surface (``undisposed_patches``)
# keeps the fact visible until an explicit apply/reject lands.
"patch_disposition": "pending"}
def _reconcile_one(drive_root: Any, gateway: Any, custody: RunCustody) -> Dict[str, Any]:
from ouroboros.gateways.claudexor import ClaudexorUnavailable
from ouroboros.tools.delegate_integration import capture_stranded_patch
try:
detail = gateway.get_run(custody.run_id)
@ -1530,7 +1529,7 @@ def _reconcile_one(drive_root: Any, gateway: Any, custody: RunCustody) -> Dict[s
"settled": settled["settled"], **output_disposition(custody)}
# The C1 half: a TERMINAL DETAIL proves the run is over, so the sweep — its
# last terminal observer — captures the diff eagerly here.
result.update(_capture_stranded_patch(drive_root, custody))
result.update(capture_stranded_patch(drive_root, custody))
else:
cancelled = cancel_and_verify(drive_root, gateway, custody, "owner_task_gone")
result = {"run_id": custody.run_id, "task_id": custody.task_id, "action": "cancelled",
@ -1541,7 +1540,7 @@ def _reconcile_one(drive_root: Any, gateway: Any, custody: RunCustody) -> Dict[s
# the run (same unknowable-state doctrine as above) — both leave the
# capture to disposition.
if cancelled["state"] in TERMINAL_STATES:
result.update(_capture_stranded_patch(drive_root, custody))
result.update(capture_stranded_patch(drive_root, custody))
emit(drive_root, RECONCILED, result)
return result
@ -1557,6 +1556,7 @@ __all__ = [
"SOURCE_RANGE_DELIVERY_CONFIRMED",
"TERMINAL_STATES",
"UNKNOWN",
"actor_decision_lock",
"cancel_and_verify",
"close_absent_run",
"custody_log_unreadable",

View file

@ -183,18 +183,30 @@ def _selected_session(task: Mapping[str, Any]) -> dict[str, Any]:
return snapshot if str(route.get("kind") or "") == "agent_session" else {}
def unsettled_start_ids(drive_root: Any, task_id: str) -> dict[str, list[str]]:
"""Durable run/start blockers which must settle before replacement."""
def unsettled_start_ids(
drive_root: Any, task_id: str, *, rows: Optional[list[dict[str, Any]]] = None,
) -> dict[str, list[str]]:
"""Durable run/start blockers from one consistent custody-log snapshot."""
mine = str(task_id or "")
snapshot = list(rows) if rows is not None else list(
custody._iter_rows(custody.event_log_path(drive_root))
)
runs = custody.replay(drive_root, rows=snapshot)
return {
"open_run_ids": [row.run_id for row in custody.open_runs(drive_root) if row.task_id == mine],
"open_run_ids": [
row.run_id for row in runs.values()
if row.task_id == mine and not row.settled
],
"pending_invocation_ids": [
str(row.get("invocation_id") or "") for row in custody.pending_invocations(drive_root)
str(row.get("invocation_id") or "")
for row in custody.pending_invocations(drive_root, rows=snapshot)
if str(row.get("task_id") or "") == mine
],
"undisposed_patch_run_ids": [
row.run_id for row in custody.undisposed_patches(drive_root) if row.task_id == mine
row.run_id for row in runs.values()
if row.task_id == mine and row.snapshot_id and row.settled
and not row.patch_disposed
],
}

View file

@ -0,0 +1,121 @@
"""Atomic pre-transport claims for delegated actor and payload starts."""
from __future__ import annotations
import pathlib
import threading
from typing import Any, Callable, Dict, Tuple
from ouroboros import delegate_custody as custody
_PAYLOAD_CLAIM_LOCK = threading.Lock()
def claimed_start_request(
drive: pathlib.Path, *, claim_target: str,
payload_busy: Callable[[pathlib.Path, pathlib.Path], str],
actor_ctx: Any = None, enforce_actor_idle: bool = False,
**request_row: Any,
) -> Tuple[bool, Dict[str, Any]]:
"""Atomically claim a fresh actor start and optional payload target."""
from ouroboros.platform_layer import (
acquire_exclusive_file_lock,
release_exclusive_file_lock,
)
def _write_request() -> Tuple[bool, Dict[str, Any]]:
if not claim_target:
return custody.record_start_requested(drive, **request_row), {}
lock_path = pathlib.Path(drive) / "state" / ".payload_delegation_claim.lock"
lock_path.parent.mkdir(parents=True, exist_ok=True)
with _PAYLOAD_CLAIM_LOCK:
fd = acquire_exclusive_file_lock(
lock_path, timeout_sec=20.0, stale_sec=120.0,
)
if fd is None:
return False, {
"reason": "payload_delegation_busy",
"holder": "payload claim lock unavailable",
"detail": "The payload start claim is currently held by another caller.",
}
try:
holder = payload_busy(drive, pathlib.Path(claim_target))
if holder:
return False, {
"reason": "payload_delegation_busy", "holder": holder,
"detail": (
"Another delegated run claimed this exact payload first. "
"Finish it before starting another assignment against the skill."
),
}
return custody.record_start_requested(drive, **request_row), {}
finally:
release_exclusive_file_lock(lock_path, fd)
if not enforce_actor_idle:
return _write_request()
task_id = str(request_row.get("task_id") or "").strip()
try:
with custody.actor_decision_lock(drive, task_id):
if custody.custody_log_unreadable(drive):
return False, {
"reason": "replacement_custody_unknown",
"custody_log_unreadable": True,
"detail": "The host cannot read the actor's custody authority.",
}
from ouroboros.delegate_recovery import unsettled_start_ids
blockers = unsettled_start_ids(drive, task_id)
if any(blockers.values()):
if claim_target:
holder = payload_busy(drive, pathlib.Path(claim_target))
if holder:
return False, {
"reason": "payload_delegation_busy",
"holder": holder,
"detail": (
"Another delegated run claimed this exact payload first. "
"Finish it before starting another assignment against the skill."
),
}
return False, {
"reason": "replacement_requires_settlement", **blockers,
"detail": (
"The actor gained an unsettled start/run or undisposed patch "
"before this fresh request could be claimed."
),
}
if actor_ctx is not None:
from ouroboros.subagent_bootstrap import _durable_zero_run_receipt
gaps: set[str] = set()
zero_run = _durable_zero_run_receipt(actor_ctx, gap_reasons=gaps)
if zero_run:
return False, {
"reason": "zero_run_already_recorded",
"detail": (
"The actor recorded a terminal delegation_zero_run decision "
"before this physical start could be claimed."
),
"zero_run_decision": str(
zero_run.get("zero_run_decision") or "unknown"
),
}
if gaps:
return False, {
"reason": "zero_run_evidence_unavailable",
"detail": (
"The host cannot prove whether the actor already recorded "
"a terminal delegation_zero_run decision."
),
"zero_run_evidence_status": "unknown",
"zero_run_evidence_gaps": sorted(gaps),
}
return _write_request()
except (TimeoutError, ValueError) as exc:
return False, {
"reason": "replacement_custody_unknown",
"detail": f"actor decision claim unavailable: {type(exc).__name__}",
}

View file

@ -67,6 +67,178 @@ def _payload(raw: str) -> dict[str, Any]:
return {"status": "fault", "detail": str(raw or "")}
def _coordination_root_id(ctx: Any) -> str:
metadata = getattr(ctx, "task_metadata", {})
metadata = metadata if isinstance(metadata, dict) else {}
return str(
metadata.get("root_task_id")
or getattr(ctx, "root_task_id", "")
or getattr(ctx, "task_id", "")
or ""
)
def _parent_intent_fact(ctx: Any) -> dict[str, Any]:
metadata = getattr(ctx, "task_metadata", {})
metadata = metadata if isinstance(metadata, dict) else {}
contract = (
getattr(ctx, "task_contract", None)
if isinstance(getattr(ctx, "task_contract", None), dict)
else metadata.get("task_contract")
if isinstance(metadata.get("task_contract"), dict)
else {}
)
budget = contract.get("delegation_budget") if isinstance(contract, dict) else {}
note = str((budget or {}).get("intent_note") or "").strip()
return {
"state": "present" if note else "absent",
"authority": "parent_authored_advisory",
"text": note,
}
def _time_fact(ctx: Any) -> dict[str, Any]:
try:
from ouroboros.task_pacing import build_budget_snapshot, resolve_budget_profile
snapshot = build_budget_snapshot(ctx, profile=resolve_budget_profile(ctx))
if not snapshot.has_deadline:
return {
"state": "not_set", "remaining_sec": None,
"reserve_sec": None, "inside_reserve": None, "expired": None,
}
return {
"state": "known",
"remaining_sec": round(max(0.0, snapshot.remaining_sec), 3),
"reserve_sec": round(max(0.0, snapshot.reserve_sec), 3),
"inside_reserve": bool(snapshot.inside_reserve),
"expired": bool(snapshot.remaining_sec <= 0),
}
except Exception as exc:
return {
"state": "unknown", "remaining_sec": None,
"reserve_sec": None, "inside_reserve": None, "expired": None,
"reason": type(exc).__name__,
}
def _settled_spend_fact(ctx: Any, root_task_id: str) -> dict[str, Any]:
try:
from ouroboros.usage_accounting import usage_breakdown
projection = usage_breakdown(
custody.custody_root(ctx), root_task_id=root_task_id,
)
integrity = bool(projection.get("integrity_degraded"))
unknown = int(projection.get("unknown_unmetered") or 0)
return {
"state": "partial" if integrity or unknown else "known",
"settled_usd": float(projection.get("settled_usd") or 0.0),
"accounted_usd": float(projection.get("accounted_usd") or 0.0),
"cost_final": bool(projection.get("cost_final")),
"unknown_unmetered": unknown,
"integrity_degraded": integrity,
}
except Exception as exc:
return {
"state": "unknown", "settled_usd": None, "accounted_usd": None,
"cost_final": None, "unknown_unmetered": None,
"integrity_degraded": None, "reason": type(exc).__name__,
}
def _active_descendants_fact(ctx: Any) -> dict[str, Any]:
"""Strict host-visible ancestry walk; vendor-internal children stay opaque."""
try:
from ouroboros.task_results import load_task_result
from ouroboros.task_status import _load_queue_snapshot, _snapshot_is_stale
root = custody.custody_root(ctx)
snapshot = _load_queue_snapshot(pathlib.Path(root))
if snapshot.get("_snapshot_missing") or snapshot.get("_snapshot_invalid"):
raise ValueError("queue_snapshot_unavailable")
if _snapshot_is_stale(snapshot):
raise ValueError("queue_snapshot_stale")
rows: dict[str, dict[str, Any]] = {}
for group, status in (("pending", "scheduled"), ("running", "running")):
for item in snapshot.get(group) or []:
if not isinstance(item, dict):
raise ValueError("queue_snapshot_row_invalid")
task = item.get("task") if isinstance(item.get("task"), dict) else item
task_id = str(item.get("id") or task.get("id") or "")
if not task_id:
raise ValueError("queue_snapshot_task_id_missing")
durable = load_task_result(root, task_id) or {}
rows[task_id] = {**durable, **task, "task_id": task_id, "status": status}
ancestor = str(getattr(ctx, "task_id", "") or "")
lineage_cache = dict(rows)
def _belongs(row: dict[str, Any]) -> bool:
parent_id = str(row.get("parent_task_id") or "")
seen = {str(row.get("task_id") or row.get("id") or "")}
while parent_id:
if parent_id == ancestor:
return True
if parent_id in seen:
raise ValueError("task_lineage_cycle")
seen.add(parent_id)
parent = lineage_cache.get(parent_id)
if parent is None:
parent = load_task_result(root, parent_id)
if not isinstance(parent, dict):
raise ValueError("task_lineage_unavailable")
lineage_cache[parent_id] = parent
parent_id = str(parent.get("parent_task_id") or "")
return False
active = [row for task_id, row in rows.items()
if task_id != ancestor and _belongs(row)]
by_status: dict[str, int] = {}
for row in active:
status = str(row.get("status") or "unknown").strip().lower() or "unknown"
by_status[status] = by_status.get(status, 0) + 1
return {
"state": "known", "count": len(active),
"by_status": dict(sorted(by_status.items())),
"scope": "host_visible_descendants",
"vendor_internal": "opaque_not_counted",
}
except Exception as exc:
return {
"state": "unknown", "count": None, "by_status": {},
"scope": "host_visible_descendants",
"vendor_internal": "opaque_not_counted",
"reason": type(exc).__name__,
}
def coordination_live_context(ctx: Any) -> dict[str, Any]:
"""One LLM-first planning snapshot for startup and meaningful nanny wakes."""
root_task_id = _coordination_root_id(ctx)
try:
from ouroboros.task_pacing import project_task_acceptance_review_capacity
review_capacity = project_task_acceptance_review_capacity(ctx)
except Exception as exc:
review_capacity = {
"state": "unknown", "reason": type(exc).__name__,
"root_task_id": root_task_id, "cap_cycles": None,
"claimed_cycles": None, "remaining_cycles": None,
"binding_seen": False, "dedupe": "task_acceptance_binding_sha256",
}
return {
"observed_at": utc_now_iso(),
"root_task_id": root_task_id,
"parent_intent": _parent_intent_fact(ctx),
"time": _time_fact(ctx),
"settled_spend": _settled_spend_fact(ctx, root_task_id),
"active_descendants": _active_descendants_fact(ctx),
"review_capacity": review_capacity,
}
def _coordination_cursor(state: dict[str, Any]) -> dict[str, Any]:
cursor = state.get("coordination_cursor")
if not isinstance(cursor, dict):
@ -345,6 +517,8 @@ def _render_wake_payload(ctx: Any, payload: dict[str, Any]) -> str:
}
envelope["supervision_wake_id"] = wake_id
envelope["wake_events"] = summaries
if isinstance(payload.get("coordination_context"), dict):
envelope["coordination_context"] = payload["coordination_context"]
envelope["wake_delivery"] = {
"complete": False,
"total_chars": len(raw),
@ -368,6 +542,23 @@ def _render_wake_payload(ctx: Any, payload: dict[str, Any]) -> str:
envelope["wake_delivery"]["wake_events_summarized"] = 0
envelope["wake_delivery"]["wake_events_omitted"] = len(events)
rendered = json.dumps(envelope, ensure_ascii=False, indent=2)
if len(rendered) > budget and isinstance(envelope.get("coordination_context"), dict):
context = envelope["coordination_context"]
envelope["coordination_context"] = {
"state": "available_in_full_wake_source",
"observed_at": str(context.get("observed_at") or ""),
"root_task_id": str(context.get("root_task_id") or ""),
}
rendered = json.dumps(envelope, ensure_ascii=False, indent=2)
if len(rendered) > budget:
envelope = {
"status": str(payload.get("status") or "wake_available")[:120],
"run_id": str(payload.get("run_id") or "")[:200],
"supervision_wake_id": wake_id,
"coordination_context": {"state": "available_in_full_wake_source"},
"wake_delivery": envelope["wake_delivery"],
}
rendered = json.dumps(envelope, ensure_ascii=False, indent=2)
return rendered
@ -628,6 +819,7 @@ def supervised_wait(
}
if wakes:
payload["wake_events"] = wakes
payload["coordination_context"] = coordination_live_context(ctx)
wake_id = uuid.uuid4().hex
payload["supervision_wake_id"] = wake_id
state["status"] = "wake_pending"
@ -704,5 +896,5 @@ def delegate_wait_entry(
__all__ = [
"acknowledge_pending_wake", "delegate_wait_entry", "supervised_wait",
"supervision_checkpoint",
"supervision_checkpoint", "coordination_live_context",
]

View file

@ -1593,6 +1593,15 @@ def _execute_task_acceptance_panel(ctx: _TaskAcceptanceContext) -> Any:
},
task_id=ctx.task_id,
)
if not slots:
return ReviewRunResult(
request={"surface": "task_acceptance", "task_id": str(ctx.task_id)},
actors=[],
parsed_findings=[],
aggregate_signal="DEGRADED",
degraded=True,
degraded_reasons=["no_review_slots"],
)
# Budget admission for the whole acceptance wave (v6.69.0): a wave that
# cannot fit the remaining root budget is declined up front as a terminal
# DEGRADED (no-quorum semantics) instead of dying mid-wave. The estimate
@ -1626,6 +1635,61 @@ def _execute_task_acceptance_panel(ctx: _TaskAcceptanceContext) -> Any:
f"${_admission.get('remaining_usd')} (no reviewer was called)"
],
)
# Q6: one strict tree-wide paid-cycle authority, claimed after every free
# launch refusal and immediately before the physical reviewer transport.
# The immutable exact binding is the dedupe key. A crash after this claim
# remains an honest unknown and cannot silently buy the same panel again.
try:
from ouroboros.task_results import (
claim_task_acceptance_review_cycle,
resolve_task_lineage,
)
tools_ctx = ctx.tools._ctx
metadata = getattr(tools_ctx, "task_metadata", {})
metadata = metadata if isinstance(metadata, dict) else {}
lineage = resolve_task_lineage(ctx.task_id, metadata=metadata)
root_task_id = str(lineage.get("root_task_id") or ctx.task_id)
accounting_root = pathlib.Path(str(
metadata.get("budget_drive_root")
or getattr(tools_ctx, "budget_drive_root", "")
or ctx.drive_root
or getattr(tools_ctx, "drive_root", ".")
))
required_blocking = bool(
ctx.mode == "required" and get_review_enforcement() == "blocking"
)
snapshot = task_pacing.build_budget_snapshot(
tools_ctx, profile=ctx.budget_profile,
)
max_cycles = task_pacing.effective_task_acceptance_review_cycles(
ctx.budget_profile,
has_deadline=snapshot.has_deadline,
required_blocking=required_blocking,
)
claim = claim_task_acceptance_review_cycle(
accounting_root,
root_task_id,
ctx.review_binding,
max_cycles=max_cycles,
claimed_by_task_id=ctx.task_id,
allow_create=bool(lineage.get("is_root_task")),
)
except Exception as exc:
claim = {
"status": "unknown",
"reason": f"review_capacity_claim_unknown:{type(exc).__name__}",
}
if claim.get("status") != "claimed":
reason = str(claim.get("reason") or "review_capacity_unknown")
return ReviewRunResult(
request={"surface": "task_acceptance", "task_id": str(ctx.task_id)},
actors=[],
parsed_findings=[],
aggregate_signal="DEGRADED",
degraded=True,
degraded_reasons=[f"{reason} (no reviewer was called)"],
)
started = time.monotonic()
result = run_review_request(
request,
@ -5738,66 +5802,15 @@ def _nanny_finalization_message(
"or state in your final answer that the delegated run failed and why "
"the remaining work ran on metered API tokens."
)
# Actor-first configured sessions may legitimately spend the episode on
# host-side coordination instead of starting a physical leaf. A typed
# coordination tool call is enough evidence for this decision; when a
# continuation has already persisted a direct child, recover that fact from
# the existing task-result SSOT as well. Neither path creates a second
# ledger or infers topology from prose.
host_coordination = bool(getattr(tools._ctx, "_nanny_coordination_activity", False))
bootstrap = getattr(tools._ctx, "_configured_actor_bootstrap", None)
if isinstance(bootstrap, dict) and bool(bootstrap.get("zero_run_receipt_recorded")):
# The durable no-leaf decision is itself the actor's explicit host-side
# coordination outcome. This marker is hydrated on resume from the exact
# receipt, so a continuation must not accuse the actor of doing nothing.
host_coordination = True
metadata = getattr(tools._ctx, "task_metadata", {})
metadata = metadata if isinstance(metadata, dict) else {}
status_root = pathlib.Path(str(
metadata.get("budget_drive_root")
or getattr(tools._ctx, "budget_drive_root", "")
or drive_root
))
if not host_coordination:
try:
from ouroboros.task_status import find_child_tasks
from ouroboros.subagent_bootstrap import (
actor_first_coordination_finalization_message,
)
host_coordination = bool(find_child_tasks(
status_root,
parent_task_id=str(task_id or ""),
exclude_task_id=str(task_id or ""),
scope="direct",
materialize_artifacts=False,
))
except Exception:
log.debug("nanny nudge: host coordination evidence read failed", exc_info=True)
# Actor-first has one additional truth obligation: a plain final answer is
# not evidence that the assigned physical leaf was intentionally skipped.
# Reuse the bootstrap/child evidence helper so a missing route or unreadable
# child store becomes a typed incomplete/unknown outcome rather than a clean
# result. This remains a one-shot advisory reminder; the durable outcome
# projection enforces the same fact after the model returns.
try:
from ouroboros.subagent_bootstrap import actor_first_unresolved_fact
actor_fact = actor_first_unresolved_fact(
tools._ctx, task_id=str(task_id or ""), drive_root=status_root,
)
except Exception:
actor_fact = None
if actor_fact:
status = str(actor_fact.get("status") or "unknown")
code = "CONFIGURED_ACTOR_INCOMPLETE" if status == "incomplete" else "CONFIGURED_ACTOR_UNKNOWN"
return (
f"⚠️ {code}: this actor-first session is finalizing before its assigned "
"physical leaf started and without a direct host child result. Call "
"verify_and_record(contract_kind=delegation_zero_run, "
f"zero_run_decision={status!r}, zero_run_basis=...) to record the "
"typed terminal decision, or start the exact assigned session now. "
"A plain prose answer cannot close this evidence gap."
)
if host_coordination:
return ""
actor_first = actor_first_coordination_finalization_message(
tools._ctx, task_id=str(task_id or ""), fallback_root=drive_root,
)
if actor_first is not None:
return actor_first
return (
"⚠️ NANNY_DID_NOT_DELEGATE: this task was dispatched onto the delegated "
"substrate (executor=harness), but you are finalizing with ZERO "

View file

@ -26,6 +26,7 @@ from ouroboros.config import get_review_models, review_model_uses_local
from ouroboros.llm import LLMClient
from ouroboros.observability import new_call_id, persist_call, redact_projection
from ouroboros.provider_models import provider_for_model
from ouroboros.task_results import review_binding_hash
# Everything below the seam. Re-exported here because the substrate is the
# historical import site for the api_chat prompt renderers; `review_execution`
# owns them now and must never import this module back.
@ -393,9 +394,7 @@ def build_review_binding(
"evidence_revision": evidence_revision,
"fence_hash": fence_hash,
}
binding_hash = hashlib.sha256(
json.dumps(binding_payload, sort_keys=True, separators=(",", ":")).encode("utf-8")
).hexdigest()
binding_hash = review_binding_hash(**binding_payload)
return {
**binding_payload,
"binding_hash": binding_hash,

View file

@ -185,7 +185,6 @@ BAND_PATHS = {
"tests/test_telegram_miniapp_lifecycle.py": None,
"tests/test_tool_api_v2_public_surface.py": None,
"tests/test_usage_accounting.py": None,
"tests/test_v647_megacommit.py": None,
"tests/test_v6730_origin_invariant.py": None,
"tests/test_v678_receipt_reconciliation.py": None,
"web/modules/api_types.js": "The shared browser contract module now includes issue 265 publication-preflight types alongside the target settings and subagent contracts.",
@ -208,7 +207,7 @@ BYTE_BASELINE_DEBT = {
}
BYTE_DEBT = {
"ouroboros/loop.py": 316096,
"ouroboros/loop.py": 315869,
"tests/test_delegated_subagent_transport.py": 320571,
"tests/test_devtools_benchmarks.py": 328282,
"web/modules/chat.js": 225761,

View file

@ -3,6 +3,7 @@
from __future__ import annotations
import json
import logging
from hashlib import sha256
from pathlib import Path
from typing import Any, Mapping
@ -14,6 +15,25 @@ from ouroboros.subagent_work_order import (
route_source_request_channel,
)
log = logging.getLogger(__name__)
def _with_coordination_context(ctx: Any, raw: str) -> str:
"""Attach one fresh planning snapshot to a startup/recovery receipt."""
if not raw:
return raw
try:
payload = json.loads(raw)
if not isinstance(payload, dict):
return raw
from ouroboros.delegate_supervision import coordination_live_context
payload["coordination_context"] = coordination_live_context(ctx)
return json.dumps(payload, ensure_ascii=False, indent=2)
except Exception:
return raw
def bootstrap_before_context(ctx: Any, task: Mapping[str, Any], dispatch: Any) -> str:
"""Freeze the selected route before the first ordinary actor episode.
@ -37,7 +57,7 @@ def bootstrap_before_context(ctx: Any, task: Mapping[str, Any], dispatch: Any) -
actor_ready = _prepare_actor_first_bootstrap(ctx, task, dispatch)
recovery = _adopt_recovery_handoff(ctx, task)
if recovery:
return recovery
return _with_coordination_context(ctx, recovery)
if bool(getattr(dispatch, "blocked", False)):
# A route fault is evidence for the ordinary host turn, never a
# reason to silently spend on native/API fallback.
@ -85,10 +105,10 @@ def bootstrap_before_context(ctx: Any, task: Mapping[str, Any], dispatch: Any) -
custody.emit(custody.custody_root(ctx), "configured_subagent_startup_fault", {
"task_id": str(getattr(ctx, "task_id", "") or ""), **startup,
})
return json.dumps({
return _with_coordination_context(ctx, json.dumps({
"status": "configured_session_actor_ready", "startup": startup,
}, ensure_ascii=False, indent=2)
return actor_ready
}, ensure_ascii=False, indent=2))
return _with_coordination_context(ctx, actor_ready)
return ""
@ -238,6 +258,58 @@ def actor_first_unresolved_fact(
}
def actor_first_coordination_finalization_message(
ctx: Any, *, task_id: str, fallback_root: Any,
) -> str | None:
"""Resolve actor-first finalization, or defer to the legacy nanny message."""
host_coordination = bool(getattr(ctx, "_nanny_coordination_activity", False))
bootstrap = getattr(ctx, "_configured_actor_bootstrap", None)
if isinstance(bootstrap, dict) and bootstrap.get("zero_run_receipt_recorded"):
host_coordination = True
metadata = getattr(ctx, "task_metadata", {})
metadata = metadata if isinstance(metadata, dict) else {}
status_root = Path(str(
metadata.get("budget_drive_root")
or getattr(ctx, "budget_drive_root", "")
or fallback_root
))
if not host_coordination:
try:
from ouroboros.task_status import find_child_tasks
host_coordination = bool(find_child_tasks(
status_root,
parent_task_id=str(task_id or ""),
exclude_task_id=str(task_id or ""),
scope="direct",
materialize_artifacts=False,
))
except Exception:
log.debug("nanny nudge: host coordination evidence read failed", exc_info=True)
try:
actor_fact = actor_first_unresolved_fact(
ctx, task_id=str(task_id or ""), drive_root=status_root,
)
except Exception:
actor_fact = None
if actor_fact:
status = str(actor_fact.get("status") or "unknown")
code = (
"CONFIGURED_ACTOR_INCOMPLETE"
if status == "incomplete" else "CONFIGURED_ACTOR_UNKNOWN"
)
return (
f"⚠️ {code}: this actor-first session is finalizing before its assigned "
"physical leaf started and without a direct host child result. Call "
"verify_and_record(contract_kind=delegation_zero_run, "
f"zero_run_decision={status!r}, zero_run_basis=...) to record the "
"typed terminal decision, or start the exact assigned session now. "
"A plain prose answer cannot close this evidence gap."
)
return "" if host_coordination else None
def _durable_zero_run_receipt(
ctx: Any, *, gap_reasons: set[str] | None = None,
) -> dict[str, Any] | None:
@ -507,4 +579,8 @@ def append_startup_receipt(
acknowledge_pending_wake(ctx, startup_wake)
__all__ = ["append_startup_receipt", "bootstrap_before_context"]
__all__ = [
"actor_first_coordination_finalization_message",
"append_startup_receipt",
"bootstrap_before_context",
]

View file

@ -547,6 +547,13 @@ def prepare_delegate_start_actor(
snapshot, route = exact_session_binding(selected_snapshot)
selected_id = str(snapshot.get("selected_subagent_id") or "")
config_fingerprint = str(snapshot.get("config_fingerprint") or "")
if custody.custody_log_unreadable(drive_root):
return {}, _fail(
"delegate_start", "replacement_custody_unknown",
"The custody event log exists but cannot be read, so the host cannot "
"prove that no physical start/run remains. Repair or reconcile the "
"canonical custody record before starting a replacement.",
)
blockers = unsettled_start_ids(
drive_root, str(getattr(ctx, "task_id", "") or "")
)

View file

@ -305,6 +305,12 @@ def effective_max_improvement_passes(
return None if cap is None else max(0, int(cap))
from ouroboros.task_results import ( # noqa: E402,F401 - compatibility re-export
effective_task_acceptance_review_cycles,
project_task_acceptance_review_capacity,
)
def improvement_pass_allowed(
snapshot: BudgetSnapshot,
passes_done: int,

View file

@ -23,6 +23,169 @@ STATUS_FAILED = "failed"
STATUS_INTERRUPTED = "interrupted"
STATUS_CANCELLED = "cancelled"
def review_binding_hash(
*, candidate_hash: str, evidence_revision: str, fence_hash: str,
) -> str:
"""Digest the immutable task-acceptance binding components."""
import hashlib
payload = {
"candidate_hash": str(candidate_hash or ""),
"evidence_revision": str(evidence_revision or ""),
"fence_hash": str(fence_hash or ""),
}
return hashlib.sha256(
json.dumps(payload, sort_keys=True, separators=(",", ":")).encode("utf-8")
).hexdigest()
def effective_task_acceptance_review_cycles(
profile: Dict[str, Any], *, has_deadline: bool = True,
required_blocking: bool = False,
) -> Optional[int]:
"""Project paid panels from the existing improvement-pass semantics."""
from ouroboros.task_pacing import effective_max_improvement_passes
passes = effective_max_improvement_passes(
profile,
has_deadline=has_deadline,
required_blocking=required_blocking,
)
return None if passes is None else max(1, int(passes) + 1)
def project_task_acceptance_review_capacity(
ctx: Any, *, binding_hash: str = "",
) -> Dict[str, Any]:
"""Read the canonical root's paid acceptance-wallet projection.
Descendants may observe but never initialize root authority. A missing or
malformed canonical result is UNKNOWN for them; a live root may begin with
the known empty state. The atomic claim remains dispatch authority.
"""
from ouroboros import config, task_pacing
from ouroboros.contracts.task_contract import normalize_budget_profile
from ouroboros.deadline_utils import parse_deadline_ts
metadata = getattr(ctx, "task_metadata", {})
metadata = metadata if isinstance(metadata, dict) else {}
task_id = str(getattr(ctx, "task_id", "") or "")
lineage = resolve_task_lineage(
task_id,
metadata=metadata,
root_task_id=getattr(ctx, "root_task_id", None),
parent_task_id=getattr(ctx, "parent_task_id", None),
delegation_role=getattr(ctx, "delegation_role", None),
original_task_id=getattr(ctx, "original_task_id", None),
timeout_retry_from=getattr(ctx, "timeout_retry_from", None),
)
root_task_id = str(lineage.get("root_task_id") or task_id)
root = pathlib.Path(str(
metadata.get("budget_drive_root")
or getattr(ctx, "budget_drive_root", "")
or getattr(ctx, "drive_root", ".")
))
base = {
"root_task_id": root_task_id,
"cap_cycles": None,
"claimed_cycles": None,
"remaining_cycles": None,
"binding_seen": False,
"dedupe": "task_acceptance_binding_sha256",
}
try:
state = load_task_acceptance_review_state(
root,
root_task_id,
require_root_result=not bool(lineage.get("is_root_task")),
)
path = task_result_path(root, root_task_id, create=False)
root_result = read_json_dict(path) if path.is_file() else {}
if path.is_file() and root_result is None:
raise ValueError("root result is malformed")
root_result = root_result if isinstance(root_result, dict) else {}
root_contract = (
root_result.get("task_contract")
if isinstance(root_result.get("task_contract"), dict)
else {}
)
if not root_contract and lineage.get("is_root_task"):
root_contract = (
getattr(ctx, "task_contract", None)
if isinstance(getattr(ctx, "task_contract", None), dict)
else metadata.get("task_contract")
if isinstance(metadata.get("task_contract"), dict)
else {}
)
profile = normalize_budget_profile(root_contract.get("budget_profile"))
deadline_value = root_result.get("deadline_at")
if deadline_value is None and lineage.get("is_root_task"):
deadline_value = metadata.get("deadline_at")
required_blocking = bool(
config.get_task_review_mode() == "required"
and config.get_review_enforcement() == "blocking"
)
cap = effective_task_acceptance_review_cycles(
profile,
has_deadline=parse_deadline_ts(deadline_value) is not None,
required_blocking=required_blocking,
)
claims = state.get("claims_by_binding") or {}
claimed = len(claims)
remaining = None if cap is None else max(0, cap - claimed)
requested_binding = str(binding_hash or "").strip().lower()
projection = {
**base,
"state": "available",
"reason": "",
"cap_cycles": cap,
"claimed_cycles": claimed,
"remaining_cycles": remaining,
"binding_seen": bool(requested_binding and requested_binding in claims),
}
try:
from ouroboros.cancel_intents import cancel_pending
if cancel_pending(root, root_task_id) or (
task_id != root_task_id and cancel_pending(root, task_id)
):
projection.update({
"state": "unavailable", "reason": "cancellation_pending",
})
return projection
except Exception as exc:
return {
**projection,
"state": "unknown",
"reason": f"cancellation_state_unknown:{type(exc).__name__}",
}
budget = task_pacing.build_budget_snapshot(
ctx, profile=task_pacing.resolve_budget_profile(ctx),
)
launch_ok, launch_reason = task_pacing.review_launch_allowed(
budget,
estimated_sec=task_pacing.acceptance_review_estimate_sec(
ctx, passes_done=claimed,
),
)
if not launch_ok:
projection.update({"state": "unavailable", "reason": launch_reason})
elif remaining == 0:
projection.update({
"state": "unavailable", "reason": "review_cycles_exhausted",
})
return projection
except Exception as exc:
return {
**base,
"state": "unknown",
"reason": f"review_capacity_unknown:{type(exc).__name__}",
}
# Intent latch: the agent/owner asked to cancel, but the supervisor has not yet
# torn the task down. Ranks above running so a late running/scheduled mirror
# cannot resurrect it, but below the truly-terminal statuses so the eventual
@ -79,6 +242,229 @@ _PLAN_REVIEW_STATE_MAX_BYTES = 1_000_000
_PLAN_REVIEW_HASH_RE = re.compile(r"^[0-9a-f]{64}$")
_PLAN_REVIEW_REASON_MAX_CHARS = 2_000
TASK_ACCEPTANCE_REVIEW_STATE_KEY = "task_acceptance_review_accounting"
_TASK_ACCEPTANCE_REVIEW_STATE_VERSION = 1
_TASK_ACCEPTANCE_REVIEW_STATE_MAX_BYTES = 1_000_000
_TASK_ACCEPTANCE_REVIEW_CLAIM_FIELDS = frozenset({
"binding_hash", "candidate_hash", "evidence_revision", "fence_hash",
"claimed_at", "claimed_by_task_id",
})
def _empty_task_acceptance_review_state(root_task_id: str) -> Dict[str, Any]:
return {
"schema_version": _TASK_ACCEPTANCE_REVIEW_STATE_VERSION,
"root_task_id": str(root_task_id),
"claims_by_binding": {},
}
def _validated_task_acceptance_review_state(
value: Any, root_task_id: str,
) -> Dict[str, Any]:
"""Strict private copy of the root tree's paid acceptance claims."""
if not isinstance(value, dict) or value.get("schema_version") != 1:
raise ValueError("TASK_ACCEPTANCE_REVIEW_STATE_INVALID: unsupported schema")
if set(value) != {"schema_version", "root_task_id", "claims_by_binding"}:
raise ValueError("TASK_ACCEPTANCE_REVIEW_STATE_INVALID: state shape is invalid")
if str(value.get("root_task_id") or "") != str(root_task_id):
raise ValueError("TASK_ACCEPTANCE_REVIEW_STATE_INVALID: root identity mismatch")
claims = value.get("claims_by_binding")
if not isinstance(claims, dict):
raise ValueError("TASK_ACCEPTANCE_REVIEW_STATE_INVALID: claims must be an object")
for binding_hash, claim in claims.items():
if not _PLAN_REVIEW_HASH_RE.fullmatch(str(binding_hash or "")):
raise ValueError("TASK_ACCEPTANCE_REVIEW_STATE_INVALID: claim key is invalid")
if not isinstance(claim, dict) or set(claim) != _TASK_ACCEPTANCE_REVIEW_CLAIM_FIELDS:
raise ValueError("TASK_ACCEPTANCE_REVIEW_STATE_INVALID: claim shape is invalid")
if str(claim.get("binding_hash") or "") != str(binding_hash):
raise ValueError("TASK_ACCEPTANCE_REVIEW_STATE_INVALID: claim identity mismatch")
for key in ("binding_hash", "candidate_hash", "evidence_revision", "fence_hash"):
if not _PLAN_REVIEW_HASH_RE.fullmatch(str(claim.get(key) or "")):
raise ValueError(
f"TASK_ACCEPTANCE_REVIEW_STATE_INVALID: {key} is invalid"
)
expected_binding = review_binding_hash(
candidate_hash=str(claim["candidate_hash"]),
evidence_revision=str(claim["evidence_revision"]),
fence_hash=str(claim["fence_hash"]),
)
if expected_binding != str(binding_hash):
raise ValueError(
"TASK_ACCEPTANCE_REVIEW_STATE_INVALID: binding digest mismatch"
)
for key in ("claimed_at", "claimed_by_task_id"):
if not isinstance(claim.get(key), str) or not str(claim.get(key) or ""):
raise ValueError(
f"TASK_ACCEPTANCE_REVIEW_STATE_INVALID: {key} must be non-empty text"
)
copied = copy.deepcopy(value)
if len(json.dumps(copied, ensure_ascii=False).encode("utf-8")) > _TASK_ACCEPTANCE_REVIEW_STATE_MAX_BYTES:
raise ValueError("TASK_ACCEPTANCE_REVIEW_STATE_INVALID: state is too large")
return copied
def load_task_acceptance_review_state(
results_drive_root: Any,
root_task_id: str,
*,
require_root_result: bool = False,
) -> Dict[str, Any]:
"""Read the canonical root's shared paid-review claims without mutation."""
path = task_result_path(results_drive_root, root_task_id, create=False)
if not path.is_file():
if require_root_result:
raise ValueError("TASK_ACCEPTANCE_REVIEW_STATE_UNKNOWN: root result is absent")
return _empty_task_acceptance_review_state(root_task_id)
result = read_json_dict(path)
if result is None:
raise ValueError("TASK_ACCEPTANCE_REVIEW_STATE_INVALID: root result is malformed")
if str(result.get("task_id") or "") != str(root_task_id):
raise ValueError("TASK_ACCEPTANCE_REVIEW_STATE_INVALID: root identity mismatch")
stored_root_id = str(result.get("root_task_id") or "")
if stored_root_id and stored_root_id != str(root_task_id):
raise ValueError("TASK_ACCEPTANCE_REVIEW_STATE_INVALID: root identity mismatch")
if TASK_ACCEPTANCE_REVIEW_STATE_KEY not in result:
return _empty_task_acceptance_review_state(root_task_id)
return _validated_task_acceptance_review_state(
result[TASK_ACCEPTANCE_REVIEW_STATE_KEY], root_task_id,
)
def _update_task_acceptance_review_state(
results_drive_root: Any,
root_task_id: str,
mutator: Callable[[Dict[str, Any]], Optional[Dict[str, Any]]],
*,
allow_create: bool,
) -> Dict[str, Any]:
"""Strict root-result update; the file lock is the tree-wide claim fence."""
path = task_result_path(results_drive_root, root_task_id, create=allow_create)
if not allow_create and not path.is_file():
raise ValueError("TASK_ACCEPTANCE_REVIEW_STATE_UNKNOWN: root result is absent")
def _merge(existing: Dict[str, Any]) -> Optional[Dict[str, Any]]:
if not allow_create and not existing:
raise ValueError("TASK_ACCEPTANCE_REVIEW_STATE_UNKNOWN: root result is absent")
if existing and str(existing.get("task_id") or "") != str(root_task_id):
raise ValueError("TASK_ACCEPTANCE_REVIEW_STATE_INVALID: root identity mismatch")
stored_root_id = str(existing.get("root_task_id") or "")
if stored_root_id and stored_root_id != str(root_task_id):
raise ValueError("TASK_ACCEPTANCE_REVIEW_STATE_INVALID: root identity mismatch")
state = (
_validated_task_acceptance_review_state(
existing[TASK_ACCEPTANCE_REVIEW_STATE_KEY], root_task_id,
)
if TASK_ACCEPTANCE_REVIEW_STATE_KEY in existing
else _empty_task_acceptance_review_state(root_task_id)
)
candidate = mutator(state)
if candidate is None:
return None
updated = _validated_task_acceptance_review_state(
candidate, root_task_id,
)
now = utc_now_iso()
return {
**existing,
TASK_ACCEPTANCE_REVIEW_STATE_KEY: updated,
"task_id": str(root_task_id),
"status": str(existing.get("status") or STATUS_RUNNING),
"ts": str(existing.get("ts") or now),
"updated_at": now,
}
try:
updated = update_json_locked(
path,
_merge,
strict_existing_dict=True,
reject_existing_empty_dict=True,
)
except ValueError as exc:
if str(exc).startswith("update_json_locked:"):
raise ValueError(
"TASK_ACCEPTANCE_REVIEW_STATE_INVALID: root result is malformed"
) from exc
raise
return _validated_task_acceptance_review_state(
updated.get(TASK_ACCEPTANCE_REVIEW_STATE_KEY), root_task_id,
)
def claim_task_acceptance_review_cycle(
results_drive_root: Any,
root_task_id: str,
review_binding: Dict[str, Any],
*,
max_cycles: Optional[int],
claimed_by_task_id: str,
allow_create: bool = False,
) -> Dict[str, Any]:
"""Atomically dedupe and claim one paid root-acceptance panel dispatch."""
binding_fields = {
key: str((review_binding or {}).get(key) or "").strip().lower()
for key in ("binding_hash", "candidate_hash", "evidence_revision", "fence_hash")
}
if any(
not _PLAN_REVIEW_HASH_RE.fullmatch(value)
for value in binding_fields.values()
):
raise ValueError("TASK_ACCEPTANCE_REVIEW_STATE_INVALID: review binding is invalid")
expected_binding = review_binding_hash(
candidate_hash=binding_fields["candidate_hash"],
evidence_revision=binding_fields["evidence_revision"],
fence_hash=binding_fields["fence_hash"],
)
if binding_fields["binding_hash"] != expected_binding:
raise ValueError("TASK_ACCEPTANCE_REVIEW_STATE_INVALID: binding digest mismatch")
binding = binding_fields["binding_hash"]
claimant = str(claimed_by_task_id or "").strip()
if not claimant:
raise ValueError("TASK_ACCEPTANCE_REVIEW_STATE_INVALID: claimant is absent")
cap = None if max_cycles is None else max(1, int(max_cycles))
decision: Dict[str, Any] = {}
def _claim(state: Dict[str, Any]) -> Optional[Dict[str, Any]]:
claims = dict(state.get("claims_by_binding") or {})
prior = claims.get(binding)
if prior is not None:
decision.update({
"status": "unknown",
"reason": "binding_dispatch_already_claimed",
})
return None
if cap is not None and len(claims) >= cap:
decision.update({"status": "unavailable", "reason": "review_cycles_exhausted"})
return None
claims[binding] = {
**binding_fields,
"claimed_at": utc_now_iso(),
"claimed_by_task_id": claimant,
}
state["claims_by_binding"] = claims
decision.update({"status": "claimed", "reason": ""})
return state
state = _update_task_acceptance_review_state(
results_drive_root,
root_task_id,
_claim,
allow_create=allow_create,
)
paid = len(state.get("claims_by_binding") or {})
return {
**decision,
"binding_hash": binding,
"cycles_paid": paid,
"max_cycles": cap,
"remaining_cycles": None if cap is None else max(0, cap - paid),
}
def cancellation_blocks_child_result(result: Any) -> bool:
"""Return whether canonical cancellation forbids child-drive promotion.

View file

@ -811,9 +811,9 @@ def _delegate_start(ctx: ToolContext, prompt: str, max_seconds: Optional[int] =
lineage = getattr(ctx, "task_metadata", {}) or {}
lineage = lineage if isinstance(lineage, dict) else {}
# Fresh payload run: busy check + durable write = ONE atomic claim (fix 5).
requested, claim_holder = claimed_start_request(
drive, claim_target=(target_root if (not recovering and
authority_source == "skill_payload") else ""),
requested, claim_refusal = claimed_start_request(
drive, claim_target=(target_root if not recovering and authority_source == "skill_payload" else ""),
actor_ctx=ctx, enforce_actor_idle=not recovering,
run_id="", task_id=str(getattr(ctx, "task_id", "") or ""),
idempotency_key=key, invocation_id=invocation_id,
max_seconds=seconds, request=request_body, project_id=project_id,
@ -837,18 +837,19 @@ def _delegate_start(ctx: ToolContext, prompt: str, max_seconds: Optional[int] =
authority_fingerprint=authority_fingerprint,
work_order_source_request=work_order_source_request,
)
if claim_holder:
if claim_refusal:
reason = str(claim_refusal.get("reason") or "replacement_custody_unknown")
detail = str(claim_refusal.get("detail") or "Actor start claim unavailable.")
facts = {key: value for key, value in claim_refusal.items()
if key not in {"reason", "detail"}}
return _fail(
"delegate_start", "payload_delegation_busy",
"Another delegated run claimed this exact payload first (the busy "
"check and the start-request write are one atomic claim). Finish "
"that run before starting another delegation against the same "
"skill.", holder=claim_holder,
"delegate_start", reason, detail,
**facts,
**_retire_orphaned_registration(ctx, gateway, owned_project_id,
definite_refusal=True,
reason="payload_delegation_busy",
invocation_id=invocation_id,
snapshot_id=snapshot_id))
definite_refusal=True,
reason=reason, invocation_id=invocation_id, snapshot_id=snapshot_id,
),
)
if not requested:
# The POST is CONDITIONAL on the durable request row: a run started
# without it is live and unfindable if this worker dies. A fresh

View file

@ -14,7 +14,6 @@ from __future__ import annotations
import json
import logging
import pathlib
import threading
from typing import TYPE_CHECKING, Any, Dict, NamedTuple, Optional, Tuple
from ouroboros import delegate_custody as custody
@ -539,6 +538,23 @@ def capture_terminal_patch_for_drive(drive: Any, entry: _RunCustody) -> Optional
return _capture_block(entry, cap_dir, manifest)
def capture_stranded_patch(drive_root: Any, run: _RunCustody) -> Dict[str, Any]:
"""Capture a reconciled dead-owner run without deciding its disposition."""
if not (run.execution_root and run.settled and not run.patch_disposed):
return {}
try:
block = capture_terminal_patch_for_drive(drive_root, run) or {}
except Exception:
log.warning("Reconcile patch capture failed for %s", run.run_id, exc_info=True)
return {"patch_capture": "failed", "patch_disposition": "pending"}
return {
"patch_capture": str(block.get("status") or ""),
"patch_artifact": block.get("patch_artifact"),
"patch_disposition": "pending",
}
# -- exact skill-payload delegation (R1) ---------------------------------------
#
# The restored delegated coding target class: an ordinary top-level task selects
@ -652,47 +668,19 @@ def payload_host_instructions(text: str, skill_name: str) -> str:
"review then re-runs before the new content is relied on.")
# In-process half of the atomic payload start claim; the cross-process half is
# the O_EXCL lockfile below. Both exist because two parallel delegate_start
# calls (same process or two workers) must produce exactly ONE started run.
_PAYLOAD_CLAIM_LOCK = threading.Lock()
def claimed_start_request(
drive: pathlib.Path, *, claim_target: str, **request_row: Any,
) -> Tuple[bool, str]:
"""Write the START_REQUESTED row, atomically fused with the payload busy check.
drive: pathlib.Path, *, claim_target: str, actor_ctx: Any = None,
enforce_actor_idle: bool = False, **request_row: Any,
) -> Tuple[bool, Dict[str, Any]]:
"""Compatibility facade for the one atomic pre-transport claim owner."""
Gate fix 5: for a payload run (``claim_target`` non-empty) the busy check
and the durable request write happen under ONE claim lock, so exactly one
of two synchronized starts wins; the loser gets the holder id back and
refuses typed. Non-payload rows pass straight through.
Returns ``(requested, busy_holder)``.
"""
if not claim_target:
return custody.record_start_requested(drive, **request_row), ""
from ouroboros.platform_layer import (
acquire_exclusive_file_lock,
release_exclusive_file_lock,
from ouroboros.delegate_start_claims import claimed_start_request as claim
return claim(
drive, claim_target=claim_target, payload_busy=_payload_delegation_busy,
actor_ctx=actor_ctx, enforce_actor_idle=enforce_actor_idle, **request_row,
)
lock_path = pathlib.Path(drive) / "state" / ".payload_delegation_claim.lock"
lock_path.parent.mkdir(parents=True, exist_ok=True)
with _PAYLOAD_CLAIM_LOCK:
fd = acquire_exclusive_file_lock(lock_path, timeout_sec=20.0, stale_sec=120.0)
if fd is None:
# Fail CLOSED: without the cross-process half the claim would be a
# plain unlocked read again — the exact race this fix removes.
return False, "(payload claim lock unavailable — another start holds it)"
try:
holder = _payload_delegation_busy(drive, pathlib.Path(claim_target))
if holder:
return False, holder
return custody.record_start_requested(drive, **request_row), ""
finally:
if fd is not None:
release_exclusive_file_lock(lock_path, fd)
def _payload_mutation_authority(
ctx: ToolContext, drive: pathlib.Path, bucket: str, skill_name: str,
@ -1538,6 +1526,7 @@ __all__ = [
"_retry_binding_refusal",
"_validated_invocation",
"capture_terminal_patch_for_drive",
"capture_stranded_patch",
"claimed_start_request",
"integrate_payload_patch",
"payload_content_hash",

View file

@ -424,25 +424,6 @@ def _record_delegation_zero_run(
"⚠️ TOOL_ARG_ERROR (verify_and_record): delegation_zero_run already "
"has a durable terminal decision for this actor."
)
try:
from ouroboros import delegate_custody as custody
from ouroboros.delegate_recovery import unsettled_start_ids
blockers = unsettled_start_ids(custody.custody_root(ctx), task_id)
except Exception as exc:
return (
"⚠️ TOOL_ERROR (verify_and_record): zero_run_custody_unknown: "
"the host could not prove that no physical start/run remains; no "
f"zero-run receipt was written ({type(exc).__name__})."
)
if any(blockers.values()):
return (
"⚠️ TOOL_ERROR (verify_and_record): zero_run_requires_settlement: "
"this task has an open run, ambiguous start invocation, or undisposed "
"physical result. Reconcile it before claiming that no physical leaf "
"started. blockers="
+ json.dumps(blockers, ensure_ascii=False, sort_keys=True)
)
decision = str(
kwargs.get("zero_run_decision") or kwargs.get("decision") or ""
).strip().lower()
@ -477,9 +458,54 @@ def _record_delegation_zero_run(
}
if crit := str(criterion_id or "").strip():
receipt["criterion_id"] = crit[:120]
# A zero-run is lifecycle authority for the actor, so keep it on the
# canonical budget root rather than an ephemeral child drive.
if not append_verification_receipt(canonical_data_root(ctx), task_id, receipt):
try:
from ouroboros import delegate_custody as custody
from ouroboros.delegate_recovery import unsettled_start_ids
from ouroboros.subagent_bootstrap import _durable_zero_run_receipt
drive_root = custody.custody_root(ctx)
with custody.actor_decision_lock(drive_root, task_id):
if custody.custody_log_unreadable(drive_root):
return (
"⚠️ TOOL_ERROR (verify_and_record): zero_run_custody_unknown: "
"the custody event log exists but cannot be read, so the host "
"cannot prove that no physical start/run remains; no zero-run "
"receipt was written (custody_log_unreadable)."
)
gaps: set[str] = set()
prior = _durable_zero_run_receipt(ctx, gap_reasons=gaps)
if prior:
return (
"⚠️ TOOL_ARG_ERROR (verify_and_record): delegation_zero_run already "
"has a durable terminal decision for this actor."
)
if gaps:
return (
"⚠️ TOOL_ERROR (verify_and_record): zero_run_custody_unknown: "
"the host cannot prove whether a prior zero-run receipt exists; "
"no new receipt was written. gaps="
+ json.dumps(sorted(gaps), ensure_ascii=False)
)
blockers = unsettled_start_ids(drive_root, task_id)
if any(blockers.values()):
return (
"⚠️ TOOL_ERROR (verify_and_record): zero_run_requires_settlement: "
"this task has an open run, ambiguous start invocation, or "
"undisposed physical result. blockers="
+ json.dumps(blockers, ensure_ascii=False, sort_keys=True)
)
# A zero-run is lifecycle authority for the actor, so keep it on the
# canonical budget root rather than an ephemeral child drive.
written = append_verification_receipt(
canonical_data_root(ctx), task_id, receipt,
)
except Exception as exc:
return (
"⚠️ TOOL_ERROR (verify_and_record): zero_run_custody_unknown: "
"the host could not atomically prove empty custody and claim the "
f"zero-run decision; no receipt was written ({type(exc).__name__})."
)
if not written:
return (
"⚠️ TOOL_ERROR (verify_and_record): the delegation_zero_run receipt "
"could not be durably written; no zero-run decision was recorded."

View file

@ -120,6 +120,12 @@ route with the same immutable canonical work order; any
`delegate_start` prompt is advisory coordination context, not a reassignment. An API
alternative or genuinely different assignment is a separately visible
`schedule_subagent` child.
The startup receipt and each newly created meaningful wake include a
`coordination_context` with my parent's advisory intent, remaining explicit-deadline
time, honest root-tree spend, active host-visible descendants, and remaining root
acceptance capacity. I use these as planning evidence rather than deterministic
thresholds. Vendor-internal descendants are opaque, and a replayed wake intentionally
shows the stored earlier snapshot until I acknowledge it.
A read-only child cannot write arbitrary local repo/data/memory state, enable tools, commit, review, change
runtime settings, run shell/skills lifecycle tools, or bypass owner resources — but it
@ -163,9 +169,11 @@ evidence-first intermediate check without starting one: publish
`tree_note(kind="review_requested", text=<why>, payload={"evidence_ref": <where>,
"evidence_sha256": <64 hex>})`. Distinct concerns remain visible even when they
reference the same bytes. The parent or root decides whether to inspect, spawn an
ordinary critic, or use its root-owned acceptance path; when it launches a check,
carry the exact hash into the existing accounting/deduplication path. The child never
blocks in a self-review loop.
ordinary critic, or use its root-owned acceptance path. Host-verify the referenced
bytes before they enter the complete root acceptance binding. The child hint itself
starts no reviewer and spends no paid review cycle; the root host atomically claims
that full binding immediately before transport. The child never blocks in a self-review
loop.
Children inside a vendor session remain opaque unless Claudexor emits a
host-visible boundary receipt.

View file

@ -3344,12 +3344,15 @@ def _handle_schedule_task(evt: Dict[str, Any], ctx: Any) -> None:
session_id = str(evt.get("session_id") or "")
actor_id = str(evt.get("actor_id") or "ouroboros")
delegation_role = str(evt.get("delegation_role") or "subagent")
if delegation_role == "subagent":
from supervisor.task_admission import subagent_schedule_owned
if subagent_schedule_owned(ctx, tid):
return
memory_mode = str(evt.get("memory_mode") or "").strip()
drive_root = str(evt.get("drive_root") or "").strip()
child_drive_root = str(evt.get("child_drive_root") or drive_root).strip()
budget_drive_root = str(evt.get("budget_drive_root") or "").strip()
# INTENT ONLY (see `_build_scheduled_task_payload`): the supervisor forwards what
# the parent ASKED for. What the child gets is resolved once, at dispatch.
# Forward parent-requested intent; dispatch resolves it once.
requested_model_lane = str(evt.get("requested_model_lane") or evt.get("model_lane") or "auto").strip() or "auto"
parent_model_lane = str(evt.get("parent_model_lane") or "").strip()
requested_executor = str(evt.get("requested_executor") or "").strip().lower() or "auto"
@ -3657,10 +3660,7 @@ def _handle_schedule_task(evt: Dict[str, Any], ctx: Any) -> None:
)
return
# Admission records an attempted depth, not achieved execution. A task
# can remain queued forever or be cancelled before a worker sees it, so
# ``achieved_depth`` stays unknown until the assignment seam proves a
# host worker actually received the child (see supervisor/workers.py).
# Assignment, not admission, proves achieved depth.
admitted_task_contract, admitted_depth_provenance = stamp_depth_provenance(
task_contract,
attempted_depth=depth,
@ -3792,6 +3792,8 @@ def _handle_schedule_task(evt: Dict[str, Any], ctx: Any) -> None:
)
return
if scheduled_failure_reason:
if scheduled_failure_reason == "scheduled_event_replay":
return
result_fields["delegation_admission"] = {
"status": "rejected",
"reason_code": scheduled_failure_reason,

View file

@ -17,6 +17,29 @@ from ouroboros.task_results import (
log = logging.getLogger(__name__)
def subagent_schedule_owned(
ctx: Any, task_id: str, *, pending_ref: Any = None,
) -> bool:
"""Return whether an exact child id already has queue/lifecycle custody."""
from supervisor import queue
tid = str(task_id or "")
with queue._queue_lock:
pending = pending_ref if isinstance(pending_ref, list) else getattr(
ctx, "PENDING", queue.PENDING,
)
running = getattr(ctx, "RUNNING", queue.RUNNING)
status = str((load_task_result(ctx.DRIVE_ROOT, tid) or {}).get("status") or "")
return (
tid in running
or any(
isinstance(row, dict) and str(row.get("id") or "") == tid
for row in pending
)
or status not in {"", STATUS_REQUESTED}
)
def enqueue_subagent_with_scheduled_result(
ctx: Any,
task: Dict[str, Any],
@ -50,10 +73,19 @@ def enqueue_subagent_with_scheduled_result(
)
with queue._queue_lock:
previous = load_task_result(ctx.DRIVE_ROOT, tid) or {}
if subagent_schedule_owned(ctx, tid, pending_ref=pending_ref):
log.info("Ignoring replayed schedule event for task %s", tid)
return (
task,
"scheduled_event_replay",
"Subagent schedule replay ignored: this task id is already owned by "
"an existing queue or durable lifecycle row.",
False,
)
admitted = ctx.enqueue_task(task)
if isinstance(admitted, dict) and admitted.get("_admission_blocked"):
return admitted, "", "", False
previous = load_task_result(ctx.DRIVE_ROOT, tid) or {}
result_fields["task_contract"] = admitted_task_contract
result_fields["depth_provenance"] = admitted_depth_provenance
result_fields["delegation_admission"] = {
@ -204,4 +236,5 @@ __all__ = [
"enqueue_subagent_with_scheduled_result",
"release_task_admission",
"reserve_task_admission",
"subagent_schedule_owned",
]

View file

@ -270,7 +270,17 @@ def test_blocked_session_bootstrap_wakes_first_turn_with_alternatives(monkeypatc
"subagent_id": "api-scout", "route_kind": "api_model",
"target_id": "google/gemini-3.7-flash", "availability": "check_at_dispatch",
}])
ctx = SimpleNamespace(task_id="child1", drive_root=tmp_path, budget_drive_root=str(tmp_path))
ctx = SimpleNamespace(
task_id="child1",
drive_root=tmp_path,
budget_drive_root=str(tmp_path),
task_metadata={
"root_task_id": "root",
"parent_task_id": "root",
"delegation_role": "subagent",
"budget_drive_root": str(tmp_path),
},
)
dispatch = SimpleNamespace(
blocked=True,
executor_resolution=SimpleNamespace(
@ -281,6 +291,8 @@ def test_blocked_session_bootstrap_wakes_first_turn_with_alternatives(monkeypatc
ctx, {"id": "child1", "configured_subagent": snapshot}, dispatch,
))
assert out["status"] == "configured_session_actor_ready"
assert out["coordination_context"]["parent_intent"]["state"] == "absent"
assert out["coordination_context"]["review_capacity"]["state"] == "unknown"
startup = out["startup"]
assert {key: startup[key] for key in (
"status", "reason", "reset_at", "selected_subagent_id", "alternatives",
@ -1105,6 +1117,11 @@ def test_one_shot_checkpoint_is_reasoned_and_consumed(monkeypatch, tmp_path):
))
wake_id = out.pop("supervision_wake_id")
assert wake_id
coordination_context = out.pop("coordination_context")
assert coordination_context["root_task_id"] == "child1"
assert coordination_context["parent_intent"]["state"] == "absent"
assert coordination_context["time"]["state"] == "not_set"
assert coordination_context["review_capacity"]["state"] == "available"
assert out == {
"status": "inspection_checkpoint",
"run_id": "run-1",
@ -1207,6 +1224,36 @@ def test_replacement_is_refused_before_gateway_or_post(monkeypatch, tmp_path):
assert out["undisposed_patch_run_ids"] == ["run-old"]
def test_replacement_refuses_unreadable_custody_before_fail_soft_scan(
monkeypatch, tmp_path,
):
from ouroboros import delegate_custody as custody
from ouroboros import delegate_recovery
import ouroboros.claudexor_daemon as daemon
import ouroboros.tools.delegate as delegate
from ouroboros.tools.registry import ToolContext
monkeypatch.setattr(custody, "custody_log_unreadable", lambda _root: True)
monkeypatch.setattr(
delegate_recovery,
"unsettled_start_ids",
lambda *_a, **_k: pytest.fail("unreadable custody must stop before scan"),
)
monkeypatch.setattr(
daemon,
"ensure_owned_gateway",
lambda: pytest.fail("unreadable custody must stop before gateway work"),
)
ctx = ToolContext(repo_dir=tmp_path, drive_root=tmp_path)
ctx.task_id = "child-unknown"
snapshot = _snapshot(_settings(_session_row()), "session-builder")
out = json.loads(delegate.exact_start(
ctx, "replacement work", {"snapshot": snapshot},
))
assert out["status"] == "refused"
assert out["reason"] == "replacement_custody_unknown"
def test_terminal_boundary_reaudits_durable_pending_starts(monkeypatch, tmp_path):
from ouroboros import delegate_custody as custody
from ouroboros import delegate_terminal
@ -1521,67 +1568,3 @@ def test_only_approved_restart_causes_reserve_and_abrupt_gap_vetoes(monkeypatch,
monkeypatch.setattr(custody, "reconcile_task_runs", lambda *_a, **_k: [])
assert recovery.pre_adopt_planned_handoffs(tmp_path, []) == set()
assert recovery._read(tmp_path, "child1")["veto_reason"] == "restart_transaction_missing"
def test_planned_restart_kill_keeps_selected_child_even_if_parent_is_interrupted(
monkeypatch, tmp_path,
):
from supervisor import queue as task_queue
from supervisor import workers
repo = tmp_path / "repo"
repo.mkdir()
workers.init(repo, tmp_path, 2, 600, 1800, 100.0)
workers.WORKERS.clear()
workers.PENDING.clear()
workers.RUNNING.clear()
workers.RUNNING.update({
"parent": {"task": {"id": "parent", "root_task_id": "parent"}, "attempt": 1},
"child": {"task": {
"id": "child", "parent_task_id": "parent", "root_task_id": "parent",
}, "attempt": 1},
})
monkeypatch.setattr(workers, "_write_failure_result", lambda *_a, **_k: "cancelled")
monkeypatch.setattr(workers, "_emit_task_done_terminal", lambda *_a, **_k: True)
monkeypatch.setattr(task_queue, "persist_queue_snapshot", lambda *a, **k: True)
workers.kill_workers(
terminal_status="cancelled", preserve_pending=True,
preserve_running_task_ids={"child"},
)
assert [task["id"] for task in workers.PENDING] == ["child"]
assert workers.PENDING[0]["_attempt"] == 2
workers.PENDING.clear()
workers.RUNNING.clear()
def test_api_row_is_refused_by_root_direct_exact_start_before_daemon(monkeypatch, tmp_path):
import ouroboros.tools.delegate as delegate
import ouroboros.claudexor_daemon as daemon
from ouroboros.tools.registry import ToolContext
settings = _settings(_api_row())
monkeypatch.setattr("ouroboros.config.load_settings", lambda: settings)
monkeypatch.setattr(daemon, "ensure_owned_gateway", lambda: (_ for _ in ()).throw(
AssertionError("API rows never POST to Claudexor")))
ctx = ToolContext(repo_dir=tmp_path, drive_root=tmp_path)
ctx.task_id = "root1"
schema = next(e.schema for e in delegate.get_tools() if e.name == "delegate_start")["parameters"]
assert schema["required"] == ["prompt"] and not ({"anyOf", "oneOf", "allOf"} & schema.keys())
missing = json.loads(delegate.exact_start(ctx, "bounded leaf"))
assert missing["reason"] == "subagent_selection_required"
out = json.loads(delegate.exact_start(
ctx, "bounded leaf", {"subagent_id": "api-builder"},
))
assert out["reason"] == "api_actor_requires_schedule_subagent"
retry = json.loads(delegate.exact_start(
ctx, "bounded leaf", {"subagent_id": "api-builder", "retry_of": "inv-old"},
))
assert retry["reason"] == "retry_selector_conflict"
def test_heavy_is_absent_from_active_runtime_and_vision_consumers(monkeypatch):
from ouroboros.llm import LLMClient
from ouroboros.tools.vision import _vision_capable_slot_candidates
monkeypatch.setenv("OUROBOROS_MODEL", "openai/main")
monkeypatch.setenv("OUROBOROS_MODEL_LIGHT", "openai/light")
monkeypatch.setenv("OUROBOROS_MODEL_HEAVY", "openai/legacy-heavy")
assert "openai/legacy-heavy" not in LLMClient().available_models()
assert "openai/legacy-heavy" not in _vision_capable_slot_candidates(LLMClient())

View file

@ -0,0 +1,88 @@
"""Small route/restart edge band for configured recursive actors."""
from __future__ import annotations
import json
def _settings(*rows):
return {
"OUROBOROS_SUBAGENTS": json.dumps({"enabled": True, "items": list(rows)}),
}
def _api_row(row_id="api-builder", model="openai/gpt-5.6-sol", effort="high"):
return {
"subagent_id": row_id,
"name": "API builder",
"recommended_use": "Exact recursive API actor.",
"route": {"kind": "api_model", "target_id": model},
"effort": effort,
}
def test_planned_restart_kill_keeps_selected_child_even_if_parent_is_interrupted(
monkeypatch, tmp_path,
):
from supervisor import queue as task_queue
from supervisor import workers
repo = tmp_path / "repo"
repo.mkdir()
workers.init(repo, tmp_path, 2, 600, 1800, 100.0)
workers.WORKERS.clear()
workers.PENDING.clear()
workers.RUNNING.clear()
workers.RUNNING.update({
"parent": {"task": {"id": "parent", "root_task_id": "parent"}, "attempt": 1},
"child": {"task": {
"id": "child", "parent_task_id": "parent", "root_task_id": "parent",
}, "attempt": 1},
})
monkeypatch.setattr(workers, "_write_failure_result", lambda *_a, **_k: "cancelled")
monkeypatch.setattr(workers, "_emit_task_done_terminal", lambda *_a, **_k: True)
monkeypatch.setattr(task_queue, "persist_queue_snapshot", lambda *a, **k: True)
workers.kill_workers(
terminal_status="cancelled", preserve_pending=True,
preserve_running_task_ids={"child"},
)
assert [task["id"] for task in workers.PENDING] == ["child"]
assert workers.PENDING[0]["_attempt"] == 2
workers.PENDING.clear()
workers.RUNNING.clear()
def test_api_row_is_refused_by_root_direct_exact_start_before_daemon(monkeypatch, tmp_path):
import ouroboros.claudexor_daemon as daemon
import ouroboros.tools.delegate as delegate
from ouroboros.tools.registry import ToolContext
settings = _settings(_api_row())
monkeypatch.setattr("ouroboros.config.load_settings", lambda: settings)
monkeypatch.setattr(daemon, "ensure_owned_gateway", lambda: (_ for _ in ()).throw(
AssertionError("API rows never POST to Claudexor")))
ctx = ToolContext(repo_dir=tmp_path, drive_root=tmp_path)
ctx.task_id = "root1"
schema = next(e.schema for e in delegate.get_tools() if e.name == "delegate_start")["parameters"]
assert schema["required"] == ["prompt"] and not ({"anyOf", "oneOf", "allOf"} & schema.keys())
missing = json.loads(delegate.exact_start(ctx, "bounded leaf"))
assert missing["reason"] == "subagent_selection_required"
out = json.loads(delegate.exact_start(
ctx, "bounded leaf", {"subagent_id": "api-builder"},
))
assert out["reason"] == "api_actor_requires_schedule_subagent"
retry = json.loads(delegate.exact_start(
ctx, "bounded leaf", {"subagent_id": "api-builder", "retry_of": "inv-old"},
))
assert retry["reason"] == "retry_selector_conflict"
def test_heavy_is_absent_from_active_runtime_and_vision_consumers(monkeypatch):
from ouroboros.llm import LLMClient
from ouroboros.tools.vision import _vision_capable_slot_candidates
monkeypatch.setenv("OUROBOROS_MODEL", "openai/main")
monkeypatch.setenv("OUROBOROS_MODEL_LIGHT", "openai/light")
monkeypatch.setenv("OUROBOROS_MODEL_HEAVY", "openai/legacy-heavy")
assert "openai/legacy-heavy" not in LLMClient().available_models()
assert "openai/legacy-heavy" not in _vision_capable_slot_candidates(LLMClient())

View file

@ -118,6 +118,111 @@ def test_physical_terminal_wake_is_not_mislabeled_as_coordination(tmp_path):
assert acknowledged[-1]["coordination"] is False
def test_meaningful_wake_carries_live_tree_planning_facts_and_replays_exactly(
tmp_path, monkeypatch,
):
from datetime import timedelta
from ouroboros import usage_accounting
from ouroboros.contracts.task_contract import build_task_contract
from ouroboros.deadline_utils import utc_now
from ouroboros.review_substrate import review_binding_hash
from ouroboros.task_results import claim_task_acceptance_review_cycle
from ouroboros.utils import atomic_write_json, utc_now_iso
monkeypatch.setenv("OUROBOROS_REVIEW_MAX_CYCLES", "3")
contract = build_task_contract({
"delegation_budget": {
"intent_note": "Prefer one strong critic and preserve time to synthesize.",
},
})
write_task_result(
tmp_path,
"root",
STATUS_RUNNING,
root_task_id="root",
delegation_role="root",
task_contract=contract,
)
_child(tmp_path)
_child(tmp_path, "grandchild", parent_task_id="child")
now = utc_now()
(tmp_path / "state").mkdir(exist_ok=True)
atomic_write_json(tmp_path / "state" / "queue_snapshot.json", {
"ts": utc_now_iso(),
"pending": [],
"running": [
{"id": "child", "task": {
"id": "child", "parent_task_id": "parent",
"root_task_id": "root", "delegation_role": "subagent",
}},
{"id": "grandchild", "task": {
"id": "grandchild", "parent_task_id": "child",
"root_task_id": "root", "delegation_role": "subagent",
}},
],
})
components = {
"candidate_hash": "2" * 64,
"evidence_revision": "3" * 64,
"fence_hash": "4" * 64,
}
binding = {**components, "binding_hash": review_binding_hash(**components)}
assert claim_task_acceptance_review_cycle(
tmp_path,
"root",
binding,
max_cycles=3,
claimed_by_task_id="root",
)["status"] == "claimed"
monkeypatch.setattr(usage_accounting, "usage_breakdown", lambda *_a, **_k: {
"settled_usd": 1.25,
"accounted_usd": 1.75,
"cost_final": False,
"unknown_unmetered": 1,
"integrity_degraded": False,
})
ctx = _ctx(tmp_path)
ctx.task_contract = contract
ctx.task_metadata.update({
"task_contract": contract,
"created_at": (now - timedelta(seconds=30)).isoformat(),
"deadline_at": (now + timedelta(minutes=10)).isoformat(),
})
raw = supervised_wait(
ctx,
"run-live",
wait_once=lambda *_args: json.dumps({
"status": "completed", "run_id": "run-live", "last_seq": 1,
}),
)
wake = json.loads(raw)
facts = wake["coordination_context"]
assert facts["parent_intent"] == {
"state": "present",
"authority": "parent_authored_advisory",
"text": "Prefer one strong critic and preserve time to synthesize.",
}
assert facts["time"]["state"] == "known"
assert 0 < facts["time"]["remaining_sec"] <= 600
assert facts["settled_spend"]["state"] == "partial"
assert facts["settled_spend"]["settled_usd"] == 1.25
assert facts["active_descendants"]["count"] == 2
assert facts["active_descendants"]["vendor_internal"] == "opaque_not_counted"
assert facts["review_capacity"]["claimed_cycles"] == 1
assert facts["review_capacity"]["remaining_cycles"] == 2
replay = supervised_wait(
ctx,
"run-live",
wait_once=lambda *_args: (_ for _ in ()).throw(
AssertionError("pending live facts must replay without recomputation")
),
)
assert replay == raw
def test_child_terminal_before_first_sleep_is_not_lost_as_cursor_baseline(tmp_path):
from ouroboros.artifacts import copy_file_to_task_artifacts
from ouroboros.task_status import load_effective_task_result
@ -172,6 +277,9 @@ def test_oversized_coordination_wake_is_valid_bounded_json_with_exact_source(tmp
from ouroboros.tool_capabilities import tool_result_limit
ctx = _ctx(tmp_path)
ctx.task_contract = {
"delegation_budget": {"intent_note": "preserve-complete-intent:" + "y" * 30_000},
}
for index in range(5):
child_id = f"child-{index}"
_child(tmp_path, task_id=child_id)
@ -193,10 +301,12 @@ def test_oversized_coordination_wake_is_valid_bounded_json_with_exact_source(tmp
assert delivered["supervision_wake_id"]
assert delivered["wake_delivery"]["complete"] is False
assert delivered["wake_delivery"]["wake_events_total"] == 5
assert delivered["coordination_context"]["state"] == "available_in_full_wake_source"
source = delivered["wake_delivery"]["source"]
full = json.loads(read_actor_source_bytes(tmp_path, "parent", source))
assert len(full["wake_events"]) == 5
assert all(len(item["beacon"]["text"]) > 3800 for item in full["wake_events"])
assert len(full["coordination_context"]["parent_intent"]["text"]) > 30_000
assert acknowledge_pending_wake(ctx, raw)
after = json.loads(supervised_wait(
@ -212,6 +322,43 @@ def test_oversized_coordination_wake_is_valid_bounded_json_with_exact_source(tmp
)
def test_live_descendant_fact_rejects_stale_queue_and_ignores_unrelated_corruption(
tmp_path,
):
from ouroboros.delegate_supervision import coordination_live_context
from ouroboros.utils import atomic_write_json, utc_now_iso
ctx = _ctx(tmp_path)
_child(tmp_path)
state = tmp_path / "state"
state.mkdir(exist_ok=True)
snapshot = {
"ts": "2000-01-01T00:00:00Z",
"pending": [],
"running": [{"id": "child", "task": {
"id": "child", "parent_task_id": "parent",
"root_task_id": "root", "delegation_role": "subagent",
}}],
}
atomic_write_json(state / "queue_snapshot.json", snapshot)
stale = coordination_live_context(ctx)["active_descendants"]
assert stale["state"] == "unknown"
assert stale["count"] is None
unrelated = tmp_path / "task_results" / "unrelated-history.json"
unrelated.write_text("{broken", encoding="utf-8")
snapshot["ts"] = utc_now_iso()
atomic_write_json(state / "queue_snapshot.json", snapshot)
current = coordination_live_context(ctx)["active_descendants"]
assert current == {
"state": "known",
"count": 1,
"by_status": {"running": 1},
"scope": "host_visible_descendants",
"vendor_internal": "opaque_not_counted",
}
def test_delegate_wait_entry_never_acks_an_undelivered_pending_wake(tmp_path, monkeypatch):
from ouroboros.tools import delegate as delegate_module

View file

@ -728,7 +728,7 @@ def test_task_acceptance_required_feeds_back_capsule(monkeypatch, tmp_path):
messages2 = [{"role": "system", "content": ""}, {"role": "user", "content": "goal"}]
tools2 = SimpleNamespace(_ctx=ctx2)
result2 = _run_task_acceptance_review_once(
tools=tools2, content="done", task_id="t", task_type="task",
tools=tools2, content="done", task_id="t-blocked", task_type="task",
llm_trace=trace2, drive_root=None, messages=messages2, emit_progress=lambda _m: None,
)
assert result2 is True # capsule -> one bounded re-loop
@ -746,7 +746,7 @@ def test_task_acceptance_required_feeds_back_capsule(monkeypatch, tmp_path):
monkeypatch.setattr(rs, "run_review_request", lambda *a, **k: solved)
replacement = _run_task_acceptance_review_once(
tools=tools2, content="revised", task_id="t", task_type="task",
tools=tools2, content="revised", task_id="t-blocked", task_type="task",
llm_trace=trace2, drive_root=None, messages=messages2, emit_progress=lambda _m: None,
)
assert replacement is False
@ -760,7 +760,7 @@ def test_task_acceptance_required_feeds_back_capsule(monkeypatch, tmp_path):
trace_ok = {"tool_calls": [{"tool": "write_file", "args": {"path": "x.py"}}]}
messages_ok = [{"role": "system", "content": ""}, {"role": "user", "content": "goal"}]
result_ok = _run_task_acceptance_review_once(
tools=tools2, content="revised", task_id="t", task_type="task",
tools=tools2, content="revised", task_id="t-ok", task_type="task",
llm_trace=trace_ok, drive_root=None, messages=messages_ok, emit_progress=lambda _m: None,
)
assert result_ok is False
@ -775,7 +775,7 @@ def test_task_acceptance_required_feeds_back_capsule(monkeypatch, tmp_path):
result3 = _run_task_acceptance_review_once(
# A changed candidate creates a fresh binding; an unchanged candidate
# must reuse the already-paid host panel under the v6.65 contract.
tools=tools2, content="revised again", task_id="t", task_type="task",
tools=tools2, content="revised again", task_id="t-blocked-alt", task_type="task",
llm_trace=trace3, drive_root=None, messages=messages3, emit_progress=lambda _m: None,
)
assert result3 is False # capsule already spent -> finalize

View file

@ -7,11 +7,12 @@ import json
from pathlib import Path
from types import SimpleNamespace
import pytest
from ouroboros import task_tree_ledger
from ouroboros.contracts.task_constraint import normalize_task_constraint
from ouroboros.contracts.task_contract import build_task_contract
from ouroboros.headless import copy_child_task_result
from ouroboros.depth_evidence import build_depth_summary
from ouroboros.outcome_receipt_store import merge_verification_receipts
from ouroboros.outcomes import latest_unreconciled_failed_receipt
from ouroboros.loop import (
@ -37,6 +38,248 @@ from supervisor import (
)
def _acceptance_binding(seed: str) -> dict[str, str]:
from ouroboros.review_substrate import review_binding_hash
components = {
"candidate_hash": chr(ord(seed) + 1) * 64,
"evidence_revision": chr(ord(seed) + 2) * 64,
"fence_hash": chr(ord(seed) + 3) * 64,
}
return {**components, "binding_hash": review_binding_hash(**components)}
def test_root_acceptance_review_claims_share_one_atomic_exact_binding_wallet(tmp_path):
from concurrent.futures import ThreadPoolExecutor
from ouroboros.task_results import (
claim_task_acceptance_review_cycle,
load_task_acceptance_review_state,
)
write_task_result(tmp_path, "root-wallet", STATUS_RUNNING, root_task_id="root-wallet")
first_binding = _acceptance_binding("1")
with ThreadPoolExecutor(max_workers=2) as pool:
outcomes = list(pool.map(
lambda _index: claim_task_acceptance_review_cycle(
tmp_path,
"root-wallet",
first_binding,
max_cycles=2,
claimed_by_task_id="root-wallet",
),
range(2),
))
assert sorted(row["status"] for row in outcomes) == ["claimed", "unknown"]
state = load_task_acceptance_review_state(tmp_path, "root-wallet")
assert list(state["claims_by_binding"]) == [first_binding["binding_hash"]]
result_path = tmp_path / "task_results" / "root-wallet.json"
before_duplicate = result_path.read_bytes()
duplicate = claim_task_acceptance_review_cycle(
tmp_path,
"root-wallet",
first_binding,
max_cycles=2,
claimed_by_task_id="root-wallet",
)
assert duplicate["status"] == "unknown"
assert duplicate["reason"] == "binding_dispatch_already_claimed"
assert result_path.read_bytes() == before_duplicate
second = claim_task_acceptance_review_cycle(
tmp_path,
"root-wallet",
_acceptance_binding("5"),
max_cycles=2,
claimed_by_task_id="root-wallet",
)
before_exhausted = result_path.read_bytes()
exhausted = claim_task_acceptance_review_cycle(
tmp_path,
"root-wallet",
_acceptance_binding("a"),
max_cycles=2,
claimed_by_task_id="root-wallet",
)
assert second["status"] == "claimed"
assert second["cycles_paid"] == 2
assert result_path.read_bytes() == before_exhausted
assert exhausted == {
"status": "unavailable",
"reason": "review_cycles_exhausted",
"binding_hash": _acceptance_binding("a")["binding_hash"],
"cycles_paid": 2,
"max_cycles": 2,
"remaining_cycles": 0,
}
def test_acceptance_review_wallet_rejects_present_empty_or_tampered_authority(tmp_path):
from ouroboros.task_results import (
TASK_ACCEPTANCE_REVIEW_STATE_KEY,
claim_task_acceptance_review_cycle,
load_task_acceptance_review_state,
)
path = tmp_path / "task_results" / "root-invalid.json"
write_task_result(
tmp_path,
"root-invalid",
STATUS_RUNNING,
root_task_id="root-invalid",
**{TASK_ACCEPTANCE_REVIEW_STATE_KEY: {}},
)
before = path.read_bytes()
with pytest.raises(ValueError, match="TASK_ACCEPTANCE_REVIEW_STATE_INVALID"):
load_task_acceptance_review_state(tmp_path, "root-invalid")
with pytest.raises(ValueError, match="TASK_ACCEPTANCE_REVIEW_STATE_INVALID"):
claim_task_acceptance_review_cycle(
tmp_path,
"root-invalid",
_acceptance_binding("1"),
max_cycles=2,
claimed_by_task_id="root-invalid",
)
assert path.read_bytes() == before
write_task_result(
tmp_path,
"root-tampered",
STATUS_RUNNING,
root_task_id="root-tampered",
)
tampered = _acceptance_binding("5")
tampered["binding_hash"] = "f" * 64
tampered_path = tmp_path / "task_results" / "root-tampered.json"
before = tampered_path.read_bytes()
with pytest.raises(ValueError, match="binding digest mismatch"):
claim_task_acceptance_review_cycle(
tmp_path,
"root-tampered",
tampered,
max_cycles=2,
claimed_by_task_id="root-tampered",
)
assert tampered_path.read_bytes() == before
def test_acceptance_review_wallet_cap_and_root_initialization_are_atomic(
tmp_path, monkeypatch,
):
from concurrent.futures import ThreadPoolExecutor
import ouroboros.task_results as task_results
with pytest.raises(ValueError, match="TASK_ACCEPTANCE_REVIEW_STATE_UNKNOWN"):
task_results.claim_task_acceptance_review_cycle(
tmp_path,
"missing-root",
_acceptance_binding("1"),
max_cycles=1,
claimed_by_task_id="child",
allow_create=False,
)
assert not (tmp_path / "task_results" / "missing-root.json").exists()
claimed = task_results.claim_task_acceptance_review_cycle(
tmp_path,
"new-root",
_acceptance_binding("1"),
max_cycles=1,
claimed_by_task_id="new-root",
allow_create=True,
)
assert claimed["status"] == "claimed"
assert load_task_result(tmp_path, "new-root")["status"] == STATUS_RUNNING
write_task_result(tmp_path, "cap-root", STATUS_RUNNING, root_task_id="cap-root")
with ThreadPoolExecutor(max_workers=2) as pool:
outcomes = list(pool.map(
lambda binding: task_results.claim_task_acceptance_review_cycle(
tmp_path,
"cap-root",
binding,
max_cycles=1,
claimed_by_task_id="cap-root",
),
(_acceptance_binding("1"), _acceptance_binding("5")),
))
assert sorted(row["status"] for row in outcomes) == ["claimed", "unavailable"]
write_task_result(tmp_path, "racy-root", STATUS_RUNNING, root_task_id="racy-root")
original_update = task_results.update_json_locked
def remove_before_locked_read(path, mutator, **kwargs):
path.unlink()
return original_update(path, mutator, **kwargs)
monkeypatch.setattr(task_results, "update_json_locked", remove_before_locked_read)
with pytest.raises(ValueError, match="TASK_ACCEPTANCE_REVIEW_STATE_UNKNOWN"):
task_results.claim_task_acceptance_review_cycle(
tmp_path,
"racy-root",
_acceptance_binding("a"),
max_cycles=1,
claimed_by_task_id="child",
allow_create=False,
)
assert not (tmp_path / "task_results" / "racy-root.json").exists()
def test_descendant_cannot_initialize_missing_root_review_authority(tmp_path, monkeypatch):
from ouroboros.task_pacing import project_task_acceptance_review_capacity
ctx = SimpleNamespace(
task_id="child",
drive_root=tmp_path,
budget_drive_root=str(tmp_path),
task_metadata={
"root_task_id": "missing-root",
"parent_task_id": "missing-root",
"delegation_role": "subagent",
"budget_drive_root": str(tmp_path),
},
)
projection = project_task_acceptance_review_capacity(ctx)
assert projection["state"] == "unknown"
assert projection["claimed_cycles"] is None
assert not (tmp_path / "task_results" / "missing-root.json").exists()
def test_root_review_capacity_uses_existing_explicit_pass_semantics(tmp_path, monkeypatch):
from ouroboros.task_pacing import project_task_acceptance_review_capacity
monkeypatch.setenv("OUROBOROS_REVIEW_MAX_CYCLES", "2")
contract = build_task_contract({
"budget_profile": {"max_improvement_passes": 4},
})
write_task_result(
tmp_path,
"root-cap",
STATUS_RUNNING,
root_task_id="root-cap",
delegation_role="root",
task_contract=contract,
)
ctx = SimpleNamespace(
task_id="root-cap",
drive_root=tmp_path,
budget_drive_root=str(tmp_path),
task_contract=contract,
task_metadata={
"root_task_id": "root-cap",
"delegation_role": "root",
"budget_drive_root": str(tmp_path),
"task_contract": contract,
},
)
projection = project_task_acceptance_review_capacity(ctx)
assert projection["state"] == "available"
assert projection["cap_cycles"] == 5
assert projection["claimed_cycles"] == 0
assert projection["remaining_cycles"] == 5
def test_depth3_control_plane_reaches_root_acceptance(tmp_path, monkeypatch):
repo = tmp_path / "repo"
repo.mkdir()
@ -716,84 +959,6 @@ def test_real_over_cap_refusal_reaches_root_acceptance_depth_summary(tmp_path, m
}
def test_depth_summary_reports_lower_cap_as_typed_reduction(monkeypatch):
# Live Settings may change after admission; persisted child provenance wins.
monkeypatch.setenv("OUROBOROS_MAX_SUBAGENT_DEPTH", "7")
root_contract = build_task_contract({"delegation_budget": {"depth_remaining": 3}})
statuses = [
{
"task_id": f"child-{depth}",
"depth_provenance": {
"requested_depth": 3,
"permitted_depth": 2,
"attempted_depth": depth,
"achieved_depth": depth,
},
}
for depth in (1, 2)
]
assert build_depth_summary(root_contract, statuses) == {
"requested_depth": 3,
"permitted_depth": 2,
"attempted_depth": 2,
"achieved_depth": 2,
"status": "capability_reduced",
"host_visible_only": True,
}
def test_depth_summary_is_order_independent_and_allows_chosen_shallower():
root_contract = build_task_contract({
"delegation_budget": {
"depth_remaining": 3,
"depth_provenance": {
"requested_depth": 3,
"permitted_depth": 3,
"attempted_depth": 0,
"achieved_depth": None,
},
},
})
mixed = [
{
"depth_provenance": {
"requested_depth": 3, "permitted_depth": 3,
"attempted_depth": 1, "achieved_depth": 1,
},
},
{
"depth_provenance": {
"requested_depth": 3, "permitted_depth": 2,
"attempted_depth": 2, "achieved_depth": 2,
},
},
]
expected = {
"requested_depth": 3, "permitted_depth": 2,
"attempted_depth": 2, "achieved_depth": 2,
"status": "capability_reduced", "host_visible_only": True,
}
assert build_depth_summary(root_contract, mixed) == expected
assert build_depth_summary(root_contract, reversed(mixed)) == expected
assert build_depth_summary(root_contract, [mixed[0]]) == {
"requested_depth": 3, "permitted_depth": 3,
"attempted_depth": 1, "achieved_depth": 1,
"status": "chosen_shallower", "host_visible_only": True,
}
def test_depth_summary_never_recomputes_missing_history_from_live_settings(monkeypatch):
root_contract = build_task_contract({"delegation_budget": {"depth_remaining": 3}})
monkeypatch.setenv("OUROBOROS_MAX_SUBAGENT_DEPTH", "7")
assert build_depth_summary(root_contract, []) == {
"requested_depth": 3, "permitted_depth": None,
"attempted_depth": 0, "achieved_depth": 0,
"status": "evidence_unknown", "host_visible_only": True,
}
def test_split_root_receipts_reconcile_in_host_timestamp_order():
old_pass = {
"criterion_id": "claim_1", "status": "pass",

View file

@ -4,6 +4,7 @@ import json
from types import SimpleNamespace
from ouroboros.contracts.task_contract import build_task_contract
from ouroboros.depth_evidence import build_depth_summary
from ouroboros.task_results import STATUS_RUNNING, write_task_result
from ouroboros.tools.control_delegation import (
check_delegation_admission,
@ -45,6 +46,70 @@ def test_explicit_rights_are_typed_and_legacy_omission_stays_permissive():
assert narrowed["may_fan_out"] is False
def test_depth_summary_reports_lower_cap_as_typed_reduction(monkeypatch):
monkeypatch.setenv("OUROBOROS_MAX_SUBAGENT_DEPTH", "7")
root_contract = build_task_contract({"delegation_budget": {"depth_remaining": 3}})
statuses = [
{
"task_id": f"child-{depth}",
"depth_provenance": {
"requested_depth": 3, "permitted_depth": 2,
"attempted_depth": depth, "achieved_depth": depth,
},
}
for depth in (1, 2)
]
assert build_depth_summary(root_contract, statuses) == {
"requested_depth": 3, "permitted_depth": 2,
"attempted_depth": 2, "achieved_depth": 2,
"status": "capability_reduced", "host_visible_only": True,
}
def test_depth_summary_is_order_independent_and_allows_chosen_shallower():
root_contract = build_task_contract({
"delegation_budget": {
"depth_remaining": 3,
"depth_provenance": {
"requested_depth": 3, "permitted_depth": 3,
"attempted_depth": 0, "achieved_depth": None,
},
},
})
mixed = [
{"depth_provenance": {
"requested_depth": 3, "permitted_depth": 3,
"attempted_depth": 1, "achieved_depth": 1,
}},
{"depth_provenance": {
"requested_depth": 3, "permitted_depth": 2,
"attempted_depth": 2, "achieved_depth": 2,
}},
]
expected = {
"requested_depth": 3, "permitted_depth": 2,
"attempted_depth": 2, "achieved_depth": 2,
"status": "capability_reduced", "host_visible_only": True,
}
assert build_depth_summary(root_contract, mixed) == expected
assert build_depth_summary(root_contract, reversed(mixed)) == expected
assert build_depth_summary(root_contract, [mixed[0]]) == {
"requested_depth": 3, "permitted_depth": 3,
"attempted_depth": 1, "achieved_depth": 1,
"status": "chosen_shallower", "host_visible_only": True,
}
def test_depth_summary_never_recomputes_missing_history_from_live_settings(monkeypatch):
root_contract = build_task_contract({"delegation_budget": {"depth_remaining": 3}})
monkeypatch.setenv("OUROBOROS_MAX_SUBAGENT_DEPTH", "7")
assert build_depth_summary(root_contract, []) == {
"requested_depth": 3, "permitted_depth": None,
"attempted_depth": 0, "achieved_depth": 0,
"status": "evidence_unknown", "host_visible_only": True,
}
def test_depth_provenance_follows_explicit_request_through_three_levels():
root = build_task_contract({"delegation_budget": {"depth_remaining": 3}})
depth_one = child_budget_for_schedule(
@ -417,6 +482,73 @@ def test_supervisor_receipt_rollback_removes_only_its_enqueue_identity(
}
def test_replayed_schedule_event_keeps_one_physical_task_and_transition(
tmp_path, monkeypatch,
):
from supervisor import events, queue, state, workers
monkeypatch.setattr(events, "_find_duplicate_task", lambda *args, **kwargs: None)
write_task_result(
tmp_path, "parent", STATUS_RUNNING,
root_task_id="parent", delegation_role="root",
task_contract=build_task_contract({"delegation_budget": {"may_fan_out": True}}),
)
write_task_result(
tmp_path, "same-child", "requested",
parent_task_id="parent", root_task_id="parent",
delegation_role="subagent", result="Awaiting supervisor acceptance.",
)
ctx = _fake_ctx(tmp_path, [])
def enqueue_task(task):
admitted = dict(task)
ctx.PENDING.append(admitted)
return admitted
ctx.enqueue_task = enqueue_task
event = _schedule_event("same-child", "parent", drive_root=tmp_path)
events._handle_schedule_task(event, ctx)
first = json.loads(
(tmp_path / "task_results" / "same-child.json").read_text(encoding="utf-8")
)
transition_id = first["delegation_admission"]["transition_id"]
def unexpected_constraint_resolution(*_args, **_kwargs):
raise AssertionError("a replay must stop before workspace provisioning")
monkeypatch.setattr(events, "_resolve_subagent_constraint", unexpected_constraint_resolution)
events._handle_schedule_task(event, ctx)
replayed = json.loads(
(tmp_path / "task_results" / "same-child.json").read_text(encoding="utf-8")
)
assert [task["id"] for task in ctx.PENDING] == ["same-child"]
assert replayed["delegation_admission"]["transition_id"] == transition_id
delivered = []
class FakeWorkerQueue:
def put(self, task):
delivered.append(dict(task))
worker_map = {
wid: SimpleNamespace(wid=wid, busy_task_id=None, in_q=FakeWorkerQueue())
for wid in (1, 2)
}
monkeypatch.setattr(workers, "DRIVE_ROOT", tmp_path)
monkeypatch.setattr(workers, "PENDING", ctx.PENDING)
monkeypatch.setattr(workers, "RUNNING", ctx.RUNNING)
monkeypatch.setattr(workers, "WORKERS", worker_map)
monkeypatch.setattr(workers, "load_state", lambda: {})
monkeypatch.setattr(state, "budget_remaining", lambda *_args, **_kwargs: 100.0)
monkeypatch.setattr(queue, "persist_queue_snapshot", lambda reason="": None)
workers.assign_tasks()
assert [task["id"] for task in delivered] == ["same-child"]
assert sum(worker.busy_task_id == "same-child" for worker in worker_map.values()) == 1
assert ctx.RUNNING["same-child"]["worker_id"] in worker_map
def test_supervisor_keeps_admission_when_scheduled_write_raises_after_commit(
tmp_path, monkeypatch,
):

View file

@ -1091,6 +1091,104 @@ def test_zero_run_refuses_ambiguous_physical_start_custody(tmp_path):
assert ctx._configured_actor_bootstrap.get("zero_run_receipt_recorded") is not True
def test_zero_run_refuses_unreadable_custody_without_fail_soft_scan(
tmp_path, monkeypatch,
):
from ouroboros import delegate_custody as custody
from ouroboros import delegate_recovery
from ouroboros.outcomes import read_verification_receipts
from ouroboros.tools.verify import _verify_and_record
ctx = _verify_ctx(tmp_path)
ctx._configured_actor_bootstrap = {
"route_id": "session-a",
"work_order_fingerprint": "a" * 64,
"physical_started": False,
}
monkeypatch.setattr(custody, "custody_log_unreadable", lambda _root: True)
monkeypatch.setattr(
delegate_recovery,
"unsettled_start_ids",
lambda *_a, **_k: pytest.fail("an unreadable authority must stop before scan"),
)
refused = _verify_and_record(
ctx,
contract_kind="delegation_zero_run",
zero_run_decision="complete",
zero_run_basis="no visible run",
)
assert "zero_run_custody_unknown" in refused
assert "custody_log_unreadable" in refused
assert read_verification_receipts(tmp_path, ctx.task_id) == []
assert ctx._configured_actor_bootstrap.get("zero_run_receipt_recorded") is not True
def test_zero_run_and_fresh_start_share_one_atomic_actor_decision(tmp_path):
from concurrent.futures import ThreadPoolExecutor
import threading
from ouroboros import delegate_custody as custody
from ouroboros.outcomes import read_verification_receipts
from ouroboros.tools.delegate_integration import claimed_start_request
from ouroboros.tools.verify import _verify_and_record
ctx = _verify_ctx(tmp_path, task_id="actor-claim")
ctx._configured_actor_bootstrap = {
"route_id": "session-a",
"work_order_fingerprint": "a" * 64,
"physical_started": False,
}
drive = custody.custody_root(ctx)
barrier = threading.Barrier(2)
def claim_start():
barrier.wait()
return claimed_start_request(
drive,
claim_target="",
actor_ctx=ctx,
enforce_actor_idle=True,
run_id="",
task_id=ctx.task_id,
idempotency_key="actor-claim-invocation",
invocation_id="actor-claim-invocation",
max_seconds=30,
request={"prompt": "exact physical assignment"},
project_id="project-1",
project_owned=False,
route="codex",
)
def claim_zero_run():
barrier.wait()
return _verify_and_record(
ctx,
contract_kind="delegation_zero_run",
zero_run_decision="complete",
zero_run_basis="host-visible work completed without a physical leaf",
)
with ThreadPoolExecutor(max_workers=2) as pool:
start_future = pool.submit(claim_start)
zero_future = pool.submit(claim_zero_run)
start = start_future.result()
zero = zero_future.result()
start_won = start[0] is True
zero_won = "typed host receipt recorded" in zero
assert start_won is not zero_won
receipts = read_verification_receipts(drive, ctx.task_id)
pending = custody.pending_invocations(drive)
if start_won:
assert "zero_run_requires_settlement" in zero
assert receipts == []
assert [row["invocation_id"] for row in pending] == ["actor-claim-invocation"]
else:
assert start[1]["reason"] == "zero_run_already_recorded"
assert len(receipts) == 1 and receipts[0]["zero_run"] is True
assert pending == []
def test_failed_direct_child_does_not_make_actor_first_terminal_clean(tmp_path):
from ouroboros.subagent_bootstrap import actor_first_unresolved_fact
from ouroboros.task_results import STATUS_FAILED, write_task_result

View file

@ -78,16 +78,31 @@ def test_acceptance_panel_persists_timing_to_canonical_root(tmp_path, monkeypatc
import ouroboros.loop as loop
import ouroboros.review_evidence as evidence_mod
import ouroboros.review_substrate as substrate
from ouroboros.task_results import STATUS_RUNNING, write_task_result
from ouroboros.tools import review_helpers
canonical = tmp_path / "canonical"
child = tmp_path / "child"
tool_ctx = SimpleNamespace(
task_id="root-timing",
drive_root=child,
budget_drive_root=str(canonical),
task_metadata={"budget_drive_root": str(canonical)},
task_metadata={
"root_task_id": "root-timing",
"delegation_role": "root",
"budget_drive_root": str(canonical),
},
)
write_task_result(
canonical, "root-timing", STATUS_RUNNING,
root_task_id="root-timing", delegation_role="root",
)
monkeypatch.setattr(evidence_mod, "build_task_acceptance_evidence", lambda *_a, **_k: {})
monkeypatch.setattr(substrate, "reviewer_slots", lambda **_k: [])
monkeypatch.setattr(
substrate, "reviewer_slots",
lambda **_k: [SimpleNamespace(model="test-reviewer")],
)
monkeypatch.setattr(review_helpers, "review_wave_budget_gate", lambda *_a, **_k: None)
monkeypatch.setattr(
substrate,
"run_review_request",
@ -106,6 +121,9 @@ def test_acceptance_panel_persists_timing_to_canonical_root(tmp_path, monkeypatc
subtree_statuses=[],
budget_profile={},
passes_done=2,
review_binding=substrate.build_review_binding(
candidate="deliverable", evidence={}, fence_token_or_state="timing-test",
),
)
loop._execute_task_acceptance_panel(ctx)
@ -585,4 +603,3 @@ def test_clean_acceptance_requires_per_criterion_evidence(tmp_path):
)
assert clean.aggregate_signal == "PASS"