rein-aharness/rein_aharness/messages_owner.py
tegwick 565b07716a
Some checks failed
Governed runtime contract / contract (push) Failing after 23s
Bootstrap one metered owner cycle against a pinned standalone runtime
Assistant: codex
Assistant-Model: gpt-6-astra
Assistant-Session: 01a07ff8-19d0-7820-b4d0-1353833cb7fc
2026-09-09 22:56:38 +02:00

125 lines
5 KiB
Python

"""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()