diff --git a/deploy/systemd/claim-loop.env.example b/deploy/systemd/claim-loop.env.example index 9722227..ee08d72 100644 --- a/deploy/systemd/claim-loop.env.example +++ b/deploy/systemd/claim-loop.env.example @@ -21,3 +21,11 @@ KUBECONFIG=/etc/rancher/k3s/k3s.yaml BAO_ADDR=https://bao.coulomb.social VAULT_ADDR=https://bao.coulomb.social EXECUTOR_APPROLE_DIR=/home/tegwick/.local/rein-aharness/approle-binky-mail + +# Monetary admission is opt-in for existing workers. A factory deployment must +# set REQUIRE=1 after owner acceptance; absent/invalid state then denies claim. +# Never mount the private policy/ledger into a workload. See docs/spend-admission.md. +# AGENT_HARNESS_REQUIRE_SPEND_ADMISSION=1 +# AGENT_HARNESS_SPEND_POLICY=/absolute/private/spend-policy.json +# AGENT_HARNESS_SPEND_LEDGER=/absolute/private/spend.sqlite3 +# AGENT_HARNESS_EXECUTION_PROJECT= diff --git a/docs/evidence/2026-09-09-spend-admission.json b/docs/evidence/2026-09-09-spend-admission.json new file mode 100644 index 0000000..35c2d4d --- /dev/null +++ b/docs/evidence/2026-09-09-spend-admission.json @@ -0,0 +1,100 @@ +{ + "schema": "rein-spend-admission-evidence/v1", + "observed_at": "2026-09-09T14:37:52.167561+00:00", + "workplan": "REINAH-WP-0003", + "task": "REINAH-WP-0003-T05", + "source_base": "4ffb8acb190df235fc7d1ba20e2643e33f596286", + "source_files_sha256": { + "rein_aharness/spend_admission.py": "334a716da6a0900bca4fdb1fd0833c2446efe08ffb59ae9ee17ce85558c67b81", + "rein_aharness/claim_loop.py": "9fc1796d9dc822f2dd93f7c3edf4fc163bc2b0e502d4c95664db741bc9408180", + "rein_aharness/glas_execution.py": "4dc90e919de113f2f7469a747062a9ba58d50be33611c4ba2bad08da8fb894d8", + "rein_aharness/ops_run_client.py": "ac9263fcac5fa7a007819987b2af11338409811a3dba42a263e81c7efc7ac282", + "rein_aharness/cli.py": "032cc133e1a7edfe4f5353cee6fe1aa55499ff526f47fc8c231e711fc39da597", + "tests/test_spend_admission.py": "6b11b3ca782789f4f91d238b6646cde755117bed50713d63e9135a3be787e3e4", + "tests/test_repository_artifact_bwrap.py": "39cc3e19f46c07e4178387cdcf70c5b32c4a75e1923da75293d8028747d48453", + "scripts/verify-runtime-contracts.sh": "5617882d62e749a8ef698f11091fd2513f87553e739f6033b1c76f43e16bb492", + "scripts/verify-recovery-contracts.sh": "b1bf99e23f33a99b03b383b70c1312bb8b49f10b583f528b4b43a3347431e794", + "docs/spend-admission.md": "d4814f79b1cb05446fd289a1f6a1237e66321d920c86204abf3ebc908d9fc36e", + "docs/sandbox-artifact-return.md": "788c0d7b04f749a9aaa3f59db04887d82f5c40082b2c4e3d94ce8a1a7a8b0673", + "deploy/systemd/claim-loop.env.example": "6cd1e8d62c59735f6f53a126af4326dc1ec15cb3409dfbacb2fe4cd76dbb4b26", + "workplans/REINAH-WP-0003-governed-runtime-integrity.md": "6c8668866a7336ee4b698f15c7bf926d8c30eef51d689d496943756bb37564c8" + }, + "validation": { + "full_suite": { + "passed": 330, + "skipped": 0, + "real_bwrap_enabled": true + }, + "spend_and_actual_bwrap": { + "passed": 41 + }, + "runtime_contract_gate": { + "passed": 90 + }, + "recovery_contract_gate": { + "passed": 52 + }, + "gate_warning": "pytest cache could not write in restricted owner checkout; test assertions passed", + "overlap_note": "Targeted suites and gates overlap the full suite; counts are not additional independent executions." + }, + "proof_scope": { + "actual_bwrap_owner_glas_worker_artifact_import": true, + "queue": "fixture", + "authoring": "deterministic fixture", + "accounting": "fixture", + "model_requests": 0, + "natural_factory_claims": 0, + "deployment_changed": false, + "factory_policy_provisioned": false, + "dispatch_enabled": false + }, + "reuse_discovery": { + "url": "https://reuse.coulomb.social/v1/federated", + "composed_at": "2026-09-09T12:18:07+00:00", + "stale": false, + "capabilities": 65, + "sources": 61, + "snapshot_sha256": "340c5b89a449bb334e6a4922710046011aaa9e302a78297bd2d61f9b8db345b9", + "result": "No reservation capability found in hosted index; inspected llm-connect reporting/token control, Railiance Fabric models and agentic-resources session memory do not provide admission." + }, + "implementation": "Private immutable-policy SQLite ledger before dispatch, conservative daily/total micro-euro reservations, full-charge success, durable unknown holds, replay denial, explicit receipt-backed reconciliation and sticky overrun freeze; existing Glas cached catalog, grants, repository transaction and close outbox reused.", + "remaining": [ + "HFACT-WP-0001-T01: demonstrated provider maximum liability and conservative FX, finalized operating acceptance", + "REINAH-WP-0003-T05/T06 and HFACT-WP-0001-T03/T04/T05: protected installation, exact identity/credential/egress/placement and natural queue/model proof" + ], + "owner_snapshot": { + "observed_at": "2026-09-09T13:13:20.768438+00:00", + "tasks": { + "secrets-replay": { + "id": "3eb9cff8-1441-5437-9e92-a2b655c82d04", + "status": "wait", + "needs_human": false, + "updated_at": "2026-09-09T07:31:58.416393Z" + }, + "native-delivery": { + "id": "f8069c8a-ad6b-5d0b-9a36-c2326699437d", + "status": "wait", + "needs_human": false, + "updated_at": "2026-09-09T07:31:58.646322Z" + }, + "client-side": { + "id": "68bff751-e48b-548b-8fb4-dfb3b16210c2", + "status": "wait", + "needs_human": false, + "updated_at": "2026-09-09T12:41:29.448507Z" + }, + "audit-custody": { + "id": "fd4a4ac3-e525-57c1-9179-e0fcd0226913", + "status": "progress", + "needs_human": false, + "updated_at": "2026-09-08T18:38:29.290026Z" + }, + "approval-deployment": { + "id": "f0aa2e6d-19e6-5b43-886c-efa4e3de5f22", + "status": "wait", + "needs_human": false, + "updated_at": "2026-09-09T12:38:47.065340Z" + } + } + } +} diff --git a/docs/sandbox-artifact-return.md b/docs/sandbox-artifact-return.md index 0ca583e..1b5414a 100644 --- a/docs/sandbox-artifact-return.md +++ b/docs/sandbox-artifact-return.md @@ -52,3 +52,7 @@ including response-lost close replay. Its queue and task authoring are fixtures; it proves neither a natural Activity Core claim nor model/provider admission. `tests/test_native_limits.py` exercises control propagation, terminal accounting, invalid/exhausted results and refusal of older CLI versions without inference. + +The subsequent [durable spend admission](spend-admission.md) return implements +private daily/total reservation and unknown-outcome recovery in the worker. +Provider liability/FX proof and final operating admission remain open. diff --git a/docs/spend-admission.md b/docs/spend-admission.md new file mode 100644 index 0000000..ca676cd --- /dev/null +++ b/docs/spend-admission.md @@ -0,0 +1,133 @@ +# Durable worker spend admission + +The claim worker can reserve a declared maximum liability before invoking Glas. +This is an opt-in control for a single host and one immutable operating envelope. +It adds no service, provider calls or credentials. llm-connect remains the usage +reporting owner; its reporting ledger and bundled FX snapshot are not admission +inputs. Reuse discovery on 2026-09-09 found no reservation capability in the fresh +hosted reuse-surface index (65 capabilities, 61 sources), or in the inspected +llm-connect, Railiance Fabric and agentic-resources implementations. This is a +bounded discovery result, not proof that every fleet capability is registered. + +No factory policy is provisioned by this change. HFACT-WP-0001-T01 still owns +provider enforcement/overrun proof, accepted conservative FX, final operating +admission and the deployment return under REINAH-WP-0003-T05/T06. + +## Admission and accounting + +The worker replays pending terminal closes first, then checks spend capacity +before claiming. A configured worker refuses profile-absent work. Immediately +before invoking Glas it checks the queue owner, ActivityDefinition, resolved +repository, project, actor, repository grant, exact profile and descriptor digests, +and the native USD/turn limits. The same cached catalog supplies the checked +profile to Glas. Policy, ledger and project come from trusted worker configuration; +queue execution references and prompt text cannot supply them. + +An SQLite `BEGIN IMMEDIATE` transaction inserts the reservation before dispatch. +The ledger opens in existing-file mode, with synchronous FULL transactions. A +missing/corrupt ledger, changed policy, invalid FX/amount, expired policy, backward +clock, exhausted daily/total capacity or unresolved reservation refuses execution. +Concurrent processes sharing the file cannot admit two active runs. A run ID or +ActivityDefinition/idempotency-key pair cannot be spent again, even on a new +claim attempt, after reconciliation or under a replacement queue row ID. + +The reservation is `ceil(max_liability_usd * eur_per_usd * 1,000,000)` integer +micro-euros; envelope ceilings round down. Monetary arithmetic uses decimal +values without reporting-FX defaults. `max_budget_usd` is the native CLI stop +threshold and must fit inside the declared maximum liability. These two values +are separate because admission must account for the provider's demonstrated +maximum exposure, including any bounded in-flight overshoot. + +On complete success with finite accounting and confirmed session cleanup and +sandbox destruction, the ledger charges the **full reserved amount**. It never +refunds capacity based on the sandbox's self-reported cost. Missing accounting, +execution failure, cancellation, lost lease or crash leaves the reservation held +and blocks further admission, including after midnight or restart. An observed +cost above declared liability is recorded conservatively and permanently marks +the envelope breached. Reconciliation does not clear that breach. + +A completed or reconciled reservation counts in total once and, conservatively, +against every local calendar day from admission through completion/reconciliation. +A run crossing midnight therefore uses its full reservation in both daily totals. +Unresolved work blocks new admission globally; no timer silently releases it. +This intentionally sacrifices some utilization for the first bounded pilot. +More exact billing reconciliation can follow measured use without weakening the +unknown-outcome boundary. + +## Operator configuration and recovery + +Use a private directory on durable local storage outside every sandbox mount and +target checkout. The directory must be worker-owned mode 0700; policy and ledger +must be worker-owned regular files mode 0600, without hardlinks or symlinks. +All admitted processes must use this one ledger. Separate copies, network-file +locking, multi-host admission, rollback to stale backups and hostile host owners +are outside v1. The protected deployment must keep both files and the launch +configuration inaccessible to the workload and preserve them across restarts. + +Configure the following only after the operating envelope is accepted: + +```text +AGENT_HARNESS_REQUIRE_SPEND_ADMISSION=1 +AGENT_HARNESS_SPEND_POLICY=/absolute/private/spend-policy.json +AGENT_HARNESS_SPEND_LEDGER=/absolute/private/spend.sqlite3 +AGENT_HARNESS_EXECUTION_PROJECT= +``` + +Any nonempty REQUIRE setting requires admission; missing policy/ledger then +refuses before claim. With all spend settings absent, the existing worker +compatibility behavior remains. The factory deployment must set REQUIRE and +verify missing-configuration denial as part of its owner acceptance. + +The strict JSON `SpendPolicy` fields are: + +| Fields | Contract | +| --- | --- | +| `version`, `envelope_id`, `authority_ref` | Version `1`, unique envelope identity, reference to the accepted owner record. Local configuration is trusted; this module does not verify signatures or grant authority. | +| `valid_from`, `expires_at`, `timezone` | Explicit timezone-aware validity interval and budget calendar, for example Europe/Berlin. | +| `worker_id`, `activity_definition_id`, `target_repo`, `project` | Exact admitted scope; target is an absolute path. Actor is `agt`. | +| `profile_ref`, `profile_sha256`, `descriptor_sha256`, `repository_grant_id` | Versioned profile plus SHA-256 of canonical `model_dump(mode="json")` for the resolved profile and descriptor; `spend_admission.digest` supplies the canonical serializer. Grant identity is the existing `RepositoryGrant.grant_id`. | +| `max_budget_usd`, `max_liability_usd`, `eur_per_usd`, `per_run_eur`, `daily_eur`, `total_eur` | Positive decimal **strings**. Explicit liability/FX must fit per-run, daily and total ceilings. No implicit exchange rate or provider price. | +| `max_turns` | Positive integer matching the resolved native profile. | + +Provision the ledger once using `rein-aharness spend init --policy +--ledger `. Initialization refuses an existing file. Worker execution never +creates an empty replacement. Keep the accepted policy immutable; replacing its +contents under the same ledger refuses, rather than resetting spent capacity. + +Inspect with `rein-aharness spend status --policy --ledger `. +If a reservation is held, the operator must first confirm provider execution has +stopped and obtain final accounting. Record the evidence in the owning work record, +then run: + +```text +rein-aharness spend reconcile --policy --ledger \ + --run-id --cost-usd --receipt \ + --provider-stopped +``` + +The flag is an operator attestation, not an automatic provider query. The local +receipt must refer to reviewed termination and final accounting; it must contain +no secrets. Reconciliation charges at least the full reservation and any higher +observed liability, cannot discard a known higher cost, is idempotent for the same +receipt/cost, and never authorizes rerunning that demand. There is no refund, +reset, policy migration or breach-clear command. A breached envelope requires a +new owner decision that accounts for existing spend before any further admission. + +## Proof and remaining boundary + +`tests/test_spend_admission.py` covers pre-dispatch denial, concurrent processes, +crash/reopen, identity/profile/grant changes, replay, unknown and overrun outcomes, +midnight, total exhaustion, decimal rounding and operator reconciliation. It is +included in both runtime-contract and recovery gates. + +The opt-in `REIN_REAL_BWRAP=1` test also runs actual bwrap, sandbox owner transport, +Glas, repository import and durable-close replay with admission enabled. Queue +responses, accounting and authoring are deterministic fixtures; no model is +called. The fixture's policy and FX have no production authority. + +This ledger enforces allocation against **declared** maximum liability. It cannot +make an opaque provider/tool loop honor that maximum, establish a hard EUR limit, +or verify the operator's FX assumption. G0 remains blocked until those semantics +are proven for the pinned provider/CLI, including retries, cache, subagents and +in-flight work. Exact credential/identity/egress/placement and protected-runtime +installation are also still required before natural factory execution. diff --git a/rein_aharness/claim_loop.py b/rein_aharness/claim_loop.py index f6683a4..49460ae 100644 --- a/rein_aharness/claim_loop.py +++ b/rein_aharness/claim_loop.py @@ -38,6 +38,7 @@ from rein_aharness.execution_cancel import ( from rein_aharness.glas_execution import ( GLAS_APPROACH, GlasExecutionError, + GlasSpendError, execute_profiled_run, normalise_execution_evidence_for_close, ) @@ -50,6 +51,7 @@ from rein_aharness.repository_transaction import ( RepositoryTransaction, RepositoryTransactionError, ) +from rein_aharness.spend_admission import SpendAdmissionError, worker_spend from rein_aharness.taskspec import TaskSpecError from rein_aharness.ops_run_client import ( ActivityCoreOpsClient, @@ -292,6 +294,16 @@ def process_one( if replay.quarantined: logger.error("quarantined %s close evidence entries", replay.quarantined) + try: + spend = worker_spend(cfg) + if spend is not None: + spend.preflight() + except SpendAdmissionError as exc: + return ProcessResult( + claimed=False, retry_full_interval=True, + ok=False, reason=f"spend admission refused: {exc}", + ) + try: claimed = client.claim(limit=1) except OpsRunError as exc: @@ -329,6 +341,16 @@ def process_one( } }, ) + if spend is not None and not run.harness_profile_ref: + reason = "spend admission requires a governed harness profile" + try: + closed = client.fail(run.id, error=reason, reopen=False, + result={"ok": False, "reason": reason}) + except OpsRunError: + return ProcessResult(claimed=True, run_id=run.id, ok=False, + reason="spend route refusal close failed") + return ProcessResult(claimed=True, run_id=run.id, ok=False, + reason=reason, ops_state=closed.state) approach = GLAS_APPROACH if run.harness_profile_ref else select_approach(run) logger.info( "claimed run_id=%s approach=%s title=%r labels=%s", @@ -590,9 +612,9 @@ def _process_profiled_run_active( execution_reason = _bounded_reason( f"refused: {exc}", default="repository transaction refused" ) - except GlasExecutionError: - execution_error = "GlasExecutionError" - execution_reason = "profiled execution failed (GlasExecutionError)" + except GlasExecutionError as exc: + execution_error = type(exc).__name__ + execution_reason = str(exc) if isinstance(exc, GlasSpendError) else "profiled execution failed (GlasExecutionError)" if tx is not None and tx.baseline is not None: tx_evidence = tx.evidence() diff --git a/rein_aharness/cli.py b/rein_aharness/cli.py index 4b6e1fe..7aaa05a 100644 --- a/rein_aharness/cli.py +++ b/rein_aharness/cli.py @@ -190,6 +190,24 @@ def _cmd_close_outbox(args: argparse.Namespace) -> int: return 1 if payload["pending"] or payload["quarantined"] else 0 +def _cmd_spend(args: argparse.Namespace) -> int: + from rein_aharness.spend_admission import SpendAdmissionError, SpendLedger, SpendPolicy + + try: + store = SpendLedger(args.ledger, SpendPolicy.load(args.policy)) + if args.action == "init": + store.initialize() + elif args.action == "reconcile": + if not args.run_id or args.cost_usd is None or not args.receipt or not args.provider_stopped: + raise SpendAdmissionError("reconcile requires run ID, final cost, receipt and --provider-stopped") + store.reconcile(args.run_id, cost_usd=args.cost_usd, receipt=args.receipt) + print(json.dumps(store.status(), indent=2, sort_keys=True)) + return 0 + except SpendAdmissionError as exc: + print(f"spend admission refused: {exc}", file=sys.stderr) + return 2 + + def _cmd_preflight(args: argparse.Namespace) -> int: from rein_aharness.readiness import run_readiness_checks @@ -584,6 +602,15 @@ def main(argv: list[str] | None = None) -> int: help="Skip the read-only Activity Core queue probe", ) + spend = sub.add_parser("spend", help="Inspect or reconcile private worker spend reservations") + spend.add_argument("action", choices=("init", "status", "reconcile")) + spend.add_argument("--policy", type=Path, required=True) + spend.add_argument("--ledger", type=Path, required=True) + spend.add_argument("--run-id") + spend.add_argument("--cost-usd") + spend.add_argument("--receipt", help="Reference to owner-verified termination and final accounting") + spend.add_argument("--provider-stopped", action="store_true") + args = parser.parse_args(argv) if args.command == "validate": @@ -601,6 +628,9 @@ def main(argv: list[str] | None = None) -> int: if args.command == "close-outbox": return _cmd_close_outbox(args) + if args.command == "spend": + return _cmd_spend(args) + if args.command == "preflight": return _cmd_preflight(args) diff --git a/rein_aharness/glas_execution.py b/rein_aharness/glas_execution.py index 9bc9cab..2420bf4 100644 --- a/rein_aharness/glas_execution.py +++ b/rein_aharness/glas_execution.py @@ -14,6 +14,7 @@ from typing import Any from rein_aharness.execution_cancel import ExecutionCancel, ExecutionCancelled, resolve_cancel from rein_aharness.ops_run_client import OpsRun, OpsRunConfig, resolve_ops_target +from rein_aharness.spend_admission import SpendAdmissionError, worker_spend GLAS_APPROACH = "glas-profile" GLAS_ACTOR = "agt" @@ -61,6 +62,10 @@ class GlasExecutionError(RuntimeError): """The authoritative Glas invocation could not produce a GatewayResult.""" +class GlasSpendError(GlasExecutionError): + """Spend refusal with bounded, operator-safe reason text.""" + + def normalise_execution_evidence_for_close(raw: Any) -> dict[str, Any]: """Retain only bounded Glas evidence fields safe for durable close state.""" if not isinstance(raw, dict): @@ -120,7 +125,7 @@ def _request_kwargs(run: OpsRun, config: OpsRunConfig, report_to_hub: bool) -> d # actor type. Sand-boxer validates governed execution actors as # adm|agt|atm, so this runtime enters the gateway as an agent. "actor": GLAS_ACTOR, - "project": "rein-aharness", + "project": config.execution_project, "request_id": run.id, "report_to_hub": report_to_hub, } @@ -168,10 +173,25 @@ def execute_profiled_run( try: request = request_factory(**_request_kwargs(run, config, report_to_hub)) - result = gateway(request, artifact_capture=transfer.capture) if transfer else gateway(request) + kwargs = {"artifact_capture": transfer.capture} if transfer else {} + spend = worker_spend(config) + if spend is not None: + from glas_harness.profiles import ProfileCatalog + catalog = ProfileCatalog() + profile, descriptor = catalog.resolve(run.harness_profile_ref) + catalog.require_operational(profile) + spend.validate_dispatch(run, config, request, profile, descriptor) + # The same cached catalog supplies the checked profile to Glas. + kwargs["catalog"] = catalog + if guard is not None: + guard.check() + spend.reserve(run) + result = gateway(request, **kwargs) raw = result.model_dump(mode="json") if hasattr(result, "model_dump") else result except ExecutionCancelled: raise + except SpendAdmissionError as exc: + raise GlasSpendError(f"spend admission refused: {exc}") from None except GlasExecutionError: raise except Exception as exc: @@ -185,6 +205,12 @@ def execute_profiled_run( raise GlasExecutionError("Glas gateway returned an invalid GatewayResult") if not isinstance(raw.get("evidence"), dict): raise GlasExecutionError("Glas GatewayResult is missing execution evidence") + if spend is not None: + try: + if not spend.observe(run.id, raw): + raise SpendAdmissionError("execution accounting requires reconciliation") + except SpendAdmissionError as exc: + raise GlasSpendError(f"spend accounting refused: {exc}") from None if transfer is not None and raw["ok"]: evidence = raw["evidence"] if evidence.get("session_cleanup") != "succeeded" or evidence.get("sandbox_destroy") != "succeeded": diff --git a/rein_aharness/ops_run_client.py b/rein_aharness/ops_run_client.py index 52ff281..ef7661e 100644 --- a/rein_aharness/ops_run_client.py +++ b/rein_aharness/ops_run_client.py @@ -137,6 +137,10 @@ class OpsRunConfig: repo_map: dict[str, str] = field(default_factory=dict) repo_roots: tuple[str, ...] = DEFAULT_REPO_ROOTS timeout: float = 30.0 + execution_project: str = "rein-aharness" + require_spend_admission: bool = False + spend_policy_path: str | None = None + spend_ledger_path: str | None = None @classmethod def from_env(cls) -> "OpsRunConfig": @@ -173,6 +177,10 @@ class OpsRunConfig: lease_seconds=lease, repo_map=repo_map, repo_roots=roots or DEFAULT_REPO_ROOTS, + execution_project=os.environ.get("AGENT_HARNESS_EXECUTION_PROJECT", "rein-aharness"), + require_spend_admission=bool(os.environ.get("AGENT_HARNESS_REQUIRE_SPEND_ADMISSION", "")), + spend_policy_path=os.environ.get("AGENT_HARNESS_SPEND_POLICY"), + spend_ledger_path=os.environ.get("AGENT_HARNESS_SPEND_LEDGER"), ) diff --git a/rein_aharness/spend_admission.py b/rein_aharness/spend_admission.py new file mode 100644 index 0000000..ee7b2d2 --- /dev/null +++ b/rein_aharness/spend_admission.py @@ -0,0 +1,489 @@ +"""Worker-owned, conservative spend reservations; no provider or secret access. + +One private SQLite file per envelope, shared by all admitted processes on one +host. Missing state is an error: only the explicit operator init creates it. +""" + +from __future__ import annotations + +import hashlib +import json +import os +import re +import sqlite3 +import stat +from collections.abc import Iterator +from contextlib import contextmanager +from dataclasses import dataclass +from datetime import UTC, datetime +from decimal import ROUND_CEILING, ROUND_FLOOR, Decimal, InvalidOperation, localcontext +from pathlib import Path +from typing import Any +from zoneinfo import ZoneInfo + + +class SpendAdmissionError(RuntimeError): + """A bounded refusal; never contains prompt, credential or provider output.""" + + +def digest(value: Any) -> str: + return hashlib.sha256( + json.dumps( + value, sort_keys=True, separators=(",", ":"), allow_nan=False + ).encode() + ).hexdigest() + + +def amount(value: Any, *, positive: bool = True) -> Decimal: + if isinstance(value, bool) or not isinstance(value, (str, int, float)): + raise SpendAdmissionError("invalid monetary amount") + if len(str(value)) > 500: + raise SpendAdmissionError("invalid monetary amount") + try: + result = Decimal(str(value)) + except InvalidOperation: + raise SpendAdmissionError("invalid monetary amount") from None + if ( + not result.is_finite() + or result < 0 + or (positive and result == 0) + or result > 1_000_000 + or abs(result.as_tuple().exponent) > 100 + ): + raise SpendAdmissionError("invalid monetary amount") + return result + + +def micros(value: Decimal) -> int: + with localcontext() as context: + context.prec = 256 + return int((value * 1_000_000).to_integral_value(rounding=ROUND_CEILING)) + + +def cap_micros(value: str) -> int: + with localcontext() as context: + context.prec = 256 + return int((amount(value) * 1_000_000).to_integral_value(rounding=ROUND_FLOOR)) + + +def converted_micros(usd: Any, eur_per_usd: str) -> int: + with localcontext() as context: + context.prec = 256 + return micros(amount(usd, positive=False) * amount(eur_per_usd)) + + +def timestamp(value: str) -> datetime: + try: + result = datetime.fromisoformat(value) + if result.tzinfo is None: + raise ValueError + return result.astimezone(UTC) + except (TypeError, ValueError): + raise SpendAdmissionError("invalid policy timestamp") from None + + +@dataclass(frozen=True) +class SpendPolicy: + version: str + envelope_id: str + authority_ref: str + valid_from: str + expires_at: str + timezone: str + worker_id: str + activity_definition_id: str + target_repo: str + project: str + profile_ref: str + profile_sha256: str + descriptor_sha256: str + repository_grant_id: str + max_budget_usd: str + max_liability_usd: str + max_turns: int + eur_per_usd: str + per_run_eur: str + daily_eur: str + total_eur: str + + def __post_init__(self) -> None: + if self.version != "1": + raise SpendAdmissionError("unsupported spend policy version") + for name in self.__dataclass_fields__: + value = getattr(self, name) + if name == "max_turns": + if type(value) is not int or not 1 <= value <= 10000: + raise SpendAdmissionError("invalid turn limit") + elif ( + not isinstance(value, str) + or not value + or len(value) > 500 + or any(ord(c) < 32 for c in value) + ): + raise SpendAdmissionError("invalid spend policy field") + for name in ("profile_sha256", "descriptor_sha256"): + if not re.fullmatch(r"[0-9a-f]{64}", getattr(self, name)): + raise SpendAdmissionError("invalid runtime digest") + if not Path(self.target_repo).is_absolute() or "@" not in self.profile_ref: + raise SpendAdmissionError( + "spend policy needs an absolute target and versioned profile" + ) + try: + ZoneInfo(self.timezone) + except (ValueError, KeyError): + raise SpendAdmissionError("invalid budget timezone") from None + if timestamp(self.valid_from) >= timestamp(self.expires_at): + raise SpendAdmissionError("empty policy validity interval") + for name in ( + "max_budget_usd", + "max_liability_usd", + "eur_per_usd", + "per_run_eur", + "daily_eur", + "total_eur", + ): + amount(getattr(self, name)) + if amount(self.max_budget_usd) > amount(self.max_liability_usd): + raise SpendAdmissionError("native limit exceeds declared liability") + if self.reservation > cap_micros(self.per_run_eur) or not amount( + self.per_run_eur + ) <= amount(self.daily_eur) <= amount(self.total_eur): + raise SpendAdmissionError("liability exceeds operating envelope") + + @property + def reservation(self) -> int: + return converted_micros(self.max_liability_usd, self.eur_per_usd) + + @property + def sha256(self) -> str: + return digest(self.__dict__) + + @classmethod + def load(cls, path: Path) -> SpendPolicy: + _private_file(path) + try: + if path.stat().st_size > 16384: + raise ValueError + + def unique(pairs): + result = {} + for key, value in pairs: + if key in result: + raise ValueError + result[key] = value + return result + + policy = cls(**json.loads(path.read_text(), object_pairs_hook=unique)) + except (OSError, ValueError, TypeError): + raise SpendAdmissionError("invalid spend policy file") from None + if path.resolve().is_relative_to(Path(policy.target_repo).resolve()): + raise SpendAdmissionError("spend policy must be outside target repository") + return policy + + +def _private_file(path: Path) -> None: + try: + info = path.lstat() + except OSError: + raise SpendAdmissionError("required private spend file unavailable") from None + if ( + not stat.S_ISREG(info.st_mode) + or info.st_nlink != 1 + or info.st_uid != os.getuid() + or info.st_mode & 0o077 + ): + raise SpendAdmissionError( + "spend file must be private, regular and worker-owned" + ) + + +class SpendLedger: + def __init__(self, path: Path, policy: SpendPolicy) -> None: + self.path, self.policy = path, policy + parent = path.parent + if not path.is_absolute() or path.resolve().is_relative_to( + Path(policy.target_repo).resolve() + ): + raise SpendAdmissionError( + "ledger must be absolute and outside target repository" + ) + try: + info = parent.lstat() + except OSError: + raise SpendAdmissionError("private ledger directory unavailable") from None + if ( + not stat.S_ISDIR(info.st_mode) + or info.st_uid != os.getuid() + or info.st_mode & 0o077 + ): + raise SpendAdmissionError( + "ledger directory must be private and worker-owned" + ) + + def initialize(self) -> None: + """Explicit provisioning only. Never truncate or recreate existing state.""" + try: + fd = os.open(self.path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600) + os.close(fd) + with self._db(check=False) as db: + db.execute( + "CREATE TABLE envelope (policy TEXT NOT NULL, last_seen TEXT NOT NULL, breached INTEGER NOT NULL CHECK(breached IN (0,1)))" + ) + db.execute("""CREATE TABLE reservations ( + run_id TEXT PRIMARY KEY, definition_id TEXT NOT NULL, + idempotency_key TEXT NOT NULL, attempt INTEGER NOT NULL, + state TEXT NOT NULL CHECK(state IN ('held','charged')), + liability INTEGER NOT NULL CHECK(liability>0), + start_day TEXT NOT NULL, end_day TEXT, + observed_usd TEXT, receipt TEXT, + UNIQUE(definition_id, idempotency_key), + CHECK((state='held' AND end_day IS NULL) OR + (state='charged' AND end_day IS NOT NULL AND end_day>=start_day))); + """) + db.execute( + "INSERT INTO envelope VALUES (?, ?, 0)", + (self.policy.sha256, self.policy.valid_from), + ) + fd = os.open(self.path.parent, os.O_RDONLY | os.O_DIRECTORY) + try: + os.fsync(fd) + finally: + os.close(fd) + except OSError: + raise SpendAdmissionError("spend ledger initialization refused") from None + + @contextmanager + def _db(self, *, check: bool = True) -> Iterator[sqlite3.Connection]: + _private_file(self.path) + db = None + try: + db = sqlite3.connect( + self.path.as_uri() + "?mode=rw", + uri=True, + timeout=10, + isolation_level=None, + ) + db.row_factory = sqlite3.Row + db.execute("PRAGMA synchronous=FULL") + db.execute("BEGIN IMMEDIATE") + if check: + rows = db.execute("SELECT * FROM envelope").fetchall() + if len(rows) != 1 or rows[0]["policy"] != self.policy.sha256: + raise SpendAdmissionError("spend ledger policy mismatch") + yield db + db.commit() + except sqlite3.Error: + raise SpendAdmissionError( + "spend ledger unavailable or inconsistent" + ) from None + finally: + if db is not None: + db.close() # rolls back any incomplete mutation + + def _clock(self, db: sqlite3.Connection, now: datetime, *, admission: bool) -> str: + if now.tzinfo is None: + raise SpendAdmissionError("budget clock must be timezone-aware") + now = now.astimezone(UTC) + previous = timestamp(db.execute("SELECT last_seen FROM envelope").fetchone()[0]) + if now < previous: + raise SpendAdmissionError("budget clock moved backwards") + if admission and not timestamp(self.policy.valid_from) <= now < timestamp( + self.policy.expires_at + ): + raise SpendAdmissionError("spend policy is not currently valid") + db.execute("UPDATE envelope SET last_seen=?", (now.isoformat(),)) + return now.astimezone(ZoneInfo(self.policy.timezone)).date().isoformat() + + def _capacity(self, db: sqlite3.Connection, day: str) -> None: + if db.execute("SELECT breached FROM envelope").fetchone()[0]: + raise SpendAdmissionError( + "spend envelope breached; new owner admission required" + ) + if db.execute("SELECT 1 FROM reservations WHERE state='held'").fetchone(): + raise SpendAdmissionError( + "unresolved spend reservation; reconcile before dispatch" + ) + total = db.execute( + "SELECT COALESCE(SUM(liability),0) FROM reservations" + ).fetchone()[0] + daily = db.execute( + "SELECT COALESCE(SUM(liability),0) FROM reservations WHERE start_day<=? AND end_day>=?", + (day, day), + ).fetchone()[0] + if total + self.policy.reservation > cap_micros(self.policy.total_eur): + raise SpendAdmissionError("total spend capacity exhausted") + if daily + self.policy.reservation > cap_micros(self.policy.daily_eur): + raise SpendAdmissionError("daily spend capacity exhausted") + + def preflight(self, *, now: datetime | None = None) -> None: + with self._db() as db: + day = self._clock(db, now or datetime.now(UTC), admission=True) + self._capacity(db, day) + + def validate_dispatch( + self, run: Any, config: Any, request: Any, profile: Any, descriptor: Any + ) -> None: + p = self.policy + if ( + config.worker_id != p.worker_id + or run.claim_owner != p.worker_id + or run.activity_definition_id != p.activity_definition_id + or request.project != p.project + or request.actor != "agt" + or Path(request.repo).resolve() != Path(p.target_repo).resolve() + or run.harness_profile_ref != p.profile_ref + or request.harness_profile_ref != p.profile_ref + or request.request_id != run.id + or not run.repository_grant + or run.repository_grant.grant_id != p.repository_grant_id + or digest(profile.model_dump(mode="json")) != p.profile_sha256 + or digest(descriptor.model_dump(mode="json")) != p.descriptor_sha256 + ): + raise SpendAdmissionError("dispatch does not match admitted spend scope") + if ( + profile.rein.id != "rein-aharness" + or amount(profile.limits.max_budget_usd) != amount(p.max_budget_usd) + or profile.limits.max_turns != p.max_turns + ): + raise SpendAdmissionError( + "resolved native limits do not match spend policy" + ) + + def reserve(self, run: Any, *, now: datetime | None = None) -> None: + for value in (run.id, run.activity_definition_id, run.idempotency_key): + if not isinstance(value, str) or not value or len(value) > 500: + raise SpendAdmissionError("invalid spend run identity") + if type(run.attempt) is not int or run.attempt < 1: + raise SpendAdmissionError("invalid spend attempt") + with self._db() as db: + day = self._clock(db, now or datetime.now(UTC), admission=True) + if db.execute( + "SELECT 1 FROM reservations WHERE run_id=? OR (definition_id=? AND idempotency_key=?)", + (run.id, run.activity_definition_id, run.idempotency_key), + ).fetchone(): + raise SpendAdmissionError( + "run already admitted; workload replay refused" + ) + self._capacity(db, day) + db.execute( + "INSERT INTO reservations VALUES (?, ?, ?, ?, 'held', ?, ?, NULL, NULL, NULL)", + ( + run.id, + run.activity_definition_id, + run.idempotency_key, + run.attempt, + self.policy.reservation, + day, + ), + ) + + def observe( + self, run_id: str, raw: dict[str, Any], *, now: datetime | None = None + ) -> bool: + """Charge full liability on complete success; never refund self-reported cost.""" + evidence = raw.get("evidence", {}) + try: + cost = amount(evidence.get("cost_usd"), positive=False) + except SpendAdmissionError: + return False # Unknown remains held, even for apparent success. + with self._db() as db: + day = self._clock(db, now or datetime.now(UTC), admission=False) + row = db.execute( + "SELECT * FROM reservations WHERE run_id=?", (run_id,) + ).fetchone() + if row is None or row["state"] != "held": + raise SpendAdmissionError("reservation observation conflict") + liability = max( + row["liability"], converted_micros(str(cost), self.policy.eur_per_usd) + ) + breached = cost > amount(self.policy.max_liability_usd) + complete = ( + raw.get("ok") is True + and evidence.get("outcome") == "succeeded" + and evidence.get("request_id") == run_id + and evidence.get("profile_ref") == self.policy.profile_ref + and evidence.get("session_cleanup") == "succeeded" + and evidence.get("sandbox_destroy") == "succeeded" + and not breached + ) + db.execute( + "UPDATE reservations SET liability=?, observed_usd=?, state=?, end_day=? WHERE run_id=?", + ( + liability, + str(cost), + "charged" if complete else "held", + day if complete else None, + run_id, + ), + ) + if breached: + db.execute("UPDATE envelope SET breached=1") + return complete + + def reconcile( + self, run_id: str, *, cost_usd: str, receipt: str, now: datetime | None = None + ) -> None: + """Operator attests provider termination + final accounting; no retry grant.""" + cost = amount(cost_usd, positive=False) + if not re.fullmatch(r"[A-Za-z0-9][A-Za-z0-9._:/@+-]{0,199}", receipt): + raise SpendAdmissionError("bounded reconciliation receipt required") + with self._db() as db: + day = self._clock(db, now or datetime.now(UTC), admission=False) + row = db.execute( + "SELECT * FROM reservations WHERE run_id=?", (run_id,) + ).fetchone() + if row is None: + raise SpendAdmissionError("unknown reservation") + if row["state"] != "held": + if row["receipt"] == receipt and row["observed_usd"] == str(cost): + return + raise SpendAdmissionError( + "reconciliation conflicts with settled reservation" + ) + if row["observed_usd"] is not None and cost < amount( + row["observed_usd"], positive=False + ): + raise SpendAdmissionError( + "reconciliation cannot discard observed liability" + ) + liability = max( + row["liability"], converted_micros(str(cost), self.policy.eur_per_usd) + ) + db.execute( + "UPDATE reservations SET state='charged', liability=?, observed_usd=?, end_day=?, receipt=? WHERE run_id=?", + (liability, str(cost), day, receipt, run_id), + ) + if cost > amount(self.policy.max_liability_usd): + db.execute("UPDATE envelope SET breached=1") + + def status(self) -> dict[str, Any]: + with self._db() as db: + return { + "envelope_id": self.policy.envelope_id, + "policy_sha256": self.policy.sha256, + "breached": bool( + db.execute("SELECT breached FROM envelope").fetchone()[0] + ), + "reservations": [ + dict(row) + for row in db.execute( + "SELECT * FROM reservations ORDER BY start_day, run_id" + ) + ], + } + + +def worker_spend(config: Any) -> SpendLedger | None: + if config.spend_policy_path is None and config.spend_ledger_path is None: + if config.require_spend_admission: + raise SpendAdmissionError("required spend admission is not configured") + return None + if not config.spend_policy_path or not config.spend_ledger_path: + raise SpendAdmissionError("both spend policy and ledger must be configured") + policy = SpendPolicy.load(Path(config.spend_policy_path)) + if ( + config.worker_id != policy.worker_id + or config.execution_project != policy.project + ): + raise SpendAdmissionError("worker identity does not match spend policy") + return SpendLedger(Path(config.spend_ledger_path), policy) diff --git a/scripts/verify-recovery-contracts.sh b/scripts/verify-recovery-contracts.sh index 75479cf..5c92383 100755 --- a/scripts/verify-recovery-contracts.sh +++ b/scripts/verify-recovery-contracts.sh @@ -9,6 +9,7 @@ PYTHON_BIN="${REIN_CONTRACT_PYTHON:-${REPO_ROOT}/.venv/bin/python}" "${PYTHON_BIN}" -c 'import glas_harness.contract, llm_connect, sandboxer.models' PYTHONPATH="${REPO_ROOT}:${REPO_ROOT}/../llm-connect" \ "${PYTHON_BIN}" -m pytest \ + tests/test_spend_admission.py \ tests/test_claim_loop.py::test_process_one_initial_heartbeat_rejection_refuses_dispatch \ tests/test_claim_loop.py::test_process_one_lease_loss_cancels_registered_adapter_process \ tests/test_claim_loop.py::test_profiled_close_failure_happens_after_repository_lock_release \ diff --git a/scripts/verify-runtime-contracts.sh b/scripts/verify-runtime-contracts.sh index ad1db9f..86de2b3 100755 --- a/scripts/verify-runtime-contracts.sh +++ b/scripts/verify-runtime-contracts.sh @@ -10,6 +10,7 @@ PYTHON_BIN="${REIN_CONTRACT_PYTHON:-${REPO_ROOT}/.venv/bin/python}" "${PYTHON_BIN}" -c 'import glas_harness.contract, llm_connect, sandboxer.models' PYTHONPATH="${REPO_ROOT}:${REPO_ROOT}/../llm-connect" \ "${PYTHON_BIN}" -m pytest \ + tests/test_spend_admission.py \ tests/test_glas_execution.py \ tests/test_ops_run_client.py \ tests/test_claim_loop.py \ diff --git a/tests/test_repository_artifact_bwrap.py b/tests/test_repository_artifact_bwrap.py index 4583c50..1bd0cc6 100644 --- a/tests/test_repository_artifact_bwrap.py +++ b/tests/test_repository_artifact_bwrap.py @@ -9,6 +9,7 @@ import pytest from glas_harness.contract import ( ExecutionSummary, + ExecutionLimits, OperationalReadiness, Rein, ToolResult, @@ -29,6 +30,7 @@ from rein_aharness.ops_run_client import ( ) from test_repository_artifact import commit, git from rein_aharness.repository_grant import RepositoryGrant +from rein_aharness.spend_admission import SpendLedger, SpendPolicy, digest pytestmark = pytest.mark.skipif( os.environ.get("REIN_REAL_BWRAP") != "1", reason="opt-in kernel namespace proof" @@ -69,7 +71,7 @@ print(git('rev-parse','HEAD')) def end_session(self, session): return ExecutionSummary( - committed=True, commit_sha=self.head, outcome="succeeded", tokens_spent=0 + committed=True, commit_sha=self.head, outcome="succeeded", tokens_spent=0, cost_usd=0.0 ) def cleanup_session(self, session): @@ -77,7 +79,8 @@ print(git('rev-parse','HEAD')) assert git(self.source, "rev-parse", "HEAD") == self.baseline -def test_real_bwrap_worker_import_and_close_replay(tmp_path, monkeypatch): +@pytest.mark.parametrize("with_spend", [False, True]) +def test_real_bwrap_worker_import_and_close_replay(tmp_path, monkeypatch, with_spend): monkeypatch.setenv("XDG_DATA_HOME", str(tmp_path / "data")) monkeypatch.setenv("REIN_AHARNESS_STATE_DIR", str(tmp_path / "state")) monkeypatch.setenv("SANDBOXER_NO_STATE_HUB", "1") @@ -129,6 +132,7 @@ def test_real_bwrap_worker_import_and_close_replay(tmp_path, monkeypatch): profile, _ = catalog.resolve(run.harness_profile_ref) catalog.profiles()[(profile.id, profile.version)] = profile.model_copy( update={ + "limits": ExecutionLimits(max_budget_usd=4.0, max_turns=8), "operational_readiness": OperationalReadiness( status="ready", reason="isolated test fixture", @@ -137,12 +141,37 @@ def test_real_bwrap_worker_import_and_close_replay(tmp_path, monkeypatch): ) } ) + spend = None + if with_spend: + private = tmp_path / "spend-state" + private.mkdir(mode=0o700) + admitted_profile, descriptor = catalog.resolve(run.harness_profile_ref) + policy = SpendPolicy( + version="1", envelope_id="bwrap-fixture", authority_ref="test:no-live-authority", + valid_from="2026-09-01T00:00:00Z", expires_at="2099-01-01T00:00:00Z", + timezone="Europe/Berlin", worker_id="fixture-worker", + activity_definition_id="fixture-definition", target_repo=str(source), + project="fixture-factory", profile_ref=run.harness_profile_ref, + profile_sha256=digest(admitted_profile.model_dump(mode="json")), + descriptor_sha256=digest(descriptor.model_dump(mode="json")), + repository_grant_id=grant.grant_id, max_budget_usd="4", max_liability_usd="5", + max_turns=8, eur_per_usd="1", per_run_eur="5", daily_eur="10", total_eur="15", + ) + policy_path = private / "policy.json" + policy_path.write_text(json.dumps(policy.__dict__)) + policy_path.chmod(0o600) + spend = SpendLedger(private / "spend.sqlite3", policy) + spend.initialize() + client.config.execution_project = "fixture-factory" + client.config.spend_policy_path = str(policy_path) + client.config.spend_ledger_path = str(spend.path) + monkeypatch.setattr("glas_harness.profiles.ProfileCatalog", lambda: catalog) manager = SandboxManager(store=SandboxStore(tmp_path / "sandboxes.json")) rein = DeterministicRein(source, baseline) monkeypatch.setattr( "glas_harness.gateway.run_execution", lambda request, **kwargs: run_execution( - request, catalog=catalog, rein=rein, manager=manager, **kwargs + request, catalog=kwargs.pop("catalog", catalog), rein=rein, manager=manager, **kwargs ), ) outbox = CloseOutbox(state_dir=tmp_path / "state") @@ -165,7 +194,14 @@ def test_real_bwrap_worker_import_and_close_replay(tmp_path, monkeypatch): client.claim.return_value = [] second = process_one(client, outbox=outbox, report_to_hub=False) assert second.empty and rein.calls == 1 + assert outbox.status() == {"pending": 0, "delivered": 1, "quarantined": 0} + if spend is not None: + rows = spend.status()["reservations"] + assert len(rows) == 1 and rows[0]["state"] == "charged" + assert rows[0]["liability"] == 5_000_000 + client.claim.return_value = [run] + replay = process_one(client, outbox=outbox, report_to_hub=False) + assert not replay.ok and rein.calls == 1 assert client.complete.call_count == 2 assert client.heartbeat.call_count >= 1 assert git(source, "rev-list", "--count", "HEAD") == "2" - assert outbox.status() == {"pending": 0, "delivered": 1, "quarantined": 0} diff --git a/tests/test_spend_admission.py b/tests/test_spend_admission.py new file mode 100644 index 0000000..6ed8f67 --- /dev/null +++ b/tests/test_spend_admission.py @@ -0,0 +1,492 @@ +"""No-provider proofs for admission, persistence, accounting and recovery.""" + +import json +from concurrent.futures import ProcessPoolExecutor +from dataclasses import replace +from datetime import UTC, datetime, timedelta +from pathlib import Path +from types import SimpleNamespace +from unittest.mock import MagicMock + +import pytest + +from rein_aharness.ops_run_client import ActivityCoreOpsClient, OpsRun, OpsRunConfig +from rein_aharness.repository_grant import RepositoryGrant +from rein_aharness.spend_admission import ( + SpendAdmissionError, + SpendLedger, + SpendPolicy, + digest, + worker_spend, +) + +NOW = datetime(2026, 9, 9, 12, tzinfo=UTC) + + +def run(key="one"): + return OpsRun( + id=f"run-{key}", + activity_definition_id="def-1", + idempotency_key=key, + target_repo="target", + title="private intent", + description="private prompt", + attempt=1, + claim_owner="worker-1", + harness_profile_ref="harness.test@1.0.0", + repository_grant=RepositoryGrant("1", ("docs/",), 1, 1, False), + ) + + +@pytest.fixture +def ledger(tmp_path): + private = tmp_path / "private" + private.mkdir(mode=0o700) + policy = SpendPolicy( + version="1", + envelope_id="test-only", + authority_ref="fixture:no-live-authority", + valid_from="2026-09-01T00:00:00Z", + expires_at="2099-01-01T00:00:00Z", + timezone="Europe/Berlin", + worker_id="worker-1", + activity_definition_id="def-1", + target_repo=str(tmp_path / "target"), + project="fixture-factory", + profile_ref="harness.test@1.0.0", + profile_sha256="a" * 64, + descriptor_sha256="b" * 64, + repository_grant_id=run().repository_grant.grant_id, + max_budget_usd="4", + max_liability_usd="5", + max_turns=8, + eur_per_usd="1", + per_run_eur="5", + daily_eur="10", + total_eur="15", + ) + store = SpendLedger(private / "spend.sqlite3", policy) + store.initialize() + return store + + +def success(store, item, cost=1): + return { + "ok": True, + "evidence": { + "request_id": item.id, + "profile_ref": store.policy.profile_ref, + "cost_usd": cost, + "outcome": "succeeded", + "session_cleanup": "succeeded", + "sandbox_destroy": "succeeded", + }, + } + + +def configured(store): + path = store.path.parent / "policy.json" + path.write_text(json.dumps(store.policy.__dict__)) + path.chmod(0o600) + return OpsRunConfig( + worker_id="worker-1", + execution_project="fixture-factory", + spend_policy_path=str(path), + spend_ledger_path=str(store.path), + repo_roots=(str(Path(store.policy.target_repo).parent),), + ) + + +def test_daily_total_and_no_self_reported_refund(ledger): + for key in ("one", "two"): + item = run(key) + ledger.reserve(item, now=NOW) + ledger.observe(item.id, success(ledger, item, 0), now=NOW) + with pytest.raises(SpendAdmissionError, match="daily"): + ledger.reserve(run("three"), now=NOW) + tomorrow = NOW + timedelta(days=1) + ledger.reserve(run("three"), now=tomorrow) + ledger.observe(run("three").id, success(ledger, run("three")), now=tomorrow) + with pytest.raises(SpendAdmissionError, match="total"): + ledger.reserve(run("four"), now=tomorrow + timedelta(days=1)) + assert sum(r["liability"] for r in ledger.status()["reservations"]) == 15_000_000 + + +def test_crash_survives_reopen_and_blocks_across_days(ledger): + ledger.reserve(run(), now=NOW) + reopened = SpendLedger(ledger.path, ledger.policy) + for day in (NOW, NOW + timedelta(days=5)): + with pytest.raises(SpendAdmissionError, match="unresolved"): + reopened.preflight(now=day) + reopened.reconcile( + run().id, + cost_usd="2", + receipt="audit:stopped-and-final", + now=NOW + timedelta(days=5), + ) + reopened.reconcile( + run().id, + cost_usd="2", + receipt="audit:stopped-and-final", + now=NOW + timedelta(days=5), + ) + with pytest.raises(SpendAdmissionError, match="replay"): + reopened.reserve(replace(run(), attempt=2), now=NOW + timedelta(days=5)) + with pytest.raises(SpendAdmissionError, match="replay"): + reopened.reserve( + replace(run(), id="different-row"), now=NOW + timedelta(days=5) + ) + assert reopened.status()["reservations"][0]["liability"] == 5_000_000 + + +def test_full_charge_on_both_sides_of_midnight(ledger): + before = datetime(2026, 9, 9, 21, 59, tzinfo=UTC) + after = before + timedelta(minutes=2) + ledger.reserve(run(), now=before) + ledger.observe(run().id, success(ledger, run()), now=after) + row = ledger.status()["reservations"][0] + assert (row["start_day"], row["end_day"]) == ("2026-09-09", "2026-09-10") + ledger.reserve(run("two"), now=after) + ledger.observe(run("two").id, success(ledger, run("two")), now=after) + with pytest.raises(SpendAdmissionError, match="daily"): + ledger.preflight(now=after) + + +@pytest.mark.parametrize("cost", [None, True, -1, "NaN", "Infinity", {}, 1000001]) +def test_unknown_accounting_never_releases_hold(ledger, cost): + ledger.reserve(run(), now=NOW) + ledger.observe(run().id, success(ledger, run(), cost), now=NOW) + with pytest.raises(SpendAdmissionError, match="unresolved"): + ledger.preflight(now=NOW) + + +@pytest.mark.parametrize( + "field,value", + [ + ("request_id", "wrong"), + ("profile_ref", "wrong"), + ("outcome", "failed"), + ("sandbox_destroy", "failed"), + ("session_cleanup", "failed"), + ], +) +def test_incomplete_or_mismatched_result_retains_hold(ledger, field, value): + ledger.reserve(run(), now=NOW) + raw = success(ledger, run()) + raw["evidence"][field] = value + ledger.observe(run().id, raw, now=NOW) + assert ledger.status()["reservations"][0]["state"] == "held" + + +def test_overrun_is_recorded_and_freezes_even_after_reconcile(ledger): + ledger.reserve(run(), now=NOW) + ledger.observe(run().id, success(ledger, run(), 7), now=NOW) + assert ledger.status()["breached"] + assert ledger.status()["reservations"][0]["liability"] == 7_000_000 + with pytest.raises(SpendAdmissionError, match="discard"): + ledger.reconcile(run().id, cost_usd="0", receipt="audit:final", now=NOW) + ledger.reconcile(run().id, cost_usd="7", receipt="audit:final", now=NOW) + with pytest.raises(SpendAdmissionError, match="breached"): + ledger.preflight(now=NOW + timedelta(days=1)) + + +def test_missing_corrupt_changed_policy_and_reinit_refuse(ledger): + with pytest.raises(SpendAdmissionError): + ledger.initialize() + with pytest.raises(SpendAdmissionError, match="policy mismatch"): + SpendLedger(ledger.path, replace(ledger.policy, total_eur="20")).preflight( + now=NOW + ) + ledger.path.unlink() + with pytest.raises(SpendAdmissionError, match="unavailable"): + ledger.preflight(now=NOW) + assert not ledger.path.exists() + ledger.path.write_bytes(b"corrupt") + ledger.path.chmod(0o600) + with pytest.raises(SpendAdmissionError, match="inconsistent"): + ledger.preflight(now=NOW) + + +def test_private_paths_and_policy(ledger): + config = configured(ledger) + assert worker_spend(config).policy.sha256 == ledger.policy.sha256 + Path(config.spend_policy_path).chmod(0o644) + with pytest.raises(SpendAdmissionError, match="private"): + worker_spend(config) + with pytest.raises(SpendAdmissionError, match="both"): + worker_spend(replace(config, spend_policy_path=None)) + with pytest.raises(SpendAdmissionError, match="outside"): + SpendLedger(Path(ledger.policy.target_repo) / "spend.db", ledger.policy) + + +def test_expiry_clock_rollback_and_upward_rounding(ledger): + ledger.preflight(now=NOW) + with pytest.raises(SpendAdmissionError, match="backwards"): + ledger.preflight(now=NOW - timedelta(seconds=1)) + with pytest.raises(SpendAdmissionError, match="valid"): + ledger.preflight(now=datetime(2099, 1, 1, tzinfo=UTC)) + policy = replace(ledger.policy, eur_per_usd="0.99999999") + assert policy.reservation == 5_000_000 + with pytest.raises(SpendAdmissionError, match="envelope"): + replace(ledger.policy, eur_per_usd="1.00000001") + + +def _reserve_process(path, policy, key): + try: + SpendLedger(Path(path), SpendPolicy(**policy)).reserve(run(key), now=NOW) + return "reserved" + except SpendAdmissionError: + return "refused" + + +def test_concurrent_processes_cannot_double_admit(ledger): + with ProcessPoolExecutor(max_workers=4) as pool: + results = list( + pool.map( + _reserve_process, + [str(ledger.path)] * 4, + [ledger.policy.__dict__] * 4, + ["a", "b", "c", "d"], + ) + ) + assert sorted(results) == ["refused", "refused", "refused", "reserved"] + assert len(ledger.status()["reservations"]) == 1 + + +def test_pending_spend_blocks_claim_but_close_replay_still_runs(ledger, monkeypatch): + from rein_aharness.claim_loop import process_one + from rein_aharness.close_outbox import CloseOutbox, CloseRequest + + ledger.reserve(run(), now=NOW) + client = MagicMock(spec=ActivityCoreOpsClient) + client.config = configured(ledger) + client.complete.return_value = SimpleNamespace(state="succeeded") + outbox = CloseOutbox() + outbox.enqueue( + CloseRequest( + run_id="run-older", + transaction_id="tx-1", + worker_id="worker-1", + action="complete", + result={"ok": True}, + ) + ) + result = process_one(client, outbox=outbox) + assert not result.claimed and "unresolved" in result.reason + client.claim.assert_not_called() + client.complete.assert_called_once() + + +def test_nonprofile_row_cannot_bypass_spend_gate(ledger, monkeypatch): + from rein_aharness.claim_loop import process_one + + client = MagicMock(spec=ActivityCoreOpsClient) + client.config = configured(ledger) + client.claim.return_value = [replace(run(), harness_profile_ref=None)] + client.fail.return_value = SimpleNamespace(state="failed") + dispatched = [] + monkeypatch.setattr( + "rein_aharness.claim_loop.execute_approach", + lambda *a, **k: dispatched.append(True), + ) + result = process_one(client) + assert not result.ok and "profile" in result.reason and not dispatched + assert ledger.status()["reservations"] == [] + + +@pytest.fixture +def dispatch_case(ledger, monkeypatch): + import subprocess + + from glas_harness.contract import ExecutionLimits, OperationalReadiness + from glas_harness.profiles import ProfileCatalog + + target = Path(ledger.policy.target_repo) + target.mkdir() + subprocess.run(["git", "init", "-q", str(target)], check=True) + catalog = ProfileCatalog() + profile, descriptor = catalog.resolve("harness.agent-dev-local@1.0.0") + profile = profile.model_copy( + update={ + "id": "harness.test", + "operational_readiness": OperationalReadiness( + status="ready", evidence_ref="test:only" + ), + "limits": ExecutionLimits(max_budget_usd=4.0, max_turns=8), + } + ) + catalog.profiles()[(profile.id, profile.version)] = profile + monkeypatch.setattr("glas_harness.profiles.ProfileCatalog", lambda: catalog) + policy = replace( + ledger.policy, + profile_sha256=digest(profile.model_dump(mode="json")), + descriptor_sha256=digest(descriptor.model_dump(mode="json")), + ) + store = SpendLedger(ledger.path.parent / "dispatch.sqlite3", policy) + store.initialize() + # Artifact transfer is exercised unmodified by the separate actual bwrap test. + transfer = MagicMock() + transfer.import_after_teardown.return_value = {"fixture": True} + monkeypatch.setattr( + "rein_aharness.repository_artifact.RepositoryArtifactTransfer", + lambda *a, **k: transfer, + ) + return store, configured(store), catalog, transfer + + +def test_reservation_precedes_gateway_and_reclaim_never_repeats(dispatch_case): + from rein_aharness.glas_execution import GlasExecutionError, execute_profiled_run + + store, config, catalog, transfer = dispatch_case + calls = [] + + def gateway(request, **kwargs): + assert kwargs["catalog"] is catalog + assert store.status()["reservations"][0]["state"] == "held" + calls.append(request) + return success(store, run()) + + execute_profiled_run(run(), config, gateway=gateway, report_to_hub=False) + assert store.status()["reservations"][0]["state"] == "charged" + assert calls[0].project == "fixture-factory" + transfer.import_after_teardown.assert_called_once() + with pytest.raises(GlasExecutionError, match="replay"): + execute_profiled_run( + replace(run(), attempt=2), config, gateway=gateway, report_to_hub=False + ) + assert len(calls) == 1 + + +@pytest.mark.parametrize( + "change", + ["worker", "definition", "profile", "grant", "project", "limits", "descriptor"], +) +def test_scope_mismatch_denied_before_gateway_or_reservation(dispatch_case, change): + from rein_aharness.glas_execution import GlasExecutionError, execute_profiled_run + + store, config, catalog, _transfer = dispatch_case + item = run() + if change == "worker": + item.claim_owner = "other" + if change == "definition": + item.activity_definition_id = "other" + if change == "profile": + item.harness_profile_ref = "harness.agent-dev-local@1.0.0" + if change == "grant": + item.repository_grant = RepositoryGrant("1", ("elsewhere/",), 1, 1, False) + if change == "project": + config.execution_project = "other" + if change == "limits": + catalog.profiles()[("harness.test", "1.0.0")].limits.max_budget_usd = 5.0 + if change == "descriptor": + catalog.reins()["rein-aharness"].version = "99.0.0" + gateway = MagicMock() + with pytest.raises(GlasExecutionError): + execute_profiled_run(item, config, gateway=gateway, report_to_hub=False) + gateway.assert_not_called() + assert store.status()["reservations"] == [] + + +@pytest.mark.parametrize( + "failure", ["exception", "missing-cost", "overrun", "lease-loss"] +) +def test_uncertain_gateway_never_imports_and_blocks_next_claim(dispatch_case, failure): + from rein_aharness.execution_cancel import ExecutionCancel, ExecutionCancelled + from rein_aharness.glas_execution import GlasExecutionError, execute_profiled_run + + store, config, _catalog, transfer = dispatch_case + cancel = ExecutionCancel() + + def gateway(request, **kwargs): + if failure == "exception": + raise RuntimeError("private provider response") + if failure == "lease-loss": + cancel.cancel("lease-loss") + return success( + store, + run(), + None if failure == "missing-cost" else 7 if failure == "overrun" else 1, + ) + + with pytest.raises((GlasExecutionError, ExecutionCancelled)): + execute_profiled_run( + run(), config, gateway=gateway, cancel=cancel, report_to_hub=False + ) + transfer.import_after_teardown.assert_not_called() + assert store.status()["reservations"][0]["state"] == "held" + assert "private" not in json.dumps(store.status()) + with pytest.raises(SpendAdmissionError): + store.preflight() + + +def test_cli_reconciliation_requires_termination_attestation(ledger, capsys): + from rein_aharness.cli import main + + config = configured(ledger) + ledger.reserve(run(), now=NOW) + args = [ + "spend", + "reconcile", + "--policy", + config.spend_policy_path, + "--ledger", + config.spend_ledger_path, + "--run-id", + run().id, + "--cost-usd", + "1", + "--receipt", + "audit:final", + ] + assert main(args) == 2 + assert ledger.status()["reservations"][0]["state"] == "held" + assert main(args + ["--provider-stopped"]) == 0 + assert ledger.status()["reservations"][0]["state"] == "charged" + capsys.readouterr() + + +def test_required_admission_missing_never_claims(): + from rein_aharness.claim_loop import process_one + + client = MagicMock(spec=ActivityCoreOpsClient) + client.config = OpsRunConfig(require_spend_admission=True) + result = process_one(client) + assert not result.claimed and "not configured" in result.reason + client.claim.assert_not_called() + + +def test_policy_file_rejects_duplicates_and_symlinks(ledger): + config = configured(ledger) + path = Path(config.spend_policy_path) + text = path.read_text() + path.write_text(text[:-1] + ', "max_budget_usd": "4"}') + with pytest.raises(SpendAdmissionError, match="invalid"): + SpendPolicy.load(path) + path.write_text(text) + link = path.with_name("symlink.json") + link.symlink_to(path) + with pytest.raises(SpendAdmissionError, match="private"): + SpendPolicy.load(link) + + +def test_reconcile_unknown_identity_or_conflicting_receipt(ledger): + with pytest.raises(SpendAdmissionError, match="unknown"): + ledger.reconcile("absent", cost_usd="1", receipt="audit:final", now=NOW) + ledger.reserve(run(), now=NOW) + with pytest.raises(SpendAdmissionError, match="receipt"): + ledger.reconcile( + run().id, cost_usd="1", receipt="private text with spaces", now=NOW + ) + ledger.reconcile(run().id, cost_usd="1", receipt="audit:final", now=NOW) + with pytest.raises(SpendAdmissionError, match="conflict"): + ledger.reconcile(run().id, cost_usd="0", receipt="audit:other", now=NOW) + + +def test_decimal_conversion_never_rounds_liability_down(ledger): + from rein_aharness.spend_admission import cap_micros, converted_micros + + assert converted_micros("1.00000000000000000000000000000001", "1") == 1_000_001 + assert cap_micros("0.99999999999999999999999999999999") == 999_999 + assert ( + ledger.observe(run().id, success(ledger, run(), "1e-99999"), now=NOW) is False + ) diff --git a/workplans/REINAH-WP-0003-governed-runtime-integrity.md b/workplans/REINAH-WP-0003-governed-runtime-integrity.md index 7ae0bf0..af266d2 100644 --- a/workplans/REINAH-WP-0003-governed-runtime-integrity.md +++ b/workplans/REINAH-WP-0003-governed-runtime-integrity.md @@ -635,6 +635,25 @@ The matching rein/Glas code must be rebuilt and admitted in the protected runtim this local proof does not close live G1/G2 or authorize a model request. See the 2026-09-09 runtime-transfer evidence and the owning runtime documentation. + +### Durable spend admission — 2026-09-09 + +The existing worker now reserves declared maximum liability in a private durable +SQLite ledger before Glas dispatch. It pins the admitted owner/definition/target/ +project/grant and resolved profile/descriptor, counts conservative micro-euro +charges against daily and total ceilings, refuses duplicate demand/claim attempts, +and blocks after unknown accounting or crashes until explicit owner reconciliation. +Close-only replay runs before admission checks; observed overruns freeze the +envelope even after reconciliation. The trusted execution project is configurable; +no queue reference can override it. See [operator contract](../docs/spend-admission.md) +and [verification evidence](../docs/evidence/2026-09-09-spend-admission.json). + +This closes the missing local reservation/recovery implementation. T05 remains +`progress` for accepted provider liability/FX, owner configuration, protected-runtime +installation and Railiance recovery proof. T06 remains `wait` for the admitted +real-model and natural queue run. No factory policy, paid execution, deployment, +profile promotion or G0 grant was created by this source change. + ## Re-prove one governed profiled run and close residuals ```task