mirror of
https://github.com/pewdiepie-archdaemon/odysseus.git
synced 2026-10-07 07:22:21 +02:00
Merge frozen lab b1666951 (Wave 3) into Wave 4 effects provenance
Integrates the merged and frozen Wave 3 lab commit b1666951faf8285054e1ca90f11533b0fb53fb57 with a normal merge, preserving every Wave 4 commit unchanged. Conflict: src/agent_runtime/resources.py. Wave 3's _control_plane_snapshot() / _control_plane_path(path, *, snapshot=None) split is kept. The snapshot adds the effect-store directories to its prefix set after the recursive job-dir inventory and no longer references path (the auto-merged prefix check would have raised NameError there). _control_plane_path calls _aliases_effect_store after its os.stat, only for multiply linked files, so single-link files never list the store. Semantic reconciliation (no textual conflict): bg_monitor keeps launch settlement right after the first successful validate_job and before the authority check, with Wave 3's post-drain revalidation intact. The deleted-session branch, terminal before linkage validation, now settles a validated launch too: that job is later pruned and its publication retired, which would otherwise leave its effect RUNNING. Regression tests cover the snapshot form of the effect-store check and both deleted-session linkage outcomes.
This commit is contained in:
@@ -523,6 +523,13 @@ def seal_task_authority(prompt, task_type, action, *, owner=None, parent_authori
|
||||
if operation is not None:
|
||||
authority = replace(authority, grants=(OperationGrant(operation.tool,
|
||||
inputs=frozenset({operation.input})),))
|
||||
if task_type == "action" and action == "cookbook_serve":
|
||||
# The direct admin scheduling ingress selects the native Cookbook
|
||||
# producer. Restore never infers this from task names/availability.
|
||||
# Any model-created task still intersects with its parent's ceiling.
|
||||
backend = NativeBackendResource("serve_model")
|
||||
authority = replace(authority, backend_resources=tuple(dict.fromkeys(
|
||||
(*authority.backend_resources, backend))))
|
||||
parent = active_request_authority() if parent_authority is MISSING_AUTHORITY else parent_authority
|
||||
if parent_authority is None:
|
||||
parent = RequestAuthority.empty(owner=owner)
|
||||
|
||||
@@ -0,0 +1,105 @@
|
||||
"""One-use transport capabilities for admitted local Cookbook producers.
|
||||
|
||||
The internal HTTP token authenticates transport only. A capability bridges one
|
||||
server-owned request/operation/backend to one exact resolved local launch body.
|
||||
It is never persisted, returned to the model, or usable for shell/job control.
|
||||
"""
|
||||
from contextlib import contextmanager
|
||||
from dataclasses import dataclass
|
||||
import hashlib
|
||||
import json
|
||||
import secrets
|
||||
import threading
|
||||
import time
|
||||
|
||||
from src.agent_runtime.resources import NativeBackendResource, ResourceIdentityError
|
||||
|
||||
CAPABILITY_HEADER = "X-Odysseus-Local-Model-Capability"
|
||||
_ROUTES = {"download_model": "/api/model/download", "serve_model": "/api/model/serve",
|
||||
"serve_preset": "/api/model/serve"}
|
||||
_PENDING = {}
|
||||
_LOCK = threading.Lock()
|
||||
|
||||
|
||||
def _digest(payload):
|
||||
return hashlib.sha256(json.dumps(payload, sort_keys=True, separators=(",", ":"),
|
||||
allow_nan=False).encode()).hexdigest()
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class _Capability:
|
||||
authority: object
|
||||
operation: object
|
||||
backend: NativeBackendResource
|
||||
path: str
|
||||
payload_digest: str
|
||||
deadline: float
|
||||
|
||||
|
||||
@contextmanager
|
||||
def model_control_headers(tool, content, owner, payload, *, scheduled=False):
|
||||
from src.tools._common import _internal_headers
|
||||
headers = _internal_headers(owner)
|
||||
if payload.get("remote_host"):
|
||||
yield headers # Remote workload authority/transport is unchanged.
|
||||
return
|
||||
from src.agent_runtime.authority import active_request_authority, ExactOperation
|
||||
from src.agent_runtime.remote_resources import active_backend_operation
|
||||
from src.tool_security import owner_is_admin_or_single_user
|
||||
authority = active_request_authority()
|
||||
operation = ExactOperation.normalize(tool, content)
|
||||
backend = active_backend_operation()
|
||||
if (authority is None or authority.owner != str(owner or "").strip().casefold()
|
||||
or tool not in _ROUTES or not owner_is_admin_or_single_user(owner)):
|
||||
raise ResourceIdentityError("Local model producer has no matching server authority")
|
||||
if scheduled:
|
||||
# Called only by the server-owned scheduled action, after restoration of
|
||||
# its immutable input ceiling. A task name or owner alone is not enough.
|
||||
if tool != "serve_model" or not authority.permits(operation):
|
||||
raise ResourceIdentityError("Scheduled local model input is outside authority")
|
||||
resource = NativeBackendResource(tool)
|
||||
if resource not in authority.backend_resources:
|
||||
raise ResourceIdentityError("Scheduled local model backend is outside authority")
|
||||
else:
|
||||
# This binding exists only after dispatch admission (including one-use
|
||||
# exact approval). A generic tool grant/header cannot create it over HTTP.
|
||||
if (backend is None or backend.resource != NativeBackendResource(tool)
|
||||
or (backend.request_id, backend.owner, backend.session_id,
|
||||
backend.transport_tool, backend.exact_input) !=
|
||||
(authority.request_id, authority.owner, authority.session_id, tool, operation.input)):
|
||||
raise ResourceIdentityError("Local model producer operation or backend changed")
|
||||
resource = backend.resource
|
||||
capability = _Capability(authority, operation, resource, _ROUTES[tool], _digest(payload), time.monotonic() + 60)
|
||||
token = secrets.token_urlsafe(32)
|
||||
headers.update({CAPABILITY_HEADER: token, "X-Odysseus-Owner": authority.owner})
|
||||
with _LOCK:
|
||||
_PENDING[token] = capability
|
||||
try:
|
||||
yield headers
|
||||
finally:
|
||||
with _LOCK:
|
||||
_PENDING.pop(token, None)
|
||||
|
||||
|
||||
def consume_model_control(request, payload):
|
||||
"""Claim exactly once at the local route, before any producer effect."""
|
||||
from core.middleware import INTERNAL_TOOL_HEADER, INTERNAL_TOOL_TOKEN, INTERNAL_TOOL_USER
|
||||
from src.auth_helpers import is_direct_loopback_request
|
||||
token = request.headers.get(CAPABILITY_HEADER)
|
||||
if not token:
|
||||
return False
|
||||
if (not is_direct_loopback_request(request)
|
||||
or not secrets.compare_digest(request.headers.get(INTERNAL_TOOL_HEADER, ""), INTERNAL_TOOL_TOKEN)):
|
||||
raise ResourceIdentityError("Local model transport is untrusted")
|
||||
with _LOCK:
|
||||
capability = _PENDING.get(token)
|
||||
if (capability is None or capability.deadline < time.monotonic()
|
||||
or request.method != "POST" or request.url.path != capability.path
|
||||
or payload.get("remote_host") or _digest(payload) != capability.payload_digest
|
||||
or request.headers.get("X-Odysseus-Owner", "") != capability.authority.owner
|
||||
or getattr(request.state, "current_user", None) not in
|
||||
(None, INTERNAL_TOOL_USER, capability.authority.owner)):
|
||||
raise ResourceIdentityError("Local model capability binding changed or expired")
|
||||
del _PENDING[token]
|
||||
request.state.local_model_authority = capability.authority
|
||||
return True
|
||||
@@ -260,6 +260,8 @@ def needs_owned_binding(operation):
|
||||
"notes", "memory", "vault", "upload", "uploads", "attachments",
|
||||
"shell", "model", "cookbook"}
|
||||
segments = path.strip("/").split("/")
|
||||
if len(segments) >= 3 and segments[:3] == ["api", "codex", "cookbook"]:
|
||||
raise ResourceIdentityError("Cookbook wrappers require a dedicated resource-bound tool")
|
||||
if len(segments) >= 2 and segments[0] == "api" and segments[1].casefold() in private:
|
||||
raise ResourceIdentityError("Owned records require a dedicated resource-bound tool")
|
||||
return False
|
||||
|
||||
@@ -8,12 +8,13 @@ from __future__ import annotations
|
||||
|
||||
from contextlib import contextmanager
|
||||
from contextvars import ContextVar
|
||||
from dataclasses import dataclass
|
||||
from dataclasses import dataclass, field
|
||||
import hashlib
|
||||
import json
|
||||
import os
|
||||
from pathlib import Path
|
||||
import re
|
||||
import threading
|
||||
from uuid import uuid4
|
||||
from core.atomic_io import store_transaction
|
||||
|
||||
@@ -220,6 +221,19 @@ def intersect_launch_scopes(parent, child):
|
||||
return tuple(dict.fromkeys(narrowed))
|
||||
|
||||
|
||||
class _LaunchUse:
|
||||
"""Non-persisted one-use producer reservation, shared by approval copies."""
|
||||
def __init__(self):
|
||||
self.used = False
|
||||
self.lock = threading.Lock()
|
||||
|
||||
def claim(self):
|
||||
with self.lock:
|
||||
if self.used:
|
||||
raise ResourceIdentityError("Launch reservation has already been used")
|
||||
self.used = True
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class BoundProcessOperation:
|
||||
operation: object
|
||||
@@ -230,6 +244,7 @@ class BoundProcessOperation:
|
||||
jobs: tuple[BackgroundJobResource, ...] = ()
|
||||
processes: tuple[ProcessResource, ...] = ()
|
||||
exact_approval: object | None = None
|
||||
_launch_use: _LaunchUse = field(default_factory=_LaunchUse, compare=False, repr=False)
|
||||
|
||||
def __post_init__(self):
|
||||
from src.agent_runtime.authority import ExactOperation
|
||||
@@ -249,7 +264,8 @@ class BoundProcessOperation:
|
||||
def validate(self):
|
||||
if self.launch is not None:
|
||||
self.launch.validate()
|
||||
guard_launch_workspace(self.launch.scope.root)
|
||||
if self._launch_use.used:
|
||||
raise ResourceIdentityError("Launch reservation has already been used")
|
||||
for job in self.jobs:
|
||||
validate_job(job, mutation=self.operation.action in {"kill", "stop", "cancel", "terminate", "ack"})
|
||||
for process in self.processes:
|
||||
@@ -321,6 +337,11 @@ def bind_process_operation(operation):
|
||||
raise TypeError("Process operation must be server-owned")
|
||||
if operation is not None:
|
||||
operation.validate()
|
||||
if operation.launch is not None:
|
||||
# One fresh authoritative scan for each execution binding. Resolution
|
||||
# and producer entry retain cheap exact identity checks; no scan is
|
||||
# reused across independent bindings or persisted in an approval.
|
||||
guard_launch_workspace(operation.launch.scope.root)
|
||||
token = _ACTIVE.set(operation)
|
||||
try:
|
||||
yield operation
|
||||
@@ -363,7 +384,7 @@ def guard_launch_workspace(root):
|
||||
"""
|
||||
from src import bg_jobs, containment, constants
|
||||
from src import browser_identity
|
||||
from src.agent_runtime.resources import _control_plane_path
|
||||
from src.agent_runtime.resources import _control_plane_path, _control_plane_snapshot
|
||||
control = (Path(bg_jobs._STORE), Path(bg_jobs._JOBS_DIR), containment._store_path(), _LAUNCH_DIR,
|
||||
Path(constants.BROWSER_RESOURCES_DIR),
|
||||
browser_identity.STATE_ROOT,
|
||||
@@ -373,12 +394,16 @@ def guard_launch_workspace(root):
|
||||
raise ResourceIdentityError("Launch boundary contains server control state")
|
||||
def unresolved(error):
|
||||
raise ResourceIdentityError("Launch workspace cannot be inspected") from error
|
||||
snapshot = None
|
||||
for directory, dirs, files in os.walk(base, followlinks=False, onerror=unresolved):
|
||||
for name in (*dirs, *files):
|
||||
path = Path(directory) / name
|
||||
info = path.lstat()
|
||||
if (path.is_symlink() or info.st_nlink > 1) and _control_plane_path(str(path.resolve())):
|
||||
raise ResourceIdentityError("Launch boundary aliases server control state")
|
||||
if path.is_symlink() or info.st_nlink > 1:
|
||||
if snapshot is None:
|
||||
snapshot = _control_plane_snapshot()
|
||||
if _control_plane_path(str(path.resolve()), snapshot=snapshot):
|
||||
raise ResourceIdentityError("Launch boundary aliases server control state")
|
||||
|
||||
|
||||
@store_transaction(lambda: _LAUNCH_DIR / "publication")
|
||||
@@ -390,11 +415,83 @@ def publish_launch(launch, authority, containment_id, *, job=None, processes=())
|
||||
path = launch_path(launch.generation)
|
||||
if path.exists():
|
||||
raise ResourceIdentityError("Launch reservation has already been used")
|
||||
bound = active_process_operation()
|
||||
if bound is not None:
|
||||
if bound.launch != launch:
|
||||
raise ResourceIdentityError("Publication differs from the bound launch")
|
||||
bound._launch_use.claim()
|
||||
atomic_write_json(path, {"launch": launch.to_dict(), "authority": authority.to_dict(),
|
||||
"containment_id": containment_id, "job": job.to_dict() if job else None,
|
||||
"processes": [p.to_dict() for p in processes]})
|
||||
|
||||
|
||||
@store_transaction(lambda: _LAUNCH_DIR / "publication")
|
||||
def retire_launch(launch, containment_id, *, job=None):
|
||||
"""Remove only this exact producer publication; never a replacement.
|
||||
|
||||
Callers establish the lifetime end (verified foreground teardown, or exact
|
||||
background history pruning). Missing/malformed/replaced state is retained.
|
||||
One-use launch reservations live in the bound operation, not this file.
|
||||
"""
|
||||
path = launch_path(launch.generation)
|
||||
try:
|
||||
published = json.loads(path.read_text())
|
||||
except FileNotFoundError:
|
||||
return False
|
||||
if (not isinstance(published, dict)
|
||||
or published.get("launch") != launch.to_dict()
|
||||
or published.get("containment_id") != containment_id
|
||||
or published.get("job") != (job.to_dict() if job else None)):
|
||||
return False
|
||||
path.unlink()
|
||||
return True
|
||||
|
||||
|
||||
@store_transaction(lambda: _LAUNCH_DIR / "publication")
|
||||
def prune_foreground_publications():
|
||||
"""Startup-only recovery: retire foreground generations without a caller.
|
||||
|
||||
A dead/replaced manager cannot resume attachment. A missing receipt also
|
||||
makes attachment impossible; publication cannot reconstruct that receipt.
|
||||
Its process tree still belongs to containment recovery; deleting a
|
||||
publication never signals or asserts tree death. Live/unverifiable managers
|
||||
retain publication even after child teardown: attachment may still need it.
|
||||
Background history stays intact.
|
||||
"""
|
||||
from src import containment
|
||||
from src import process_ownership
|
||||
try:
|
||||
receipts = json.loads(containment._store_path().read_text())
|
||||
except FileNotFoundError:
|
||||
receipts = {}
|
||||
except (OSError, ValueError):
|
||||
return 0 # Unreadable state is not evidence that consumers are gone.
|
||||
if not isinstance(receipts, dict) or any(not isinstance(r, dict) for r in receipts.values()):
|
||||
return 0
|
||||
retired = 0
|
||||
for path in _LAUNCH_DIR.glob("*.json"):
|
||||
try:
|
||||
published = json.loads(path.read_text())
|
||||
launch = ProcessLaunchResource.from_dict(published["launch"])
|
||||
receipt = receipts.get(published["containment_id"])
|
||||
abandoned = (receipt is not None
|
||||
and type(receipt.get("manager_pid")) is int and receipt["manager_pid"] > 0
|
||||
and isinstance(receipt.get("manager_token"), str) and bool(receipt["manager_token"])
|
||||
and process_ownership.verify(receipt["manager_pid"], receipt["manager_token"]) in {
|
||||
process_ownership.GONE, process_ownership.FOREIGN})
|
||||
if (published.get("job") is None and path == launch_path(launch.generation)
|
||||
and (receipt is None or (
|
||||
receipt.get("launch_generation") == launch.generation
|
||||
and receipt.get("id") == published["containment_id"]
|
||||
and abandoned))):
|
||||
# Already under the publication lock; no nested file lock.
|
||||
path.unlink()
|
||||
retired += 1
|
||||
except (ValueError, TypeError, KeyError, OSError):
|
||||
continue
|
||||
return retired
|
||||
|
||||
|
||||
@store_transaction(lambda: _LAUNCH_DIR / "publication")
|
||||
def attach_containment_processes(launch, containment_id):
|
||||
"""Attach producer-frozen lifecycle records; never capture a current PID."""
|
||||
@@ -410,10 +507,13 @@ def attach_containment_processes(launch, containment_id):
|
||||
processes = []
|
||||
for role, pid_key, token_key, group_key in (("leader", "pid", "start_token", "pgid"),
|
||||
("namespace_init", "namespace_pid", "namespace_start_token", None)):
|
||||
if record.get(pid_key):
|
||||
processes.append(ProcessResource("native:containment", launch.owner, launch.request_id,
|
||||
launch.thread_id, ProcessIdentity(record[pid_key], record.get(token_key), record.get(group_key) if group_key else None),
|
||||
role, "", containment_id))
|
||||
pid = record.get(pid_key)
|
||||
token = record.get(token_key)
|
||||
if not pid or not token:
|
||||
continue
|
||||
processes.append(ProcessResource("native:containment", launch.owner, launch.request_id,
|
||||
launch.thread_id, ProcessIdentity(pid, token, record.get(group_key) if group_key else None),
|
||||
role, "", containment_id))
|
||||
from core.atomic_io import atomic_write_json
|
||||
published["processes"] = [p.to_dict() for p in processes]
|
||||
atomic_write_json(path, published)
|
||||
|
||||
@@ -67,7 +67,7 @@ def _aliases_effect_store(candidate, directories):
|
||||
return False
|
||||
|
||||
|
||||
def _control_plane_path(path):
|
||||
def _control_plane_snapshot():
|
||||
# Execution snapshots/receipts are server state, even if a workspace root
|
||||
# contains the data directory. A writable user file cannot mint authority.
|
||||
from src import constants
|
||||
@@ -85,11 +85,9 @@ def _control_plane_path(path):
|
||||
if processes is not None:
|
||||
job_dirs.add(canonical_root(processes._LAUNCH_DIR))
|
||||
# Durable effect claims/outcomes/observations are server evidence state.
|
||||
# The store is not inventoried here: it grows with every run. Aliases are
|
||||
# caught below by ``_aliases_effect_store`` instead.
|
||||
# They are prefix-protected below, but never inventoried: the store grows
|
||||
# with every run. Hardlink aliases are caught by ``_aliases_effect_store``.
|
||||
effect_dirs = _effect_store_dirs()
|
||||
if any(Path(path).is_relative_to(directory) for directory in effect_dirs):
|
||||
return True
|
||||
# Producers may have configured paths different from the default constants.
|
||||
# Inspect already-loaded server metadata without initializing a store here.
|
||||
bg = sys.modules.get("src.bg_jobs")
|
||||
@@ -118,8 +116,6 @@ def _control_plane_path(path):
|
||||
protected.add(canonical_root(Path(uploader.upload_dir) / "uploads.json"))
|
||||
for directory in job_dirs:
|
||||
jobs = Path(directory)
|
||||
if Path(path).is_relative_to(jobs):
|
||||
return True
|
||||
if jobs.exists():
|
||||
# Uninspectable state fails closed; hardlinks retain object identity.
|
||||
protected.update(canonical_root(p) for p in jobs.rglob("*") if p.is_file())
|
||||
@@ -128,22 +124,32 @@ def _control_plane_path(path):
|
||||
for suffix in ("-wal", "-shm", "-journal"))
|
||||
protected.add(canonical_root(Path(constants.DATA_DIR) / ".app_key"))
|
||||
protected.add(canonical_root(Path(constants.UPLOAD_DIR) / "uploads.json"))
|
||||
if path in protected:
|
||||
return True
|
||||
try:
|
||||
candidate = os.stat(path)
|
||||
except FileNotFoundError:
|
||||
return False
|
||||
if _aliases_effect_store(candidate, effect_dirs):
|
||||
return True
|
||||
identities = set()
|
||||
for control in protected:
|
||||
try:
|
||||
observed = os.stat(control)
|
||||
except FileNotFoundError:
|
||||
continue
|
||||
if (candidate.st_dev, candidate.st_ino) == (observed.st_dev, observed.st_ino):
|
||||
return True
|
||||
return False
|
||||
identities.add((observed.st_dev, observed.st_ino))
|
||||
# Effect directories join the prefix set only after the recursive inventory.
|
||||
return frozenset(job_dirs | effect_dirs), frozenset(protected), frozenset(identities)
|
||||
|
||||
|
||||
def _control_plane_path(path, *, snapshot=None):
|
||||
# A scan-local snapshot bounds repeated hardlink checks. Ordinary resource
|
||||
# resolution always observes fresh state. Neither form is an atomic kernel
|
||||
# access policy, and snapshots must never survive a workspace guard call.
|
||||
directories, protected, identities = _control_plane_snapshot() if snapshot is None else snapshot
|
||||
if any(Path(path).is_relative_to(directory) for directory in directories) or path in protected:
|
||||
return True
|
||||
try:
|
||||
candidate = os.stat(path)
|
||||
except FileNotFoundError:
|
||||
return False
|
||||
if (candidate.st_dev, candidate.st_ino) in identities:
|
||||
return True
|
||||
# Only a multiply linked file can alias the (uninventoried) effect store.
|
||||
return candidate.st_nlink > 1 and _aliases_effect_store(candidate, directories & _effect_store_dirs())
|
||||
|
||||
|
||||
class FilesystemScope(str, Enum):
|
||||
|
||||
@@ -513,6 +513,7 @@ async def _run_owned_command(command, ctx: dict, *, tool: str, timeout: int, arg
|
||||
from src.tool_execution import agent_cwd, _truncate
|
||||
|
||||
grant = None
|
||||
launch = None
|
||||
result = None
|
||||
try:
|
||||
from src.agent_runtime.process_resources import require_launch, publish_launch, validate_launch_spec
|
||||
@@ -554,6 +555,16 @@ async def _run_owned_command(command, ctx: dict, *, tool: str, timeout: int, arg
|
||||
**({"failure_kind": "resource_linkage_unavailable",
|
||||
"teardown": result.release.to_dict() if result.release else {"dead": False}}
|
||||
if result is not None else {})}
|
||||
finally:
|
||||
if launch is not None and grant is not None:
|
||||
record = containment._load_records().get(grant.id, {})
|
||||
if (record.get("launch_generation") == launch.generation
|
||||
and (record.get("release") or {}).get("dead") is True):
|
||||
from src.agent_runtime.process_resources import retire_launch
|
||||
try:
|
||||
retire_launch(launch, grant.id)
|
||||
except (OSError, ValueError, TypeError):
|
||||
logger.warning("Foreground launch publication retirement failed", exc_info=True)
|
||||
|
||||
boundary = result.grant.to_dict()
|
||||
boundary["executed"] = True
|
||||
|
||||
@@ -7,6 +7,27 @@ from fastapi import Request, HTTPException
|
||||
from src.owner_identity import auth_disabled, effective_storage_owner
|
||||
|
||||
|
||||
def is_direct_loopback_request(request: Request) -> bool:
|
||||
"""Local operator transport, excluding reverse proxies and cross-site calls.
|
||||
|
||||
Locality supplies no model/tool authority. Native administration uses this
|
||||
only in the operator's explicit auth-disabled single-user mode.
|
||||
"""
|
||||
client = getattr(request, "client", None)
|
||||
if not client or client.host not in {"127.0.0.1", "::1"}:
|
||||
return False
|
||||
forwarding = ("cf-connecting-ip", "cf-ray", "cf-visitor", "x-forwarded-for",
|
||||
"x-forwarded-host", "x-forwarded-proto", "x-real-ip", "forwarded")
|
||||
if any(request.headers.get(name) for name in forwarding):
|
||||
return False
|
||||
if request.headers.get("sec-fetch-site") in {"cross-site", "same-site"}:
|
||||
return False
|
||||
origin = request.headers.get("origin")
|
||||
if origin and origin != str(request.base_url).rstrip("/"):
|
||||
return False
|
||||
return True
|
||||
|
||||
|
||||
def get_current_user(request: Request) -> Optional[str]:
|
||||
"""Get current username from request state (set by auth middleware)."""
|
||||
return getattr(request.state, 'current_user', None)
|
||||
|
||||
+35
-8
@@ -213,10 +213,22 @@ def _prune(jobs: Dict[str, Dict[str, Any]], now: float) -> bool:
|
||||
"""Drop records (and their on-disk files) for jobs that finished, were
|
||||
followed up, and are older than the retention window. Mutates `jobs`."""
|
||||
stale = [jid for jid, rec in jobs.items()
|
||||
if rec.get("followed_up") and rec.get("ended_at")
|
||||
if rec.get("status") in {"done", "failed"}
|
||||
and (rec.get("followed_up") or rec.get("followup_state") == "terminal_unfollowable")
|
||||
and rec.get("ended_at")
|
||||
and (rec.get("teardown") or {}).get("dead") is not False
|
||||
and (now - rec["ended_at"]) > _RETENTION_S]
|
||||
for jid in stale:
|
||||
jobs.pop(jid, None)
|
||||
rec = jobs.pop(jid)
|
||||
from src.agent_runtime.process_resources import job_from_record, retire_launch
|
||||
from src.agent_runtime.resources import ProcessLaunchResource
|
||||
try:
|
||||
resource = job_from_record(rec)
|
||||
retire_launch(ProcessLaunchResource.from_dict(rec["launch_resource"]),
|
||||
resource.containment_id, job=resource)
|
||||
except (ValueError, TypeError, OSError):
|
||||
# Malformed/replaced publications never become deletion authority.
|
||||
pass
|
||||
for p in _JOBS_DIR.glob(f"{jid}.*"): # .sh .cmd.sh .log .exit
|
||||
try:
|
||||
p.unlink()
|
||||
@@ -333,10 +345,29 @@ def _kill_record(rec):
|
||||
|
||||
def pending_followups() -> List[Dict[str, Any]]:
|
||||
"""Finished jobs the agent hasn't been re-invoked for yet. The monitor
|
||||
drains these; mark_followed_up() flips the flag only on success."""
|
||||
drains these; valid continuations acknowledge success, invalid immutable
|
||||
linkage receives a terminal disposition without fabricating delivery."""
|
||||
jobs = refresh()
|
||||
return [r for r in jobs.values()
|
||||
if r.get("status") in ("done", "failed") and not r.get("followed_up")]
|
||||
if r.get("status") in ("done", "failed") and not r.get("followed_up")
|
||||
and r.get("followup_state") != "terminal_unfollowable"]
|
||||
|
||||
|
||||
@store_transaction(lambda: _STORE)
|
||||
def mark_unfollowable(job_id: str, *, expected_record) -> bool:
|
||||
"""Suppress only the exact completed snapshot inspected by the monitor.
|
||||
|
||||
This conveys no read/signal/continuation authority and cannot renew a PID.
|
||||
It deliberately needs no invalid/missing authority sidecar to suppress it.
|
||||
"""
|
||||
jobs = _load()
|
||||
record = jobs.get(job_id)
|
||||
if (record is None or record != expected_record or record.get("id") != job_id
|
||||
or record.get("status") not in {"done", "failed"}):
|
||||
return False
|
||||
record["followup_state"] = "terminal_unfollowable"
|
||||
_save(jobs)
|
||||
return True
|
||||
|
||||
|
||||
@store_transaction(lambda: _STORE)
|
||||
@@ -373,10 +404,6 @@ def get(job_id: str, *, expected) -> Optional[Dict[str, Any]]:
|
||||
return rec
|
||||
|
||||
|
||||
def list_for_session(session_id: str) -> List[Dict[str, Any]]:
|
||||
return [r for r in _load().values() if r.get("session_id") == session_id]
|
||||
|
||||
|
||||
@store_transaction(lambda: _STORE)
|
||||
def kill(job_id: str, *, expected) -> Optional[Dict[str, Any]]:
|
||||
"""Terminate a running job's process tree and mark it killed. Returns the
|
||||
|
||||
+46
-16
@@ -13,6 +13,7 @@ from __future__ import annotations
|
||||
import asyncio
|
||||
import json
|
||||
import logging
|
||||
from enum import Enum, auto
|
||||
|
||||
from src import bg_jobs
|
||||
from src.prompt_security import untrusted_context_message
|
||||
@@ -26,6 +27,12 @@ POLL_INTERVAL_S = 5
|
||||
_FOLLOWUP_MAX_ROUNDS = 12
|
||||
|
||||
|
||||
class FollowupResult(Enum):
|
||||
RETRYABLE_LATER = auto()
|
||||
COMPLETED = auto()
|
||||
TERMINAL_UNFOLLOWABLE = auto()
|
||||
|
||||
|
||||
def _background_result_message(rec):
|
||||
inject = (
|
||||
f"[Background job {rec['id']} finished]\n\n"
|
||||
@@ -120,22 +127,30 @@ async def _drain_agent(sess, messages, request_authority=None):
|
||||
return full, tool_events
|
||||
|
||||
|
||||
async def _run_followup(rec: dict) -> bool:
|
||||
"""Re-invoke the agent in the job's session with the result. Returns True
|
||||
if the follow-up completed (or there's nothing to do) — i.e. it's safe to
|
||||
mark followed_up. Returns False to retry on the next tick."""
|
||||
async def _run_followup(rec: dict) -> FollowupResult:
|
||||
"""Continue only an exactly linked result; distinguish retry from terminal."""
|
||||
from src.ai_interaction import get_session_manager
|
||||
from core.models import ChatMessage
|
||||
|
||||
sm = get_session_manager()
|
||||
if not sm:
|
||||
return False # not ready yet — retry
|
||||
return FollowupResult.RETRYABLE_LATER
|
||||
sess = sm.get_session(rec["session_id"])
|
||||
if not sess:
|
||||
# Session was deleted — nothing to continue. Consider it handled so we
|
||||
# don't retry forever.
|
||||
logger.info("bg-followup: session %s gone for job %s — skipping", rec.get("session_id"), rec.get("id"))
|
||||
return True
|
||||
# The job is retired without a continuation, then pruned with its
|
||||
# publication. Settle its launch effect first so it is not left RUNNING.
|
||||
from src.agent_runtime.process_resources import job_from_record, validate_job
|
||||
try:
|
||||
resource = job_from_record(rec)
|
||||
validate_job(resource)
|
||||
except (ValueError, TypeError, OSError, RuntimeError):
|
||||
pass # no validated linkage: nothing may be settled
|
||||
else:
|
||||
_settle_launch_effect(resource, rec)
|
||||
return FollowupResult.TERMINAL_UNFOLLOWABLE
|
||||
|
||||
# Don't write into a session that's mid-stream. The followup appends to
|
||||
# history + save_sessions(); a concurrent live turn does the same, and with
|
||||
@@ -145,13 +160,10 @@ async def _run_followup(rec: dict) -> bool:
|
||||
from src import agent_runs
|
||||
if agent_runs.is_active(sess.id):
|
||||
logger.info("bg-followup: session %s busy (live turn) — deferring job %s", sess.id, rec.get("id"))
|
||||
return False
|
||||
return FollowupResult.RETRYABLE_LATER
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
context = sess.get_context_messages()
|
||||
context.append(_background_result_message(rec))
|
||||
|
||||
from src.agent_runtime.authority import restore_background_authority
|
||||
from src.settings import get_setting
|
||||
authority = restore_background_authority(
|
||||
@@ -165,11 +177,19 @@ async def _run_followup(rec: dict) -> bool:
|
||||
_settle_launch_effect(resource, rec)
|
||||
if not authority.grants or (resource.owner, resource.thread_id, resource.request_id) != (
|
||||
str(getattr(sess, "owner", None) or "").strip().casefold(), sess.id, authority.request_id):
|
||||
return False
|
||||
return FollowupResult.TERMINAL_UNFOLLOWABLE
|
||||
except (ValueError, TypeError, OSError, RuntimeError):
|
||||
return False
|
||||
return FollowupResult.TERMINAL_UNFOLLOWABLE
|
||||
context = sess.get_context_messages()
|
||||
context.append(_background_result_message(rec))
|
||||
authority = authority.restrict(disabled_tools=get_setting("disabled_tools", []) or ())
|
||||
full, tool_events = await _drain_agent(sess, context, request_authority=authority)
|
||||
# An awaited continuation must not deliver a result after its immutable
|
||||
# linkage disappears or is replaced. This check grants no new authority.
|
||||
try:
|
||||
validate_job(resource)
|
||||
except (ValueError, TypeError, OSError, RuntimeError):
|
||||
return FollowupResult.TERMINAL_UNFOLLOWABLE
|
||||
|
||||
# Persist ONLY the assistant continuation so it renders as a normal agent
|
||||
# turn — a standard chat bubble plus `tool_events` that the frontend
|
||||
@@ -188,7 +208,19 @@ async def _run_followup(rec: dict) -> bool:
|
||||
sm.save_sessions()
|
||||
logger.info("bg-followup: auto-continued session %s for job %s (%d chars, %d tools)",
|
||||
sess.id, rec["id"], len(full), len(tool_events))
|
||||
return True
|
||||
return FollowupResult.COMPLETED
|
||||
|
||||
|
||||
async def _process_followup(rec):
|
||||
outcome = await _run_followup(rec)
|
||||
if outcome is FollowupResult.COMPLETED:
|
||||
from src.agent_runtime.process_resources import job_from_record
|
||||
bg_jobs.mark_followed_up(rec["id"], expected=job_from_record(rec))
|
||||
elif outcome is FollowupResult.TERMINAL_UNFOLLOWABLE:
|
||||
if not bg_jobs.mark_unfollowable(rec["id"], expected_record=rec):
|
||||
return FollowupResult.RETRYABLE_LATER
|
||||
logger.warning("bg-followup: job %s has no valid continuation linkage; retired from pending", rec.get("id"))
|
||||
return outcome
|
||||
|
||||
|
||||
async def _loop():
|
||||
@@ -196,9 +228,7 @@ async def _loop():
|
||||
try:
|
||||
for rec in bg_jobs.pending_followups():
|
||||
try:
|
||||
if await _run_followup(rec):
|
||||
from src.agent_runtime.process_resources import job_from_record
|
||||
bg_jobs.mark_followed_up(rec["id"], expected=job_from_record(rec))
|
||||
await _process_followup(rec)
|
||||
except Exception as e:
|
||||
# Idempotent: leave followed_up=False so the next tick retries.
|
||||
logger.warning("bg-followup failed for %s (will retry): %s", rec.get("id"), e)
|
||||
|
||||
@@ -30,6 +30,8 @@ from src.process_lifecycle import ProcessIdentity, observe
|
||||
from src.constants import BROWSER_RESOURCES_DIR
|
||||
|
||||
PRODUCER_VERSION = "0.35.0"
|
||||
# Wave 3 session metadata supports only these observed glibc Linux artifacts.
|
||||
# macOS/Windows and other architectures fail closed before any producer call.
|
||||
PRODUCER_HASHES = {
|
||||
"linux-x64": "b7a28c3a43a7008dd02585e2e60c391c08983f7a099149caed63c9f13f57b752",
|
||||
"linux-arm64": "92cd7d0897837ac648b9a6ab1965c69c5920e0f54df57e4295cdb1143b0541c8",
|
||||
|
||||
@@ -3387,10 +3387,12 @@ async def action_cookbook_serve(
|
||||
if srv.get("platform"): body["platform"] = srv["platform"]
|
||||
|
||||
try:
|
||||
async with httpx.AsyncClient(timeout=30) as client:
|
||||
r = await client.post(f"{internal_api_base()}/api/model/serve",
|
||||
json=body, headers=headers)
|
||||
data = r.json() if r.content else {}
|
||||
from src.agent_runtime.local_model_control import model_control_headers
|
||||
with model_control_headers("serve_model", command, owner, body, scheduled=True) as launch_headers:
|
||||
async with httpx.AsyncClient(timeout=30) as client:
|
||||
r = await client.post(f"{internal_api_base()}/api/model/serve",
|
||||
json=body, headers=launch_headers)
|
||||
data = r.json() if r.content else {}
|
||||
except Exception as e:
|
||||
return f"Launch HTTP failed: {e}", False
|
||||
if not data.get("ok"):
|
||||
|
||||
@@ -284,11 +284,21 @@ def reap_orphans() -> Dict[str, Any]:
|
||||
Blocking: a teardown escalates SIGTERM → grace → SIGKILL and waits for the
|
||||
process to actually go. Call it off the event loop.
|
||||
"""
|
||||
# Observe publication consumers before receipt recovery can forget a dead
|
||||
# manager's record. Publication retirement itself neither signals nor
|
||||
# asserts successful teardown; containment remains the recovery authority.
|
||||
from src.agent_runtime.process_resources import prune_foreground_publications
|
||||
try:
|
||||
publications_retired = prune_foreground_publications()
|
||||
except (OSError, ValueError, TypeError):
|
||||
publications_retired = 0
|
||||
logger.warning("process_reaper: foreground publication retirement failed", exc_info=True)
|
||||
report = {
|
||||
"mechanism": process_ownership.inspection_mechanism(),
|
||||
"grants": reap_containment_grants(),
|
||||
"bg_jobs": reap_bg_jobs(),
|
||||
"agent_tmux": reap_legacy_agent_tmux(),
|
||||
"foreground_publications_retired": publications_retired,
|
||||
}
|
||||
if report["mechanism"] == process_ownership.MECHANISM_NONE:
|
||||
logger.error(
|
||||
|
||||
+3
-13
@@ -1246,8 +1246,6 @@ def _split_bg_marker(content: str):
|
||||
return False, content
|
||||
|
||||
|
||||
import re as _re
|
||||
|
||||
# Variables a legitimate agent bash/python subprocess needs from the host.
|
||||
# Anything not listed here is never inherited.
|
||||
_SAFE_SUBPROCESS_VARS = frozenset({
|
||||
@@ -1267,19 +1265,11 @@ _SAFE_SUBPROCESS_VARS = frozenset({
|
||||
"LD_LIBRARY_PATH",
|
||||
})
|
||||
|
||||
# Defence-in-depth: reject any allowlisted variable whose *name* matches
|
||||
# a credential-bearing pattern (e.g. a user who sets PATH_TOKEN=...).
|
||||
_SENSITIVE_PATTERN = _re.compile(
|
||||
r"(?:KEY|TOKEN|SECRET|PASSW|AUTH|CREDENTIAL|PRIVATE|DATABASE_URL)",
|
||||
_re.IGNORECASE,
|
||||
)
|
||||
|
||||
|
||||
def _agent_subprocess_env() -> dict:
|
||||
base = {
|
||||
key: os.environ[key]
|
||||
for key in _SAFE_SUBPROCESS_VARS
|
||||
if key in os.environ and not _SENSITIVE_PATTERN.search(key)
|
||||
if key in os.environ
|
||||
}
|
||||
base.setdefault("PATH", os.environ.get("PATH") or os.defpath or "/usr/local/bin:/usr/bin:/bin")
|
||||
base.setdefault("LANG", "C.UTF-8")
|
||||
@@ -1425,7 +1415,7 @@ async def execute_tool_block(
|
||||
owner=owner, session_id=session_id, workspace=workspace,
|
||||
tool_name=getattr(block, "tool_type", None), content=getattr(block, "content", None)))
|
||||
admitted = valid and (authority.permits(operation) or exact_admission)
|
||||
except (ValueError, TypeError, AttributeError) as error:
|
||||
except (ValueError, TypeError) as error:
|
||||
return f"{getattr(block, 'tool_type', '')}: invalid arguments", {
|
||||
"error": (f"Tool arguments are not valid JSON: {error}"
|
||||
if isinstance(error, json.JSONDecodeError) else str(error)),
|
||||
@@ -1503,7 +1493,7 @@ async def execute_tool_block(
|
||||
authority, operation, document_id=active_document_id,
|
||||
approved=pending.owned_operation if pending is not None else None,
|
||||
exact_admission=exact_admission)
|
||||
except (ValueError, TypeError, OSError, RuntimeError, AttributeError) as error:
|
||||
except (ValueError, TypeError, OSError) as error:
|
||||
return f"{transport}: BLOCKED", {
|
||||
"error": str(error), "exit_code": 1, "blocked": True,
|
||||
"failure_kind": "resource_identity_denied",
|
||||
|
||||
+15
-9
@@ -772,9 +772,11 @@ async def do_download_model(content: str, owner: Optional[str] = None) -> Dict:
|
||||
if env_cfg.get("platform"): payload["platform"] = env_cfg["platform"]
|
||||
if env_cfg.get("ssh_port"): payload["ssh_port"] = env_cfg["ssh_port"]
|
||||
try:
|
||||
async with httpx.AsyncClient(timeout=30) as client:
|
||||
resp = await client.post(f"{_INTERNAL_BASE}/api/model/download",
|
||||
json=payload, headers=_internal_headers())
|
||||
from src.agent_runtime.local_model_control import model_control_headers
|
||||
with model_control_headers("download_model", content, owner, payload) as launch_headers:
|
||||
async with httpx.AsyncClient(timeout=30) as client:
|
||||
resp = await client.post(f"{_INTERNAL_BASE}/api/model/download",
|
||||
json=payload, headers=launch_headers)
|
||||
data = resp.json()
|
||||
if data.get("ok"):
|
||||
sid = data.get("session_id", "?")
|
||||
@@ -857,9 +859,11 @@ async def do_serve_model(content: str, owner: Optional[str] = None) -> Dict:
|
||||
if env_cfg.get("platform"): payload["platform"] = env_cfg["platform"]
|
||||
if env_cfg.get("ssh_port"): payload["ssh_port"] = env_cfg["ssh_port"]
|
||||
try:
|
||||
async with httpx.AsyncClient(timeout=30) as client:
|
||||
resp = await client.post(f"{_INTERNAL_BASE}/api/model/serve",
|
||||
json=payload, headers=_internal_headers())
|
||||
from src.agent_runtime.local_model_control import model_control_headers
|
||||
with model_control_headers("serve_model", content, owner, payload) as launch_headers:
|
||||
async with httpx.AsyncClient(timeout=30) as client:
|
||||
resp = await client.post(f"{_INTERNAL_BASE}/api/model/serve",
|
||||
json=payload, headers=launch_headers)
|
||||
data = resp.json()
|
||||
if data.get("ok"):
|
||||
sid = data.get("session_id", "?")
|
||||
@@ -1908,9 +1912,11 @@ async def do_serve_preset(content: str, owner: Optional[str] = None) -> Dict:
|
||||
payload["ssh_port"] = env_cfg["ssh_port"]
|
||||
|
||||
try:
|
||||
async with httpx.AsyncClient(timeout=30) as client:
|
||||
resp = await client.post(f"{_INTERNAL_BASE}/api/model/serve",
|
||||
json=payload, headers=_internal_headers())
|
||||
from src.agent_runtime.local_model_control import model_control_headers
|
||||
with model_control_headers("serve_preset", content, owner, payload) as launch_headers:
|
||||
async with httpx.AsyncClient(timeout=30) as client:
|
||||
resp = await client.post(f"{_INTERNAL_BASE}/api/model/serve",
|
||||
json=payload, headers=launch_headers)
|
||||
data = resp.json()
|
||||
if data.get("ok"):
|
||||
sid = data.get("session_id", "?")
|
||||
|
||||
Reference in New Issue
Block a user