ouroboros/supervisor/state.py

1418 lines
69 KiB
Python

"""Supervisor persistent state, atomic writes, locks, and budget accounting."""
from __future__ import annotations
import contextlib
import dataclasses
import json
import logging
import os
import pathlib
import time
import uuid
from typing import Any, Dict, List, Optional, Tuple
from ouroboros.config import DATA_DIR
from ouroboros.contracts.schema_versions import SCHEMA_VERSION_KEY
from ouroboros.platform_layer import acquire_exclusive_file_lock, release_exclusive_file_lock
from ouroboros.utils import append_jsonl, assert_test_data_path, utc_now_iso, write_bytes_atomic # noqa: F401 -- assert_test_data_path re-exported: the guard moved to the utils leaf so append_jsonl shares it; state.assert_test_data_path callers keep working
log = logging.getLogger(__name__)
# ABI 7.0 (Q8=B): the durable ``state.json`` snapshot names its schema on
# every write. Stamp-on-write ONLY — readers do not require the stamp (no
# compat branching), so a pre-7.0 file and a post-rollback stamped file both
# load unchanged.
STATE_SCHEMA_VERSION = 1
DRIVE_ROOT: pathlib.Path = pathlib.Path(DATA_DIR)
STATE_PATH: pathlib.Path = DRIVE_ROOT / "state" / "state.json"
STATE_LAST_GOOD_PATH: pathlib.Path = DRIVE_ROOT / "state" / "state.last_good.json"
STATE_LOCK_PATH: pathlib.Path = DRIVE_ROOT / "locks" / "state.lock"
# Explicit marker a benchmark/evolution driver writes into its THROWAWAY data root. A live
# data root (the default ~/Ouroboros/data OR a custom/Drive-backed OUROBOROS_DATA_DIR) never
# has it, so reset_per_task_budget can refuse on it regardless of how the path resolves —
# closing the budget-reset guard for custom-data-root installs (BIBLE P8).
ISOLATED_BENCHMARK_SENTINEL = ".ouroboros_isolated_benchmark"
def init(drive_root: pathlib.Path, total_budget_limit: float = 0.0) -> None:
global DRIVE_ROOT, STATE_PATH, STATE_LAST_GOOD_PATH, STATE_LOCK_PATH
DRIVE_ROOT = drive_root
STATE_PATH = drive_root / "state" / "state.json"
STATE_LAST_GOOD_PATH = drive_root / "state" / "state.last_good.json"
STATE_LOCK_PATH = drive_root / "locks" / "state.lock"
set_budget_limit(total_budget_limit)
def atomic_write_text(path: pathlib.Path, content: str) -> None:
"""Durable state write: byte-exact, fsync'd, every byte landed.
Rides the utils atomic SSOT — its write loop survives a short ``os.write``
(the old single call could publish a truncated ``state.json`` behind a
successful rename), and the pytest live-data guard now sits on that same
seam, so every writer through it is guarded rather than this one alone.
"""
write_bytes_atomic(path, content.encode("utf-8"), fsync=True)
def json_load_file(path: pathlib.Path) -> Optional[Dict[str, Any]]:
"""Legacy best-effort read (None for every failure). Authority reads use ``read_state``."""
status, obj, _raw, _detail = _read_json_file(path)
return obj if status == "ok" else None
def read_state_copy(path: pathlib.Path) -> Tuple[str, Optional[Dict[str, Any]], str]:
"""``(status, object, detail)`` of one state-shaped file, never written (display readers)."""
status, obj, _raw, detail = _read_json_file(path)
return status, obj, detail
def _read_json_file(path: pathlib.Path) -> Tuple[str, Optional[Dict[str, Any]], bytes, str]:
"""``(status, object, raw bytes, detail)`` for one state copy, classified by the
operation itself (#1307): ``missing`` only for a not-found below a real directory
(``confirm_absent``), ``unreadable`` for any other OSError (EACCES, ENFILE, EIO, and
ENOTDIR — a file where ``state/`` belongs is not absence, on Windows too),
``invalid`` for bytes that are not a non-empty JSON object. No ``exists()``
pre-check: it can fail the same way."""
from supervisor.state_initialization import confirm_absent
try:
try:
raw = pathlib.Path(path).read_bytes()
except FileNotFoundError:
confirm_absent(path)
return "missing", None, b"", ""
except OSError as exc:
return "unreadable", None, b"", f"{type(exc).__name__} errno={exc.errno}"
try:
obj = json.loads(raw.decode("utf-8"))
except (UnicodeDecodeError, ValueError) as exc:
return "invalid", None, raw, type(exc).__name__
if not isinstance(obj, dict) or not obj:
return "invalid", None, raw, "not a non-empty JSON object"
return "ok", obj, raw, ""
def acquire_file_lock(lock_path: pathlib.Path, timeout_sec: float = 4.0,
stale_sec: float = 90.0) -> Optional[int]:
return acquire_exclusive_file_lock(
lock_path,
timeout_sec=timeout_sec,
stale_sec=stale_sec,
metadata=f"pid={os.getpid()} ts={utc_now_iso()}\n",
owner_aware_stale=True,
)
# Direct alias: the platform helper already has the exact signature.
release_file_lock = release_exclusive_file_lock
def ensure_state_defaults(st: Dict[str, Any]) -> Dict[str, Any]:
st.setdefault("created_at", utc_now_iso())
st.setdefault("owner_id", None)
st.setdefault("owner_chat_id", None)
# Separate slot authorizing owner slash commands from external transports
# (e.g. Telegram), so the local web owner never locks out a real chat owner.
st.setdefault("owner_external_id", None)
st.setdefault("owner_external_chat_id", None)
st.setdefault("owner_external_bound_at", None)
st.setdefault("message_offset", 0)
if "tg_offset" in st:
st.setdefault("message_offset", st.pop("tg_offset"))
st.setdefault("spent_usd", 0.0)
st.setdefault("spent_calls", 0)
st.setdefault("spent_tokens_prompt", 0)
st.setdefault("spent_tokens_completion", 0)
st.setdefault("spent_tokens_cached", 0)
st.setdefault("session_id", uuid.uuid4().hex)
st.setdefault("current_branch", None)
st.setdefault("current_sha", None)
st.setdefault("last_owner_message_at", "")
st.setdefault("last_evolution_task_at", "")
st.setdefault("budget_messages_since_report", 0)
st.setdefault("evolution_mode_enabled", False)
# Durable owner-stop sentinel: set True by the owner-stop sites, cleared by an
# owner-authorized start (/evolve start or the owner-directed toggle_evolution(True)
# tool). apply_pending_request refuses to autonomously re-arm while True.
st.setdefault("evolution_owner_stopped", False)
st.setdefault("evolution_cycle", 0)
st.setdefault("session_total_snapshot", None)
st.setdefault("session_spent_snapshot", None)
# Drift compares like with like: the OpenRouter-only settled ledger total vs
# the queried OpenRouter key's usage. The all-provider spent_usd delta is NOT
# comparable once direct-provider lanes (e.g. Anthropic advisory) carry real
# spend — that shape kept budget_drift_alert latched at ~88% while nothing
# was wrong with the ledger.
st.setdefault("session_openrouter_settled_snapshot", None)
st.setdefault("session_openrouter_key_fp", "")
st.setdefault("openrouter_ledger_settled_usd", None)
st.setdefault("budget_drift_pct", None)
st.setdefault("budget_drift_alert", False)
st.setdefault("evolution_consecutive_failures", 0)
st.setdefault("bg_consciousness_enabled", False)
for legacy_key in ("approvals", "idle_cursor", "idle_stats", "last_idle_task_at",
"last_auto_review_at", "last_review_task_id", "session_daily_snapshot"):
st.pop(legacy_key, None)
return st
# --- #1307: typed read quality, writer-owned recovery, one explicit initializer ---
#
# ``state.json`` is a small service file (owner binding, evolution/consciousness
# controls, session identity, legacy money projection). Absent, corrupt and
# unreadable-right-now are different facts: only the explicit ``init_state`` may
# create a first state, a readable backup restores DATA but never current control
# authority (its controls stay ``unconfirmed`` in the writer-owned ``_recovery``
# block until an actual decision confirms each one; only a set owner binding of the
# same initialization identity is proven), and an unreadable primary is never
# overwritten. Money never reads this file (ledger authority).
RECOVERY_KEY = "_recovery"
STATE_READ_KEY = "_state_read" # projection-only read quality; never persisted
CONTROL_KEYS = (
"evolution_mode_enabled", "evolution_owner_stopped", "evolution_stop_source",
"post_task_autostop", "bg_consciousness_enabled",
"owner_id", "owner_chat_id", "owner_external_id", "owner_external_chat_id",
)
OPTIONAL_CONTROL_KEYS = frozenset({"evolution_stop_source", "post_task_autostop"})
CURRENT_QUALITIES = frozenset({"current", "recovered"})
class StateUnavailable(RuntimeError):
"""State authority or persistence is unavailable; partial writes are explicit."""
def __init__(self, reason: str, detail: str = "", *, primary_written: bool = False) -> None:
super().__init__(f"state unavailable: {reason}" + (f" ({detail})" if detail else ""))
self.reason = reason
self.primary_written = primary_written
@dataclasses.dataclass(frozen=True)
class StateRead:
"""One state observation: its values WITH their source and quality.
``current``: the primary copy, no pending recovery. ``recovered``: the primary
after a backup recovery; ``unconfirmed`` controls are unknown. ``recovered_transient``:
the primary cannot be read right now; backup values are display-only and every
control but a proven owner binding is unknown. ``uninitialized``/``unavailable``: no values."""
quality: str
source: str
values: Dict[str, Any]
unconfirmed: Tuple[str, ...] = ()
reason: str = ""
def projection(self) -> Dict[str, Any]:
"""The legacy dict view for DISPLAY readers; authority uses ``control_value``."""
out = dict(self.values)
if self.quality != "current":
out[STATE_READ_KEY] = {"quality": self.quality, "source": self.source,
"unconfirmed": list(self.unconfirmed), "reason": self.reason}
return out
def control_value(st: Dict[str, Any], key: str) -> Tuple[bool, Any]:
"""``(known, value)`` of one control in a state dict (a projection or a live
mutator dict). Unknown when the read was not current or the key awaits
confirmation after a recovery: a missing fact is never a default."""
meta = st.get(STATE_READ_KEY) if isinstance(st, dict) else None
if isinstance(meta, dict) and (meta.get("quality") not in CURRENT_QUALITIES | {"recovered_transient"}
or key in (meta.get("unconfirmed") or ())):
return False, None
recovery = st.get(RECOVERY_KEY) if isinstance(st, dict) else None
if isinstance(recovery, dict) and key in (recovery.get("unconfirmed") or ()):
return False, None
return (True, st.get(key)) if isinstance(st, dict) and (key in st or key in OPTIONAL_CONTROL_KEYS) else (False, None)
def control_in_copy(path: pathlib.Path, key: str) -> Tuple[bool, Any]:
"""``control_value`` of one control in the primary copy at ``path`` (a lock-free read
for a caller that addresses a drive by path): unknown unless that copy is readable."""
status, obj, _raw, _detail = _read_json_file(path)
from supervisor.state_initialization import authority_reason
if status != "ok" or authority_reason(pathlib.Path(path).parent.parent, str(obj.get("initialization_id") or "")):
return False, None
return control_value(obj, key)
def mark_unconfirmed(live: Dict[str, Any], key: str) -> None:
"""Inside an ``update_state`` mutator: restore a control to unknown (a failed
decision puts back what it could not prove)."""
recovery = live.setdefault(RECOVERY_KEY, {"source": "restored_unknown"})
if key not in (recovery.setdefault("unconfirmed", [])):
recovery["unconfirmed"].append(key)
def control_is(st: Dict[str, Any], key: str, expected: Any) -> bool:
"""True only when the control is KNOWN to equal ``expected``."""
known, value = control_value(st, key)
return known and value == expected
def _backup_unconfirmed(backup: Dict[str, Any], drive_root=None) -> Tuple[str, ...]:
"""The controls a backup copy cannot prove. A SET owner binding is proven when the
backup carries the completed initialization identity: its only writers fill a
known-empty slot and only an owner Reset (a new identity, both copies deleted)
clears one, so within an identity a set binding never changes. Switches and an
empty binding slot may have changed after the backup was written: unknown."""
from supervisor import state_initialization as witness
identity = str(backup.get("initialization_id") or "")
status, record = witness.read_witness(drive_root or DRIVE_ROOT) if identity else ("missing", {})
same = status == "ok" and record.get("phase") == "complete" and record.get("initialization_id") == identity
return tuple(key for key in CONTROL_KEYS
if not (same and key.startswith("owner_") and backup.get(key) is not None
and key not in _recovery_unconfirmed(backup)))
def _recovery_unconfirmed(st: Dict[str, Any]) -> Tuple[str, ...]:
recovery = st.get(RECOVERY_KEY)
return tuple(recovery.get("unconfirmed") or ()) if isinstance(recovery, dict) else ()
def read_state(drive_root=None) -> StateRead:
"""Classify both copies WITHOUT writing (a GET, a display, a boot probe)."""
from supervisor.state_initialization import authority_reason, read_witness
root = pathlib.Path(drive_root) if drive_root is not None else DRIVE_ROOT
p_status, primary, _raw, p_detail = _read_json_file(root / "state" / "state.json")
if p_status == "ok":
reason = authority_reason(root, str(primary.get("initialization_id") or ""))
if reason:
return StateRead("unavailable", "primary", dict(primary), CONTROL_KEYS, reason)
unconfirmed = tuple(set(_recovery_unconfirmed(primary)) | (set(CONTROL_KEYS) - primary.keys() - OPTIONAL_CONTROL_KEYS))
return StateRead("recovered" if unconfirmed else "current", "primary",
ensure_state_defaults(dict(primary)), unconfirmed)
b_status, backup, _braw, b_detail = _read_json_file(root / "state" / "state.last_good.json")
reason = f"primary {p_status}{f' ({p_detail})' if p_detail else ''}; backup {b_status}"
if b_status == "ok":
return StateRead("recovered_transient", "backup", ensure_state_defaults(dict(backup)),
_backup_unconfirmed(backup, root), reason)
w_status, _witness = read_witness(root)
quality = "uninitialized" if p_status == b_status == w_status == "missing" else "unavailable"
return StateRead(quality, "none", {}, CONTROL_KEYS, reason + (f" ({b_detail})" if b_detail else ""))
def _preserve_corrupt_primary(raw: bytes) -> str:
"""Keep the damaged primary's bytes under an exclusive new name before any
recovery write replaces them; refusal is typed, never a silent overwrite."""
stamp = time.strftime("%Y%m%dT%H%M%S", time.gmtime())
for suffix in range(100):
target = STATE_PATH.with_name(f"state.corrupt-{stamp}{f'-{suffix}' if suffix else ''}.json")
try:
fd = os.open(str(target), os.O_WRONLY | os.O_CREAT | os.O_EXCL | getattr(os, "O_BINARY", 0), 0o600)
except FileExistsError:
continue
except OSError as exc:
raise StateUnavailable("corrupt_primary_unpreserved", f"{type(exc).__name__}") from exc
try:
offset = 0
while offset < len(raw):
written = os.write(fd, raw[offset:])
if written <= 0:
raise OSError("corrupt copy made no progress")
offset += written
os.fsync(fd)
if os.fstat(fd).st_size != len(raw):
raise OSError("corrupt copy length mismatch")
except OSError as exc:
raise StateUnavailable("corrupt_primary_unpreserved", type(exc).__name__) from exc
finally:
os.close(fd)
return target.name
raise StateUnavailable("corrupt_primary_unpreserved", "no free name")
def _state_for_write(*, initializing: bool = False) -> Dict[str, Any]:
"""The CURRENT dict a locked writer may mutate (caller holds STATE_LOCK).
A readable primary is returned as is. A missing/invalid primary with a readable
backup is durably RECOVERED first: the invalid bytes are preserved, the backup's
values become the primary with its unprovable controls ``unconfirmed`` (#1144: the ledger
freshness marker is dropped so money re-derives). Anything else raises."""
p_status, primary, raw, p_detail = _read_json_file(STATE_PATH)
if p_status == "ok":
from supervisor.state_initialization import authority_reason
reason = authority_reason(DRIVE_ROOT, str(primary.get("initialization_id") or ""))
if reason and not initializing:
# Ordinary bookkeeping may retain legacy data; it cannot adopt or grant.
if reason == "initialization_witness_missing" and not primary.get("initialization_id"):
for key in CONTROL_KEYS:
mark_unconfirmed(primary, key)
else:
raise StateUnavailable(reason)
for key in set(CONTROL_KEYS) - primary.keys() - OPTIONAL_CONTROL_KEYS:
mark_unconfirmed(primary, key)
return primary
if p_status == "unreadable":
raise StateUnavailable("primary_unreadable", p_detail)
b_status, backup, _braw, b_detail = _read_json_file(STATE_LAST_GOOD_PATH)
if b_status != "ok":
raise StateUnavailable("uninitialized" if p_status == b_status == "missing" else "no_readable_copy",
f"primary {p_status}; backup {b_status} {b_detail}".strip())
unconfirmed = list(_backup_unconfirmed(backup))
recovery: Dict[str, Any] = {"source": "backup", "primary": p_status, "recovered_at": utc_now_iso(),
"unconfirmed": unconfirmed}
if p_status == "invalid":
recovery["corrupt_copy"] = _preserve_corrupt_primary(raw)
recovered = {key: value for key, value in backup.items() if key != "usage_ledger_high_water_seq"}
recovered[RECOVERY_KEY] = recovery
log.error("state.json %s; recovered values from state.last_good.json with %s unconfirmed",
p_status, ", ".join(unconfirmed))
_save_state_unlocked(recovered)
append_jsonl(DRIVE_ROOT / "logs" / "events.jsonl", {
"ts": utc_now_iso(), "type": "state_recovered_from_backup", "primary": p_status,
"unconfirmed": unconfirmed, **({"corrupt_copy": recovery["corrupt_copy"]}
if "corrupt_copy" in recovery else {})})
return recovered
def _load_state_unlocked() -> Dict[str, Any]:
"""Locked CURRENT dict for the legacy in-module writers (raises when unavailable)."""
return ensure_state_defaults(_state_for_write())
def _save_state_unlocked(st: Dict[str, Any]) -> None:
"""Save state; caller must hold STATE_LOCK."""
st.pop(STATE_READ_KEY, None)
st = ensure_state_defaults(st)
st[SCHEMA_VERSION_KEY] = STATE_SCHEMA_VERSION
payload = json.dumps(st, ensure_ascii=False, indent=2)
primary_written = False
try:
if atomic_write_text(STATE_PATH, payload) is False:
raise OSError("primary writer returned False")
primary_written = True
if atomic_write_text(STATE_LAST_GOOD_PATH, payload) is False:
raise OSError("backup writer returned False")
except OSError as exc:
raise StateUnavailable("backup_write_failed" if primary_written else "primary_write_failed",
f"{type(exc).__name__} errno={exc.errno}",
primary_written=primary_written) from exc
@contextlib.contextmanager
def _state_lock(op: str, timeout_sec: float = 4.0):
"""STATE_LOCK or a typed refusal: a writer never proceeds unlocked (#1307)."""
assert_test_data_path(STATE_PATH)
try:
lock_fd = acquire_file_lock(STATE_LOCK_PATH, timeout_sec=timeout_sec)
except OSError as exc:
raise StateUnavailable("lock_unavailable", f"{type(exc).__name__} errno={exc.errno}") from exc
if lock_fd is None:
log.error("state.json %s refused: lock timeout on %s", op, STATE_LOCK_PATH)
raise StateUnavailable("lock_timeout", op)
try:
yield
finally:
release_file_lock(STATE_LOCK_PATH, lock_fd)
def load_state() -> Dict[str, Any]:
"""Display projection of ``read_state``: never writes, never mints defaults for a
missing/unreadable file; a non-current read carries ``_state_read``."""
assert_test_data_path(STATE_PATH)
return read_state().projection()
def save_state(st: Dict[str, Any]) -> None:
"""Author one WHOLE state (fixtures and isolated tooling). Production changes
fields through ``update_state``. It never overwrites an unreadable primary,
preserves an invalid one's bytes, and keeps the writer-owned ``_recovery``."""
from supervisor import state_initialization as witness
with _state_lock("save"):
p_status, primary, raw, p_detail = _read_json_file(STATE_PATH)
if p_status == "unreadable":
raise StateUnavailable("primary_unreadable", p_detail)
st = {key: value for key, value in st.items() if key not in (RECOVERY_KEY, STATE_READ_KEY)}
created = ""
if p_status == "missing":
# A whole-state write mints identity ONLY where ``init_state`` would: never
# over a lost initialized state, a history, or a recoverable backup.
b_status, _backup, _braw, b_detail = _read_json_file(STATE_LAST_GOOD_PATH)
if b_status != "missing":
raise StateUnavailable("primary_missing", f"backup {b_status} {b_detail}".strip())
decision = witness.initialization_decision(DRIVE_ROOT)
if not decision.get("create"):
raise StateUnavailable(str(decision.get("reason") or "refused"), str(decision.get("detail") or ""))
created = st["initialization_id"] = str(decision["initialization_id"])
if p_status == "invalid":
primary = _state_for_write()
p_status = "ok"
if p_status == "ok":
for key in (set(CONTROL_KEYS) - primary.keys() - OPTIONAL_CONTROL_KEYS) | (set(CONTROL_KEYS) - st.keys() - OPTIONAL_CONTROL_KEYS):
mark_unconfirmed(primary, key)
if p_status == "ok" and isinstance(primary.get(RECOVERY_KEY), dict):
st[RECOVERY_KEY] = primary[RECOVERY_KEY]
if p_status == "ok":
try:
st["initialization_id"] = witness.prepare_adoption(DRIVE_ROOT, primary)
except ValueError as exc:
raise StateUnavailable(str(exc)) from exc
_save_state_unlocked(st)
if not witness.complete(DRIVE_ROOT, st["initialization_id"], adopted=not created):
raise StateUnavailable("initialization_incomplete", "the witness could not be completed")
def update_state(mutator, *, confirm: Tuple[str, ...] = (), lock_timeout_sec: float = 4.0) -> Dict[str, Any]:
"""Atomically read-modify-write state under a single held lock.
Rereads CURRENT under STATE_LOCK, applies ``mutator(st)`` in place, and persists
the result while holding the lock for the WHOLE operation, so concurrent updates
cannot lose each other. Returns the saved state.
The ``_recovery`` block is writer-owned: a mutator cannot clear it, and only the
control keys a real decision names in ``confirm`` leave its ``unconfirmed`` list
(bookkeeping never launders a backup into authority). Lock timeout, an
unreadable primary or a missing state raise ``StateUnavailable`` with nothing
written — bounded, typed, never an unlocked write.
``mutator`` must NOT call ``load_state``/``save_state``/``update_state`` itself:
STATE_LOCK is not re-entrant within a process, so re-entering would block.
"""
with _state_lock("update", lock_timeout_sec):
st = ensure_state_defaults(_state_for_write())
before = _recovery_unconfirmed(st)
recovery = st.get(RECOVERY_KEY)
mutator(st)
# A mutator may ADD an unknown (``mark_unconfirmed``), never remove one.
unconfirmed = [key for key in before if key not in set(confirm)] + [
key for key in _recovery_unconfirmed(st) if key not in before]
if unconfirmed:
st[RECOVERY_KEY] = {**(recovery if isinstance(recovery, dict) else {}), "unconfirmed": unconfirmed}
else:
st.pop(RECOVERY_KEY, None)
_save_state_unlocked(st)
return st
def _set_aside_copies_older_than_a_reset(witness: Any) -> None:
"""An owner Reset's ``pending`` witness is the owner's explicit fresh start: a state
copy of another identity (a writer that raced the Reset's delete) is moved aside
under a new name — never adopted as current, never deleted. Caller holds STATE_LOCK."""
w_status, record = witness.read_witness(DRIVE_ROOT)
if not (w_status == "ok" and record.get("phase") == "pending" and record.get("origin") == "owner_reset"):
return
identity, stamp = str(record.get("initialization_id") or ""), time.strftime("%Y%m%dT%H%M%S", time.gmtime())
for path in (STATE_PATH, STATE_LAST_GOOD_PATH):
status, obj, _raw, detail = _read_json_file(path)
if status == "missing" or (status == "ok" and obj.get("initialization_id") == identity):
continue
if status == "unreadable":
raise StateUnavailable("primary_unreadable", detail)
target = path.with_name(f"{path.stem}.pre-reset-{stamp}-{uuid.uuid4().hex[:6]}.json")
os.replace(path, target)
log.warning("owner reset pending: %s of another identity set aside as %s", path.name, target.name)
def _recover_stopped_controls(st: Dict[str, Any]) -> None:
"""An existing owner Stop intent proves disabled controls, never an enable grant."""
status, campaign, _raw, _detail = _read_json_file(DRIVE_ROOT / "state" / "evolution_campaign.json")
intent = campaign.get("stop_intent") if status == "ok" else None
if not isinstance(intent, dict) or intent.get("source") not in {"owner", "owner_chat", "panic"}:
return
facts = {"evolution_mode_enabled": False, "evolution_owner_stopped": True,
"evolution_stop_source": None, "post_task_autostop": False}
st.update(facts)
recovery = st.get(RECOVERY_KEY)
if isinstance(recovery, dict):
recovery["unconfirmed"] = [key for key in recovery.get("unconfirmed", []) if key not in facts]
def init_state(*, origin: str = "first_boot") -> StateRead:
"""The ONE explicit state initializer, run by supervisor boot before any
owner registration, autonomy admission or chat ingress.
A readable (or backup-recoverable) state is adopted and its initialization
witness completed. Both copies absent create a first state ONLY on positive
evidence (``state_initialization``); otherwise the answer is ``unavailable``
and nothing is minted — the supervisor keeps serving independent work.
Money/network facts for drift detection are read before taking the lock."""
from supervisor import state_initialization as witness
or_settled = _openrouter_ledger_settled()
key_fp = _openrouter_key_fingerprint()
ground_truth = check_openrouter_ground_truth()
try:
with _state_lock("init"):
created = ""
_set_aside_copies_older_than_a_reset(witness)
try:
st = _state_for_write(initializing=True)
except StateUnavailable as exc:
if exc.reason != "uninitialized":
raise
decision = witness.initialization_decision(DRIVE_ROOT, origin=origin)
if not decision.get("create"):
raise StateUnavailable(str(decision.get("reason") or "refused"),
str(decision.get("detail") or ""))
created = str(decision["initialization_id"])
st = ensure_state_defaults({"initialization_id": created})
if not created:
try:
st["initialization_id"] = witness.prepare_adoption(DRIVE_ROOT, st)
except ValueError as exc:
raise StateUnavailable(str(exc)) from exc
for key in set(CONTROL_KEYS) - st.keys() - OPTIONAL_CONTROL_KEYS:
mark_unconfirmed(st, key)
st = ensure_state_defaults(st)
_recover_stopped_controls(st)
st["session_spent_snapshot"] = float(st.get("spent_usd") or 0.0)
st["session_openrouter_settled_snapshot"] = or_settled
st["openrouter_ledger_settled_usd"] = or_settled
st["session_openrouter_key_fp"] = key_fp
if ground_truth is not None:
st["session_total_snapshot"] = ground_truth["total_usd"]
st["openrouter_total_usd"] = ground_truth["total_usd"]
st["openrouter_daily_usd"] = ground_truth["daily_usd"]
st["openrouter_last_check_at"] = utc_now_iso()
else:
st["session_total_snapshot"] = 0.0
st["budget_drift_pct"] = None
st["budget_drift_alert"] = False
_save_state_unlocked(st)
if not witness.complete(DRIVE_ROOT, st["initialization_id"], adopted=not created):
raise StateUnavailable("initialization_incomplete", "the exact witness did not complete")
except (StateUnavailable, OSError) as exc: # a failed witness/set-aside write is typed too, never fatal
log.error("State initialization refused: %s", exc)
try:
append_jsonl(DRIVE_ROOT / "logs" / "events.jsonl", {
"ts": utc_now_iso(), "type": "state_unavailable_at_boot",
"reason": getattr(exc, "reason", type(exc).__name__), "detail": str(exc)})
except Exception:
log.warning("state_unavailable_at_boot could not be recorded", exc_info=True)
read = read_state()
return StateRead("unavailable" if read.quality in CURRENT_QUALITIES else read.quality,
read.source, read.values, CONTROL_KEYS, str(exc))
return read_state()
TOTAL_BUDGET_LIMIT: float = 0.0
EVOLUTION_BUDGET_RESERVE: float = 2.0 # Stop evolution when remaining < this
def set_budget_limit(limit: float) -> None:
"""Set total budget limit for budget_pct."""
global TOTAL_BUDGET_LIMIT
TOTAL_BUDGET_LIMIT = limit
def refresh_budget_from_settings(settings: Dict[str, Any]) -> None:
"""Hot-reload TOTAL_BUDGET; bad/missing values mean no limit."""
try:
raw = settings.get("TOTAL_BUDGET")
value = float(raw) if raw is not None else 0.0
set_budget_limit(value)
except (TypeError, ValueError):
pass
def budget_remaining(
st: Dict[str, Any],
*,
strict: bool = False,
projection: Optional[Dict[str, Any]] = None,
allow_stale: bool = False,
refuse_below: float = 0.0,
) -> float:
"""Return ledger-derived remaining budget in USD.
``state.json`` is only a compatibility projection. A corrupt or
unavailable monetary ledger fails closed while a configured limit is in
force, so the supervisor cannot dispatch against stale counters.
``projection`` is an optional pre-computed global usage projection (same limit and drive
root) so a caller that already replayed the ledger — e.g. ``/api/state`` — does not replay
it again. It is accepted only when its ``limit_usd`` equals the limit this function reads
itself (what ``usage_projection`` stamps); a mismatch falls through to the read below.
``allow_stale`` is for a loop-thread pre-check that must not wait on money: it rides the
last validated snapshot against the LIVE limit. A snapshot may only ADMIT (every paid
attempt still passes ``reserve_attempt``): an answer at or below ``refuse_below`` and a cold
memo are decided on the exact locked read; a passed ``projection`` is display, never re-read.
"""
total = float(TOTAL_BUDGET_LIMIT or 0.0)
if total <= 0:
return float('inf')
if projection is not None and projection.get("limit_usd") != round(max(0.0, total), 6):
projection = None
try:
if projection is None:
from ouroboros.usage_accounting import ensure_legacy_imported, usage_projection
from ouroboros.usage_ledger import UsageLockUnavailable
ensure_legacy_imported(DRIVE_ROOT)
with contextlib.suppress(*((UsageLockUnavailable,) if allow_stale else ())):
projection = usage_projection(DRIVE_ROOT, global_limit_usd=total, allow_stale=allow_stale)
if projection is None or (
allow_stale and float(projection.get("remaining_known_usd") or 0.0) <= refuse_below):
projection = usage_projection(DRIVE_ROOT, global_limit_usd=total)
return float(projection.get("remaining_known_usd") or 0.0)
except Exception:
log.exception("Budget ledger unavailable; refusing new model dispatch")
if strict:
raise
return 0.0
def reset_per_task_budget(data_root: Any, *, confirm_isolated: bool = False) -> bool:
"""Zero legacy budget *projection* fields in an isolated benchmark root.
The append-only physical-attempt ledger is deliberately untouched. New
runs receive their allowance through a new root task id/root limit, while a
campaign limit continues to cover every prior physical attempt. This
compatibility helper only prevents stale pre-ledger state fields from being
mistaken for a current per-task counter by old benchmark tooling.
CRITICAL safety guard (BIBLE P8): the live TOTAL_BUDGET / Emergency-Stop
contract must never be defeated by a reset. This refuses unless ALL hold:
the target is NOT the live ``~/Ouroboros/data`` dir, the caller passes
``confirm_isolated=True`` (explicit bench intent), and ``OUROBOROS_DATA_DIR``
is set (a non-default, isolated data dir). Evolutionary drivers call this
between tasks so each instance starts with a fresh per-task allowance while
learned knowledge/identity/code carry forward. Returns True only when a reset
was actually written.
"""
try:
target = pathlib.Path(str(data_root)).resolve(strict=False)
except Exception:
return False
live = (pathlib.Path.home() / "Ouroboros" / "data").resolve(strict=False)
if target == live:
return False
if not confirm_isolated:
return False
env_dir = str(os.environ.get("OUROBOROS_DATA_DIR", "") or "").strip()
if not env_dir:
return False
try:
if pathlib.Path(env_dir).resolve(strict=False) != target:
return False
except Exception:
return False
# Final guard: the target MUST carry the isolated-benchmark sentinel. A live root (default
# or custom/Drive-backed) never has it, so this reset can never zero a live budget even if
# the home-path comparison above does not match a non-default live data root (BIBLE P8).
if not (target / ISOLATED_BENCHMARK_SENTINEL).exists():
return False
state_path = target / "state" / "state.json"
# Lock on the TARGET root's own state.lock (the isolated server holds the same
# path as its STATE_LOCK), so this between-instance reset and a concurrent server
# save_state cannot lost-update each other in the B-full server-driven model.
lock_path = target / "locks" / "state.lock"
try:
lock_path.parent.mkdir(parents=True, exist_ok=True)
except OSError:
return False
lock_fd = acquire_file_lock(lock_path)
if lock_fd is None:
# Lock acquisition timed out (the isolated server is actively writing state):
# skip rather than run an UNLOCKED read-modify-write that would race save_state.
log.warning("reset_per_task_budget: could not acquire state lock for %s; skipping reset", state_path)
return False
budget_keys = ("spent_usd", "spent_calls", "spent_tokens_prompt",
"spent_tokens_completion", "spent_tokens_cached")
try:
st = json_load_file(state_path) or {}
st["spent_usd"] = 0.0
st["spent_calls"] = 0
st["spent_tokens_prompt"] = 0
st["spent_tokens_completion"] = 0
st["spent_tokens_cached"] = 0
atomic_write_text(state_path, json.dumps(st, ensure_ascii=False, indent=2))
# Also zero the budget counters in the last-good snapshot. _load_state
# falls back to it when state.json is missing/corrupt; leaving stale
# spend there could re-inflate the per-task ledger after a mid-run
# crash+recovery, defeating the reset (the kit reset both files).
lg_path = target / "state" / "state.last_good.json"
lg = json_load_file(lg_path)
if isinstance(lg, dict):
for key in budget_keys:
if key in lg:
lg[key] = 0 if key != "spent_usd" else 0.0
atomic_write_text(lg_path, json.dumps(lg, ensure_ascii=False, indent=2))
except Exception:
log.warning("reset_per_task_budget: failed to write %s", state_path, exc_info=True)
return False
finally:
release_file_lock(lock_path, lock_fd)
return True
def _openrouter_key_fingerprint() -> str:
"""Non-secret identity of the currently configured OpenRouter key.
Drift comparison is only meaningful while the ledger baseline and the
``/auth/key`` ground truth describe the SAME key; a settings hot-reload can
swap the key mid-session. Returns a short sha256 prefix (never the key)."""
import hashlib
api_key = os.environ.get("OPENROUTER_API_KEY", "").strip()
if not api_key:
return ""
return hashlib.sha256(api_key.encode("utf-8")).hexdigest()[:16]
def _openrouter_ledger_settled(breakdown: Optional[Dict[str, Any]] = None) -> Optional[float]:
"""Cumulative settled USD attributed to provider=openrouter in the attempt
ledger, or None when the ledger is unavailable. Settled-only on purpose:
reservations/unresolved bounds are conservative estimates and would inflate
the tracked side of the drift comparison."""
try:
if breakdown is None:
from ouroboros.usage_accounting import ensure_legacy_imported, usage_breakdown
ensure_legacy_imported(DRIVE_ROOT)
breakdown = usage_breakdown(DRIVE_ROOT)
bucket = dict(breakdown.get("by_provider") or {}).get("openrouter") or {}
return float(bucket.get("settled_usd") or 0.0)
except Exception:
log.debug("OpenRouter ledger settled total unavailable", exc_info=True)
return None
def check_openrouter_ground_truth() -> Optional[Dict[str, float]]:
"""Return OpenRouter total/daily usage, or None on error."""
try:
import urllib.request
api_key = os.environ.get("OPENROUTER_API_KEY", "").strip()
if not api_key:
return None
from ouroboros.net_transport import trust_ssl_context
req = urllib.request.Request(
"https://openrouter.ai/api/v1/auth/key",
headers={"Authorization": f"Bearer {api_key}"},
)
# A provider call: it verifies against the owner's trust bundle like every other one.
with urllib.request.urlopen(req, timeout=10, context=trust_ssl_context()) as resp:
data = json.loads(resp.read().decode("utf-8"))
# OpenRouter usage is dollars, not cents.
usage_total = data.get("data", {}).get("usage", 0)
usage_daily = data.get("data", {}).get("usage_daily", 0)
return {
"total_usd": float(usage_total),
"daily_usd": float(usage_daily),
}
except Exception:
log.warning("Failed to fetch OpenRouter ground truth", exc_info=True)
return None
def budget_pct(st: Dict[str, Any]) -> float:
"""Return ledger-derived budget percent used."""
total = float(TOTAL_BUDGET_LIMIT or 0.0)
if total <= 0:
return 0.0
try:
from ouroboros.usage_accounting import ensure_legacy_imported, usage_projection
ensure_legacy_imported(DRIVE_ROOT)
projection = usage_projection(DRIVE_ROOT, global_limit_usd=total)
return (float(projection.get("accounted_usd") or 0.0) / total) * 100.0
except Exception:
log.exception("Budget ledger unavailable while calculating percent")
return 100.0
def update_budget_from_usage(usage: Dict[str, Any]) -> bool:
"""Refresh the legacy state projection from the physical-attempt ledger.
``usage`` is retained for caller compatibility but is never added to the
monetary total: every core-mediated provider attempt has already been
persisted by the transport wrapper. This prevents logical usage events,
retries, and review aggregation from charging the same attempt twice.
The persisted projection carries totals only; the per-root map is never written.
The ledger read is the writer's slim snapshot (``usage_writer_snapshot``): only what this
function persists is rendered; the loop's llm_usage path writes once per turn, direct callers on call.
"""
def _to_float(v: Any, default: float = 0.0) -> float:
try:
return float(v)
except Exception:
log.debug(f"Failed to convert value to float: {v!r}", exc_info=True)
return default
def _to_int(v: Any, default: int = 0) -> int:
try:
return int(v)
except Exception:
log.debug(f"Failed to convert value to int: {v!r}", exc_info=True)
return default
def _ledger_high_water_marker(breakdown: Dict[str, Any]) -> Optional[tuple[int, int]]:
"""Return the ``(compaction_epoch, seq)`` fact from this read.
``usage_breakdown`` supplies this field while rendering the same
validated rows used for the monetary buckets. Do not reconstruct it
from a second read or a private accounting cache: a marker that cannot
be parsed is unknown, never zero.
"""
marker = breakdown.get("_ledger_high_water_seq")
if not isinstance(marker, (list, tuple)) or len(marker) != 2:
return None
epoch, seq = marker
if any(isinstance(value, bool) or not isinstance(value, int) or value < 0
for value in (epoch, seq)):
return None
return int(epoch), int(seq)
from ouroboros.usage_accounting import (
UsageLedgerCorrupt,
ensure_legacy_imported,
usage_projection,
usage_writer_snapshot,
)
# Ledger I/O is deliberately OUTSIDE STATE_LOCK: the lock stays
# short-lived, and the validated-snapshot marker below preserves the old
# serialization invariant without holding STATE_LOCK across a long read.
# A DISPLAY read (``allow_stale``: this runs on the supervisor loop once per turn with
# ``llm_usage`` events): a lagging snapshot carries its own lower marker, so it never regresses money.
try:
ensure_legacy_imported(DRIVE_ROOT)
breakdown = usage_writer_snapshot(DRIVE_ROOT, allow_stale=True)
total_limit = float(TOTAL_BUDGET_LIMIT or 0.0)
projection_snapshot = breakdown.pop("_usage_projection", None)
if total_limit > 0 and isinstance(projection_snapshot, dict):
from ouroboros._usage_rows import _with_limit
# Totals only (issue #1002): per-root money is a ledger render nothing reads back from here.
projection_snapshot.pop("by_root", None)
projection = _with_limit(projection_snapshot, total_limit)
else:
projection = (
usage_projection(DRIVE_ROOT, global_limit_usd=total_limit, include_roots=False, allow_stale=True)
if total_limit > 0
else {key: breakdown.get(key) for key in (
"settled_usd", "confirmed_usd", "estimated_usd", "reserved_usd",
"unresolved_upper_bound_usd", "accounted_usd", "unknown_unmetered",
"cost_final", "attempt_counts", "integrity_degraded",
)}
)
except UsageLedgerCorrupt:
# A damaged ledger is unknown, never zero. Leave the prior projection
# in place, report the refusal to the caller, and let the next event
# retry: paid usage itself is already persisted in the ledger.
log.warning("Skipping legacy budget projection: usage ledger is corrupt", exc_info=True)
return False
ledger_high_water_marker = (
None if breakdown.get("integrity_degraded") else _ledger_high_water_marker(breakdown)
)
openrouter_ledger_settled = _openrouter_ledger_settled(breakdown)
should_check_ground_truth = False
lock_fd = acquire_file_lock(STATE_LOCK_PATH)
if lock_fd is None: # the ledger stays the money authority; this projection waits for the next event
log.warning("legacy budget projection skipped: state lock timeout")
return False
try:
try:
st = _load_state_unlocked()
except StateUnavailable as exc:
log.warning("legacy budget projection skipped: %s", exc)
return False
previous_marker = st.get("usage_ledger_high_water_seq")
previous_known = (
isinstance(previous_marker, (list, tuple))
and len(previous_marker) == 2
and all(isinstance(value, int) and not isinstance(value, bool) and value >= 0
for value in previous_marker)
)
previous_marker_present = "usage_ledger_high_water_seq" in st
if ledger_high_water_marker is None or (previous_marker_present and not previous_known):
# An unreadable/missing marker is unknown, never zero. Keep the
# existing projection untouched: writing money without ordering
# evidence could reintroduce the stale-snapshot regression.
log.warning(
"legacy budget projection FRESHNESS MARKER UNKNOWN: preserving prior projection"
)
return False
if previous_known:
saved_marker = (int(previous_marker[0]), int(previous_marker[1]))
if ledger_high_water_marker < saved_marker:
# ANY lower marker is refused, epoch or seq: a delayed writer
# holding a pre-compaction snapshot must never overwrite money
# a newer snapshot already saved. Equal/higher markers keep the
# normal positive update path below.
log.warning(
"legacy budget projection STALE SNAPSHOT REJECTED: ledger marker %s < saved %s",
ledger_high_water_marker,
saved_marker,
)
return False
st["spent_usd"] = _to_float(breakdown.get("accounted_usd"))
st["spent_calls"] = _to_int(breakdown.get("physical_calls"))
st["spent_tokens_prompt"] = _to_int(breakdown.get("prompt_tokens"))
st["spent_tokens_completion"] = _to_int(breakdown.get("completion_tokens"))
st["spent_tokens_cached"] = _to_int(breakdown.get("cached_tokens"))
st["usage_accounting"] = projection
st["openrouter_ledger_settled_usd"] = openrouter_ledger_settled
# Historical key retained for state.json compatibility; its value
# is now the ordered ``[compaction_epoch, seq]`` pair.
st["usage_ledger_high_water_seq"] = list(ledger_high_water_marker)
previous_check_call = _to_int(st.get("openrouter_last_check_call"), -1)
# Every 50th call by CROSSING (a coalesced write may jump 49 -> 51), deduped by the last check.
should_check_ground_truth = st["spent_calls"] > 0 and st["spent_calls"] // 50 > max(previous_check_call, 0) // 50
if should_check_ground_truth:
st["openrouter_last_check_call"] = st["spent_calls"]
_save_state_unlocked(st)
finally:
release_file_lock(STATE_LOCK_PATH, lock_fd)
if should_check_ground_truth:
ground_truth = check_openrouter_ground_truth()
if ground_truth is not None:
lock_fd = acquire_file_lock(STATE_LOCK_PATH)
if lock_fd is None:
return True
try:
try:
st = _load_state_unlocked()
except StateUnavailable:
return True
st["openrouter_total_usd"] = ground_truth["total_usd"]
st["openrouter_daily_usd"] = ground_truth["daily_usd"]
st["openrouter_last_check_at"] = utc_now_iso()
# Drift compares the OpenRouter-only settled ledger delta with
# the queried key's usage delta — the only like-for-like pair.
# Direct-provider spend is invisible to /auth/key by
# construction and must not count as "drift".
session_total_snap = st.get("session_total_snapshot")
session_or_settled_snap = st.get("session_openrouter_settled_snapshot")
or_ledger_settled = st.get("openrouter_ledger_settled_usd")
current_fp = _openrouter_key_fingerprint()
baseline_fp = str(st.get("session_openrouter_key_fp") or "")
integrity_degraded = bool(breakdown.get("integrity_degraded"))
key_changed = bool(current_fp) and bool(baseline_fp) and current_fp != baseline_fp
if integrity_degraded:
# A quarantined ledger tail makes the tracked side
# non-final; a confident percentage would be dishonest.
# Comparison is suppressed, not zeroed.
st["budget_drift_pct"] = None
st["budget_drift_alert"] = False
elif (
key_changed
or session_total_snap is None
or session_or_settled_snap is None
or or_ledger_settled is None
):
# Rebaseline (key swapped mid-session, or pre-upgrade state
# lacks the OpenRouter-only snapshot) and skip this cycle:
# the old baseline describes a different key/metric.
st["session_total_snapshot"] = ground_truth["total_usd"]
st["session_openrouter_settled_snapshot"] = or_ledger_settled
st["session_openrouter_key_fp"] = current_fp
st["budget_drift_pct"] = None
st["budget_drift_alert"] = False
else:
or_delta = ground_truth["total_usd"] - _to_float(session_total_snap)
our_delta = _to_float(or_ledger_settled) - _to_float(session_or_settled_snap)
if or_delta > 0.001:
drift_pct = abs(or_delta - our_delta) / max(abs(or_delta), 0.01) * 100.0
st["budget_drift_pct"] = drift_pct
abs_diff = abs(or_delta - our_delta)
if drift_pct > 50.0 and abs_diff > 5.0:
st["budget_drift_alert"] = True
all_provider_delta = _to_float(st.get("spent_usd") or 0.0) - _to_float(
st.get("session_spent_snapshot") or 0.0
)
append_jsonl(
DRIVE_ROOT / "logs" / "events.jsonl",
{
"ts": utc_now_iso(),
# "type" is the events.jsonl schema key every
# other event uses; type-keyed aggregations
# lost this row when it was written as "event".
"type": "budget_drift_warning",
"drift_pct": round(drift_pct, 2),
"our_delta": round(our_delta, 4),
"or_delta": round(or_delta, 4),
"abs_diff": round(abs_diff, 4),
"all_provider_delta": round(all_provider_delta, 4),
"spent_calls": st["spent_calls"],
"note": (
"OpenRouter-only ledger delta vs /auth/key usage delta. "
"High drift usually means a shared OR key or missing ledger rows."
),
}
)
else:
st["budget_drift_alert"] = False
else:
st["budget_drift_pct"] = 0.0
st["budget_drift_alert"] = False
_save_state_unlocked(st)
finally:
release_file_lock(STATE_LOCK_PATH, lock_fd)
return True
def budget_breakdown(st: Dict[str, Any]) -> Dict[str, float]:
"""Aggregate accounted physical-attempt cost by category."""
breakdown: Dict[str, float] = {}
try:
from ouroboros.usage_accounting import ensure_legacy_imported, usage_breakdown
ensure_legacy_imported(DRIVE_ROOT)
ledger = usage_breakdown(DRIVE_ROOT, allow_stale=True)
for category, bucket in dict(ledger.get("by_category") or {}).items():
breakdown[str(category)] = float(bucket.get("accounted_usd") or 0.0)
unattributed = dict(ledger.get("unattributed") or {}).get("category") or {}
if float(unattributed.get("accounted_usd") or 0.0) > 0:
breakdown["(unattributed)"] = float(unattributed.get("accounted_usd") or 0.0)
except Exception:
log.error("Failed to calculate ledger budget breakdown", exc_info=True)
return breakdown
def model_breakdown(st: Dict[str, Any]) -> Dict[str, Dict[str, float]]:
"""Aggregate physical calls/tokens/accounted cost by model."""
breakdown: Dict[str, Dict[str, float]] = {}
try:
from ouroboros.usage_accounting import ensure_legacy_imported, usage_breakdown
ensure_legacy_imported(DRIVE_ROOT)
ledger = usage_breakdown(DRIVE_ROOT, allow_stale=True)
buckets = dict(ledger.get("by_model") or {})
unattributed = dict(ledger.get("unattributed") or {}).get("model") or {}
if int(unattributed.get("physical_calls") or 0) or float(unattributed.get("accounted_usd") or 0.0):
buckets["(unattributed)"] = unattributed
for model, bucket in buckets.items():
breakdown[str(model)] = {
"cost": float(bucket.get("accounted_usd") or 0.0),
"calls": int(bucket.get("physical_calls") or 0),
"prompt_tokens": int(bucket.get("prompt_tokens") or 0),
"completion_tokens": int(bucket.get("completion_tokens") or 0),
"cached_tokens": int(bucket.get("cached_tokens") or 0),
}
except Exception:
log.error("Failed to calculate ledger model breakdown", exc_info=True)
return breakdown
def per_task_cost_summary(max_tasks: int = 10, tail_bytes: int = 512_000) -> List[Dict[str, Any]]:
"""Return task cost summary from ledger-attributed physical attempts."""
del tail_bytes # compatibility-only; the append-only ledger is replayed in full
tasks: Dict[str, Dict[str, Any]] = {}
try:
from ouroboros.usage_accounting import ensure_legacy_imported, usage_breakdown
ensure_legacy_imported(DRIVE_ROOT)
ledger = usage_breakdown(DRIVE_ROOT)
for task_id, bucket in dict(ledger.get("by_task") or {}).items():
tasks[str(task_id)] = {
"task_id": str(task_id),
"cost": float(bucket.get("accounted_usd") or 0.0),
"rounds": int(bucket.get("physical_calls") or 0),
"model": "",
}
except Exception:
log.error("Failed to calculate ledger per-task cost summary", exc_info=True)
raise
sorted_tasks = sorted(tasks.values(), key=lambda x: x["cost"], reverse=True)
return sorted_tasks[:max_tasks]
def reconstruct_task_cost(
task_id: str, *, fields: bool = False, drive_root: Optional[pathlib.Path] = None,
breakdown: Optional[Dict[str, Any]] = None,
) -> Any:
"""Reconstruct cost; an indexed breakdown belongs to this drive after legacy import."""
want = str(task_id or "")
if not want:
projection = {
"cost_accounting_status": "available", "accounted_upper_bound_usd": 0.0,
"total_rounds": 0, "prompt_tokens": 0, "completion_tokens": 0,
"cost_final": True, "reserved_usd": 0.0,
"unresolved_upper_bound_usd": 0.0, "unknown_unmetered": 0,
"non_final_rows": 0,
# No task was named, so no ledger bucket was summed: there is nothing
# for a carrier to explain (#498), and None says exactly that.
"cost_presentation": None,
}
else:
try:
from ouroboros.cost_projection import (
COST_SCOPE_OWN, build_cost_presentation, honest_accounted_amount,
)
from ouroboros.usage_accounting import ensure_legacy_imported, usage_breakdown
authority_root = pathlib.Path(drive_root) if drive_root is not None else DRIVE_ROOT
if breakdown is None:
ensure_legacy_imported(authority_root)
bucket = usage_breakdown(authority_root, task_id=want)
else:
from ouroboros._usage_rows import _breakdown_bucket, _with_integrity
bucket = breakdown["by_task"].get(want)
if bucket is None:
bucket = _with_integrity(_breakdown_bucket(()), bool(breakdown.get("integrity_degraded")))
projection = {
"cost_accounting_status": "available",
"accounted_upper_bound_usd": (
round(amount, 6)
if (amount := honest_accounted_amount(bucket)) is not None
else None
),
"total_rounds": int(bucket.get("physical_calls") or 0),
"prompt_tokens": int(bucket.get("prompt_tokens") or 0),
"completion_tokens": int(bucket.get("completion_tokens") or 0),
"cost_final": bool(bucket.get("cost_final")),
"reserved_usd": float(bucket.get("reserved_usd") or 0.0),
"unresolved_upper_bound_usd": float(
bucket.get("unresolved_upper_bound_usd") or 0.0
),
"unknown_unmetered": int(bucket.get("unknown_unmetered") or 0),
# The disclosed CAUSE of cost_final=false, carried with the flag.
"non_final_rows": int(bucket.get("non_final_rows") or 0),
"ledger_integrity_degraded": bool(bucket.get("integrity_degraded")),
# #498: the same bucket's own explanation of its own amount.
"cost_presentation": build_cost_presentation(bucket, scope=COST_SCOPE_OWN),
}
except Exception:
log.error("Failed to reconstruct ledger task cost for %s", task_id, exc_info=True)
projection = {
"cost_accounting_status": "unavailable", "cost_final": False,
"cost_accounting_error": "ledger_unavailable",
"cost_presentation": None,
"accounted_upper_bound_usd": None, "total_rounds": None,
"prompt_tokens": None, "completion_tokens": None,
"reserved_usd": None, "unresolved_upper_bound_usd": None,
"unknown_unmetered": None, "non_final_rows": None,
"ledger_integrity_degraded": True,
}
if fields:
# SSOT cost naming (C2/ABI-3): the authority assembles the honest
# name directly (Ф3.1 fix-round — no producer touches the retired
# alias); the seam stays as the idempotent amount-normalization and
# would strip any retired key a future mutation leaked.
from ouroboros.cost_projection import with_cost_aliases
return with_cost_aliases(projection)
if projection.get("cost_accounting_status") != "available":
from ouroboros.usage_accounting import UsageAccountingError
raise UsageAccountingError(f"task cost authority unavailable for {task_id}")
return (
float(projection["accounted_upper_bound_usd"]), int(projection["total_rounds"]),
int(projection["prompt_tokens"]), int(projection["completion_tokens"]),
)
def status_text(workers_dict: Dict[int, Any], pending_list: list,
running_dict: Dict[str, Dict[str, Any]]) -> str:
"""Build status text from worker and queue state."""
st = load_state()
now = time.time()
lines = []
lines.append(f"owner_id: {st.get('owner_id')}")
lines.append(f"session_id: {st.get('session_id')}")
lines.append(f"version: {st.get('current_branch')}@{(st.get('current_sha') or '')[:8]}")
busy_count = sum(1 for w in workers_dict.values() if getattr(w, 'busy_task_id', None) is not None)
lines.append(f"workers: {len(workers_dict)} (busy: {busy_count})")
lines.append(f"pending: {len(pending_list)}")
lines.append(f"running: {len(running_dict)}")
if pending_list:
preview = []
for t in pending_list[:10]:
preview.append(
f"{t.get('id')}:{t.get('type')}:pr{t.get('priority')}:a{int(t.get('_attempt') or 1)}")
lines.append("pending_queue: " + ", ".join(preview))
if running_dict:
lines.append("running_ids: " + ", ".join(list(running_dict.keys())[:10]))
busy = [f"{getattr(w, 'wid', '?')}:{getattr(w, 'busy_task_id', '?')}"
for w in workers_dict.values() if getattr(w, 'busy_task_id', None)]
if busy:
lines.append("busy: " + ", ".join(busy))
if running_dict:
details = []
for task_id, meta in list(running_dict.items())[:10]:
task = meta.get("task") if isinstance(meta, dict) else {}
started = float(meta.get("started_at") or 0.0) if isinstance(meta, dict) else 0.0
hb = float(meta.get("last_heartbeat_at") or 0.0) if isinstance(meta, dict) else 0.0
runtime_sec = int(max(0.0, now - started)) if started > 0 else 0
hb_lag_sec = int(max(0.0, now - hb)) if hb > 0 else -1
details.append(
f"{task_id}:type={task.get('type')} pr={task.get('priority')} "
f"attempt={meta.get('attempt')} runtime={runtime_sec}s hb_lag={hb_lag_sec}s")
if details:
lines.append("running_details:")
lines.extend([f" - {d}" for d in details])
if running_dict and busy_count == 0:
lines.append("queue_warning: running>0 while busy=0")
accounting_available = True
try:
from ouroboros.usage_accounting import ensure_legacy_imported, usage_breakdown, usage_projection
ensure_legacy_imported(DRIVE_ROOT)
ledger_breakdown = usage_breakdown(DRIVE_ROOT, allow_stale=True) # /status renders on the loop
ledger_projection = (
usage_projection(DRIVE_ROOT, global_limit_usd=TOTAL_BUDGET_LIMIT, allow_stale=True)
if TOTAL_BUDGET_LIMIT > 0
else ledger_breakdown
)
spent = float(ledger_projection.get("accounted_usd") or 0.0)
pct = (spent / TOTAL_BUDGET_LIMIT * 100.0) if TOTAL_BUDGET_LIMIT > 0 else 0.0
budget_remaining_usd = (
float(ledger_projection.get("remaining_known_usd") or 0.0)
if TOTAL_BUDGET_LIMIT > 0
else float("inf")
)
spent_calls = int(ledger_breakdown.get("physical_calls") or 0)
prompt_tokens = int(ledger_breakdown.get("prompt_tokens") or 0)
completion_tokens = int(ledger_breakdown.get("completion_tokens") or 0)
cached_tokens = int(ledger_breakdown.get("cached_tokens") or 0)
except Exception:
log.exception("Budget ledger unavailable while building status")
accounting_available = False
spent = pct = budget_remaining_usd = None
spent_calls = prompt_tokens = completion_tokens = cached_tokens = None
lines.append(f"budget_total: ${TOTAL_BUDGET_LIMIT:.0f}")
if not accounting_available:
lines.append("accounting: unavailable (physical-attempt ledger read failed; dispatch is fail-closed)")
lines.append("budget_remaining: unavailable")
lines.append("spent_usd: unavailable")
lines.append("spent_calls: unavailable")
lines.append("prompt_tokens: unavailable, completion_tokens: unavailable, cached_tokens: unavailable")
else:
if ledger_projection.get("integrity_degraded"):
lines.append("accounting_integrity: DEGRADED (quarantined ledger tail; cost is not final)")
lines.append(f"budget_remaining: ${budget_remaining_usd:.0f}")
if pct > 0:
lines.append(f"spent_usd: ${spent:.2f} ({pct:.1f}% of budget)")
else:
lines.append(f"spent_usd: ${spent:.2f}")
lines.append(f"spent_calls: {spent_calls}")
lines.append(
f"prompt_tokens: {prompt_tokens}, completion_tokens: {completion_tokens}, "
f"cached_tokens: {cached_tokens}"
)
breakdown = budget_breakdown(st)
if breakdown:
sorted_categories = sorted(breakdown.items(), key=lambda x: x[1], reverse=True)
breakdown_parts = [f"{cat}=${cost:.2f}" for cat, cost in sorted_categories if cost > 0]
if breakdown_parts:
lines.append(f"budget_breakdown: {', '.join(breakdown_parts)}")
drift_pct = st.get("budget_drift_pct")
if accounting_available and drift_pct is not None:
# Same fields as the update_budget computation: OpenRouter-only settled
# ledger delta vs the key's usage delta. Rendering the all-provider
# spent delta here used to contradict the percentage next to it.
session_total_snap = st.get("session_total_snapshot")
session_or_settled_snap = st.get("session_openrouter_settled_snapshot")
or_ledger_settled = st.get("openrouter_ledger_settled_usd")
or_total = st.get("openrouter_total_usd")
if (
session_total_snap is not None
and session_or_settled_snap is not None
and or_ledger_settled is not None
and or_total is not None
):
or_delta = or_total - session_total_snap
our_delta = float(or_ledger_settled) - float(session_or_settled_snap)
drift_icon = " ⚠️" if st.get("budget_drift_alert") else ""
lines.append(
f"budget_drift: {drift_pct:.1f}%{drift_icon} "
f"(openrouter tracked: ${our_delta:.2f} vs OpenRouter key: ${or_delta:.2f})"
)
models = model_breakdown(st)
if models:
sorted_models = sorted(models.items(), key=lambda x: x[1]["cost"], reverse=True)
lines.append("model_breakdown:")
for model_name, stats in sorted_models:
if stats["cost"] > 0 or stats["calls"] > 0:
cost = stats["cost"]
calls = int(stats["calls"])
pt = int(stats["prompt_tokens"])
ct = int(stats["completion_tokens"])
lines.append(f" {model_name}: ${cost:.2f} ({calls} calls, {pt:,}p/{ct:,}c tok)")
lines.append(
"evolution: "
+ f"enabled={int(bool(st.get('evolution_mode_enabled')))}, "
+ f"cycle={int(st.get('evolution_cycle') or 0)}")
lines.append(f"last_owner_message_at: {st.get('last_owner_message_at') or '-'}")
lines.append("active_liveness: idle+deadline+absolute_ceiling+reaper")
return "\n".join(lines)
def rotate_jsonl_log_if_needed(
drive_root: pathlib.Path,
name: str,
archive_prefix: str,
max_bytes: int = 800_000,
) -> None:
"""Rotate ``logs/<name>`` to ``archive/<archive_prefix>_<ts>.jsonl`` when it
exceeds ``max_bytes``.
Rotation is an atomic ``os.replace`` rename performed under the SAME sidecar
lock that ``append_jsonl`` writers take — the old copy+truncate destroyed any
line appended between the read and the truncate.
Suppressed in isolated benchmark data roots (``ISOLATED_BENCHMARK_SENTINEL``):
bench harnesses read trial-local logs as one file from birth, the roots are
throwaway, and not every harness reader is archive-chain-aware.
Archives are durable history, NOT GC targets: no retention sweep touches
``archive/`` (retention.py governs subagent worktrees, task drives, and
service logs only) and none may be added — readers backfill from these
segments, so pruning them would silently erase visible history (BIBLE P1).
"""
if (drive_root / ISOLATED_BENCHMARK_SENTINEL).exists():
return
path = drive_root / "logs" / name
if not path.exists():
return
if path.stat().st_size < max_bytes:
return
ts = utc_now_iso().replace("-", "").replace(":", "").split(".")[0]
archive_path = drive_root / "archive" / f"{archive_prefix}_{ts}.jsonl"
# Second-resolution names can collide when a fast writer forces two rotations
# within one second; os.replace onto an existing archive would destroy it. The
# "_<n>" suffix sorts lexicographically AFTER "<ts>.jsonl" ("_" > "."), so
# name-ordered readers keep the true chronological chain.
suffix = 0
while archive_path.exists():
suffix += 1
archive_path = drive_root / "archive" / f"{archive_prefix}_{ts}_{suffix}.jsonl"
archive_path.parent.mkdir(parents=True, exist_ok=True)
from ouroboros.utils import jsonl_append_lock_path
lock_path = jsonl_append_lock_path(path)
lock_fd = acquire_exclusive_file_lock(lock_path, timeout_sec=2.0, stale_sec=10.0, owner_aware_stale=True)
if lock_fd is None:
log.warning("%s rotation skipped: append lock busy", name)
return
try:
if not path.exists() or path.stat().st_size < max_bytes:
return
os.replace(path, archive_path)
path.touch()
finally:
release_exclusive_file_lock(lock_path, lock_fd)
def rotate_chat_log_if_needed(drive_root: pathlib.Path, max_bytes: int = 800_000) -> None:
"""Compatibility wrapper: chat.jsonl rotation via the generalized rotator."""
rotate_jsonl_log_if_needed(drive_root, "chat.jsonl", "chat", max_bytes)