feat(runtime): bind process and job resources to authority

This commit is contained in:
Alexandre Teixeira
2026-10-02 18:54:09 +01:00
parent 7b8ac6f631
commit db41d7e822
36 changed files with 2251 additions and 174 deletions
+67 -15
View File
@@ -13,6 +13,7 @@ from uuid import uuid4
from src.agent_runtime.resources import (
FilesystemRoot, ExternalResource, NativeBackendResource, OwnedScope,
ProcessLaunchScope, ProcessResource, BackgroundJobResource,
backend_from_dict, intersect_roots, seal_owned_scopes,
)
from src.tool_policy import ToolPolicy, build_effective_tool_policy
@@ -120,6 +121,9 @@ class RequestAuthority:
resource_roots: tuple[FilesystemRoot, ...] | None = None
backend_resources: tuple[ExternalResource | NativeBackendResource, ...] | None = None
owned_scopes: tuple[OwnedScope, ...] | None = None
launch_scopes: tuple[ProcessLaunchScope, ...] | None = None
process_resources: tuple[ProcessResource, ...] = ()
job_resources: tuple[BackgroundJobResource, ...] | None = None
def __post_init__(self):
if (not isinstance(self.request_id, str) or not self.request_id
@@ -156,11 +160,29 @@ class RequestAuthority:
or any(not isinstance(s, OwnedScope) or (s.owner, s.thread_id) != (self.owner, self.session_id)
for s in self.owned_scopes)):
raise ValueError("Malformed backend or owned resource scope")
from src.agent_runtime.process_resources import seal_launch_scopes, seal_jobs
if self.launch_scopes is None:
object.__setattr__(self, "launch_scopes", seal_launch_scopes(self))
if self.job_resources is None:
object.__setattr__(self, "job_resources", seal_jobs(self))
for field, kind in (("launch_scopes", ProcessLaunchScope), ("process_resources", ProcessResource),
("job_resources", BackgroundJobResource)):
values = getattr(self, field)
if not isinstance(values, tuple) or any(not isinstance(r, kind) for r in values):
raise ValueError("Malformed process resource scope")
if any(r.owner != self.owner for r in (*self.process_resources, *self.job_resources)):
raise ValueError("Process resource owner changed")
if any(s.root.owner and s.root.owner != self.owner for s in self.launch_scopes):
raise ValueError("Launch resource owner changed")
if any(r.thread_id != self.session_id for r in self.job_resources):
raise ValueError("Job resource thread changed")
if any(r.thread_id != (self.session_id or "request:" + self.request_id) for r in self.process_resources):
raise ValueError("Process resource thread changed")
@classmethod
def empty(cls, *, owner=None, session_id=None, workspace=None):
return cls(uuid4().hex, _owner(owner), str(session_id or ""), str(workspace or ""),
resource_roots=(), backend_resources=(), owned_scopes=())
resource_roots=(), backend_resources=(), owned_scopes=(), launch_scopes=(), job_resources=())
def bound_to(self, *, owner=None, session_id=None, workspace=None):
return (self.owner == _owner(owner) and self.session_id == str(session_id or "")
@@ -188,6 +210,7 @@ class RequestAuthority:
roots = ()
backends = ()
owned = ()
launches = processes = jobs = ()
if (self.owner, self.session_id, self.workspace) == (child.owner, child.session_id, child.workspace):
theirs = {g.tool: g for g in child.grants}
grants = [g.intersect(theirs[g.tool]) for g in self.grants if g.tool in theirs]
@@ -195,10 +218,15 @@ class RequestAuthority:
backends = tuple(r for r in self.backend_resources if r in child.backend_resources)
owned = tuple(s for left in self.owned_scopes for right in child.owned_scopes
if (s := left.intersect(right)) is not None)
from src.agent_runtime.process_resources import intersect_observed, intersect_launch_scopes, validate_job
launches = intersect_launch_scopes(self.launch_scopes, child.launch_scopes)
processes = intersect_observed(self.process_resources, child.process_resources, lambda r: r.validate())
jobs = intersect_observed(self.job_resources, child.job_resources, validate_job)
return replace(self, grants=tuple(grants), denied=self.denied | child.denied,
block_all=self.block_all or child.block_all,
disable_mcp=self.disable_mcp or child.disable_mcp, inherited=True,
resource_roots=roots, backend_resources=backends, owned_scopes=owned)
resource_roots=roots, backend_resources=backends, owned_scopes=owned,
launch_scopes=launches, process_resources=processes, job_resources=jobs)
def continuation(self, *, owner=None, session_id=None):
"""A server continuation may rebind a session, never change owner/grants."""
@@ -206,10 +234,12 @@ class RequestAuthority:
return RequestAuthority.empty(owner=owner, session_id=session_id)
rebound = str(session_id or "")
return replace(self, session_id=rebound, inherited=True,
owned_scopes=tuple(replace(s, thread_id=rebound) for s in self.owned_scopes) if rebound else ())
owned_scopes=tuple(replace(s, thread_id=rebound) for s in self.owned_scopes) if rebound else (),
process_resources=tuple(r for r in self.process_resources if r.thread_id == rebound),
job_resources=tuple(r for r in self.job_resources if r.thread_id == rebound))
def to_dict(self):
return {"version": 3, "request_id": self.request_id, "owner": self.owner,
return {"version": 4, "request_id": self.request_id, "owner": self.owner,
"session_id": self.session_id, "workspace": self.workspace,
"grants": [{"tool": g.tool,
"actions": None if g.actions is None else sorted(g.actions),
@@ -218,12 +248,15 @@ class RequestAuthority:
"disable_mcp": self.disable_mcp, "inherited": self.inherited,
"resource_roots": [r.to_dict() for r in self.resource_roots],
"backend_resources": [r.to_dict() for r in self.backend_resources],
"owned_scopes": [s.to_dict() for s in self.owned_scopes]}
"owned_scopes": [s.to_dict() for s in self.owned_scopes],
"launch_scopes": [s.to_dict() for s in self.launch_scopes],
"process_resources": [r.to_dict() for r in self.process_resources],
"job_resources": [r.to_dict() for r in self.job_resources]}
@classmethod
def from_dict(cls, value):
if (not isinstance(value, dict) or type(value.get("version")) is not int
or value["version"] not in {1, 2, 3}):
or value["version"] not in {1, 2, 3, 4}):
raise ValueError("Unsupported authority snapshot")
def limits(value):
if value is None:
@@ -234,8 +267,12 @@ class RequestAuthority:
roots = value["resource_roots"] if value["version"] >= 2 else []
if not isinstance(roots, list):
raise ValueError("Malformed request resource snapshot")
backends = value["backend_resources"] if value["version"] == 3 else []
owned = value["owned_scopes"] if value["version"] == 3 else []
backends = value["backend_resources"] if value["version"] >= 3 else []
owned = value["owned_scopes"] if value["version"] >= 3 else []
process_fields = {name: value[name] if value["version"] >= 4 else []
for name in ("launch_scopes", "process_resources", "job_resources")}
if any(not isinstance(v, list) for v in process_fields.values()):
raise ValueError("Malformed process resource snapshot")
if not isinstance(backends, list) or not isinstance(owned, list):
raise ValueError("Malformed request resource scope snapshot")
return cls(value["request_id"], value["owner"], value["session_id"], value["workspace"],
@@ -243,7 +280,10 @@ class RequestAuthority:
for g in value["grants"]), limits(value["denied"]),
value["block_all"], value["disable_mcp"], value["inherited"],
tuple(FilesystemRoot.from_dict(r) for r in roots),
tuple(backend_from_dict(r) for r in backends), tuple(OwnedScope.from_dict(s) for s in owned))
tuple(backend_from_dict(r) for r in backends), tuple(OwnedScope.from_dict(s) for s in owned),
tuple(ProcessLaunchScope.from_dict(s) for s in process_fields["launch_scopes"]),
tuple(ProcessResource.from_dict(r) for r in process_fields["process_resources"]),
tuple(BackgroundJobResource.from_dict(r) for r in process_fields["job_resources"]))
_BROWSER_READ_ACTIONS = frozenset({"open", "navigate", "snapshot", "text", "read", "find",
@@ -462,7 +502,10 @@ def seal_task_authority(prompt, task_type, action, *, owner=None, parent_authori
workspace=parent.workspace,
resource_roots=parent.resource_roots,
backend_resources=parent.backend_resources,
owned_scopes=parent.owned_scopes))
owned_scopes=parent.owned_scopes,
launch_scopes=parent.launch_scopes,
process_resources=parent.process_resources,
job_resources=parent.job_resources))
return _json({"task_input": [prompt, task_type, action], "authority": authority.to_dict()})
@@ -479,18 +522,27 @@ def restore_task_authority(snapshot, prompt, task_type, action, *, owner=None, s
def _background_path(job_id):
if not isinstance(job_id, str) or not re.fullmatch(r"[A-Za-z0-9_-]+", job_id):
raise ValueError("Invalid background authority identity")
from src.constants import BG_JOBS_DIR
return Path(BG_JOBS_DIR) / (job_id + ".authority.json")
from src.bg_jobs import _JOBS_DIR
return Path(_JOBS_DIR) / (job_id + ".authority.json")
def save_background_authority(job_id, authority):
def save_background_authority(job_id, authority, *, resource=None):
from core.atomic_io import atomic_write_json
atomic_write_json(_background_path(job_id), authority.to_dict())
if resource is None or resource.job_id != job_id:
raise ValueError("Background authority requires exact job linkage")
atomic_write_json(_background_path(job_id), {"authority": authority.to_dict(), "job": resource.to_dict()})
def restore_background_authority(job_id, *, owner=None, session_id=None):
try:
authority = RequestAuthority.from_dict(json.loads(_background_path(job_id).read_text()))
value = json.loads(_background_path(job_id).read_text())
resource = BackgroundJobResource.from_dict(value["job"])
from src.agent_runtime.process_resources import validate_job
validate_job(resource)
authority = RequestAuthority.from_dict(value["authority"])
if (resource.job_id, resource.owner, resource.thread_id, resource.request_id) != (
job_id, authority.owner, authority.session_id, authority.request_id):
raise ValueError("Background authority linkage changed")
if authority.session_id != str(session_id or ""):
raise ValueError("Background session changed")
return authority.continuation(owner=owner, session_id=session_id)
+2 -1
View File
@@ -257,7 +257,8 @@ def needs_owned_binding(operation):
raise ResourceIdentityError("Unresolved internal resource selector")
path = posixpath.normpath(urlsplit(path).path)
private = {"document", "documents", "session", "sessions", "history", "chat", "chats",
"notes", "memory", "vault", "upload", "uploads", "attachments"}
"notes", "memory", "vault", "upload", "uploads", "attachments",
"shell", "model", "cookbook"}
segments = path.strip("/").split("/")
if len(segments) >= 2 and segments[0] == "api" and segments[1].casefold() in private:
raise ResourceIdentityError("Owned records require a dedicated resource-bound tool")
+418
View File
@@ -0,0 +1,418 @@
"""Process/job admission. Lifecycle mechanics remain in process_lifecycle.
Only trusted launch producers publish observations. Persisted legacy records
are never enrolled by looking at their PID. Receipts identify boundaries, not
application authority. Resource snapshots contain no command or environment.
"""
from __future__ import annotations
from contextlib import contextmanager
from contextvars import ContextVar
from dataclasses import dataclass
import hashlib
import json
import os
from pathlib import Path
import re
from uuid import uuid4
from core.atomic_io import store_transaction
from src.agent_runtime.resources import (
BackgroundJobResource, NativeBackendResource, ProcessLaunchResource,
ProcessLaunchScope, ProcessResource, ResourceIdentityError,
)
from src.constants import PROCESS_RESOURCES_DIR
_LAUNCH_DIR = Path(PROCESS_RESOURCES_DIR)
LAUNCH_TOOLS = frozenset({"bash", "python"})
JOB_TOOL = "manage_bg_jobs"
_ACTIVE = ContextVar("process_resource_operation", default=None)
def digest(value):
return hashlib.sha256(value.encode("utf-8")).hexdigest()
def _thread(authority):
return authority.session_id or "request:" + authority.request_id
def launch_path(generation):
if not isinstance(generation, str) or not re.fullmatch(r"[a-f0-9]{32}", generation):
raise ResourceIdentityError("Malformed launch generation")
return _LAUNCH_DIR / (generation + ".json")
def seal_launch_scopes(authority):
return tuple(seal_launch_scope(backend, root)
for backend in authority.backend_resources
if isinstance(backend, NativeBackendResource) and backend.tool_id in LAUNCH_TOOLS
for root in authority.resource_roots)
def seal_launch_scope(backend, root, *, env=None):
from src.agent_tools.subprocess_tools import _owned_spec
from src.tool_execution import _agent_subprocess_env
from src.agent_runtime.resources import PathObservation, FileObjectIdentity
env = _agent_subprocess_env() if env is None else env
extra = tuple(Path(p).resolve().as_posix() for p in str(env.get("ODYSSEUS_PYTHON_TOOL_SITE_PACKAGES", "")).split(os.pathsep)
if p and os.path.isabs(p)) if backend.tool_id == "python" else ()
spec = _owned_spec(root.path, env, 3600, extra)
return ProcessLaunchScope(backend, root, spec.required,
tuple(PathObservation(str(Path(p).resolve()), FileObjectIdentity.observe(Path(p).resolve())) for p in spec.readonly_extra),
spec.network, spec.wall_clock_s)
def validate_launch_spec(launch, spec):
scope = launch.scope
scope.validate()
if (spec.workspace != scope.root.path or spec.required != scope.required or spec.network != scope.network
or spec.wall_clock_s > scope.max_runtime_s or spec.writable_extra
or tuple(spec.readonly_extra) != tuple(r.path for r in scope.runtime_roots)):
raise ResourceIdentityError("Producer launch boundary exceeds the sealed reservation")
def job_from_record(record):
if not isinstance(record, dict):
raise ResourceIdentityError("Missing authoritative job")
try:
resource = BackgroundJobResource.from_dict(record["resource_identity"])
if (resource.namespace != "native:bg_jobs"
or (record["id"], record["session_id"], record["containment_id"])
!= (resource.job_id, resource.thread_id, resource.containment_id)):
raise ValueError("Job linkage changed")
supervisor = next(p for p in resource.processes if p.role == "supervisor")
if (record.get("pid"), record.get("start_token"), record.get("pgid")) != (
supervisor.identity.pid, supervisor.identity.start_token, supervisor.identity.pgid):
raise ValueError("Supervisor linkage changed")
launch = ProcessLaunchResource.from_dict(record["launch_resource"])
if (launch.generation, launch.owner, launch.request_id, launch.thread_id) != (
resource.generation, resource.owner, resource.request_id, resource.thread_id):
raise ValueError("Launch/job linkage changed")
return resource
except (ValueError, TypeError, KeyError, StopIteration, AttributeError) as error:
raise ResourceIdentityError("Malformed or unowned background job") from error
def validate_job(resource, *, mutation=False):
try:
return _validate_job(resource, mutation=mutation)
except ResourceIdentityError:
raise
except (ValueError, TypeError, OSError, KeyError, AttributeError) as error:
raise ResourceIdentityError("Background job linkage is missing or malformed") from error
def validate_job_receipt(resource, receipt):
from src import containment
supervisor = resource.processes[0]
if (not isinstance(receipt, dict) or receipt.get("id") != resource.containment_id
or receipt.get("launch_generation") != resource.generation
or receipt.get("owner") != "bg:" + resource.thread_id
or (receipt.get("supervisor_pid"), receipt.get("supervisor_token")) !=
(supervisor.identity.pid, supervisor.identity.start_token)
or receipt.get("mechanism") not in {m.name for m in containment.MECHANISMS}
or receipt.get("external") is True):
raise ResourceIdentityError("Containment receipt linkage changed")
def _validate_job(resource, *, mutation=False):
from src import bg_jobs, containment
if not isinstance(resource, BackgroundJobResource):
raise ResourceIdentityError("Missing exact background job identity")
record = bg_jobs.peek(resource.job_id)
if job_from_record(record) != resource:
raise ResourceIdentityError("Background job resource changed")
if record.get("status") not in {"running", "done", "failed"}:
raise ResourceIdentityError("Unknown job lifecycle")
launch = ProcessLaunchResource.from_dict(record["launch_resource"])
persisted = json.loads(launch_path(resource.generation).read_text())
if (persisted.get("launch") != launch.to_dict()
or persisted.get("job") != resource.to_dict()
or persisted.get("containment_id") != resource.containment_id):
raise ResourceIdentityError("Job/launch publication changed")
sidecar = json.loads((bg_jobs._JOBS_DIR / (resource.job_id + ".authority.json")).read_text())
origin = persisted.get("authority", {})
if (sidecar.get("job") != resource.to_dict() or sidecar.get("authority") != origin
or (origin.get("owner"), origin.get("request_id"), origin.get("session_id")) !=
(resource.owner, resource.request_id, resource.thread_id)):
raise ResourceIdentityError("Background authority linkage changed")
receipt = containment._load_records().get(resource.containment_id)
# Lifecycle receipts have a shorter retention than job results. A finished
# exact generation needs only its durable application linkage for history;
# it never regains signalling authority when its receipt has been pruned.
historical = record.get("status") in {"done", "failed"}
if receipt is None and not historical:
raise ResourceIdentityError("Missing active containment receipt")
if receipt is not None:
validate_job_receipt(resource, receipt)
if record.get("status") == "running":
for process in resource.processes:
try:
process.validate()
except ResourceIdentityError:
# Publication can precede store reconciliation. That exact
# completed generation is readable, but never signallable.
if mutation or not Path(record["exit_path"]).is_file():
raise
report = json.loads(Path(record["result_path"]).read_text())
if report.get("resource_identity") != resource.to_dict() or report.get("containment", {}).get("id") != resource.containment_id:
raise ResourceIdentityError("Historical result linkage changed")
# A completed record is readable history, never a new process observation.
return record
def seal_jobs(authority):
if not any(g.tool == JOB_TOOL for g in authority.grants) or not authority.session_id:
return ()
from src import bg_jobs
admitted = []
for record in bg_jobs._load().values():
try:
resource = job_from_record(record)
if (resource.owner, resource.thread_id) == (authority.owner, authority.session_id):
validate_job(resource)
admitted.append(resource)
except (ValueError, TypeError, OSError, RuntimeError):
continue
return tuple(admitted)
def intersect_observed(parent, child, validate):
# Validate both sides before equality. Seeing a replacement cannot renew a
# stale parent observation, even when the child has just sealed it.
for resource in (*parent, *child):
validate(resource)
return tuple(resource for resource in parent if resource in child)
def intersect_launch_scopes(parent, child):
from src.agent_runtime.resources import FilesystemResource
for scope in (*parent, *child):
scope.validate()
narrowed = []
for left in parent:
for right in child:
if (left.backend != right.backend or not left.required <= right.required
or right.max_runtime_s > left.max_runtime_s
or not set(right.runtime_roots) <= set(left.runtime_roots)
or (left.network == "none" and right.network != "none")):
continue
if Path(right.root.path).is_relative_to(left.root.path):
observation = FilesystemResource.resolve(left.root, right.root.path)
if observation.identity == right.root.identity:
narrowed.append(right)
return tuple(dict.fromkeys(narrowed))
@dataclass(frozen=True)
class BoundProcessOperation:
operation: object
request_id: str
owner: str
thread_id: str
launch: ProcessLaunchResource | None = None
jobs: tuple[BackgroundJobResource, ...] = ()
processes: tuple[ProcessResource, ...] = ()
exact_approval: object | None = None
def __post_init__(self):
from src.agent_runtime.authority import ExactOperation
if (not isinstance(self.operation, ExactOperation) or not isinstance(self.request_id, str) or not self.request_id
or not isinstance(self.owner, str) or not isinstance(self.thread_id, str) or not self.thread_id
or (self.launch is not None and not isinstance(self.launch, ProcessLaunchResource))
or not isinstance(self.jobs, tuple) or any(not isinstance(j, BackgroundJobResource) for j in self.jobs)
or not isinstance(self.processes, tuple) or any(not isinstance(p, ProcessResource) for p in self.processes)):
raise ValueError("Malformed process-bound operation")
if self.launch is not None and (
(self.launch.owner, self.launch.request_id, self.launch.thread_id, self.launch.tool, self.launch.input_digest)
!= (self.owner, self.request_id, self.thread_id, self.operation.tool, digest(self.operation.input))):
raise ValueError("Launch operation/application binding changed")
if any((r.owner, r.thread_id) != (self.owner, self.thread_id) for r in (*self.jobs, *self.processes)):
raise ValueError("Observed resource application binding changed")
def validate(self):
if self.launch is not None:
self.launch.validate()
guard_launch_workspace(self.launch.scope.root)
for job in self.jobs:
validate_job(job, mutation=self.operation.action in {"kill", "stop", "cancel", "terminate", "ack"})
for process in self.processes:
process.validate()
def to_dict(self):
return {"tool": self.operation.transport_tool, "input_digest": digest(self.operation.input),
"request_id": self.request_id, "owner": self.owner, "thread_id": self.thread_id,
"launch": self.launch.to_dict() if self.launch else None,
"jobs": [r.to_dict() for r in self.jobs], "processes": [r.to_dict() for r in self.processes]}
def needs_process_binding(operation, backend):
return isinstance(backend, NativeBackendResource) and operation.tool in LAUNCH_TOOLS | {JOB_TOOL}
def resolve_process_operation(authority, operation, backend, *, approved=None, exact_admission=False):
if not needs_process_binding(operation, backend):
raise ResourceIdentityError("No native process adapter for this backend")
if approved is not None:
if (approved.operation != operation or (approved.request_id, approved.owner, approved.thread_id)
!= (authority.request_id, authority.owner, _thread(authority))):
raise ResourceIdentityError("Approved process operation binding changed")
bound = approved
elif operation.tool in LAUNCH_TOOLS:
scopes = [s for s in authority.launch_scopes if s.backend == backend]
if len(scopes) != 1:
raise ResourceIdentityError("Process creation requires a sealed workspace and launch scope")
launch = ProcessLaunchResource("native:containment", authority.owner, authority.request_id,
_thread(authority), uuid4().hex, operation.tool, digest(operation.input), scopes[0],
digest(json.dumps(authority.to_dict(), sort_keys=True)))
bound = BoundProcessOperation(operation, authority.request_id, authority.owner, _thread(authority), launch)
else:
try:
args = json.loads(operation.input)
action = str(args.get("action", "list")).strip().lower()
job_id = args.get("job_id", args.get("id", ""))
except (ValueError, TypeError, AttributeError) as error:
raise ResourceIdentityError("Malformed job operation") from error
if action in {"list", "ls", "jobs"}:
jobs = authority.job_resources
elif action in {"output", "get", "read", "tail", "status", "show", "kill", "stop", "cancel", "terminate", "ack"}:
if not isinstance(job_id, str) or not job_id:
raise ResourceIdentityError("An exact job selector is required")
jobs = tuple(r for r in authority.job_resources if r.job_id == job_id)
if len(jobs) != 1:
raise ResourceIdentityError("Job is outside admitted resource scope")
else:
raise ResourceIdentityError("Unsupported job operation")
bound = BoundProcessOperation(operation, authority.request_id, authority.owner, _thread(authority), jobs=jobs)
if not (approved is not None and exact_admission and not authority.inherited):
if bound.launch is not None and bound.launch.scope not in authority.launch_scopes:
raise ResourceIdentityError("Launch exceeds inherited creation scope")
if any(j not in authority.job_resources for j in bound.jobs) or any(p not in authority.process_resources for p in bound.processes):
raise ResourceIdentityError("Process/job exceeds inherited resource scope")
if bound.launch is not None and bound.launch.scope.backend != backend:
raise ResourceIdentityError("Launch backend changed")
bound.validate()
return bound
def active_process_operation():
return _ACTIVE.get()
@contextmanager
def bind_process_operation(operation):
if operation is not None and not isinstance(operation, BoundProcessOperation):
raise TypeError("Process operation must be server-owned")
if operation is not None:
operation.validate()
token = _ACTIVE.set(operation)
try:
yield operation
finally:
_ACTIVE.reset(token)
def require_launch(tool, *, cwd, content=None):
bound = active_process_operation()
if bound is None or bound.launch is None or bound.operation.tool != tool:
raise ResourceIdentityError("Native process producer has no bound launch reservation")
require_process_admission(bound)
bound.validate()
if Path(cwd).resolve() != Path(bound.launch.scope.root.path):
raise ResourceIdentityError("Launch workspace changed")
if content is not None and content.strip() != bound.operation.input.strip():
raise ResourceIdentityError("Launch operation changed at producer entry")
return bound.launch
def require_process_admission(bound):
from src.agent_runtime.authority import active_request_authority
authority = active_request_authority()
if authority is None or (authority.owner, authority.request_id, _thread(authority)) != (
bound.owner, bound.request_id, bound.thread_id):
raise ResourceIdentityError("Producer application authority changed")
if not authority.permits(bound.operation):
approval = bound.exact_approval
if (authority.inherited or approval is None or not approval._claimed
or approval.pending.process_operation is None
or approval.pending.process_operation.to_dict() != bound.to_dict()):
raise ResourceIdentityError("Producer operation has no request admission or exact claim")
def guard_launch_workspace(root):
"""Reject a boundary containing execution control state or its aliases.
These are pathname/inode observations, not an atomic kernel access policy.
They do not claim freedom from concurrent link replacement after checking.
"""
from src import bg_jobs, containment, constants
from src.agent_runtime.resources import _control_plane_path
control = (Path(bg_jobs._STORE), Path(bg_jobs._JOBS_DIR), containment._store_path(), _LAUNCH_DIR,
Path(constants.APP_DB), Path(constants.AUTH_FILE), Path(constants.SETTINGS_FILE))
base = Path(root.path)
if any(Path(p).resolve().is_relative_to(base) for p in control):
raise ResourceIdentityError("Launch boundary contains server control state")
def unresolved(error):
raise ResourceIdentityError("Launch workspace cannot be inspected") from error
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")
@store_transaction(lambda: _LAUNCH_DIR / "publication")
def publish_launch(launch, authority, containment_id, *, job=None, processes=()):
from core.atomic_io import atomic_write_json
launch.validate()
if authority is None or (authority.owner, authority.request_id) != (launch.owner, launch.request_id):
raise ResourceIdentityError("Launch authority linkage changed")
path = launch_path(launch.generation)
if path.exists():
raise ResourceIdentityError("Launch reservation has already been used")
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 attach_containment_processes(launch, containment_id):
"""Attach producer-frozen lifecycle records; never capture a current PID."""
from src import containment
from src.process_lifecycle import ProcessIdentity
record = containment._load_records().get(containment_id, {})
path = launch_path(launch.generation)
published = json.loads(path.read_text())
if (published.get("launch") != launch.to_dict() or published.get("containment_id") != containment_id
or record.get("id") != containment_id or record.get("launch_generation") != launch.generation
or record.get("workspace") != launch.scope.root.path):
raise ResourceIdentityError("Launch/receipt changed during publication")
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))
from core.atomic_io import atomic_write_json
published["processes"] = [p.to_dict() for p in processes]
atomic_write_json(path, published)
def expected_job(job_id, *, action):
bound = active_process_operation()
if bound is None or bound.operation.tool != JOB_TOOL:
raise ResourceIdentityError("Job producer has no bound operation")
require_process_admission(bound)
# The caller's actual action must agree with the normalized proposal.
args = json.loads(bound.operation.input)
proposed = str(args.get("action", "list")).strip().lower()
if action != proposed:
raise ResourceIdentityError("Job action changed at producer entry")
target = next((j for j in bound.jobs if j.job_id == job_id), None)
if target is None:
raise ResourceIdentityError("Job selector is outside the bound operation")
validate_job(target, mutation=action in {"kill", "stop", "cancel", "terminate", "ack"})
return target
+155 -12
View File
@@ -37,7 +37,10 @@ def _control_plane_path(path):
"SETTINGS_FILE", "SESSIONS_FILE", "USER_PREFS_FILE", "VAULT_FILE",
"SCHEDULED_EMAILS_DB", "EMAIL_CACHE_DB", "MEMORY_FILE", "INTEGRATIONS_FILE",
)}
job_dirs = {canonical_root(constants.BG_JOBS_DIR)}
job_dirs = {canonical_root(constants.BG_JOBS_DIR), canonical_root(constants.PROCESS_RESOURCES_DIR)}
processes = sys.modules.get("src.agent_runtime.process_resources")
if processes is not None:
job_dirs.add(canonical_root(processes._LAUNCH_DIR))
# 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")
@@ -275,25 +278,165 @@ def intersect_roots(parent, child):
@dataclass(frozen=True)
class ProcessResource:
namespace: str
incarnation: str
owner: str
pid: int
start_token: str
request_id: str
thread_id: str
identity: "ProcessIdentity"
role: str
job_id: str = ""
containment_id: str = ""
namespace_pid: int | None = None
namespace_start_token: str = ""
def __post_init__(self):
for name in ("namespace", "incarnation", "owner", "start_token"):
from src.process_lifecycle import ProcessIdentity
for name in ("namespace", "request_id", "thread_id"):
_text(getattr(self, name), name)
for name in ("job_id", "containment_id", "namespace_start_token"):
for name in ("owner", "job_id", "containment_id"):
_text(getattr(self, name), name, optional=True)
if (type(self.pid) is not int or self.pid <= 0
or (self.namespace_pid is not None and
(type(self.namespace_pid) is not int or self.namespace_pid <= 0))
or bool(self.namespace_pid) != bool(self.namespace_start_token)):
if (not isinstance(self.identity, ProcessIdentity)
or type(self.identity.pid) is not int or self.identity.pid <= 0
or (self.identity.pgid is not None and (type(self.identity.pgid) is not int or self.identity.pgid <= 0))
or self.role not in {"supervisor", "leader", "namespace_init", "manager", "pty", "service"}):
raise ValueError("Malformed process resource identity")
supported_roles = {"native:containment": {"leader", "namespace_init"},
"native:bg_jobs": {"supervisor"}}
if self.role not in supported_roles.get(self.namespace, set()):
raise ValueError("Unsupported process producer or role")
_text(self.identity.start_token, "process start token")
def validate(self):
if not self.identity.owned() or self.identity.exited():
raise ResourceIdentityError("Process resource is stale or unverifiable")
def to_dict(self):
return {"namespace": self.namespace, "owner": self.owner, "request_id": self.request_id,
"thread_id": self.thread_id, "identity": self.identity.to_record(), "role": self.role,
"job_id": self.job_id, "containment_id": self.containment_id}
@classmethod
def from_dict(cls, value):
from src.process_lifecycle import ProcessIdentity
if not isinstance(value, dict) or set(value) != {"namespace", "owner", "request_id", "thread_id", "identity", "role", "job_id", "containment_id"}:
raise ValueError("Malformed process resource snapshot")
identity = value["identity"]
if not isinstance(identity, dict) or set(identity) != {"pid", "start_token", "pgid"}:
raise ValueError("Malformed lifecycle identity snapshot")
return cls(**{**value, "identity": ProcessIdentity(**identity)})
@dataclass(frozen=True)
class ProcessLaunchScope:
backend: "NativeBackendResource"
root: FilesystemRoot
required: frozenset[str]
runtime_roots: tuple[PathObservation, ...] = ()
network: str = "inherit"
max_runtime_s: int = 3600
def __post_init__(self):
if (not isinstance(self.backend, NativeBackendResource) or not isinstance(self.root, FilesystemRoot)
or not isinstance(self.required, frozenset) or not self.required
or any(not isinstance(v, str) or not v for v in self.required)):
raise ValueError("Malformed process launch scope")
if self.backend.tool_id not in {"bash", "python"}:
raise ValueError("Unsupported native launch producer")
if (not isinstance(self.runtime_roots, tuple) or any(not isinstance(r, PathObservation) for r in self.runtime_roots)
or self.network not in {"inherit", "none"}
or type(self.max_runtime_s) is not int or self.max_runtime_s <= 0):
raise ValueError("Malformed launch boundary selectors")
def validate(self):
self.root.validate()
for runtime in self.runtime_roots:
if canonical_root(runtime.path) != runtime.path or FileObjectIdentity.observe(runtime.path) != runtime.identity:
raise ResourceIdentityError("Launch runtime root changed")
def to_dict(self):
return {"backend": self.backend.to_dict(), "root": self.root.to_dict(), "required": sorted(self.required),
"runtime_roots": [{"path": r.path, "identity": asdict(r.identity)} for r in self.runtime_roots],
"network": self.network, "max_runtime_s": self.max_runtime_s}
@classmethod
def from_dict(cls, value):
if not isinstance(value, dict) or set(value) != {"backend", "root", "required", "runtime_roots", "network", "max_runtime_s"} or not isinstance(value["required"], list) or not isinstance(value["runtime_roots"], list):
raise ValueError("Malformed launch scope snapshot")
return cls(backend_from_dict(value["backend"]), FilesystemRoot.from_dict(value["root"]), frozenset(value["required"]),
tuple(PathObservation(r["path"], FileObjectIdentity(**r["identity"])) for r in value["runtime_roots"]),
value["network"], value["max_runtime_s"])
@dataclass(frozen=True)
class ProcessLaunchResource:
namespace: str
owner: str
request_id: str
thread_id: str
generation: str
tool: str
input_digest: str
scope: ProcessLaunchScope
ceiling_digest: str
def __post_init__(self):
for name in ("namespace", "request_id", "thread_id", "generation", "tool", "input_digest", "ceiling_digest"):
_text(getattr(self, name), name)
_text(self.owner, "owner", optional=True)
if not isinstance(self.scope, ProcessLaunchScope) or self.tool != self.scope.backend.tool_id:
raise ValueError("Malformed launch resource")
import re
if (self.namespace != "native:containment" or not re.fullmatch(r"[a-f0-9]{32}", self.generation)
or any(not re.fullmatch(r"[a-f0-9]{64}", v) for v in (self.input_digest, self.ceiling_digest))):
raise ValueError("Malformed native launch producer or generation")
def validate(self):
self.scope.validate()
def to_dict(self):
return {**{k: getattr(self, k) for k in ("namespace", "owner", "request_id", "thread_id", "generation", "tool", "input_digest", "ceiling_digest")},
"scope": self.scope.to_dict()}
@classmethod
def from_dict(cls, value):
if not isinstance(value, dict) or set(value) != {"namespace", "owner", "request_id", "thread_id", "generation", "tool", "input_digest", "scope", "ceiling_digest"}:
raise ValueError("Malformed launch resource snapshot")
return cls(**{**value, "scope": ProcessLaunchScope.from_dict(value["scope"])})
@dataclass(frozen=True)
class BackgroundJobResource:
namespace: str
job_id: str
generation: str
owner: str
request_id: str
thread_id: str
containment_id: str
processes: tuple[ProcessResource, ...]
def __post_init__(self):
for name in ("namespace", "job_id", "generation", "request_id", "thread_id", "containment_id"):
_text(getattr(self, name), name)
_text(self.owner, "owner", optional=True)
import re
if (not re.fullmatch(r"[A-Za-z0-9_-]+", self.job_id)
or not re.fullmatch(r"[a-f0-9]{32}", self.generation)):
raise ValueError("Malformed job selector or launch generation")
if (not isinstance(self.processes, tuple) or not self.processes
or any(not isinstance(p, ProcessResource) or (p.owner, p.request_id, p.thread_id, p.job_id, p.containment_id)
!= (self.owner, self.request_id, self.thread_id, self.job_id, self.containment_id) for p in self.processes)
or len({p.role for p in self.processes}) != len(self.processes)):
raise ValueError("Malformed background job resource")
if self.namespace != "native:bg_jobs" or any(p.namespace != "native:bg_jobs" or p.role != "supervisor" for p in self.processes):
raise ValueError("Unsupported job producer or process role")
def to_dict(self):
return {**{k: getattr(self, k) for k in ("namespace", "job_id", "generation", "owner", "request_id", "thread_id", "containment_id")},
"processes": [p.to_dict() for p in self.processes]}
@classmethod
def from_dict(cls, value):
if not isinstance(value, dict) or set(value) != {"namespace", "job_id", "generation", "owner", "request_id", "thread_id", "containment_id", "processes"} or not isinstance(value["processes"], list):
raise ValueError("Malformed background resource snapshot")
return cls(**{**value, "processes": tuple(ProcessResource.from_dict(p) for p in value["processes"])})
@dataclass(frozen=True)
+19 -3
View File
@@ -67,8 +67,20 @@ class ManageBgJobsTool:
if not session_id:
return {"error": "manage_bg_jobs: no active chat session; background jobs are scoped to a chat.", "exit_code": 1}
from src.agent_runtime.process_resources import active_process_operation, expected_job, require_process_admission
from src.agent_runtime.resources import ResourceIdentityError
bound = active_process_operation()
if bound is None or (bound.owner, bound.thread_id) != (str(ctx.get("owner") or "").strip().casefold(), session_id):
return {"error": "manage_bg_jobs: no exact server resource binding", "exit_code": 1,
"blocked": True, "failure_kind": "resource_identity_denied"}
from src.agent_runtime.authority import ExactOperation
if bound.operation != ExactOperation.normalize("manage_bg_jobs", raw or "{}"):
return {"error": "Job operation changed at producer entry", "exit_code": 1, "blocked": True}
require_process_admission(bound)
if action in _LIST_ACTIONS:
jobs: List[Dict[str, Any]] = bg_jobs.list_for_session(session_id)
bound.validate()
jobs: List[Dict[str, Any]] = [bg_jobs.peek(j.job_id) for j in bound.jobs]
if not jobs:
return {"output": "No background jobs in this chat.", "exit_code": 0}
jobs.sort(key=lambda r: r.get("started_at") or 0, reverse=True)
@@ -78,7 +90,11 @@ class ManageBgJobsTool:
if action in _OUTPUT_ACTIONS or action in _KILL_ACTIONS:
if not job_id:
return {"error": f"manage_bg_jobs: action '{action}' requires a job_id (see action='list').", "exit_code": 1}
rec = bg_jobs.get(job_id)
try:
resource = expected_job(job_id, action=action)
rec = bg_jobs.get(job_id, expected=resource)
except (ResourceIdentityError, OSError, ValueError) as error:
return {"error": str(error), "exit_code": 1, "blocked": True, "failure_kind": "resource_identity_denied"}
# Scope: only the chat that launched a job may see or control it.
if rec is None or rec.get("session_id") != session_id:
return {"error": f"manage_bg_jobs: no background job '{job_id}' in this chat.", "exit_code": 1}
@@ -86,7 +102,7 @@ class ManageBgJobsTool:
if action in _KILL_ACTIONS:
if rec.get("status") != "running":
return {"output": f"Job `{job_id}` already {_status_label(rec)}; nothing to kill.", "exit_code": 0}
killed = bg_jobs.kill(job_id)
killed = bg_jobs.kill(job_id, expected=resource)
if not killed or not killed.get("killed"):
return {"error": f"Could not verify termination of background job `{job_id}`.",
"exit_code": 1, "teardown": (killed or {}).get("teardown")}
+38 -5
View File
@@ -511,26 +511,47 @@ async def _run_owned_command(command, ctx: dict, *, tool: str, timeout: int, arg
from src.tool_execution import agent_cwd, _truncate
grant = None
result = None
try:
from src.agent_runtime.process_resources import require_launch, publish_launch, validate_launch_spec
from src.agent_runtime.authority import active_request_authority
launch = require_launch(tool, cwd=agent_cwd())
authority = active_request_authority()
if (str(ctx.get("owner") or "").strip().casefold(), str(ctx.get("session_id") or "")) != (
authority.owner, authority.session_id):
raise ValueError("Native producer owner or session changed")
spec = _owned_spec(agent_cwd(), ctx.get("subproc_env"), timeout, readonly_extra)
validate_launch_spec(launch, spec)
grant = containment.acquire(
_owned_spec(agent_cwd(), ctx.get("subproc_env"), timeout, readonly_extra),
spec,
owner=str(ctx.get("session_id") or ctx.get("owner") or tool),
)
containment._update_record(grant.id, launch_generation=launch.generation)
publish_launch(launch, authority, grant.id)
if containment.FILESYSTEM not in grant.enforced:
if argv:
command = [*command[:-1], _replace_workspace_alias(command[-1], grant.workspace)]
else:
command = _replace_workspace_alias(command, grant.workspace)
result = await containment.run(grant, command, argv=argv, progress_cb=ctx.get("progress_cb"))
from src.agent_runtime.process_resources import attach_containment_processes
attach_containment_processes(launch, grant.id)
except containment.ContainmentUnavailable as exc:
return containment.unavailable_tool_result(exc, tool=tool)
except (OSError, RuntimeError, ValueError) as exc:
boundary = grant.to_dict() if grant else {}
boundary["executed"] = bool(getattr(exc, "containment_executed", False))
if not getattr(exc, "containment_established", False):
if grant is not None:
record = containment._load_records().get(grant.id, {})
if not record.get("pid") and not record.get("release"):
containment.release(grant, grace_s=0)
boundary = result.grant.to_dict() if result is not None else grant.to_dict() if grant else {}
boundary["executed"] = result is not None or bool(getattr(exc, "containment_executed", False))
if result is None and not getattr(exc, "containment_established", False):
boundary.update(contained=False, enforced=[])
return {"error": f"{tool}: execution failed: {exc}", "exit_code": 1,
"containment": boundary}
"containment": boundary,
**({"failure_kind": "resource_linkage_unavailable",
"teardown": result.release.to_dict() if result.release else {"dead": False}}
if result is not None else {})}
boundary = result.grant.to_dict()
boundary["executed"] = True
@@ -590,6 +611,12 @@ class BashTool:
),
"exit_code": 1,
}
from src.agent_runtime.process_resources import require_launch
from src.agent_runtime.resources import ResourceIdentityError
try:
require_launch("bash", cwd=agent_cwd(), content=content)
except ResourceIdentityError as error:
return {"error": str(error), "exit_code": 1, "blocked": True, "failure_kind": "resource_identity_denied"}
if _ffmpeg_unicode_drawtext_needs_fontfile(content):
resolved_font = _resolve_fontfile_for_text(content)
resolved_hint = (
@@ -879,6 +906,12 @@ class PythonTool:
),
"exit_code": 1,
}
from src.agent_runtime.process_resources import require_launch
from src.agent_runtime.resources import ResourceIdentityError
try:
require_launch("python", cwd=agent_cwd(), content=content)
except ResourceIdentityError as error:
return {"error": str(error), "exit_code": 1, "blocked": True, "failure_kind": "resource_identity_denied"}
if "/tmp/" in content:
isolated_tmp = _isolated_tmp_dir(agent_cwd())
content = content.replace("/tmp/", isolated_tmp.rstrip("/") + "/")
+76 -11
View File
@@ -86,6 +86,20 @@ def launch(command: str, session_id: str, cwd: Optional[str] = None,
A trusted detached supervisor owns the shared containment runner, output,
wall clock and exit metadata, independently of the request/server lifetime.
"""
from src.agent_runtime.process_resources import require_launch, active_process_operation, publish_launch, launch_path, validate_launch_spec
from src.agent_runtime.authority import active_request_authority, save_background_authority
from src.agent_runtime.resources import ProcessResource, BackgroundJobResource
from src.process_lifecycle import ProcessIdentity
cwd = cwd or os.getcwd()
launch_resource = require_launch("bash", cwd=cwd)
bound = active_process_operation()
from src.tool_execution import _split_bg_marker
marked, proposed = _split_bg_marker(bound.operation.input)
if command != (proposed if marked else bound.operation.input).strip() or session_id != launch_resource.thread_id:
raise ValueError("Background launch operation or session changed")
authority = active_request_authority()
if authority is None or (authority.owner, authority.request_id) != (launch_resource.owner, launch_resource.request_id):
raise ValueError("Background launch authority changed")
_JOBS_DIR.mkdir(parents=True, exist_ok=True)
job_id = uuid.uuid4().hex[:12]
log_path = _JOBS_DIR / f"{job_id}.log"
@@ -94,6 +108,7 @@ def launch(command: str, session_id: str, cwd: Optional[str] = None,
from src import containment
from src.agent_tools.subprocess_tools import _owned_spec, _replace_workspace_alias
spec = _owned_spec(cwd or os.getcwd(), env, max_runtime_s)
validate_launch_spec(launch_resource, spec)
grant = containment.acquire(spec, owner=f"bg:{session_id}")
bounded_command = command
if containment.FILESYSTEM not in grant.enforced:
@@ -147,16 +162,33 @@ def launch(command: str, session_id: str, cwd: Optional[str] = None,
"start_token": process_ownership.capture(proc.pid)["start_token"],
}
try:
supervisor = ProcessResource("native:bg_jobs", launch_resource.owner, launch_resource.request_id,
launch_resource.thread_id, ProcessIdentity(proc.pid, rec["start_token"], rec["pgid"]),
"supervisor", job_id, grant.id)
supervisor.validate()
resource = BackgroundJobResource("native:bg_jobs", job_id, launch_resource.generation,
launch_resource.owner, launch_resource.request_id, launch_resource.thread_id, grant.id, (supervisor,))
rec["resource_identity"] = resource.to_dict()
rec["launch_resource"] = launch_resource.to_dict()
containment._update_record(grant.id, lifetime="background", supervisor_pid=proc.pid,
supervisor_token=rec["start_token"])
supervisor_token=rec["start_token"], launch_generation=resource.generation)
jobs = _load()
jobs[job_id] = rec
_save(jobs)
publish_launch(launch_resource, authority, grant.id, job=resource, processes=(supervisor,))
save_background_authority(job_id, authority, resource=resource)
payload.update(job_store=str(_STORE.resolve()), job_id=job_id,
launch_path=str(launch_path(resource.generation)),
authority_path=str(_JOBS_DIR / (job_id + ".authority.json")),
resource_identity=resource.to_dict(), launch_resource=launch_resource.to_dict())
# The supervisor cannot execute until the identity and job record are durable.
proc.stdin.write(json.dumps(payload).encode("utf-8"))
proc.stdin.close()
except BaseException:
kill_process_tree(proc.pid)
# EOF closes the unreleased worker even if identity observation failed.
if proc.stdin is not None and not proc.stdin.closed:
proc.stdin.close()
kill_process_tree(proc.pid, start_token=rec["start_token"], pgid=rec["pgid"], require_identity=True)
proc.wait(timeout=5)
containment.release(grant, grace_s=0)
raise
@@ -194,16 +226,20 @@ def _prune(jobs: Dict[str, Dict[str, Any]], now: float) -> bool:
@store_transaction(lambda: _STORE)
def refresh() -> Dict[str, Dict[str, Any]]:
def refresh(job_id=None) -> Dict[str, Dict[str, Any]]:
"""Reconcile every running job against disk. Marks done/failed (incl.
timeout). Idempotent — safe to call from a poll loop. Returns the store."""
jobs = _load()
for pid, proc in list(_LIVE_PROCS.items()):
if job_id is not None and pid != jobs.get(job_id, {}).get("pid"):
continue
if proc.poll() is not None:
_LIVE_PROCS.pop(pid, None)
changed = False
now = time.time()
for rec in jobs.values():
for jid, rec in jobs.items():
if job_id is not None and jid != job_id:
continue
if rec.get("status") != "running":
continue
exit_path = Path(rec.get("exit_path", ""))
@@ -218,7 +254,15 @@ def refresh() -> Dict[str, Dict[str, Any]]:
if rec.get("result_path"):
try:
report = json.loads(Path(rec["result_path"]).read_text(encoding="utf-8"))
rec.update(report)
# Result publication is not an identity producer. It cannot
# overwrite ownership, generations, PIDs, paths or authority.
if rec.get("resource_identity") and report.get("resource_identity") != rec["resource_identity"]:
raise ValueError("Result/job linkage mismatch")
if report.get("containment", {}).get("id") != rec.get("containment_id"):
raise ValueError("Result/receipt linkage mismatch")
for key in ("containment", "teardown", "output_truncated", "timed_out", "error", "failure_kind"):
if key in report:
rec[key] = report[key]
except (OSError, ValueError):
rec["status"], rec["exit_code"] = "failed", 1
rec["result_unavailable"] = True
@@ -243,7 +287,7 @@ def refresh() -> Dict[str, Dict[str, Any]]:
rec["ended_at"] = now
rec["died"] = True
changed = True
if _prune(jobs, now):
if job_id is None and _prune(jobs, now):
changed = True
if changed:
_save(jobs)
@@ -288,28 +332,45 @@ def pending_followups() -> List[Dict[str, Any]]:
@store_transaction(lambda: _STORE)
def mark_followed_up(job_id: str) -> None:
def mark_followed_up(job_id: str, *, expected) -> None:
jobs = _load()
if job_id in jobs:
from src.agent_runtime.process_resources import validate_job
if expected.job_id != job_id:
raise ValueError("Acknowledgement job resource changed")
validate_job(expected, mutation=True)
jobs[job_id]["followed_up"] = True
_save(jobs)
def get(job_id: str) -> Optional[Dict[str, Any]]:
refresh() # reconcile against disk so status/exit_code are current
def peek(job_id: str) -> Optional[Dict[str, Any]]:
"""Resolve one record without reaping or changing any job."""
return _load().get(job_id)
def get(job_id: str, *, expected) -> Optional[Dict[str, Any]]:
from src.agent_runtime.process_resources import validate_job
if expected.job_id != job_id:
raise ValueError("Output job selector changed")
validate_job(expected)
refresh(job_id)
validate_job(expected)
rec = _load().get(job_id)
if rec:
from src.agent_runtime.process_resources import job_from_record
if job_from_record(rec) != expected:
raise ValueError("Output job resource changed")
rec = dict(rec)
rec["output"] = _read_output(rec)
return rec
def list_for_session(session_id: str) -> List[Dict[str, Any]]:
return [r for r in refresh().values() if r.get("session_id") == session_id]
return [r for r in _load().values() if r.get("session_id") == session_id]
@store_transaction(lambda: _STORE)
def kill(job_id: str) -> Optional[Dict[str, Any]]:
def kill(job_id: str, *, expected) -> Optional[Dict[str, Any]]:
"""Terminate a running job's process tree and mark it killed. Returns the
updated record, or None if the id is unknown. Idempotent: a job that already
finished is returned unchanged. Sets followed_up so the monitor does not also
@@ -318,6 +379,10 @@ def kill(job_id: str) -> Optional[Dict[str, Any]]:
rec = jobs.get(job_id)
if rec is None:
return None
from src.agent_runtime.process_resources import validate_job
if expected.job_id != job_id:
raise ValueError("Job selector changed")
validate_job(expected, mutation=True)
if rec.get("status") == "running":
outcome = _kill_record(rec)
rec["teardown"] = outcome.to_dict()
+13 -1
View File
@@ -140,6 +140,17 @@ async def _run_followup(rec: dict) -> bool:
from src.settings import get_setting
authority = restore_background_authority(
rec["id"], owner=getattr(sess, "owner", None), session_id=sess.id)
# A result can trigger a continuation only through the immutable producer
# linkage, never merely because it names an existing chat.
from src.agent_runtime.process_resources import job_from_record, validate_job
try:
resource = job_from_record(rec)
validate_job(resource)
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
except (ValueError, TypeError, OSError, RuntimeError):
return False
authority = authority.restrict(disabled_tools=get_setting("disabled_tools", []) or ())
full, tool_events = await _drain_agent(sess, context, request_authority=authority)
@@ -169,7 +180,8 @@ async def _loop():
for rec in bg_jobs.pending_followups():
try:
if await _run_followup(rec):
bg_jobs.mark_followed_up(rec["id"])
from src.agent_runtime.process_resources import job_from_record
bg_jobs.mark_followed_up(rec["id"], expected=job_from_record(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)
+21 -15
View File
@@ -878,22 +878,28 @@ async def action_consolidate_memory(owner: str, **kwargs) -> Tuple[str, bool]:
async def _run_subprocess(argv, *, shell: bool = False, timeout: int = 120, label: str = "Command") -> Tuple[str, bool]:
"""Shared subprocess runner. Wraps the blocking subprocess.run in
asyncio.to_thread so the event loop stays responsive."""
import asyncio
import subprocess
"""Scheduled local work consumes the request's sealed launch ceiling."""
from src.agent_runtime.authority import active_request_authority, ExactOperation
from src.agent_runtime.process_resources import resolve_process_operation, bind_process_operation
from src.agent_runtime.resources import NativeBackendResource
from src.agent_tools.subprocess_tools import _run_owned_command
authority = active_request_authority()
if authority is None:
return "Scheduled process launch has no server authority.", False
if isinstance(argv, list) and argv and argv[0] == "ssh":
return "Remote scheduled workload requires an exact external backend binding.", False
command = argv[-1] if isinstance(argv, list) else argv
operation = ExactOperation.normalize("bash", command)
if not authority.permits(operation):
return "Scheduled launch differs from the sealed operation.", False
try:
result = await asyncio.to_thread(
subprocess.run, argv, shell=shell, capture_output=True, text=True, timeout=timeout,
)
output = (result.stdout or "").strip()
if result.returncode != 0 and result.stderr:
output += "\nSTDERR: " + result.stderr.strip()
return output or "(no output)", result.returncode == 0
except subprocess.TimeoutExpired:
return f"{label} timed out ({timeout}s)", False
except Exception as e:
return str(e), False
bound = resolve_process_operation(authority, operation, NativeBackendResource("bash"))
with bind_process_operation(bound):
result = await _run_owned_command(command, {"owner": authority.owner,
"session_id": authority.session_id}, tool="bash", timeout=timeout)
return result.get("output") or result.get("error") or "(no output)", result.get("exit_code") == 0
except (ValueError, OSError, RuntimeError) as error:
return str(error), False
async def action_ssh_command(owner: str, command: str = "", host: str = "localhost", **kwargs) -> Tuple[str, bool]:
+1
View File
@@ -89,6 +89,7 @@ EMOJI_CACHE_DIR = os.path.join(DATA_DIR, "emoji_cache")
RAG_DIR = os.path.join(DATA_DIR, "rag")
CHROMA_DIR = os.path.join(DATA_DIR, "chroma")
BG_JOBS_DIR = os.path.join(DATA_DIR, "bg_jobs")
PROCESS_RESOURCES_DIR = os.path.join(DATA_DIR, "process_resources")
DEEP_RESEARCH_DIR = os.path.join(DATA_DIR, "deep_research")
MCP_OAUTH_DIR = os.path.join(DATA_DIR, "mcp_oauth")
GENERATED_IMAGES_DIR = os.path.join(DATA_DIR, "generated_images")
+36
View File
@@ -6,6 +6,7 @@ import json
import signal
import sys
import types
import os
from pathlib import Path
# Launch by absolute script path, so a task workspace cannot shadow src.
@@ -39,6 +40,40 @@ async def supervise(payload: dict) -> None:
loop.add_signal_handler(signal.SIGTERM, task.cancel)
loop.add_signal_handler(signal.SIGINT, task.cancel)
try:
# The supervisor is held on stdin until *all* publication succeeds.
# No legacy payload can reconstruct ownership from its PID or receipt.
job = json.loads(Path(payload["job_store"]).read_text())[payload["job_id"]]
published = json.loads(Path(payload["launch_path"]).read_text())
sidecar = json.loads(Path(payload["authority_path"]).read_text())
resource = payload["resource_identity"]
launch = payload["launch_resource"]
from src.agent_runtime.resources import ProcessLaunchResource, BackgroundJobResource
from src.agent_runtime.process_resources import validate_launch_spec, validate_job_receipt
typed_launch = ProcessLaunchResource.from_dict(launch)
typed_job = BackgroundJobResource.from_dict(resource)
typed_launch.validate()
validate_launch_spec(typed_launch, spec)
supervisor = typed_job.processes[0]
supervisor.validate()
receipt = containment._load_records().get(grant.id)
validate_job_receipt(typed_job, receipt)
if (supervisor.identity.pid != os.getpid()
or (typed_job.owner, typed_job.request_id, typed_job.thread_id) !=
(typed_launch.owner, typed_launch.request_id, typed_launch.thread_id)
or (published["authority"]["owner"], published["authority"]["request_id"], published["authority"]["session_id"]) !=
(typed_job.owner, typed_job.request_id, typed_job.thread_id)):
raise ValueError("Detached producer ownership changed")
if (job.get("resource_identity") != resource or job.get("launch_resource") != launch
or published.get("job") != resource or published.get("launch") != launch
or sidecar.get("job") != resource or sidecar.get("authority") != published.get("authority")
or published.get("containment_id") != grant.id
or (receipt.get("owner"), receipt.get("mechanism"), receipt.get("mode"), receipt.get("workspace")) !=
(grant.owner, grant.mechanism, grant.mode, spec.workspace)
or info.get("external") is True
or resource["containment_id"] != grant.id
or resource["generation"] != launch["generation"]
or receipt.get("launch_generation") != launch["generation"]):
raise ValueError("Detached launch authority linkage mismatch")
with open(payload["log_path"], "w", encoding="utf-8") as log:
def capture(text):
log.write(text)
@@ -75,6 +110,7 @@ async def supervise(payload: dict) -> None:
except OSError:
# A failed log initialization must not hide completion metadata.
sys.stderr.write(output)
report["resource_identity"] = payload.get("resource_identity")
atomic_write_json(payload["result_path"], report)
# Publish completion last: refresh must never see an exit without metadata.
atomic_write_text(payload["exit_path"], str(code if code is not None else 1))
+11
View File
@@ -31,6 +31,7 @@ if TYPE_CHECKING:
from src.agent_runtime.resource_binding import BoundFilesystemOperation
from src.agent_runtime.remote_resources import BoundBackendOperation
from src.agent_runtime.owned_resources import BoundOwnedOperation
from src.agent_runtime.process_resources import BoundProcessOperation
DEFAULT_APPROVAL_TTL_SECONDS = 10 * 60
@@ -127,6 +128,7 @@ def _binding_payload(
resource_operation=None,
backend_operation=None,
owned_operation=None,
process_operation=None,
) -> dict[str, Any]:
return {
"owner": _normalized_owner(owner),
@@ -151,6 +153,7 @@ def _binding_payload(
"resource_operation": resource_operation.to_dict() if resource_operation is not None else None,
"backend_operation": backend_operation.to_dict() if backend_operation is not None else None,
"owned_operation": owned_operation.to_dict() if owned_operation is not None else None,
"process_operation": process_operation.to_dict() if process_operation is not None else None,
}
@@ -184,6 +187,7 @@ class PendingToolApproval:
resource_operation: BoundFilesystemOperation | None = None
backend_operation: BoundBackendOperation | None = None
owned_operation: BoundOwnedOperation | None = None
process_operation: BoundProcessOperation | None = None
def public_payload(self, *, reason: str | None = None) -> dict[str, Any]:
return {
@@ -296,6 +300,7 @@ class ExactToolApproval:
resource_operation=self.pending.resource_operation,
backend_operation=self.pending.backend_operation,
owned_operation=self.pending.owned_operation,
process_operation=self.pending.process_operation,
)
return _canonical_digest(expected) == self.pending.digest
@@ -391,6 +396,7 @@ class ToolApprovalStore:
resource_operation = None
backend_operation = None
owned_operation = None
process_operation = None
from src.agent_runtime.remote_resources import BoundBackendOperation, resolve_backend
from src.agent_runtime.owned_resources import needs_owned_binding, resolve_owned_operation
from src.agent_runtime.resources import NativeBackendResource
@@ -403,6 +409,9 @@ class ToolApprovalStore:
backend_operation = BoundBackendOperation(backend,
request_authority.request_id if request_authority is not None else "",
_normalized_owner(owner), str(session_id or ""), operation.transport_tool, operation.input)
from src.agent_runtime.process_resources import needs_process_binding, resolve_process_operation
if request_authority is not None and needs_process_binding(operation, backend):
process_operation = resolve_process_operation(request_authority, operation, backend)
if isinstance(backend, NativeBackendResource) and needs_owned_binding(operation):
resolved_owned = resolve_owned_operation(operation, owner=_normalized_owner(owner),
thread_id=str(session_id or ""), request_id=backend_operation.request_id,
@@ -450,6 +459,7 @@ class ToolApprovalStore:
resource_operation=resource_operation,
backend_operation=backend_operation,
owned_operation=owned_operation,
process_operation=process_operation,
)
pending = PendingToolApproval(
approval_id=secrets.token_urlsafe(32),
@@ -477,6 +487,7 @@ class ToolApprovalStore:
resource_operation=resource_operation,
backend_operation=backend_operation,
owned_operation=owned_operation,
process_operation=process_operation,
)
with self._lock:
self._purge_expired_locked(now)
+27 -4
View File
@@ -978,7 +978,10 @@ def vet_workspace(raw: str) -> Optional[str]:
def agent_cwd() -> str:
"""Working directory for agent subprocesses (bash/python/background jobs):
the active workspace when set, else the persistent data dir."""
return get_active_workspace() or _AGENT_WORKDIR
from src.agent_runtime.process_resources import active_process_operation
bound = active_process_operation()
return (bound.launch.scope.root.path if bound is not None and bound.launch is not None
else get_active_workspace() or _AGENT_WORKDIR)
def get_mcp_manager():
@@ -1319,7 +1322,10 @@ async def _document_tool_dispatch(
from src.agent_runtime.journal import dispatched, mark_authorized, mark_dispatch, record_action
from src.agent_runtime.authority import (
MISSING_AUTHORITY, ExactOperation, RequestAuthority, active_request_authority,
bind_request_authority, save_background_authority,
bind_request_authority,
)
from src.agent_runtime.process_resources import (
active_process_operation, bind_process_operation, needs_process_binding, resolve_process_operation,
)
@@ -1415,6 +1421,12 @@ async def execute_tool_block(
exact_admission=exact_admission)
external_resource_call = isinstance(backend_operation.resource, ExternalResource)
owned_operation = None
process_operation = None
if needs_process_binding(operation, backend_operation.resource):
if pending is not None and pending.process_operation is None:
raise ResourceIdentityError("Approved action has no sealed process/job identity")
process_operation = resolve_process_operation(authority, operation, backend_operation.resource,
approved=pending.process_operation if pending is not None else None, exact_admission=exact_admission)
if needs_owned_binding(operation) and not external_resource_call:
if pending is not None and pending.owned_operation is None:
raise ResourceIdentityError("Approved action has no sealed owned resource identity")
@@ -1541,10 +1553,13 @@ async def execute_tool_block(
token = _active_workspace.set(workspace or None)
try:
backend_operation.validate(client_runtime_context)
if process_operation is not None and approval_claimed:
process_operation = replace(process_operation, exact_approval=exact_approval)
normalized = resource_operation or owned_operation
sealed_document = owned_operation or (exact_approval.pending if approval_claimed else None)
with (bind_request_authority(authority), bind_resource_operation(resource_operation),
bind_backend_operation(backend_operation), bind_owned_operation(owned_operation)):
bind_backend_operation(backend_operation), bind_owned_operation(owned_operation),
bind_process_operation(process_operation)):
output = await _execute_tool_block_impl(
ToolBlock(transport, normalized.execution_input) if normalized is not None else block,
session_id=session_id,
@@ -1790,7 +1805,6 @@ async def _execute_tool_block_impl(
return "bash (background): containment unavailable", containment.unavailable_tool_result(exc, tool="bash")
# Only this server launch may seal detached-job authority; a
# handler/bridge output carrying a job id is not a grant source.
save_background_authority(rec["id"], active_request_authority())
short = _bg_cmd.strip().split(chr(10))[0][:80]
desc = f"bash (background): {short}"
result = {
@@ -1833,6 +1847,15 @@ async def _execute_tool_block_impl(
or {"error": f"{tool}: execution failed", "exit_code": 1}
if tool == "edit_file":
desc = result.get("output") or result.get("error") or "edit_file"
elif tool in {"bash", "python"} and backend is not None and isinstance(backend.resource, NativeBackendResource):
# Native reservations are pinned to the native producer. Pass the
# application binding explicitly rather than the MCP fallback's empty
# owner/session context.
first_line = content.split(chr(10))[0][:80]
desc = f"{tool}: {first_line}"
result = await dispatched(_direct_fallback(tool, content, progress_cb=progress_cb,
owner=owner, session_id=session_id, client_runtime_context=client_runtime_context)) \
or {"error": f"{tool}: execution failed", "exit_code": 1}
elif tool in _MCP_TOOL_MAP:
first_line = content.split(chr(10))[0][:80]
desc = f"{tool}: {first_line}"
+2 -2
View File
@@ -1227,8 +1227,8 @@ async def _cookbook_kill_session(session_id: str, *, remote_host: str = "",
)
target_label = f"{session_id} on {remote}"
else:
cmd = f"tmux kill-session -t {shlex.quote(session_id)}"
target_label = session_id
return {"error": "Local Cookbook control has no admitted process resource; session discovery is not ownership",
"exit_code": 1, "blocked": True, "failure_kind": "resource_identity_denied"}
# Capture what this session owns BEFORE the kill. Once tmux tears the
# session down the pane is gone, and with it the only evidence linking a