From 264df8d5f5521e68f24df2c12c2552cf29150863 Mon Sep 17 00:00:00 2001 From: Jiri Puc Date: Tue, 21 Jul 2026 15:59:39 +0200 Subject: [PATCH] feat: enforce worker cost budgets --- plan/designs/orchestration_design.md | 19 +- plan/plans/phase-7-scale-ops.md | 26 +- plan/plans/roadmap.md | 2 +- src/tests/adapters/test_openrouter.py | 42 ++++ src/tests/adapters/test_selfhost_queue.py | 7 +- src/tests/spine/test_backfill.py | 10 +- src/tests/spine/test_work_ledger.py | 176 ++++++++++++- .../test_port_inventory_and_conformance.py | 27 +- src/tests/test_port_model_values.py | 39 ++- src/tests/workers/test_e0_chain.py | 7 +- src/tests/workers/test_e1_chain.py | 7 +- src/tests/workers/test_e2_chain.py | 26 +- src/tests/workers/test_e3_chain.py | 7 +- .../workers/test_lifecycle_reconciliation.py | 13 +- src/ultimate_memory/adapters/openrouter.py | 47 +++- .../adapters/testing/__init__.py | 8 +- .../adapters/testing/cost_meter.py | 13 + .../adapters/testing/model_provider.py | 41 ++- src/ultimate_memory/eval/consumption.py | 2 +- src/ultimate_memory/model/__init__.py | 14 ++ src/ultimate_memory/model/model_provider.py | 30 +++ src/ultimate_memory/model/processing.py | 55 +++- src/ultimate_memory/ports/cost_meter.py | 17 ++ src/ultimate_memory/ports/model_provider.py | 5 +- .../spine/observation_adjudication.py | 71 +++++- src/ultimate_memory/spine/resolver.py | 70 +++++- src/ultimate_memory/spine/supersession.py | 30 ++- src/ultimate_memory/spine/work_ledger.py | 235 +++++++++++++++++- src/ultimate_memory/surfaces/cli.py | 45 +++- src/ultimate_memory/workers/base.py | 50 +++- src/ultimate_memory/workers/e0.py | 17 +- src/ultimate_memory/workers/e1.py | 21 +- src/ultimate_memory/workers/e2.py | 26 +- src/ultimate_memory/workers/e3.py | 48 +++- .../workers/knowledge_authored.py | 4 +- src/ultimate_memory/workers/p1.py | 17 +- src/ultimate_memory/workers/p2_analytics.py | 2 +- src/ultimate_memory/workers/reconcile.py | 4 +- website/src/app/docs/project-status/page.mdx | 1 + website/src/app/docs/reference/cli/page.mdx | 30 ++- 40 files changed, 1185 insertions(+), 126 deletions(-) create mode 100644 src/tests/adapters/test_openrouter.py create mode 100644 src/ultimate_memory/adapters/testing/cost_meter.py create mode 100644 src/ultimate_memory/ports/cost_meter.py diff --git a/plan/designs/orchestration_design.md b/plan/designs/orchestration_design.md index ba7796a5..4edb4cfd 100644 --- a/plan/designs/orchestration_design.md +++ b/plan/designs/orchestration_design.md @@ -113,10 +113,12 @@ document" — requirements). Rules: - **Declaration**: budgets are per `(deployment, stage, lane, window)` — e.g. "extract_claims / steady / $N per day" — in deployment config. -- **Enforcement is a deterministic pre-flight check in the worker**: read the deduplicated - `cost_ledger` total for the row's authoritative lane and window (cached, refreshed on the - order of minutes — enforcement is a dam, not a scalpel; minutes of lag are priced in). If the - budget is exhausted, the worker +- **Enforcement is a deterministic pre-flight check in the worker**: after locking one due row, + read the deduplicated `cost_ledger` total for its authoritative lane and aligned window in the + same claim transaction. The single-deployment library uses the existing range index directly; + it adds no cache or second budget-state service. Concurrent calls may still finish after another + worker passed pre-flight — enforcement is a dam, not a reservation system. If the budget is + exhausted, the worker **parks** by setting `status='pending'`, `defer_reason='budget'`, and `not_before` to the window roll, then re-announces the row and exits *without* starting the handler (`processing_state.attempts` and `last_error` are untouched). @@ -142,6 +144,15 @@ splitting. `cost_ledger.lane` is copied from `lane IS NULL`, but they do not silently join either plane-E lane. A delivery envelope or Cloud Tasks header cannot choose the attribution. +The model-provider port returns every successful generation or embedding together with required +provider accounting (resolved model name, input/output tokens, USD cost, and latency). A worker +binds a small cost-meter port to the running `processing_id`; each model-using handler writes that +usage under its deterministic call key immediately after the provider returns, before it consumes +the output. Missing or malformed provider accounting is an adapter error rather than an implicit +zero, so configured ceilings cannot be silently bypassed. Deterministic and non-worker consumers +may use the provider port without inventing a processing row; enforced worker-route budgets remain +anchored solely to `processing_state`. + ## 5. Provider-neutral I/O discipline (portable core of design-review F9) The spine or another provider may be remote from a worker. The library therefore avoids chatty diff --git a/plan/plans/phase-7-scale-ops.md b/plan/plans/phase-7-scale-ops.md index 973ef1a3..ff44fc9f 100644 --- a/plan/plans/phase-7-scale-ops.md +++ b/plan/plans/phase-7-scale-ops.md @@ -30,7 +30,7 @@ their round trips. |---|---|---|---|---|---|---| | WP-7.1 | Backfill lanes + seeding + reprocessing orchestration (version bumps) | orchestration §3–4 | Phase 6 | lane machinery | steady-state unaffected during backfill test | done | | WP-7.2 | Reproducible scale battery: D23 partitions/indexes, hub entities/lineages, recount cost, and provider-neutral read/write batching | schema §12; D23; lifecycle §11.5; orchestration §5; retrieval §13.7 | WP-7.1 | fixed synthetic profiles + report | shapes and batching invariants recorded; timings remain measurements, not hosted SLAs | done | -| WP-7.3 | Cost metering + configurable budget enforcement | orchestration §4; schema §2 `cost_ledger` | WP-7.1 | enforcement + admin inspection | explicit fixture ceiling parks and resumes an over-budget lane; attribution is visible | planned | +| WP-7.3 | Cost metering + configurable budget enforcement | orchestration §4; schema §2 `cost_ledger` | WP-7.1 | enforcement + admin inspection | explicit fixture ceiling parks and resumes an over-budget lane; attribution is visible | done | | WP-7.4 | Operational correctness surfaces + drills: typed telemetry, pipeline/DLQ inspection and replay, P2/P3 rebuild, currency-ledger audit | orchestration §6–7; D7, D60–D61 | WP-7.1 | telemetry/admin surfaces + deterministic drills | failures remain visible and drills pass without a dashboard or hosted control plane | planned | | WP-7.5 | **Hard-delete end-to-end**: design first, then purge active P1/P2/P3/K surfaces and prevent restore resurrection through the portable purge record/adapter contract | new design (gate #24); lifecycle §8; k_layers §10; S55 | gate #24 | forget pipeline | **S55 CI gate ON and green** across library-controlled surfaces + restore canary | blocked(#24) | | WP-7.6 | **Release engineering**: semver across PyPI + GHCR images + pinned compose; migrations-before-workers upgrade drill; quickstart cold-start release gate | packaging §1, §5–6; D62 | WP-7.1, rename/CLA gate | release pipeline | tagged release produces all artifacts; upgrade drill green; quickstart under target | blocked(rename-gate) | @@ -73,3 +73,27 @@ The provider-neutral measurement uses the real SQLAlchemy engine with explicit i statements, while currency writes and hub recount retain their constant statement counts. Only shape, correctness, query-count, and transaction-count invariants gate acceptance; every elapsed time is recorded as a machine-specific measurement, never an OSS SLA or topology commitment. + +## WP-7.3 implementation + +Operators declare no implicit monetary policy. An optional typed `UGM_WORK_BUDGETS` list supplies +explicit ceilings keyed by deployment, stage, lane, and aligned fixed window; an omitted route is +unlimited. After locking one due row, the existing claim transaction sums the deduplicated +`cost_ledger` range for that route. Exhaustion moves the row to durable `pending` / `budget` state +until the window boundary without starting a handler, consuming an attempt, changing the last +error, or creating a second scheduling ledger. The worker re-announces that existing row through +the delivery port with the stored resume time. + +`WorkLedger.budget_status` and `ugm budget inspect` read the same two authoritative Postgres +tables. They expose configured ceiling, current-window spend, remaining amount, tier attribution, +aligned bounds, and parked-work count; they do not add a dashboard, hosted billing policy, cache, +or control plane. The PostgreSQL acceptance fixture records two attributed calls, proves an +over-budget handler never starts, then crosses the fixture window and proves the exact row resumes +and completes normally. + +Successful generation and embedding responses carry mandatory provider-reported usage through the +existing model port. The worker binds that usage to its running processing row, and every +model-using stage records a deterministic logical call key and cascade tier before consuming the +response. OpenRouter responses without cost/token accounting fail visibly instead of degrading to +zero spend; the deterministic test provider emits zero-cost usage so end-to-end worker tests prove +the same production attribution path without network calls. diff --git a/plan/plans/roadmap.md b/plan/plans/roadmap.md index 287f9d9f..51181145 100644 --- a/plan/plans/roadmap.md +++ b/plan/plans/roadmap.md @@ -102,7 +102,7 @@ needs to restate it: | 4 | Projections | P2 (spikes → views → rebuild → snapshots), P3 (tree + mounts incl. raw), communities | p2_graph, e0 §6, `p3_agent_navigation.md` | done (exit criteria met 2026-07-19; PRs #100-#104 — see the phase file) | | 5 | Retrieval complete | full primitives + recipe registry, envelope contract CI, MCP/CLI, batch scan, **consumption skill + S58** | retrieval | done (exit criteria met 2026-07-20; PRs #105–#111 — see the phase file) | | 6 | Plane K | planner/writer/driver, fact-sheet → prose bands, citations/staleness, authored + sidecars, triggers + subscriptions, K1 + K2 purpose scopes | k_layers | done (exit criteria met 2026-07-21; PRs #112–#117; former WP-6.7 removed by D73) | -| 7 | Operational correctness + portability | backfill/reprocessing, fixed scale batteries, configurable budgets, failure inspection/drills, hard-delete, release, export/import | orchestration, packaging, schema §12–13 | planned | +| 7 | Operational correctness + portability | backfill/reprocessing, fixed scale batteries, configurable budgets, failure inspection/drills, hard-delete, release, export/import | orchestration, packaging, schema §12–13 | in progress (WP-7.1–7.3 done; see the phase file) | | 8 | Competitive benchmarks | external benchmark harness, adapters, baselines (Mem0/Zep-class), capability benchmark, published methodology + results | D22 (internal) + `phase-8` survey | planned | Sequencing calls already argued (see the phase files for the rest): **K after retrieval** diff --git a/src/tests/adapters/test_openrouter.py b/src/tests/adapters/test_openrouter.py new file mode 100644 index 00000000..073879c6 --- /dev/null +++ b/src/tests/adapters/test_openrouter.py @@ -0,0 +1,42 @@ +"""Provider-accounting proofs for the shipped OpenRouter adapter.""" + +from decimal import Decimal + +import pytest + +from ultimate_memory.adapters.openrouter import _usage +from ultimate_memory.model import ProviderAccountingError + + +def test_usage_keeps_exact_cost_and_defaults_embedding_output_tokens() -> None: + """Parse required accounting without introducing float rounding.""" + usage = _usage( + body={ + "model": "resolved/provider-model", + "usage": {"prompt_tokens": 17, "cost": "0.000123"}, + }, + requested_model="requested/model", + latency_ms=9, + ) + + assert usage.model_name == "resolved/provider-model" + assert usage.tokens_in == 17 + assert usage.tokens_out == 0 + assert usage.cost_usd == Decimal("0.000123") + assert usage.latency_ms == 9 + + +@pytest.mark.parametrize( + "body", + ( + {}, + {"usage": {"prompt_tokens": 1}}, + {"usage": {"prompt_tokens": 1, "cost": "not-a-number"}}, + ), +) +def test_usage_fails_closed_when_required_accounting_is_unusable( + body: dict[str, object], +) -> None: + """Never let a worker interpret absent or malformed provider cost as zero.""" + with pytest.raises(ProviderAccountingError): + _usage(body=body, requested_model="requested/model", latency_ms=1) diff --git a/src/tests/adapters/test_selfhost_queue.py b/src/tests/adapters/test_selfhost_queue.py index 87bb9271..f68a002d 100644 --- a/src/tests/adapters/test_selfhost_queue.py +++ b/src/tests/adapters/test_selfhost_queue.py @@ -27,6 +27,7 @@ from ultimate_memory.model import ProcessingTarget from ultimate_memory.model import QueueRoute from ultimate_memory.model import RunResultOutcome +from ultimate_memory.ports.cost_meter import CostMeterPort from ultimate_memory.ports.queue import TaskQueuePort from ultimate_memory.spine import DeploymentBootstrapper from ultimate_memory.spine import WorkLedger @@ -161,7 +162,7 @@ def test_announce_reannounces_an_existing_row( stage=PipelineStage.EXTRACT_CLAIMS, lane=ProcessingLane.STEADY, ) - assert claimed is not None + assert isinstance(claimed, ClaimedWork) ledger.fail( processing_id=claimed.processing_id, error="Traceback: transient", @@ -178,9 +179,9 @@ def test_announce_reannounces_an_existing_row( class _NoOpHandler: """Succeed without work — the demo chain's terminal stage handler.""" - def handle(self, *, work: ClaimedWork) -> HandlerOutcome: + def handle(self, *, work: ClaimedWork, meter: CostMeterPort) -> HandlerOutcome: """Do nothing and chain nothing.""" - del work + del work, meter return HandlerOutcome() diff --git a/src/tests/spine/test_backfill.py b/src/tests/spine/test_backfill.py index b23f3c6d..ff4f1306 100644 --- a/src/tests/spine/test_backfill.py +++ b/src/tests/spine/test_backfill.py @@ -14,6 +14,7 @@ from ultimate_memory.model import BackfillNotDrainedError from ultimate_memory.model import BackfillSeedRequest +from ultimate_memory.model import ClaimedWork from ultimate_memory.model import DeploymentBootstrapInput from ultimate_memory.model import EnqueueWork from ultimate_memory.model import LaneRouteError @@ -101,7 +102,7 @@ def _complete_prior_work(*, ledger: WorkLedger, target_ids: tuple[UUID, ...]) -> stage=PipelineStage.EXTRACT_CLAIMS, lane=ProcessingLane.STEADY, ) - assert claimed is not None + assert isinstance(claimed, ClaimedWork) ledger.complete(processing_id=claimed.processing_id) @@ -187,13 +188,14 @@ def test_version_bump_is_bounded_resumable_and_cannot_starve_steady_work( stage=PipelineStage.EXTRACT_CLAIMS, lane=ProcessingLane.STEADY, ) - assert claimed_live is not None and claimed_live.target_id == target_ids[3] + assert isinstance(claimed_live, ClaimedWork) + assert claimed_live.target_id == target_ids[3] claimed_backfill = ledger.claim_one( deployment_id=_DEPLOYMENT_ID, stage=PipelineStage.EXTRACT_CLAIMS, lane=ProcessingLane.BACKFILL, ) - assert claimed_backfill is not None + assert isinstance(claimed_backfill, ClaimedWork) assert claimed_backfill.target_id in target_ids[:3] @@ -227,7 +229,7 @@ def test_search_indexes_build_only_after_backfill_has_drained( stage=PipelineStage.EXTRACT_CLAIMS, lane=ProcessingLane.BACKFILL, ) - assert claimed is not None + assert isinstance(claimed, ClaimedWork) ledger.complete(processing_id=claimed.processing_id) finalizer.build_search_indexes(deployment_id=_DEPLOYMENT_ID) diff --git a/src/tests/spine/test_work_ledger.py b/src/tests/spine/test_work_ledger.py index 2a5873fd..fb39c8ba 100644 --- a/src/tests/spine/test_work_ledger.py +++ b/src/tests/spine/test_work_ledger.py @@ -4,6 +4,8 @@ from datetime import datetime from datetime import timedelta from datetime import UTC +from decimal import Decimal +import json from pathlib import Path from uuid import UUID from uuid import uuid4 @@ -16,7 +18,9 @@ from sqlalchemy import text from sqlalchemy.engine import Engine +from ultimate_memory.adapters.testing import RecordingTaskQueue from ultimate_memory.model import ClaimedWork +from ultimate_memory.model import CostBudget from ultimate_memory.model import DeploymentBootstrapInput from ultimate_memory.model import EnqueueWork from ultimate_memory.model import LaneRouteError @@ -27,10 +31,12 @@ from ultimate_memory.model import RunResultOutcome from ultimate_memory.model import UnknownStageHandlerError from ultimate_memory.model import WorkNotRunningError +from ultimate_memory.ports.cost_meter import CostMeterPort from ultimate_memory.spine import DeploymentBootstrapper from ultimate_memory.spine import WorkLedger from ultimate_memory.spine import WorkLedgerSettings from ultimate_memory.spine.settings import load_database_settings +from ultimate_memory.surfaces import cli_main from ultimate_memory.workers import HandlerOutcome from ultimate_memory.workers import HandlerRegistry from ultimate_memory.workers import Worker @@ -134,7 +140,8 @@ def test_enqueue_is_idempotent_and_promotes_backfill_to_steady( stage=PipelineStage.EXTRACT_CLAIMS, lane=ProcessingLane.STEADY, ) - assert claimed is not None and claimed.processing_id == first.processing_id + assert isinstance(claimed, ClaimedWork) + assert claimed.processing_id == first.processing_id def test_lane_pairing_is_enforced_at_enqueue_and_claim(ledger: WorkLedger) -> None: @@ -165,10 +172,12 @@ def test_cost_attribution_is_copied_from_the_running_row( stage=PipelineStage.EXTRACT_CLAIMS, lane=ProcessingLane.STEADY, ) - assert claimed is not None + assert isinstance(claimed, ClaimedWork) call = RecordCall( - processing_id=claimed.processing_id, call_key="selection", cost_usd=0.01 + processing_id=claimed.processing_id, + call_key="selection", + cost_usd=Decimal("0.01"), ) assert ledger.record_call(call=call) is True assert ledger.record_call(call=call) is False # ack-lost retry cannot double-bill @@ -239,7 +248,7 @@ def test_running_work_can_never_be_budget_parked(ledger: WorkLedger) -> None: stage=PipelineStage.EXTRACT_CLAIMS, lane=ProcessingLane.STEADY, ) - assert claimed is not None + assert isinstance(claimed, ClaimedWork) with pytest.raises(WorkNotRunningError): ledger.park_for_budget( processing_id=claimed.processing_id, @@ -247,11 +256,162 @@ def test_running_work_can_never_be_budget_parked(ledger: WorkLedger) -> None: ) +class _CountingHandler: + """A successful handler that exposes whether budget pre-flight let it execute.""" + + def __init__(self) -> None: + """Start with no handler executions.""" + self.calls = 0 + + def handle(self, *, work: ClaimedWork, meter: CostMeterPort) -> HandlerOutcome: + """Count one execution and complete without follow-up work.""" + del work, meter + self.calls += 1 + return HandlerOutcome() + + +def test_configured_budget_parks_reports_and_resumes_without_losing_work( + database_engine: Engine, +) -> None: + """A fixture ceiling parks before an attempt and the next window resumes normally.""" + budget = CostBudget( + deployment_id=_DEPLOYMENT_ID, + stage=PipelineStage.EXTRACT_CLAIMS, + lane=ProcessingLane.STEADY, + window_seconds=86_400, + ceiling_usd=Decimal("1.00"), + ) + budgeted = WorkLedger( + engine=database_engine, + settings=WorkLedgerSettings( + retry_backoff_base_s=0.0, retry_backoff_max_s=0.0, budgets=(budget,) + ), + ) + + billed = budgeted.enqueue(work=_work()) + claimed = budgeted.claim_one( + deployment_id=_DEPLOYMENT_ID, + stage=PipelineStage.EXTRACT_CLAIMS, + lane=ProcessingLane.STEADY, + ) + assert isinstance(claimed, ClaimedWork) + assert claimed.processing_id == billed.processing_id + assert budgeted.record_call( + call=RecordCall( + processing_id=claimed.processing_id, + call_key="selection", + tier="selection", + cost_usd=Decimal("0.75"), + ) + ) + assert budgeted.record_call( + call=RecordCall( + processing_id=claimed.processing_id, + call_key="decontextualize", + tier="frontier", + cost_usd=Decimal("0.50"), + ) + ) + budgeted.complete(processing_id=claimed.processing_id) + + waiting = budgeted.enqueue(work=_work()) + handler = _CountingHandler() + registry = HandlerRegistry() + registry.register(stage=PipelineStage.EXTRACT_CLAIMS, handler=handler) + queue = RecordingTaskQueue() + worker = Worker(ledger=budgeted, registry=registry, queue=queue) + + parked_result = worker.run_one( + deployment_id=_DEPLOYMENT_ID, + stage=PipelineStage.EXTRACT_CLAIMS, + lane=ProcessingLane.STEADY, + ) + assert parked_result.processing_id == waiting.processing_id + assert parked_result.outcome is RunResultOutcome.BUDGET_PARKED + assert handler.calls == 0 + (announcement,) = queue.announcements + assert announcement.processing_id == waiting.processing_id + + with database_engine.connect() as connection: + parked = ( + connection.execute( + text( + "SELECT status, defer_reason, attempts, last_error, not_before" + " FROM processing_state WHERE processing_id = :processing_id" + ), + {"processing_id": waiting.processing_id}, + ) + .mappings() + .one() + ) + assert parked["status"] == "pending" + assert parked["defer_reason"] == "budget" + assert parked["attempts"] == 0 + assert parked["last_error"] is None + assert announcement.not_before_snapshot == parked["not_before"] + + (status,) = budgeted.budget_status(deployment_id=_DEPLOYMENT_ID) + assert status.spent_usd == Decimal("1.250000") + assert status.remaining_usd == Decimal(0) + assert status.exhausted + assert status.parked_work == 1 + assert {tier.tier: tier.cost_usd for tier in status.tiers} == { + "frontier": Decimal("0.500000"), + "selection": Decimal("0.750000"), + } + + with database_engine.begin() as connection: + connection.execute( + text("UPDATE cost_ledger SET occurred_at = occurred_at - interval '2 days'") + ) + connection.execute( + text( + "UPDATE processing_state SET not_before = now()" + " WHERE processing_id = :processing_id" + ), + {"processing_id": waiting.processing_id}, + ) + + resumed = worker.run_one( + deployment_id=_DEPLOYMENT_ID, + stage=PipelineStage.EXTRACT_CLAIMS, + lane=ProcessingLane.STEADY, + ) + assert resumed.processing_id == waiting.processing_id + assert resumed.outcome is RunResultOutcome.SUCCEEDED + assert handler.calls == 1 + + +def test_budget_settings_are_unique_and_cli_inspection_uses_them( + monkeypatch: pytest.MonkeyPatch, capsys: pytest.CaptureFixture[str] +) -> None: + """The environment declares one unambiguous route ceiling visible through the CLI.""" + configured = { + "deployment_id": str(_DEPLOYMENT_ID), + "stage": PipelineStage.EXTRACT_CLAIMS.value, + "lane": ProcessingLane.STEADY.value, + "window_seconds": 3600, + "ceiling_usd": "2.50", + } + monkeypatch.setenv("UGM_WORK_BUDGETS", json.dumps([configured])) + settings = WorkLedgerSettings() + assert settings.budgets[0].ceiling_usd == Decimal("2.50") + with pytest.raises(ValidationError, match="only one cost budget"): + WorkLedgerSettings(budgets=(settings.budgets[0], settings.budgets[0])) + + assert cli_main(["budget", "inspect", "--deployment", str(_DEPLOYMENT_ID)]) == 0 + payload = json.loads(capsys.readouterr().out) + assert payload["stage"] == PipelineStage.EXTRACT_CLAIMS.value + assert payload["lane"] == ProcessingLane.STEADY.value + assert payload["ceiling_usd"] == "2.50" + + class _ChainingNoOpHandler: """The demo no-op handler: succeeds and chains the next stage for its target.""" - def handle(self, *, work: ClaimedWork) -> HandlerOutcome: + def handle(self, *, work: ClaimedWork, meter: CostMeterPort) -> HandlerOutcome: """Produce the chain follow-up without doing any real work.""" + del meter return HandlerOutcome( follow_up=( EnqueueWork( @@ -270,8 +430,9 @@ def handle(self, *, work: ClaimedWork) -> HandlerOutcome: class _AlwaysFailingHandler: """A handler whose every execution raises — exercising retry then dead-letter.""" - def handle(self, *, work: ClaimedWork) -> HandlerOutcome: + def handle(self, *, work: ClaimedWork, meter: CostMeterPort) -> HandlerOutcome: """Fail unconditionally with a real traceback.""" + del meter raise RuntimeError(f"deliberate failure for {work.processing_id}") @@ -352,7 +513,8 @@ def test_promotion_clears_backfill_budget_parking( stage=PipelineStage.EXTRACT_CLAIMS, lane=ProcessingLane.STEADY, ) - assert claimed is not None and claimed.processing_id == enqueued.processing_id + assert isinstance(claimed, ClaimedWork) + assert claimed.processing_id == enqueued.processing_id def test_unregistered_stage_never_strands_a_claimed_row( diff --git a/src/tests/test_port_inventory_and_conformance.py b/src/tests/test_port_inventory_and_conformance.py index 293a81d5..e7a222ed 100644 --- a/src/tests/test_port_inventory_and_conformance.py +++ b/src/tests/test_port_inventory_and_conformance.py @@ -2,6 +2,7 @@ from datetime import datetime from datetime import timezone +from decimal import Decimal from pathlib import Path from types import ModuleType from typing import TypeVar @@ -14,10 +15,12 @@ from ultimate_memory.model import AuthenticatedContext from ultimate_memory.model import EmbeddingRequest from ultimate_memory.model import EmbeddingResponse +from ultimate_memory.model import GeneratedResponse from ultimate_memory.model import KRevision from ultimate_memory.model import ModelRequest from ultimate_memory.model import ObjectKey from ultimate_memory.model import PerimeterCredential +from ultimate_memory.model import ProviderCallUsage from ultimate_memory.model import PublishedMounts from ultimate_memory.model import QueueRoute from ultimate_memory.model import StructuredResponseModel @@ -119,17 +122,32 @@ class FakeModelProvider: def generate( self, *, request: ModelRequest, response_type: type[ResponseT] - ) -> ResponseT: + ) -> GeneratedResponse[ResponseT]: """Construct the caller's declared response type from the rendered prompt.""" - return response_type.model_validate({"answer": request.prompt}) + return GeneratedResponse( + output=response_type.model_validate({"answer": request.prompt}), + usage=_usage(model_name=request.model), + ) def embed(self, *, request: EmbeddingRequest) -> EmbeddingResponse: """Return one two-dimensional vector for every input text.""" return EmbeddingResponse( - vectors=tuple((float(index), 1.0) for index, _ in enumerate(request.texts)) + vectors=tuple((float(index), 1.0) for index, _ in enumerate(request.texts)), + usage=_usage(model_name=request.model), ) +def _usage(*, model_name: str) -> ProviderCallUsage: + """Return deterministic accounting for this structural fake.""" + return ProviderCallUsage( + model_name=model_name, + tokens_in=0, + tokens_out=0, + cost_usd=Decimal(0), + latency_ms=0, + ) + + class FakeTelemetry: """Minimal recording fake preserving original exception identities.""" @@ -256,7 +274,8 @@ def test_model_fake_validates_caller_declared_response_schema() -> None: ) ) - assert generated == ExampleResponse(answer="typed answer") + assert generated.output == ExampleResponse(answer="typed answer") + assert generated.usage.model_name == "configured-model" assert len(embedded.vectors) == 2 diff --git a/src/tests/test_port_model_values.py b/src/tests/test_port_model_values.py index 084e2562..f573646e 100644 --- a/src/tests/test_port_model_values.py +++ b/src/tests/test_port_model_values.py @@ -1,19 +1,56 @@ """Meaningful invariants on shared immutable provider-boundary values.""" +from decimal import Decimal + +from pydantic import BaseModel from pydantic import SecretBytes from pydantic import ValidationError import pytest from ultimate_memory.model import EmbeddingResponse +from ultimate_memory.model import GeneratedResponse from ultimate_memory.model import ObjectKey from ultimate_memory.model import PerimeterCredential +from ultimate_memory.model import ProviderCallUsage from ultimate_memory.model import PublishedMounts +class _Output(BaseModel): + """Small structured output used to prove response/usage pairing.""" + + answer: str + + +def test_generated_response_keeps_exact_decimal_provider_cost() -> None: + """Carry provider accounting beside a validated structured output.""" + response = GeneratedResponse( + output=_Output(answer="ok"), + usage=ProviderCallUsage( + model_name="generation-model", + tokens_in=7, + tokens_out=2, + cost_usd=Decimal("0.000123"), + latency_ms=4, + ), + ) + + assert response.output.answer == "ok" + assert response.usage.cost_usd == Decimal("0.000123") + + def test_embedding_response_rejects_mixed_dimensions() -> None: """Reject malformed provider batches before vectors reach application logic.""" with pytest.raises(ValidationError): - EmbeddingResponse(vectors=((1.0, 2.0), (3.0,))) + EmbeddingResponse( + vectors=((1.0, 2.0), (3.0,)), + usage=ProviderCallUsage( + model_name="embedding-model", + tokens_in=1, + tokens_out=0, + cost_usd=Decimal(0), + latency_ms=0, + ), + ) def test_object_key_is_non_empty_and_frozen() -> None: diff --git a/src/tests/workers/test_e0_chain.py b/src/tests/workers/test_e0_chain.py index eee522ba..9abee84d 100644 --- a/src/tests/workers/test_e0_chain.py +++ b/src/tests/workers/test_e0_chain.py @@ -21,6 +21,7 @@ from ultimate_memory.adapters import MarkitdownConverter from ultimate_memory.adapters.selfhost import LocalFSObjectStore from ultimate_memory.adapters.testing import FakeModelProvider +from ultimate_memory.adapters.testing import NoopCostMeter from ultimate_memory.core import blockize from ultimate_memory.core import ConversionRouter from ultimate_memory.core import MarkdownPassthroughConverter @@ -239,7 +240,7 @@ def test_initial_bulk_ingest_can_enter_the_backfill_lane(rig: _E0Rig) -> None: stage=PipelineStage.CONVERT, lane=ProcessingLane.BACKFILL, ) - assert claimed is not None + assert isinstance(claimed, ClaimedWork) assert claimed.target_id == ingested.version_id @@ -359,8 +360,8 @@ def test_retried_convert_replays_the_stored_representation(rig: _E0Rig) -> None: attempt=1, payload={"version_id": str(ingested.version_id)}, ) - first = handler.handle(work=work) - replay = handler.handle(work=work) # the retried attempt + first = handler.handle(work=work, meter=NoopCostMeter()) + replay = handler.handle(work=work, meter=NoopCostMeter()) # the retried attempt assert replay.follow_up[0].payload == first.follow_up[0].payload count = rig.row( diff --git a/src/tests/workers/test_e1_chain.py b/src/tests/workers/test_e1_chain.py index a428ac88..7840da77 100644 --- a/src/tests/workers/test_e1_chain.py +++ b/src/tests/workers/test_e1_chain.py @@ -20,6 +20,7 @@ from ultimate_memory.adapters.selfhost import LanceChunkIndex from ultimate_memory.adapters.selfhost import LocalFSObjectStore from ultimate_memory.adapters.testing import FakeModelProvider +from ultimate_memory.adapters.testing import NoopCostMeter from ultimate_memory.core import chunker_version from ultimate_memory.core import ChunkerParams from ultimate_memory.core import ConversionRouter @@ -274,7 +275,8 @@ def test_rerunning_the_chunk_stage_replays_the_stored_packing( "version_id": str(ingested.version_id), "representation_id": str(representation), }, - ) + ), + meter=NoopCostMeter(), ) # the replay never re-read artifacts (nonexistent store) and kept the rows: assert replay.follow_up[0].stage is PipelineStage.EMBED_CHUNK @@ -351,6 +353,7 @@ def test_embed_retry_replays_stored_prefixes(rig: _E1Rig, tmp_path: Path) -> Non "version_id": str(ingested.version_id), "representation_id": str(representation), }, - ) + ), + meter=NoopCostMeter(), ) assert len(rig.provider.generated_prompts) == calls_after_first diff --git a/src/tests/workers/test_e2_chain.py b/src/tests/workers/test_e2_chain.py index 00cd1ed1..ba77ed38 100644 --- a/src/tests/workers/test_e2_chain.py +++ b/src/tests/workers/test_e2_chain.py @@ -16,6 +16,7 @@ from ultimate_memory.adapters.selfhost import LanceChunkIndex from ultimate_memory.adapters.selfhost import LocalFSObjectStore from ultimate_memory.adapters.testing import FakeModelProvider +from ultimate_memory.adapters.testing import NoopCostMeter from ultimate_memory.core import chunker_version from ultimate_memory.core import ChunkerParams from ultimate_memory.core import ConversionRouter @@ -275,6 +276,17 @@ def test_claims_land_grounded_with_drops_ledgered_and_stance_kept(rig: _E2Rig) - links = connection.execute( text("SELECT count(*) FROM chunk_claims") ).scalar_one() + metered_calls = ( + connection.execute( + text( + "SELECT call_key, model_name, tier, tokens_in, cost_usd" + " FROM cost_ledger WHERE stage = 'extract_claims'" + " ORDER BY call_key" + ) + ) + .mappings() + .all() + ) # the two grounded claims landed; both fabrications were rejected: assert [claim["claim_text"] for claim in claims] == [ @@ -309,6 +321,13 @@ def test_claims_land_grounded_with_drops_ledgered_and_stance_kept(rig: _E2Rig) - assert by_kind["selection_keep_flagged"]["claim_id"] == flagged_claim assert links == len(claims) + assert [call["call_key"].split(":", 1)[0] for call in metered_calls] == [ + "decontextualize", + "selection", + ] + assert all(call["model_name"] and call["tokens_in"] > 0 for call in metered_calls) + assert {call["tier"] for call in metered_calls} == {"decontextualize", "selection"} + assert all(call["cost_usd"] == 0 for call in metered_calls) def test_rerunning_extraction_replays_without_model_calls(rig: _E2Rig) -> None: @@ -346,7 +365,8 @@ def test_rerunning_extraction_replays_without_model_calls(rig: _E2Rig) -> None: "version_id": str(version), "representation_id": str(representation), }, - ) + ), + meter=NoopCostMeter(), ) assert len(rig.provider.generated_prompts) == calls_after_first with rig.engine.connect() as connection: @@ -410,7 +430,7 @@ def test_empty_extraction_is_terminal_and_replays_without_calls( attempt=1, payload={"version_id": str(version), "representation_id": str(representation)}, ) - handler.handle(work=work) + handler.handle(work=work, meter=NoopCostMeter()) calls_after_first = len(empty_provider.generated_prompts) assert calls_after_first == 1 # one Selection call, no fused call @@ -427,5 +447,5 @@ def test_empty_extraction_is_terminal_and_replays_without_calls( ) assert marker["decision_type"] == "selection_drop" - handler.handle(work=work.model_copy(update={"attempt": 2})) + handler.handle(work=work.model_copy(update={"attempt": 2}), meter=NoopCostMeter()) assert len(empty_provider.generated_prompts) == calls_after_first diff --git a/src/tests/workers/test_e3_chain.py b/src/tests/workers/test_e3_chain.py index 6b076192..5fc2660f 100644 --- a/src/tests/workers/test_e3_chain.py +++ b/src/tests/workers/test_e3_chain.py @@ -16,6 +16,7 @@ from ultimate_memory.adapters.selfhost import LanceChunkIndex from ultimate_memory.adapters.selfhost import LocalFSObjectStore from ultimate_memory.adapters.testing import FakeModelProvider +from ultimate_memory.adapters.testing import NoopCostMeter from ultimate_memory.core import chunker_version from ultimate_memory.core import ChunkerParams from ultimate_memory.core import ConversionRouter @@ -450,7 +451,8 @@ def test_rerunning_normalization_replays_without_model_calls(rig: _E3Rig) -> Non "version_id": str(version), "representation_id": str(representation), }, - ) + ), + meter=NoopCostMeter(), ) assert len(rig.provider.generated_prompts) == calls_after_first with rig.engine.connect() as connection: @@ -623,6 +625,7 @@ def test_p1_channels_carry_claims_and_labeled_facts(rig: _E3Rig) -> None: lane=ProcessingLane.STEADY, attempt=2, payload={"doc_id": str(ingested.doc_id)}, - ) + ), + meter=NoopCostMeter(), ) assert len(rig.provider.generated_prompts) == calls diff --git a/src/tests/workers/test_lifecycle_reconciliation.py b/src/tests/workers/test_lifecycle_reconciliation.py index 8f388a1d..c0f90af6 100644 --- a/src/tests/workers/test_lifecycle_reconciliation.py +++ b/src/tests/workers/test_lifecycle_reconciliation.py @@ -28,6 +28,7 @@ from ultimate_memory.adapters.selfhost import LanceChunkIndex from ultimate_memory.adapters.selfhost import LocalFSObjectStore from ultimate_memory.adapters.testing import FakeModelProvider +from ultimate_memory.adapters.testing import NoopCostMeter from ultimate_memory.core import chunker_version from ultimate_memory.core import ChunkerParams from ultimate_memory.core import ConversionRouter @@ -479,7 +480,8 @@ def test_worked_example_edit_retracts_solely_supported_fact(rig: _LifecycleRig) "version_id": str(version_id), "representation_id": str(representation_id), }, - ) + ), + meter=NoopCostMeter(), ) assert ( rig.scalar("SELECT count(*) FROM testimony_currency_events") == event @@ -533,7 +535,8 @@ def test_extractor_bump_without_rederivation_flags_support_withdrawn( "version_id": str(ingested.version_id), "representation_id": str(representation_id), }, - ) + ), + meter=NoopCostMeter(), ) fact = rig.relation() assert fact["evidence_count"] == 0 @@ -824,7 +827,8 @@ def test_finalization_never_closes_a_flagged_fact(rig: _LifecycleRig) -> None: "version_id": str(ingested.version_id), "representation_id": str(representation_id), }, - ) + ), + meter=NoopCostMeter(), ) assert ( rig.scalar( @@ -915,7 +919,8 @@ def test_restore_support_plants_a_canary_the_pack_rechecks(rig: _LifecycleRig) - "version_id": str(ingested.version_id), "representation_id": str(representation_id), }, - ) + ), + meter=NoopCostMeter(), ) review_id = rig.scalar("SELECT review_id FROM review_queue") rig.review.decide_support_withdrawn( diff --git a/src/ultimate_memory/adapters/openrouter.py b/src/ultimate_memory/adapters/openrouter.py index 7040d4d1..5b515136 100644 --- a/src/ultimate_memory/adapters/openrouter.py +++ b/src/ultimate_memory/adapters/openrouter.py @@ -1,6 +1,9 @@ """The OpenRouter model-provider adapter (D63/D70): the shipped default binding.""" +from decimal import Decimal +from decimal import InvalidOperation import json +import time from typing import Any from typing import TypeVar @@ -11,7 +14,10 @@ from ultimate_memory.model import EmbeddingRequest from ultimate_memory.model import EmbeddingResponse +from ultimate_memory.model import GeneratedResponse from ultimate_memory.model import ModelRequest +from ultimate_memory.model import ProviderAccountingError +from ultimate_memory.model import ProviderCallUsage from ultimate_memory.model import StructuredResponseModel ResponseT = TypeVar("ResponseT", bound=StructuredResponseModel) @@ -45,8 +51,9 @@ def __init__(self, *, settings: OpenRouterSettings) -> None: def generate( self, *, request: ModelRequest, response_type: type[ResponseT] - ) -> ResponseT: + ) -> GeneratedResponse[ResponseT]: """One chat completion constrained to the caller's declared JSON schema.""" + started_ns = time.monotonic_ns() body = self._post( path="/chat/completions", payload={ @@ -62,24 +69,36 @@ def generate( }, }, ) + usage = _usage( + body=body, + requested_model=request.model, + latency_ms=(time.monotonic_ns() - started_ns) // 1_000_000, + ) try: content = body["choices"][0]["message"]["content"] - return response_type.model_validate(json.loads(content)) + output = response_type.model_validate(json.loads(content)) except (KeyError, IndexError, ValueError) as err: raise OpenRouterProviderError( f"unusable completion body for {response_type.__name__}" ) from err + return GeneratedResponse(output=output, usage=usage) def embed(self, *, request: EmbeddingRequest) -> EmbeddingResponse: """One embeddings call for the caller's batch.""" + started_ns = time.monotonic_ns() body = self._post( path="/embeddings", payload={"model": request.model, "input": list(request.texts)}, ) + usage = _usage( + body=body, + requested_model=request.model, + latency_ms=(time.monotonic_ns() - started_ns) // 1_000_000, + ) try: ordered = sorted(body["data"], key=lambda item: item["index"]) return EmbeddingResponse( - vectors=tuple(tuple(item["embedding"]) for item in ordered) + vectors=tuple(tuple(item["embedding"]) for item in ordered), usage=usage ) except (KeyError, TypeError, ValueError) as err: raise OpenRouterProviderError("unusable embeddings body") from err @@ -93,3 +112,25 @@ def _post(self, *, path: str, payload: dict[str, object]) -> dict[str, Any]: f"{response.text[:500]}" ) return response.json() + + +def _usage( + *, body: dict[str, Any], requested_model: str, latency_ms: int +) -> ProviderCallUsage: + """Validate OpenRouter accounting; missing usage must not silently disable budgets.""" + raw = body.get("usage") + if not isinstance(raw, dict): + raise ProviderAccountingError("OpenRouter response carries no usage accounting") + model_name = body.get("model", requested_model) + try: + return ProviderCallUsage( + model_name=model_name, + tokens_in=raw["prompt_tokens"], + tokens_out=raw.get("completion_tokens", 0), + cost_usd=Decimal(str(raw["cost"])), + latency_ms=latency_ms, + ) + except (InvalidOperation, KeyError, TypeError, ValueError) as err: + raise ProviderAccountingError( + "OpenRouter response carries invalid usage accounting" + ) from err diff --git a/src/ultimate_memory/adapters/testing/__init__.py b/src/ultimate_memory/adapters/testing/__init__.py index 17623cdc..81fb63fa 100644 --- a/src/ultimate_memory/adapters/testing/__init__.py +++ b/src/ultimate_memory/adapters/testing/__init__.py @@ -1,7 +1,13 @@ """Test-tier adapters: in-memory doubles outside the two-maintained-adapter set.""" +from ultimate_memory.adapters.testing.cost_meter import NoopCostMeter from ultimate_memory.adapters.testing.model_provider import FakeModelProvider from ultimate_memory.adapters.testing.queue import RecordedAnnouncement from ultimate_memory.adapters.testing.queue import RecordingTaskQueue -__all__ = ("FakeModelProvider", "RecordedAnnouncement", "RecordingTaskQueue") +__all__ = ( + "FakeModelProvider", + "NoopCostMeter", + "RecordedAnnouncement", + "RecordingTaskQueue", +) diff --git a/src/ultimate_memory/adapters/testing/cost_meter.py b/src/ultimate_memory/adapters/testing/cost_meter.py new file mode 100644 index 00000000..220486be --- /dev/null +++ b/src/ultimate_memory/adapters/testing/cost_meter.py @@ -0,0 +1,13 @@ +"""A no-op cost meter for tests that invoke stage handlers directly.""" + +from ultimate_memory.model import ProviderCallUsage + + +class NoopCostMeter: + """Accept provider accounting without persisting it.""" + + def record( + self, *, call_key: str, tier: str | None, usage: ProviderCallUsage + ) -> None: + """Discard one test-only call record.""" + del call_key, tier, usage diff --git a/src/ultimate_memory/adapters/testing/model_provider.py b/src/ultimate_memory/adapters/testing/model_provider.py index 610be5f7..ae933f87 100644 --- a/src/ultimate_memory/adapters/testing/model_provider.py +++ b/src/ultimate_memory/adapters/testing/model_provider.py @@ -1,10 +1,13 @@ """A deterministic in-memory model provider for behavior tests (no network).""" +from decimal import Decimal import hashlib from ultimate_memory.model import EmbeddingRequest from ultimate_memory.model import EmbeddingResponse +from ultimate_memory.model import GeneratedResponse from ultimate_memory.model import ModelRequest +from ultimate_memory.model import ProviderCallUsage from ultimate_memory.model import StructuredResponseModel _EMBEDDING_DIMENSION = 8 @@ -30,28 +33,50 @@ def __init__( def generate[ResponseT: StructuredResponseModel]( self, *, request: ModelRequest, response_type: type[ResponseT] - ) -> ResponseT: + ) -> GeneratedResponse[ResponseT]: """Return the canned payload validated as the caller's declared type.""" self.generated_prompts.append(request.prompt) if callable(self._generate_router): - return response_type.model_validate( + output = response_type.model_validate( self._generate_router(request.prompt, response_type.__name__) ) - payload = self._generate_payloads.get( - response_type.__name__, self._generate_payload + else: + payload = self._generate_payloads.get( + response_type.__name__, self._generate_payload + ) + if payload is None: + raise AssertionError(f"no canned payload for {response_type.__name__}") + output = response_type.model_validate(payload) + return GeneratedResponse( + output=output, + usage=_fake_usage( + model_name=request.model, tokens_in=len(request.prompt.split()) + ), ) - if payload is None: - raise AssertionError(f"no canned payload for {response_type.__name__}") - return response_type.model_validate(payload) def embed(self, *, request: EmbeddingRequest) -> EmbeddingResponse: """Return one deterministic content-derived vector per input text.""" self.embedded_texts.extend(request.texts) return EmbeddingResponse( - vectors=tuple(_vector_for(text=text) for text in request.texts) + vectors=tuple(_vector_for(text=text) for text in request.texts), + usage=_fake_usage( + model_name=request.model, + tokens_in=sum(len(text.split()) for text in request.texts), + ), ) +def _fake_usage(*, model_name: str, tokens_in: int) -> ProviderCallUsage: + """Return deterministic zero-cost accounting for one in-memory provider call.""" + return ProviderCallUsage( + model_name=model_name, + tokens_in=tokens_in, + tokens_out=0, + cost_usd=Decimal(0), + latency_ms=0, + ) + + def _vector_for(*, text: str) -> tuple[float, ...]: """Derive a stable pseudo-vector from the text content.""" digest = hashlib.sha256(text.encode("utf-8")).digest() diff --git a/src/ultimate_memory/eval/consumption.py b/src/ultimate_memory/eval/consumption.py index b8486c03..e6d670af 100644 --- a/src/ultimate_memory/eval/consumption.py +++ b/src/ultimate_memory/eval/consumption.py @@ -82,7 +82,7 @@ def evaluate(case: CanaryCase) -> bool: request=ModelRequest(model=model, prompt=_prompt(skill=skill, task=task)), response_type=S58Answer, ) - return answer == expected + return answer.output == expected return evaluate diff --git a/src/ultimate_memory/model/__init__.py b/src/ultimate_memory/model/__init__.py index a1089b5a..7102945e 100644 --- a/src/ultimate_memory/model/__init__.py +++ b/src/ultimate_memory/model/__init__.py @@ -207,7 +207,10 @@ from ultimate_memory.model.lifecycle import ReconciliationDelta from ultimate_memory.model.model_provider import EmbeddingRequest from ultimate_memory.model.model_provider import EmbeddingResponse +from ultimate_memory.model.model_provider import GeneratedResponse from ultimate_memory.model.model_provider import ModelRequest +from ultimate_memory.model.model_provider import ProviderAccountingError +from ultimate_memory.model.model_provider import ProviderCallUsage from ultimate_memory.model.model_provider import StructuredResponseModel from ultimate_memory.model.mounts import PublishedMounts from ultimate_memory.model.object_store import ObjectAlreadyExistsError @@ -218,7 +221,11 @@ from ultimate_memory.model.processing import BackfillNotDrainedError from ultimate_memory.model.processing import BackfillSeedRequest from ultimate_memory.model.processing import BackfillSeedResult +from ultimate_memory.model.processing import BudgetParked from ultimate_memory.model.processing import ClaimedWork +from ultimate_memory.model.processing import CostBudget +from ultimate_memory.model.processing import CostBudgetStatus +from ultimate_memory.model.processing import CostTierSpend from ultimate_memory.model.processing import DeferReason from ultimate_memory.model.processing import EnqueueOutcome from ultimate_memory.model.processing import EnqueueWork @@ -268,6 +275,7 @@ "BackfillNotDrainedError", "BackfillSeedRequest", "BackfillSeedResult", + "BudgetParked", "AddedContext", "AdjudicationVerdict", "AggregateBucket", @@ -298,6 +306,9 @@ "ComponentVersionRecord", "ConnectorCreate", "ConnectorDescriptor", + "CostBudget", + "CostBudgetStatus", + "CostTierSpend", "ConnectorNotFoundError", "ContextPrefix", "ConversionError", @@ -319,6 +330,7 @@ "DocumentVersionNotFoundError", "EmbeddingRequest", "EmbeddingResponse", + "GeneratedResponse", "EmbeddingUpdate", "EnqueueOutcome", "EnqueueWork", @@ -373,6 +385,8 @@ "ProcessingLane", "ProcessingStatus", "ProcessingTarget", + "ProviderAccountingError", + "ProviderCallUsage", "PublishedMounts", "QueueRoute", "PageRef", diff --git a/src/ultimate_memory/model/model_provider.py b/src/ultimate_memory/model/model_provider.py index 629bf20f..c5021660 100644 --- a/src/ultimate_memory/model/model_provider.py +++ b/src/ultimate_memory/model/model_provider.py @@ -1,8 +1,11 @@ """Typed, provider-neutral model and embedding call values for the LLM boundary.""" +from decimal import Decimal from typing import Annotated +from typing import Generic from typing import Self from typing import TypeAlias +from typing import TypeVar from pydantic import BaseModel from pydantic import ConfigDict @@ -12,6 +15,32 @@ _NonEmptyText = Annotated[str, Field(min_length=1)] _EmbeddingVector = Annotated[tuple[float, ...], Field(min_length=1)] StructuredResponseModel: TypeAlias = BaseModel +ResponseT = TypeVar("ResponseT", bound=StructuredResponseModel) + + +class ProviderAccountingError(Exception): + """A provider response omitted or malformed required usage accounting.""" + + +class ProviderCallUsage(BaseModel): + """Provider-reported accounting for one successful generation or embedding call.""" + + model_config = ConfigDict(frozen=True, extra="forbid") + + model_name: _NonEmptyText + tokens_in: int = Field(ge=0) + tokens_out: int = Field(ge=0) + cost_usd: Decimal = Field(ge=Decimal(0)) + latency_ms: int = Field(ge=0) + + +class GeneratedResponse(BaseModel, Generic[ResponseT]): + """A validated structured output paired with its provider accounting.""" + + model_config = ConfigDict(frozen=True, extra="forbid") + + output: ResponseT + usage: ProviderCallUsage class ModelRequest(BaseModel): @@ -38,6 +67,7 @@ class EmbeddingResponse(BaseModel): model_config = ConfigDict(frozen=True, extra="forbid") vectors: Annotated[tuple[_EmbeddingVector, ...], Field(min_length=1)] + usage: ProviderCallUsage @model_validator(mode="after") def require_one_dimension(self) -> Self: diff --git a/src/ultimate_memory/model/processing.py b/src/ultimate_memory/model/processing.py index 39f517c4..f4cfa224 100644 --- a/src/ultimate_memory/model/processing.py +++ b/src/ultimate_memory/model/processing.py @@ -1,5 +1,6 @@ """Typed records for the D67 work ledger: enqueue, claim, attempt, and cost rows.""" +from decimal import Decimal from enum import StrEnum from uuid import UUID @@ -111,6 +112,57 @@ class ClaimedWork(BaseModel): payload: dict[str, object] | None +class CostBudget(BaseModel): + """One explicit spend ceiling for a deployment, stage, lane, and fixed window.""" + + model_config = ConfigDict(frozen=True, extra="forbid") + + deployment_id: UUID + stage: PipelineStage + lane: ProcessingLane | None + window_seconds: int = Field(gt=0) + ceiling_usd: Decimal = Field(gt=Decimal(0)) + + +class BudgetParked(BaseModel): + """A due work row parked before its next handler attempt because spend is exhausted.""" + + model_config = ConfigDict(frozen=True, extra="forbid") + + processing_id: UUID + resume_at: UTCDateTime + spent_usd: Decimal = Field(ge=Decimal(0)) + ceiling_usd: Decimal = Field(gt=Decimal(0)) + + +class CostTierSpend(BaseModel): + """Current-window spend attributed to one recorded cascade tier.""" + + model_config = ConfigDict(frozen=True, extra="forbid") + + tier: str | None + cost_usd: Decimal = Field(ge=Decimal(0)) + + +class CostBudgetStatus(BaseModel): + """Admin-visible state derived from one configured ceiling and the durable ledgers.""" + + model_config = ConfigDict(frozen=True, extra="forbid") + + deployment_id: UUID + stage: PipelineStage + lane: ProcessingLane | None + window_seconds: int = Field(gt=0) + window_started_at: UTCDateTime + window_ends_at: UTCDateTime + ceiling_usd: Decimal = Field(gt=Decimal(0)) + spent_usd: Decimal = Field(ge=Decimal(0)) + remaining_usd: Decimal = Field(ge=Decimal(0)) + exhausted: bool + parked_work: int = Field(ge=0) + tiers: tuple[CostTierSpend, ...] + + class RecordCall(BaseModel): """One billed model/provider call to attribute to the claimed row's attempt. @@ -126,7 +178,7 @@ class RecordCall(BaseModel): tier: str | None = None tokens_in: int | None = None tokens_out: int | None = None - cost_usd: float | None = None + cost_usd: Decimal | None = None latency_ms: int | None = None @@ -134,6 +186,7 @@ class RunResultOutcome(StrEnum): """How one worker pass ended for the row it claimed (or that none was due).""" NO_WORK = "no_work" + BUDGET_PARKED = "budget_parked" SUCCEEDED = "succeeded" RETRY_SCHEDULED = "retry_scheduled" DEAD_LETTERED = "dead_lettered" diff --git a/src/ultimate_memory/ports/cost_meter.py b/src/ultimate_memory/ports/cost_meter.py new file mode 100644 index 00000000..e9a63b24 --- /dev/null +++ b/src/ultimate_memory/ports/cost_meter.py @@ -0,0 +1,17 @@ +"""Provider-neutral sink for attributing one worker attempt's model calls.""" + +from typing import Protocol +from typing import runtime_checkable + +from ultimate_memory.model import ProviderCallUsage + + +@runtime_checkable +class CostMeterPort(Protocol): + """Record provider usage under a deterministic call key and cascade tier.""" + + def record( + self, *, call_key: str, tier: str | None, usage: ProviderCallUsage + ) -> None: + """Persist one successful provider call for the bound processing attempt.""" + ... diff --git a/src/ultimate_memory/ports/model_provider.py b/src/ultimate_memory/ports/model_provider.py index 43141434..196900d8 100644 --- a/src/ultimate_memory/ports/model_provider.py +++ b/src/ultimate_memory/ports/model_provider.py @@ -6,6 +6,7 @@ from ultimate_memory.model import EmbeddingRequest from ultimate_memory.model import EmbeddingResponse +from ultimate_memory.model import GeneratedResponse from ultimate_memory.model import ModelRequest from ultimate_memory.model import StructuredResponseModel @@ -18,8 +19,8 @@ class ModelProviderPort(Protocol): def generate( self, *, request: ModelRequest, response_type: type[ResponseT] - ) -> ResponseT: - """Return a response validated as the caller's declared structured type.""" + ) -> GeneratedResponse[ResponseT]: + """Return validated output plus the provider-reported usage for this call.""" ... def embed(self, *, request: EmbeddingRequest) -> EmbeddingResponse: diff --git a/src/ultimate_memory/spine/observation_adjudication.py b/src/ultimate_memory/spine/observation_adjudication.py index 41430825..1eaed0ec 100644 --- a/src/ultimate_memory/spine/observation_adjudication.py +++ b/src/ultimate_memory/spine/observation_adjudication.py @@ -32,6 +32,7 @@ from ultimate_memory.model import ObservationAssertion from ultimate_memory.model import ObservationOutcome from ultimate_memory.model import ObservationVerdict +from ultimate_memory.ports.cost_meter import CostMeterPort from ultimate_memory.ports.model_provider import ModelProviderPort OBSERVATION_ADJUDICATOR_VERSION: Final = "obs-adjudicator-2026.07" @@ -94,6 +95,8 @@ def add_observation( statement: str, claim_id: UUID, doc_id: UUID, + meter: CostMeterPort | None = None, + call_key: str = "observation", ) -> UUID: """Compatibility wrapper for a one-assertion entity batch.""" return self.add_observations( @@ -104,6 +107,8 @@ def add_observation( statement=statement, claim_id=claim_id, doc_id=doc_id ), ), + meter=meter, + call_key=call_key, )[0] def add_observations( @@ -112,6 +117,8 @@ def add_observations( deployment_id: UUID, subject_entity_id: UUID, assertions: tuple[ObservationAssertion, ...], + meter: CostMeterPort | None = None, + call_key: str = "observation", ) -> tuple[UUID, ...]: """Adjudicate one document/entity batch against one front-loaded block. @@ -153,8 +160,10 @@ def add_observations( assertion=assertion, asserted_at=asserted_by_claim.get(assertion.claim_id), candidates=candidates, + meter=meter, + call_key=f"{call_key}:{assertion_index}", ) - for assertion in assertions + for assertion_index, assertion in enumerate(assertions) ) def _add_with_block( @@ -166,6 +175,8 @@ def _add_with_block( assertion: ObservationAssertion, asserted_at: object, candidates: list[dict[str, object]], + meter: CostMeterPort | None, + call_key: str, ) -> UUID: """Apply one assertion while keeping the front-loaded block current.""" exact = next( @@ -235,7 +246,12 @@ def _add_with_block( statement=assertion.statement, ) return observation_id - ranked = self._rank(statement=assertion.statement, candidates=open_candidates) + ranked = self._rank( + statement=assertion.statement, + candidates=open_candidates, + meter=meter, + call_key=f"{call_key}:rank", + ) if ranked[0][1] < self._settings.novelty_floor: observation_id = self._insert_new( connection=connection, @@ -267,6 +283,8 @@ def _add_with_block( asserted_at=asserted_at, ranked=ranked[: self._settings.hub_top_k], candidates=candidates, + meter=meter, + call_key=call_key, ) def judge_statements( @@ -289,11 +307,16 @@ def _adjudicate_residue( asserted_at: object, ranked: list[tuple[dict[str, object], float]], candidates: list[dict[str, object]], + meter: CostMeterPort | None, + call_key: str, ) -> UUID: """Ladder the similar candidates; apply the first decisive outcome.""" for candidate, similarity in ranked: verdict, method = self._ladder( - existing=str(candidate["statement"]), new=statement + existing=str(candidate["statement"]), + new=statement, + meter=meter, + call_key=f"{call_key}:verdict:{candidate['observation_id']}", ) features: dict[str, object] = { "similarity": similarity, @@ -469,30 +492,58 @@ def _adjudicate_residue( ) return new_id - def _ladder(self, *, existing: str, new: str) -> tuple[ObservationVerdict, str]: + def _ladder( + self, + *, + existing: str, + new: str, + meter: CostMeterPort | None = None, + call_key: str = "observation:verdict", + ) -> tuple[ObservationVerdict, str]: """Small-model verdict, escalating to frontier below the floor.""" prompt = _VERDICT_PROMPT.format(existing=existing, new=new) - verdict = self._model_provider.generate( + verdict_call = self._model_provider.generate( request=ModelRequest(model=self._settings.small_model, prompt=prompt), response_type=ObservationVerdict, ) + if meter is not None: + meter.record( + call_key=f"{call_key}:small", + tier="small_model", + usage=verdict_call.usage, + ) + verdict = verdict_call.output if verdict.confidence >= self._settings.confidence_floor: return verdict, "small_model" - frontier = self._model_provider.generate( + frontier_call = self._model_provider.generate( request=ModelRequest(model=self._settings.frontier_model, prompt=prompt), response_type=ObservationVerdict, ) - return frontier, "frontier_llm" + if meter is not None: + meter.record( + call_key=f"{call_key}:frontier", + tier="frontier_llm", + usage=frontier_call.usage, + ) + return frontier_call.output, "frontier_llm" def _rank( - self, *, statement: str, candidates: Sequence[dict[str, object]] + self, + *, + statement: str, + candidates: Sequence[dict[str, object]], + meter: CostMeterPort | None = None, + call_key: str = "observation:rank", ) -> list[tuple[dict[str, object], float]]: """Similarity-rank candidates (ordering only — the block is already exhaustive, so a skipped candidate can never cause a wrong cap).""" texts = (statement, *(str(c["statement"]) for c in candidates)) - vectors = self._model_provider.embed( + response = self._model_provider.embed( request=EmbeddingRequest(model=self._settings.embedding_model, texts=texts) - ).vectors + ) + if meter is not None: + meter.record(call_key=call_key, tier="embedding", usage=response.usage) + vectors = response.vectors query = vectors[0] scored = [ (candidate, _cosine(query, vector)) diff --git a/src/ultimate_memory/spine/resolver.py b/src/ultimate_memory/spine/resolver.py index ae9b99b1..04899bac 100644 --- a/src/ultimate_memory/spine/resolver.py +++ b/src/ultimate_memory/spine/resolver.py @@ -28,6 +28,7 @@ from ultimate_memory.model import ResolutionCandidate from ultimate_memory.model import ResolvedEntity from ultimate_memory.model import ResolverConfig +from ultimate_memory.ports.cost_meter import CostMeterPort from ultimate_memory.ports.model_provider import ModelProviderPort from ultimate_memory.ports.p1_index import EntityIndexPort from ultimate_memory.spine.entity_registry import normalized_lemma @@ -81,7 +82,13 @@ def __init__( self._last_rejection: tuple[str, float, dict[str, object]] | None = None def resolve( - self, *, deployment_id: UUID, reference: EntityRef, claim: ClaimForNormalization + self, + *, + deployment_id: UUID, + reference: EntityRef, + claim: ClaimForNormalization, + meter: CostMeterPort | None = None, + call_key: str = "resolve", ) -> ResolvedEntity: """Run the cascade for one reference; mint when nothing matches. @@ -125,6 +132,8 @@ def resolve( reference=reference, claim=claim, candidates=candidates, + meter=meter, + call_key=call_key, ) if decision is not None: candidate, method, confidence, features = decision @@ -148,6 +157,8 @@ def resolve( claim=claim, lemma=lemma, considered=candidates, + meter=meter, + call_key=call_key, ) def _ensure_registered(self, *, deployment_id: UUID) -> None: @@ -216,13 +227,13 @@ def judge_pair( request=ModelRequest(model=self._small_model, prompt=prompt), response_type=AdjudicationVerdict, ) - if verdict.confidence >= thresholds.t4_small_confidence_floor: - return verdict.match, "T4_small" + if verdict.output.confidence >= thresholds.t4_small_confidence_floor: + return verdict.output.match, "T4_small" frontier = self._model_provider.generate( request=ModelRequest(model=self._frontier_model, prompt=prompt), response_type=AdjudicationVerdict, ) - return frontier.match, "T4_frontier" + return frontier.output.match, "T4_frontier" def _blocked_candidates( self, *, connection: Connection, deployment_id: UUID, lemma: str @@ -259,13 +270,19 @@ def _decide( reference: EntityRef, claim: ClaimForNormalization, candidates: tuple[ResolutionCandidate, ...], + meter: CostMeterPort | None, + call_key: str, ) -> tuple[ResolutionCandidate, str, float, dict[str, object]] | None: """T3 embedding bands, then T4 adjudication for the ambiguous band.""" if not candidates: return None thresholds = self._config.thresholds_for(entity_type=reference.type) scored = self._t3_scores( - deployment_id=deployment_id, reference=reference, candidates=candidates + deployment_id=deployment_id, + reference=reference, + candidates=candidates, + meter=meter, + call_key=f"{call_key}:t3", ) ordered = sorted( scored, @@ -293,7 +310,11 @@ def _decide( break adjudicated += 1 verdict, seat, model = self._t4( - reference=reference, claim=claim, candidate=candidate + reference=reference, + claim=claim, + candidate=candidate, + meter=meter, + call_key=f"{call_key}:t4:{candidate.entity_id}", ) if verdict.match: return ( @@ -320,13 +341,17 @@ def _t3_scores( deployment_id: UUID, reference: EntityRef, candidates: tuple[ResolutionCandidate, ...], + meter: CostMeterPort | None, + call_key: str, ) -> tuple[tuple[ResolutionCandidate, float | None], ...]: """Cosine similarity against candidate profiles; None = no profile. A missing/stale profile vector is AMBIGUITY (route to T4), never a confident non-match (Codex review). """ - query_vector = self._embed(surface=reference.name) + query_vector = self._embed( + surface=reference.name, meter=meter, call_key=call_key + ) by_id = self._entity_index.entity_vectors( deployment_id=str(deployment_id), entity_ids=tuple(str(candidate.entity_id) for candidate in candidates), @@ -347,6 +372,8 @@ def _t4( reference: EntityRef, claim: ClaimForNormalization, candidate: ResolutionCandidate, + meter: CostMeterPort | None, + call_key: str, ) -> tuple[AdjudicationVerdict, str, str]: """T4 small-model adjudication, escalating to frontier below the floor.""" prompt = _T4_PROMPT.format( @@ -356,18 +383,29 @@ def _t4( candidate=candidate.canonical_name, candidate_type=candidate.type, ) - verdict = self._model_provider.generate( + verdict_call = self._model_provider.generate( request=ModelRequest(model=self._small_model, prompt=prompt), response_type=AdjudicationVerdict, ) + if meter is not None: + meter.record( + call_key=f"{call_key}:small", tier="T4_small", usage=verdict_call.usage + ) + verdict = verdict_call.output thresholds = self._config.thresholds_for(entity_type=reference.type) if verdict.confidence >= thresholds.t4_small_confidence_floor: return verdict, "T4_small", self._small_model - frontier = self._model_provider.generate( + frontier_call = self._model_provider.generate( request=ModelRequest(model=self._frontier_model, prompt=prompt), response_type=AdjudicationVerdict, ) - return frontier, "T4_frontier", self._frontier_model + if meter is not None: + meter.record( + call_key=f"{call_key}:frontier", + tier="T4_frontier", + usage=frontier_call.usage, + ) + return frontier_call.output, "T4_frontier", self._frontier_model def _mint( self, @@ -378,6 +416,8 @@ def _mint( claim: ClaimForNormalization, lemma: str, considered: tuple[ResolutionCandidate, ...], + meter: CostMeterPort | None, + call_key: str, ) -> ResolvedEntity: """Create the canonical entity + alias and index its T3 profile.""" entity_id = uuid4() @@ -414,7 +454,9 @@ def _mint( deployment_id=deployment_id, type=reference.type, canonical_name=reference.name, - vector=self._embed(surface=reference.name), + vector=self._embed( + surface=reference.name, meter=meter, call_key=f"{call_key}:mint" + ), ), ) ) @@ -493,11 +535,15 @@ def _record( entity_id=entity_id, created=created, entity_type=entity_type ) - def _embed(self, *, surface: str) -> tuple[float, ...]: + def _embed( + self, *, surface: str, meter: CostMeterPort | None, call_key: str + ) -> tuple[float, ...]: """One profile/query embedding through the configured port (D63).""" response = self._model_provider.embed( request=EmbeddingRequest(model=self._embedding_model, texts=(surface,)) ) + if meter is not None: + meter.record(call_key=call_key, tier="T3", usage=response.usage) return response.vectors[0] diff --git a/src/ultimate_memory/spine/supersession.py b/src/ultimate_memory/spine/supersession.py index b01465dd..90295365 100644 --- a/src/ultimate_memory/spine/supersession.py +++ b/src/ultimate_memory/spine/supersession.py @@ -27,6 +27,7 @@ from ultimate_memory.model import ModelRequest from ultimate_memory.model import SupersessionOutcome from ultimate_memory.model import SupersessionVerdict +from ultimate_memory.ports.cost_meter import CostMeterPort from ultimate_memory.ports.model_provider import ModelProviderPort ADJUDICATOR_VERSION: Final = "adjudicator-2026.07" @@ -75,7 +76,12 @@ def __init__( self._settings = settings def adjudicate_new_relation( - self, *, deployment_id: UUID, relation_id: UUID + self, + *, + deployment_id: UUID, + relation_id: UUID, + meter: CostMeterPort | None = None, + call_key: str = "supersession", ) -> None: """Run the cascade for one new relation (idempotent per generation). @@ -181,6 +187,8 @@ def adjudicate_new_relation( new=dict(subject), new_relation_id=relation_id, old=dict(candidate), + meter=meter, + call_key=f"{call_key}:{candidate['relation_id']}", ) def _adjudicate_pair( @@ -191,6 +199,8 @@ def _adjudicate_pair( new: dict[str, object], new_relation_id: UUID, old: dict[str, object], + meter: CostMeterPort | None, + call_key: str, ) -> None: """Climb the ladder for one blocked pair and apply the outcome.""" prompt = _ADJUDICATION_PROMPT.format( @@ -201,19 +211,33 @@ def _adjudicate_pair( new_evidence=new["evidence_text"], new_asserted=new["asserted_at"] or "unknown", ) - verdict = self._model_provider.generate( + verdict_call = self._model_provider.generate( request=ModelRequest(model=self._settings.small_model, prompt=prompt), response_type=SupersessionVerdict, ) + if meter is not None: + meter.record( + call_key=f"{call_key}:small", + tier="small_model", + usage=verdict_call.usage, + ) + verdict = verdict_call.output method = "small_model" model = self._settings.small_model if verdict.confidence < self._settings.confidence_floor: - verdict = self._model_provider.generate( + verdict_call = self._model_provider.generate( request=ModelRequest( model=self._settings.frontier_model, prompt=prompt ), response_type=SupersessionVerdict, ) + if meter is not None: + meter.record( + call_key=f"{call_key}:frontier", + tier="frontier_llm", + usage=verdict_call.usage, + ) + verdict = verdict_call.output method = "frontier_llm" model = self._settings.frontier_model features: dict[str, object] = {"model": model, "rationale": verdict.rationale} diff --git a/src/ultimate_memory/spine/work_ledger.py b/src/ultimate_memory/spine/work_ledger.py index 449d76dd..1203e8bc 100644 --- a/src/ultimate_memory/spine/work_ledger.py +++ b/src/ultimate_memory/spine/work_ledger.py @@ -8,10 +8,15 @@ running row and callers can never supply it. """ +from dataclasses import dataclass from datetime import datetime +from decimal import Decimal +from typing import cast +from typing import Self from uuid import UUID from uuid import uuid4 +from pydantic import model_validator from pydantic_settings import BaseSettings from pydantic_settings import SettingsConfigDict from sqlalchemy import bindparam @@ -21,7 +26,11 @@ from sqlalchemy.engine import Engine from sqlalchemy.engine import RowMapping +from ultimate_memory.model import BudgetParked from ultimate_memory.model import ClaimedWork +from ultimate_memory.model import CostBudget +from ultimate_memory.model import CostBudgetStatus +from ultimate_memory.model import CostTierSpend from ultimate_memory.model import EnqueueOutcome from ultimate_memory.model import EnqueueWork from ultimate_memory.model import LaneRouteError @@ -34,12 +43,42 @@ class WorkLedgerSettings(BaseSettings): - """Retry-backoff configuration for failed handler attempts (D67 starting points).""" + """Retry backoff plus explicit route budgets for one worker deployment.""" model_config = SettingsConfigDict(env_prefix="UGM_WORK_", extra="ignore") retry_backoff_base_s: float = 2.0 retry_backoff_max_s: float = 60.0 + budgets: tuple[CostBudget, ...] = () + + @model_validator(mode="after") + def require_unique_valid_budget_routes(self) -> Self: + """Reject ambiguous ceilings and stage/lane pairs that cannot be queued.""" + routes: set[tuple[UUID, PipelineStage, ProcessingLane | None]] = set() + for budget in self.budgets: + if not lane_is_valid( + stage=budget.stage, + lane=None if budget.lane is None else budget.lane.value, + ): + raise ValueError( + f"stage {budget.stage} does not accept budget lane {budget.lane!r}" + ) + route = (budget.deployment_id, budget.stage, budget.lane) + if route in routes: + raise ValueError( + "only one cost budget may be configured per deployment, stage, and lane" + ) + routes.add(route) + return self + + +@dataclass(frozen=True) +class _BudgetWindowSpend: + """The database-clock window and deduplicated spend used by one pre-flight.""" + + started_at: datetime + ends_at: datetime + spent_usd: Decimal class WorkLedger: @@ -64,12 +103,14 @@ def enqueue(self, *, work: EnqueueWork) -> EnqueueOutcome: def claim_one( self, *, deployment_id: UUID, stage: PipelineStage, lane: ProcessingLane | None - ) -> ClaimedWork | None: - """Claim the next due row on one route, or return None when nothing is due. - - Claiming locks the row (SKIP LOCKED), clears any defer reason, moves it to - running, and increments attempts exactly once immediately before the - handler begins — delivery wake-ups without a claim never consume attempts. + ) -> ClaimedWork | BudgetParked | None: + """Claim, budget-park, or find no due row on one route. + + The due row is locked before the current route-window spend is checked. + Exhaustion durably parks it until the aligned window rolls, without + consuming an attempt or touching its last error. Otherwise claiming + clears any defer reason, moves it to running, and increments attempts + exactly once immediately before the handler begins. """ _require_valid_lane(stage=stage, lane=lane) with self._engine.begin() as connection: @@ -83,6 +124,25 @@ def claim_one( ) if row is None: return None + budget = self._budget_for( + deployment_id=deployment_id, stage=stage, lane=lane + ) + if budget is not None: + spend = _budget_window_spend(connection=connection, budget=budget) + if spend.spent_usd >= budget.ceiling_usd: + connection.execute( + _PARK_BUDGET, + { + "processing_id": row["processing_id"], + "resume_at": spend.ends_at, + }, + ) + return BudgetParked( + processing_id=row["processing_id"], + resume_at=spend.ends_at, + spent_usd=spend.spent_usd, + ceiling_usd=budget.ceiling_usd, + ) started = ( connection.execute( _CLAIM_START, {"processing_id": row["processing_id"]} @@ -92,6 +152,60 @@ def claim_one( ) return _claimed_work(row=started) + def budget_status(self, *, deployment_id: UUID) -> tuple[CostBudgetStatus, ...]: + """Return current spend and parked work for every configured deployment budget.""" + statuses: list[CostBudgetStatus] = [] + with self._engine.connect() as connection: + for budget in self._settings.budgets: + if budget.deployment_id != deployment_id: + continue + spend = _budget_window_spend(connection=connection, budget=budget) + tier_rows = connection.execute( + _BUDGET_TIER_SPEND, + { + "deployment_id": budget.deployment_id, + "stage": budget.stage, + "lane": budget.lane, + "window_started_at": spend.started_at, + "window_ends_at": spend.ends_at, + }, + ).mappings() + tiers = tuple( + CostTierSpend( + tier=cast(str | None, row["tier"]), + cost_usd=_decimal(row["cost_usd"]), + ) + for row in tier_rows + ) + parked_work = int( + connection.execute( + _BUDGET_PARKED_COUNT, + { + "deployment_id": budget.deployment_id, + "stage": budget.stage, + "lane": budget.lane, + }, + ).scalar_one() + ) + remaining = max(Decimal(0), budget.ceiling_usd - spend.spent_usd) + statuses.append( + CostBudgetStatus( + deployment_id=budget.deployment_id, + stage=budget.stage, + lane=budget.lane, + window_seconds=budget.window_seconds, + window_started_at=spend.started_at, + window_ends_at=spend.ends_at, + ceiling_usd=budget.ceiling_usd, + spent_usd=spend.spent_usd, + remaining_usd=remaining, + exhausted=spend.spent_usd >= budget.ceiling_usd, + parked_work=parked_work, + tiers=tiers, + ) + ) + return tuple(statuses) + def complete( self, *, processing_id: UUID, follow_up: tuple[EnqueueWork, ...] = () ) -> tuple[EnqueueOutcome, ...]: @@ -162,10 +276,10 @@ def fail( return None def park_for_budget(self, *, processing_id: UUID, resume_at: datetime) -> None: - """Park healthy pending work until its budget window rolls (D67). + """Park queued work until its budget window rolls (D67). Parking happens at claim-time pre-flight, before an attempt starts: it - applies only to pending rows (a running attempt is never parked — that + applies only to pending/failed rows (a running attempt is never parked — that would allow a second concurrent claim), sets defer_reason budget with a future not_before, consumes no attempt, and touches no error state, so it can never cause dead-lettering. @@ -176,7 +290,7 @@ def park_for_budget(self, *, processing_id: UUID, resume_at: datetime) -> None: ).rowcount if updated == 0: raise WorkNotRunningError( - f"processing row {processing_id} is not pending; only queued " + f"processing row {processing_id} is not queued; only queued " "work can be budget-parked" ) @@ -239,6 +353,50 @@ def record_call(self, *, call: RecordCall) -> bool: ).rowcount return inserted == 1 + def _budget_for( + self, *, deployment_id: UUID, stage: PipelineStage, lane: ProcessingLane | None + ) -> CostBudget | None: + """Return the one validated ceiling for a route, if the operator configured it.""" + return next( + ( + budget + for budget in self._settings.budgets + if budget.deployment_id == deployment_id + and budget.stage == stage + and budget.lane == lane + ), + None, + ) + + +def _budget_window_spend( + *, connection: Connection, budget: CostBudget +) -> _BudgetWindowSpend: + """Read one aligned window and its deduplicated cost using the database clock.""" + row = ( + connection.execute( + _BUDGET_WINDOW_SPEND, + { + "deployment_id": budget.deployment_id, + "stage": budget.stage, + "lane": budget.lane, + "window_seconds": budget.window_seconds, + }, + ) + .mappings() + .one() + ) + return _BudgetWindowSpend( + started_at=cast(datetime, row["window_started_at"]), + ends_at=cast(datetime, row["window_ends_at"]), + spent_usd=_decimal(row["spent_usd"]), + ) + + +def _decimal(value: object) -> Decimal: + """Normalize a PostgreSQL numeric aggregate without introducing float rounding.""" + return value if isinstance(value, Decimal) else Decimal(str(value)) + def _require_valid_lane(*, stage: PipelineStage, lane: ProcessingLane | None) -> None: """Reject a lane value that is illegal for the stage's route (D67 pairing).""" @@ -449,8 +607,61 @@ def _claimed_work(*, row: RowMapping) -> ClaimedWork: _PARK_BUDGET = text( """ UPDATE processing_state - SET defer_reason = 'budget', not_before = :resume_at - WHERE processing_id = :processing_id AND status = 'pending' + SET status = 'pending', defer_reason = 'budget', not_before = :resume_at + WHERE processing_id = :processing_id AND status IN ('pending', 'failed') + """ +) + +_BUDGET_WINDOW_SPEND = text( + """ + WITH bounds AS ( + SELECT + to_timestamp( + floor(extract(epoch FROM now()) / :window_seconds) + * :window_seconds + ) AS window_started_at, + to_timestamp( + (floor(extract(epoch FROM now()) / :window_seconds) + 1) + * :window_seconds + ) AS window_ends_at + ) + SELECT bounds.window_started_at, + bounds.window_ends_at, + COALESCE(sum(cost_ledger.cost_usd), 0) AS spent_usd + FROM bounds + LEFT JOIN cost_ledger + ON cost_ledger.deployment_id = :deployment_id + AND cost_ledger.stage = :stage + AND cost_ledger.lane IS NOT DISTINCT FROM :lane + AND cost_ledger.occurred_at >= bounds.window_started_at + AND cost_ledger.occurred_at < bounds.window_ends_at + GROUP BY bounds.window_started_at, bounds.window_ends_at + """ +) + +_BUDGET_TIER_SPEND = text( + """ + SELECT tier, COALESCE(sum(cost_usd), 0) AS cost_usd + FROM cost_ledger + WHERE deployment_id = :deployment_id + AND stage = :stage + AND lane IS NOT DISTINCT FROM :lane + AND occurred_at >= :window_started_at + AND occurred_at < :window_ends_at + GROUP BY tier + ORDER BY tier NULLS FIRST + """ +) + +_BUDGET_PARKED_COUNT = text( + """ + SELECT count(*) + FROM processing_state + WHERE deployment_id = :deployment_id + AND stage = :stage + AND lane IS NOT DISTINCT FROM :lane + AND status = 'pending' + AND defer_reason = 'budget' """ ) diff --git a/src/ultimate_memory/surfaces/cli.py b/src/ultimate_memory/surfaces/cli.py index 6229183e..8b1ac979 100644 --- a/src/ultimate_memory/surfaces/cli.py +++ b/src/ultimate_memory/surfaces/cli.py @@ -1,7 +1,8 @@ -"""The ``ugm`` CLI: a dependency-light client plus an optional local review UI. +"""The ``ugm`` CLI: a dependency-light client plus optional local admin commands. Query, ingest, connector management, and MCP all talk to the deployment HTTP -API. Only ``ugm review`` imports the server extra and connects to the spine. +API. ``ugm review`` and ``ugm budget`` import the server extra and connect to +the spine. """ from __future__ import annotations @@ -26,6 +27,7 @@ if TYPE_CHECKING: from ultimate_memory.spine.review import ReviewQueue + from ultimate_memory.spine.work_ledger import WorkLedger _MERGE_VERDICTS = ("merge", "not_merge") _TRIAGE_VERDICTS = ("restore_support", "invalidate_fact", "uncertain") @@ -38,6 +40,8 @@ def main(argv: list[str] | None = None) -> int: try: if args.command == "review": return _run_review(args) + if args.command == "budget": + return _run_budget(args) if args.command == "query": return _run_query(args) if args.command == "ingest": @@ -84,6 +88,29 @@ def _run_review(args: argparse.Namespace) -> int: engine.dispose() +def _run_budget(args: argparse.Namespace) -> int: + """Compose the local WorkLedger and print configured budget state.""" + try: + from sqlalchemy import create_engine + + from ultimate_memory.spine.settings import load_database_settings + from ultimate_memory.spine.work_ledger import WorkLedger + from ultimate_memory.spine.work_ledger import WorkLedgerSettings + except ModuleNotFoundError: + print( + "error: budget commands require the 'ultimate-memory[server]' extra", + file=sys.stderr, + ) + return 1 + + engine = create_engine(load_database_settings().sqlalchemy_url()) + try: + ledger = WorkLedger(engine=engine, settings=WorkLedgerSettings()) + return _inspect_budgets(ledger=ledger, deployment_id=args.deployment) + finally: + engine.dispose() + + def _run_query(args: argparse.Namespace) -> int: """Run a query command through the typed remote SDK.""" with MemoryClient.from_settings() as client: @@ -218,6 +245,13 @@ def _list(*, queue: ReviewQueue, deployment_id: UUID) -> int: return 0 +def _inspect_budgets(*, ledger: WorkLedger, deployment_id: UUID) -> int: + """Print one current-window JSON record per configured deployment budget.""" + for status in ledger.budget_status(deployment_id=deployment_id): + print(status.model_dump_json()) + return 0 + + def _decide( *, queue: ReviewQueue, @@ -275,6 +309,13 @@ def _build_parser() -> argparse.ArgumentParser: decide.add_argument("--reviewer", required=True) decide.add_argument("--note", default=None) + budget = commands.add_parser("budget", help="inspect configured spend ceilings") + budget_commands = budget.add_subparsers(dest="budget_command", required=True) + inspect = budget_commands.add_parser( + "inspect", help="current spend, tier attribution, and parked work" + ) + inspect.add_argument("--deployment", type=UUID, required=True) + query = commands.add_parser("query", help="query deployment recipes") query_commands = query.add_subparsers(dest="query_command", required=True) query_commands.add_parser("list", help="list the remote recipe tools") diff --git a/src/ultimate_memory/workers/base.py b/src/ultimate_memory/workers/base.py index 3c7f5f7b..1acd3e5b 100644 --- a/src/ultimate_memory/workers/base.py +++ b/src/ultimate_memory/workers/base.py @@ -15,15 +15,19 @@ from pydantic import BaseModel from pydantic import ConfigDict +from ultimate_memory.model import BudgetParked from ultimate_memory.model import ClaimedWork from ultimate_memory.model import EnqueueWork from ultimate_memory.model import HandlerAlreadyRegisteredError from ultimate_memory.model import NonRetryableHandlerError from ultimate_memory.model import PipelineStage from ultimate_memory.model import ProcessingLane +from ultimate_memory.model import ProviderCallUsage from ultimate_memory.model import QueueRoute +from ultimate_memory.model import RecordCall from ultimate_memory.model import RunResultOutcome from ultimate_memory.model import UnknownStageHandlerError +from ultimate_memory.ports.cost_meter import CostMeterPort from ultimate_memory.ports.queue import TaskQueuePort from ultimate_memory.spine.work_ledger import WorkLedger @@ -42,7 +46,7 @@ class HandlerOutcome(BaseModel): class StageHandler(Protocol): """One stage's transformation over one claimed unit of work.""" - def handle(self, *, work: ClaimedWork) -> HandlerOutcome: + def handle(self, *, work: ClaimedWork, meter: CostMeterPort) -> HandlerOutcome: """Process the claimed work and return its chain follow-ups. Raise `NonRetryableHandlerError` for permanent failures (the work @@ -51,6 +55,32 @@ def handle(self, *, work: ClaimedWork) -> HandlerOutcome: ... +class _LedgerCostMeter: + """Bind provider accounting to one claimed processing attempt.""" + + def __init__(self, *, ledger: WorkLedger, processing_id: UUID) -> None: + """Bind every call record to the authoritative running row.""" + self._ledger = ledger + self._processing_id = processing_id + + def record( + self, *, call_key: str, tier: str | None, usage: ProviderCallUsage + ) -> None: + """Persist one provider-reported call through the spine attribution path.""" + self._ledger.record_call( + call=RecordCall( + processing_id=self._processing_id, + call_key=call_key, + model_name=usage.model_name, + tier=tier, + tokens_in=usage.tokens_in, + tokens_out=usage.tokens_out, + cost_usd=usage.cost_usd, + latency_ms=usage.latency_ms, + ) + ) + + class RunResult(BaseModel): """What one runner pass did: which row it ran and how the attempt ended.""" @@ -119,8 +149,24 @@ def run_one( ) if claimed is None: return RunResult(processing_id=None, outcome=RunResultOutcome.NO_WORK) + if isinstance(claimed, BudgetParked): + if self._queue is not None: + self._queue.announce( + processing_id=claimed.processing_id, + route_snapshot=QueueRoute( + deployment_id=deployment_id, stage=stage, lane=lane + ), + not_before_snapshot=claimed.resume_at, + ) + return RunResult( + processing_id=claimed.processing_id, + outcome=RunResultOutcome.BUDGET_PARKED, + ) + meter = _LedgerCostMeter( + ledger=self._ledger, processing_id=claimed.processing_id + ) try: - outcome = handler.handle(work=claimed) + outcome = handler.handle(work=claimed, meter=meter) except NonRetryableHandlerError: _logger.exception( "non-retryable failure in stage %s for %s", diff --git a/src/ultimate_memory/workers/e0.py b/src/ultimate_memory/workers/e0.py index e26f2617..9e518505 100644 --- a/src/ultimate_memory/workers/e0.py +++ b/src/ultimate_memory/workers/e0.py @@ -47,6 +47,7 @@ from ultimate_memory.model import ObjectKey from ultimate_memory.model import PipelineStage from ultimate_memory.model import ProcessingLane +from ultimate_memory.model import ProviderAccountingError from ultimate_memory.model import RepresentationRecord from ultimate_memory.model import SectionTreeRecord from ultimate_memory.model import SnappedSection @@ -54,6 +55,7 @@ from ultimate_memory.model import StructureSource from ultimate_memory.model import UnroutableMimeError from ultimate_memory.model import UploadRecord +from ultimate_memory.ports.cost_meter import CostMeterPort from ultimate_memory.ports.model_provider import ModelProviderPort from ultimate_memory.ports.object_store import ObjectStorePort from ultimate_memory.spine.document_catalog import DocumentCatalog @@ -201,13 +203,14 @@ def __init__( self._artifact_store = artifact_store self._router = router - def handle(self, *, work: ClaimedWork) -> HandlerOutcome: + def handle(self, *, work: ClaimedWork, meter: CostMeterPort) -> HandlerOutcome: """Convert one document version and record its representation. Replay before regenerate (D65/D7): a representation this toolchain already produced for the version is re-chained as-is — the converter is never re-called on a retried or replayed attempt. """ + del meter source = self._catalog.convert_source( version_id=_payload_uuid(work=work, field="version_id") ) @@ -393,7 +396,7 @@ def __init__( self._model_provider = model_provider self._settings = settings or StructurerSettings() - def handle(self, *, work: ClaimedWork) -> HandlerOutcome: + def handle(self, *, work: ClaimedWork, meter: CostMeterPort) -> HandlerOutcome: """Structure one representation and flip currency.""" source = self._catalog.structure_source( representation_id=_payload_uuid(work=work, field="representation_id") @@ -404,7 +407,7 @@ def handle(self, *, work: ClaimedWork) -> HandlerOutcome: blocks = tuple( Block.model_validate(payload) for payload in blocks_doc["blocks"] ) - response = self._propose(source=source, block_count=len(blocks)) + response = self._propose(source=source, block_count=len(blocks), meter=meter) proposed = response.sections if response is not None else () placement = (response.placement or None) if response is not None else None sections = snap_sections( @@ -453,7 +456,7 @@ def handle(self, *, work: ClaimedWork) -> HandlerOutcome: ) def _propose( - self, *, source: StructureSource, block_count: int + self, *, source: StructureSource, block_count: int, meter: CostMeterPort ) -> StructureResponse | None: """Ask the structurer LLM for a tree; every failure degrades to None.""" if self._model_provider is None: @@ -469,12 +472,16 @@ def _propose( document=markdown[: self._settings.max_prompt_chars], ) try: - return self._model_provider.generate( + generated = self._model_provider.generate( request=ModelRequest(model=self._settings.model, prompt=prompt), response_type=StructureResponse, ) + except ProviderAccountingError: + raise # budget enforcement must never degrade missing usage to zero except Exception: # noqa: BLE001 — a document never fails structuring return None + meter.record(call_key="structure", tier="structure", usage=generated.usage) + return generated.output def _write_sidecar( self, diff --git a/src/ultimate_memory/workers/e1.py b/src/ultimate_memory/workers/e1.py index f79447d3..7fdf438f 100644 --- a/src/ultimate_memory/workers/e1.py +++ b/src/ultimate_memory/workers/e1.py @@ -37,6 +37,7 @@ from ultimate_memory.model import P1ChunkRow from ultimate_memory.model import PackedChunk from ultimate_memory.model import PipelineStage +from ultimate_memory.ports.cost_meter import CostMeterPort from ultimate_memory.ports.model_provider import ModelProviderPort from ultimate_memory.ports.object_store import ObjectStorePort from ultimate_memory.ports.p1_index import ChunkIndexPort @@ -88,12 +89,13 @@ def __init__( self._params = params self._chunker_version = chunker_version(params=params) - def handle(self, *, work: ClaimedWork) -> HandlerOutcome: + def handle(self, *, work: ClaimedWork, meter: CostMeterPort) -> HandlerOutcome: """Pack one representation into chunks and chain the embed stage. Replay before regenerate (D7): rows this chunker generation already packed for the version are kept as-is and the stage just re-chains. """ + del meter source = self._catalog.chunk_source( representation_id=_payload_uuid(work=work, field="representation_id") ) @@ -159,7 +161,7 @@ def __init__( self._settings = settings self._chunker_version = chunker_version(params=params) - def handle(self, *, work: ClaimedWork) -> HandlerOutcome: + def handle(self, *, work: ClaimedWork, meter: CostMeterPort) -> HandlerOutcome: """Prefix, embed, and index every chunk of one document version. The D56/A3 carry-forward runs here: an unchanged chunk (same content @@ -191,7 +193,11 @@ def handle(self, *, work: ClaimedWork) -> HandlerOutcome: carried_vectors = self._carried_vectors(work=work, chunks=chunks, carry=carry) prefixes = tuple( self._resolve_prefix( - source=source, chunk=chunk, document_md=document_md, carry=carry + source=source, + chunk=chunk, + document_md=document_md, + carry=carry, + meter=meter, ) for chunk in chunks ) @@ -211,6 +217,9 @@ def handle(self, *, work: ClaimedWork) -> HandlerOutcome: texts=tuple(texts[index] for index in fresh), ) ) + meter.record( + call_key="embed_chunks", tier="embedding", usage=response.usage + ) fresh_vectors = dict( zip( (chunks[index].chunk_id for index in fresh), @@ -301,6 +310,7 @@ def _resolve_prefix( chunk: ChunkForEmbedding, document_md: str, carry: dict[str, CarryForwardSource], + meter: CostMeterPort, ) -> str: """One chunk's "where this sits" sentence: replayed if already stored. @@ -327,7 +337,10 @@ def _resolve_prefix( request=ModelRequest(model=self._settings.prefix_model, prompt=prompt), response_type=ContextPrefix, ) - return response.prefix + meter.record( + call_key=f"prefix:{chunk.chunk_id}", tier="prefix", usage=response.usage + ) + return response.output.prefix def _chunk_record( diff --git a/src/ultimate_memory/workers/e2.py b/src/ultimate_memory/workers/e2.py index 19a40343..a85931ae 100644 --- a/src/ultimate_memory/workers/e2.py +++ b/src/ultimate_memory/workers/e2.py @@ -33,6 +33,7 @@ from ultimate_memory.model import SelectionCandidate from ultimate_memory.model import SelectionResponse from ultimate_memory.model import SelectionVerdict +from ultimate_memory.ports.cost_meter import CostMeterPort from ultimate_memory.ports.model_provider import ModelProviderPort from ultimate_memory.ports.object_store import ObjectStorePort from ultimate_memory.spine.chunk_catalog import ChunkCatalog @@ -102,7 +103,7 @@ def __init__( self._settings = settings self._chunker_version = chunker_version - def handle(self, *, work: ClaimedWork) -> HandlerOutcome: + def handle(self, *, work: ClaimedWork, meter: CostMeterPort) -> HandlerOutcome: """Extract claims for one document version, chunk by chunk (D12 replay).""" source = self._chunk_catalog.chunk_source( representation_id=_payload_uuid(work=work, field="representation_id") @@ -124,7 +125,11 @@ def handle(self, *, work: ClaimedWork) -> HandlerOutcome: if self._reuse_prior_extraction(source=source, chunk=chunk): continue # D56: the prior version's claims are re-attached self._extract_chunk( - source=source, chunks=chunks, index=index, document_md=document_md + source=source, + chunks=chunks, + index=index, + document_md=document_md, + meter=meter, ) return HandlerOutcome( follow_up=( @@ -185,19 +190,26 @@ def _extract_chunk( chunks: tuple[ChunkForEmbedding, ...], index: int, document_md: str, + meter: CostMeterPort, ) -> None: """Run the two Claimify calls for one chunk and land the results.""" chunk = chunks[index] bundle = _bundle_text( source=source, chunks=chunks, index=index, document_md=document_md ) - selection = self._model_provider.generate( + selection_call = self._model_provider.generate( request=ModelRequest( model=self._settings.extract_model, prompt=_SELECTION_PROMPT.format(bundle=bundle), ), response_type=SelectionResponse, ) + meter.record( + call_key=f"selection:{chunk.chunk_id}", + tier="selection", + usage=selection_call.usage, + ) + selection = selection_call.output decisions = list( _selection_decisions(source=source, chunk=chunk, selection=selection) ) @@ -216,7 +228,7 @@ def _extract_chunk( for candidate in keeps if candidate.verdict is SelectionVerdict.KEEP_FLAGGED } - response = self._model_provider.generate( + response_call = self._model_provider.generate( request=ModelRequest( model=self._settings.extract_model, prompt=_CLAIMIFY_PROMPT.format( @@ -226,6 +238,12 @@ def _extract_chunk( ), response_type=ClaimifyResponse, ) + meter.record( + call_key=f"decontextualize:{chunk.chunk_id}", + tier="decontextualize", + usage=response_call.usage, + ) + response = response_call.output for candidate in response.claims: record = _grounded_claim( candidate=candidate, diff --git a/src/ultimate_memory/workers/e3.py b/src/ultimate_memory/workers/e3.py index b85c9128..4900a97a 100644 --- a/src/ultimate_memory/workers/e3.py +++ b/src/ultimate_memory/workers/e3.py @@ -25,6 +25,7 @@ from ultimate_memory.model import NormalizationResponse from ultimate_memory.model import ObservationAssertion from ultimate_memory.model import PipelineStage +from ultimate_memory.ports.cost_meter import CostMeterPort from ultimate_memory.ports.model_provider import ModelProviderPort from ultimate_memory.spine.chunk_catalog import ChunkCatalog from ultimate_memory.spine.claim_catalog import ClaimCatalog @@ -104,7 +105,7 @@ def __init__( self._settings = settings self._chunker_version = chunker_version - def handle(self, *, work: ClaimedWork) -> HandlerOutcome: + def handle(self, *, work: ClaimedWork, meter: CostMeterPort) -> HandlerOutcome: """Normalize one document version's claims into relations/observations. Newly-created relations chain the supersession adjudicator (D3/D4) @@ -160,12 +161,15 @@ def handle(self, *, work: ClaimedWork) -> HandlerOutcome: prompt_lines=prompt_lines, signatures=signatures, type_parents=type_parents, + meter=meter, ) for entity_id, assertions in observations_by_entity.items(): self._observation_adjudicator.add_observations( deployment_id=deployment_id, subject_entity_id=entity_id, assertions=tuple(assertions), + meter=meter, + call_key=f"observation:{entity_id}", ) return HandlerOutcome( follow_up=( @@ -213,9 +217,10 @@ def _normalize_claim( prompt_lines: str, signatures: dict[str, tuple[tuple[str, str], ...]], type_parents: dict[str, str | None], + meter: CostMeterPort, ) -> None: """One claim through the normalizer call and the deterministic gates.""" - response = self._model_provider.generate( + response_call = self._model_provider.generate( request=ModelRequest( model=self._settings.normalize_model, prompt=_NORMALIZE_PROMPT.format( @@ -227,7 +232,13 @@ def _normalize_claim( ), response_type=NormalizationResponse, ) - for relation in response.relations: + meter.record( + call_key=f"normalize:{claim.claim_id}", + tier="normalize", + usage=response_call.usage, + ) + response = response_call.output + for relation_index, relation in enumerate(response.relations): if _OTHER_PREDICATE.fullmatch(relation.predicate): # the D5 escape funnel: register tier=other, unconstrained # by signatures, ranked by usage for periodic promotion @@ -258,10 +269,20 @@ def _normalize_claim( ) continue subject = self._resolver.resolve( - deployment_id=deployment_id, reference=relation.subject, claim=claim + deployment_id=deployment_id, + reference=relation.subject, + claim=claim, + meter=meter, + call_key=( + f"resolve:{claim.claim_id}:relation:{relation_index}:subject" + ), ) object_ = self._resolver.resolve( - deployment_id=deployment_id, reference=relation.object, claim=claim + deployment_id=deployment_id, + reference=relation.object, + claim=claim, + meter=meter, + call_key=f"resolve:{claim.claim_id}:relation:{relation_index}:object", ) if not _signature_allows( predicate=relation.predicate, @@ -292,9 +313,15 @@ def _normalize_claim( ) if upserted.created: created_relations.append(str(upserted.relation_id)) - for observation in response.observations: + for observation_index, observation in enumerate(response.observations): subject = self._resolver.resolve( - deployment_id=deployment_id, reference=observation.subject, claim=claim + deployment_id=deployment_id, + reference=observation.subject, + claim=claim, + meter=meter, + call_key=( + f"resolve:{claim.claim_id}:observation:{observation_index}:subject" + ), ) # the one write path for observations is the D43 adjudicator: # block on the entity, gate cheaply, ladder the residue, @@ -363,7 +390,7 @@ def __init__(self, *, adjudicator: SupersessionAdjudicator) -> None: """Bind the handler to the composed adjudicator.""" self._adjudicator = adjudicator - def handle(self, *, work: ClaimedWork) -> HandlerOutcome: + def handle(self, *, work: ClaimedWork, meter: CostMeterPort) -> HandlerOutcome: """Adjudicate every relation the normalize stage created (idempotent). Chains the reconcile stage — the truth machinery for this version's @@ -379,7 +406,10 @@ def handle(self, *, work: ClaimedWork) -> HandlerOutcome: ) for raw in relation_ids: self._adjudicator.adjudicate_new_relation( - deployment_id=work.deployment_id, relation_id=UUID(str(raw)) + deployment_id=work.deployment_id, + relation_id=UUID(str(raw)), + meter=meter, + call_key=f"supersession:{raw}", ) version_id = payload.get("version_id") representation_id = payload.get("representation_id") diff --git a/src/ultimate_memory/workers/knowledge_authored.py b/src/ultimate_memory/workers/knowledge_authored.py index 19f04431..8f58cb42 100644 --- a/src/ultimate_memory/workers/knowledge_authored.py +++ b/src/ultimate_memory/workers/knowledge_authored.py @@ -15,6 +15,7 @@ from ultimate_memory.model import NonRetryableHandlerError from ultimate_memory.model import PipelineStage from ultimate_memory.model import ProcessingTarget +from ultimate_memory.ports.cost_meter import CostMeterPort from ultimate_memory.spine.knowledge import KnowledgeCompilationError from ultimate_memory.spine.knowledge import KnowledgeControlPlane from ultimate_memory.spine.knowledge import KnowledgeDispatchUnavailableError @@ -98,8 +99,9 @@ def __init__( self._control_plane = control_plane self._dispatcher = dispatcher - def handle(self, *, work: ClaimedWork) -> HandlerOutcome: + def handle(self, *, work: ClaimedWork, meter: CostMeterPort) -> HandlerOutcome: """Deliver at least once and keep both delivery/mirror failures visible.""" + del meter if ( work.target_kind is not ProcessingTarget.KNOWLEDGE_DISPATCH or work.stage is not PipelineStage.DISPATCH_KNOWLEDGE diff --git a/src/ultimate_memory/workers/p1.py b/src/ultimate_memory/workers/p1.py index dc673109..4bc57f0c 100644 --- a/src/ultimate_memory/workers/p1.py +++ b/src/ultimate_memory/workers/p1.py @@ -21,6 +21,7 @@ from ultimate_memory.model import NonRetryableHandlerError from ultimate_memory.model import P1ClaimRow from ultimate_memory.model import P1FactRow +from ultimate_memory.ports.cost_meter import CostMeterPort from ultimate_memory.ports.model_provider import ModelProviderPort from ultimate_memory.ports.p1_index import ClaimIndexPort from ultimate_memory.ports.p1_index import FactIndexPort @@ -71,7 +72,7 @@ def __init__( self._settings = settings self._chunker_version = chunker_version - def handle(self, *, work: ClaimedWork) -> HandlerOutcome: + def handle(self, *, work: ClaimedWork, meter: CostMeterPort) -> HandlerOutcome: """Embed the version's not-yet-embedded claims as one document batch.""" source = self._chunk_catalog.chunk_source( representation_id=_payload_uuid(work=work, field="representation_id") @@ -92,6 +93,7 @@ def handle(self, *, work: ClaimedWork) -> HandlerOutcome: texts=tuple(claim.claim_text for claim in claims), ) ) + meter.record(call_key="embed_claims", tier="embedding", usage=response.usage) self._claim_index.upsert_claims( rows=tuple( P1ClaimRow( @@ -132,7 +134,7 @@ def __init__( self._fact_index = fact_index self._settings = settings - def handle(self, *, work: ClaimedWork) -> HandlerOutcome: + def handle(self, *, work: ClaimedWork, meter: CostMeterPort) -> HandlerOutcome: """Label and embed the document's facts still lacking this generation. Ordering is the invariant (Codex review): the index write lands @@ -151,7 +153,7 @@ def handle(self, *, work: ClaimedWork) -> HandlerOutcome: doc_id=doc_id, label_version=generation, ): - label = self._model_provider.generate( + label_call = self._model_provider.generate( request=ModelRequest( model=self._settings.label_model, prompt=_FACT_LABEL_PROMPT.format( @@ -161,7 +163,13 @@ def handle(self, *, work: ClaimedWork) -> HandlerOutcome: ), ), response_type=FactLabelResponse, - ).label + ) + meter.record( + call_key=f"label_relation:{relation.relation_id}", + tier="label", + usage=label_call.usage, + ) + label = label_call.output.label rows.append( P1FactRow( fact_id=relation.relation_id, @@ -195,6 +203,7 @@ def handle(self, *, work: ClaimedWork) -> HandlerOutcome: texts=tuple(row.label for row in rows), ) ) + meter.record(call_key="embed_facts", tier="embedding", usage=response.usage) self._fact_index.upsert_facts( rows=tuple( row.model_copy(update={"vector": vector}) diff --git a/src/ultimate_memory/workers/p2_analytics.py b/src/ultimate_memory/workers/p2_analytics.py index 703b77f3..9ecf77a6 100644 --- a/src/ultimate_memory/workers/p2_analytics.py +++ b/src/ultimate_memory/workers/p2_analytics.py @@ -230,7 +230,7 @@ def _labels( return {} return { ordered[item.index]: item.label - for item in response.labels + for item in response.output.labels if 0 <= item.index < len(ordered) and item.label } diff --git a/src/ultimate_memory/workers/reconcile.py b/src/ultimate_memory/workers/reconcile.py index 377e04c4..0fb5b383 100644 --- a/src/ultimate_memory/workers/reconcile.py +++ b/src/ultimate_memory/workers/reconcile.py @@ -28,6 +28,7 @@ from ultimate_memory.model import CurrencyTransition from ultimate_memory.model import NonRetryableHandlerError from ultimate_memory.model import ReconciliationDelta +from ultimate_memory.ports.cost_meter import CostMeterPort from ultimate_memory.spine.lifecycle import LifecycleCatalog from ultimate_memory.spine.review import ReviewQueue from ultimate_memory.workers.base import HandlerOutcome @@ -61,13 +62,14 @@ def __init__( params=ChunkerParams() ) - def handle(self, *, work: ClaimedWork) -> HandlerOutcome: + def handle(self, *, work: ClaimedWork, meter: CostMeterPort) -> HandlerOutcome: """Diff → transition → recount → policy → emit, idempotently. The work row's processing_id is the run's `reconciliation_id`: a retried attempt re-emits every ledger row, closure, flag, and trigger as a no-op. """ + del meter version_id = _payload_uuid(work=work, field="version_id") representation_id = _payload_uuid(work=work, field="representation_id") context = self._catalog.reconciliation_context(version_id=version_id) diff --git a/website/src/app/docs/project-status/page.mdx b/website/src/app/docs/project-status/page.mdx index 3141e30a..5afa3aec 100644 --- a/website/src/app/docs/project-status/page.mdx +++ b/website/src/app/docs/project-status/page.mdx @@ -22,6 +22,7 @@ Ultimate Memory is being built **design-first**: the complete system — require - **Phase 5, complete** — retrieval complete. The full zero-LLM primitive set is implemented (`fuse`, `rerank`, `transcript`, `delta`, `pages_about`, enumerated `aggregate`, and a streaming batch `scan`, alongside resolve/lookup/search/graph); recipes are registry rows whose declared grain is enforced mechanically; and the self-accounting envelope surfaces contradictions, withdrawn support, mixed-grain parts, identity regime, belief horizons, freshness, and typed negatives. API, CLI, and MCP all render the same deployment registry. The consumption skill is guarded by the repeatable S58 cold-agent eval, while the retrieval spike battery records filtered Lance search at 10 million rows and graph pagination at a 100,000-edge hub. The client-first base wheel now contains the typed SDK, remote CLI and MCP, lineage-aware E0 ingest, and typed remote connector-management commands plus their deployment composition port, without server dependencies; `[server]`, `[connectors-watched-directory]`, and `[k]` name the heavier install surfaces. The public `remember-dev`/`remember` rename remains deliberately gated with release engineering. - **Phase 6, complete** — Plane K now includes the deterministic control plane and crash-safe single-committer driver, exact rule routing and staleness, deterministic fact-sheet and agent-written prose bands, planner/reflection decisions with quarantine and adoption, and authored-page frontmatter, citations, watches, review flags, and debounced workflow dispatch. D73 removed the proposed K3 tier: personal or organizational principles are authored K2 content, supported by compiled scope pages and never machine-promoted or rewritten. - **Phase 3, complete** — the evidence lifecycle: documents that change. Watched sources poll as recorded sync cycles (a local-directory watcher ships; revision and content no-ops, debounce, source-deletion detection); an edited document becomes a new version of its lineage and the full structure route (an LLM-proposed section tree normalized by a deterministic snap) re-reads it; unchanged chunks reuse their prior claims, prefixes, and vectors so cost is proportional to the edit (measured hit rate 0.79 on the spike corpus); reconciliation transitions testimony currency on an append-only ledger, recounts evidence by distinct current lineages, closes solely-supported facts when the source withdrew them (at the sync-cycle barrier, so a moved section is a support swap — never a retract flicker) and flags them for review when only the toolchain changed; deletion removes a document's contribution uniformly while keeping its claims as history; and a standing `lifecycle` eval suite guards the cache/ledger/count invariants with planted regression canaries. Those lifecycle events now feed the graph, corpus filesystem, and Plane-K control plane described above. +- **Phase 7, in progress** — operational correctness now includes resumable initial-load and version-bump backfill on the same work ledger, a reproducible PostgreSQL scale battery for the designed partitions, indexes, hubs, and batching invariants, and optional per-deployment/stage/lane cost ceilings. Budget exhaustion parks healthy work durably without consuming an attempt; the local CLI reports current spend, remaining budget, tier attribution, and parked work from the same authoritative rows. Timings remain measurements rather than hosted SLAs, and no dashboard, billing policy, HA topology, or control plane enters the library. ## The build order diff --git a/website/src/app/docs/reference/cli/page.mdx b/website/src/app/docs/reference/cli/page.mdx index aa6e99a2..5a0990cf 100644 --- a/website/src/app/docs/reference/cli/page.mdx +++ b/website/src/app/docs/reference/cli/page.mdx @@ -1,7 +1,7 @@ export const metadata = { title: "CLI Reference", description: - "The client-first ugm command line: query, ingest, manage deployment connectors, serve MCP, and drive the optional local review queue.", + "The client-first ugm command line: query, ingest, manage connectors, inspect budgets, serve MCP, and drive the local review queue.", }; # CLI Reference @@ -83,6 +83,34 @@ Runs the base wheel's MCP stdio server. It discovers recipes from the remote deployment and proxies `tools/call` to the same HTTP recipe endpoint used by `ugm query`. +## `ugm budget inspect` + +Budget inspection requires the `[server]` extra and connects directly to the +deployment's PostgreSQL spine through `UGM_DATABASE_URL`: + +```bash +ugm budget inspect --deployment +``` + +Each output line is one configured route ceiling with its aligned window, +current spend, remaining amount, exhaustion state, parked-work count, and spend +grouped by the recorded cascade tier. Both enforcement and inspection read the +same `processing_state` and `cost_ledger` rows; the CLI is not a billing system +or a second source of pipeline truth. + +Ceilings are optional typed configuration in `UGM_WORK_BUDGETS`. The value is a +JSON list; each item explicitly names the deployment, stage, lane, window size, +and USD ceiling: + +```bash +export UGM_WORK_BUDGETS='[{"deployment_id":"","stage":"extract_claims","lane":"steady","window_seconds":86400,"ceiling_usd":"10.00"}]' +``` + +Use `null` for an unlaned K/P route. A route absent from the list is unlimited. +When spend reaches its ceiling, due work is parked until the aligned window +ends without consuming a handler attempt or becoming a failure; the same row +then resumes normally. + ## `ugm review` The review commands require the `[server]` extra. They connect directly to the