From 4a4d509a9769fe795b7ffaef50dbf0fa869a5bbf Mon Sep 17 00:00:00 2001 From: Zhang Handuo Date: Mon, 5 Oct 2026 21:36:59 +0800 Subject: [PATCH 1/2] fix(spill): one store and one scope key for spill storage, mount and reads Handoff item 3 (Harness #494 / #521, recovery and isolation only): - Persistence goes through AgentCore's SpillStore; Frontier's hand-written writer/reader/registry is gone, so there is one store and one created-stores registry. - current_store_scope() (task:llm_session, or a per-pid unscoped key) is the single identity: the store directory, the per-command bwrap mount (only that scope and its sub-agents under /spill) and _path_auth read authorization all derive from it. Siblings and other conversations no longer see each other's stores; the shared root is never mounted. - _visible_root() follows the sandbox actually running commands, then config (auto+E2B and unusable bwrap name nothing). maybe_overflow writes nothing then, and default_compaction_spill() withholds the compaction callback. - Compacting a preview this module produced returns the ref of the full body instead of storing the preview as "[Full text]"; re-compaction keeps the ref. - The recover_result footer is suppressed when no trajectory JSONL exists. - Symlinked spill roots fall back to one private directory; stateful runs clean their scope tree at the end; cleanup only touches stores this process created. Container mode without the inner bwrap jail remains unisolated (shared tool uid); documented in spill_root(). Co-Authored-By: Claude Opus 5.5 (1M context) --- CHANGELOG.md | 13 + .../core/runtime/loop/agent_loop.py | 16 +- plugins/tools/_overflow.py | 495 ++++++++++-------- plugins/tools/_path_auth.py | 38 +- plugins/tools/_sandbox.py | 100 ++-- plugins/tools/recover_result.py | 10 + tests/test_agent_team_workflow.py | 68 ++- tests/test_long_run_compaction.py | 13 +- tests/test_path_authorization.py | 46 +- tests/test_site3_recovery.py | 29 +- tests/test_spill_scope_recovery.py | 317 +++++++++++ tests/test_tool_result_truncation.py | 4 +- workflows/agent_team/nodes/main_agent.py | 4 +- workflows/agent_team/subagent_runtime.py | 4 +- .../stateful_react_agent/nodes/main_agent.py | 12 +- 15 files changed, 840 insertions(+), 329 deletions(-) create mode 100644 tests/test_spill_scope_recovery.py diff --git a/CHANGELOG.md b/CHANGELOG.md index 1255b4d..08f0146 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -62,6 +62,19 @@ Initial open-source release of FrontierAgent. ### Fixed +- Spill recovery: oversized and compacted tool results are stored through + AgentCore's `SpillStore` (one store, one registry). A single scope key + (`task:llm_session`) now decides the store directory, the bwrap `/spill` mount + and read authorization: a jail sees only its own store and its sub-agents', + siblings and other conversations see nothing, and `_path_auth` no longer + authorizes every store the process created. A path is advertised only when + the backend actually running commands can open it (`auto` going to E2B, or + bwrap unusable, gets none), and the compaction spill callback is withheld + then. Compacting a truncated preview points at the original full body + instead of storing the preview as "[Full text]". The `recover_result` footer + is only shown when a trajectory JSONL exists. Known limitation: `container` + mode without the inner bwrap jail shares one tool uid, so model commands can + still read other scopes there. - Surface finalize-gate bypasses on the final turn: an answer delivered despite open task-board items now carries an unfinished-work note and a `finalize_gate_bypassed` marker instead of reading as a clean success. diff --git a/frontier_agent/core/runtime/loop/agent_loop.py b/frontier_agent/core/runtime/loop/agent_loop.py index cf80bc5..c36114f 100644 --- a/frontier_agent/core/runtime/loop/agent_loop.py +++ b/frontier_agent/core/runtime/loop/agent_loop.py @@ -21,6 +21,7 @@ from frontier_agent.core.execution_context import ( ExecutionScope, chain_fallback_active, + get_current_execution_scope, reset_current_execution_scope, set_current_execution_scope, ) @@ -49,13 +50,26 @@ def _enter_scope( phase_id=phase_id, metadata=metadata, ) + # A loop entered from inside another (an in-process sub-agent) is that + # loop's child: the parent may read the child's spill store, siblings not. + from plugins.tools._overflow import register_child_scope + + register_child_scope(get_current_execution_scope(), scope) return scope, set_current_execution_scope(scope) def _body_has_spill_reference(body: str) -> bool: + """Whether the trajectory recovery footer should stay quiet for ``body``. + + True when the body already names a spill file (that file holds the full + output), and also when this agent has no trajectory JSONL: the footer names + ``recover_result``, which reads exactly that file, so without one it would + hand the model a handle that can only answer "unavailable". + """ from plugins.tools._overflow import body_names_a_spill_file + from plugins.tools.recover_result import trajectory_recovery_available - return body_names_a_spill_file(body) + return body_names_a_spill_file(body) or not trajectory_recovery_available() def _with_recovery_handle(body: str, result: Any, turn: int, *, enabled: bool) -> str: diff --git a/plugins/tools/_overflow.py b/plugins/tools/_overflow.py index 32282a5..6166e01 100644 --- a/plugins/tools/_overflow.py +++ b/plugins/tools/_overflow.py @@ -4,66 +4,121 @@ import contextlib import hashlib -import json import logging import os -import time -import uuid +import threading +from collections import OrderedDict +from collections.abc import Callable from pathlib import Path +from agent_core.runtime.spill import SpillStore, scope_component + from plugins.tools.meta import get_tool_meta logger = logging.getLogger(__name__) -# The store's physical root comes from ``_sandbox.spill_root()``; see there for -# why it lives outside every root the agent can write. What is left here is the -# per-conversation partitioning under it. -_RUN_SUBDIR = "spill" +# Persistence is AgentCore's ``SpillStore`` — one implementation, one registry +# of created stores — so nothing here writes, reads or deletes spill files by +# hand any more. What stays product-side is the policy around it: +# +# * the scope KEY (:func:`current_store_scope`), which is the one identity the +# store directory, the bwrap mount (``spill_bind_args``) and read +# authorization (``_path_auth``) all derive from; +# * the VISIBLE root (:func:`_visible_root`), i.e. whether the backend that +# actually runs model commands can name the store at all; +# * preview shaping and the footer, through AgentCore's ``budgeted_preview``. _SPILL_SEPARATOR = "\n---\n\n" -# Backends whose commands always run on this process's own filesystem, so the -# physical path IS the path a model command can name. Container mode normally -# does too, but may opt into an inner bwrap jail; that case is resolved at run -# time in :func:`_overflow_dir`. -_SAME_FILESYSTEM_BACKENDS = frozenset({"native"}) -# Backends whose commands run on another machine entirely. Nothing on this -# filesystem is nameable there, so spill advertises no path and the footer says -# the remainder is unreadable rather than pointing somewhere that cannot resolve. -_REMOTE_BACKENDS = frozenset({"e2b"}) -#: Stores this process created, so a discarded session can drop exactly its own -#: recovery files. Replaces walking a directory tree looking for them: we delete -#: only paths we made, which is why this needs no symlink or filesystem-root -#: defence — the previous implementation walked a tree inside the agent's own -#: workspace and had to assume it was hostile. -_created_stores: set[Path] = set() - - -def _scope_component(task_id: str) -> str: - """Map an arbitrary task id to one safe, stable directory component.""" - return hashlib.sha256(task_id.encode("utf-8")).hexdigest()[:16] if task_id else "" +#: Unscoped writes (no ExecutionScope) go to one store per PROCESS. The pid is +#: part of the key because the default root is shared by every process of a +#: uid: a constant key would merge two processes' unscoped stores, and either +#: one's cleanup would delete the other's files. +UNSCOPED_PROCESS_STORE = "__process__" +#: Registry of the stores this process created — AgentCore's, shared by +#: reference so a reset in tests affects the store's own bookkeeping too. +_created_stores: set[Path] = SpillStore._created_stores +#: parent scope key -> child scope keys entered from inside it (in-process +#: sub-agents). A parent may read its children's stores — a fan-in report can +#: carry a child's spill path back — but siblings never see each other's. +_child_scopes: dict[str, set[str]] = {} +_child_lock = threading.Lock() +#: digest(preview text) -> ref of the FULL body behind it. Lets a later +#: compaction of that preview point at the original full text instead of +#: storing the preview and labelling it "[Full text]". Bounded; exact-match only, +#: so an agent echoing part of a preview cannot claim someone else's ref. +_preview_refs: OrderedDict[str, str] = OrderedDict() +_PREVIEW_REFS_MAX = 4_096 + + +def spill_scope_key(task_id: str, llm_session_id: str = "") -> str: + """The store key for one conversation: ``task:session`` (or ``task``). + + Built only from stable ids — never ``metadata["session_id"]``, which trace + plumbing may replace mid-run. + """ + task_id = str(task_id or "") + session = str(llm_session_id or "") + if not task_id: + return "" + return f"{task_id}:{session}" if session else task_id + + +def _scope_key_of(scope: object | None) -> str: + if scope is None: + return "" + metadata = getattr(scope, "metadata", None) or {} + return spill_scope_key( + str(getattr(scope, "task_id", "") or ""), + str(metadata.get("llm_session_id") or ""), + ) def _current_task_id() -> str: + """The current ExecutionScope's store key, or "" outside any scope.""" from frontier_agent.core.execution_context import get_current_execution_scope - scope = get_current_execution_scope() - if scope is None: - return "" - task_id = str(scope.task_id or "") - session_id = str(scope.metadata.get("llm_session_id") or "") - # This is a per-conversation overflow cache, not cross-session memory. - return f"{task_id}:{session_id}" if session_id else task_id + return _scope_key_of(get_current_execution_scope()) + + +def current_store_scope() -> str: + """The ONE scope key writers, the mount and read authorization share.""" + return _current_task_id() or f"{UNSCOPED_PROCESS_STORE}:{os.getpid()}" + + +def register_child_scope(parent: object | None, child: object | None) -> None: + """Record that ``child``'s loop was entered from inside ``parent``'s.""" + parent_key, child_key = _scope_key_of(parent), _scope_key_of(child) + if parent_key and child_key and parent_key != child_key: + with _child_lock: + _child_scopes.setdefault(parent_key, set()).add(child_key) + + +def readable_scopes(scope: str | None = None) -> list[str]: + """Scope keys an agent may read: its own and its descendants'.""" + root_key = current_store_scope() if scope is None else scope + out: list[str] = [] + pending = [root_key] + with _child_lock: + while pending: + key = pending.pop() + if not key or key in out: + continue + out.append(key) + pending.extend(_child_scopes.get(key, ())) + return out + + +def _physical_root() -> Path: + from plugins.tools._sandbox import spill_root + + return spill_root().expanduser().resolve() def _resolved_backend() -> str: """The active backend as ``_sandbox`` resolves it, or "" when unresolvable. Reading ``SANDBOX_BACKEND`` straight from the environment misses a backend - supplied only through ``config.yaml`` — the resolver consults ``get_config()`` - for exactly that case — and would then advertise ``/workspace/.spill/...`` - for a run whose commands execute directly on the host filesystem, where that - literal path may name an unrelated directory. A misconfigured backend must - not take spill down with it: spill is a diagnostic aid, so fall back to the - conservative canonical mount rather than raising. + supplied only through ``config.yaml``. A misconfigured backend must not take + spill down with it, so an error resolves to "" (decided by bwrap below). """ from plugins.tools._sandbox import _get_sandbox_backend @@ -73,52 +128,103 @@ def _resolved_backend() -> str: return "" -def _overflow_dir(task_id: str = "", *, create: bool = True) -> tuple[Path, str]: - """Return the physical write directory and the path visible to the agent. - - ``create=False`` resolves the same location without touching the filesystem, - for callers that only want to inspect or remove an existing store. +def _visible_root() -> str | None: + """How model commands can name the store, or ``None`` when they cannot. - This used to branch four ways over the workspace, the configured mount dir, - a run directory and a legacy host-only path, each with its own rule for what - the agent could name. The store now has one root outside every write root, so - the only remaining question is how a model command reaches it: directly under - ``native`` and ordinary ``container`` mode, through the read-only ``/spill`` - mount under bwrap, and not at all from a remote backend. + Decided by the sandbox that ACTUALLY runs commands when one exists, then by + configuration: ``auto`` with an E2B key goes off-host, and a bwrap backend on + a host where bwrap cannot run has no ``/spill`` mount. Advertising a path in + either case would hand the model a ref it can never open — and would make + the compaction callback promise a recovery that does not exist. """ - from plugins.tools._sandbox import _DEFAULT_SPILL_DIR, spill_root - - scope = _scope_component(task_id) - target = spill_root() / scope if scope else spill_root() - if create: - target.mkdir(parents=True, exist_ok=True) - # Readable and traversable by others, writable only by the harness. Set - # explicitly because this is the WHOLE enforcement under ``container``, - # where model commands are dropped to an unprivileged uid: they may read - # a 0644 spill file but cannot create or unlink inside a directory they - # do not own. Leaving it to the ambient umask would make that guarantee - # depend on whoever launched the process. bwrap gets a read-only mount - # instead (uid 0 in a user namespace ignores DAC); ``native`` has no - # isolation to enforce anything with. - with contextlib.suppress(OSError): - target.chmod(0o755) - _created_stores.add(target) + from plugins.tools import _sandbox as sb + + live = sb.get_existing_sandbox() + if isinstance(live, sb.BwrapSandbox): + return sb._DEFAULT_SPILL_DIR + if isinstance(live, sb.CurrentSandbox): + return sb._DEFAULT_SPILL_DIR if live._inner is not None else str(_physical_root()) + if live is not None and type(live).__module__.split(".", 1)[0].startswith("e2b"): + return None backend = _resolved_backend() - if backend in _REMOTE_BACKENDS: - return target, "" - if backend in _SAME_FILESYSTEM_BACKENDS: - return target, str(target) + if backend == "e2b": + return None + if backend == "native": + return str(_physical_root()) if backend == "container": - # Production CurrentSandbox explicitly runs without a mount namespace; - # advertising /spill there points at nothing because the physical store - # is normally under /tmp or the run directory. Only the optional inner - # bwrap path creates the canonical read-only /spill mount. - from plugins.tools._sandbox import container_uses_inner_bwrap + return sb._DEFAULT_SPILL_DIR if sb.container_uses_inner_bwrap() else str(_physical_root()) + if backend == "auto": + try: + use_e2b = sb._resolve_use_e2b()[0] + except Exception: + use_e2b = False + if use_e2b: + return None + return sb._DEFAULT_SPILL_DIR if sb.bwrap_available() else None + + +def _store(scope: str | None = None) -> SpillStore: + """The SpillStore for ``scope`` (default: the current one).""" + return SpillStore( + _physical_root(), + current_store_scope() if scope is None else scope, + visible_root=_visible_root(), + ) + - if not container_uses_inner_bwrap(): - return target, str(target) - return target, f"{_DEFAULT_SPILL_DIR}/{scope}" if scope else _DEFAULT_SPILL_DIR +def spill_is_recoverable() -> bool: + """Whether a compaction spill would produce a ref the agent can open.""" + return _visible_root() is not None + + +def default_compaction_spill() -> Callable[[str, str], str | None] | None: + """The compaction spill callback, or ``None`` when nothing is recoverable. + + AgentCore cannot introspect a plain function, so it takes any callback as a + promise that discarded bodies stay recoverable and drops its guard that + keeps the newest unseen results verbatim. Withholding the callback is how a + backend without a readable store keeps that guard. + """ + return spill_compacted_body if spill_is_recoverable() else None + + +def readable_store_dirs() -> list[Path]: + """Physical directories of :func:`readable_scopes` this process created.""" + root = _physical_root() + dirs: list[Path] = [] + for key in readable_scopes(): + directory = root / scope_component(key) + if directory in _created_stores and directory.is_dir(): + dirs.append(directory) + return dirs + + +def spill_bind_args() -> list[str]: + """bwrap args mounting ONLY the readable scopes under ``/spill``. + + Built per command from the current scope, so a shared or pooled jail still + shows each agent just its own store (and its sub-agents'). ``--ro-bind-try`` + because a store is created lazily on the first spill. + """ + from plugins.tools._sandbox import _DEFAULT_SPILL_DIR + + root = _physical_root() + args = ["--dir", _DEFAULT_SPILL_DIR] + for key in readable_scopes(): + component = scope_component(key) + args.extend([ + "--ro-bind-try", str(root / component), f"{_DEFAULT_SPILL_DIR}/{component}", + ]) + return args + + +def _overflow_dir(task_id: str = "", *, create: bool = True) -> tuple[Path, str]: + """Physical directory and agent-visible path of one scope's store.""" + store = _store(task_id or current_store_scope()) + if create: + store.ensure() + return store.directory, store.visible_directory def body_names_a_spill_file(body: str) -> bool: @@ -126,19 +232,8 @@ def body_names_a_spill_file(body: str) -> bool: A presence test against the two roots the store can be named by — the canonical mount and the physical path — not a parse of the pointer's prose. - The wording differs per backend and per caller, and ``7cf9188`` moved - deliberately away from recognising spill refs by shape. - - Deliberately does NOT go through :func:`agent_visible_spill_dir`, which - resolves via ``_overflow_dir`` and would CREATE the store as a side effect of - asking a read-only question. - - Exists so the site-3 recovery footer can stay quiet when it would be - redundant. Measured on a live agent-team run: every result site 3 shortened - was a ``bash`` result that already carried a spill pointer, and the spill file - behind it held the FULL pre-gate-① output — 42,770 chars against the 8,000 the - model saw. The agent read those files with ``cat`` and never called the tool - the footer named. Two routes to the same bytes, and the footer lost. + Exists so the trajectory recovery footer can stay quiet when it would be + redundant: the spill file already holds the full output. """ if not body: return False @@ -148,8 +243,7 @@ def body_names_a_spill_file(body: str) -> bool: with contextlib.suppress(Exception): roots.append(str(spill_root())) # The separator is not cosmetic: a bare ``"/spill" in body`` also fires on - # ``/spillover``, and a pointer always names a FILE under the store, so the - # trailing slash is both stricter and exactly what a real pointer contains. + # ``/spillover``, and a pointer always names a FILE under the store. return any( root and f"{root.rstrip('/')}/" in body for root in roots ) @@ -157,7 +251,7 @@ def body_names_a_spill_file(body: str) -> bool: def agent_visible_spill_dir() -> str: """Return the spill directory as tools should name it, or empty if unreadable.""" - return _overflow_dir(_current_task_id())[1] + return _overflow_dir(current_store_scope())[1] # Smallest inline preview worth keeping. A cap tighter than @@ -284,52 +378,48 @@ def budgeted_preview( def _write_spill( tool_name: str, body: str, *, require_visible: bool, task_id: str = "", ) -> tuple[Path, str] | None: - """Persist ``body`` once and return its path plus the agent-visible ref. - - Named by ``sha256(tool_name, body)`` rather than a fresh uuid so re-spilling - the SAME body is idempotent. Tier 2 re-spills every protected fan-in result - on each pass that wins, which under a uuid name left one identical copy per - compaction on disk and burned a manifest slot each time. + """Persist ``body`` once through the scope's SpillStore. - Returns ``None`` when the store is unreachable, or when the caller requires - an agent-visible path and this backend cannot name one. + Content-addressed (``sha256(tool_name, body)``) so re-spilling the SAME body + is idempotent. Returns ``None`` when the store is unreachable, or when the + caller requires an agent-visible path and this backend cannot name one. """ - write_dir, visible_dir = _overflow_dir(task_id or _current_task_id()) - if require_visible and not visible_dir: + try: + store = _store(task_id or current_store_scope()) + return store.write(tool_name, body, require_visible=require_visible) + except OSError as exc: + logger.warning("Failed to spill %s result: %s", tool_name, exc) return None - digest = hashlib.sha256( - tool_name.encode("utf-8") + b"\x00" + body.encode("utf-8", "replace"), - ).hexdigest()[:16] - path = write_dir / f"{digest}.md" - if not path.exists(): - # Write-then-rename: a crash mid-write would otherwise leave a truncated - # file under the name the digest resolves to, which ``path.exists()`` - # then treats as a complete spill forever. - tmp = path.with_name(f".{digest}.{uuid.uuid4().hex[:8]}.tmp") - try: - tmp.write_text(_spill_document(tool_name, digest, body), encoding="utf-8") - os.replace(tmp, path) - except OSError as exc: - logger.warning("Failed to spill %s result: %s", tool_name, exc) - with contextlib.suppress(OSError): - tmp.unlink() - return None - return path, (f"{visible_dir}/{path.name}" if visible_dir else "") -def _spill_document(tool_name: str, spill_id: str, result: str) -> str: - """Use grep-friendly markdown and keep the captured body verbatim. +def _preview_key(text: str) -> str: + return hashlib.sha256(text.encode("utf-8", "replace")).hexdigest() - ``spill_id`` is the content digest, not a call id: the same body spilled - twice is one file, so no single call owns it. - """ - return ( - f"# {tool_name} — spilled tool result\n\n" - f"- id: `{spill_id}`\n" - f"- captured: {time.strftime('%Y-%m-%dT%H:%M:%S%z')}\n" - f"- length: {len(result):,} chars" - f"{_SPILL_SEPARATOR}{result}" - ) + +def remember_preview(preview: str, ref: str) -> None: + """Record that ``preview`` stands for the full body stored at ``ref``.""" + if not preview or not ref: + return + key = _preview_key(preview) + with _child_lock: + _preview_refs[key] = ref + _preview_refs.move_to_end(key) + while len(_preview_refs) > _PREVIEW_REFS_MAX: + _preview_refs.popitem(last=False) + + +def _full_body_ref(body: str) -> str: + """The ref of the full body ``body`` is a preview of, if it is still readable.""" + with _child_lock: + ref = _preview_refs.get(_preview_key(body), "") + if not ref: + return "" + store = _store() + for key in readable_scopes(): + candidate = SpillStore(store.root, key, visible_root=store.visible_root) + if candidate.contains_path(ref) and candidate.read(ref) is not None: + return ref + return "" def _spill_footer(ref: str, *, full_len: int, note: str = "") -> str: @@ -360,9 +450,17 @@ def _spill_footer(ref: str, *, full_len: int, note: str = "") -> str: def spill_compacted_body(tool_name: str, body: str) -> str | None: - """Persist a result immediately before compaction discards its inline body.""" + """Persist a result immediately before compaction discards its inline body. + + When ``body`` is itself a preview this module produced, the ref of the full + body behind it is returned instead: storing the preview would label a cut + copy "[Full text]" and orphan the real one. + """ if len(body) < _SPILL_MIN_CHARS: return None + existing = _full_body_ref(body) + if existing: + return existing spilled = _write_spill(tool_name, body, require_visible=True) return spilled[1] if spilled else None @@ -379,17 +477,17 @@ def maybe_overflow( Args: tool_name: Name of the tool that produced the result. result: The full tool result string. - task_id: Optional task ID for organizing overflow files. Defaults to - the current execution scope. Pass the SAME composite form the scope - uses (``f"{task_id}:{llm_session_id}"``) or the store will not be - the one ``cleanup_overflow`` removes. + task_id: Optional store key. Defaults to :func:`current_store_scope`; + pass the SAME ``spill_scope_key`` form or the store will not be the + one the mount and read authorization expose. call_id: Accepted for backwards compatibility and no longer used to name - the file — see :func:`_write_spill` on content-hash naming. + the file — the store is content-addressed. Returns: The original result if within limits, or a head-and-tail preview with a reference to the overflow file, together no longer than the cap. """ + del call_id meta = get_tool_meta(tool_name) # 0 means no limit @@ -399,8 +497,10 @@ def maybe_overflow( if len(result) <= meta.max_result_chars: return result + # ``require_visible``: a backend that cannot name the store (E2B) gets no + # file at all — one written on this host would be unreachable forever. spilled = _write_spill( - tool_name, result, require_visible=False, task_id=task_id, + tool_name, result, require_visible=True, task_id=task_id, ) if spilled is not None: logger.info( @@ -410,9 +510,11 @@ def maybe_overflow( # A failed write only costs the pointer: the footer then says the remainder # is unreadable instead of naming a path that does not exist. ref = spilled[1] if spilled else "" - return budgeted_preview( + preview = budgeted_preview( result, cap=meta.max_result_chars, ref=ref, tool_name=tool_name, ) + remember_preview(preview, ref) + return preview # ── Aggregate budget (per-turn total) ─────────────────────────────────── @@ -480,6 +582,9 @@ def check_aggregate_budget( note="Cut further to fit the per-turn tool-result budget.", tool_name=name, ) + # A re-cut preview still stands for the ORIGINAL full body when the + # input was already one of ours. + remember_preview(replacement, _full_body_ref(result) or (spilled[1] if spilled else "")) total -= len(result) - len(replacement) adjusted[idx] = replacement @@ -487,28 +592,18 @@ def check_aggregate_budget( def get_overflow_content(overflow_path: str) -> str | None: - """Read the full content from an overflow file. - - Args: - overflow_path: Path to the overflow JSON file. + """Read the full content from a spill file this conversation may read. - Returns: - The full tool result content, or None if not found. + Accepts a physical path or the agent-visible ref. Only the current scope's + store, its sub-agents' stores, or (for a physical path) a store this + process created are consulted — never an arbitrary file. """ - path = Path(overflow_path) - if not path.is_file(): - return None - - try: - text = path.read_text(encoding="utf-8") - if _SPILL_SEPARATOR in text: - return text.split(_SPILL_SEPARATOR, 1)[1] - # Backward compatibility for sessions holding pointers to old JSON spills. - data = json.loads(text) - return data.get("content") - except Exception as e: - logger.warning("Failed to read overflow file %s: %s", overflow_path, e) - return None + store = _store() + for key in readable_scopes(): + candidate = SpillStore(store.root, key, visible_root=store.visible_root) + if candidate.contains_path(overflow_path): + return candidate.read(overflow_path) + return SpillStore.read_created(overflow_path) def cleanup_overflow( @@ -519,15 +614,10 @@ def cleanup_overflow( """Remove the spilled tool results of one finished conversation. Args: - scope: The store to remove, in the SAME composite form the writers use — - ``f"{task_id}:{llm_session_id}"``, which is what - :func:`_current_task_id` returns. A bare ``task_id`` hashes to a - different directory and would silently match nothing, since - ``llm_session_id`` defaults to ``task_id`` rather than staying empty. - Omit it to clean up the caller's own current scope. An explicitly - empty string is always a safe no-op. + scope: The store key, in :func:`spill_scope_key` form. Omit it to clean + up the caller's own current scope. An explicitly empty string is + always a safe no-op. workspace: Ignored, kept so existing teardown calls still type-check. - The store no longer lives under a workspace. Returns: Number of files removed. @@ -535,48 +625,41 @@ def cleanup_overflow( del workspace if scope == "": return 0 - resolved_scope = _current_task_id() if scope is None else scope + resolved_scope = current_store_scope() if scope is None else scope if not resolved_scope: return 0 + removed = SpillStore(_physical_root(), resolved_scope).cleanup() + with _child_lock: + _child_scopes.pop(resolved_scope, None) + return removed - # ``create=False``: resolving a store in order to delete it must not first - # bring it into existence, which would also leave a stray empty directory - # behind for any scope that never spilled. - return _remove_store(_overflow_dir(resolved_scope, create=False)[0]) +def cleanup_overflow_tree(scope: str | None = None) -> int: + """Remove one conversation's store and every sub-agent store under it. -def _remove_store(store: Path) -> int: - """Delete one store directory's files, then the directory. Count the files.""" - if not store.is_dir(): - return 0 - count = 0 - for entry in store.iterdir(): - try: - entry.unlink() - count += 1 - except OSError: - pass - with contextlib.suppress(OSError): - store.rmdir() - _created_stores.discard(store) - return count + Only stores this process created are touched, so another process — or an + unrelated session in this one — keeps its files. + """ + keys = readable_scopes(scope) + root = _physical_root() + removed = 0 + for key in keys: + if root / scope_component(key) in _created_stores: + removed += SpillStore(root, key).cleanup() + with _child_lock: + _child_scopes.pop(key, None) + return removed def cleanup_overflow_process() -> int: """Remove every store THIS process created. - For a discarded conversation: the TUI's ``/clear`` and ``/mode`` drop all - history that could reference a spill path, so the files it named are dead. - - PRECONDITION: one conversation per process, which holds for the only caller — - the terminal app — and for the benchmark runner's subprocess-per-question. A - server multiplexing concurrent sessions in one process must NOT use this; it - would delete a live session's recovery files. Such a caller wants - :func:`cleanup_overflow` per scope instead. - Scoped to what this process created rather than to a directory tree, which - is both safer — another session's store is not ours to delete, whatever it - is next to — and simpler, since deleting only paths we made needs none of - the symlink and filesystem-root defences that walking the agent's workspace - required. + For a discarded conversation in a one-conversation process: the TUI's + ``/clear`` and ``/mode`` drop all history that could reference a spill path. + A server multiplexing concurrent sessions in one process must use + :func:`cleanup_overflow_tree` per conversation instead. """ - return sum(_remove_store(store) for store in list(_created_stores)) + with _child_lock: + _child_scopes.clear() + _preview_refs.clear() + return SpillStore.cleanup_process() diff --git a/plugins/tools/_path_auth.py b/plugins/tools/_path_auth.py index 31041e4..277e46f 100644 --- a/plugins/tools/_path_auth.py +++ b/plugins/tools/_path_auth.py @@ -253,39 +253,21 @@ def _resolve_spill_dirs() -> list[Path]: Authorized for READ so ``read_file`` / ``grep_search`` can recover a body compaction dropped, and never for write — the same shape as ``/inputs``. The - canonical ``/spill`` path a model sees is rewritten to this by - ``resolve_runtime_path`` before it reaches here. Gating on existence means a - run that never spilled adds no prefix. Imported lazily to avoid an import - cycle with ``_sandbox``. + canonical ``/spill`` path a model sees is rewritten to the physical root by + ``resolve_runtime_path`` before it reaches here. Imported lazily to avoid an + import cycle with ``_sandbox``. + + The set is exactly what ``spill_bind_args`` mounts into a bwrap jail: the + current scope's store plus its own sub-agents' stores, each only if this + process created it. A sibling sub-agent's store, another conversation in + this process, or another process's store matches none of them. """ try: - from plugins.tools._overflow import _created_stores, _current_task_id - from plugins.tools._overflow import _scope_component as scope_of - from plugins.tools._sandbox import spill_root + from plugins.tools._overflow import readable_store_dirs - root = spill_root() + return readable_store_dirs() except Exception: return [] - if not root.is_dir(): - return [] - - # Narrower than the root on purpose. The root is shared — a temp directory, - # or a run directory — so authorizing it would let one conversation read - # another's spilled tool results, which the old in-workspace layout made - # impossible. Two things are authorized instead: - # - # * this conversation's own scope, which is what its recovery index names; - # * every store THIS process created, because in-process sub-agents spill - # under their own scope and a fan-in report can carry one of those paths - # back to the parent. - # - # A different session in a different process matches neither. - allowed: list[Path] = [] - scope = scope_of(_current_task_id()) - if scope and (root / scope).is_dir(): - allowed.append(root / scope) - allowed.extend(store for store in _created_stores if store.is_dir()) - return allowed def _allowed_local_prefixes( diff --git a/plugins/tools/_sandbox.py b/plugins/tools/_sandbox.py index 5e23b86..1eeb40d 100644 --- a/plugins/tools/_sandbox.py +++ b/plugins/tools/_sandbox.py @@ -863,11 +863,6 @@ def bwrap_available() -> bool: # ``_venv_bind_args``. ``/proc`` is deliberately NOT here — see ``--proc`` below. _BWRAP_SYSTEM_PATHS = ("/usr", "/bin", "/lib", "/lib64", "/etc") -# Recovery store below the workspace. Kept in sync with ``_overflow``'s -# ``_WORKSPACE_SUBDIR`` by name rather than by import: ``_overflow`` imports this -# module (lazily, from inside its functions), so a module-level import back the -# other way would be a cycle. -_SPILL_SUBDIR = ".spill" def _interpreter_bind_args() -> list[str]: @@ -1008,6 +1003,27 @@ def _bwrap_mem_limit_mb() -> int: return _mem_limit_mb("sandbox_bwrap_mem_mb", "SANDBOX_BWRAP_MEM_MB", 12 * 1024) +def _spill_mount_args() -> list[str]: + """Mount the CURRENT scope's spill stores read-only under ``/spill``. + + Resolved per command, not when the jail is built: one jail (the shared + singleton, a pool lease, the container's inner jail) serves many agents, and + each must see only its own store and its sub-agents' — the same set + ``_path_auth`` authorizes, from the same scope key the writer used. + + Read-only because model commands run as uid 0 inside the user namespace, so + file modes are not consulted; a mount is. The harness writes spill through + the host filesystem, so this constrains model commands only. + """ + try: + from plugins.tools._overflow import spill_bind_args + + return spill_bind_args() + except Exception as exc: + logger.warning("spill mount unavailable for this command: %s", exc) + return [] + + class _BwrapCommands: """Command executor that wraps each command in a bubblewrap sandbox.""" @@ -1043,6 +1059,7 @@ def run(self, command: str, timeout: int = 60, _BWRAP_PATH, *base_args, *self._bind_args, + *_spill_mount_args(), "--chdir", self._chdir, "--", @@ -1173,32 +1190,6 @@ def __init__( + [str(src_path), str(dst)] ) - # Mount the spill store read-only at its own top-level path. - # - # Without a mount the read-only store is only a lexical promise: - # ``_deliverable_policy`` refuses every natural way to write there — - # including bash, which IS token-scanned — but shell expansion can hide a - # path from any scanner (a glob, a brace, a ``$VAR`` assembled in pieces, - # a ``$(…)`` substitution), and the files are ordinary 0644 files. File - # modes cannot close it either — model commands run as uid 0 inside the - # user namespace, so DAC is not consulted. A mount is, which is why this - # is the layer that actually holds. - # - # This no longer has to be the LAST bind to win: the source now sits - # outside the workspace, so nothing above it overlaps and the ordering - # constraint that used to be load-bearing is gone. - # - # ``--ro-bind-try``, not ``--ro-bind``: the store is created lazily on the - # first spill, and bwrap aborts the whole jail when a ``--ro-bind`` source - # is missing. Args are rebuilt per command, so the mount appears as soon - # as the directory does — and until then there is nothing to protect. - # - # The harness writes spill through the host filesystem, not through the - # jail, so this constrains model commands only. - bind_args.extend([ - "--ro-bind-try", str(spill_root()), _DEFAULT_SPILL_DIR, - ]) - self._workdir = str(workspace_path) self._tmpdir = str(tmp_path) self.commands = _BwrapCommands( @@ -2365,18 +2356,53 @@ def spill_root() -> Path: """ explicit = os.environ.get(_SPILL_DIR_ENV, "").strip() if explicit: - return Path(explicit) + return _refuse_symlinked_spill_root(Path(explicit)) run_dir = os.environ.get("APODEX_RUN_DIR", "").strip() if run_dir: - return Path(run_dir) / "spill" + return _refuse_symlinked_spill_root(Path(run_dir) / "spill") # The uid is in the NAME so two accounts on one host do not share a store. # It is not a permission boundary and cannot be: under ``container`` the # model runs as a DIFFERENT uid than the harness and has to read these # files, so the directory must stay traversable by others (0755). What keeps - # one conversation out of another's recovery files is ``_path_auth``, which - # authorizes only the current scope plus the stores this process created — - # a local user with their own shell is outside that model either way. - return Path(tempfile.gettempdir()) / f"apodex-spill-{os.getuid()}" + # one conversation out of another's recovery files is the scope key + # (``_overflow.current_store_scope``): bwrap mounts only that scope's store + # and its sub-agents', and ``_path_auth`` authorizes the same set. + # + # KNOWN LIMITATION — ``container`` without the inner bwrap jail: every + # agent's commands run as the same unprivileged tool uid with no mount + # namespace, so a model command can ``cat`` any scope under this root. + # File modes cannot separate agents that share a uid; closing it needs the + # inner jail (``FRONTIER_AGENT_CONTAINER_INNER_BWRAP``). In-process file + # tools are still scope-checked by ``_path_auth``. + return _refuse_symlinked_spill_root( + Path(tempfile.gettempdir()) / f"apodex-spill-{os.getuid()}", + ) + + +_private_spill_root: Path | None = None +_private_spill_lock = threading.Lock() + + +def _refuse_symlinked_spill_root(path: Path) -> Path: + """``path``, unless it is a symlink — then one private temp directory. + + Every user of the root resolves it, so a symlink planted at a shared + location (``/tmp/apodex-spill-``) would redirect the store, the + ``/spill`` mount and read authorization to wherever it points. Falling back + to a private directory keeps spill working instead of failing the run; it is + created once so the writer, the mount and the reader still agree. + """ + global _private_spill_root + if not path.is_symlink(): + return path + with _private_spill_lock: + if _private_spill_root is None: + _private_spill_root = Path(tempfile.mkdtemp(prefix="apodex-spill-")) + logger.warning( + "spill root %s is a symlink; using private %s instead", + path, _private_spill_root, + ) + return _private_spill_root def is_spill_path(path: str) -> bool: diff --git a/plugins/tools/recover_result.py b/plugins/tools/recover_result.py index 407e881..e846289 100644 --- a/plugins/tools/recover_result.py +++ b/plugins/tools/recover_result.py @@ -69,6 +69,16 @@ def _trajectory_path() -> Path | None: return Path(raw) if raw else None +def trajectory_recovery_available() -> bool: + """Whether a handle minted now could ever resolve. + + The trajectory observer advertises this agent's JSONL at loop start and + withdraws it when JSONL output is off; without it every ``recover_result`` + call answers "unavailable", so the loop must not advertise one. + """ + return _trajectory_path() is not None + + def _find_record(path: Path, turn: int, call_id: str) -> dict[str, Any] | None: """Last ``t:"result"`` record matching ``(turn, call_id)`` in this run. diff --git a/tests/test_agent_team_workflow.py b/tests/test_agent_team_workflow.py index 29d4c69..3561791 100644 --- a/tests/test_agent_team_workflow.py +++ b/tests/test_agent_team_workflow.py @@ -343,37 +343,61 @@ def test_spill_cd_refusal_explains_the_real_reason() -> None: assert "read-only recovery store" in error -def test_bwrap_mounts_the_spill_store_read_only(tmp_path, monkeypatch) -> None: +def test_bwrap_mounts_only_the_current_scopes_spill_store_read_only(tmp_path, monkeypatch) -> None: """The lexical gate cannot be the only layer protecting recovery files. - The bash token scan refuses every parseable way to write into the store, but - shell expansion can hide a path from any scanner (a glob, a brace, a ``$VAR`` - assembled in pieces, a ``$(…)`` substitution), and the files are ordinary 0644 - files. File modes cannot close it either — model commands run as uid 0 inside - the user namespace, so DAC is never consulted. A mount is. - - The mount is now a sibling of ``/workspace`` rather than a remount inside it, - which is what let the ordering requirement go: the source no longer overlaps - any writable bind, so it does not have to come last to win. + Shell expansion can hide a path from any scanner, and model commands run as + uid 0 inside the user namespace, so DAC is never consulted. A read-only + mount is. It is resolved PER COMMAND from the current scope — one jail serves + many agents — and exposes that scope's store and its sub-agents' only, never + the shared root (where ``ls /spill`` would enumerate every conversation). """ - from plugins.tools import _sandbox + from frontier_agent.core.execution_context import ( + ExecutionScope, + reset_current_execution_scope, + set_current_execution_scope, + ) + from plugins.tools import _overflow, _sandbox monkeypatch.setattr(_sandbox, "bwrap_available", lambda: True) monkeypatch.setenv("APODEX_SPILL_DIR", str(tmp_path / "store")) workspace = tmp_path / "ws" sandbox = _sandbox.BwrapSandbox(workspace=workspace) - args = sandbox.commands._bind_args - - spill_at = args.index("--ro-bind-try") - assert args[spill_at + 1] == str(tmp_path / "store") - assert args[spill_at + 2] == "/spill" - # Outside the workspace, so nothing above it overlaps. - assert not str(tmp_path / "store").startswith(str(workspace)) - # ``--ro-bind-try``, not ``--ro-bind``: the store is created lazily on the - # first spill and bwrap aborts the jail when a --ro-bind source is missing. - assert "--ro-bind" not in args[spill_at:spill_at + 1] - assert not (workspace / ".spill").exists() + assert str(tmp_path / "store") not in sandbox.commands._bind_args + + parent = ExecutionScope(task_id="T", metadata={"llm_session_id": "main"}) + child = ExecutionScope(task_id="T", metadata={"llm_session_id": "sub-a"}) + sibling = ExecutionScope(task_id="T", metadata={"llm_session_id": "sub-b"}) + _overflow.register_child_scope(parent, child) + _overflow.register_child_scope(parent, sibling) + root = (tmp_path / "store").resolve() + + def mounts_as(scope: ExecutionScope) -> dict[str, str]: + token = set_current_execution_scope(scope) + try: + args = _sandbox._spill_mount_args() + finally: + reset_current_execution_scope(token) + assert args[:2] == ["--dir", "/spill"] + assert "--bind" not in args and "--ro-bind" not in args + return { + args[i + 2]: args[i + 1] + for i, a in enumerate(args) if a == "--ro-bind-try" + } + + def comp(scope: ExecutionScope) -> str: + return _overflow.scope_component(_overflow._scope_key_of(scope)) + + try: + assert mounts_as(child) == {f"/spill/{comp(child)}": str(root / comp(child))} + assert set(mounts_as(parent)) == { + f"/spill/{comp(s)}" for s in (parent, child, sibling) + } + assert str(root) not in mounts_as(child).values() + assert not str(root).startswith(str(workspace)) + finally: + _overflow._child_scopes.clear() def test_the_shipped_writer_agrees_with_the_authoritative_spill_rule( diff --git a/tests/test_long_run_compaction.py b/tests/test_long_run_compaction.py index 6fe7d24..82107ca 100644 --- a/tests/test_long_run_compaction.py +++ b/tests/test_long_run_compaction.py @@ -784,18 +784,17 @@ def test_a_spilled_body_is_readable_but_not_writable_by_the_file_tools( from plugins.tools._path_auth import _is_path_allowed monkeypatch.setenv("APODEX_SPILL_DIR", str(tmp_path / "store")) + from plugins.tools._sandbox import resolve_runtime_path + token = set_current_execution_scope(ExecutionScope(task_id="t", metadata={})) try: visible = _overflow.spill_compacted_body("bash", "body " * 400) + assert visible + physical = resolve_runtime_path(visible) + readable, _ = _is_path_allowed(physical) + writable, reason = _is_path_allowed(physical, write_access=True) finally: reset_current_execution_scope(token) - assert visible - - from plugins.tools._sandbox import resolve_runtime_path - - physical = resolve_runtime_path(visible) - readable, _ = _is_path_allowed(physical) - writable, reason = _is_path_allowed(physical, write_access=True) assert readable, "recovery cannot work if the store is unreadable" assert not writable, reason diff --git a/tests/test_path_authorization.py b/tests/test_path_authorization.py index 2e0f157..e0a061a 100644 --- a/tests/test_path_authorization.py +++ b/tests/test_path_authorization.py @@ -443,10 +443,11 @@ def test_another_conversations_store_is_not_readable(tmp_path, monkeypatch) -> N _overflow._created_stores.update(saved) -def test_an_in_process_subagents_store_stays_readable(tmp_path, monkeypatch) -> None: +def test_a_parent_reads_its_subagents_store_but_siblings_do_not(tmp_path, monkeypatch) -> None: """A sub-agent spills under its OWN scope, and a fan-in report can carry that - path back to the parent — so scope alone is too narrow. Stores this process - created are authorized too.""" + path back to the parent — so the parent may read its children's stores. A + sibling sub-agent, or an unrelated conversation in this process, may not: + authorization follows the same scope tree the bwrap mount exposes.""" from frontier_agent.core.execution_context import ( ExecutionScope, reset_current_execution_scope, @@ -454,33 +455,46 @@ def test_an_in_process_subagents_store_stays_readable(tmp_path, monkeypatch) -> ) from plugins.tools import _overflow from plugins.tools._path_auth import _is_path_allowed + from plugins.tools._sandbox import resolve_runtime_path monkeypatch.setenv("APODEX_SPILL_DIR", str(tmp_path / "store")) + monkeypatch.setattr(_overflow, "_visible_root", lambda: "/spill") saved = set(_overflow._created_stores) _overflow._created_stores.clear() - try: - token = set_current_execution_scope( - ExecutionScope(task_id="sub", metadata={"llm_session_id": "sub-s"}), - ) + parent = ExecutionScope(task_id="T", metadata={"llm_session_id": "main"}) + sub_a = ExecutionScope(task_id="T", metadata={"llm_session_id": "sub-a"}) + sub_b = ExecutionScope(task_id="T", metadata={"llm_session_id": "sub-b"}) + other = ExecutionScope(task_id="U", metadata={"llm_session_id": "main"}) + _overflow.register_child_scope(parent, sub_a) + _overflow.register_child_scope(parent, sub_b) + + def spill_as(scope: ExecutionScope) -> str: + token = set_current_execution_scope(scope) try: - sub_ref = _overflow.spill_compacted_body("collect_reports", "sub " * 400) + ref = _overflow.spill_compacted_body("collect_reports", f"{scope.metadata} " * 400) finally: reset_current_execution_scope(token) - assert sub_ref + assert ref + return resolve_runtime_path(ref) - # Back in the parent's scope, the sub-agent's path is still readable. - token = set_current_execution_scope( - ExecutionScope(task_id="parent", metadata={"llm_session_id": "p-s"}), - ) + def readable_as(scope: ExecutionScope, path: str) -> bool: + token = set_current_execution_scope(scope) try: - from plugins.tools._sandbox import resolve_runtime_path - - assert _is_path_allowed(resolve_runtime_path(sub_ref))[0] + return _is_path_allowed(path)[0] finally: reset_current_execution_scope(token) + + try: + a_path, b_path = spill_as(sub_a), spill_as(sub_b) + assert readable_as(parent, a_path) and readable_as(parent, b_path) + assert readable_as(sub_a, a_path) + assert not readable_as(sub_a, b_path), "a sibling's store must stay closed" + assert not readable_as(sub_b, a_path) + assert not readable_as(other, a_path), "another conversation's store must stay closed" finally: _overflow._created_stores.clear() _overflow._created_stores.update(saved) + _overflow._child_scopes.clear() def test_every_write_guard_refuses_the_store(tmp_path, monkeypatch) -> None: diff --git a/tests/test_site3_recovery.py b/tests/test_site3_recovery.py index b3c1c6a..f8a7f2b 100644 --- a/tests/test_site3_recovery.py +++ b/tests/test_site3_recovery.py @@ -62,6 +62,21 @@ def scope(): reset_current_execution_scope(token) +@pytest.fixture +def traced(tmp_path): + """A scope whose trajectory JSONL is advertised, as the observer does at + loop start — the precondition for the footer to name ``recover_result``.""" + sc = ExecutionScope( + task_id="t1", role_id="react", + metadata={TrajectoryFileObserver.SCOPE_KEY: str(tmp_path / "t1.jsonl")}, + ) + token = set_current_execution_scope(sc) + try: + yield sc + finally: + reset_current_execution_scope(token) + + def _ctx(turn: int) -> TurnContext: return TurnContext( turn=turn, max_turns=4, task_id="t1", role_id="react", ai_text="", @@ -116,7 +131,7 @@ def _result(body: str, call_id: str = "call_7") -> ToolResult: ) -def test_footer_appears_only_when_content_was_actually_cut() -> None: +def test_footer_appears_only_when_content_was_actually_cut(traced) -> None: full = "x" * 1_000 assert "recover_result" in _with_recovery_handle( "x" * 400, _result(full), 3, enabled=True, @@ -139,7 +154,7 @@ def test_no_footer_without_a_call_id() -> None: assert "recover_result" not in out -def test_footer_names_the_turn_and_id_without_call_syntax() -> None: +def test_footer_names_the_turn_and_id_without_call_syntax(traced) -> None: """The values must be exact, and the shape must not look like source. This asserted ``recover_result(turn=17, call_id="call_abc")`` verbatim until a @@ -175,7 +190,7 @@ def test_no_footer_when_the_body_already_names_a_spill_file() -> None: assert out == body, "footer competed with a pointer that already covers the cut" -def test_the_footer_still_fires_when_nothing_else_covers_the_cut() -> None: +def test_the_footer_still_fires_when_nothing_else_covers_the_cut(traced) -> None: """The suppression must not swallow the case the tool exists for: a result cut below gate ① never reached the store, so no path names it.""" body = "kept output with no pointer at all" @@ -184,7 +199,7 @@ def test_the_footer_still_fires_when_nothing_else_covers_the_cut() -> None: assert "recover_result" in out -def test_a_lookalike_directory_does_not_suppress_the_footer() -> None: +def test_a_lookalike_directory_does_not_suppress_the_footer(traced) -> None: """``"/spill" in body`` also fires on ``/spillover``; a real pointer always names a file UNDER the store, so the separator is what distinguishes them.""" body = "see /spillover/notes.md for context" @@ -452,3 +467,9 @@ def test_a_subagent_recovers_the_body_its_own_scope_names(tmp_path) -> None: assert "SUBAGENT-BODY" in out assert "COORDINATOR-BODY" not in out + + +def test_no_footer_when_no_trajectory_jsonl_is_written(scope) -> None: + """A handle into a JSONL nobody writes can only answer "unavailable".""" + full = "x" * 1_000 + assert _with_recovery_handle("x" * 400, _result(full), 3, enabled=True) == "x" * 400 diff --git a/tests/test_spill_scope_recovery.py b/tests/test_spill_scope_recovery.py new file mode 100644 index 0000000..db1cafc --- /dev/null +++ b/tests/test_spill_scope_recovery.py @@ -0,0 +1,317 @@ +"""Spill recovery and scope isolation (handoff item 3; Harness #494 / #521). + +The rules pinned here: + +* one scope key decides the store directory, the bwrap mount and read + authorization (``current_store_scope``); +* a ref is advertised only when the backend that runs commands can open it, + and the compaction callback is withheld otherwise; +* a truncated body stays recoverable through compaction, twice over; +* nothing crosses scopes, and cleanup touches only this process's stores. + +Persistence is AgentCore's ``SpillStore``; there is no second store. +""" + +from __future__ import annotations + +import os +import shutil +import subprocess + +import pytest + +from frontier_agent.core.execution_context import ( + ExecutionScope, + reset_current_execution_scope, + set_current_execution_scope, +) +from frontier_agent.core.runtime.loop.compact import KeepLastNToolResultsCompactor +from plugins.tools import _overflow, _sandbox +from plugins.tools.meta import get_tool_meta + + +@pytest.fixture(autouse=True) +def _isolated(tmp_path, monkeypatch): + monkeypatch.setenv("APODEX_SPILL_DIR", str(tmp_path / "store")) + monkeypatch.delenv("SANDBOX_BACKEND", raising=False) + saved = set(_overflow._created_stores) + _overflow._created_stores.clear() + try: + yield + finally: + _overflow._created_stores.clear() + _overflow._created_stores.update(saved) + _overflow._child_scopes.clear() + _overflow._preview_refs.clear() + + +class _scoped: + def __init__(self, task: str, session: str = "") -> None: + meta = {"llm_session_id": session} if session else {} + self.scope = ExecutionScope(task_id=task, metadata=meta) + + def __enter__(self) -> ExecutionScope: + self._token = set_current_execution_scope(self.scope) + return self.scope + + def __exit__(self, *exc: object) -> None: + reset_current_execution_scope(self._token) + + +def _body(n: int = 40_000) -> str: + return "".join(f"line {i}: evidence {i * 7}\n" for i in range(n // 20)) + + +def _mounted(monkeypatch) -> None: + monkeypatch.setattr(_overflow, "_visible_root", lambda: "/spill") + + +# ── one identity ───────────────────────────────────────────────────────── + + +def test_store_mount_and_read_auth_share_one_scope_key(monkeypatch) -> None: + from plugins.tools._path_auth import _resolve_spill_dirs + + _mounted(monkeypatch) + with _scoped("T", "sess") as scope: + ref = _overflow.spill_compacted_body("bash", _body()) + assert ref + key = _overflow.current_store_scope() + assert key == _overflow.spill_scope_key("T", "sess") == "T:sess" + component = _overflow.scope_component(key) + store_dir = (_sandbox.spill_root().resolve() / component) + assert ref.startswith(f"/spill/{component}/") + assert _resolve_spill_dirs() == [store_dir] + mounts = _sandbox._spill_mount_args() + assert mounts == [ + "--dir", "/spill", "--ro-bind-try", str(store_dir), f"/spill/{component}", + ] + assert scope.metadata["llm_session_id"] == "sess" + + +def test_unscoped_store_is_per_process() -> None: + key = _overflow.current_store_scope() + assert key == f"{_overflow.UNSCOPED_PROCESS_STORE}:{os.getpid()}" + + +def test_spill_files_are_read_only_and_symlinks_refused(monkeypatch, tmp_path) -> None: + _mounted(monkeypatch) + with _scoped("T"): + path, _ref = _overflow._write_spill("bash", _body(), require_visible=True) + assert oct(path.stat().st_mode & 0o777) == "0o444" + + real = tmp_path / "elsewhere" + real.mkdir() + link = tmp_path / "linked-root" + link.symlink_to(real) + monkeypatch.setenv("APODEX_SPILL_DIR", str(link)) + monkeypatch.setattr(_sandbox, "_private_spill_root", None) + first, second = _sandbox.spill_root(), _sandbox.spill_root() + assert first == second, "writer, mount and reader must agree on one fallback" + assert not first.is_symlink() and first != link + + +# ── refs only where the backend can open them ──────────────────────────── + + +@pytest.mark.parametrize(("backend", "bwrap", "expected"), [ + ("e2b", True, None), + ("bwrap", False, None), + ("bwrap", True, "/spill"), + ("local", True, "/spill"), + ("native", False, "physical"), + ("container", False, "physical"), +]) +def test_visible_root_follows_the_backend(monkeypatch, backend, bwrap, expected) -> None: + monkeypatch.setenv("SANDBOX_BACKEND", backend) + monkeypatch.setenv("E2B_API_KEY", "k") + monkeypatch.setattr(_sandbox, "bwrap_available", lambda: bwrap) + monkeypatch.setattr(_sandbox, "get_existing_sandbox", lambda: None) + monkeypatch.setattr(_sandbox, "_container_inner_bwrap_enabled", lambda: False) + want = str(_sandbox.spill_root().resolve()) if expected == "physical" else expected + assert _overflow._visible_root() == want + + +def test_auto_with_an_e2b_key_names_no_host_path(monkeypatch) -> None: + monkeypatch.setenv("SANDBOX_BACKEND", "auto") + monkeypatch.setattr(_sandbox, "get_existing_sandbox", lambda: None) + monkeypatch.setattr(_sandbox, "bwrap_available", lambda: True) + monkeypatch.setattr(_sandbox, "_resolve_use_e2b", lambda: (True, "k", "t", 60)) + assert _overflow._visible_root() is None + + +def test_a_live_sandbox_wins_over_configuration(monkeypatch) -> None: + monkeypatch.setenv("SANDBOX_BACKEND", "e2b") + + class FakeBwrap(_sandbox.BwrapSandbox): + def __init__(self) -> None: # no real jail needed + pass + + monkeypatch.setattr(_sandbox, "get_existing_sandbox", lambda: FakeBwrap()) + assert _overflow._visible_root() == "/spill" + + e2b_like = type("Sandbox", (), {"__module__": "e2b_code_interpreter.main"})() + monkeypatch.setenv("SANDBOX_BACKEND", "bwrap") + monkeypatch.setattr(_sandbox, "bwrap_available", lambda: True) + monkeypatch.setattr(_sandbox, "get_existing_sandbox", lambda: e2b_like) + assert _overflow._visible_root() is None + + +def test_unrecoverable_backend_withholds_callback_and_writes_nothing(monkeypatch, tmp_path) -> None: + monkeypatch.setattr(_overflow, "_visible_root", lambda: None) + assert _overflow.default_compaction_spill() is None + cap = get_tool_meta("bash").max_result_chars + with _scoped("T"): + out = _overflow.maybe_overflow("bash", _body(cap * 4)) + assert "not readable from this backend" in out + assert "saved read-only at" not in out + store = tmp_path / "store" + assert not store.exists() or not any(store.rglob("*.md")) + + +def test_recoverable_backend_gets_the_callback(monkeypatch) -> None: + _mounted(monkeypatch) + assert _overflow.default_compaction_spill() is _overflow.spill_compacted_body + + +# ── recoverable through truncation and repeated compaction ─────────────── + + +def _history(content: str) -> list[dict]: + return [ + {"role": "system", "content": "system"}, + {"role": "user", "content": "task"}, + {"role": "assistant", "content": "", + "tool_calls": [{"id": "c1", "name": "bash", "args": {}}]}, + {"role": "tool", "tool_call_id": "c1", "content": content}, + {"role": "user", "content": "continue"}, + ] + + +def test_full_text_survives_truncation_and_two_compactions(monkeypatch) -> None: + _mounted(monkeypatch) + full = _body(60_000) + with _scoped("T", "s"): + preview = _overflow.maybe_overflow("bash", full) + assert len(preview) < len(full) + compactor = KeepLastNToolResultsCompactor( + keep_tool_result=0, spill=_overflow.default_compaction_spill(), + ) + once = compactor.compact(_history(preview), keep_recent=1) + twice = compactor.compact(once, keep_recent=1) + ref = once[3]["spill_refs"][0] + assert twice[3]["spill_refs"] == [ref], "re-compaction must keep the same ref" + assert f"[Full text] {ref}" in once[3]["content"] + # The card points at the ORIGINAL body, not at a copy of the preview. + assert _overflow.get_overflow_content(ref) == full + files = list(_sandbox.spill_root().rglob("*.md")) + assert len(files) == 1, "the preview must not be stored a second time" + + +def test_a_preview_re_cut_by_the_aggregate_budget_still_maps_to_the_full_body(monkeypatch) -> None: + _mounted(monkeypatch) + full = _body(60_000) + with _scoped("T"): + previews = [_overflow.maybe_overflow("bash", full + str(i)) for i in range(40)] + adjusted = _overflow.check_aggregate_budget(previews, ["bash"] * 40) + recut = next(a for a, p in zip(adjusted, previews, strict=True) if a != p) + ref = _overflow.spill_compacted_body("bash", recut) + assert ref and _overflow.get_overflow_content(ref).startswith(full) + + +def test_an_echoed_fragment_of_a_preview_claims_no_ref(monkeypatch) -> None: + _mounted(monkeypatch) + with _scoped("T"): + preview = _overflow.maybe_overflow("bash", _body(60_000)) + fragment = preview[: len(preview) // 2] + "x" * 2_000 + assert _overflow._full_body_ref(fragment) == "" + + +# ── cross-scope access ─────────────────────────────────────────────────── + + +def test_a_siblings_ref_is_not_readable(monkeypatch) -> None: + _mounted(monkeypatch) + parent, a, b = (ExecutionScope(task_id="T", metadata={"llm_session_id": s}) + for s in ("main", "a", "b")) + _overflow.register_child_scope(parent, a) + _overflow.register_child_scope(parent, b) + with _scoped("T", "a"): + ref_a = _overflow.spill_compacted_body("bash", _body()) + with _scoped("T", "b"): + assert _overflow.get_overflow_content(ref_a) is None + assert _overflow.readable_store_dirs() == [] + with _scoped("T", "main"): + assert _overflow.get_overflow_content(ref_a) is not None + + +def test_entering_a_nested_loop_registers_it_as_a_child() -> None: + from frontier_agent.core.loop_types import LoopConfig + from frontier_agent.core.runtime.loop.agent_loop import _enter_scope + + with _scoped("T", "main"): + child, token = _enter_scope( + LoopConfig(task_id="T", llm_session_id="sub"), "p", {"llm_session_id": "sub"}, + ) + reset_current_execution_scope(token) + assert "T:sub" in _overflow.readable_scopes("T:main") + assert _overflow.readable_scopes("T:sub") == ["T:sub"] + assert child.metadata["llm_session_id"] == "sub" + + +# ── cleanup ────────────────────────────────────────────────────────────── + + +def test_tree_cleanup_removes_only_this_conversation(monkeypatch, tmp_path) -> None: + _mounted(monkeypatch) + parent = ExecutionScope(task_id="T", metadata={"llm_session_id": "main"}) + child = ExecutionScope(task_id="T", metadata={"llm_session_id": "sub"}) + _overflow.register_child_scope(parent, child) + with _scoped("T", "main"): + _overflow.spill_compacted_body("bash", _body()) + with _scoped("T", "sub"): + _overflow.spill_compacted_body("bash", _body()) + with _scoped("U", "main"): + keep = _overflow.spill_compacted_body("bash", _body()) + # A store another process created, sitting under the same root. + foreign = _sandbox.spill_root() / _overflow.scope_component("T:other-process") + foreign.mkdir(parents=True) + (foreign / "x.md").write_text("theirs") + + assert _overflow.cleanup_overflow_tree("T:main") == 2 + with _scoped("U", "main"): + assert _overflow.get_overflow_content(keep) is not None + assert (foreign / "x.md").exists() + + +# ── real jail (needs a host where bwrap can create namespaces) ─────────── + + +@pytest.mark.skipif( + not _sandbox.bwrap_available(), + reason="bwrap cannot create namespaces here; isolation is NOT verified by a skip", +) +def test_real_jail_shows_only_the_current_scope(tmp_path, monkeypatch) -> None: + monkeypatch.setenv("SANDBOX_BACKEND", "bwrap") + sandbox = _sandbox.BwrapSandbox(workspace=tmp_path / "ws") + with _scoped("T", "a"): + ref_a = _overflow.spill_compacted_body("bash", _body()) + with _scoped("T", "b"): + ref_b = _overflow.spill_compacted_body("bash", _body() + "b") + listing = sandbox.commands.run("ls /spill").stdout.split() + assert listing == [_overflow.scope_component("T:b")] + assert sandbox.commands.run(f"cat {ref_a}").exit_code != 0 + assert sandbox.commands.run(f"cat {ref_b}").exit_code == 0 + assert sandbox.commands.run(f"touch {ref_b}").exit_code != 0 + sandbox.kill() + + +def test_bwrap_is_really_unavailable_here_when_skipped() -> None: + """Records WHY the real-jail test skips on this host, so a skip is never + read as a pass.""" + if _sandbox.bwrap_available() or shutil.which("bwrap") is None: + pytest.skip("not applicable") + probe = subprocess.run( + ["bwrap", "--ro-bind", "/", "/", "true"], capture_output=True, text=True, + ) + assert probe.returncode != 0 diff --git a/tests/test_tool_result_truncation.py b/tests/test_tool_result_truncation.py index f6bd515..0116c30 100644 --- a/tests/test_tool_result_truncation.py +++ b/tests/test_tool_result_truncation.py @@ -170,8 +170,8 @@ def test_overflow_degrades_to_a_pointerless_footer_when_the_write_fails( from plugins.tools.meta import get_tool_meta monkeypatch.setattr( - _overflow, "_spill_document", - lambda *a, **k: (_ for _ in ()).throw(OSError("read-only fs")), + _overflow.SpillStore, "_document", + staticmethod(lambda *a, **k: (_ for _ in ()).throw(OSError("read-only fs"))), ) cap = get_tool_meta("bash").max_result_chars body = _pytest_output(cap * 6) diff --git a/workflows/agent_team/nodes/main_agent.py b/workflows/agent_team/nodes/main_agent.py index 0ada121..af309be 100644 --- a/workflows/agent_team/nodes/main_agent.py +++ b/workflows/agent_team/nodes/main_agent.py @@ -1089,7 +1089,7 @@ async def main_agent_node( main_keep_recent_msgs = max(6, int(agent_cfg.get("keep_recent_turns", 5)) * 3) main_compaction_policy: Any = None main_gauge: InputTokenGauge | None = None - from plugins.tools._overflow import spill_compacted_body + from plugins.tools._overflow import default_compaction_spill if context_compaction == "tiered" and max_len > 0: main_gauge = InputTokenGauge() @@ -1102,7 +1102,7 @@ async def main_agent_node( relief_target=int(max_len * 0.6), protect_tool_names=PROTECTED_FANIN_TOOLS, gauge=main_gauge, # calibrate relief to real tokens (unit-match trigger) - spill=spill_compacted_body if compaction_spill else None, + spill=default_compaction_spill() if compaction_spill else None, # Bound the whole retry sequence by what ONE summariser call was # already allowed to spend, so retrying costs no extra worst case. summary_retry_timeout_s=llm_timeout, diff --git a/workflows/agent_team/subagent_runtime.py b/workflows/agent_team/subagent_runtime.py index d51ead2..e905ab3 100644 --- a/workflows/agent_team/subagent_runtime.py +++ b/workflows/agent_team/subagent_runtime.py @@ -1117,7 +1117,7 @@ def build_swarm_session_runtime_spec( # is already reset by then and the list has to be handed over explicitly. _denials: dict[str, list[tuple[str, str]]] = {} _tiered = runtime.context_compaction == "tiered" and runtime.max_len > 0 - from plugins.tools._overflow import spill_compacted_body + from plugins.tools._overflow import default_compaction_spill def _config(_job_id: str, item: SubTask, max_turns: int) -> LoopConfig: if _tiered: @@ -1127,7 +1127,7 @@ def _config(_job_id: str, item: SubTask, max_turns: int) -> LoopConfig: summary_llm=runtime.sub_agent_llm, relief_target=int(runtime.max_len * 0.6), gauge=gauge, # calibrate relief to real tokens (unit-match trigger) - spill=spill_compacted_body if runtime.compaction_spill else None, + spill=default_compaction_spill() if runtime.compaction_spill else None, summary_retry_timeout_s=( runtime.sub_agent_llm_timeout or LoopConfig.llm_timeout ), diff --git a/workflows/stateful_react_agent/nodes/main_agent.py b/workflows/stateful_react_agent/nodes/main_agent.py index 8fcaa4b..e5e77cf 100644 --- a/workflows/stateful_react_agent/nodes/main_agent.py +++ b/workflows/stateful_react_agent/nodes/main_agent.py @@ -1050,14 +1050,14 @@ async def react_agent_node(state: dict[str, Any], ctx: NodeContext) -> dict[str, if context_compaction == "tiered" and max_len > 0: gauge = InputTokenGauge() observers.append(gauge) - from plugins.tools._overflow import spill_compacted_body + from plugins.tools._overflow import default_compaction_spill compactor: Any = TieredCompactor( keep_tool_result=tier1_keep_tool_result, summary_llm=llm, relief_target=int(max_len * 0.6), gauge=gauge, # calibrate relief to real tokens (unit-match trigger) - spill=spill_compacted_body if compaction_spill else None, + spill=default_compaction_spill() if compaction_spill else None, summary_retry_timeout_s=llm_timeout, ) compaction_policy = InputTokenThresholdPolicy( @@ -1125,6 +1125,14 @@ async def react_agent_node(state: dict[str, Any], ctx: NodeContext) -> dict[str, # Drop this run's task board (no-op when task_board is off) so boards # don't leak across trials in a long-lived worker process. clear_board(ctx.task_id) + # Same for its spilled tool results: finalization ran inside the loop, + # so nothing can name them any more. Only stores this process created + # under this run's scope are removed. + from plugins.tools._overflow import cleanup_overflow_tree, spill_scope_key + + cleanup_overflow_tree(spill_scope_key( + ctx.task_id, str(metadata.get("session_id") or ctx.task_id), + )) kill = getattr(sandbox, "kill", None) if callable(kill): kill() From 89875b2b248d165735643e1c3f4f714aa4a95d81 Mon Sep 17 00:00:00 2001 From: zhanghanduo Date: Tue, 6 Oct 2026 14:57:10 +0800 Subject: [PATCH 2/2] Keep spill changelog separate from path gate entry --- CHANGELOG.md | 26 +++++++++++++------------- 1 file changed, 13 insertions(+), 13 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 7540752..98ddd05 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -62,19 +62,6 @@ Initial open-source release of FrontierAgent. ### Fixed -- Spill recovery: oversized and compacted tool results are stored through - AgentCore's `SpillStore` (one store, one registry). A single scope key - (`task:llm_session`) now decides the store directory, the bwrap `/spill` mount - and read authorization: a jail sees only its own store and its sub-agents', - siblings and other conversations see nothing, and `_path_auth` no longer - authorizes every store the process created. A path is advertised only when - the backend actually running commands can open it (`auto` going to E2B, or - bwrap unusable, gets none), and the compaction spill callback is withheld - then. Compacting a truncated preview points at the original full body - instead of storing the preview as "[Full text]". The `recover_result` footer - is only shown when a trajectory JSONL exists. Known limitation: `container` - mode without the inner bwrap jail shares one tool uid, so model commands can - still read other scopes there. - Bash policy: privilege escalation (`sudo`/`su`/…), remote/exfil clients (`ssh`/`nc`/`rsync`/…) and signal senders (`kill`/`pkill`/`killall`) are now refused in every allowlist mode, including the default `off`. The local CLI @@ -89,6 +76,19 @@ Initial open-source release of FrontierAgent. `reboot` are no longer refused, while `bash -c`, shell heredocs, pipes into a shell, evaluators (`watch`/`tmux`/…) and `systemctl` shutdown units still are. `DROP TABLE` keeps screening the whole text. +- Spill recovery: oversized and compacted tool results are stored through + AgentCore's `SpillStore` (one store, one registry). A single scope key + (`task:llm_session`) now decides the store directory, the bwrap `/spill` mount + and read authorization: a jail sees only its own store and its sub-agents', + siblings and other conversations see nothing, and `_path_auth` no longer + authorizes every store the process created. A path is advertised only when + the backend actually running commands can open it (`auto` going to E2B, or + bwrap unusable, gets none), and the compaction spill callback is withheld + then. Compacting a truncated preview points at the original full body + instead of storing the preview as "[Full text]". The `recover_result` footer + is only shown when a trajectory JSONL exists. Known limitation: `container` + mode without the inner bwrap jail shares one tool uid, so model commands can + still read other scopes there. - Surface finalize-gate bypasses on the final turn: an answer delivered despite open task-board items now carries an unfinished-work note and a `finalize_gate_bypassed` marker instead of reading as a clean success.