mirror of
https://github.com/razzant/ouroboros.git
synced 2026-10-02 19:58:46 +00:00
Merge branch 'ou2-waiting-abc' into claude/steer-20260926
# Conflicts: # tests/test_reference_book_budgets.py
This commit is contained in:
commit
5be5231629
16 changed files with 677 additions and 72 deletions
|
|
@ -303,7 +303,7 @@ A definite refusal needs a typed cause, no custody handle and either producer `d
|
|||
|
||||
**Work orders.** The compiler (`subagent_work_order.py`) sends the entire chosen assignment, preserving context and instruction roles without an arbitrary host-size cutoff: direct starts carry the normalized host contract once in `instructions` (`delegate_start_instructions.py`; the separately hashed coordination appendix is absent from the host pre-start) and the chosen assignment separately in `prompt`, and coordination context stays complete. Real native/HTTP limits return their actual failure with the original input and any pending invocation retained. Exact-source readers stay optional capabilities; incomplete source coverage never authorizes a terminal PASS or apply (`delegate_source_coverage.py`). Legacy partial starts keep their exact renderer digest, source-interval validation and stored-body retry; removing partial-start production certifies no incomplete old run and launches no duplicate after ambiguous dispatch.
|
||||
|
||||
**Supervision.** `delegate_wait` is model-visible as an event-only sleep, not a caller-sized poll: `delegate_supervision.supervised_wait` renews bounded transport windows in host code at zero LLM calls, and only a meaningful event (settlement, interaction, fault, addressed message, child signal, control, recovery judgment) becomes a coalesced durable wake, replayed across worker interruption until acknowledged; deadline, ceiling, budget and cancellation stay outer bounds. Every receipt and wake carries one host-rendered `coordination_context` (intent, time remaining, root-tree spend, active descendants, remaining paid acceptance capacity) (facts for LLM judgment, never thresholds), observed READ-ONLY, so a metadata-poor task reports `time.state = "not_set"` instead of latching an anchor from a poll. Polling writes nothing of its own beyond the canonical usage-ledger reader's bounded maintenance, identical for every reader: the torn-tail quarantine after a SINGLE crash mid-append (one verbatim row in `state/usage_attempts.quarantine.jsonl`, one `usage_ledger_tail_quarantined` event), the empty `state/` lock directory on a never-initialized root, and owner-aware `usage_attempts.lock` recovery (§1 Platform substrate); an absent ledger is that reader's known-zero; a crash inside the quarantine repair can leave the sink torn (a disclosed residual). A requested future inspection (`checkpoint_after_sec`) wakes once and is consumed by any earlier real event: no cadence, stall classifier or hidden polling. An observation the transport could not complete is a quiet renewal carrying its typed reason, never a refusal that spends a model round: `observation_read_timeout` is our own read bound expiring against a live daemon (quiet, nothing more), `daemon_unreachable` a socket that carried no answer, the only half worth an owner line, delivered once per episode with one recovery line on the supervising task's own progress surface. The beat stays three seconds with no backoff, durable counter or outage latch, and deadline, ceiling, budget and cancellation still cut a long unobserved stretch. On a wake the nanny holds its full tool surface and the parent's captured model/effort. Nanny economics are structurally quiet: only `BASELINE_RESET_TOOLS` (`delegate_start`/`schedule_subagent`) reset the burn baseline, and coordination verbs never buy metered silence (`nanny_pacing.py`).
|
||||
**Supervision.** `delegate_wait` is model-visible as an event-only sleep, not a caller-sized poll: `delegate_supervision.supervised_wait` renews bounded transport windows in host code at zero LLM calls, and only a meaningful event (settlement, interaction, fault, addressed message, child signal, control, recovery judgment) becomes a coalesced durable wake, replayed across worker interruption until acknowledged; deadline, ceiling, budget and cancellation stay outer bounds. Every receipt and wake carries one host-rendered `coordination_context` (intent, time remaining, root-tree spend, active descendants, remaining paid acceptance capacity) (facts for LLM judgment, never thresholds), observed READ-ONLY, so a metadata-poor task reports `time.state = "not_set"` instead of latching an anchor from a poll. Polling writes nothing of its own beyond the canonical usage-ledger reader's bounded maintenance, identical for every reader: the torn-tail quarantine after a SINGLE crash mid-append (one verbatim row in `state/usage_attempts.quarantine.jsonl`, one `usage_ledger_tail_quarantined` event), the empty `state/` lock directory on a never-initialized root, and owner-aware `usage_attempts.lock` recovery (§1 Platform substrate); an absent ledger is that reader's known-zero; a crash inside the quarantine repair can leave the sink torn (a disclosed residual). A requested future inspection (`checkpoint_after_sec`) wakes once and is consumed by any earlier real event: no cadence, stall classifier or hidden polling. An observation the transport could not complete is a quiet renewal carrying its typed reason, never a refusal that spends a model round: `observation_read_timeout` is our own read bound expiring against a live daemon (quiet, nothing more), `daemon_unreachable` a socket that carried no answer, the only half worth an owner line, delivered once per episode with one recovery line on the supervising task's own progress surface. The beat stays three seconds with no backoff, durable counter or outage latch, and deadline, ceiling, budget and cancellation still cut a long unobserved stretch. On a wake the nanny holds its full tool surface and the parent's captured model/effort. The child-delivery cursor is task-scoped: its acknowledged part survives a new run id and the recovery handoff/restore, so a delivered child terminal or beacon is never re-announced, while an unacknowledged one re-emits once. Every wake, terminal included, carries host-measured `sleep` facts over the whole call (entry, seconds slept, quiet renewals, journal advances; each call opens a new sleep, labelled after an adoption or an interrupted call), stamped at the one durable publication point so a replay repeats them; a quiet-status wake drops the last tick's `waited_sec`/`quiet_for_sec` and keep-watching/cancel note but keeps the PAUSED note, and `cache_horizon_note` is computed once per wake from the last recorded model response (a lower bound on cache age; unknown tiers stay silent). Each wake also carries `leaf_live_input`, the route's declared live-input capability read once at entry (`unknown` when unread; a capability, not a liveness verdict), and the state records dated observation facts (`last_answered_observation_at`, `observation_failure`), never file mtime. The shared wait window and compact child projection apply to every `wait_task` caller (Presence, review executors, owner_wait and budget_pause alike); `budget_pause` is a dispatch fence with an owner-granted resume, never a sleep. Nanny economics are structurally quiet: only `BASELINE_RESET_TOOLS` (`delegate_start`/`schedule_subagent`) reset the burn baseline, and coordination verbs never buy metered silence (`nanny_pacing.py`); the reminder's money sentence follows the rounds' own cost evidence (priced = metered, no provider price = unknown cash, never zero) and names a fallback only when a chain is configured.
|
||||
|
||||
**A run's question is the nanny's to answer.** Supervision wakes immediately with typed `status="waiting_on_user"` on a NEW `pendingInteractions` entry instead of burning the engine's answer timeout in dead polling; answer keys are echoed verbatim into `delegate_answer` (custody-gated like cancel, relaying the engine's typed outcomes: `subscription_window_exhausted` carries `reset_at`; a transport death or 5xx is `delivery_unknown`), and delivered interaction ids are acknowledged only after transcript injection, so a question neither re-triggers a round nor disappears across recovery; every waiting payload states `continuation: same_session` (an answer resumes THIS session; each turn paid). A question above the nanny's authority escalates to the nearest live ancestor. A route without a mid-run question channel ends needing input instead: a terminal with `outcome_facts.reason=input_required` (`continuation: new_physical_run`) is answered by a plain new `delegate_start(subagent_id=..., prompt=...)`, never the engine's rerun verb, which would start a run outside this task's custody trail.
|
||||
|
||||
|
|
|
|||
|
|
@ -820,6 +820,8 @@ and what enforces each.
|
|||
an addressed task/owner message, a direct-child signal, control/recovery judgment or a
|
||||
model-requested one-shot checkpoint wakes it. No caller-visible `wait_sec`, repeating
|
||||
timers, progress wakes or host semantic stall detector.
|
||||
- Wake facts are measured over the interval the actor experienced (whole-call `sleep` stamped at
|
||||
the one publication point, never a tick's `waited_sec`); the acked child cursor is task-scoped.
|
||||
- Wait/continue/stop is a structured fact — terminal status plus heartbeat freshness
|
||||
from `queue_snapshot.json` via `task_status.py` — never a keyword or regex over
|
||||
content (BIBLE P5). Fixed kill-timeouts (hard task/tool ceilings, watchdog) stay the
|
||||
|
|
|
|||
|
|
@ -624,6 +624,19 @@ def emit(ctx: Any, run_id: str, advance: _Advance, *,
|
|||
log.debug("delegated progress emit failed", exc_info=True)
|
||||
|
||||
|
||||
def paused_note(pending_interactions: Optional[List[Dict[str, Any]]]) -> str:
|
||||
"""The note of a run PAUSED on its own question(s): never "stuck" (owner 7=A /
|
||||
F13). Shared by the per-window payload and the supervising wake, which keeps it
|
||||
when it drops the per-tick keep-watching/cancel note."""
|
||||
return ("The run is alive and PAUSED on the question(s) it already asked "
|
||||
"(waiting_on_user; see pending_interactions). Decide: answer with "
|
||||
"delegate_answer, escalate an above-authority question with the "
|
||||
"escalate verb (parent-first; the reply reaches your mailbox on a "
|
||||
"later round), or "
|
||||
f"keep waiting (call again) — {waiting_expiry_clause(pending_interactions)}. "
|
||||
"Do not cancel a run merely because it asked a question.")
|
||||
|
||||
|
||||
def window_payload(
|
||||
*,
|
||||
run_id: str,
|
||||
|
|
@ -673,14 +686,7 @@ def window_payload(
|
|||
# generic delegate_cancel hint invited cancelling a run that is simply
|
||||
# waiting to be answered. The waiting state gets its own note.
|
||||
payload["note"] = (
|
||||
("The run is alive and PAUSED on the question(s) it already asked "
|
||||
"(waiting_on_user; see pending_interactions). Decide: answer with "
|
||||
"delegate_answer, escalate an above-authority question with the "
|
||||
"escalate verb (parent-first; the reply reaches your mailbox on a "
|
||||
"later round), or "
|
||||
f"keep waiting (call again) — {waiting_expiry_clause(pending_interactions)}. "
|
||||
"Do not cancel a run merely because it asked a question.")
|
||||
if waiting_on_user else
|
||||
paused_note(pending_interactions) if waiting_on_user else
|
||||
("The run is alive but silent. Decide: keep waiting (call again), "
|
||||
"or delegate_cancel if it is stuck."))
|
||||
_fitted_pending(payload, list(pending_interactions or []), budget)
|
||||
|
|
|
|||
|
|
@ -407,6 +407,10 @@ def _restore_wait_checkpoint(drive_root: Any, row: Mapping[str, Any]) -> None:
|
|||
state["last_wake"] = dict(payload)
|
||||
if isinstance(row.get("checkpoint"), dict):
|
||||
state["checkpoint"] = dict(row["checkpoint"])
|
||||
# A run-id mismatch discarded the state above; the acked child cursor is
|
||||
# task-scoped, so the successor must not re-announce delivered child events.
|
||||
if isinstance(row.get("coordination_cursor"), dict) and not isinstance(state.get("coordination_cursor"), dict):
|
||||
state["coordination_cursor"] = dict(row["coordination_cursor"])
|
||||
path.parent.mkdir(parents=True, exist_ok=True)
|
||||
atomic_write_json(path, state)
|
||||
interaction_ids = frozenset(
|
||||
|
|
@ -572,6 +576,9 @@ def prepare_handoff(
|
|||
supervision.get("interaction_acknowledged_ids") or []
|
||||
),
|
||||
"pending_wake": dict(pending_wake),
|
||||
# The task-scoped COMMITTED child-delivery cursor (acked events only).
|
||||
"coordination_cursor": dict(supervision["coordination_cursor"])
|
||||
if isinstance(supervision.get("coordination_cursor"), dict) else {},
|
||||
"checkpoint": supervision.get("checkpoint") if isinstance(supervision.get("checkpoint"), dict) else {},
|
||||
"no_resume_veto_causes": list(NO_RESUME_CAUSES),
|
||||
"created_at": utc_now_iso(),
|
||||
|
|
|
|||
|
|
@ -98,10 +98,13 @@ def _load_state(ctx: Any, run_id: str) -> dict[str, Any]:
|
|||
# supervised_wait's own entry persisted the reset, so a worker crash
|
||||
# during the (hours-long) wait lost the hold and the recovered
|
||||
# successor dispatched the unknown transcript again (final-pair F2).
|
||||
hold = data.get("unknown_provider_hold") if isinstance(data, dict) else None
|
||||
data = {"schema": 1, "run_id": str(run_id), "journal_cursor": 0}
|
||||
if isinstance(hold, dict):
|
||||
data["unknown_provider_hold"] = hold
|
||||
# The COMMITTED coordination cursor is task-scoped too: a new run id must
|
||||
# not re-announce every already-acknowledged child terminal/beacon. Only
|
||||
# the acked cursor rides; an unacked ``pending_wake`` (and its uncommitted
|
||||
# cursor) stays run-scoped, so an undelivered child event re-emits once.
|
||||
carried = {key: data[key] for key in ("unknown_provider_hold", "coordination_cursor")
|
||||
if isinstance(data, dict) and isinstance(data.get(key), dict)}
|
||||
data = {"schema": 1, "run_id": str(run_id), "journal_cursor": 0, **carried}
|
||||
return data
|
||||
|
||||
|
||||
|
|
@ -639,7 +642,8 @@ def _render_wake_payload(ctx: Any, payload: dict[str, Any]) -> ToolResult:
|
|||
envelope: dict[str, Any] = {
|
||||
key: (str(value)[:600] if isinstance(value, str) else value)
|
||||
for key in ("status", "ok", "host_code", "run_id", "state", "last_seq", "reason",
|
||||
"continuation", "continuation_note")
|
||||
"continuation", "continuation_note", "sleep", "cache_horizon_note",
|
||||
"leaf_live_input")
|
||||
if (value := payload.get(key)) not in (None, "")
|
||||
}
|
||||
envelope["supervision_wake_id"] = wake_id
|
||||
|
|
@ -683,7 +687,8 @@ def _render_wake_payload(ctx: Any, payload: dict[str, Any]) -> ToolResult:
|
|||
**({"ok": False, "host_code": str(payload.get("host_code") or "")}
|
||||
if payload.get("ok") is False else {}),
|
||||
"run_id": str(payload.get("run_id") or "")[:200],
|
||||
**{key: payload[key] for key in ("continuation", "continuation_note") if key in payload},
|
||||
**{key: payload[key] for key in ("continuation", "continuation_note", "sleep",
|
||||
"leaf_live_input") if key in payload},
|
||||
"supervision_wake_id": wake_id,
|
||||
"coordination_context": {"state": "available_in_full_wake_source"},
|
||||
"wake_delivery": envelope["wake_delivery"],
|
||||
|
|
@ -864,6 +869,126 @@ def _control_wakes(ctx: Any) -> list[dict[str, Any]]:
|
|||
return wakes
|
||||
|
||||
|
||||
def _open_sleep_entry(state: dict[str, Any], run_id: str) -> None:
|
||||
"""Open this call's sleep: every ``supervised_wait`` call is a NEW sleep.
|
||||
|
||||
The prior status only labels it, it never carries elapsed time over: an
|
||||
``adopted`` state is a worker-loss recovery, a ``sleeping`` one is a previous
|
||||
call that died without publishing a wake (tool timeout, exception)."""
|
||||
prior = str(state.get("status") or "")
|
||||
previous = state.get("sleep_entry") if isinstance(state.get("sleep_entry"), dict) else {}
|
||||
entry: dict[str, Any] = {
|
||||
"run_id": str(run_id), "entered_at": utc_now_iso(), "entered_at_unix": time.time(),
|
||||
"quiet_renewals_at_entry": int(state.get("quiet_renewals") or 0),
|
||||
"journal_cursor_at_entry": int(state.get("journal_cursor") or 0),
|
||||
}
|
||||
if prior == "adopted":
|
||||
entry["opened_after"] = "worker_loss_adoption"
|
||||
elif prior == "sleeping":
|
||||
entry["opened_after"] = "interrupted_sleep"
|
||||
if previous.get("entered_at"):
|
||||
entry["previous_entered_at"] = str(previous["entered_at"])
|
||||
state["sleep_entry"] = entry
|
||||
|
||||
|
||||
def _record_observation(state: dict[str, Any], payload: dict[str, Any], answered: bool) -> None:
|
||||
"""Dated observation facts (run-scoped): the supervisor's own ``status`` and the
|
||||
file's ``updated_at`` move on every quiet renewal, failed reads included, so
|
||||
freshness is the last ANSWERED observation plus the current failure, never mtime."""
|
||||
at = utc_now_iso()
|
||||
state["observation"] = {
|
||||
"at": at, "answered": bool(answered), "status": str(payload.get("status") or ""),
|
||||
"run_state": str(payload.get("state") or ""), "reason": str(payload.get("reason") or ""),
|
||||
}
|
||||
if answered:
|
||||
state["last_answered_observation_at"] = at
|
||||
state["observation_failure"] = "" if answered else str(
|
||||
payload.get("reason") or payload.get("status") or "unknown")
|
||||
|
||||
|
||||
def _leaf_live_input(ctx: Any, run_id: str, gateway: Any) -> str:
|
||||
"""The leaf route's DECLARED live-input capability (engine-static), read ONCE per
|
||||
call at entry through the loop's transport (a one-read transport when the caller
|
||||
observes with its own) and stamped on every wake: the route row's ``liveInput``
|
||||
in the engine catalog. A capability, never a liveness verdict; an unreadable
|
||||
catalog, an unknown route or a read slower than one beat is ``unknown`` (never a
|
||||
refusal, and nothing is read on the wake path). A row without the field is an
|
||||
engine that has no live-input channel at all: ``none``."""
|
||||
try:
|
||||
_status, entry = custody.lookup(
|
||||
custody.custody_root(ctx), str(getattr(ctx, "task_id", "") or ""), run_id)
|
||||
route = str(getattr(entry, "route_id", "") or "")
|
||||
except Exception:
|
||||
return "unknown"
|
||||
reader = (gateway if gateway is not None else _loop_gateway()) if route else None
|
||||
if reader is None:
|
||||
return "unknown"
|
||||
try:
|
||||
catalog = reader.agent_capabilities(timeout_sec=_TICK_SEC)
|
||||
row = next((item for item in (catalog.get("harnesses") or [])
|
||||
if isinstance(item, dict) and item.get("id") == route), None)
|
||||
except Exception:
|
||||
row = None
|
||||
finally:
|
||||
if reader is not gateway:
|
||||
_drop_gateway(reader)
|
||||
return "unknown" if row is None else str(row.get("liveInput") or "none")
|
||||
|
||||
|
||||
def _since_last_model_response_sec(ctx: Any) -> Optional[float]:
|
||||
"""Seconds since this task's latest model response was recorded
|
||||
(``_last_llm_call_meta.ts``). A LOWER bound on prompt-cache age: the provider
|
||||
read/wrote the cached prefix before that response finished, so the horizon note
|
||||
built on it can only fire late, never early. Unknown (silent) after a worker
|
||||
loss, whose successor has no recorded send yet."""
|
||||
usage = getattr(ctx, "_accumulated_usage", None)
|
||||
meta = usage.get("_last_llm_call_meta") if isinstance(usage, dict) else None
|
||||
try:
|
||||
import datetime as _dt
|
||||
|
||||
stamp = _dt.datetime.fromisoformat(str((meta or {}).get("ts") or ""))
|
||||
return max(0.0, (_dt.datetime.now(tz=_dt.timezone.utc) - stamp).total_seconds())
|
||||
except (TypeError, ValueError):
|
||||
return None
|
||||
|
||||
|
||||
def _publish_wake_facts(ctx: Any, state: dict[str, Any], payload: dict[str, Any],
|
||||
live_input: str) -> None:
|
||||
"""Whole-sleep facts at the ONE durable wake-publication point, every status
|
||||
(terminal included), before ``pending_wake`` is stored, so a replay repeats the
|
||||
recorded facts exactly. A quiet-status wake (an event arrived during a quiet
|
||||
tick) drops the tick's own ``waited_sec``/``quiet_for_sec`` and its generic
|
||||
keep-watching/cancel note; a run paused on its question keeps the PAUSED note."""
|
||||
entry = state.get("sleep_entry") if isinstance(state.get("sleep_entry"), dict) else {}
|
||||
if entry:
|
||||
sleep: dict[str, Any] = {
|
||||
"entered_at": str(entry.get("entered_at") or ""),
|
||||
"slept_sec": round(max(0.0, time.time() - float(entry.get("entered_at_unix") or 0)), 1),
|
||||
"quiet_renewals": max(0, int(state.get("quiet_renewals") or 0)
|
||||
- int(entry.get("quiet_renewals_at_entry") or 0)),
|
||||
"journal_advances": max(0, int(state.get("journal_cursor") or 0)
|
||||
- int(entry.get("journal_cursor_at_entry") or 0)),
|
||||
}
|
||||
sleep.update({key: entry[key] for key in ("opened_after", "previous_entered_at") if entry.get(key)})
|
||||
payload["sleep"] = sleep
|
||||
if str(payload.get("status") or "") in _QUIET_STATUSES:
|
||||
payload.pop("waited_sec", None)
|
||||
payload.pop("quiet_for_sec", None)
|
||||
if payload.get("waiting_on_user"):
|
||||
from ouroboros.delegate_progress import paused_note
|
||||
|
||||
payload["note"] = paused_note(payload.get("pending_interactions"))
|
||||
else:
|
||||
payload.pop("note", None)
|
||||
payload.pop("cache_horizon_note", None)
|
||||
from ouroboros.tools.control import cache_horizon_note
|
||||
|
||||
horizon = cache_horizon_note(ctx, _since_last_model_response_sec(ctx))
|
||||
if horizon:
|
||||
payload["cache_horizon_note"] = horizon
|
||||
payload["leaf_live_input"] = live_input
|
||||
|
||||
|
||||
def supervised_wait(
|
||||
ctx: Any,
|
||||
run_id: str,
|
||||
|
|
@ -913,9 +1038,10 @@ def supervised_wait(
|
|||
snapshot = snapshot.get("configured_subagent") if isinstance(snapshot, dict) else {}
|
||||
snapshot = snapshot if isinstance(snapshot, dict) else {}
|
||||
state["config_fingerprint"] = str(snapshot.get("config_fingerprint") or "")
|
||||
state["status"] = "sleeping"
|
||||
if since_seq is not None:
|
||||
state["journal_cursor"] = max(int(state.get("journal_cursor") or 0), int(since_seq))
|
||||
_open_sleep_entry(state, run_id)
|
||||
state["status"] = "sleeping"
|
||||
if checkpoint_after_sec is not None:
|
||||
delay = max(1, min(604_800, int(checkpoint_after_sec)))
|
||||
state["checkpoint"] = {
|
||||
|
|
@ -936,6 +1062,7 @@ def supervised_wait(
|
|||
_emit(ctx, "delegate_supervision_wait_entered", {
|
||||
"run_id": str(run_id), "journal_cursor": int(state.get("journal_cursor") or 0),
|
||||
"checkpoint_scheduled": bool(checkpoint_after_sec is not None),
|
||||
"sleep_entry": dict(state["sleep_entry"]),
|
||||
})
|
||||
|
||||
# The OPEN unreachable-daemon episode: its UTC start (empty when none) plus the
|
||||
|
|
@ -943,7 +1070,9 @@ def supervised_wait(
|
|||
# dedup surfaces (the loop_transport.incident precedent).
|
||||
outage_since = ""
|
||||
outage_stamp = ""
|
||||
gateway = None # the loop's own transport, only when it built the observing wait
|
||||
# The loop's own transport, only when it built the observing wait.
|
||||
gateway = _loop_gateway() if owns_transport else None
|
||||
live_input = _leaf_live_input(ctx, run_id, gateway)
|
||||
try:
|
||||
while True:
|
||||
# A control already present does not wait behind another HTTP read.
|
||||
|
|
@ -957,6 +1086,8 @@ def supervised_wait(
|
|||
**({"gateway": gateway} if gateway is not None else {}))
|
||||
payload = _payload(raw)
|
||||
answered = str(payload.get("status") or "") not in _NO_DAEMON_ANSWER_STATUSES
|
||||
if observed:
|
||||
_record_observation(state, payload, answered)
|
||||
if not answered:
|
||||
gateway = _drop_gateway(gateway)
|
||||
unreachable = (
|
||||
|
|
@ -1010,6 +1141,7 @@ def supervised_wait(
|
|||
if wakes:
|
||||
payload["wake_events"] = wakes
|
||||
payload["coordination_context"] = coordination_live_context(ctx)
|
||||
_publish_wake_facts(ctx, state, payload, live_input)
|
||||
if ignored_note:
|
||||
payload["ignored_arguments"] = [ignored_note]
|
||||
wake_id = uuid.uuid4().hex
|
||||
|
|
|
|||
|
|
@ -554,8 +554,9 @@ class ClaudexorGateway:
|
|||
self._engine_build_sha = str(engine.get("sha") or "")
|
||||
return body
|
||||
|
||||
def agent_capabilities(self) -> Dict[str, Any]:
|
||||
body = self._request("GET", "/v2/agent-capabilities")
|
||||
def agent_capabilities(self, *, timeout_sec: Optional[float] = None) -> Dict[str, Any]:
|
||||
body = self._request("GET", "/v2/agent-capabilities",
|
||||
**({"timeout_sec": timeout_sec} if timeout_sec is not None else {}))
|
||||
return body if isinstance(body, dict) else {}
|
||||
|
||||
# Model operations use the same private control transport, not Agent runs
|
||||
|
|
|
|||
|
|
@ -334,6 +334,37 @@ def _maybe_inject_cost_budget_milestone(
|
|||
return True
|
||||
|
||||
|
||||
def _own_round_cost_sentence(ctx: Any, cost: float) -> tuple[str, str]:
|
||||
"""The reminder's money sentence from COST EVIDENCE first: a positive priced
|
||||
delta (provider-reported or estimated cost of the nanny's own rounds) is
|
||||
metered money; no price is an unknown cash cost, never zero. The latest round's
|
||||
route is named only as context, and a fallback clause appears only when a
|
||||
fallback chain is actually configured (a paid round is no proof the next is)."""
|
||||
usage = getattr(ctx, "_accumulated_usage", None)
|
||||
meta = usage.get("_last_llm_call_meta") if isinstance(usage, dict) else None
|
||||
meta = meta if isinstance(meta, dict) else {}
|
||||
model = str(meta.get("resolved_model") or meta.get("model") or "")
|
||||
route = " / ".join(item for item in (str(meta.get("provider") or ""), model) if item)
|
||||
if cost > 0:
|
||||
cost_class, sentence = "priced", (
|
||||
"Your own rounds here reported a price, so that spend is metered money")
|
||||
else:
|
||||
cost_class, sentence = "unpriced", (
|
||||
"No provider price reached this task for your own rounds here: an unknown "
|
||||
"cash cost, not a zero one")
|
||||
sentence += f" (latest round: {route})." if route else "."
|
||||
try:
|
||||
from ouroboros.model_slots import get_fallback_models
|
||||
|
||||
fallbacks = get_fallback_models(str(meta.get("model") or ""))
|
||||
except Exception:
|
||||
fallbacks = []
|
||||
if fallbacks:
|
||||
sentence += (f" If this route stops serving, the configured fallback ({', '.join(fallbacks[:3])}) "
|
||||
"may take over, and a fallback round may be metered.")
|
||||
return cost_class, sentence
|
||||
|
||||
|
||||
def _maybe_inject_nanny_economics_reminder(
|
||||
round_idx: int,
|
||||
messages: List[Dict[str, Any]],
|
||||
|
|
@ -378,13 +409,13 @@ def _maybe_inject_nanny_economics_reminder(
|
|||
# BR1-3: never an unconditional "$0" claim — the owner's wording law is
|
||||
# typed cost classes: known-zero only on a settled $0 spend, never "free"
|
||||
# unqualified (estimated/undisclosed spend is never zero).
|
||||
cost_class, cost_sentence = _own_round_cost_sentence(ctx, cost)
|
||||
reminder = (
|
||||
"[NANNY ECONOMICS REMINDER]\n"
|
||||
f"You are a harness-dispatched NANNY and you have spent {_nanny_burn_phrase(rounds, cost)} "
|
||||
f"{since_phrase}. A subscription-lane delegated run has known-zero "
|
||||
"marginal cost only when its settled spend reports $0 (estimated or "
|
||||
"undisclosed spend is never zero); every round you think yourself is "
|
||||
"metered API money.\n"
|
||||
f"undisclosed spend is never zero). {cost_sentence}\n"
|
||||
"This is a reminder, not a stop. Consider: delegate the remaining work "
|
||||
"(delegate_start / delegate_wait — follow-up work and fixes are delegated too), "
|
||||
"and keep your own rounds for judgment: acceptance, integration, honest "
|
||||
|
|
@ -400,6 +431,7 @@ def _maybe_inject_nanny_economics_reminder(
|
|||
"round": round_idx,
|
||||
"metered_rounds_since_delegate_activity": rounds,
|
||||
"metered_cost_since_delegate_activity_usd": round(cost, 4),
|
||||
"own_round_cost_class": cost_class,
|
||||
})
|
||||
return True
|
||||
|
||||
|
|
|
|||
|
|
@ -132,7 +132,7 @@ def nanny_burn_phrase(rounds: int, cost: float) -> str:
|
|||
return f"~${cost:.2f} of your own metered spend"
|
||||
if cost > 0:
|
||||
return f"{rounds} of your own metered LLM rounds (~${cost:.2f})"
|
||||
return f"{rounds} of your own metered LLM rounds"
|
||||
return f"{rounds} of your own LLM rounds (no provider price reported: cash cost unknown, not zero)"
|
||||
|
||||
|
||||
# Compatibility spellings retained on ``ouroboros.loop`` through imports.
|
||||
|
|
|
|||
|
|
@ -1146,10 +1146,6 @@ def _delegate_wait(ctx: ToolContext, run_id: str, wait_sec: Optional[int] = None
|
|||
# PREVIEW path too: a payload big enough to spill is exactly the one whose
|
||||
# containment block a reader is least likely to reach.
|
||||
_record_containment(ctx, entry, payload)
|
||||
from ouroboros.tools.control import cache_horizon_note
|
||||
_horizon = cache_horizon_note(ctx, time.monotonic() - started)
|
||||
if _horizon:
|
||||
payload["cache_horizon_note"] = _horizon
|
||||
return json.dumps(payload, ensure_ascii=False, indent=2)
|
||||
if last_seq > baseline:
|
||||
# The STREAM is not collapsed — the TIMER is. Every advance reaches the
|
||||
|
|
@ -1177,10 +1173,9 @@ def _delegate_wait(ctx: ToolContext, run_id: str, wait_sec: Optional[int] = None
|
|||
if entry.work_order_source_request else {}
|
||||
))
|
||||
def _expired() -> str:
|
||||
# The window payload is a DICT here, and the cache-horizon note is a
|
||||
# field in it: appending the note after the rendered JSON left the
|
||||
# result unparseable for every reader of this family — the supervising
|
||||
# loop included, which then read the whole window as a `fault`.
|
||||
# The window payload is a DICT; the supervising wake adds its own
|
||||
# whole-sleep facts (``sleep``, one cache-horizon note) as fields at
|
||||
# publication, never per tick and never after the rendered JSON.
|
||||
payload = progress.window_payload(
|
||||
run_id=rid, state=state, last_seq=last_seq,
|
||||
window=(time.monotonic() - started) if observation_only else window,
|
||||
|
|
@ -1191,9 +1186,6 @@ def _delegate_wait(ctx: ToolContext, run_id: str, wait_sec: Optional[int] = None
|
|||
pending_interactions=_bounded_interactions(pending) if pending else None,
|
||||
detail=detail, seen=seen,
|
||||
budget=tool_result_limit("delegate_wait"))
|
||||
from ouroboros.tools.control import cache_horizon_note
|
||||
if _horizon := cache_horizon_note(ctx, time.monotonic() - started):
|
||||
payload["cache_horizon_note"] = _horizon
|
||||
return json.dumps(payload, ensure_ascii=False, indent=2)
|
||||
|
||||
if observation_only or time.monotonic() >= deadline:
|
||||
|
|
@ -1434,7 +1426,10 @@ def get_tools() -> List[ToolEntry]:
|
|||
"Terminal settlement, a new interaction, "
|
||||
"fault, addressed owner/task message, a direct-child attention/terminal event, "
|
||||
"cancel/deadline control, recovery judgment, or an explicit one-shot checkpoint "
|
||||
"wakes exactly once. A run that asks its "
|
||||
"wakes exactly once. Every wake carries `sleep` (measured over this whole call), "
|
||||
"`leaf_live_input` (the route's declared live-input capability, 'unknown' when "
|
||||
"unread) and, once the prompt-cache horizon has passed since your last model "
|
||||
"response, one cache_horizon_note. A run that asks its "
|
||||
"user a question returns IMMEDIATELY as status='waiting_on_user' with "
|
||||
"the full question set (interaction/question ids ride WHOLE, never "
|
||||
"truncated): answer it with delegate_answer, or raise it with the "
|
||||
|
|
|
|||
|
|
@ -957,11 +957,10 @@ def test_one_shot_checkpoint_is_reasoned_and_consumed(monkeypatch, tmp_path):
|
|||
assert coordination_context["parent_intent"]["state"] == "absent"
|
||||
assert coordination_context["time"]["state"] == "not_set"
|
||||
assert coordination_context["review_capacity"]["state"] == "available"
|
||||
assert out.pop("sleep")["slept_sec"] == 2.0 and out.pop("leaf_live_input") == "unknown"
|
||||
assert out == {
|
||||
"status": "inspection_checkpoint",
|
||||
"run_id": "run-1",
|
||||
"reason": "inspect a promised artifact",
|
||||
"last_seq": 0,
|
||||
"status": "inspection_checkpoint", "run_id": "run-1",
|
||||
"reason": "inspect a promised artifact", "last_seq": 0,
|
||||
}
|
||||
state = supervision.supervision_checkpoint(ctx)
|
||||
assert state["checkpoint"]["consumed"] is True
|
||||
|
|
|
|||
|
|
@ -379,17 +379,22 @@ def test_an_unsupported_engine_build_refuses_the_answer_typed(tmp_path, monkeypa
|
|||
assert json.loads(result.text)["reason"] == "interaction_answers_unsupported"
|
||||
|
||||
|
||||
def test_the_expired_wait_window_and_its_cache_horizon_note_stay_valid_json(tmp_path, monkeypatch):
|
||||
"""The note is a FIELD of the window payload; appended after the JSON it made
|
||||
the whole result unparseable for every reader of this family."""
|
||||
def test_the_supervising_wake_and_its_cache_horizon_note_stay_valid_json(tmp_path, monkeypatch):
|
||||
"""The note is a FIELD of the wake payload, stamped once at the wake's publication
|
||||
(never per 3 s tick, never appended after the rendered JSON, which once left the
|
||||
whole result unparseable for every reader of this family). Drives the REAL
|
||||
observing tick through the supervising wait: an event during a quiet tick wakes
|
||||
with the note, and without the per-tick window's waited_sec or cancel advice."""
|
||||
import json
|
||||
|
||||
import ouroboros.delegate_supervision as supervision
|
||||
import ouroboros.tools.control as control
|
||||
import ouroboros.tools.delegate as delegate
|
||||
from ouroboros.gateways import claudexor as gw
|
||||
|
||||
_own_run(tmp_path)
|
||||
monkeypatch.setattr(control, "cache_horizon_note", lambda *_a, **_k: "the prompt cache expires soon")
|
||||
monkeypatch.setattr(delegate.time, "sleep", lambda _sec: None)
|
||||
|
||||
class _Alive:
|
||||
engine_version = "3.10.2"
|
||||
|
|
@ -400,11 +405,18 @@ def test_the_expired_wait_window_and_its_cache_horizon_note_stay_valid_json(tmp_
|
|||
def close(self): pass
|
||||
|
||||
monkeypatch.setattr(gw, "ClaudexorGateway", lambda *a, **k: _Alive())
|
||||
ctx = SimpleNamespace(repo_dir=tmp_path, drive_root=tmp_path, task_id="t-family")
|
||||
raw = delegate._delegate_wait(ctx, "run-1", wait_sec=1, since_seq=0)
|
||||
checks = []
|
||||
# First check lets the tick observe; the second (after it) delivers a control.
|
||||
monkeypatch.setattr(supervision, "_control_wakes",
|
||||
lambda _ctx: checks.append(1) or ([{"type": "deadline"}] if len(checks) > 1 else []))
|
||||
ctx = SimpleNamespace(repo_dir=tmp_path, drive_root=tmp_path, task_id="t-family", task_attempt=1)
|
||||
raw = supervision.supervised_wait(ctx, "run-1").text
|
||||
payload = json.loads(raw) # the contract: still ONE JSON object
|
||||
assert payload["cache_horizon_note"] == "the prompt cache expires soon"
|
||||
assert payload["status"] in {"progress", "no_progress"}
|
||||
assert payload["wake_events"] == [{"type": "deadline"}]
|
||||
assert "waited_sec" not in payload and "delegate_cancel" not in str(payload.get("note") or "")
|
||||
assert payload["sleep"]["quiet_renewals"] == 0
|
||||
|
||||
|
||||
def _wait_ctx(tmp_path, task_id="t-nanny"):
|
||||
|
|
|
|||
|
|
@ -466,46 +466,84 @@ def test_wait_for_task_appends_cache_horizon_note(tmp_path, monkeypatch):
|
|||
)
|
||||
|
||||
|
||||
def test_cache_horizon_reachability_matches_the_wait_clamps():
|
||||
"""Honest reachability of the wait disclosure, derived from the REAL clamps.
|
||||
def test_cache_horizon_reachability_matches_the_wait_clamps(tmp_path, monkeypatch):
|
||||
"""Honest reachability of the wait disclosure, measured on BEHAVIOR.
|
||||
|
||||
Each wait tool caps its own window, so the line's availability is per-tool and
|
||||
per-TTL: at the shipped '1h' default only wait_tasks (7200s) can genuinely emit
|
||||
it, wait_task sits exactly on the 3600s horizon and delegate_wait (1800s window
|
||||
max — the wait's REAL clamp since F5a, not the 2100s ToolEntry kill timeout)
|
||||
cannot reach it at all; at '5m' all three do. The stream report advertised it as
|
||||
a live capability of all three at the default — this pin makes the truth loud and
|
||||
fails if a clamp, the ceiling, or the tier scale moves without revisiting it."""
|
||||
import inspect
|
||||
import re
|
||||
The root waits measure their own elapsed window, so their line's availability is
|
||||
per-tool and per-TTL and follows the window each one actually hands the poller: at
|
||||
the shipped '1h' default only wait_tasks (7200s) can genuinely emit it and
|
||||
wait_task sits exactly on the 3600s horizon; at '5m' both do. delegate_wait no
|
||||
longer measures its 3 s tick: its wake measures the time since the task's last
|
||||
model response, so it can emit at ANY tier (pinned in the supervising-wake test
|
||||
below). Fails if a clamp, the ceiling or the tier scale moves without revisiting
|
||||
the disclosure."""
|
||||
import json
|
||||
from types import SimpleNamespace
|
||||
|
||||
from ouroboros.config import DELEGATE_WAIT_WINDOW_MAX_SEC
|
||||
from ouroboros.llm import cache_ttl_seconds
|
||||
from ouroboros.task_results import STATUS_COMPLETED, write_task_result
|
||||
from ouroboros.tools import control_task_results as control_mod
|
||||
from ouroboros.tools.control import cache_horizon_note
|
||||
|
||||
def _clamp(fn):
|
||||
found = re.findall(r"min\(int\(timeout_sec\),\s*(\d+)\)", inspect.getsource(fn))
|
||||
assert len(found) == 1, f"{fn.__name__}: expected one timeout clamp, found {found}"
|
||||
return int(found[0])
|
||||
write_task_result(tmp_path, "child42", STATUS_COMPLETED, result="done")
|
||||
handed = []
|
||||
|
||||
ceilings = {
|
||||
"wait_task": _clamp(control_mod._wait_for_task),
|
||||
"wait_tasks": _clamp(control_mod._wait_for_tasks),
|
||||
"delegate_wait": DELEGATE_WAIT_WINDOW_MAX_SEC,
|
||||
}
|
||||
assert ceilings == {"wait_task": 3600, "wait_tasks": 7200, "delegate_wait": 1800}
|
||||
def _record(*_args, timeout_sec=None, **_kwargs):
|
||||
handed.append(timeout_sec)
|
||||
return {"all_terminal": True, "elapsed_sec": 0.0, "tasks": {}}
|
||||
|
||||
monkeypatch.setattr(control_mod, "wait_for_effective_tasks", _record)
|
||||
ctx = SimpleNamespace(drive_root=tmp_path)
|
||||
control_mod._wait_for_task(ctx, "child42", timeout_sec=10 ** 6)
|
||||
json.loads(control_mod._wait_for_tasks(ctx, ["child42"], timeout_sec=10 ** 6))
|
||||
ceilings = {"wait_task": handed[0], "wait_tasks": handed[1]}
|
||||
assert ceilings == {"wait_task": control_mod._WAIT_TASK_CLAMP_SEC,
|
||||
"wait_tasks": control_mod._WAIT_TASKS_CLAMP_SEC}
|
||||
assert ceilings == {"wait_task": 3600, "wait_tasks": 7200}
|
||||
|
||||
def _emits(tier, ceiling):
|
||||
ctx = SimpleNamespace(_accumulated_usage={"_last_prompt_cache_ttl": tier})
|
||||
return bool(cache_horizon_note(ctx, float(ceiling)))
|
||||
|
||||
assert cache_ttl_seconds("1h") == 3600 and cache_ttl_seconds("5m") == 300
|
||||
at_1h = {name: _emits("1h", sec) for name, sec in ceilings.items()}
|
||||
assert at_1h == {"wait_task": False, "wait_tasks": True, "delegate_wait": False}
|
||||
at_5m = {name: _emits("5m", sec) for name, sec in ceilings.items()}
|
||||
assert at_5m == {"wait_task": True, "wait_tasks": True, "delegate_wait": True}
|
||||
assert {name: _emits("1h", sec) for name, sec in ceilings.items()} == {
|
||||
"wait_task": False, "wait_tasks": True}
|
||||
assert {name: _emits("5m", sec) for name, sec in ceilings.items()} == {
|
||||
"wait_task": True, "wait_tasks": True}
|
||||
|
||||
doc = cache_horizon_note.__doc__ or ""
|
||||
assert "REACHABILITY" in doc # the limitation is stated where the note is built
|
||||
|
||||
|
||||
def _response_meta(seconds_ago):
|
||||
import datetime as _dt
|
||||
|
||||
stamp = _dt.datetime.now(tz=_dt.timezone.utc) - _dt.timedelta(seconds=seconds_ago)
|
||||
return {"ts": stamp.isoformat(), "provider": "anthropic", "model": "m"}
|
||||
|
||||
|
||||
def test_supervising_wake_horizon_fires_only_past_the_applied_ttl(tmp_path):
|
||||
"""delegate_wait's note is computed ONCE per wake from the task's last model
|
||||
response (a lower bound on cache age): it fires past the APPLIED horizon, stays
|
||||
silent below it and on an unknown tier, and a replay repeats the recorded wake."""
|
||||
import json
|
||||
from types import SimpleNamespace
|
||||
|
||||
from ouroboros.delegate_supervision import supervised_wait
|
||||
|
||||
def _wake(ttl, seconds_ago):
|
||||
ctx = SimpleNamespace(
|
||||
task_id=f"t-{ttl or 'none'}-{seconds_ago}", task_attempt=1, drive_root=tmp_path,
|
||||
_accumulated_usage={"_last_prompt_cache_ttl": ttl,
|
||||
"_last_llm_call_meta": _response_meta(seconds_ago)})
|
||||
raw = supervised_wait(ctx, "run-h", wait_once=lambda *_a: json.dumps(
|
||||
{"status": "completed", "run_id": "run-h", "last_seq": 1})).text
|
||||
return ctx, json.loads(raw)
|
||||
|
||||
_ctx, past = _wake("1h", 3700)
|
||||
assert "1h" in past["cache_horizon_note"] and "may be cold" in past["cache_horizon_note"]
|
||||
assert "cache_horizon_note" not in _wake("1h", 3500)[1]
|
||||
assert "cache_horizon_note" not in _wake("", 10_000)[1]
|
||||
replay = json.loads(supervised_wait(_ctx, "run-h", wait_once=lambda *_a: (_ for _ in ()).throw(
|
||||
AssertionError("an unacknowledged wake replays before any new observation"))).text)
|
||||
assert replay == past
|
||||
|
|
|
|||
|
|
@ -42,6 +42,9 @@ def test_public_wait_reads_response_slower_than_five_seconds(tmp_path, monkeypat
|
|||
|
||||
def do_GET(self):
|
||||
requests.append((self.command, self.path))
|
||||
if self.path == "/v2/agent-capabilities":
|
||||
self.answer({"harnesses": [{"id": "fixture", "liveInput": "mid_turn"}]})
|
||||
return
|
||||
assert self.path == "/v2/runs/run-slow"
|
||||
# Deliberate network-latency reproduction: the previous five-second
|
||||
# HTTP timeout fails before this real socket sends its headers.
|
||||
|
|
@ -71,9 +74,12 @@ def test_public_wait_reads_response_slower_than_five_seconds(tmp_path, monkeypat
|
|||
assert time.monotonic() - started >= 5.0
|
||||
assert result["status"] == "terminal", result
|
||||
assert result["state"] == "succeeded"
|
||||
# One handshake per supervision loop, one GET per tick (S2): the loop holds the
|
||||
# transport across quiet ticks instead of rebuilding it every 3 s.
|
||||
assert requests == [("POST", "/v2/handshake")] + [("GET", "/v2/runs/run-slow")] * (2 if initially_queued else 1)
|
||||
# The route's live-input capability is read ONCE at entry on the same transport
|
||||
# and stamped on the wake; then one handshake per supervision loop, one GET per
|
||||
# tick (S2): the loop holds the transport across quiet ticks.
|
||||
assert result["leaf_live_input"] == "mid_turn"
|
||||
assert requests == [("GET", "/v2/agent-capabilities"), ("POST", "/v2/handshake")] + [
|
||||
("GET", "/v2/runs/run-slow")] * (2 if initially_queued else 1)
|
||||
finally:
|
||||
server.shutdown()
|
||||
server.server_close()
|
||||
|
|
|
|||
328
tests/test_delegate_wake_truthful_facts.py
Normal file
328
tests/test_delegate_wake_truthful_facts.py
Normal file
|
|
@ -0,0 +1,328 @@
|
|||
"""Truthful supervising-wait facts: the task-scoped child cursor across run ids and
|
||||
recovery, whole-sleep facts at the single wake-publication point, dated observation
|
||||
facts, and the leaf route's declared live-input capability."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
from types import SimpleNamespace
|
||||
|
||||
import pytest
|
||||
|
||||
import ouroboros.delegate_supervision as supervision
|
||||
from ouroboros import task_tree_ledger
|
||||
from ouroboros.delegate_supervision import acknowledge_pending_wake, supervised_wait
|
||||
from ouroboros.task_results import STATUS_COMPLETED, STATUS_RUNNING, write_task_result
|
||||
|
||||
|
||||
def _ctx(tmp_path, task_id="parent"):
|
||||
return SimpleNamespace(
|
||||
task_id=task_id, task_attempt=1, drive_root=tmp_path, budget_drive_root=str(tmp_path),
|
||||
task_metadata={"root_task_id": "root", "delegation_role": "subagent"},
|
||||
)
|
||||
|
||||
|
||||
def _child(tmp_path, task_id="child", status=STATUS_RUNNING):
|
||||
return write_task_result(tmp_path, task_id, status, parent_task_id="parent",
|
||||
root_task_id="root", delegation_role="subagent")
|
||||
|
||||
|
||||
def _settled(run_id, seq=1):
|
||||
return lambda *_a, **_k: json.dumps({"status": "completed", "run_id": run_id, "last_seq": seq})
|
||||
|
||||
|
||||
def _child_events(wake, kind="child_terminal"):
|
||||
return [event for event in wake.get("wake_events", []) if event.get("type") == kind]
|
||||
|
||||
|
||||
def _state(ctx):
|
||||
return json.loads(supervision._state_path(ctx).read_text(encoding="utf-8"))
|
||||
|
||||
|
||||
# -- A: the committed child-delivery cursor is task-scoped ---------------------
|
||||
|
||||
|
||||
@pytest.mark.parametrize("carry", [True, False])
|
||||
def test_acked_child_events_are_not_reannounced_under_a_new_run_id(tmp_path, monkeypatch, carry):
|
||||
"""After a run-id switch an ACKED settled child and an acked beacon stay
|
||||
delivered, while a child settling after the switch emits exactly once. With the
|
||||
carry removed (the old rebuild) the delivered terminal is re-announced."""
|
||||
if not carry:
|
||||
def _old_rebuild(ctx, run_id): # the pre-fix rebuild, verbatim in effect
|
||||
try:
|
||||
data = json.loads(supervision._state_path(ctx).read_text(encoding="utf-8"))
|
||||
except (OSError, ValueError):
|
||||
data = {}
|
||||
if not isinstance(data, dict) or (str(run_id) and str(data.get("run_id") or "") != str(run_id)):
|
||||
data = {"schema": 1, "run_id": str(run_id), "journal_cursor": 0}
|
||||
return data
|
||||
|
||||
monkeypatch.setattr(supervision, "_load_state", _old_rebuild)
|
||||
_child(tmp_path, "done", status=STATUS_COMPLETED)
|
||||
_child(tmp_path, "late")
|
||||
assert task_tree_ledger.tree_ledger_append(
|
||||
"root", "blocker", "done needs input", task_id="done", data_root=tmp_path,
|
||||
).startswith("OK:")
|
||||
ctx = _ctx(tmp_path)
|
||||
first = json.loads(supervised_wait(ctx, "run-1", wait_once=_settled("run-1")).text)
|
||||
assert [e["child_task_id"] for e in _child_events(first)] == ["done"]
|
||||
assert len(_child_events(first, "child_attention_beacon")) == 1
|
||||
assert acknowledge_pending_wake(ctx, first)
|
||||
|
||||
def _late_settles(*_a, **_k):
|
||||
write_task_result(tmp_path, "late", STATUS_COMPLETED, result="late result")
|
||||
return json.dumps({"status": "completed", "run_id": "run-2", "last_seq": 1})
|
||||
|
||||
second = json.loads(supervised_wait(ctx, "run-2", wait_once=_late_settles).text)
|
||||
announced = [e["child_task_id"] for e in _child_events(second)]
|
||||
if carry:
|
||||
assert announced == ["late"]
|
||||
assert _child_events(second, "child_attention_beacon") == []
|
||||
else:
|
||||
assert "done" in announced # the defect the carry repairs
|
||||
assert acknowledge_pending_wake(ctx, second)
|
||||
if carry:
|
||||
third = json.loads(supervised_wait(ctx, "run-3", wait_once=_settled("run-3")).text)
|
||||
assert _child_events(third) == [] and _child_events(third, "child_attention_beacon") == []
|
||||
|
||||
|
||||
def test_an_unacked_child_terminal_survives_the_run_switch_and_emits_once(tmp_path):
|
||||
_child(tmp_path, "done", status=STATUS_COMPLETED)
|
||||
ctx = _ctx(tmp_path)
|
||||
first = json.loads(supervised_wait(ctx, "run-1", wait_once=_settled("run-1")).text)
|
||||
assert [e["child_task_id"] for e in _child_events(first)] == ["done"]
|
||||
# Never acknowledged: its cursor never committed, so the new run re-detects it once.
|
||||
second = json.loads(supervised_wait(ctx, "run-2", wait_once=_settled("run-2")).text)
|
||||
assert [e["child_task_id"] for e in _child_events(second)] == ["done"]
|
||||
assert acknowledge_pending_wake(ctx, second)
|
||||
third = json.loads(supervised_wait(ctx, "run-3", wait_once=_settled("run-3")).text)
|
||||
assert _child_events(third) == []
|
||||
|
||||
|
||||
@pytest.mark.parametrize("carried", [True, False])
|
||||
def test_recovery_handoff_and_restore_keep_the_committed_cursor(tmp_path, carried):
|
||||
import ouroboros.delegate_recovery as recovery
|
||||
|
||||
cursor = {"attention_after_ts": "2026-09-26T00:00:00+00:00", "attention_seen": ["b1"],
|
||||
"children": {"done": {"status": "completed", "updated_at": "t", "result_sha256": "x"}}}
|
||||
path = supervision._state_path(_ctx(tmp_path))
|
||||
path.parent.mkdir(parents=True, exist_ok=True)
|
||||
path.write_text(json.dumps({"schema": 1, "run_id": "run-old", "coordination_cursor": cursor}),
|
||||
encoding="utf-8")
|
||||
row = {"task_id": "parent", "run_id": "run-new", "journal_cursor": 0}
|
||||
if carried:
|
||||
row["coordination_cursor"] = dict(cursor)
|
||||
recovery._restore_wait_checkpoint(tmp_path, row)
|
||||
restored = json.loads(path.read_text(encoding="utf-8"))
|
||||
assert restored["run_id"] == "run-new" and restored["status"] == "adopted"
|
||||
assert (restored.get("coordination_cursor") == cursor) is carried
|
||||
|
||||
|
||||
def test_prepare_handoff_row_snapshots_the_committed_cursor(tmp_path):
|
||||
from tests.test_delegate_pending_recovery import _pending_handoff
|
||||
|
||||
cursor = {"children": {"done": {"status": "completed"}}}
|
||||
path = tmp_path / "state" / "delegate_supervision" / "t-cursor.json"
|
||||
path.parent.mkdir(parents=True, exist_ok=True)
|
||||
path.write_text(json.dumps({"schema": 1, "run_id": "", "coordination_cursor": cursor}),
|
||||
encoding="utf-8")
|
||||
_custody, recovery, *_rest = _pending_handoff(tmp_path, "t-cursor")
|
||||
assert recovery._read(tmp_path, "t-cursor")["coordination_cursor"] == cursor
|
||||
|
||||
|
||||
# -- B: whole-sleep facts at the single publication point ---------------------
|
||||
|
||||
|
||||
class _Clock:
|
||||
def __init__(self, start=1_000_000.0):
|
||||
self.now = start
|
||||
|
||||
def time(self):
|
||||
return self.now
|
||||
|
||||
|
||||
def _fake_time(monkeypatch, clock):
|
||||
import time as real_time
|
||||
|
||||
monkeypatch.setattr(supervision, "time", SimpleNamespace(
|
||||
time=clock.time, sleep=lambda _s: None, strftime=real_time.strftime, gmtime=real_time.gmtime))
|
||||
|
||||
|
||||
def test_every_wake_measures_the_whole_sleep_and_replays_it_exactly(tmp_path, monkeypatch):
|
||||
clock = _Clock()
|
||||
_fake_time(monkeypatch, clock)
|
||||
ticks = []
|
||||
|
||||
def wait_once(_ctx, run_id, _window, _cursor, **_k):
|
||||
ticks.append(1)
|
||||
clock.now += 3.0
|
||||
if len(ticks) <= 4:
|
||||
return json.dumps({"status": "no_progress", "run_id": run_id, "last_seq": len(ticks),
|
||||
"waited_sec": 3, "note": "or delegate_cancel if it is stuck"})
|
||||
return json.dumps({"status": "completed", "run_id": run_id, "last_seq": 7})
|
||||
|
||||
ctx = _ctx(tmp_path)
|
||||
wake = json.loads(supervised_wait(ctx, "run-1", since_seq=2, wait_once=wait_once).text)
|
||||
sleep = wake["sleep"]
|
||||
assert sleep["slept_sec"] == 15.0 # the whole call, not the last 3 s tick
|
||||
assert sleep["quiet_renewals"] == 4
|
||||
assert sleep["journal_advances"] == 5 # cursor 2 at entry -> 7 at the wake
|
||||
assert "opened_after" not in sleep
|
||||
clock.now += 500.0
|
||||
replay = json.loads(supervised_wait(ctx, "run-1", wait_once=lambda *_a, **_k: (_ for _ in ()).throw(
|
||||
AssertionError("a pending wake replays before any new observation"))).text)
|
||||
assert replay == wake # recorded facts, nothing added since
|
||||
|
||||
|
||||
def test_a_quiet_status_wake_drops_the_tick_advice_but_keeps_the_paused_note(tmp_path, monkeypatch):
|
||||
from ouroboros.delegate_progress import paused_note
|
||||
|
||||
calls: list = []
|
||||
monkeypatch.setattr(supervision, "_control_wakes", lambda _ctx: calls.append(1) or (
|
||||
[{"type": "deadline"}] if len(calls) > 1 else []))
|
||||
tick = {"status": "no_progress", "run_id": "run-q", "last_seq": 1, "waited_sec": 3,
|
||||
"quiet_for_sec": 3, "reason": "non_terminal", "note": "or delegate_cancel if it is stuck"}
|
||||
quiet = json.loads(supervised_wait(_ctx(tmp_path, "t-quiet"), "run-q",
|
||||
wait_once=lambda *_a, **_k: json.dumps(tick)).text)
|
||||
assert quiet["wake_events"] == [{"type": "deadline"}]
|
||||
assert not {"waited_sec", "quiet_for_sec", "note"} & set(quiet)
|
||||
assert quiet["sleep"]["quiet_renewals"] == 0
|
||||
|
||||
pending = [{"interaction_id": "i1", "timeout_at": None}]
|
||||
calls.clear()
|
||||
paused = json.loads(supervised_wait(_ctx(tmp_path, "t-paused"), "run-q", wait_once=lambda *_a, **_k: json.dumps({
|
||||
**tick, "waiting_on_user": True, "pending_interactions": pending})).text)
|
||||
assert paused["note"] == paused_note(pending)
|
||||
assert "delegate_cancel" not in paused["note"] and "waited_sec" not in paused
|
||||
# A terminal wake is not quiet: its own fields stay, and it carries the sleep too.
|
||||
calls.clear()
|
||||
terminal = json.loads(supervised_wait(_ctx(tmp_path, "t-term"), "run-q", wait_once=lambda *_a, **_k: json.dumps({
|
||||
"status": "completed", "run_id": "run-q", "last_seq": 1, "note": "terminal note"})).text)
|
||||
assert terminal["note"] == "terminal note" and "sleep" in terminal
|
||||
|
||||
|
||||
def test_an_old_pending_wake_without_sleep_replays_unchanged(tmp_path):
|
||||
ctx = _ctx(tmp_path)
|
||||
stored = {"status": "completed", "run_id": "run-1", "supervision_wake_id": "w-old"}
|
||||
path = supervision._state_path(ctx)
|
||||
path.parent.mkdir(parents=True, exist_ok=True)
|
||||
path.write_text(json.dumps({"schema": 1, "run_id": "run-1", "pending_wake": {
|
||||
"wake_id": "w-old", "attempt_key": "1", "payload": stored}}), encoding="utf-8")
|
||||
replay = json.loads(supervised_wait(ctx, "run-1", wait_once=_settled("run-1")).text)
|
||||
assert replay == stored # nothing invented for an old record
|
||||
|
||||
|
||||
@pytest.mark.parametrize("prior,label", [
|
||||
("adopted", "worker_loss_adoption"), ("sleeping", "interrupted_sleep"), ("awake", None),
|
||||
])
|
||||
def test_a_new_call_opens_a_new_sleep_labelled_by_the_prior_status(tmp_path, prior, label):
|
||||
ctx = _ctx(tmp_path)
|
||||
path = supervision._state_path(ctx)
|
||||
path.parent.mkdir(parents=True, exist_ok=True)
|
||||
path.write_text(json.dumps({"schema": 1, "run_id": "run-1", "status": prior, "sleep_entry": {
|
||||
"entered_at": "2026-09-26T00:00:00+00:00", "entered_at_unix": 1.0}}), encoding="utf-8")
|
||||
wake = json.loads(supervised_wait(ctx, "run-1", wait_once=_settled("run-1")).text)
|
||||
assert wake["sleep"].get("opened_after") == label
|
||||
assert wake["sleep"]["slept_sec"] < 60 # never carried over from the old entry
|
||||
assert ("previous_entered_at" in wake["sleep"]) is (prior == "sleeping")
|
||||
|
||||
|
||||
def test_a_spilled_wake_keeps_the_sleep_and_live_input_facts(tmp_path, monkeypatch):
|
||||
import ouroboros.tool_capabilities as capabilities
|
||||
|
||||
monkeypatch.setattr(capabilities, "tool_result_limit", lambda _name: 900)
|
||||
ctx = _ctx(tmp_path)
|
||||
ctx.task_contract = {"delegation_budget": {"intent_note": "x" * 5000}}
|
||||
wake = json.loads(supervised_wait(ctx, "run-s", wait_once=_settled("run-s")).text)
|
||||
assert wake["wake_delivery"]["complete"] is False
|
||||
assert wake["sleep"]["quiet_renewals"] == 0 and wake["leaf_live_input"] == "unknown"
|
||||
|
||||
|
||||
# -- D1: dated observation facts (freshness is never the file's mtime) ---------
|
||||
|
||||
|
||||
def test_observation_facts_separate_answered_reads_from_failures(tmp_path, monkeypatch):
|
||||
_fake_time(monkeypatch, _Clock())
|
||||
replies = iter([
|
||||
{"status": "no_progress", "run_id": "run-o", "state": "running", "last_seq": 1},
|
||||
{"status": "observation_pending", "run_id": "run-o", "reason": "daemon_unreachable"},
|
||||
])
|
||||
seen = []
|
||||
|
||||
def wait_once(ctx, *_a, **_k):
|
||||
try:
|
||||
return json.dumps(next(replies))
|
||||
except StopIteration:
|
||||
seen.append(_state(ctx))
|
||||
return json.dumps({"status": "completed", "run_id": "run-o", "last_seq": 2})
|
||||
|
||||
ctx = _ctx(tmp_path)
|
||||
supervised_wait(ctx, "run-o", wait_once=wait_once)
|
||||
before_wake = seen[0]
|
||||
assert before_wake["observation_failure"] == "daemon_unreachable"
|
||||
assert before_wake["observation"]["answered"] is False
|
||||
answered_at = before_wake["last_answered_observation_at"]
|
||||
assert answered_at # the earlier ANSWERED read, not advanced
|
||||
assert "state_written_at" not in before_wake
|
||||
after = _state(ctx)
|
||||
assert after["observation_failure"] == "" and after["last_answered_observation_at"] >= answered_at
|
||||
|
||||
|
||||
# -- leaf_live_input: read once at entry, stamped on every wake ----------------
|
||||
|
||||
|
||||
class _Catalog:
|
||||
def __init__(self, rows=None, fail=False):
|
||||
self.rows, self.fail, self.reads, self.closed = rows or [], fail, [], False
|
||||
|
||||
def agent_capabilities(self, *, timeout_sec=None):
|
||||
self.reads.append(timeout_sec)
|
||||
if self.fail:
|
||||
raise RuntimeError("catalog unreadable")
|
||||
return {"harnesses": self.rows}
|
||||
|
||||
def close(self):
|
||||
self.closed = True
|
||||
|
||||
|
||||
def _owned(tmp_path, monkeypatch, run_id="run-l", route="fixture"):
|
||||
from ouroboros import delegate_custody as custody
|
||||
|
||||
monkeypatch.setitem(custody._CUSTODY, run_id, custody.RunCustody(
|
||||
run_id=run_id, task_id="parent", route_id=route, model="m"))
|
||||
|
||||
|
||||
@pytest.mark.parametrize("rows,fail,expected", [
|
||||
([{"id": "fixture", "liveInput": "mid_turn"}], False, "mid_turn"),
|
||||
([{"id": "fixture"}], False, "none"), # an engine without the field
|
||||
([{"id": "other", "liveInput": "mid_turn"}], False, "unknown"),
|
||||
([], True, "unknown"), # failure is a fact, never a refusal
|
||||
])
|
||||
def test_leaf_live_input_is_read_once_at_entry_and_stamped_on_every_wake(
|
||||
tmp_path, monkeypatch, rows, fail, expected,
|
||||
):
|
||||
_owned(tmp_path, monkeypatch)
|
||||
catalog = _Catalog(rows, fail)
|
||||
monkeypatch.setattr(supervision, "_loop_gateway", lambda: catalog)
|
||||
monkeypatch.setattr(supervision.time, "sleep", lambda _s: None)
|
||||
ticks = []
|
||||
|
||||
def wait_once(_ctx, run_id, *_a, **_k):
|
||||
ticks.append(1)
|
||||
status = "no_progress" if len(ticks) < 4 else "completed"
|
||||
return json.dumps({"status": status, "run_id": run_id, "last_seq": len(ticks)})
|
||||
|
||||
ctx = _ctx(tmp_path)
|
||||
wake = json.loads(supervised_wait(ctx, "run-l", wait_once=wait_once).text)
|
||||
assert wake["leaf_live_input"] == expected
|
||||
assert catalog.reads == [supervision._TICK_SEC] # once per call, never per tick
|
||||
assert catalog.closed # the one-read transport is dropped
|
||||
replay = json.loads(supervised_wait(ctx, "run-l", wait_once=wait_once).text)
|
||||
assert replay["leaf_live_input"] == expected and catalog.reads == [supervision._TICK_SEC]
|
||||
|
||||
|
||||
def test_an_unowned_run_reads_no_catalog(tmp_path, monkeypatch):
|
||||
monkeypatch.setattr(supervision, "_loop_gateway", lambda: (_ for _ in ()).throw(
|
||||
AssertionError("no route, no catalog read")))
|
||||
wake = json.loads(supervised_wait(_ctx(tmp_path), "run-x", wait_once=_settled("run-x")).text)
|
||||
assert wake["leaf_live_input"] == "unknown"
|
||||
|
|
@ -904,3 +904,47 @@ def test_burn_phrase_never_claims_zero_rounds_with_real_dollars():
|
|||
assert "0 of your own metered LLM rounds" not in _nanny_burn_phrase(0, 2.45)
|
||||
assert "$2.45" in _nanny_burn_phrase(0, 2.45)
|
||||
assert "3 of your own metered LLM rounds" in _nanny_burn_phrase(3, 2.45)
|
||||
|
||||
|
||||
# -- the money sentence follows the round's own cost evidence ------------------
|
||||
|
||||
|
||||
def _fire_reminder(cost_per_round, meta, monkeypatch, fallbacks=""):
|
||||
from ouroboros.loop import _maybe_inject_nanny_economics_reminder, _note_nanny_delegate_activity
|
||||
|
||||
monkeypatch.setenv("OUROBOROS_MODEL_FALLBACKS", fallbacks)
|
||||
monkeypatch.delenv("OUROBOROS_MODEL_FALLBACK", raising=False)
|
||||
ctx = _nanny_ctx(_accumulated_usage={"_last_llm_call_meta": meta})
|
||||
tools = SimpleNamespace(_ctx=ctx)
|
||||
_note_nanny_delegate_activity(ctx, 1, {"cost": 0.0}, [_delegate_call()])
|
||||
round_idx, msgs = 1, []
|
||||
while not msgs:
|
||||
round_idx += 1
|
||||
_note_nanny_delegate_activity(ctx, round_idx, {"cost": cost_per_round * round_idx}, [])
|
||||
_maybe_inject_nanny_economics_reminder(round_idx, msgs, tools, lambda *_: None)
|
||||
return "\n".join(m.get("content", "") for m in msgs)
|
||||
|
||||
|
||||
def test_an_unpriced_round_is_an_unknown_cash_cost_never_metered_or_zero(monkeypatch):
|
||||
text = _fire_reminder(0.0, {"provider": "claudexor", "model": "codex=gpt"}, monkeypatch)
|
||||
assert "metered API money" not in text and "metered LLM rounds" not in text
|
||||
assert "cash cost unknown, not zero" in text # burn phrase, cost class 'unpriced'
|
||||
assert "claudexor / codex=gpt" in text # the route, named as context only
|
||||
assert "only when its settled spend reports $0" in text # the pinned conditional stays
|
||||
assert "fallback" not in text # no chain configured, no claim
|
||||
|
||||
|
||||
def test_a_priced_round_stays_metered_and_names_a_configured_fallback(monkeypatch):
|
||||
text = _fire_reminder(0.05, {"provider": "openrouter", "model": "anthropic/claude"},
|
||||
monkeypatch, fallbacks="openai/gpt-x, google/gem")
|
||||
assert "metered LLM rounds" in text and "that spend is metered money" in text
|
||||
assert "cash cost unknown" not in text
|
||||
assert "configured fallback (openai/gpt-x, google/gem)" in text
|
||||
|
||||
|
||||
def test_the_unpriced_burn_phrase_never_calls_rounds_metered():
|
||||
from ouroboros.loop import _nanny_burn_phrase
|
||||
|
||||
assert _nanny_burn_phrase(4, 0.0) == (
|
||||
"4 of your own LLM rounds (no provider price reported: cash cost unknown, not zero)")
|
||||
assert "4 of your own metered LLM rounds" in _nanny_burn_phrase(4, 0.5)
|
||||
|
|
|
|||
|
|
@ -294,6 +294,9 @@ CHAPTER_BYTE_BUDGETS: dict[str, int] = {
|
|||
# delegated-lane bullet for the live-message verb (capability gate, mirrored outcomes,
|
||||
# message_id custody, no retry loop); the base sat 104 bytes under the previous budget.
|
||||
"docs/development/06-rules-by-change-class.md": 97200,
|
||||
# 96700 -> 96800 (truthful waiting A-C, measured 96791): one Timeout & Wait Control bullet —
|
||||
# wake facts measured over the interval the actor experienced; no older text to displace.
|
||||
"docs/development/06-rules-by-change-class.md": 96800,
|
||||
"docs/development/07-managed-update-rule.md": 4166,
|
||||
"docs/development/08-mutation-attribution-rule.md": 2899,
|
||||
"docs/development/09-process-custody-rule.md": 10028,
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue