mirror of
https://github.com/razzant/ouroboros.git
synced 2026-10-03 04:07:04 +00:00
fix(state): read the usage ledger outside STATE_LOCK
update_budget_from_usage held STATE_LOCK (4.0s timeout) across ensure_legacy_imported, usage_breakdown and usage_projection, which contend for the monetary lock (45.0s timeout). A short-timeout lock held across a long-timeout one starved every state.json reader, /api/state among them, and stalled the supervisor loop. The ledger read now happens outside the critical section; only load, compare, mutate and save remain inside. The lock was held deliberately: reading the ledger outside it lets an older concurrent snapshot acquire STATE_LOCK later and regress the monetary fields. That invariant is preserved by an ordering marker instead of by lock duration. The marker is the lexicographic pair (compaction_epoch, max live row seq), derived by _usage_rows from the same validated rows the buckets come from and exposed as the private _ledger_high_water_seq. Neither component works alone: compaction starts a fresh dense epoch, so the live seq restarts at 1, header seq is always 1, pre_compaction_seq is per-row provenance into the archived range, and source_last_seq is the just-folded file's length. Only compaction_epoch is globally monotonic, and only the pair orders snapshots both inside an epoch and across a fold. Any marker strictly lower than the saved one is refused, whether the epoch or only the sequence differs. An earlier revision of this change wrote a lower epoch on the theory that it could only mean a restore. That is wrong, and no restore is needed to break it: a writer can read at epoch N, a reserve_attempt can compact to N+1, a second writer can save the post-compaction projection, and the delayed writer would then overwrite state.json with stale money. Equal and higher markers take the normal write path. A genuine restore can therefore freeze the compatibility projection until usage_ledger_high_water_seq is cleared; that belongs to the restore path, not to this writer, and the refusal is logged rather than silent. Each decision carries a greppable substring (FRESHNESS MARKER UNKNOWN, STALE SNAPSHOT REJECTED). The marker and every persisted accounting field now come from one validated snapshot: usage_breakdown renders the projection from the same final rows, so two writers cannot share a marker while saving projections read at different sequences. An unreadable ledger is unknown, never zero: the previous projection is preserved, update_budget_from_usage returns a typed false, and _handle_llm_usage records projection_update_status=unavailable instead of claiming an availability it does not have. usage_breakdown's private provenance is filtered at the gateway boundary, where the unbounded-budget branch reuses the whole mapping as its /api/state accounting projection. The filter drops every leading-underscore key, so the class is closed rather than this one key. Measured on a copy of a real 33 MB ledger in a throwaway data root, one writer against six readers for 25s: reader lock-timeout fallbacks 21 to 0, reader max latency 4.05s to 1.07s (the 4.0s timeout ceiling was being hit), reader throughput 5019 to 20503 calls. With six writers, 18 to 0 across three runs. Writers alone are faster, not slower: 431 to 855 calls, p50 29.3ms to 16.6ms.
This commit is contained in:
parent
499e5609fc
commit
2ed63f8aa9
14 changed files with 714 additions and 101 deletions
|
|
@ -4,6 +4,8 @@ This chapter owns the single scheduler for pooled work: what a healthy tick does
|
|||
|
||||
`server.py::_run_supervisor()` is the single scheduler for pooled tasks. A healthy tick publishes liveness, rotates runtime logs, checks worker health, drains worker and direct-chat events (a wake-up's frames are a direct turn's), accepts owner bridge input, enforces deadlines and schedules, runs throttled reconciliation and evolution admission, assigns eligible work, persists `state/queue_snapshot.json`, and last ticks the consciousness alarm clock, which owns no thread — this pass is the only thing that can start a wake-up. Bridge intake precedes timeout, maintenance, evolution and assignment work so a slow control-plane step cannot hide a new owner message. Three consecutive loop failures clear supervisor readiness, stop its watchdog generation and notify the owner, rather than leave a healthy-looking server that no longer assigns work; a failure while a shutdown or restart is in progress (the teardown sets a process-local stop event first) is not a crash and never holds the shutdown.
|
||||
|
||||
The legacy `state.json` budget projection reads the validated usage ledger before taking `STATE_LOCK`, so a multi-megabyte replay cannot hold the short reader lock while waiting on the monetary lock. The same read returns a ledger provenance marker `(compaction_epoch, seq)`: the epoch comes from the lock-free leading baseline header and `seq` is the live file's validated high-water sequence. A compaction increments the epoch while renumbering live rows, so the pair remains ordered even when the file gets shorter; timestamps, row counts, file sizes, and dollar totals are not ordering evidence. Inside `STATE_LOCK`, a same-epoch lower marker is proven stale and leaves the compatibility projection untouched. An equal marker is written harmlessly, a higher marker is written, and a lower epoch is written with a warning because restore lineage is unknown. An unavailable or malformed marker leaves the prior projection untouched; this fail-safe also applies when lock acquisition times out and the writer deliberately proceeds without the lock. On that no-lock path the marker rejects a demonstrably stale same-epoch snapshot, but the compare-and-save sequence remains non-atomic for concurrent no-lock writers, as it was before this protection.
|
||||
|
||||
`PENDING` and `RUNNING`, guarded by `supervisor.queue._queue_lock`, are the live task-lifecycle authority. Admission reserves identity before project, workspace, attachment or routing side effects can create a duplicate; refuses a disabled pool, duplicate task, project deletion, accepted or sealed root, exhausted root budget, or an unusable or unprovisionable workspace (typed, with the repair in `detail`); attaches the task contract; and preserves stable priority order. Assignment runs on the same locked state and skips reaping slots, budget-paused work, closed project roots, conflicting project writers, and tasks exceeding the root's subagent capacity or depth reservation. Evolution alone is fenced by runtime mode — blocked in Light — three times (`supervisor/evolution_lifecycle.py`): entry points refuse a campaign start, `enqueue_evolution_task_if_needed()` pauses and disables a carried campaign, and assignment drops what slipped through while `evolution_block_reason()` is set; generic `supervisor.queue.enqueue_task()` has no runtime-mode predicate. Configured worker count is therefore not available capacity: the truthful value is the assignable idle count after custody, reaping and admission fences.
|
||||
|
||||
A headless task is ADDRESSED when it is admitted, not when it is displayed (`log_addressing.ingress_chat_id`). A registered project's run has exactly ONE destination: an explicit `chat_id` may only agree with that thread, and any other value — the hidden partition included — is refused with a typed 400 rather than honoured or silently overridden, because a run addressed away from its room puts a card in Main whose project holds none of its work. Without a Project, ordinary API tasks default to `HIDDEN_CHAT_ID` (0); the confirmed browser Publish flow carries `source="web"` and `WEB_UI_CHAT_ID` to request Main (caller-declared addressing, not an authentication proof), and other non-Project conversation addresses stay refused. A run scoped to a REGISTERED, active project is admitted into that project's thread, and Main receives the one host-stamped completion row only when the work is actually in that room: addressed there at admission or BOUND to the project. Registration alone does not qualify. Every other run stays in the hidden partition, silent in every chat, read back through the terminal, `--result-json-out`, the chat-blind Logs panel and `GET /api/tasks/<id>`. A reserved but inactive project keeps its chat acceptable so the queue's lifecycle fence refuses with its own typed reason.
|
||||
|
|
|
|||
|
|
@ -2,7 +2,7 @@
|
|||
|
||||
AST-derived inventory of compatibility facades, regenerated by `python scripts/regenerate_inventories.py`. Do not edit. A facade row is any runtime module whose top-level `from <population module> import ...` statements carry the `noqa: F401` re-export marker — the codebase's declared "this binding exists for its binding, not for this module's own use" convention (reference FACADE_CONSUMERS method). Leaf domains come from `ouroboros/domains.toml`; a leaf outside the facade's domain is marked ✗ (that edge also appears in the manifest's pinned direction matrix). `tests/test_generated_inventories.py` pins byte-identity, so any re-export surface change must regenerate this file.
|
||||
|
||||
- facade modules: **58**; marked re-export bindings: **2284**; cross-domain facade→leaf pairs: **131**
|
||||
- facade modules: **58**; marked re-export bindings: **2285**; cross-domain facade→leaf pairs: **131**
|
||||
|
||||
| facade | domain | bindings | leaves |
|
||||
|---|---|---:|---|
|
||||
|
|
@ -54,7 +54,7 @@ AST-derived inventory of compatibility facades, regenerated by `python scripts/r
|
|||
| `ouroboros/tools/shell_guards.py` | D13 | 14 | `ouroboros/tools/write_shape.py` (14) |
|
||||
| `ouroboros/tools/shell_process.py` | D05 | 1 | `ouroboros/tools/process_facts.py` (1 ✗D04) |
|
||||
| `ouroboros/tools/subagent_integration.py` | D07 | 13 | `ouroboros/headless.py` (2 ✗D17)<br>`ouroboros/tools/subagent_integration_delegated.py` (11) |
|
||||
| `ouroboros/usage_accounting.py` | D16 | 42 | `ouroboros/_usage_cache_splits.py` (4)<br>`ouroboros/_usage_rows.py` (6)<br>`ouroboros/_usage_rows_memo.py` (6)<br>`ouroboros/usage_ledger.py` (18)<br>`ouroboros/usage_legacy_import.py` (5)<br>`ouroboros/utils.py` (3 ✗D18) |
|
||||
| `ouroboros/usage_accounting.py` | D16 | 43 | `ouroboros/_usage_cache_splits.py` (4)<br>`ouroboros/_usage_rows.py` (7)<br>`ouroboros/_usage_rows_memo.py` (6)<br>`ouroboros/usage_ledger.py` (18)<br>`ouroboros/usage_legacy_import.py` (5)<br>`ouroboros/utils.py` (3 ✗D18) |
|
||||
| `server.py` | D11 | 51 | `ouroboros/server_liveness.py` (4)<br>`ouroboros/server_maintenance.py` (14)<br>`ouroboros/server_owner_routing.py` (5)<br>`ouroboros/server_process.py` (6)<br>`ouroboros/server_restart.py` (8)<br>`ouroboros/server_routing_context.py` (14) |
|
||||
| `supervisor/events.py` | D08 | 94 | `ouroboros/config.py` (1 ✗D12)<br>`ouroboros/contracts/task_constraint.py` (1 ✗D19)<br>`ouroboros/cost_projection.py` (3 ✗D16)<br>`ouroboros/subagent_messages.py` (1 ✗D07)<br>`ouroboros/task_results.py` (2 ✗D17)<br>`ouroboros/tool_capabilities.py` (2 ✗D04)<br>`ouroboros/utils.py` (4 ✗D18)<br>`supervisor/cognitive_operations.py` (2)<br>`supervisor/events_budget.py` (4)<br>`supervisor/events_chat_delivery.py` (9)<br>`supervisor/events_coop_checkpoint.py` (6)<br>`supervisor/events_evolution_done.py` (1)<br>`supervisor/events_project_routing.py` (10)<br>`supervisor/events_runtime_controls.py` (6)<br>`supervisor/events_schedule_task.py` (4)<br>`supervisor/events_subagent_admission.py` (15)<br>`supervisor/events_task_done.py` (8)<br>`supervisor/events_worker_reports.py` (7)<br>`supervisor/log_addressing.py` (4)<br>`supervisor/queue_transitions.py` (1)<br>`supervisor/steering.py` (2 ✗D09)<br>`supervisor/task_dispatch.py` (1) |
|
||||
| `supervisor/git_ops.py` | D10 | 37 | `ouroboros/utils.py` (1 ✗D18)<br>`supervisor/git_ops_remotes.py` (4)<br>`supervisor/git_ops_rescue.py` (8)<br>`supervisor/git_ops_reset.py` (10)<br>`supervisor/git_ops_updates.py` (8)<br>`supervisor/state.py` (4 ✗D08)<br>`supervisor/update_recovery.py` (2) |
|
||||
|
|
|
|||
|
|
@ -200,6 +200,7 @@ def _summary(rows: Sequence[Dict[str, Any]]) -> Dict[str, Any]:
|
|||
|
||||
|
||||
def _with_limit(summary: Dict[str, Any], limit: Optional[float]) -> Dict[str, Any]:
|
||||
"""Decorate a summary with its configured limit and remaining headroom."""
|
||||
if limit is None:
|
||||
return summary
|
||||
summary["limit_usd"] = round(max(0.0, float(limit)), 6)
|
||||
|
|
@ -215,6 +216,36 @@ def _with_integrity(summary: Dict[str, Any], degraded: bool) -> Dict[str, Any]:
|
|||
return summary
|
||||
|
||||
|
||||
def _marker_from_final(final: Sequence[Dict[str, Any]]) -> Optional[list]:
|
||||
"""The ordered ``[compaction_epoch, seq]`` fact of these validated rows.
|
||||
|
||||
Compaction advances the epoch in its leading ``usage_baseline`` header
|
||||
while renumbering live rows, so the PAIR stays ordered even when the file
|
||||
gets shorter. ``None`` means unknown ordering — never zero — so a
|
||||
compatibility writer can fail safe instead of writing money it cannot
|
||||
place in time.
|
||||
"""
|
||||
try:
|
||||
baselines = [row for row in final
|
||||
if isinstance(row, dict) and str(row.get("kind") or "") == "usage_baseline"]
|
||||
if len(baselines) > 1:
|
||||
raise ValueError("multiple usage baseline headers")
|
||||
epoch = baselines[0].get("compaction_epoch", 0) if baselines else 0
|
||||
if isinstance(epoch, bool) or not isinstance(epoch, int) or epoch < 0:
|
||||
raise ValueError("invalid usage baseline compaction epoch")
|
||||
seqs = []
|
||||
for row in final:
|
||||
value = row.get("seq") if isinstance(row, dict) else None
|
||||
if isinstance(value, bool) or not isinstance(value, int) or value < 0:
|
||||
raise ValueError("invalid usage ledger sequence marker")
|
||||
seqs.append(value)
|
||||
return [epoch, max(seqs, default=0)]
|
||||
except (OSError, TypeError, ValueError):
|
||||
return None
|
||||
|
||||
|
||||
|
||||
|
||||
def _physical_call_count(row: Dict[str, Any]) -> int:
|
||||
kind = str(row.get("kind") or "attempt")
|
||||
if kind == "legacy_metadata":
|
||||
|
|
|
|||
|
|
@ -152,7 +152,15 @@ def _state_snapshot(request: Request) -> Dict[str, Any]:
|
|||
accounting_available = True
|
||||
try:
|
||||
ensure_legacy_imported(drive_root)
|
||||
breakdown = usage_breakdown(drive_root)
|
||||
# ``usage_breakdown`` also carries private provenance used by the
|
||||
# compatibility writer. Keep that internal vocabulary at this
|
||||
# boundary even when the unbounded-budget branch reuses the mapping
|
||||
# directly as its accounting projection.
|
||||
breakdown = {
|
||||
key: value
|
||||
for key, value in usage_breakdown(drive_root).items()
|
||||
if not str(key).startswith("_")
|
||||
}
|
||||
# include_roots=False: /api/state serializes named scalars only, so the
|
||||
# per-root map would be built per poll and thrown away (O(N×roots) work
|
||||
# with zero readers on this path). The slim projection still carries
|
||||
|
|
|
|||
|
|
@ -168,6 +168,7 @@ BAND_PATHS = {
|
|||
"skills/telegram/scripts/sidecar.py": None,
|
||||
"supervisor/message_bus.py": "push_log gained the A2A frame suppression and explicit-addressing contract comments in the live log routing fix (fix/main-chat-leak sprint); the module sat at exactly 1000 lines before it",
|
||||
"supervisor/queue_transitions.py": "F2.2 cancel/custody organ: the owner-stop campaign closure (_close_campaign_after_owner_stop, reference row 970) moved in from the hot events monolith to live beside stop_evolution_tasks - one honesty rule, one owner; 1016 lines, shrink-only from here",
|
||||
"supervisor/state.py": "Entered the band from 967 lines with the STATE_LOCK read-inversion fix (issue #1102): the ledger snapshot is now read OUTSIDE the lock, so the freshness-marker comparison deciding whether a snapshot may overwrite money must sit with the atomic state write it guards. One decision owner - lock, marker comparison, projection write - instead of a second budget authority beside it; its rationale comments are why this regression class stays visible. Shrink-only from here.",
|
||||
"supervisor/update_merge.py": "Entered the band from above (1593 lines) by extraction: the F2.4 update-engine re-split moved the planner, the clean-plan commit builder and the live materializer \u2014 the carrier engine's three insertion points \u2014 into supervisor/update_merge_plan.py (D34 return, owner answers 5.12-5.14=A); shrink-only.",
|
||||
"tests/system_e2e/harness.py": "system_e2e harness: waves 3a+3b grew the one scenario-suite machinery module into the band \u2014 skill-review stub branch, review-organ verdict scripting (ReviewScript), plan-review/native-episode classification markers, the advisory reviewer-slot row and the S11-S17 manifest rows; split when the next wave lands new actors.",
|
||||
"tests/system_e2e/test_system_scenarios_w4.py": "system_e2e wave-4 scenario module: six scenarios (S18-S23 - update carrier/conflict/crash variants, chat-lineage cancel, absorb kill-recovery, delegated interactive answer) plus the interactive fake-daemon contract pin; one module per wave is the suite convention - split only if a later wave extends THIS module instead of adding its own.",
|
||||
|
|
|
|||
|
|
@ -52,6 +52,7 @@ from ouroboros.utils import append_jsonl, atomic_write_json, utc_now_iso # noqa
|
|||
from ouroboros._usage_rows import ( # noqa: F401 (re-exported substrate vocabulary)
|
||||
REVIEW_ATTRIBUTION_KEYS,
|
||||
_breakdown_bucket,
|
||||
_marker_from_final,
|
||||
_physical_call_count,
|
||||
_summary,
|
||||
_with_integrity,
|
||||
|
|
@ -452,6 +453,35 @@ from ouroboros._usage_rows_memo import ( # noqa: F401,E402 (re-exported seam)
|
|||
)
|
||||
|
||||
|
||||
def _projection_from_final(
|
||||
final: list, integrity_degraded: bool, configured_limit: Optional[float] = None,
|
||||
*, root_task_id: str = "", include_roots: bool = True,
|
||||
) -> Dict[str, Any]:
|
||||
"""Render the money projection from ALREADY-VALIDATED final rows: one
|
||||
snapshot, one projection, so a caller deriving the ordering marker from
|
||||
the SAME rows writes both under one authority instead of pairing a marker
|
||||
with a second, later ledger read."""
|
||||
def limit_of(rows: list) -> Optional[float]:
|
||||
known = [v for v in (_number(row.get("root_limit_usd")) for row in rows) if v is not None]
|
||||
return min(known) if known else None
|
||||
if root_task_id:
|
||||
rows = [row for row in final if str(row.get("root_task_id") or "") == root_task_id]
|
||||
return _with_integrity(_with_limit(_summary(rows), limit_of(rows)), integrity_degraded)
|
||||
result = _with_limit(_summary(final), configured_limit)
|
||||
if include_roots:
|
||||
grouped: Dict[str, list] = {}
|
||||
for row in final:
|
||||
rid = str(row.get("root_task_id") or "")
|
||||
if rid:
|
||||
grouped.setdefault(rid, []).append(row)
|
||||
result["by_root"] = {
|
||||
rid: _with_integrity(_with_limit(_summary(grouped[rid]), limit_of(grouped[rid])),
|
||||
integrity_degraded)
|
||||
for rid in sorted(grouped)
|
||||
}
|
||||
return _with_integrity(result, integrity_degraded)
|
||||
|
||||
|
||||
def usage_projection(
|
||||
drive_root: pathlib.Path | str | None = None,
|
||||
*,
|
||||
|
|
@ -460,61 +490,24 @@ def usage_projection(
|
|||
include_roots: bool = True,
|
||||
) -> Dict[str, Any]:
|
||||
"""Return a replayed global projection, or one root/subtree projection.
|
||||
|
||||
``include_roots=False`` skips building the per-root ``by_root`` map for
|
||||
hot-path readers that never consume it (``/api/state``); the slim result
|
||||
still carries ``limit_usd``/``remaining_known_usd`` — the two fields
|
||||
``budget_remaining`` consumes. The default keeps the full contract."""
|
||||
``include_roots=False`` skips the per-root ``by_root`` map for hot-path
|
||||
readers that never consume it (``/api/state``); the slim result still
|
||||
carries the two fields ``budget_remaining`` reads."""
|
||||
root = _drive_root(drive_root)
|
||||
if root_task_id:
|
||||
cache_key = ("usage_projection", root_task_id, "", None, True)
|
||||
|
||||
def render_root(final: list, integrity_degraded: bool) -> Dict[str, Any]:
|
||||
rows = [row for row in final if str(row.get("root_task_id") or "") == root_task_id]
|
||||
limits = [_number(row.get("root_limit_usd")) for row in rows]
|
||||
known_limits = [value for value in limits if value is not None]
|
||||
result = _with_limit(_summary(rows), min(known_limits) if known_limits else None)
|
||||
return _with_integrity(result, integrity_degraded)
|
||||
|
||||
return _render_cached(root, cache_key, render_root)
|
||||
return _render_cached(
|
||||
root, ("usage_projection", root_task_id, "", None, True),
|
||||
lambda f, degraded: _projection_from_final(f, degraded, root_task_id=root_task_id))
|
||||
if global_limit_usd is not None:
|
||||
configured_limit = max(0.0, float(global_limit_usd))
|
||||
else:
|
||||
from ouroboros.settings_setup_contract import resolve_total_budget_usd
|
||||
|
||||
configured_limit = resolve_total_budget_usd() or 0.0
|
||||
apply_limit = global_limit_usd is not None or configured_limit > 0
|
||||
cache_key = (
|
||||
"usage_projection", "", "",
|
||||
configured_limit if apply_limit else None,
|
||||
include_roots,
|
||||
)
|
||||
|
||||
def render_global(final: list, integrity_degraded: bool) -> Dict[str, Any]:
|
||||
result = (
|
||||
_with_limit(_summary(final), configured_limit) if apply_limit else _summary(final)
|
||||
)
|
||||
if include_roots:
|
||||
grouped_rows: Dict[str, list] = {}
|
||||
for row in final:
|
||||
rid = str(row.get("root_task_id") or "")
|
||||
if rid:
|
||||
grouped_rows.setdefault(rid, []).append(row)
|
||||
result["by_root"] = {}
|
||||
for rid in sorted(grouped_rows):
|
||||
root_rows = grouped_rows[rid]
|
||||
known_limits = [
|
||||
value
|
||||
for value in (_number(row.get("root_limit_usd")) for row in root_rows)
|
||||
if value is not None
|
||||
]
|
||||
result["by_root"][rid] = _with_integrity(
|
||||
_with_limit(_summary(root_rows), min(known_limits) if known_limits else None),
|
||||
integrity_degraded,
|
||||
)
|
||||
return _with_integrity(result, integrity_degraded)
|
||||
|
||||
return _render_cached(root, cache_key, render_global)
|
||||
limit = configured_limit if (global_limit_usd is not None or configured_limit > 0) else None
|
||||
return _render_cached(
|
||||
root, ("usage_projection", "", "", limit, include_roots),
|
||||
lambda final, degraded: _projection_from_final(final, degraded, limit,
|
||||
include_roots=include_roots))
|
||||
|
||||
|
||||
def usage_breakdown(
|
||||
|
|
@ -523,11 +516,18 @@ def usage_breakdown(
|
|||
root_task_id: str = "",
|
||||
task_id: str = "",
|
||||
) -> Dict[str, Any]:
|
||||
"""Read-only physical-call/token/cost buckets from validated ledger finals."""
|
||||
"""Read-only physical-call/token/cost buckets from validated ledger finals.
|
||||
Both private compatibility fields — the ordered ``[compaction_epoch, seq]``
|
||||
marker in ``_ledger_high_water_seq`` and the money projection in
|
||||
``_usage_projection`` — are rendered from THIS one validated read, so a
|
||||
writer authorizes its projection with the marker of the same snapshot."""
|
||||
root = _drive_root(drive_root)
|
||||
cache_key = ("usage_breakdown", root_task_id, task_id, None, True)
|
||||
|
||||
def render(final: list, integrity_degraded: bool) -> Dict[str, Any]:
|
||||
# Marker and money are two renderings of THESE rows; unknown marker =>
|
||||
# buckets stay readable while a compatibility writer fails safe.
|
||||
ledger_marker = _marker_from_final(final)
|
||||
rows = final
|
||||
if root_task_id:
|
||||
rows = [row for row in rows if str(row.get("root_task_id") or "") == root_task_id]
|
||||
|
|
@ -556,20 +556,20 @@ def usage_breakdown(
|
|||
|
||||
result = {
|
||||
**_with_integrity(_breakdown_bucket(rows), integrity_degraded),
|
||||
"_ledger_high_water_seq": ledger_marker,
|
||||
"by_model": by_model,
|
||||
"by_provider": by_provider,
|
||||
"by_category": by_category,
|
||||
"by_task": by_task,
|
||||
"by_root": by_root,
|
||||
# Execution-axis filter (v6.91): delegated (subscription-harness) rows only — a
|
||||
# VIEW for "where did the money go" readers, never a third monetary sum or
|
||||
# authority; disclosed-free settles $0, undisclosed stays `unknown`.
|
||||
# v6.91 execution-axis VIEW of delegated (subscription-harness) rows for
|
||||
# "where did the money go" readers; never a third monetary sum or authority.
|
||||
"delegated": _with_integrity(
|
||||
_breakdown_bucket([row for row in rows if str(row.get("kind") or "") == "subscription_session"]),
|
||||
integrity_degraded,
|
||||
),
|
||||
# Legacy call-count metadata and monetary delta stay explicit; neither
|
||||
# is fabricated into a model/provider/category identity.
|
||||
# Legacy call-count metadata and monetary delta stay explicit; neither is
|
||||
# fabricated into a model/provider/category identity.
|
||||
"unattributed": {
|
||||
"model": model_unattributed,
|
||||
"provider": provider_unattributed,
|
||||
|
|
@ -585,6 +585,7 @@ def usage_breakdown(
|
|||
):
|
||||
for bucket in grouped_buckets.values():
|
||||
_with_integrity(bucket, True)
|
||||
result["_usage_projection"] = _projection_from_final(final, integrity_degraded)
|
||||
return result
|
||||
|
||||
return _render_cached(root, cache_key, render)
|
||||
|
|
@ -605,7 +606,6 @@ def _reservation_cost(request: AttemptRequest) -> Optional[float]:
|
|||
prompt_tokens = max(0, int(request.prompt_tokens_estimate or 0))
|
||||
# OpenAI-family chars/4 estimates keep the measured 1.10 reservation envelope.
|
||||
from ouroboros.provider_models import normalize_model_identity
|
||||
|
||||
normalized_model = normalize_model_identity(str(request.model or "").lstrip("~"))
|
||||
if (
|
||||
str(request.provider or "").strip().lower() in {"openai", "openrouter"}
|
||||
|
|
|
|||
|
|
@ -106,7 +106,8 @@ def _handle_llm_usage(evt: Dict[str, Any], ctx: Any) -> None:
|
|||
}
|
||||
projection_update_status = "available"
|
||||
try:
|
||||
ctx.update_budget_from_usage(usage_for_budget)
|
||||
if ctx.update_budget_from_usage(usage_for_budget) is False:
|
||||
projection_update_status = "unavailable"
|
||||
except Exception:
|
||||
projection_update_status = "unavailable"
|
||||
log.error("Paid llm_usage retained but compatibility projection update failed", exc_info=True)
|
||||
|
|
|
|||
|
|
@ -469,7 +469,7 @@ def budget_pct(st: Dict[str, Any]) -> float:
|
|||
return 100.0
|
||||
|
||||
|
||||
def update_budget_from_usage(usage: Dict[str, Any]) -> None:
|
||||
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
|
||||
|
|
@ -491,42 +491,128 @@ def update_budget_from_usage(usage: Dict[str, Any]) -> None:
|
|||
log.debug(f"Failed to convert value to int: {v!r}", exc_info=True)
|
||||
return default
|
||||
|
||||
from ouroboros.usage_accounting import ensure_legacy_imported, usage_breakdown, usage_projection
|
||||
def _ledger_high_water_marker(breakdown: Dict[str, Any]) -> Optional[tuple[int, int]]:
|
||||
"""Return the ``(compaction_epoch, seq)`` fact from this read.
|
||||
|
||||
# Serialize the ledger snapshot with its compatibility write. Otherwise an
|
||||
# older concurrent reader can acquire STATE_LOCK later and regress state.json.
|
||||
lock_fd = acquire_file_lock(STATE_LOCK_PATH)
|
||||
_warn_state_unlocked("budget-update", lock_fd)
|
||||
``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_breakdown,
|
||||
usage_projection,
|
||||
)
|
||||
|
||||
# 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.
|
||||
try:
|
||||
ensure_legacy_imported(DRIVE_ROOT)
|
||||
breakdown = usage_breakdown(DRIVE_ROOT)
|
||||
total_limit = float(TOTAL_BUDGET_LIMIT or 0.0)
|
||||
projection = (
|
||||
usage_projection(DRIVE_ROOT, global_limit_usd=total_limit)
|
||||
if total_limit > 0
|
||||
else {key: breakdown.get(key) for key in (
|
||||
projection_snapshot = breakdown.pop("_usage_projection", None)
|
||||
if total_limit > 0 and isinstance(projection_snapshot, dict):
|
||||
from ouroboros._usage_rows import _with_limit
|
||||
roots = projection_snapshot.pop("by_root", None)
|
||||
projection = _with_limit(projection_snapshot, total_limit)
|
||||
if roots is not None:
|
||||
projection["by_root"] = roots
|
||||
else:
|
||||
projection = (
|
||||
usage_projection(DRIVE_ROOT, global_limit_usd=total_limit)
|
||||
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)
|
||||
_warn_state_unlocked("budget-update", lock_fd)
|
||||
try:
|
||||
st = _load_state_unlocked()
|
||||
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(breakdown)
|
||||
previous_check_call = _to_int(st.get("openrouter_last_check_call"), -1)
|
||||
should_check_ground_truth = bool(
|
||||
st["spent_calls"] > 0
|
||||
and st["spent_calls"] % 50 == 0
|
||||
and st["spent_calls"] != previous_check_call
|
||||
previous_marker = st.get("usage_ledger_high_water_seq")
|
||||
marker_known = (
|
||||
ledger_high_water_marker is not None
|
||||
)
|
||||
if should_check_ground_truth:
|
||||
st["openrouter_last_check_call"] = st["spent_calls"]
|
||||
_save_state_unlocked(st)
|
||||
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
|
||||
marker_to_write: Optional[tuple[int, int]] = None
|
||||
if not marker_known 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
|
||||
elif previous_known:
|
||||
current_marker = ledger_high_water_marker
|
||||
saved_marker = (int(previous_marker[0]), int(previous_marker[1]))
|
||||
if current_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",
|
||||
current_marker,
|
||||
saved_marker,
|
||||
)
|
||||
return False
|
||||
marker_to_write = current_marker
|
||||
else:
|
||||
marker_to_write = ledger_high_water_marker
|
||||
|
||||
if marker_to_write is not None:
|
||||
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(marker_to_write)
|
||||
previous_check_call = _to_int(st.get("openrouter_last_check_call"), -1)
|
||||
should_check_ground_truth = bool(
|
||||
st["spent_calls"] > 0
|
||||
and st["spent_calls"] % 50 == 0
|
||||
and st["spent_calls"] != previous_check_call
|
||||
)
|
||||
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)
|
||||
|
||||
|
|
@ -540,10 +626,10 @@ def update_budget_from_usage(usage: Dict[str, Any]) -> None:
|
|||
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 OpenRouter key's usage delta — the only like-for-like
|
||||
# pair. Direct-provider spend (Anthropic/OpenAI/local) is invisible
|
||||
# to /auth/key by construction and must not count as "drift".
|
||||
# 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")
|
||||
|
|
@ -553,9 +639,9 @@ def update_budget_from_usage(usage: Dict[str, Any]) -> None:
|
|||
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.
|
||||
# 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 (
|
||||
|
|
@ -589,9 +675,9 @@ def update_budget_from_usage(usage: Dict[str, Any]) -> None:
|
|||
DRIVE_ROOT / "logs" / "events.jsonl",
|
||||
{
|
||||
"ts": utc_now_iso(),
|
||||
# "type" is the events.jsonl schema key every other
|
||||
# event uses; type-keyed aggregations used to lose
|
||||
# this row when it was written under "event".
|
||||
# "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),
|
||||
|
|
@ -615,6 +701,8 @@ def update_budget_from_usage(usage: Dict[str, Any]) -> None:
|
|||
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."""
|
||||
|
|
|
|||
|
|
@ -290,6 +290,7 @@ class TestBudgetDriftOpenRouterOnly:
|
|||
"attempt_counts": {"settled": calls},
|
||||
"integrity_degraded": integrity_degraded,
|
||||
"by_provider": {"openrouter": {"settled_usd": or_settled}},
|
||||
"_ledger_high_water_seq": [0, calls],
|
||||
}
|
||||
import ouroboros.usage_accounting as ua
|
||||
|
||||
|
|
|
|||
|
|
@ -101,6 +101,57 @@ def test_llm_usage_preserves_unknown_cost_as_null(tmp_path):
|
|||
assert ctx.last_usage["cost"] is None
|
||||
|
||||
|
||||
def test_llm_usage_reports_corrupt_projection_unavailable_and_keeps_paid_usage(tmp_path):
|
||||
from supervisor import events as ev_module
|
||||
(tmp_path / "logs").mkdir()
|
||||
|
||||
class FakeCtx:
|
||||
DRIVE_ROOT = tmp_path
|
||||
|
||||
def update_budget_from_usage(self, usage):
|
||||
self.last_usage = usage
|
||||
return False
|
||||
|
||||
ctx = FakeCtx()
|
||||
ev_module._handle_llm_usage(
|
||||
{"type": "llm_usage", "task_id": "paid", "usage": {"prompt_tokens": 4, "cost": 0.75}},
|
||||
ctx,
|
||||
)
|
||||
written = json.loads((tmp_path / "logs" / "events.jsonl").read_text(encoding="utf-8"))
|
||||
assert written["projection_update_status"] == "unavailable"
|
||||
assert written["cost"] == 0.75
|
||||
assert ctx.last_usage["cost"] == 0.75
|
||||
|
||||
|
||||
def test_llm_usage_real_corrupt_ledger_is_unavailable_and_paid_event_survives(tmp_path):
|
||||
from supervisor import events as ev_module
|
||||
from supervisor import state
|
||||
from ouroboros.usage_ledger import LEDGER_REL
|
||||
|
||||
(tmp_path / "logs").mkdir()
|
||||
state.init(tmp_path, total_budget_limit=0.0)
|
||||
state.save_state({"spent_usd": 1.25})
|
||||
ledger = tmp_path / LEDGER_REL
|
||||
ledger.parent.mkdir(parents=True, exist_ok=True)
|
||||
ledger.write_text("{broken}\n{}\n", encoding="utf-8")
|
||||
|
||||
class Ctx:
|
||||
DRIVE_ROOT = tmp_path
|
||||
|
||||
@staticmethod
|
||||
def update_budget_from_usage(usage):
|
||||
return state.update_budget_from_usage(usage)
|
||||
|
||||
ev_module._handle_llm_usage(
|
||||
{"type": "llm_usage", "task_id": "paid-real", "usage": {"prompt_tokens": 2, "cost": 0.5}},
|
||||
Ctx(),
|
||||
)
|
||||
written = json.loads((tmp_path / "logs" / "events.jsonl").read_text(encoding="utf-8"))
|
||||
assert written["projection_update_status"] == "unavailable"
|
||||
assert written["cost"] == 0.5
|
||||
assert state.load_state()["spent_usd"] == 1.25
|
||||
|
||||
|
||||
def test_cost_breakdown_aggregates_cache_tokens_and_ttl(tmp_path):
|
||||
import asyncio
|
||||
import json
|
||||
|
|
|
|||
|
|
@ -194,6 +194,69 @@ def test_api_state_money_and_call_count_are_ledger_projections(tmp_path, monkeyp
|
|||
assert payload["accounting"]["remaining_known_usd"] == 5.75
|
||||
|
||||
|
||||
@pytest.mark.serial
|
||||
def test_api_state_unbounded_budget_does_not_expose_private_breakdown_keys(tmp_path, monkeypatch):
|
||||
"""The no-limit branch keeps internal ledger provenance off the wire."""
|
||||
from ouroboros.gateway.state import api_state
|
||||
from supervisor import queue, state, workers
|
||||
|
||||
root = _data_root(tmp_path, monkeypatch)
|
||||
monkeypatch.setattr(state, "TOTAL_BUDGET_LIMIT", 0.0)
|
||||
monkeypatch.setattr(state, "load_state", lambda: {"current_branch": "ouroboros"})
|
||||
monkeypatch.setattr(workers, "WORKERS", {})
|
||||
monkeypatch.setattr(workers, "PENDING", [])
|
||||
monkeypatch.setattr(workers, "RUNNING", {})
|
||||
monkeypatch.setattr(queue, "get_evolution_status_snapshot", lambda **_kwargs: {})
|
||||
monkeypatch.setattr(
|
||||
ua,
|
||||
"usage_breakdown",
|
||||
lambda *_args, **_kwargs: {
|
||||
"accounted_usd": 2.0,
|
||||
"physical_calls": 1,
|
||||
"settled_usd": 2.0,
|
||||
"confirmed_usd": 2.0,
|
||||
"estimated_usd": 0.0,
|
||||
"reserved_usd": 0.0,
|
||||
"unresolved_upper_bound_usd": 0.0,
|
||||
"unknown_unmetered": 0,
|
||||
"cost_final": True,
|
||||
"attempt_counts": {},
|
||||
"integrity_degraded": False,
|
||||
"_ledger_high_water_seq": [0, 7],
|
||||
"_root_accounting_snapshot": {"accounted_usd": 2.0},
|
||||
},
|
||||
)
|
||||
monkeypatch.setattr(
|
||||
ua,
|
||||
"usage_projection",
|
||||
lambda *_args, **_kwargs: pytest.fail("unbounded /api/state must not call usage_projection"),
|
||||
)
|
||||
request = Request({
|
||||
"type": "http", "method": "GET", "path": "/api/state", "headers": [],
|
||||
"query_string": b"", "scheme": "http", "server": ("test", 80),
|
||||
"client": ("test", 1),
|
||||
"app": types.SimpleNamespace(state=types.SimpleNamespace(drive_root=root, app_start=0.0)),
|
||||
})
|
||||
|
||||
response = asyncio.run(api_state(request))
|
||||
payload = json.loads(response.body)
|
||||
|
||||
assert response.status_code == 200
|
||||
assert payload["budget_limit"] == 0.0
|
||||
|
||||
def walk_keys(value):
|
||||
if isinstance(value, dict):
|
||||
for key, child in value.items():
|
||||
yield key
|
||||
yield from walk_keys(child)
|
||||
elif isinstance(value, list):
|
||||
for child in value:
|
||||
yield from walk_keys(child)
|
||||
|
||||
assert all(not str(key).startswith("_") for key in walk_keys(payload))
|
||||
assert payload["accounting"]["accounted_usd"] == 2.0
|
||||
|
||||
|
||||
def test_cost_breakdown_fails_loudly_when_authoritative_history_is_corrupt(tmp_path, monkeypatch):
|
||||
from ouroboros.gateway.history import make_cost_breakdown_endpoint
|
||||
|
||||
|
|
|
|||
353
tests/test_state_budget_projection_lock_order.py
Normal file
353
tests/test_state_budget_projection_lock_order.py
Normal file
|
|
@ -0,0 +1,353 @@
|
|||
"""Deterministic ordering and freshness tests for the legacy budget projection."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import threading
|
||||
import time
|
||||
import json
|
||||
|
||||
import pytest
|
||||
|
||||
|
||||
def _breakdown(value: float, marker: int) -> dict:
|
||||
return {
|
||||
"accounted_usd": value,
|
||||
"physical_calls": int(value),
|
||||
"prompt_tokens": int(value),
|
||||
"completion_tokens": 0,
|
||||
"cached_tokens": 0,
|
||||
"settled_usd": value,
|
||||
"confirmed_usd": value,
|
||||
"estimated_usd": 0.0,
|
||||
"reserved_usd": 0.0,
|
||||
"unresolved_upper_bound_usd": 0.0,
|
||||
"unknown_unmetered": 0,
|
||||
"cost_final": True,
|
||||
"attempt_counts": {"settled": int(value)},
|
||||
"integrity_degraded": False,
|
||||
"_ledger_high_water_seq": [0, marker],
|
||||
}
|
||||
|
||||
|
||||
@pytest.mark.serial
|
||||
def test_budget_projection_reads_ledger_before_state_lock(tmp_path, monkeypatch):
|
||||
from supervisor import state
|
||||
import ouroboros.usage_accounting as accounting
|
||||
|
||||
state.init(tmp_path, total_budget_limit=1.0)
|
||||
holding_state_lock = False
|
||||
operations = []
|
||||
|
||||
def acquire(_path, *args, **kwargs):
|
||||
nonlocal holding_state_lock
|
||||
assert not holding_state_lock
|
||||
holding_state_lock = True
|
||||
operations.append("acquire")
|
||||
return 1
|
||||
|
||||
def release(_path, _fd):
|
||||
nonlocal holding_state_lock
|
||||
assert holding_state_lock
|
||||
holding_state_lock = False
|
||||
operations.append("release")
|
||||
|
||||
def read_ledger(*_args, **_kwargs):
|
||||
assert not holding_state_lock, "ledger I/O must precede STATE_LOCK"
|
||||
operations.append("ledger")
|
||||
|
||||
monkeypatch.setattr(state, "acquire_file_lock", acquire)
|
||||
monkeypatch.setattr(state, "release_file_lock", release)
|
||||
monkeypatch.setattr(state, "_load_state_unlocked", lambda: (operations.append("load") or {}))
|
||||
monkeypatch.setattr(
|
||||
state,
|
||||
"_save_state_unlocked",
|
||||
lambda _st: (operations.append("save"), assert_not_holding(holding_state_lock)),
|
||||
)
|
||||
monkeypatch.setattr(state, "_openrouter_ledger_settled", lambda *_a, **_k: 1.0)
|
||||
monkeypatch.setattr(accounting, "ensure_legacy_imported", read_ledger)
|
||||
monkeypatch.setattr(accounting, "usage_breakdown", lambda *_a, **_k: (read_ledger() or _breakdown(1.0, 1)))
|
||||
monkeypatch.setattr(accounting, "usage_projection", lambda *_a, **_k: (read_ledger() or {"accounted_usd": 1.0}))
|
||||
|
||||
state.update_budget_from_usage({})
|
||||
|
||||
assert operations[:3] == ["ledger", "ledger", "ledger"]
|
||||
assert operations[3:] == ["acquire", "load", "save", "release"]
|
||||
assert not holding_state_lock
|
||||
|
||||
|
||||
def assert_not_holding(value):
|
||||
assert value, "state load/mutate/save must stay inside STATE_LOCK"
|
||||
|
||||
|
||||
@pytest.mark.serial
|
||||
def test_older_budget_snapshot_cannot_regress_state(tmp_path, monkeypatch):
|
||||
from supervisor import state
|
||||
import ouroboros.usage_accounting as accounting
|
||||
|
||||
state.init(tmp_path, total_budget_limit=1.0)
|
||||
first_started = threading.Event()
|
||||
release_first = threading.Event()
|
||||
calls = []
|
||||
|
||||
def breakdown(_root):
|
||||
calls.append(len(calls) + 1)
|
||||
if len(calls) == 1:
|
||||
first_started.set()
|
||||
assert release_first.wait(2.0)
|
||||
value = 1.0
|
||||
else:
|
||||
value = 2.0
|
||||
snapshot = _breakdown(value, int(value))
|
||||
snapshot["_usage_projection"] = {
|
||||
"accounted_usd": value, "integrity_degraded": False, "cost_final": True,
|
||||
}
|
||||
return snapshot
|
||||
|
||||
monkeypatch.setattr(accounting, "ensure_legacy_imported", lambda *_a, **_k: None)
|
||||
monkeypatch.setattr(accounting, "usage_breakdown", breakdown)
|
||||
older = threading.Thread(target=state.update_budget_from_usage, args=({},))
|
||||
newer = threading.Thread(target=state.update_budget_from_usage, args=({},))
|
||||
older.start()
|
||||
assert first_started.wait(2.0)
|
||||
newer.start()
|
||||
newer.join(2.0)
|
||||
release_first.set()
|
||||
older.join(2.0)
|
||||
|
||||
assert state.load_state()["spent_usd"] == 2.0
|
||||
assert state.load_state()["usage_ledger_high_water_seq"] == [0, 2]
|
||||
|
||||
|
||||
@pytest.mark.serial
|
||||
def test_limited_projection_uses_breakdown_snapshot(tmp_path, monkeypatch):
|
||||
from supervisor import state
|
||||
import ouroboros.usage_accounting as accounting
|
||||
|
||||
state.init(tmp_path, total_budget_limit=5.0)
|
||||
snapshot = _breakdown(1.0, 1)
|
||||
snapshot["_usage_projection"] = {
|
||||
"accounted_usd": 1.0, "integrity_degraded": False, "cost_final": True,
|
||||
}
|
||||
monkeypatch.setattr(accounting, "ensure_legacy_imported", lambda *_a, **_k: None)
|
||||
monkeypatch.setattr(accounting, "usage_breakdown", lambda *_a, **_k: dict(snapshot))
|
||||
monkeypatch.setattr(
|
||||
accounting, "usage_projection",
|
||||
lambda *_a, **_k: pytest.fail("projection must come from the breakdown snapshot"),
|
||||
)
|
||||
|
||||
assert state.update_budget_from_usage({}) is True
|
||||
stored = state.load_state()
|
||||
assert stored["spent_usd"] == 1.0
|
||||
assert stored["usage_accounting"]["accounted_usd"] == 1.0
|
||||
|
||||
|
||||
@pytest.mark.serial
|
||||
def test_compaction_provenance_keeps_high_water_marker(tmp_path, monkeypatch):
|
||||
from supervisor import state
|
||||
import ouroboros.usage_accounting as accounting
|
||||
|
||||
state.init(tmp_path, total_budget_limit=0.0)
|
||||
from ouroboros import usage_compaction as compaction
|
||||
request = accounting.AttemptRequest(
|
||||
model="test/model", provider="test", drive_root=tmp_path,
|
||||
task_id="task", root_task_id="root", reservation_usd=1.0,
|
||||
)
|
||||
hold = accounting.reserve_attempt(request)
|
||||
accounting.mark_dispatched(hold)
|
||||
accounting.settle_attempt(hold, {"prompt_tokens": 1}, cost_usd=1.0, cost_final=True)
|
||||
before_rows = [line for line in (tmp_path / accounting.LEDGER_REL).read_text().splitlines() if line]
|
||||
before = accounting.usage_breakdown(tmp_path)["_ledger_high_water_seq"]
|
||||
state.update_budget_from_usage({})
|
||||
monkeypatch.setattr(compaction, "_fold_clock", lambda: time.time() + 2 * compaction.USAGE_LEDGER_FOLD_MIN_AGE_SEC)
|
||||
with accounting._locked(tmp_path) as heartbeat:
|
||||
assert compaction.compact_usage_ledger_locked(tmp_path, heartbeat=heartbeat)
|
||||
after_rows = [line for line in (tmp_path / accounting.LEDGER_REL).read_text().splitlines() if line]
|
||||
after = accounting.usage_breakdown(tmp_path)["_ledger_high_water_seq"]
|
||||
state.update_budget_from_usage({})
|
||||
|
||||
stored = state.load_state()
|
||||
assert len(after_rows) < len(before_rows)
|
||||
assert max(row["seq"] for row in map(json.loads, after_rows)) < max(
|
||||
row["seq"] for row in map(json.loads, before_rows)
|
||||
)
|
||||
assert tuple(after) > tuple(before)
|
||||
assert stored["spent_usd"] == 1.0
|
||||
assert stored["usage_ledger_high_water_seq"] == after
|
||||
|
||||
|
||||
@pytest.mark.serial
|
||||
def test_reordered_writers_across_real_compaction_reject_lower_epoch(tmp_path, monkeypatch):
|
||||
from supervisor import state
|
||||
import ouroboros.usage_accounting as accounting
|
||||
from ouroboros import usage_compaction as compaction
|
||||
|
||||
state.init(tmp_path, total_budget_limit=0.0)
|
||||
|
||||
def settle(task):
|
||||
request = accounting.AttemptRequest(
|
||||
model="test/model", provider="test", drive_root=tmp_path,
|
||||
task_id=task, root_task_id="root", reservation_usd=1.0,
|
||||
)
|
||||
hold = accounting.reserve_attempt(request)
|
||||
accounting.mark_dispatched(hold)
|
||||
accounting.settle_attempt(hold, {"prompt_tokens": 1}, cost_usd=1.0, cost_final=True)
|
||||
|
||||
settle("before")
|
||||
state.update_budget_from_usage({})
|
||||
started, release = threading.Event(), threading.Event()
|
||||
real_breakdown = accounting.usage_breakdown
|
||||
calls = 0
|
||||
|
||||
def delayed_breakdown(root):
|
||||
nonlocal calls
|
||||
snapshot = real_breakdown(root)
|
||||
calls += 1
|
||||
if calls == 1:
|
||||
started.set()
|
||||
assert release.wait(2.0)
|
||||
return snapshot
|
||||
|
||||
monkeypatch.setattr(accounting, "usage_breakdown", delayed_breakdown)
|
||||
older = threading.Thread(target=state.update_budget_from_usage, args=({},))
|
||||
older.start()
|
||||
assert started.wait(2.0)
|
||||
|
||||
monkeypatch.setattr(
|
||||
compaction, "_fold_clock",
|
||||
lambda: time.time() + 2 * compaction.USAGE_LEDGER_FOLD_MIN_AGE_SEC,
|
||||
)
|
||||
with accounting._locked(tmp_path) as heartbeat:
|
||||
assert compaction.compact_usage_ledger_locked(tmp_path, heartbeat=heartbeat)
|
||||
settle("after")
|
||||
newer = threading.Thread(target=state.update_budget_from_usage, args=({},))
|
||||
newer.start()
|
||||
newer.join(2.0)
|
||||
release.set()
|
||||
older.join(2.0)
|
||||
|
||||
final_marker = real_breakdown(tmp_path)["_ledger_high_water_seq"]
|
||||
stored = state.load_state()
|
||||
assert stored["spent_usd"] == 2.0
|
||||
assert stored["usage_ledger_high_water_seq"] == final_marker
|
||||
|
||||
|
||||
@pytest.mark.serial
|
||||
def test_stale_snapshot_is_rejected_without_state_lock(tmp_path, monkeypatch, caplog):
|
||||
"""No-lock comparison rejects proven staleness, without claiming atomicity."""
|
||||
from supervisor import state
|
||||
import ouroboros.usage_accounting as accounting
|
||||
|
||||
state.init(tmp_path, total_budget_limit=0.0)
|
||||
state.save_state({"spent_usd": 2.0, "usage_ledger_high_water_seq": [0, 2]})
|
||||
monkeypatch.setattr(state, "acquire_file_lock", lambda *_a, **_k: None)
|
||||
monkeypatch.setattr(state, "release_file_lock", lambda *_a, **_k: None)
|
||||
monkeypatch.setattr(accounting, "ensure_legacy_imported", lambda *_a, **_k: None)
|
||||
monkeypatch.setattr(accounting, "usage_breakdown", lambda *_a, **_k: _breakdown(1.0, 1))
|
||||
|
||||
state.update_budget_from_usage({})
|
||||
|
||||
# The marker check preserves the newer projection even after a lock
|
||||
# timeout. It does not serialize two writers that both continue without
|
||||
# STATE_LOCK; that pre-existing compare/save race remains outside this fix.
|
||||
assert state.load_state()["spent_usd"] == 2.0
|
||||
assert state.load_state()["usage_ledger_high_water_seq"] == [0, 2]
|
||||
assert "STALE SNAPSHOT REJECTED" in caplog.text
|
||||
|
||||
|
||||
@pytest.mark.serial
|
||||
def test_equal_marker_writes_newer_projection_even_when_spend_decreases(tmp_path, monkeypatch):
|
||||
from supervisor import state
|
||||
import ouroboros.usage_accounting as accounting
|
||||
|
||||
state.init(tmp_path, total_budget_limit=0.0)
|
||||
state.save_state({"spent_usd": 9.0, "usage_ledger_high_water_seq": [2, 4]})
|
||||
monkeypatch.setattr(accounting, "ensure_legacy_imported", lambda *_a, **_k: None)
|
||||
monkeypatch.setattr(
|
||||
accounting,
|
||||
"usage_breakdown",
|
||||
lambda *_a, **_k: _breakdown(1.0, 4) | {"_ledger_high_water_seq": [2, 4]},
|
||||
)
|
||||
|
||||
state.update_budget_from_usage({})
|
||||
|
||||
assert state.load_state()["spent_usd"] == 1.0
|
||||
assert state.load_state()["usage_ledger_high_water_seq"] == [2, 4]
|
||||
|
||||
|
||||
@pytest.mark.serial
|
||||
def test_lower_epoch_is_rejected_and_preserves_money(tmp_path, monkeypatch, caplog):
|
||||
from supervisor import state
|
||||
import ouroboros.usage_accounting as accounting
|
||||
|
||||
state.init(tmp_path, total_budget_limit=0.0)
|
||||
state.save_state({"spent_usd": 9.0, "usage_ledger_high_water_seq": [3, 4]})
|
||||
monkeypatch.setattr(accounting, "ensure_legacy_imported", lambda *_a, **_k: None)
|
||||
monkeypatch.setattr(
|
||||
accounting,
|
||||
"usage_breakdown",
|
||||
lambda *_a, **_k: _breakdown(1.0, 2) | {"_ledger_high_water_seq": [2, 2]},
|
||||
)
|
||||
|
||||
state.update_budget_from_usage({})
|
||||
|
||||
assert state.load_state()["spent_usd"] == 9.0
|
||||
assert state.load_state()["usage_ledger_high_water_seq"] == [3, 4]
|
||||
assert "STALE SNAPSHOT REJECTED" in caplog.text
|
||||
|
||||
|
||||
@pytest.mark.serial
|
||||
def test_missing_marker_fails_safe_without_fabricating_zero(tmp_path, monkeypatch, caplog):
|
||||
from supervisor import state
|
||||
import ouroboros.usage_accounting as accounting
|
||||
|
||||
state.init(tmp_path, total_budget_limit=0.0)
|
||||
monkeypatch.setattr(accounting, "ensure_legacy_imported", lambda *_a, **_k: None)
|
||||
monkeypatch.setattr(accounting, "usage_breakdown", lambda *_a, **_k: {"accounted_usd": 9.0})
|
||||
|
||||
state.update_budget_from_usage({})
|
||||
|
||||
stored = state.load_state()
|
||||
assert stored["spent_usd"] == 0.0
|
||||
assert "usage_ledger_high_water_seq" not in stored
|
||||
assert "FRESHNESS MARKER UNKNOWN" in caplog.text
|
||||
|
||||
state.save_state({"spent_usd": 7.0, "usage_ledger_high_water_seq": "corrupt"})
|
||||
state.update_budget_from_usage({})
|
||||
assert state.load_state()["spent_usd"] == 7.0
|
||||
|
||||
|
||||
@pytest.mark.serial
|
||||
def test_malformed_current_marker_is_unknown(tmp_path, monkeypatch):
|
||||
from supervisor import state
|
||||
import ouroboros.usage_accounting as accounting
|
||||
|
||||
state.init(tmp_path, total_budget_limit=0.0)
|
||||
state.save_state({"spent_usd": 7.0, "usage_ledger_high_water_seq": [1, 3]})
|
||||
monkeypatch.setattr(accounting, "ensure_legacy_imported", lambda *_a, **_k: None)
|
||||
monkeypatch.setattr(
|
||||
accounting,
|
||||
"usage_breakdown",
|
||||
lambda *_a, **_k: _breakdown(1.0, 4) | {"_ledger_high_water_seq": ["bad", 4]},
|
||||
)
|
||||
|
||||
state.update_budget_from_usage({})
|
||||
|
||||
assert state.load_state()["spent_usd"] == 7.0
|
||||
assert state.load_state()["usage_ledger_high_water_seq"] == [1, 3]
|
||||
|
||||
|
||||
@pytest.mark.serial
|
||||
def test_corrupt_ledger_header_is_unknown_not_zero(tmp_path):
|
||||
from supervisor import state
|
||||
from ouroboros.usage_ledger import LEDGER_REL
|
||||
|
||||
state.init(tmp_path, total_budget_limit=0.0)
|
||||
state.save_state({"spent_usd": 7.0, "usage_ledger_high_water_seq": [1, 3]})
|
||||
ledger_path = tmp_path / LEDGER_REL
|
||||
ledger_path.parent.mkdir(parents=True, exist_ok=True)
|
||||
ledger_path.write_bytes(b"{not-json}\n")
|
||||
|
||||
state.update_budget_from_usage({})
|
||||
|
||||
assert state.load_state()["spent_usd"] == 7.0
|
||||
assert state.load_state()["usage_ledger_high_water_seq"] == [1, 3]
|
||||
|
|
@ -944,6 +944,7 @@ def test_legacy_state_projection_cannot_regress_under_reordered_writers(
|
|||
|
||||
state.init(data_root, total_budget_limit=0.0)
|
||||
first_started = threading.Event()
|
||||
second_started = threading.Event()
|
||||
release_first = threading.Event()
|
||||
calls = []
|
||||
|
||||
|
|
@ -954,6 +955,7 @@ def test_legacy_state_projection_cannot_regress_under_reordered_writers(
|
|||
assert release_first.wait(2.0)
|
||||
value = 1.0
|
||||
else:
|
||||
second_started.set()
|
||||
value = 2.0
|
||||
return {
|
||||
"accounted_usd": value, "physical_calls": int(value),
|
||||
|
|
@ -961,6 +963,7 @@ def test_legacy_state_projection_cannot_regress_under_reordered_writers(
|
|||
"settled_usd": value, "confirmed_usd": value, "estimated_usd": 0.0,
|
||||
"reserved_usd": 0.0, "unresolved_upper_bound_usd": 0.0,
|
||||
"unknown_unmetered": 0, "cost_final": True, "attempt_counts": {},
|
||||
"_ledger_high_water_seq": [0, int(value)],
|
||||
}
|
||||
|
||||
monkeypatch.setattr(ua, "ensure_legacy_imported", lambda *_args, **_kwargs: {})
|
||||
|
|
@ -970,8 +973,11 @@ def test_legacy_state_projection_cannot_regress_under_reordered_writers(
|
|||
older.start()
|
||||
assert first_started.wait(2.0)
|
||||
newer.start()
|
||||
time.sleep(0.1)
|
||||
assert calls == [1]
|
||||
# Ledger snapshots are intentionally taken before STATE_LOCK, so the
|
||||
# newer writer can read while the older one is paused; the sequence marker
|
||||
# below prevents the paused snapshot from regressing state afterward.
|
||||
assert second_started.wait(2.0)
|
||||
assert calls == [1, 2]
|
||||
release_first.set()
|
||||
older.join(2.0)
|
||||
newer.join(2.0)
|
||||
|
|
|
|||
|
|
@ -113,13 +113,21 @@ def _snapshot_looks(monkeypatch, on_look=lambda looks: None):
|
|||
|
||||
|
||||
def _projection_snapshot(data_root):
|
||||
def breakdown(**kwargs):
|
||||
result = ua.usage_breakdown(data_root, **kwargs)
|
||||
# The compatibility writer's freshness marker advances across a real
|
||||
# compaction; monetary/non-money projection equality intentionally
|
||||
# excludes that ordering fact.
|
||||
result.pop("_ledger_high_water_seq", None)
|
||||
return result
|
||||
|
||||
return (
|
||||
ua.usage_projection(data_root),
|
||||
ua.usage_projection(data_root, root_task_id="root"),
|
||||
ua.usage_projection(data_root, root_task_id="root2"),
|
||||
ua.usage_breakdown(data_root),
|
||||
ua.usage_breakdown(data_root, root_task_id="root"),
|
||||
ua.usage_breakdown(data_root, task_id="t2"),
|
||||
breakdown(),
|
||||
breakdown(root_task_id="root"),
|
||||
breakdown(task_id="t2"),
|
||||
)
|
||||
|
||||
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue