"""Explicit same-host request owner, activated only after worker admission. The bootstrap supplies the accepted Messages policy and provider credential. No environment discovery, custody acquisition or default production activation. """ from __future__ import annotations from contextlib import contextmanager from dataclasses import dataclass, field from datetime import UTC, datetime import tempfile import threading from pathlib import Path from llm_connect.messages_gate import MessagesPolicy, MessagesServer from rein_aharness.execution_cancel import ExecutionCancel from rein_aharness.request_admission import RequestLedger from rein_aharness.spend_admission import SpendAdmissionError, digest, timestamp def accepted_lease(run, accepted, worker_id): """Accept a queue response, never a caller-declared lease duration.""" if ( accepted.id != run.id or accepted.claim_owner != worker_id or accepted.attempt != run.attempt or type(accepted.attempt) is not int or accepted.attempt < 1 or accepted.state != "claimed" or not isinstance(accepted.lease_until, str) or timestamp(accepted.lease_until) <= datetime.now(UTC) ): raise SpendAdmissionError("accepted queue lease required for metered route") return accepted.lease_until @dataclass(frozen=True) class MessagesOwner: policy: MessagesPolicy provider_key: str = field(repr=False) upstream_url: str = "https://api.anthropic.com" allow_test_http: bool = False runtime_path: Path | None = None runtime_sha256: str | None = None @contextmanager def activate(self, run, config, spend, profile, cancel: ExecutionCancel): from sandboxer.core.manager import SandboxManager from sandboxer.extensions.messages_route import OwnerMessagesRoute if spend is None or cancel is None: raise SpendAdmissionError("metered owner requires spend and cancellation") # The claim loop overwrites this only with its accepted initial heartbeat. expiry = accepted_lease(run, run, config.worker_id) if profile.credential_route_refs or profile.model.model != self.policy.model: raise SpendAdmissionError("metered owner model or credential route mismatch") deadline = min(timestamp(expiry), timestamp(spend.policy.expires_at)) meter = RequestLedger(spend) cancel.check() token = meter.bind_route( run.id, self.policy.sha256, lease_id=digest({"run_id": run.id, "worker_id": config.worker_id, "attempt": run.attempt, "lease_until": expiry}), expires_at=deadline.isoformat(), ) stopped = threading.Event() def revoke(): # This latch denies local forwards even if durable revocation fails. stopped.set() meter.revoke_route(run.id) class ActiveMeter: def reserve_request(self, *args): if stopped.is_set() or cancel.cancelled: raise SpendAdmissionError("metered owner stopped") return meter.reserve_request(*args) def request_active(self, receipt): return not stopped.is_set() and not cancel.cancelled and meter.request_active(receipt) def complete_request(self, *args): return meter.complete_request(*args) cancel.register_stop(revoke) timer = threading.Timer( max(0, (deadline - datetime.now(UTC)).total_seconds()), lambda: cancel.cancel("timeout"), ) timer.daemon = True server = None try: # Keep the socket with the private ledger, outside source/runtime/workspace. with tempfile.TemporaryDirectory(prefix="messages-", dir=spend.path.parent) as directory: socket_path = Path(directory) / "route.sock" server = MessagesServer( self.policy, ActiveMeter(), provider_key=self.provider_key, upstream_url=self.upstream_url, allow_test_http=self.allow_test_http, unix_path=socket_path, ) route = OwnerMessagesRoute( socket_path, token, profile.sandbox_profile, "agt", config.execution_project, run.id, private_paths=(spend.path.parent,), runtime_path=self.runtime_path, runtime_sha256=self.runtime_sha256, ) server.start() timer.start() try: cancel.check() yield SandboxManager(messages_route=route) finally: # Before teardown of the listener, even on a refused/failed gateway. try: revoke() finally: server.stop() server = None finally: timer.cancel() try: revoke() finally: if server is not None: server.stop()