Bind the metered owner route to worker leases and sandbox lifecycle
Some checks failed
Governed runtime contract / contract (push) Failing after 23s
Some checks failed
Governed runtime contract / contract (push) Failing after 23s
Assistant: codex Assistant-Model: gpt-6-astra Assistant-Session: 01a07ff8-19d0-7820-b4d0-1353833cb7fc
This commit is contained in:
parent
c30806b968
commit
3e4c976090
13 changed files with 502 additions and 9 deletions
|
|
@ -296,6 +296,10 @@ def process_one(
|
|||
|
||||
try:
|
||||
spend = worker_spend(cfg)
|
||||
if cfg.require_request_admission and cfg.messages_owner is None:
|
||||
raise SpendAdmissionError("required request owner is not configured")
|
||||
if cfg.messages_owner is not None and spend is None:
|
||||
raise SpendAdmissionError("request owner requires parent spend admission")
|
||||
if spend is not None:
|
||||
spend.preflight()
|
||||
except SpendAdmissionError as exc:
|
||||
|
|
@ -320,7 +324,15 @@ def process_one(
|
|||
try:
|
||||
# Make active lease ownership observable even for normal short jobs,
|
||||
# and fail closed before dispatch if Activity Core no longer accepts it.
|
||||
client.heartbeat(run.id, lease_seconds=cfg.lease_seconds)
|
||||
accepted = client.heartbeat(run.id, lease_seconds=cfg.lease_seconds)
|
||||
if cfg.messages_owner is not None:
|
||||
from rein_aharness.messages_owner import accepted_lease
|
||||
try:
|
||||
run.lease_until = accepted_lease(run, accepted, cfg.worker_id)
|
||||
except SpendAdmissionError:
|
||||
return ProcessResult(claimed=True, run_id=run.id, ok=False,
|
||||
retry_full_interval=True,
|
||||
reason="accepted queue lease required for metered route")
|
||||
logger.info("accepted initial heartbeat run_id=%s", run.id)
|
||||
except OpsRunError as exc:
|
||||
rejected = exc.status_code == 409
|
||||
|
|
|
|||
|
|
@ -10,6 +10,7 @@ from __future__ import annotations
|
|||
|
||||
import math
|
||||
from collections.abc import Callable
|
||||
from contextlib import nullcontext
|
||||
from typing import Any
|
||||
|
||||
from rein_aharness.execution_cancel import ExecutionCancel, ExecutionCancelled, resolve_cancel
|
||||
|
|
@ -175,6 +176,11 @@ def execute_profiled_run(
|
|||
request = request_factory(**_request_kwargs(run, config, report_to_hub))
|
||||
kwargs = {"artifact_capture": transfer.capture} if transfer else {}
|
||||
spend = worker_spend(config)
|
||||
owner = config.messages_owner
|
||||
if config.require_request_admission and owner is None:
|
||||
raise SpendAdmissionError("required request owner is not configured")
|
||||
if owner is not None and (spend is None or guard is None):
|
||||
raise SpendAdmissionError("request owner requires spend and cancellation")
|
||||
if spend is not None:
|
||||
from glas_harness.profiles import ProfileCatalog
|
||||
catalog = ProfileCatalog()
|
||||
|
|
@ -186,7 +192,11 @@ def execute_profiled_run(
|
|||
if guard is not None:
|
||||
guard.check()
|
||||
spend.reserve(run)
|
||||
result = gateway(request, **kwargs)
|
||||
route = owner.activate(run, config, spend, profile, guard) if owner else nullcontext(None)
|
||||
with route as owner_manager:
|
||||
if owner_manager is not None:
|
||||
kwargs["manager"] = owner_manager
|
||||
result = gateway(request, **kwargs)
|
||||
raw = result.model_dump(mode="json") if hasattr(result, "model_dump") else result
|
||||
except ExecutionCancelled:
|
||||
raise
|
||||
|
|
|
|||
122
rein_aharness/messages_owner.py
Normal file
122
rein_aharness/messages_owner.py
Normal file
|
|
@ -0,0 +1,122 @@
|
|||
"""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
|
||||
|
||||
@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,),
|
||||
)
|
||||
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()
|
||||
|
|
@ -141,6 +141,8 @@ class OpsRunConfig:
|
|||
require_spend_admission: bool = False
|
||||
spend_policy_path: str | None = None
|
||||
spend_ledger_path: str | None = None
|
||||
require_request_admission: bool = False
|
||||
messages_owner: Any = field(default=None, repr=False)
|
||||
|
||||
@classmethod
|
||||
def from_env(cls) -> "OpsRunConfig":
|
||||
|
|
@ -181,6 +183,7 @@ class OpsRunConfig:
|
|||
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"),
|
||||
require_request_admission=bool(os.environ.get("AGENT_HARNESS_REQUIRE_REQUEST_ADMISSION", "")),
|
||||
)
|
||||
|
||||
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue