From 8d5ff852de84c466eee2aa0c4170814d907ddcc7 Mon Sep 17 00:00:00 2001 From: Alexandre Teixeira <111787685+alteixeira20@users.noreply.github.com> Date: Fri, 2 Oct 2026 23:22:41 +0100 Subject: [PATCH] perf(runtime): bound process launch validation cost --- src/agent_runtime/process_resources.py | 95 +++++++++++- src/agent_runtime/resources.py | 30 ++-- src/agent_tools/subprocess_tools.py | 11 ++ src/bg_jobs.py | 15 +- src/process_reaper.py | 10 ++ tests/test_runtime_resource_integration.py | 12 +- tests/test_wave3_launch_cost_lifecycle.py | 170 +++++++++++++++++++++ 7 files changed, 322 insertions(+), 21 deletions(-) create mode 100644 tests/test_wave3_launch_cost_lifecycle.py diff --git a/src/agent_runtime/process_resources.py b/src/agent_runtime/process_resources.py index ee16b88c8..68856ea66 100644 --- a/src/agent_runtime/process_resources.py +++ b/src/agent_runtime/process_resources.py @@ -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,71 @@ 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 (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 and background history stay. + """ + from src import containment + from src import process_ownership + receipts = containment._load_records() + 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 process_ownership.verify(receipt.get("manager_pid"), receipt.get("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 ((receipt.get("release") or {}).get("dead") is True or 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.""" diff --git a/src/agent_runtime/resources.py b/src/agent_runtime/resources.py index 8513d087f..03006024c 100644 --- a/src/agent_runtime/resources.py +++ b/src/agent_runtime/resources.py @@ -28,7 +28,7 @@ def _absolute(value): raise ValueError("Resource path must be canonical and absolute") -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 @@ -73,8 +73,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()) @@ -83,20 +81,28 @@ 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 + 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)) + return frozenset(job_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 + return (candidate.st_dev, candidate.st_ino) in identities class FilesystemScope(str, Enum): diff --git a/src/agent_tools/subprocess_tools.py b/src/agent_tools/subprocess_tools.py index 399755caf..b6170c878 100644 --- a/src/agent_tools/subprocess_tools.py +++ b/src/agent_tools/subprocess_tools.py @@ -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 diff --git a/src/bg_jobs.py b/src/bg_jobs.py index 9a258af35..6669c091a 100644 --- a/src/bg_jobs.py +++ b/src/bg_jobs.py @@ -213,10 +213,21 @@ 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") 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() diff --git a/src/process_reaper.py b/src/process_reaper.py index 875c24ad7..f29f689ba 100644 --- a/src/process_reaper.py +++ b/src/process_reaper.py @@ -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( diff --git a/tests/test_runtime_resource_integration.py b/tests/test_runtime_resource_integration.py index 86a590937..b5d0799e6 100644 --- a/tests/test_runtime_resource_integration.py +++ b/tests/test_runtime_resource_integration.py @@ -433,6 +433,14 @@ async def test_end_to_end_fast_exit_preserves_command_result(workspace, monkeypa return {"pid": pid, "start_token": None} monkeypatch.setattr(process_ownership, "capture", mocked_capture) + # Observe the real attachment before foreground lifecycle retirement. + published = [] + attach = resources.attach_containment_processes + def observe_attachment(launch, containment_id): + attach(launch, containment_id) + published.append(json.loads(resources.launch_path(launch.generation).read_text())) + monkeypatch.setattr(resources, "attach_containment_processes", observe_attachment) + auth = authority(workspace, tool=tool) approval = approval_for(auth, tool, command) _, result = await dispatch(auth, tool, command, approval) @@ -443,10 +451,10 @@ async def test_end_to_end_fast_exit_preserves_command_result(workspace, monkeypa launches_dir = resources._LAUNCH_DIR launch_files = list(launches_dir.glob("*.json")) - assert launch_files + assert not launch_files cid = result.get("containment", {}).get("id") assert cid - matching = [json.loads(p.read_text()) for p in launch_files if json.loads(p.read_text()).get("containment_id") == cid] + matching = [record for record in published if record.get("containment_id") == cid] assert len(matching) == 1 assert matching[0]["processes"] == [] diff --git a/tests/test_wave3_launch_cost_lifecycle.py b/tests/test_wave3_launch_cost_lifecycle.py new file mode 100644 index 000000000..ef4329139 --- /dev/null +++ b/tests/test_wave3_launch_cost_lifecycle.py @@ -0,0 +1,170 @@ +"""Structural dispatch cost and exact publication lifetime regressions.""" +import asyncio +from dataclasses import replace +import json +import os +import time + +import pytest +from core.atomic_io import atomic_write_json +from src import bg_jobs, containment, process_ownership +from src.agent_runtime import resources as identities +from src.agent_runtime.authority import ExactOperation, bind_request_authority +from src.agent_runtime.resources import NativeBackendResource, ResourceIdentityError +from src.agent_tools.subprocess_tools import BashTool +from src.process_lifecycle import ProcessIdentity +from tests.test_runtime_resource_integration import workspace, authority, dispatch +from tests.test_background_resource_identity import seed +from src.agent_runtime import process_resources as resources + + +@pytest.mark.parametrize('tool,content', [('bash', 'printf guarded'), ('python', 'print("guarded")')]) +async def test_real_dispatch_scans_workspace_once_per_binding(workspace, monkeypatch, tool, content): + (workspace / 'child').mkdir() + (workspace / 'child' / 'link').symlink_to(workspace / 'child') + calls = [] + walk = os.walk + def counted(*args, **kwargs): + calls.append(args[0]) + return walk(*args, **kwargs) + monkeypatch.setattr(os, 'walk', counted) + admitted = authority(workspace, tool) + for _ in range(2): + calls.clear() + _, result = await dispatch(admitted, tool, content) + assert result['exit_code'] == 0, result + assert calls == [workspace] + assert not list(resources._LAUNCH_DIR.glob('*.json')) + + +async def test_alias_created_after_resolution_is_denied_at_binding(workspace): + admitted = authority(workspace) + op = ExactOperation.normalize('bash', 'printf safe') + bound = resources.resolve_process_operation(admitted, op, NativeBackendResource('bash')) + resources._LAUNCH_DIR.mkdir(parents=True) + state = resources._LAUNCH_DIR / ('a' * 32 + '.json') + state.write_text('{}') + (workspace / 'alias').symlink_to(state) + with bind_request_authority(admitted), pytest.raises(ResourceIdentityError): + with resources.bind_process_operation(bound): + pytest.fail('New control-plane alias admitted') + + +async def test_publication_retained_during_launch_and_retired_after_teardown(workspace, monkeypatch): + entered, resume = asyncio.Event(), asyncio.Event() + run = containment.run + paths = [] + async def held(grant, command, **kwargs): + launch = resources.active_process_operation().launch + path = resources.launch_path(launch.generation) + assert path.is_file() + paths.append(path) + entered.set() + await resume.wait() + return await run(grant, command, **kwargs) + monkeypatch.setattr(containment, 'run', held) + task = asyncio.create_task(dispatch(authority(workspace), 'bash', 'printf foreground')) + await asyncio.wait_for(entered.wait(), 5) + assert paths[0].is_file() + resume.set() + _, result = await task + assert result['exit_code'] == 0 and result['teardown']['dead'] + assert not paths[0].exists() + + +async def test_retired_publication_cannot_replay_bound_reservation(workspace): + admitted = authority(workspace) + op = ExactOperation.normalize('bash', 'printf once') + bound = resources.resolve_process_operation(admitted, op, NativeBackendResource('bash')) + from src import tool_execution + token = tool_execution._active_workspace.set(str(workspace)) + try: + with bind_request_authority(admitted), resources.bind_process_operation(bound): + ctx = {'owner': 'alice', 'session_id': 'thread'} + first = await BashTool().execute(op.input, ctx) + assert first['exit_code'] == 0 + assert not resources.launch_path(bound.launch.generation).exists() + second = await BashTool().execute(op.input, ctx) + assert second['failure_kind'] == 'resource_identity_denied' + copy = replace(bound, exact_approval=None) + with pytest.raises(ResourceIdentityError): + with resources.bind_process_operation(copy): + pytest.fail('Approval copy renewed a consumed launch') + finally: + tool_execution._active_workspace.reset(token) + + +@pytest.mark.parametrize('status,followed_up,old,removed', [ + ('running', True, True, False), ('done', False, True, False), + ('done', True, False, False), ('done', True, True, True), ('failed', True, True, True), +]) +def test_background_publication_tracks_supported_history_lifetime(workspace, status, followed_up, old, removed): + resource, rec = seed(workspace, status=status) + rec.update(followed_up=followed_up, ended_at=time.time() - (bg_jobs._RETENTION_S + 10 if old else 0)) + jobs = {'job': rec} + bg_jobs._save(jobs) + assert resources.launch_path(resource.generation).exists() + bg_jobs._prune(jobs, time.time()) + assert resources.launch_path(resource.generation).exists() is not removed + assert ('job' not in jobs) is removed + if removed: + bg_jobs._save(jobs) + with pytest.raises(ResourceIdentityError): + resources.validate_job(resource) + + +def test_old_generation_retirement_cannot_delete_replacement(workspace): + old, rec = seed(workspace, status='done') + new, _ = seed(workspace, status='done') + old_launch = identities.ProcessLaunchResource.from_dict(rec['launch_resource']) + assert resources.retire_launch(old_launch, old.containment_id, job=old) + assert resources.launch_path(new.generation).is_file() + # Even a replaced file at the old generation's slot is not deletable by old linkage. + replacement = json.loads(resources.launch_path(new.generation).read_text()) + atomic_write_json(resources.launch_path(old.generation), replacement) + assert not resources.retire_launch(old_launch, old.containment_id, job=old) + assert resources.launch_path(old.generation).is_file() + + +@pytest.mark.parametrize('manager,release,retired', [ + (process_ownership.OWNED, False, False), (process_ownership.UNVERIFIABLE, False, False), + (process_ownership.GONE, False, True), (process_ownership.FOREIGN, False, True), + (process_ownership.GONE, True, True), +]) +def test_startup_retirement_does_not_invent_process_death(workspace, monkeypatch, manager, release, retired): + admitted = authority(workspace) + launch = resources.resolve_process_operation(admitted, ExactOperation.normalize('bash', 'printf recovery'), NativeBackendResource('bash')).launch + cid = 'receipt' + resources.publish_launch(launch, admitted, cid) + atomic_write_json(containment._store_path(), {cid: {'id': cid, 'launch_generation': launch.generation, + 'manager_pid': 123, 'manager_token': 'old-manager', 'release': {'dead': release}}}) + monkeypatch.setattr(process_ownership, 'verify', lambda *args: manager) + assert resources.prune_foreground_publications() == int(retired) + assert resources.launch_path(launch.generation).exists() is not retired + assert containment._load_records()[cid]['release']['dead'] is release + + +def test_restart_never_prunes_background_linkage(workspace, monkeypatch): + resource, _ = seed(workspace, status='done') + monkeypatch.setattr(process_ownership, 'verify', lambda *args: process_ownership.GONE) + assert resources.prune_foreground_publications() == 0 + assert resources.launch_path(resource.generation).is_file() + + +def test_snapshot_is_rebuilt_for_each_guard(workspace): + resources._LAUNCH_DIR.mkdir(parents=True) + target = workspace / 'data'; target.write_text('ordinary') + (workspace / 'link').symlink_to(target) + resources.guard_launch_workspace(identities.FilesystemRoot.seal(workspace)) + os.link(target, resources._LAUNCH_DIR / ('b' * 32 + '.json')) + with pytest.raises(ResourceIdentityError): + resources.guard_launch_workspace(identities.FilesystemRoot.seal(workspace)) + + +def test_missing_receipt_publication_cannot_recover_authority(workspace): + admitted = authority(workspace) + launch = resources.resolve_process_operation(admitted, ExactOperation.normalize('bash', 'printf recovery'), NativeBackendResource('bash')).launch + resources.publish_launch(launch, admitted, 'missing-receipt') + assert resources.prune_foreground_publications() == 1 + assert not resources.launch_path(launch.generation).exists() + assert not containment._load_records()