From 34f01c0b58c29993cdecd3ad90163882034e638f Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?L=C3=A9o?= Date: Thu, 1 Oct 2026 16:03:03 +0200 Subject: [PATCH 1/4] fix(agent): capture Windows Bash output and pass the subprocess env MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The Windows branch of `_create_bash_subprocess` spawned Git Bash with neither pipes nor the env it was handed. `proc.stdout` and `proc.stderr` came back `None`, so `_run_subprocess_streaming`'s reader returned immediately and the Bash tool reported `"(no output)"` alongside the real exit code — while the child inherited the server's own stdout/stderr and wrote agent command output into the console and the launchd/Docker logs. The `env` parameter was accepted and never used, so `PATH`, `VIRTUAL_ENV`, `HOME`, `TMPDIR` and the configured import paths carried in `ctx["subproc_env"]` never reached the child on Windows, even though every POSIX path applies them. Spawn it the way the POSIX path at `:688` already does: `stdin=DEVNULL`, `stdout=PIPE`, `stderr=PIPE`, `env=env`. `website/configuration-reference.md` is generated from source line numbers, so the four added lines shift one entry; regenerated with `scripts/generate_env_reference.py`. --- src/agent_tools/subprocess_tools.py | 4 ++ tests/test_agent_bash_windows.py | 78 +++++++++++++++++++++++++++++ website/configuration-reference.md | 2 +- 3 files changed, 83 insertions(+), 1 deletion(-) diff --git a/src/agent_tools/subprocess_tools.py b/src/agent_tools/subprocess_tools.py index 7537c82b5..b4b8b5405 100644 --- a/src/agent_tools/subprocess_tools.py +++ b/src/agent_tools/subprocess_tools.py @@ -117,6 +117,10 @@ async def _create_bash_subprocess( bash, "-c", str(command or ""), + stdin=asyncio.subprocess.DEVNULL, + stdout=asyncio.subprocess.PIPE, + stderr=asyncio.subprocess.PIPE, + env=env, cwd=cwd, ) kwargs = {"cwd": cwd} if cwd is not None else {} diff --git a/tests/test_agent_bash_windows.py b/tests/test_agent_bash_windows.py index 0145b4b28..fabcee974 100644 --- a/tests/test_agent_bash_windows.py +++ b/tests/test_agent_bash_windows.py @@ -1,5 +1,7 @@ """Windows execution contract for the agent Bash tool.""" +import asyncio + import pytest from types import SimpleNamespace @@ -38,6 +40,82 @@ async def test_windows_bash_uses_git_bash_with_structural_cwd(monkeypatch): assert captured["kwargs"]["cwd"] == workspace +@pytest.mark.asyncio +async def test_windows_bash_captures_output_instead_of_inheriting_server_handles(monkeypatch): + captured = {} + + monkeypatch.setattr(subprocess_tools, "IS_WINDOWS", True) + monkeypatch.setattr( + subprocess_tools, "find_bash", lambda: r"C:\Program Files\Git\bin\bash.exe" + ) + + async def fake_exec(*_argv, **kwargs): + captured.update(kwargs) + return object() + + monkeypatch.setattr(subprocess_tools.asyncio, "create_subprocess_exec", fake_exec) + + await subprocess_tools._create_bash_subprocess("pwd", cwd=r"C:\Work") + + assert captured["stdout"] == asyncio.subprocess.PIPE + assert captured["stderr"] == asyncio.subprocess.PIPE + assert captured["stdin"] == asyncio.subprocess.DEVNULL + + +@pytest.mark.asyncio +async def test_windows_bash_applies_the_subprocess_env(monkeypatch): + captured = {} + env = {"PATH": r"C:\Odysseus\venv\Scripts", "HOME": r"C:\Odysseus\data"} + + monkeypatch.setattr(subprocess_tools, "IS_WINDOWS", True) + monkeypatch.setattr( + subprocess_tools, "find_bash", lambda: r"C:\Program Files\Git\bin\bash.exe" + ) + + async def fake_exec(*_argv, **kwargs): + captured.update(kwargs) + return object() + + monkeypatch.setattr(subprocess_tools.asyncio, "create_subprocess_exec", fake_exec) + + await subprocess_tools._create_bash_subprocess("pwd", cwd=r"C:\Work", env=env) + + assert captured["env"] == env + + +@pytest.mark.asyncio +async def test_windows_bash_tool_passes_ctx_env_through_to_the_child(monkeypatch): + captured = {} + env = {"PATH": r"C:\Odysseus\venv\Scripts", "VIRTUAL_ENV": r"C:\Odysseus\venv"} + + monkeypatch.setattr(subprocess_tools, "IS_WINDOWS", True) + monkeypatch.setattr( + subprocess_tools, "find_bash", lambda: r"C:\Program Files\Git\bin\bash.exe" + ) + monkeypatch.setattr("src.tool_execution.agent_cwd", lambda: r"D:\Workspaces\Project") + + async def fake_exec(*argv, **kwargs): + captured["argv"] = argv + captured["kwargs"] = kwargs + return SimpleNamespace(pid=4242) + + async def fake_stream(_process, **_kwargs): + return "ok", "", 0, False + + monkeypatch.setattr(subprocess_tools.asyncio, "create_subprocess_exec", fake_exec) + monkeypatch.setattr(subprocess_tools, "_run_subprocess_streaming", fake_stream) + + result = await subprocess_tools.BashTool().execute( + "pwd", + {"subproc_env": env, "session_id": "chat-1"}, + ) + + assert result == {"output": "ok", "exit_code": 0} + assert captured["kwargs"]["env"] == env + assert captured["kwargs"]["stdout"] == asyncio.subprocess.PIPE + assert captured["kwargs"]["stderr"] == asyncio.subprocess.PIPE + + @pytest.mark.asyncio async def test_windows_bash_without_git_bash_fails_clearly(monkeypatch): monkeypatch.setattr(subprocess_tools, "IS_WINDOWS", True) diff --git a/website/configuration-reference.md b/website/configuration-reference.md index fc2a3bdeb..ba0541539 100644 --- a/website/configuration-reference.md +++ b/website/configuration-reference.md @@ -75,7 +75,7 @@ The source tree reads **108** `ODYSSEUS_*` variables: 78 an operator may want to | `ODYSSEUS_MAX_VISUAL_EVIDENCE_FRAMES` | `'3'` | `src/agent_loop.py:15361` | How many video frames one tool result may contribute. Clamped to 1-8. | | `ODYSSEUS_MAX_VISUAL_EVIDENCE_IMAGES` | `'1'` | `src/agent_loop.py:15329` | How many images one tool result may contribute to the model turn. Clamped to 1-8. | | `ODYSSEUS_MCP_ALLOWED_COMMANDS` | `''` | `src/agent_tools/admin_tools.py:140` | Security-relevant. Comma-separated allowlist of MCP launcher basenames the agent may start. Empty by default, and the deny list still wins. | -| `ODYSSEUS_PYTHON_TOOL_SITE_PACKAGES` | `''` | `src/agent_tools/subprocess_tools.py:927` | Security-relevant. Absolute package roots, separated by the platform path separator, exposed to the sandboxed Python tool. Empty exposes none. | +| `ODYSSEUS_PYTHON_TOOL_SITE_PACKAGES` | `''` | `src/agent_tools/subprocess_tools.py:931` | Security-relevant. Absolute package roots, separated by the platform path separator, exposed to the sandboxed Python tool. Empty exposes none. | | `ODYSSEUS_SCRIPT_HOST` | `'localhost'` | `src/builtin_actions.py:919` | Default host for the run-script action. `localhost`, `127.0.0.1`, `local` and empty run locally; any other value runs over SSH. | | `ODYSSEUS_TOOL_APPROVAL_GATE` | `'0'` | `src/tool_capabilities.py:645` | Security-relevant. Truthy makes tool calls pass through the approval gate. Off by default. | From f49e09e59a7b81b57d541e4602732d07235e5874 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?L=C3=A9o?= Date: Thu, 1 Oct 2026 16:03:13 +0200 Subject: [PATCH 2/4] fix(shell): kill the PTY command's whole session on timeout MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit /api/shell/stream starts its PTY child under os.setsid, so the child leads its own session and process group. The timeout, client-disconnect and error paths all called proc.kill(), which signals only the group leader. Creating a group and then signalling only its leader is strictly worse than never creating one: the descendants are detached from the server's group as well, so nothing else will ever reach them, while the route reports "Command timed out after Ns" and exit_code -1 as if the command were gone. The kernel's controlling-terminal SIGHUP hid this for well-behaved children, which is why it reads as working. Anything that ignores SIGHUP — a nohup'ed job, a daemon, a process that means to outlive its terminal — survives the kill indefinitely. Signal the whole group instead, escalate to SIGKILL if it outlives the grace period, and confirm it is actually gone. The timeout response now says so when containment could not be established rather than claiming a clean kill it did not get. --- routes/shell_routes.py | 145 +++++++++++++++++++++-- tests/test_shell_routes.py | 229 +++++++++++++++++++++++++++++++++++++ 2 files changed, 364 insertions(+), 10 deletions(-) diff --git a/routes/shell_routes.py b/routes/shell_routes.py index d63e80ee6..900ac4125 100644 --- a/routes/shell_routes.py +++ b/routes/shell_routes.py @@ -1,6 +1,7 @@ """Shell routes — user-facing command execution endpoint.""" import asyncio +import contextlib import importlib import json import logging @@ -8,6 +9,7 @@ import os import re import shlex import shutil +import signal import subprocess import uuid import tempfile @@ -49,6 +51,7 @@ from core.platform_compat import ( detached_popen_kwargs, find_bash, git_bash_path, + pid_alive, ) @@ -558,6 +561,13 @@ STREAM_TIMEOUT = 120 # default for short commands MAX_OUTPUT = 200_000 # truncate limit TMUX_LOG_DIR = Path(tempfile.gettempdir()) / "odysseus-tmux" PTY_UNSUPPORTED_ERROR = "pty_unsupported" +# PTY teardown. The PTY child leads its own session (os.setsid), so killing it +# has to signal the whole process group and then confirm the group is gone — +# see _terminate_pty_session. +PTY_KILL_ESCALATION = (signal.SIGTERM, signal.SIGKILL) +PTY_KILL_GRACE = 1.0 # seconds a signalled session gets to exit +PTY_KILL_POLL_INTERVAL = 0.05 # re-check interval while waiting for it +PTY_KILL_FAILED_HINT = "; processes it started survived the kill and are still running" class ShellExecRequest(BaseModel): @@ -662,6 +672,124 @@ async def _exec_shell(command: str, timeout: int = EXEC_TIMEOUT) -> Dict[str, An return {"stdout": "", "stderr": str(e), "exit_code": -1} +def _session_pgid(pid: int) -> int | None: + """Process-group id of the session ``pid`` leads, or None if unavailable. + + Read this *before* the leader is reaped: once it is, ``getpgid`` fails and + the group id can no longer be recovered from the pid. + """ + getpgid = getattr(os, "getpgid", None) + if getpgid is None: # no process groups (native Windows) + return None + try: + return getpgid(pid) + except OSError: + return None + + +def _signal_session(pgid: int | None, pid: int, sig: int) -> bool: + """Send ``sig`` to the whole process group, or to the lone process. + + Returns whether anything was signalled, so a caller can tell "the session + is already gone" from "the signal landed". The single-pid fallback + matters: if ``setsid`` did not take effect, or the platform has no process + groups, teardown must still reach the child rather than do nothing. + """ + killpg = getattr(os, "killpg", None) + if pgid is not None and killpg is not None: + try: + killpg(pgid, sig) + return True + except ProcessLookupError: + return False + except OSError: + pass # group signalling refused — fall through to the single pid + try: + os.kill(pid, sig) + return True + except OSError: + return False + + +def _session_alive(pgid: int | None, pid: int) -> bool: + """True while any member of the process group still exists. + + An unreaped zombie is still signallable, so a True here can also mean the + leader has exited but not yet been collected. Without a group id this can + only speak for the child itself, not for anything it spawned. + """ + killpg = getattr(os, "killpg", None) + if pgid is not None and killpg is not None: + try: + killpg(pgid, 0) + return True + except OSError: + return False + return pid_alive(pid) + + +async def _await_session_exit(proc, pgid: int | None, pid: int) -> bool: + """Wait up to the grace period for the leader and its group to go away.""" + loop = asyncio.get_running_loop() + deadline = loop.time() + PTY_KILL_GRACE + while True: + remaining = deadline - loop.time() + if proc.returncode is None and remaining > 0: + # Reap the leader, otherwise its own zombie keeps the group alive + # and the liveness probe below can never come back clean. + with contextlib.suppress(asyncio.TimeoutError): + await asyncio.wait_for(proc.wait(), remaining) + if not _session_alive(pgid, pid): + return True + if loop.time() >= deadline: + return False + await asyncio.sleep(PTY_KILL_POLL_INTERVAL) + + +async def _terminate_pty_session(proc) -> bool: + """Kill the PTY child and every process in the session it leads. + + The child is spawned under ``os.setsid``, so it leads its own session and + process group. Signalling only the leader is strictly worse than never + calling ``setsid`` at all: the descendants are detached from the server's + group too, so nothing will ever reach them, while the caller reports the + command as terminated. Signal the group instead, escalate to SIGKILL if it + outlives the grace period, and return whether the session is actually gone + so the caller can say so rather than assume it. + """ + pid = getattr(proc, "pid", None) + if pid is None: + return True + pgid = _session_pgid(pid) + + for sig in PTY_KILL_ESCALATION: + if not _signal_session(pgid, pid, sig): + break # nothing left to signal + if await _await_session_exit(proc, pgid, pid): + return True + return not _session_alive(pgid, pid) + + +async def _terminate_pty_session_quietly(proc) -> None: + """Best-effort :func:`_terminate_pty_session` for paths with no reader. + + The client-disconnect and exception paths have nowhere left to report a + containment failure to, so they log it instead of raising over the top of + whatever is already going wrong. + """ + pid = getattr(proc, "pid", None) + try: + contained = await _terminate_pty_session(proc) + except Exception: + logger.exception("PTY session teardown failed for pid %s", pid) + return + if not contained: + logger.warning( + "PTY session for pid %s survived teardown; it may still be running", + pid, + ) + + async def _generate_pty(cmd: str, timeout: int, request: Request): """Run command in a pseudo-TTY so tqdm/progress bars work natively.""" if not PTY_SUPPORTED: @@ -702,16 +830,17 @@ async def _generate_pty(cmd: str, timeout: int, request: Request): try: while not process_done.is_set(): if deadline and loop.time() > deadline: - proc.kill() - await proc.wait() - yield f"data: {json.dumps({'stream': 'stderr', 'data': f'Command timed out after {timeout}s'})}\n\n" + contained = await _terminate_pty_session(proc) + msg = f"Command timed out after {timeout}s" + if not contained: + msg += PTY_KILL_FAILED_HINT + yield f"data: {json.dumps({'stream': 'stderr', 'data': msg})}\n\n" yield f"data: {json.dumps({'exit_code': -1})}\n\n" return # Check client disconnect if await request.is_disconnected(): - proc.kill() - await proc.wait() + await _terminate_pty_session_quietly(proc) return # Read available data from PTY @@ -773,11 +902,7 @@ async def _generate_pty(cmd: str, timeout: int, request: Request): yield f"data: {json.dumps({'exit_code': proc.returncode})}\n\n" except Exception as e: - try: - proc.kill() - await proc.wait() - except ProcessLookupError: - pass + await _terminate_pty_session_quietly(proc) yield f"data: {json.dumps({'stream': 'stderr', 'data': str(e)})}\n\n" yield f"data: {json.dumps({'exit_code': -1})}\n\n" finally: diff --git a/tests/test_shell_routes.py b/tests/test_shell_routes.py index 072a13d96..281cc8ea6 100644 --- a/tests/test_shell_routes.py +++ b/tests/test_shell_routes.py @@ -1,16 +1,20 @@ """Tests for shell_routes.py helpers.""" +import asyncio import builtins import importlib import importlib.util import json import os +import shlex +import signal import sys from pathlib import Path from types import SimpleNamespace import pytest +from core.platform_compat import pid_alive from routes.shell_routes import ( _find_line_break, _host_docker_access_enabled, @@ -85,6 +89,231 @@ async def test_generate_pty_reports_explicit_unsupported_error(monkeypatch): ] +pty_session = pytest.mark.skipif( + not hasattr(os, "setsid"), reason="process sessions are POSIX-only" +) + + +async def _spawn_pty_style_session(script: str): + """Spawn `script` the way _generate_pty does: its own session via setsid.""" + return await asyncio.create_subprocess_shell( + script, + stdout=asyncio.subprocess.DEVNULL, + stderr=asyncio.subprocess.DEVNULL, + preexec_fn=os.setsid, + ) + + +def _stubborn_child(pid_file: Path, ignore: tuple[str, ...]) -> str: + """Shell snippet that starts a child ignoring `ignore`, then waits for it. + + Killing a PTY session leader makes the kernel send SIGHUP to the + terminal's foreground process group, so a plain `sleep` child looks + contained even when nothing ever signalled the group. A child that ignores + SIGHUP is what an admin actually runs into — a `nohup`ed job, a daemon, + anything meant to outlive its terminal. + + The child publishes its own pid only after installing the handlers, and + the snippet blocks until it does, so a test can never signal it while it + is still starting up and read that as teardown having worked. + """ + ignores = "".join( + f"signal.signal(signal.{name}, signal.SIG_IGN); " for name in ignore + ) + script = ( + f"import os, signal, time; {ignores}" + f"open({str(pid_file)!r}, 'w').write(str(os.getpid())); " + "time.sleep(120)" + ) + return ( + f"{sys.executable} -c {shlex.quote(script)} & " + f"while [ ! -s {pid_file} ]; do sleep 0.02; done" + ) + + +async def _never_disconnected() -> bool: + return False + + +def _reap_if_alive(pid: int) -> None: + """Clean up a descendant the code under test was supposed to have killed.""" + if pid_alive(pid): + try: + os.kill(pid, signal.SIGKILL) + except OSError: + pass + + +async def _read_pid(path: Path, timeout: float = 5.0) -> int: + """Wait for a child to publish its pid, then return it.""" + loop = asyncio.get_running_loop() + deadline = loop.time() + timeout + while loop.time() < deadline: + if path.exists(): + text = path.read_text().strip() + if text: + return int(text) + await asyncio.sleep(0.01) + raise AssertionError(f"child never wrote its pid to {path}") + + +@pty_session +async def test_terminate_pty_session_kills_descendants(tmp_path): + """Tearing down a PTY command takes its children, not only the shell.""" + import routes.shell_routes as shell_routes + + pid_file = tmp_path / "child.pid" + proc = await _spawn_pty_style_session( + f"sleep 120 & echo $! > {pid_file}; sleep 120" + ) + try: + child_pid = await _read_pid(pid_file) + assert pid_alive(child_pid) + + assert await shell_routes._terminate_pty_session(proc) is True + + assert proc.returncode is not None + assert not pid_alive(child_pid) + finally: + await shell_routes._terminate_pty_session(proc) + + +@pty_session +async def test_terminate_pty_session_escalates_past_ignored_sigterm(tmp_path): + """A child that ignores SIGTERM is still gone when teardown returns.""" + import routes.shell_routes as shell_routes + + pid_file = tmp_path / "child.pid" + child = _stubborn_child(pid_file, ("SIGHUP", "SIGTERM")) + proc = await _spawn_pty_style_session(f"{child}; sleep 120") + try: + child_pid = await _read_pid(pid_file) + assert pid_alive(child_pid) + + assert await shell_routes._terminate_pty_session(proc) is True + + assert not pid_alive(child_pid) + finally: + await shell_routes._terminate_pty_session(proc) + + +@pty_session +async def test_generate_pty_timeout_kills_the_whole_session(tmp_path): + """A timed-out PTY command leaves none of its children running.""" + import routes.shell_routes as shell_routes + + pid_file = tmp_path / "child.pid" + child = _stubborn_child(pid_file, ("SIGHUP",)) + cmd = f"{child}; echo ready; sleep 120" + request = SimpleNamespace(is_disconnected=_never_disconnected) + + events = [ + json.loads(chunk.removeprefix("data: ").strip()) + async for chunk in shell_routes._generate_pty(cmd, 1, request) + ] + + child_pid = await _read_pid(pid_file) + try: + assert events[-1] == {"exit_code": -1} + assert events[-2]["data"].startswith("Command timed out after 1s") + assert not pid_alive(child_pid) + finally: + _reap_if_alive(child_pid) + + +@pty_session +async def test_generate_pty_disconnect_kills_the_whole_session(tmp_path): + """Abandoning the stream kills the command's children too.""" + import routes.shell_routes as shell_routes + + pid_file = tmp_path / "child.pid" + child = _stubborn_child(pid_file, ("SIGHUP",)) + cmd = f"{child}; echo ready; sleep 120" + + polls = [] + + async def disconnect_after_first_poll() -> bool: + polls.append(None) + return len(polls) > 1 + + request = SimpleNamespace(is_disconnected=disconnect_after_first_poll) + + async for _ in shell_routes._generate_pty(cmd, 0, request): + pass + + child_pid = await _read_pid(pid_file) + try: + assert not pid_alive(child_pid) + finally: + _reap_if_alive(child_pid) + + +async def test_terminate_pty_session_reports_a_session_it_could_not_kill( + monkeypatch, +): + """Teardown returns False rather than claiming a surviving session died.""" + import routes.shell_routes as shell_routes + + monkeypatch.setattr(shell_routes, "PTY_KILL_GRACE", 0.01) + monkeypatch.setattr(shell_routes, "_signal_session", lambda *_: True) + monkeypatch.setattr(shell_routes, "_session_alive", lambda *_: True) + monkeypatch.setattr(shell_routes, "_session_pgid", lambda _: 4242) + + proc = SimpleNamespace(pid=4242, returncode=0, wait=None) + assert await shell_routes._terminate_pty_session(proc) is False + + +async def test_terminate_pty_session_escalates_before_giving_up(monkeypatch): + """SIGTERM then SIGKILL — the group is never signalled only once.""" + import routes.shell_routes as shell_routes + + sent = [] + monkeypatch.setattr(shell_routes, "PTY_KILL_GRACE", 0.01) + monkeypatch.setattr(shell_routes, "_session_pgid", lambda _: 4242) + monkeypatch.setattr(shell_routes, "_session_alive", lambda *_: True) + monkeypatch.setattr( + shell_routes, + "_signal_session", + lambda pgid, pid, sig: sent.append(sig) or True, + ) + + proc = SimpleNamespace(pid=4242, returncode=0, wait=None) + await shell_routes._terminate_pty_session(proc) + + assert sent == [signal.SIGTERM, signal.SIGKILL] + + +@pty_session +async def test_generate_pty_timeout_says_so_when_the_session_survives( + monkeypatch, +): + """A timed-out command no longer reports clean termination it didn't get.""" + import routes.shell_routes as shell_routes + + real_terminate = shell_routes._terminate_pty_session + + async def terminate_but_report_failure(proc): + await real_terminate(proc) + return False + + monkeypatch.setattr( + shell_routes, "_terminate_pty_session", terminate_but_report_failure + ) + + request = SimpleNamespace(is_disconnected=_never_disconnected) + events = [ + json.loads(chunk.removeprefix("data: ").strip()) + async for chunk in shell_routes._generate_pty("echo ready; sleep 30", 1, request) + ] + + assert events[-1] == {"exit_code": -1} + timed_out = events[-2] + assert timed_out["stream"] == "stderr" + assert timed_out["data"] == ( + "Command timed out after 1s" + shell_routes.PTY_KILL_FAILED_HINT + ) + + class TestFindLineBreak: """Test line-break detection in byte buffers.""" From 2429805a450dd22d4e88fed6a224be866973889c Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?L=C3=A9o?= Date: Thu, 1 Oct 2026 18:46:19 +0200 Subject: [PATCH 3/4] fix(shell): only treat ESRCH as proof a PTY session is gone MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit _session_alive collapsed every OSError from killpg(pgid, 0) into "the group is gone". EPERM means the opposite — the group answered the probe but holds a process we may not signal — so a session we could not touch was reported as contained, and a timed-out command that left children running said it had terminated cleanly. Resolving PTY_KILL_ESCALATION also named signal.SIGKILL unconditionally, which does not exist on native Windows. app.py imports this module at start-up, so that turned a POSIX-only teardown detail into the whole app failing to import there. --- routes/shell_routes.py | 20 ++++++++++-- tests/test_shell_routes.py | 66 ++++++++++++++++++++++++++++++++++++++ 2 files changed, 83 insertions(+), 3 deletions(-) diff --git a/routes/shell_routes.py b/routes/shell_routes.py index 900ac4125..73e055c18 100644 --- a/routes/shell_routes.py +++ b/routes/shell_routes.py @@ -564,7 +564,14 @@ PTY_UNSUPPORTED_ERROR = "pty_unsupported" # PTY teardown. The PTY child leads its own session (os.setsid), so killing it # has to signal the whole process group and then confirm the group is gone — # see _terminate_pty_session. -PTY_KILL_ESCALATION = (signal.SIGTERM, signal.SIGKILL) +# ``signal.SIGKILL`` does not exist on native Windows, and this module is +# imported unconditionally by app.py, so resolve the escalation defensively +# rather than at the cost of the whole app failing to start there. +PTY_KILL_ESCALATION = tuple( + sig + for sig in (getattr(signal, "SIGTERM", None), getattr(signal, "SIGKILL", None)) + if sig is not None +) PTY_KILL_GRACE = 1.0 # seconds a signalled session gets to exit PTY_KILL_POLL_INTERVAL = 0.05 # re-check interval while waiting for it PTY_KILL_FAILED_HINT = "; processes it started survived the kill and are still running" @@ -722,9 +729,16 @@ def _session_alive(pgid: int | None, pid: int) -> bool: if pgid is not None and killpg is not None: try: killpg(pgid, 0) - return True + except ProcessLookupError: + return False # ESRCH — no member of the group is left except OSError: - return False + # Anything else (EPERM when the group holds a process we may not + # signal, EINVAL) answers the probe without proving the group is + # gone. Only ESRCH does that, so treat the rest as still running: + # reporting a surviving session as contained is the one outcome + # teardown must never produce. + return True + return True return pid_alive(pid) diff --git a/tests/test_shell_routes.py b/tests/test_shell_routes.py index 281cc8ea6..69e7b6ab1 100644 --- a/tests/test_shell_routes.py +++ b/tests/test_shell_routes.py @@ -2,6 +2,7 @@ import asyncio import builtins +import errno import importlib import importlib.util import json @@ -64,6 +65,29 @@ def test_shell_routes_import_without_posix_pty_modules(monkeypatch): assert module._find_line_break(b"ok\n") == (2, 1) +def test_shell_routes_import_without_sigkill(monkeypatch): + """Native Windows has no signal.SIGKILL; app.py imports this module anyway. + + The teardown escalation is resolved at import time, so naming SIGKILL + unconditionally would stop the whole app from starting on Windows rather + than only degrading PTY teardown there. + """ + monkeypatch.delattr(signal, "SIGKILL", raising=False) + + module_path = Path(__file__).resolve().parents[1] / "routes" / "shell_routes.py" + spec = importlib.util.spec_from_file_location( + "_shell_routes_without_sigkill", module_path + ) + module = importlib.util.module_from_spec(spec) + sys.modules[spec.name] = module + try: + spec.loader.exec_module(module) + finally: + sys.modules.pop(spec.name, None) + + assert module.PTY_KILL_ESCALATION == (signal.SIGTERM,) + + async def test_generate_pty_reports_explicit_unsupported_error(monkeypatch): """Clients can distinguish unsupported PTY mode from process failures.""" import routes.shell_routes as shell_routes @@ -314,6 +338,48 @@ async def test_generate_pty_timeout_says_so_when_the_session_survives( ) +def test_session_alive_treats_a_refused_probe_as_alive(monkeypatch): + """EPERM says the group exists but we may not signal it, not that it died. + + Only ESRCH proves a process group is gone. Collapsing every OSError into + "gone" is the one error that makes teardown report a surviving session as + contained. + """ + import routes.shell_routes as shell_routes + + def refuse(_pgid, _sig): + raise PermissionError(errno.EPERM, "Operation not permitted") + + monkeypatch.setattr(shell_routes.os, "killpg", refuse) + assert shell_routes._session_alive(4242, 4242) is True + + def gone(_pgid, _sig): + raise ProcessLookupError(errno.ESRCH, "No such process") + + monkeypatch.setattr(shell_routes.os, "killpg", gone) + assert shell_routes._session_alive(4242, 4242) is False + + +async def test_terminate_pty_session_reports_a_group_it_may_not_signal(monkeypatch): + """A session we cannot signal at all is reported as not contained. + + Both the signal and the liveness probe are refused, so teardown has done + nothing and must say so rather than infer death from its own failure. + """ + import routes.shell_routes as shell_routes + + def refuse(*_args): + raise PermissionError(errno.EPERM, "Operation not permitted") + + monkeypatch.setattr(shell_routes, "PTY_KILL_GRACE", 0.01) + monkeypatch.setattr(shell_routes, "_session_pgid", lambda _: 4242) + monkeypatch.setattr(shell_routes.os, "killpg", refuse) + monkeypatch.setattr(shell_routes.os, "kill", refuse) + + proc = SimpleNamespace(pid=4242, returncode=0, wait=None) + assert await shell_routes._terminate_pty_session(proc) is False + + class TestFindLineBreak: """Test line-break detection in byte buffers.""" From 9d0257134f343578b84979a3f1333b9c95fe6f04 Mon Sep 17 00:00:00 2001 From: Alexandre Teixeira <111787685+alteixeira20@users.noreply.github.com> Date: Thu, 1 Oct 2026 18:35:41 +0100 Subject: [PATCH 4/4] feat(runtime): enforce server request authority --- core/database.py | 11 + .../wave-2-request-authority.md | 248 ++++++++ .../wave-2-validation-full.md | 80 +++ routes/chat_routes.py | 39 +- routes/skills_routes.py | 17 +- routes/task/task_routes.py | 13 + src/agent_loop.py | 7 + src/agent_runtime/authority.py | 433 +++++++++++++ src/bg_monitor.py | 13 +- src/clean_agent_preview.py | 6 +- src/task_scheduler.py | 34 +- src/teacher_escalation.py | 3 + src/tool_approvals.py | 10 + src/tool_execution.py | 94 ++- src/tools/system.py | 7 + tests/runtime_evidence_helpers.py | 27 + tests/test_agent_external_tool_schemas.py | 5 + tests/test_ask_user_tool.py | 3 + tests/test_client_tool_routing.py | 3 +- tests/test_edit_file.py | 8 + tests/test_execution_bridge.py | 4 + tests/test_external_context_tool_gate.py | 8 + tests/test_image_creation_routing.py | 6 +- tests/test_request_authority.py | 584 ++++++++++++++++++ tests/test_review_regressions.py | 3 +- tests/test_runtime_evidence_contract.py | 8 + tests/test_task_cookbook_admin_gate.py | 2 +- tests/test_task_scheduler_cancel.py | 5 +- tests/test_tool_approvals.py | 8 + tests/test_tool_path_confinement.py | 8 + tests/test_tool_policy.py | 3 + tests/test_turn_contract.py | 2 + tests/test_turn_contract_integration.py | 8 + tests/test_update_plan_tool.py | 3 + tests/test_weather_search_recovery.py | 3 + tests/test_workspace_confine.py | 4 + website/configuration-reference.md | 2 +- 37 files changed, 1683 insertions(+), 39 deletions(-) create mode 100644 docs/runtime-decomposition/wave-2-request-authority.md create mode 100644 docs/runtime-decomposition/wave-2-validation-full.md create mode 100644 src/agent_runtime/authority.py create mode 100644 tests/test_request_authority.py diff --git a/core/database.py b/core/database.py index 6addc95c4..1621492c7 100644 --- a/core/database.py +++ b/core/database.py @@ -777,6 +777,7 @@ class ScheduledTask(TimestampMixin, Base): owner = Column(String, nullable=True, index=True) name = Column(String, nullable=False, default="Untitled Task") prompt = Column(Text, nullable=True) # LLM prompt (for task_type="llm") + request_authority_json = Column(Text, nullable=True) # server-only admitted request snapshot task_type = Column(String, default="llm") # "llm" | "action" action = Column(String, nullable=True) # builtin action name (for task_type="action") schedule = Column(String, nullable=True) # "once", "daily", "weekly", "monthly" @@ -2335,6 +2336,15 @@ def _migrate_seed_email_account(): # Any future migrations or schema changes that temporarily violate foreign-key # constraints will fail. To perform such operations, foreign_keys must be # temporarily disabled around the migration workflow. +def _migrate_add_task_authority_column(): + """Retain snapshots after legacy task-table rebuilds; support all DBs.""" + from sqlalchemy import inspect + with engine.begin() as conn: + columns = {column["name"] for column in inspect(conn).get_columns("scheduled_tasks")} + if "request_authority_json" not in columns: + conn.execute(text("ALTER TABLE scheduled_tasks ADD COLUMN request_authority_json TEXT")) + + def init_db(): """ Initialize the database by creating all tables. @@ -2412,6 +2422,7 @@ def init_db(): _migrate_add_oauth_config() _migrate_add_email_oauth_columns() _migrate_add_task_automation_columns() + _migrate_add_task_authority_column() _migrate_add_disabled_tools() _migrate_add_mcp_oauth_tokens_column() _migrate_add_task_v2_columns() diff --git a/docs/runtime-decomposition/wave-2-request-authority.md b/docs/runtime-decomposition/wave-2-request-authority.md new file mode 100644 index 000000000..44a72c350 --- /dev/null +++ b/docs/runtime-decomposition/wave-2-request-authority.md @@ -0,0 +1,248 @@ +# Wave 2: request authority + +Base: `d6c3c98c75e03f70c05ebe4058c6fa12e0395f62`, branch +`feature/runtime-request-authority`. Discovery and this plan precede production +changes. No later runtime waves are included. + +## Discovered call paths + +`routes/chat_routes.py` parses mode, toggles, workspace, approval decisions and +runtime context. User intent can promote Chat to Agent. Owner privileges, +global disabled tools, compare/incognito and plan restrictions produce +`ToolPolicy`. Compact/native routes resolve `TurnContract`; regular/full models +can receive the full enabled schema inventory. The route calls +`_stream_agent_with_execution_bridge` and `stream_agent_loop`. Detached runs +retain this generator; reconnecting subscribes to it rather than creating a new +invocation. Their stream IDs are distinct from journal IDs. + +`src/turn_contract.py` classifies request families and selected tools, resolves +exact safe reads, and filters schema availability. Empty-family routing has a +legacy core inventory. Warm tools and editor availability may enlarge offers. +Transcription, OCR and tasks have narrow selection; static web retrieval may +offer private_browser for fallback. These routing choices are not grants. + +`src/agent_loop.py` selects provider/profile transports, parses native or textual +tool blocks, repairs calls, performs deterministic preflights and retries, and +calls `src/tool_execution.py:execute_tool_block`. Compact preview uses +`src/clean_agent_preview.py` but reaches the same dispatcher. The dispatcher +checks run security, exact approval, contract membership, disabled tools, +ToolPolicy, owner restrictions and bridges before MCP/dynamic/built-in handlers. +It forwards policy to dynamic handlers. Legacy loop reconciliation removes +disabled names found in a contract's offered inventory. This must not erase a +request-authority denial. + +Approvals use `src/tool_approvals.py`. A server record binds tool/content, owner, +session, workspace, document id/version/digest, origin run and continuation +state. Consume is destructive; claim is one-use. Task/chat scopes bypass an +existing run-security gate; they do not define the requested operation classes. +Approval continuation executes the sealed action in round zero. Denial exits +the route without execution. + +Generic app_api forwards both the internal token and the caller's owner to +loopback HTTP. Its blocklist does not exclude Chat/skill approval ingress. +Matching owner/session/input bindings alone therefore cannot distinguish a +model-produced HTTP decision from a user approval. Those existing ingress +points need an explicit internal-tool rejection before consuming approval. +Internal HTTP skill-test task bodies likewise cannot mint fresh authority. +The same origin rule applies to generic Chat HTTP entry: a loopback generated +message is not a new trusted user request, even with correct owner attribution. +Both Chat entry points use the existing non-persistence switch for these +messages and append explicitly untrusted transient context instead. Later +referential turns cannot inherit their operation class as prior user intent. + +Teacher takeover is queued by the student, then owned by the outer adapter in +`src/teacher_escalation.py`. It invokes a child loop after the student gate closes +and forwards policy, contract, workspace and runtime context. The teacher's +synthetic user message is model context, not a new authority source. + +`src/task_scheduler.py:_execute_assistant` composes crew/global restrictions and +RAG/default shell availability. `_run_agent_loop` supplies task.prompt or a +synthetic override as a user message, with background provider fallback. Exact +approval pauses are retired because there is no interactive approver. +`_execute_action` invokes BUILTIN_ACTIONS directly, with a separate admin gate. +`src/tools/system.py:do_manage_tasks` and `routes/task/task_routes.py` create/edit +persisted tasks. No authority snapshot currently survives scheduling. + +Detached Bash dispatch launches `bg_jobs.launch` and returns bg_job_id. +`src/bg_monitor.py:_run_followup` appends an explicitly untrusted result to session +context and re-enters the loop. It currently forwards neither the originating +authority nor its request restrictions. Skill tests/audits in +`routes/skills_routes.py` also invoke the loop with task/user messages; generated +audit context must not manufacture grants. + +| Question | Current source | +| --- | --- | +| Requested operation | User intent classifiers, exact safe-read resolver; ultimately parsed/repaired model tool block | +| Available capabilities | Registry/MCP inventory, profiles, RAG, TurnContract and request-specific schema filters | +| Authorized capabilities | Fragmented policy, privileges, run security and approval checks; no independent envelope | +| Restrictions | Route toggles, owner/global policy, plan/compare/incognito, dispatcher owner/workspace checks | +| Approval required | Deterministic run-security decision; model output can propose the action but cannot consume approval | +| Approval input scope | Server-sealed exact tool/content and owner/session/workspace/document binding | +| Nested state | Explicit policy/contract/workspace/context forwarding and journal lineage; no authority snapshot | +| Model influence | Tool/input proposals, repairs, recovery choices, generated task/audit prompts; availability currently participates in execution gating | + +## Implementation plan and contract + +1. Add immutable `ExactOperation`, `OperationGrant` and `RequestAuthority` in + `src/agent_runtime/authority.py`. Normalize canonical tool identity and JSON + inputs (reject duplicate keys/non-finite values); retain exact raw text for + Bash/Python, built-in scheduled actions and non-JSON inputs. Grants contain an operation class/tool identity, optional + action limits and exact input limits. Authority has its own request id, + owner/session/workspace binding, immutable grants and hard denials. It is + independent of schema presence, model/profile, stream/journal/receipt IDs. +2. Create authority from trusted request text/history and deterministic policy + at the chat route before availability reconciliation. The general loop + boundary creates it for other trusted direct callers, without consulting + schemas, relevant_tools, forced_tools or model output. Authority family + inheritance reads only trusted user history. Tool-history exact reads may + narrow an already admitted class, never create a class. Unknown intent grants + no execution floor. Neutral interaction/planning controls remain explicit. +3. Keep semantic classification and availability in TurnContract. Resolve + authority grants separately from those semantic facts and hard policy. + Exact safe reads restrict action/identifiers. Static web fallback authorizes + browser reading/navigation, not arbitrary click/evaluate/form operations. + Media/task families do not inherit the shell inventory. +4. Bind authority around the whole logical stream, including teacher takeover; + forward it explicitly to teacher children and approval records. Children + inherit the parent or intersect explicit authority with it. Policy denials + union; grants intersect; a child cannot replace the parent scope. Restore + the parent on close/error/cancellation. Capture restrictions before legacy + offered-tool reconciliation can erase them. +5. Enforce at `execute_tool_block`, before approvals are claimed or handlers, + bridges/MCP/process dispatch begin. Current policy/disabled gates still win. + Missing/malformed dispatcher state fails closed. Standalone callers/tests + must supply explicit server authority. Journal ownership remains unchanged; + denied calls produce no authoritative execution receipt. +6. Existing approvals remain one-use exact claims. Seal the originating + authority in the approval digest. Resumption keeps original class limits and + current hard restrictions. The approved exact operation may cross its + original class boundary only through the consumed, matching server record + at that call; it does not mutate authority for subsequent calls. Nested + execution cannot use an approval to exceed its parent ceiling. Existing + task/chat UI and run-security scope semantics are unchanged. + Chat/skill approval ingress rejects validated internal-tool requests before + consumption; identity impersonation is not a user approval decision. + A shared HTTP factory admits trusted user requests and produces an empty, + policy-restricted envelope for known internal-tool Chat/skill requests. +7. Persist a server-only authority snapshot and task-input binding on scheduled + records. Direct authenticated task ingress can admit its user-supplied task; + task creation inside model execution intersects with parent authority. + Scheduler overrides, retries and provider fallbacks reuse that snapshot. + Missing/stale snapshots grant no tool authority. Newly seeded server-owned + housekeeping jobs receive exact snapshots at their static creation point; + existing rows are not retrospectively authorized by their names/actions. + Internal tool HTTP task payloads cannot become fresh user requests across an + ASGI context boundary. Built-in actions receive + an exact admission check. Persist detached-job authority in a separate + authority sidecar at dispatch; monitor continuations reuse it and current + denials. Do not edit bg_jobs/process containment implementation. +8. Production files: new authority module; routes/chat_routes.py; + src/agent_loop.py; src/tool_execution.py; src/teacher_escalation.py; + src/tool_approvals.py; core/database.py; routes/task/task_routes.py; + src/tools/system.py; src/task_scheduler.py; src/bg_monitor.py; + routes/skills_routes.py. Change preview only if direct-entry binding is + required by validation. No TurnContract/profile/schema redesign. +9. Shared hotspots: route/loop/dispatch, approvals and task/database integration. + One coordinator writes all production files. Keep changes confined to + authority creation, forwarding, persistence and admission. Do not modify + containment, provenance/effect classification or egress implementation. +10. Focused regressions: available schema/bridge/dynamic handler without grants; + model-selected unrelated tool/action; explicit class admission; exact read + arguments; narrow transcription/OCR/tasks/browser fallback; hard denials + despite offered-tool reconciliation; retry/fallback stability; child and + teacher non-widening and restoration; malformed/missing state; exact + approval mismatch/replay and continuation scope; scheduled snapshot/input + binding and synthetic override; detached followup inheritance; journal + denial evidence. Preserve existing policy-forwarding and Ajax assertions. + +Validation: new focused tests; existing contract/policy/capability/profile +tests; scripts/validate_runtime_wave1.sh; broad affected runtime tests; full +pytest; compileall; JS/MJS syntax; diff check and conflict-marker scan. Any +production edit after full pytest requires affected tests and full pytest again. + +## Implemented boundaries and remaining limits + +The preview entry also binds authority because it supports direct callers. +Research task admission binds the snapshot around the researcher, so nested +execution cannot infer grants from generated research context. LAN lookup +intent has a narrow host_shell-only admission rule; it adds neither Bash nor +Python and does not alter Ajax schemas or profiles. + +Scheduled loop entry explicitly forwards the restored workspace as well as +the envelope; rebinding the continuation session never drops confinement to +the original workspace. Only the actual server Bash launch seals a detached +job sidecar. A handler/bridge result claiming a job id cannot create one. + +Snapshots are trusted server state, stored in the task database and detached +job authority sidecars. Missing, malformed, changed-input, wrong-owner or +wrong-session snapshots fail closed. Legacy tasks need a trusted task-input +save to obtain a snapshot; legacy detached jobs have no execution grants on +followup. No broad backfill, authority-mode UI, containment, effect/egress or +receipt/journal redesign is included. Sidecars follow the detached job's server +storage trust assumptions; retention/integrity hardening is outside this slice. + +Class admission deliberately reuses the deterministic semantic classifiers. +Unrecognized intent has only explicit ask_user/update_plan controls. This can +deny unsupported phrasing and generated default skill tests/audits; model +prompts and tool inventory cannot repair that denial. Existing exact approvals +can admit one sealed root operation, never widen subsequent calls or nested +authority. They still require the existing armed security context, matching +bindings, one-use claim, document checks and current hard restrictions. + +Standalone dispatcher test fixtures now supply explicit registry grants to +continue exercising their original handler/policy/confinement assertions. +New authority tests use the raw dispatcher and prove denial before dispatch. + +## File ownership and reasons + +| Production file | Wave 2 change | +| --- | --- | +| src/agent_runtime/authority.py | Immutable intent/admission/operation API, trusted factory, intersection/context binding, task/job snapshots | +| routes/chat_routes.py | Capture authority before availability reconciliation; pass it into execution; guard approval ingress | +| src/agent_loop.py | Bind logical-invocation authority; capture it in approvals and teacher takeover | +| src/tool_execution.py | Normalize/check operations before dispatch and approval claims; bind handler context; seal actual detached launch | +| src/teacher_escalation.py | Explicit child/approval inheritance without synthetic-prompt grants | +| src/tool_approvals.py | Bind immutable originating authority into exact approval digest | +| src/clean_agent_preview.py | Bind authority at the supported direct preview entry | +| core/database.py | Add nullable server-only scheduled snapshot column and additive migration | +| routes/task/task_routes.py | Seal direct user task inputs; deny fresh grants to internal-tool HTTP payloads | +| src/tools/system.py | Cap model-created/edited task snapshots by active authority | +| src/task_scheduler.py | Restore original scope/workspace for loops, admit exact built-ins/research, seal new static defaults | +| src/bg_monitor.py | Restore original detached-job scope and current hard restrictions | +| routes/skills_routes.py | Separate explicit user task authority from generated/internal skill prompts; guard approval ingress | + +Shared hotspots touched: chat routes, agent loop, central dispatcher, preview, +teacher escalation, approvals, task CRUD/scheduler/system handlers, database, +background monitor and skill entry routes. All production edits have one writer. +TurnContract, tool schemas, model profiles, journal/completion foundations, +bg_jobs/process containment and effect/egress implementations are untouched. + +`tests/test_request_authority.py` adds the focused authority regressions. +`tests/runtime_evidence_helpers.py` adds explicit standalone server fixture +grants. Original assertions are preserved in these adapted fixture suites: + +- tests/test_agent_external_tool_schemas.py +- tests/test_ask_user_tool.py +- tests/test_client_tool_routing.py +- tests/test_edit_file.py +- tests/test_execution_bridge.py +- tests/test_external_context_tool_gate.py +- tests/test_image_creation_routing.py +- tests/test_review_regressions.py +- tests/test_runtime_evidence_contract.py +- tests/test_task_cookbook_admin_gate.py +- tests/test_task_scheduler_cancel.py +- tests/test_tool_approvals.py +- tests/test_tool_path_confinement.py +- tests/test_tool_policy.py +- tests/test_turn_contract.py +- tests/test_turn_contract_integration.py +- tests/test_update_plan_tool.py +- tests/test_weather_search_recovery.py +- tests/test_workspace_confine.py + +`website/configuration-reference.md` is regenerated solely to update the +chat-route environment-read line number. This document records discovery, +the pre-edit plan, implementation boundaries and file ownership. The validation +report records final commands/results. No production files in parallel lanes +are claimed. diff --git a/docs/runtime-decomposition/wave-2-validation-full.md b/docs/runtime-decomposition/wave-2-validation-full.md new file mode 100644 index 000000000..579b12c5a --- /dev/null +++ b/docs/runtime-decomposition/wave-2-validation-full.md @@ -0,0 +1,80 @@ +# Wave 2 final validation + +Worktree: `odysseus-runtime-request-authority`; branch: +`feature/runtime-request-authority`. +Starting SHA: `d6c3c98c75e03f70c05ebe4058c6fa12e0395f62`. +The final SHA is the local commit containing this report, returned in the final +implementation report. No rebase, merge, push or PR was performed. + +All results below apply to the final production code. The last production +changes addressed internal HTTP request/approval origin and transient untrusted +Chat context. Focused, Wave 1.1, broad runtime and full pytest were rerun after +those changes. Subsequent edits only recorded results and removed temporary +validation logs. + +| Gate | Final result | +| --- | --- | +| New Wave 2 authority tests | 58 passed, 1 warning; 1.23s | +| Relevant contract/policy/approval/capability/Ajax/task/background tests | 1500 passed, 28 skipped, 1 warning; 30.55s | +| Wave 1.1 validation script | 2292 passed, 1 warning; 65.75s | +| Broad affected runtime suite | 3079 passed, 28 skipped, 1 warning; 92.83s | +| Full pytest | 11644 passed, 54 skipped, 2 xfailed, 182 warnings, 6 subtests passed; 444.40s | +| Python compileall | Passed | +| JS/MJS syntax | Passed for all 361 tracked files | +| Git whitespace gate | Passed | +| Conflict-marker scan | Passed | + +The existing release smoke hook skipped because `APP_PORT` was unset; no live +instance was driven. Full pytest includes its existing skips and expected +failures. Warnings are retained in the local raw log. Missing development test +dependencies and Playwright Chromium were installed locally, without changing +project dependency declarations. No global dotenv-disable override was used. + +## Commands + +```sh +ODYSSEUS_TEST_STATIC_PORT=0 .venv/bin/python -m pytest -q tests/test_request_authority.py + +ODYSSEUS_TEST_STATIC_PORT=0 .venv/bin/python -m pytest -q tests/test_request_authority.py tests/test_turn_contract*.py tests/test_tool_policy.py tests/test_tool_approval*.py tests/test_execution_capabilities.py tests/test_ajax*.py tests/test_task_*.py tests/test_bg_*.py + +ODYSSEUS_TEST_PYTHON="$PWD/.venv/bin/python" bash scripts/validate_runtime_wave1.sh + +ODYSSEUS_TEST_STATIC_PORT=0 .venv/bin/python -m pytest -q tests/test_request_authority.py tests/test_agent_*.py tests/test_turn_contract*.py tests/test_tool_policy.py tests/test_tool_approval*.py tests/test_task_*.py tests/test_bg_*.py tests/test_*completion*.py tests/test_foreground_model_routing.py tests/test_client_tool_routing.py tests/test_workspace_confine.py tests/test_product_turn_contract_route.py tests/test_execution_bridge.py tests/test_execution_capabilities.py tests/test_ajax*.py tests/test_external_context_tool_gate.py tests/test_tool_path_confinement.py tests/test_edit_file.py tests/test_runtime_evidence_contract.py tests/test_review_regressions.py tests/test_image_creation_routing.py tests/test_ask_user_tool.py tests/test_update_plan_tool.py tests/test_weather_search_recovery.py tests/test_clean_agent_preview.py tests/test_skill_audit*.py tests/test_preview_execution_evidence.py + +ODYSSEUS_TEST_STATIC_PORT=0 .venv/bin/python -m pytest -q + +.venv/bin/python -m compileall -q -x '(^|/)(\.venv|\.git|node_modules|data|logs|uploads)/' . +git ls-files -z '*.js' '*.mjs' | xargs -0 -n 1 node --check +git diff --check +# Staged whitespace check used --cached --check with all 37 changed paths explicit. +git grep --cached -l -E '^(<<<<<<< |=======$|>>>>>>> )' -- '*.py' '*.js' '*.mjs' '*.html' '*.css' '*.json' '*.md' '*.sh' +``` + +Conflict-marker grep returns exit 1 with no matches on success. +The context firewall rejected the unbounded staged whitespace command before +execution; the exact-path check passed. No admitted source inspection was +blocked by staging. +Local raw validation outputs are archived under the ignored +`.venv/wave2-validation/` directory; they are not committed. + +## Regression scope and limits + +The 58 authority tests cover schema/handler/model-selection non-authority, +narrow media/tasks/browser behavior, exact reads, deterministic grants, hard +denials, malformed/missing state, retry and nested inheritance, teacher +forwarding, exact approval scope/replay/digest, scheduled input sealing and +workspace restoration, detached followups and actual-launch-only sealing, +internal HTTP origin, untrusted Chat persistence, and denied-call journal +completion evidence. Existing fixture assertions remain intact; standalone +dispatch fixtures now provide explicit server authority. + +Remaining limits: class admission uses deterministic request classifiers and +can reject unsupported phrasing; legacy task/job snapshots fail closed until +trusted resealing; snapshots assume trusted server database/job storage; +sidecar retention hardening is deferred. Existing approvals can admit one exact +root operation without granting subsequent or nested operations. + +No Wave 3, 3-S, 4, 5 or 6 work was started. No containment, effect/egress, +provenance, authority-mode UI, journal or completion-foundation redesign is +included. File ownership and the discovery/implementation contract are recorded +in [wave-2-request-authority.md](wave-2-request-authority.md). diff --git a/routes/chat_routes.py b/routes/chat_routes.py index 967e36b16..09596ccbc 100644 --- a/routes/chat_routes.py +++ b/routes/chat_routes.py @@ -86,6 +86,7 @@ from src.model_profiles import ( tool_schema_profile, ) from src.tool_execution import AgentExecutionBridge, bind_execution_bridge +from src.agent_runtime.authority import is_internal_tool_request, request_authority_for_http from src.turn_contract import ( FAMILY_TOOLS, bind_turn_contract, preserve_bound_editor_selected_tools, requested_capabilities, resolve_turn_contract, @@ -95,6 +96,15 @@ from src.turn_contract import ( logger = logging.getLogger(__name__) + +def _append_internal_chat_context(ctx, message): + tagged = untrusted_context_message("internal tool request", message) + ctx.messages.append(tagged) + routed = getattr(ctx, "route_messages", None) + if routed is not None and routed is not ctx.messages: + routed.append(tagged) + + # Track active streams for partial-save safety net _active_streams: Dict[str, dict] = {} @@ -2196,7 +2206,10 @@ def setup_chat_routes( webhook_manager=webhook_manager, allow_tool_preprocessing=allow_tool_preprocessing, defer_context_shaping=foreground_policy.enabled, + persist_user_message=not is_internal_tool_request(request), ) + if is_internal_tool_request(request): + _append_internal_chat_context(ctx, message) # Research injection research_blocked_by_policy = ( @@ -2648,6 +2661,8 @@ def setup_chat_routes( ) owner = effective_user(request) if tool_approval_id: + from src.agent_runtime.authority import require_user_approval_request + require_user_approval_request(request) pending_tool_approval = tool_approval_store.peek(tool_approval_id) normalized_owner = str(owner or "").strip().casefold() if ( @@ -2907,10 +2922,12 @@ def setup_chat_routes( and pending_tool_approval.continuation_query else None ), - persist_user_message=not tool_approval_continuation, + persist_user_message=not tool_approval_continuation and not is_internal_tool_request(request), interaction_mode=chat_mode, auto_escalated=auto_escalated, ) + if is_internal_tool_request(request): + _append_internal_chat_context(ctx, message) _research_flags = {"do": do_research} # Mutable container for generator scope @@ -3349,6 +3366,24 @@ def setup_chat_routes( "manage_documents", "create_document", "edit_document", "update_document", }.issubset(disabled_tools), } + # Capture permission state before schema selection/reconciliation. + # Only deterministic request intent supplies grants, never inventory. + _request_authority = request_authority_for_http( + request, message, owner=_user, session_id=session, workspace=workspace, + history=_turn_history, policy=tool_policy, + active_document=bool(active_doc), + image_attachment=any(str(a.get('mime') or '').startswith('image/') + for a in (ctx.preprocessed.attachment_meta or [])), + capabilities=({'search_browser'} if ( + _explicit_browser_intent or _external_discovery_intent + ) else ()), + ) + if exact_tool_approval is not None: + from src.agent_runtime.authority import RequestAuthority + _request_authority = ( + exact_tool_approval.pending.request_authority + or RequestAuthority.empty(owner=_user, session_id=session, workspace=workspace) + ).restrict(tool_policy) _turn_contract = None # Image models execute directly, not through the text-agent inventory. # Keep the permission policy above, but do not apply routing omissions @@ -3477,6 +3512,7 @@ def setup_chat_routes( warm_tools=_warm_tools, message=message, history=getattr(sess, "history", []) or [], ) + _request_authority = _request_authority.restrict(_contract_policy) # Resolution already applies user, owner, and global policy. An # admitted tool must not later be rejected by the stale # pre-contract disabled snapshot during execution. @@ -4415,6 +4451,7 @@ def setup_chat_routes( cwd=_agent_turn_cwd(sess, client_runtime_context), forced_tools=_forced_tools, turn_contract=_turn_contract, + request_authority=_request_authority, uploaded_files=ctx.uploaded_files, defer_context_shaping=_foreground_policy.enabled, external_untrusted_context_seen=external_untrusted_context_seen, diff --git a/routes/skills_routes.py b/routes/skills_routes.py index ef5f65047..ec5a49876 100644 --- a/routes/skills_routes.py +++ b/routes/skills_routes.py @@ -534,11 +534,13 @@ async def _run_skill_test_job( messages=None, transcript=None, exact_approval=None, + request_authority=None, ): """Background coroutine: run the skill in an agent loop, capture a condensed log + transcript, then have the judge grade it. Writes into _skill_test_jobs.""" import json as _json from src.agent_loop import stream_agent_loop + from src.agent_runtime.authority import RequestAuthority job = _skill_test_jobs.get(key) if job is None: @@ -559,6 +561,9 @@ async def _run_skill_test_job( url, model, messages, headers=headers, temperature=0.3, max_tokens=0, max_rounds=8, owner=owner, exact_approval=exact_approval, + request_authority=(request_authority or ( + exact_approval.pending.request_authority if exact_approval is not None else None + ) or RequestAuthority.empty(owner=owner)), ): if not chunk.startswith("data: ") or chunk.strip() == "data: [DONE]": continue @@ -1039,6 +1044,7 @@ async def _run_skill_audit_arm(messages: list[dict], url, model, headers, owner, """Run one audit arm in the agent loop; return transcript, stats, approval.""" import json as _json from src.agent_loop import stream_agent_loop + from src.agent_runtime.authority import RequestAuthority, active_request_authority transcript = [] approval_required = None stats = {"turns": 0, "tool_calls": 0} @@ -1051,6 +1057,7 @@ async def _run_skill_audit_arm(messages: list[dict], url, model, headers, owner, url, model, messages, headers=headers, temperature=0.3, max_tokens=4096, max_rounds=8, owner=owner, workload=workload, suppress_skills=True, + request_authority=(active_request_authority() or RequestAuthority.empty(owner=owner)), ): # Streams can include an SSE event line before the data line, # notably `event: error`. Do not silently discard those failures. @@ -2001,6 +2008,9 @@ def setup_skills_routes(skills_manager: SkillsManager) -> APIRouter: user = _owner(request) body = await request.json() task = (body.get("task") or "").strip() + from src.agent_runtime.authority import RequestAuthority, request_authority_for_http + request_authority = (request_authority_for_http(request, task, owner=user) + if task else RequestAuthority.empty(owner=user)) skills = skills_manager.load(owner=user) match = next((s for s in skills if s.get("name") == skill_id or s.get("id") == skill_id), None) @@ -2064,9 +2074,12 @@ def setup_skills_routes(skills_manager: SkillsManager) -> APIRouter: "model": model, "headers": headers, "owner": user, + "request_authority": request_authority, }, } - _asyncio.create_task(_run_skill_test_job(key, name, md, task, url, model, headers, user, skills_manager)) + _asyncio.create_task(_run_skill_test_job( + key, name, md, task, url, model, headers, user, skills_manager, + request_authority=request_authority)) return {"ok": True, "status": "running", "skill": name, "model": model} @router.post("/{skill_id}/test-approval") @@ -2074,6 +2087,8 @@ def setup_skills_routes(skills_manager: SkillsManager) -> APIRouter: """Resume a manual skill test with one exact server-sealed action.""" import asyncio as _asyncio from src.tool_approvals import tool_approval_store + from src.agent_runtime.authority import require_user_approval_request + require_user_approval_request(request) user = _owner(request) skills = skills_manager.load(owner=user) diff --git a/routes/task/task_routes.py b/routes/task/task_routes.py index c19e73ac9..0749cb548 100644 --- a/routes/task/task_routes.py +++ b/routes/task/task_routes.py @@ -25,6 +25,14 @@ from routes.prefs_routes import _load_for_user, _save_for_user logger = logging.getLogger(__name__) +def _seal_request_task_authority(request, prompt, task_type, action, owner): + from src.agent_runtime.authority import MISSING_AUTHORITY, is_internal_tool_request, seal_task_authority + # A tool HTTP call starts another ASGI context. Its model-produced body is + # not a fresh user request, even though the internal token authenticates it. + parent = None if is_internal_tool_request(request) else MISSING_AUTHORITY + return seal_task_authority(prompt, task_type, action, owner=owner, parent_authority=parent) + + def _maybe_cascade_calendar_event(task) -> None: """Delete the linked calendar event when a cookbook_serve task is removed. Two lookup strategies: @@ -530,6 +538,8 @@ def setup_task_routes(task_scheduler) -> APIRouter: owner=user, name=name, prompt=req.prompt, + request_authority_json=_seal_request_task_authority( + request, req.prompt, req.task_type, req.action, user), task_type=req.task_type, action=req.action, schedule=req.schedule, @@ -737,6 +747,9 @@ def setup_task_routes(task_scheduler) -> APIRouter: task.task_type = req.task_type if req.action is not None: task.action = req.action + if any(value is not None for value in (req.prompt, req.task_type, req.action)): + task.request_authority_json = _seal_request_task_authority( + request, task.prompt, task.task_type, task.action, user) if req.output_target is not None: task.output_target = req.output_target if req.model is not None: diff --git a/src/agent_loop.py b/src/agent_loop.py index 988b454f0..5de35898a 100644 --- a/src/agent_loop.py +++ b/src/agent_loop.py @@ -20375,6 +20375,10 @@ def _blocks_before_inference(turn_contract) -> bool: ) +from src.agent_runtime.authority import MISSING_AUTHORITY, active_request_authority, with_request_authority + + +@with_request_authority @with_turn_contract @with_teacher_takeover @with_completion_gate @@ -20420,6 +20424,7 @@ async def stream_agent_loop( suppress_skills: bool = False, reasoning_effort: Optional[str] = None, _parent_run_id: Optional[str] = None, + request_authority=MISSING_AUTHORITY, ) -> AsyncGenerator[str, None]: """Streaming agent loop generator. @@ -32758,6 +32763,7 @@ async def stream_agent_loop( block.tool_type, block.content ), request_text=_last_user, + request_authority=active_request_authority(), ) desc = f"{block.tool_type}: APPROVAL REQUIRED" result = { @@ -37372,6 +37378,7 @@ async def stream_agent_loop( active_document=active_document, active_email=active_email, turn_contract=turn_contract, + request_authority=active_request_authority(), external_untrusted_context_seen=run_security.external_untrusted_context_seen, client_runtime_context=client_runtime_context, plan_mode=plan_mode, diff --git a/src/agent_runtime/authority.py b/src/agent_runtime/authority.py new file mode 100644 index 000000000..ef8552903 --- /dev/null +++ b/src/agent_runtime/authority.py @@ -0,0 +1,433 @@ +"""Server-owned request admission, independent of model tool availability.""" +from __future__ import annotations + +from contextlib import aclosing, contextmanager +from contextvars import ContextVar +from dataclasses import dataclass, replace +from functools import wraps +from inspect import signature +import json +from pathlib import Path +import re +from uuid import uuid4 + +from src.tool_policy import ToolPolicy, build_effective_tool_policy +from src.turn_contract import ( + FAMILY_TOOLS, canonical_tool, requested_capabilities, + RequiredReadOperation, required_read_operation_for_request, selected_tools_for_request, +) + + +def _owner(value): + return str(value or "").strip().casefold() + + +def _pairs(pairs): + result = {} + for key, value in pairs: + if key in result: + raise ValueError("Duplicate operation argument") + result[key] = value + return result + + +def _invalid_constant(value): + raise ValueError("Non-finite operation argument") + + +def _json(value): + return json.dumps(value, sort_keys=True, separators=(",", ":"), allow_nan=False) + + +@dataclass(frozen=True) +class ExactOperation: + tool: str + input: str + action: str | None = None + transport_tool: str = "" + + @classmethod + def normalize(cls, tool, content): + if not isinstance(tool, str) or not tool.strip() or not isinstance(content, str): + raise ValueError("Operation requires a tool name and string input") + transport_tool = tool.strip() + tool = canonical_tool(transport_tool) + normalized = content + payload = None + raw_input = tool in {"bash", "python"} or tool.startswith("scheduled__") + if not raw_input and content.lstrip().startswith("{"): + payload = json.loads(content, object_pairs_hook=_pairs, parse_constant=_invalid_constant) + if not isinstance(payload, dict): + raise ValueError("Structured tool input must be an object") + normalized = _json(payload) + # Reuse the runtime's existing multiplexed-action normalization; this + # classifies input and never grants permission or changes the input. + from src.tool_capabilities import _action_from_content + action = _action_from_content(tool, content) + if tool == "private_browser" and isinstance(payload, dict): + action = payload.get("action") + if action is not None and not isinstance(action, str): + raise ValueError("Browser action must be a string") + action = action.strip().casefold() if action else None + return cls(tool, normalized, action, transport_tool) + + +@dataclass(frozen=True) +class OperationGrant: + tool: str + actions: frozenset[str] | None = None + inputs: frozenset[str] | None = None + + def __post_init__(self): + if not isinstance(self.tool, str) or not self.tool or canonical_tool(self.tool) != self.tool: + raise ValueError("Grant requires a canonical tool identity") + for values in (self.actions, self.inputs): + if values is not None and (not isinstance(values, frozenset) + or any(not isinstance(v, str) for v in values)): + raise TypeError("Grant limits must be immutable string sets") + + def permits(self, operation): + return (self.tool == operation.tool + and (self.actions is None or operation.action in self.actions) + and (self.inputs is None or operation.input in self.inputs)) + + def intersect(self, other): + if self.tool != other.tool: + raise ValueError("Cannot intersect different operation classes") + def limits(left, right): + return right if left is None else left if right is None else left & right + return OperationGrant(self.tool, limits(self.actions, other.actions), + limits(self.inputs, other.inputs)) + + +@dataclass(frozen=True) +class RequestAuthority: + request_id: str + owner: str + session_id: str + workspace: str + grants: tuple[OperationGrant, ...] = () + denied: frozenset[str] = frozenset() + block_all: bool = False + disable_mcp: bool = False + inherited: bool = False + + def __post_init__(self): + if (not isinstance(self.request_id, str) or not self.request_id + or any(not isinstance(v, str) for v in (self.owner, self.session_id, self.workspace)) + or not isinstance(self.grants, tuple) + or any(not isinstance(g, OperationGrant) for g in self.grants) + or len({g.tool for g in self.grants}) != len(self.grants) + or not isinstance(self.denied, frozenset) + or any(not isinstance(n, str) or canonical_tool(n) != n for n in self.denied) + or any(type(v) is not bool for v in (self.block_all, self.disable_mcp, self.inherited))): + raise ValueError("Malformed request authority") + + @classmethod + def empty(cls, *, owner=None, session_id=None, workspace=None): + return cls(uuid4().hex, _owner(owner), str(session_id or ""), str(workspace or "")) + + def bound_to(self, *, owner=None, session_id=None, workspace=None): + return (self.owner == _owner(owner) and self.session_id == str(session_id or "") + and self.workspace == str(workspace or "")) + + def restricted(self, operation): + return (self.block_all or operation.tool in self.denied + or (self.disable_mcp and (operation.tool.startswith("mcp__") + or operation.transport_tool.startswith("mcp__")))) + + def permits(self, operation): + return not self.restricted(operation) and any(g.permits(operation) for g in self.grants) + + def restrict(self, policy=None, disabled_tools=()): + policy = policy or ToolPolicy() + return replace(self, denied=self.denied | frozenset( + canonical_tool(n) for n in set(disabled_tools or ()) | policy.all_disabled_names()), + block_all=self.block_all or policy.block_all_tool_calls, + disable_mcp=self.disable_mcp or policy.disable_mcp) + + def intersect(self, child): + if not isinstance(child, RequestAuthority): + raise TypeError("Child authority must be server-owned RequestAuthority") + grants = [] + if (self.owner, self.session_id, self.workspace) == (child.owner, child.session_id, child.workspace): + theirs = {g.tool: g for g in child.grants} + grants = [g.intersect(theirs[g.tool]) for g in self.grants if g.tool in theirs] + return replace(self, grants=tuple(grants), denied=self.denied | child.denied, + block_all=self.block_all or child.block_all, + disable_mcp=self.disable_mcp or child.disable_mcp, inherited=True) + + def continuation(self, *, owner=None, session_id=None): + """A server continuation may rebind a session, never change owner/grants.""" + if self.owner != _owner(owner): + return RequestAuthority.empty(owner=owner, session_id=session_id) + return replace(self, session_id=str(session_id or ""), inherited=True) + + def to_dict(self): + return {"version": 1, "request_id": self.request_id, "owner": self.owner, + "session_id": self.session_id, "workspace": self.workspace, + "grants": [{"tool": g.tool, + "actions": None if g.actions is None else sorted(g.actions), + "inputs": None if g.inputs is None else sorted(g.inputs)} for g in self.grants], + "denied": sorted(self.denied), "block_all": self.block_all, + "disable_mcp": self.disable_mcp, "inherited": self.inherited} + + @classmethod + def from_dict(cls, value): + if (not isinstance(value, dict) or type(value.get("version")) is not int + or value["version"] != 1): + raise ValueError("Unsupported authority snapshot") + def limits(value): + if value is None: + return None + if not isinstance(value, list) or any(not isinstance(v, str) for v in value): + raise ValueError("Malformed authority limits") + return frozenset(value) + return cls(value["request_id"], value["owner"], value["session_id"], value["workspace"], + tuple(OperationGrant(g["tool"], limits(g["actions"]), limits(g["inputs"])) + for g in value["grants"]), limits(value["denied"]), + value["block_all"], value["disable_mcp"], value["inherited"]) + + +_BROWSER_READ_ACTIONS = frozenset({"open", "navigate", "snapshot", "text", "read", "find", + "screenshot", "scroll", "back", "forward", "wait", "status", "close", "tabs"}) + + +@dataclass(frozen=True) +class SemanticIntent: + """Routing facts, with no execution permission or provider inventory.""" + capabilities: frozenset[str] + selected_tools: frozenset[str] | None + required_read: RequiredReadOperation | None + + +def interpret_request(request_text, *, history=(), workspace=None, active_document=False, + image_attachment=False): + if not isinstance(request_text, str): + raise TypeError("Intent requires request text") + # Routing may use model/tool history. Admission may only inherit intent + # from trusted user requests; a model's proposal or attempted tool call + # cannot establish a new authorized operation class. + history = tuple(history or ()) + trusted_history = [] + for row in history: + get = row.get if isinstance(row, dict) else lambda key, default=None: getattr(row, key, default) + metadata = get("metadata") or {} + if isinstance(metadata, str): + try: + metadata = json.loads(metadata) + except ValueError: + metadata = {} + if (get("role") == "user" and isinstance(metadata, dict) + and metadata.get("trusted") is not False and not metadata.get("tool_gate_untrusted")): + trusted_history.append({"role": "user", "content": get("content", "")}) + families = requested_capabilities(request_text, trusted_history, + active_document=active_document, workspace=bool(workspace), image_attachment=image_attachment) + selected = selected_tools_for_request(request_text) + if (families <= {"unknown"} and selected is None + and re.search(r"\b(?:lan|local\s+(?:network|ip)|tailscale|arp|ip\s+route|default\s+route|subnet|network\s+interface|neighbor\s+table|wifi|ethernet)\b", request_text, re.I) + and re.search(r"\b(?:find|check|inspect|show|list|lookup|locate)\b", request_text, re.I) + and not re.search(r"\b(?:web|internet|online)\b", request_text, re.I)): + # An explicit local-network lookup is a host operation. The existing + # router already chooses host_shell; neither its schema nor bridge + # availability grants Bash/Python alongside this request. + return SemanticIntent(frozenset({"shell_files"}), frozenset({"host_shell"}), None) + return SemanticIntent(families, selected, required_read_operation_for_request(request_text, history)) + + +def create_request_authority(request_text, *, owner=None, session_id=None, workspace=None, + history=(), policy=None, active_document=False, + image_attachment=False, capabilities=None): + """Deterministic server policy over semantic facts, never schema inventory.""" + if not isinstance(request_text, str): + raise TypeError("Authority requires trusted request text") + intent = interpret_request(request_text, history=history, active_document=active_document, + workspace=workspace, image_attachment=image_attachment) + families = intent.capabilities + if capabilities is not None: + families |= frozenset(capabilities) + tools = set().union(*(FAMILY_TOOLS.get(f, ()) for f in families)) + selected = intent.selected_tools + if selected is not None: + tools = tools & set(selected) if families else set(selected) + if tools & {"web_search", "web_fetch"}: + tools.add("private_browser") + operation = intent.required_read + if operation is not None and canonical_tool(operation.tool) in {canonical_tool(n) for n in tools}: + tools = {canonical_tool(operation.tool)} + else: + operation = None + grants = [] + for name in sorted(tools | {"ask_user", "update_plan"}): + name = canonical_tool(name) + actions = inputs = None + if name == "private_browser": + actions = _BROWSER_READ_ACTIONS + # Explicit interaction intent admits its operation class. A + # browser offered only as static-Web fallback gets no such grant. + if re.search(r"\b(?:browser|browse|private_browser)\b", request_text, re.I): + actions |= frozenset(action for action in ("click", "fill", "type", "press", "evaluate", "select") + if re.search(r"\b" + action + r"\b", request_text, re.I)) + if operation is not None and name == canonical_tool(operation.tool): + inputs = frozenset({ExactOperation.normalize(name, _json(dict(operation.args))).input}) + grants.append(OperationGrant(name, actions, inputs)) + authority = RequestAuthority(uuid4().hex, _owner(owner), str(session_id or ""), + str(workspace or ""), tuple(grants)) + return authority.restrict(policy or build_effective_tool_policy(last_user_message=request_text)) + + +_ACTIVE: ContextVar[RequestAuthority | None] = ContextVar("request_authority", default=None) +MISSING_AUTHORITY = object() + + +def active_request_authority(): + return _ACTIVE.get() + + +def is_internal_tool_request(request): + """HTTP authentication/owner attribution does not make a tool payload user intent.""" + from core.middleware import INTERNAL_TOOL_HEADER, INTERNAL_TOOL_TOKEN + return (request.headers.get(INTERNAL_TOOL_HEADER) == INTERNAL_TOOL_TOKEN + or getattr(request.state, "current_user", None) == "internal-tool") + + +def require_user_approval_request(request): + if is_internal_tool_request(request): + from fastapi import HTTPException + raise HTTPException(403, "Tool requests cannot submit user approval decisions.") + + +def request_authority_for_http(request, request_text, **context): + """Known tool loopback is a continuation, never a fresh user grant source.""" + if is_internal_tool_request(request): + return RequestAuthority.empty(owner=context.get("owner"), + session_id=context.get("session_id"), workspace=context.get("workspace")).restrict(context.get("policy")) + return create_request_authority(request_text, **context) + + +@contextmanager +def bind_request_authority(authority): + if not isinstance(authority, RequestAuthority): + raise TypeError("Authority must be server-owned RequestAuthority") + parent = _ACTIVE.get() + authority = parent.intersect(authority) if parent is not None else authority + token = _ACTIVE.set(authority) + try: + yield authority + finally: + _ACTIVE.reset(token) + + +def _request_text(messages): + for message in reversed(messages or ()): + metadata = message.get("metadata") or {} + if (message.get("role") != "user" or metadata.get("trusted") is False + or metadata.get("tool_gate_untrusted")): + continue + content = message.get("content", "") + if isinstance(content, str): + return content + if isinstance(content, list): + return "\n".join(p.get("text", "") for p in content + if isinstance(p, dict) and p.get("type") == "text") + return "" + + +def with_request_authority(func): + """Bind once per invocation; model rounds/fallbacks never recreate grants.""" + call_signature = signature(func) + @wraps(func) + async def wrapped(*args, **kwargs): + bound = call_signature.bind(*args, **kwargs) + bound.apply_defaults() + parameters = bound.arguments + parent = active_request_authority() + authority = parameters.get("request_authority", MISSING_AUTHORITY) + if authority is MISSING_AUTHORITY: + approval = parameters.get("exact_approval") + if parent is not None: + authority = parent + elif approval is not None: + authority = approval.pending.request_authority or RequestAuthority.empty( + owner=parameters.get("owner"), session_id=parameters.get("session_id"), + workspace=parameters.get("workspace")) + elif (parameters.get("_parent_run_id") or parameters.get("_is_teacher_run") + or parameters.get("workload") == "background"): + authority = RequestAuthority.empty(owner=parameters.get("owner"), + session_id=parameters.get("session_id"), workspace=parameters.get("workspace")) + else: + authority = create_request_authority(_request_text(parameters.get("messages")), + owner=parameters.get("owner"), session_id=parameters.get("session_id"), + workspace=parameters.get("workspace"), + history=getattr(parameters.get("history_session"), "history", ()) or (), + active_document=bool(parameters.get("active_document"))) + if not isinstance(authority, RequestAuthority): + raise TypeError("Missing or malformed server request authority") + if parent is None and parameters.get("exact_approval") is not None: + authority = replace(authority, inherited=False) + authority = authority.restrict(parameters.get("tool_policy"), parameters.get("disabled_tools")) + with bind_request_authority(authority) as effective: + if "request_authority" in parameters: + parameters["request_authority"] = effective + async with aclosing(func(*bound.args, **bound.kwargs)) as stream: + async for chunk in stream: + yield chunk + return wrapped + + +def task_operation(task_type, action, prompt): + if task_type == "action": + tool = ("bash" if action in {"run_local", "run_script", "ssh_command"} + else "serve_model" if action == "cookbook_serve" else "scheduled__" + str(action)) + return ExactOperation.normalize(tool, str(prompt or "")) + if task_type == "research": + return ExactOperation.normalize("trigger_research", str(prompt or "")) + return None + + +def seal_task_authority(prompt, task_type, action, *, owner=None, parent_authority=MISSING_AUTHORITY): + """Only direct ingress grants; a model-created task is capped by its parent.""" + operation = task_operation(task_type, action, prompt) + authority = create_request_authority(str(prompt or ""), owner=owner) + if operation is not None: + authority = replace(authority, grants=(OperationGrant(operation.tool, + inputs=frozenset({operation.input})),)) + parent = active_request_authority() if parent_authority is MISSING_AUTHORITY else parent_authority + if parent_authority is None: + parent = RequestAuthority.empty(owner=owner) + if parent is not None: + authority = parent.intersect(replace(authority, session_id=parent.session_id, + workspace=parent.workspace)) + return _json({"task_input": [prompt, task_type, action], "authority": authority.to_dict()}) + + +def restore_task_authority(snapshot, prompt, task_type, action, *, owner=None, session_id=None): + try: + value = json.loads(snapshot) + if value["task_input"] != [prompt, task_type, action]: + raise ValueError("Scheduled request changed") + return RequestAuthority.from_dict(value["authority"]).continuation(owner=owner, session_id=session_id) + except (ValueError, TypeError, KeyError, AttributeError): + return RequestAuthority.empty(owner=owner, session_id=session_id) + + +def _background_path(job_id): + if not isinstance(job_id, str) or not re.fullmatch(r"[A-Za-z0-9_-]+", job_id): + raise ValueError("Invalid background authority identity") + from src.constants import BG_JOBS_DIR + return Path(BG_JOBS_DIR) / (job_id + ".authority.json") + + +def save_background_authority(job_id, authority): + from core.atomic_io import atomic_write_json + atomic_write_json(_background_path(job_id), authority.to_dict()) + + +def restore_background_authority(job_id, *, owner=None, session_id=None): + try: + authority = RequestAuthority.from_dict(json.loads(_background_path(job_id).read_text())) + if authority.session_id != str(session_id or ""): + raise ValueError("Background session changed") + return authority.continuation(owner=owner, session_id=session_id) + except (OSError, ValueError, TypeError, KeyError, AttributeError): + return RequestAuthority.empty(owner=owner, session_id=session_id) diff --git a/src/bg_monitor.py b/src/bg_monitor.py index 2c17c3a1b..086faae19 100644 --- a/src/bg_monitor.py +++ b/src/bg_monitor.py @@ -36,11 +36,12 @@ def _background_result_message(rec): return untrusted_context_message("background job output", inject) -async def _drain_agent(sess, messages): +async def _drain_agent(sess, messages, request_authority=None): """Run the agent loop headless against a session. Returns (final_prose, tool_events) — tool_events in the same shape the live chat saves, so the frontend rebuilds them as standard agent-thread tool cards.""" from src.agent_loop import stream_agent_loop + from src.agent_runtime.authority import RequestAuthority full = "" final_replaced = False tool_events = [] @@ -52,6 +53,9 @@ async def _drain_agent(sess, messages): session_id=sess.id, max_rounds=_FOLLOWUP_MAX_ROUNDS, owner=getattr(sess, "owner", None), + workspace=request_authority.workspace or None if request_authority is not None else None, + request_authority=(request_authority or RequestAuthority.empty( + owner=getattr(sess, "owner", None), session_id=sess.id)), ): if not chunk.startswith("data: "): continue @@ -132,7 +136,12 @@ async def _run_followup(rec: dict) -> bool: context = sess.get_context_messages() context.append(_background_result_message(rec)) - full, tool_events = await _drain_agent(sess, context) + from src.agent_runtime.authority import restore_background_authority + from src.settings import get_setting + authority = restore_background_authority( + rec["id"], owner=getattr(sess, "owner", None), session_id=sess.id) + authority = authority.restrict(disabled_tools=get_setting("disabled_tools", []) or ()) + full, tool_events = await _drain_agent(sess, context, request_authority=authority) # Persist ONLY the assistant continuation so it renders as a normal agent # turn — a standard chat bubble plus `tool_events` that the frontend diff --git a/src/clean_agent_preview.py b/src/clean_agent_preview.py index 83a448db1..4c0302eb7 100644 --- a/src/clean_agent_preview.py +++ b/src/clean_agent_preview.py @@ -4984,6 +4984,10 @@ async def preview_lines_until_finish(response, finish_event=None): await asyncio.gather(finish_task, return_exceptions=True) +from src.agent_runtime.authority import MISSING_AUTHORITY, with_request_authority + + +@with_request_authority async def stream_preview(*, endpoint_url, model, messages, headers, turn_contract, session_id, owner, disabled_tools, tool_policy, history_session=None, external_untrusted_context_seen=False, @@ -4991,7 +4995,7 @@ async def stream_preview(*, endpoint_url, model, messages, headers, turn_contrac client_runtime_context=None, max_tokens=768, max_rounds=8, max_tool_calls=0, external_tool_schemas=None, temperature=0.0, - **ignored): + request_authority=MISSING_AUTHORITY, **ignored): from src.generation_sampling import validate_temperature temperature = validate_temperature(temperature) # This path sends requests directly with httpx and therefore bypasses diff --git a/src/task_scheduler.py b/src/task_scheduler.py index ce0103f48..02fa9970f 100644 --- a/src/task_scheduler.py +++ b/src/task_scheduler.py @@ -1316,6 +1316,17 @@ class TaskScheduler: async def _execute_action(self, task, run_id: str | None = None) -> tuple: """Execute a built-in action (no LLM needed).""" from src.builtin_actions import BUILTIN_ACTIONS + from src.agent_runtime.authority import ( + bind_request_authority, restore_task_authority, task_operation, + ) + authority = restore_task_authority( + getattr(task, "request_authority_json", None), task.prompt, task.task_type, + task.action, owner=task.owner) + from src.settings import get_setting + authority = authority.restrict(disabled_tools=get_setting("disabled_tools", []) or ()) + operation = task_operation(task.task_type, task.action, task.prompt) + if operation is None or not authority.permits(operation): + return "Scheduled action has no matching server request authority.", False action_fn = BUILTIN_ACTIONS.get(task.action) if not action_fn: @@ -1342,7 +1353,8 @@ class TaskScheduler: if getattr(task, "model", None): kwargs["model"] = task.model kwargs["endpoint_url"] = getattr(task, "endpoint_url", None) - result, success = await action_fn(**kwargs) + with bind_request_authority(authority): + result, success = await action_fn(**kwargs) if getattr(task, "model", None): self._last_run_model = task.model return result, success @@ -1945,6 +1957,7 @@ class TaskScheduler: datetime_context_msg: dict | None = None) -> str: """Run the full agent loop with tool access, collecting the final text.""" from src.agent_loop import stream_agent_loop + from src.agent_runtime.authority import restore_task_authority system_content = system_prompt or "You are a helpful assistant executing a scheduled task. Use available tools to complete the task thoroughly." user_content = override_user_message or task.prompt @@ -2000,6 +2013,10 @@ class TaskScheduler: _task_fallbacks = [] # Close the stream in this task on every exit, including the # approval-pause break, so the agent run's context state unwinds here. + request_authority = restore_task_authority( + getattr(task, "request_authority_json", None), task.prompt, + getattr(task, "task_type", "llm"), getattr(task, "action", None), + owner=task.owner, session_id=session_id) async with contextlib.aclosing(stream_agent_loop( endpoint_url=endpoint_url, model=model, @@ -2007,11 +2024,13 @@ class TaskScheduler: max_rounds=_task_max_rounds, session_id=session_id, owner=task.owner, + workspace=request_authority.workspace or None, headers=headers, disabled_tools=disabled_tools, relevant_tools=relevant_tools, fallbacks=_task_fallbacks, workload="background", + request_authority=request_authority, )) as agent_stream: async for event_str in agent_stream: if event_str.startswith("data: ") and not event_str.startswith("data: [DONE]"): @@ -2109,6 +2128,14 @@ class TaskScheduler: async def _execute_research_task(self, task, db) -> str: """Execute a deep research task using DeepResearcher.""" + from src.agent_runtime.authority import bind_request_authority, restore_task_authority, task_operation + from src.settings import get_setting + authority = restore_task_authority( + getattr(task, "request_authority_json", None), task.prompt, task.task_type, + getattr(task, "action", None), owner=task.owner) + authority = authority.restrict(disabled_tools=get_setting("disabled_tools", []) or ()) + if not authority.permits(task_operation(task.task_type, getattr(task, "action", None), task.prompt)): + raise PermissionError("Scheduled research has no matching server request authority.") from core.database import Session as DbSession, ChatMessage from src.deep_research import DeepResearcher from src.research_handler import RESEARCH_DATA_DIR, ResearchHandler @@ -2180,7 +2207,8 @@ class TaskScheduler: ) started_ts = time.time() - report = await researcher.research(task.prompt) + with bind_request_authority(authority): + report = await researcher.research(task.prompt) completed_ts = time.time() try: stats = researcher.get_stats() or {} @@ -2603,6 +2631,7 @@ class TaskScheduler: if (task.output_target or "session") == "session": task.output_target = defs.get("output_target", "none") seeded = [] + from src.agent_runtime.authority import seal_task_authority for action, defs in HOUSEKEEPING_DEFAULTS.items(): if action in existing_actions: continue @@ -2620,6 +2649,7 @@ class TaskScheduler: name=defs["name"], task_type="action", action=action, + request_authority_json=seal_task_authority(None, "action", action, owner=owner), trigger_type=trigger_type, trigger_event=defs.get("trigger_event"), trigger_count=defs.get("trigger_count"), diff --git a/src/teacher_escalation.py b/src/teacher_escalation.py index 1f646b583..fda3f4cec 100644 --- a/src/teacher_escalation.py +++ b/src/teacher_escalation.py @@ -585,6 +585,7 @@ async def run_teacher_inline( external_untrusted_context_seen: bool = False, client_runtime_context: Optional[Dict[str, Any]] = None, plan_mode: bool = False, + request_authority=None, ): """Async generator. Yields SSE event strings. @@ -700,6 +701,7 @@ async def run_teacher_inline( active_email=active_email, turn_contract=turn_contract, _parent_run_id=parent_run_id, + request_authority=request_authority, external_untrusted_context_seen=external_untrusted_context_seen, client_runtime_context=deepcopy(client_runtime_context), plan_mode=plan_mode, @@ -822,6 +824,7 @@ async def run_teacher_inline( workspace=workspace, external_untrusted_context_seen=True, capabilities=capabilities_for_action("manage_skills", skill_content), + request_authority=request_authority, ) approval = pending.public_payload( reason=( diff --git a/src/tool_approvals.py b/src/tool_approvals.py index 5e688172b..7144fc3cf 100644 --- a/src/tool_approvals.py +++ b/src/tool_approvals.py @@ -25,6 +25,7 @@ from src.tool_approval_scopes import ( scope_for_decision, ) from src.tool_capabilities import ToolCapabilities, capabilities_for_action +from src.agent_runtime.authority import RequestAuthority DEFAULT_APPROVAL_TTL_SECONDS = 10 * 60 @@ -117,6 +118,7 @@ def _binding_payload( continuation_query: Any, effects: tuple[str, ...], result_integrity: str, + request_authority: RequestAuthority | None = None, ) -> dict[str, Any]: return { "owner": _normalized_owner(owner), @@ -137,6 +139,7 @@ def _binding_payload( "continuation_query": _normalized_continuation_query(continuation_query), "effects": list(effects), "result_integrity": str(result_integrity), + "request_authority": request_authority.to_dict() if request_authority is not None else None, } @@ -165,6 +168,7 @@ class PendingToolApproval: # The originating user request is internal continuation context only; it # is never displayed or treated as authorization for the sealed action. request_text: str = "" + request_authority: RequestAuthority | None = None def public_payload(self, *, reason: str | None = None) -> dict[str, Any]: return { @@ -273,6 +277,7 @@ class ExactToolApproval: continuation_query=self.pending.continuation_query, effects=effects, result_integrity=result_integrity, + request_authority=self.pending.request_authority, ) return _canonical_digest(expected) == self.pending.digest @@ -356,7 +361,10 @@ class ToolApprovalStore: external_untrusted_context_seen: bool, capabilities: ToolCapabilities, request_text: Any = "", + request_authority: RequestAuthority | None = None, ) -> PendingToolApproval: + if request_authority is not None and not isinstance(request_authority, RequestAuthority): + raise TypeError("Approval authority must be server-owned RequestAuthority") now = time.time() effects = tuple(sorted(effect.value for effect in capabilities.effects)) result_integrity = capabilities.result_integrity.value @@ -375,6 +383,7 @@ class ToolApprovalStore: continuation_query=continuation_query, effects=effects, result_integrity=result_integrity, + request_authority=request_authority, ) pending = PendingToolApproval( approval_id=secrets.token_urlsafe(32), @@ -398,6 +407,7 @@ class ToolApprovalStore: selected_tools=tuple(payload["selected_tools"]), continuation_query=payload["continuation_query"], request_text=str(request_text or ""), + request_authority=request_authority, ) with self._lock: self._purge_expired_locked(now) diff --git a/src/tool_execution.py b/src/tool_execution.py index e72594cba..b9cd65c61 100644 --- a/src/tool_execution.py +++ b/src/tool_execution.py @@ -1219,6 +1219,7 @@ async def _direct_fallback( "client_runtime_context": client_runtime_context, "disabled_tools": frozenset(disabled_tools or ()), "tool_policy": tool_policy, + "request_authority": active_request_authority(), } from src.agent_tools import TOOL_HANDLERS @@ -1259,6 +1260,10 @@ async def _document_tool_dispatch( # --------------------------------------------------------------------------- from src.agent_runtime.journal import dispatched, mark_authorized, mark_dispatch, record_action +from src.agent_runtime.authority import ( + MISSING_AUTHORITY, ExactOperation, RequestAuthority, active_request_authority, + bind_request_authority, save_background_authority, +) @record_action @@ -1278,6 +1283,7 @@ async def execute_tool_block( exact_approval: Optional[ExactToolApproval] = None, active_document_id: Optional[str] = None, client_runtime_context: Optional[Dict[str, Any]] = None, + request_authority=MISSING_AUTHORITY, ) -> Tuple[str, Dict]: """Execute a single tool block. Returns (description, result_dict). @@ -1299,6 +1305,40 @@ async def execute_tool_block( "NO_TOOL_SECURITY_CONTEXT" ) + authority = active_request_authority() if request_authority is MISSING_AUTHORITY else request_authority + parent = active_request_authority() + if isinstance(authority, RequestAuthority) and parent is not None and authority is not parent: + authority = parent.intersect(authority) + try: + operation = ExactOperation.normalize(getattr(block, "tool_type", None), getattr(block, "content", None)) + valid = isinstance(authority, RequestAuthority) and authority.bound_to( + owner=owner, session_id=session_id, workspace=workspace) + if valid: + authority = authority.restrict(tool_policy, disabled_tools) + exact_admission = bool( + valid and not authority.inherited and not authority.restricted(operation) + and exact_approval is not None and exact_approval.matches( + owner=owner, session_id=session_id, workspace=workspace, + tool_name=getattr(block, "tool_type", None), content=getattr(block, "content", None))) + admitted = valid and (authority.permits(operation) or exact_admission) + except (ValueError, TypeError, AttributeError) as error: + return f"{getattr(block, 'tool_type', '')}: invalid arguments", { + "error": (f"Tool arguments are not valid JSON: {error}" + if isinstance(error, json.JSONDecodeError) else str(error)), + "exit_code": 1, "blocked": True, + "failure_kind": "request_authority_denied", + } + if not admitted: + reason = "The exact operation is outside server request authority." + if tool_policy and any(tool_policy.blocks(name) for name in email_tool_policy_names(getattr(block, "tool_type", ""))): + reason = f"Execution of tool '{getattr(block, 'tool_type', '')}' is forbade by the active tool policy." + elif isinstance(authority, RequestAuthority) and authority.restricted(operation): + reason = "The exact operation is disabled by user or server request authority policy." + return f"{getattr(block, 'tool_type', '')}: BLOCKED", { + "error": reason, + "exit_code": 1, "blocked": True, "failure_kind": "request_authority_denied", + } + from src.turn_contract import active_turn_contract contract = active_turn_contract() if contract is not None and not contract.permits(getattr(block, "tool_type", "")): @@ -1393,31 +1433,32 @@ async def execute_tool_block( token = _active_workspace.set(workspace or None) try: - output = await _execute_tool_block_impl( - block, - session_id=session_id, - disabled_tools=disabled_tools, - owner=owner, - progress_cb=progress_cb, - tool_policy=tool_policy, - approved_document_id=( - exact_approval.pending.document_id - if approval_claimed - else None - ), - approved_document_version=( - exact_approval.pending.document_version - if approval_claimed - else None - ), - approved_document_digest=( - exact_approval.pending.document_digest - if approval_claimed - else None - ), - active_document_id=active_document_id, - client_runtime_context=client_runtime_context, - ) + with bind_request_authority(authority): + output = await _execute_tool_block_impl( + block, + session_id=session_id, + disabled_tools=disabled_tools, + owner=owner, + progress_cb=progress_cb, + tool_policy=tool_policy, + approved_document_id=( + exact_approval.pending.document_id + if approval_claimed + else None + ), + approved_document_version=( + exact_approval.pending.document_version + if approval_claimed + else None + ), + approved_document_digest=( + exact_approval.pending.document_digest + if approval_claimed + else None + ), + active_document_id=active_document_id, + client_runtime_context=client_runtime_context, + ) if isinstance(security_context, ToolRunSecurityContext): security_context.observe_tool_result( getattr(block, "tool_type", None), @@ -1628,6 +1669,9 @@ async def _execute_tool_block_impl( from src import bg_jobs mark_dispatch() rec = bg_jobs.launch(_bg_cmd, session_id=session_id, cwd=agent_cwd()) + # Only this server launch may seal detached-job authority; a + # handler/bridge output carrying a job id is not a grant source. + save_background_authority(rec["id"], active_request_authority()) short = _bg_cmd.strip().split(chr(10))[0][:80] desc = f"bash (background): {short}" result = { diff --git a/src/tools/system.py b/src/tools/system.py index 76783badc..d60d8f585 100644 --- a/src/tools/system.py +++ b/src/tools/system.py @@ -486,12 +486,15 @@ async def do_manage_tasks(content: str, owner: Optional[str] = None) -> Dict: # Guard each fallback with `or`: args.get("prompt", default) returns # None when the key is present but null, and None[:50] raises. name = args.get("name") or (args.get("prompt") or args.get("action_name") or "Task")[:50] + from src.agent_runtime.authority import seal_task_authority task = ScheduledTask( id=task_id, owner=owner, name=name, prompt=args.get("prompt"), + request_authority_json=seal_task_authority( + args.get("prompt"), task_type, args.get("action_name"), owner=owner), task_type=task_type, action=args.get("action_name"), schedule=args.get("schedule", "daily") if trigger_type == "schedule" else None, @@ -557,6 +560,10 @@ async def do_manage_tasks(content: str, owner: Optional[str] = None) -> Dict: if args.get("action_name") is not None: task.action = args["action_name"] changed.append("action") + if any(args.get(field) is not None for field in ("prompt", "task_type", "action_name")): + from src.agent_runtime.authority import seal_task_authority + task.request_authority_json = seal_task_authority( + task.prompt, task.task_type, task.action, owner=owner) if args.get("trigger_type") is not None: task.trigger_type = args["trigger_type"] changed.append("trigger_type") diff --git a/tests/runtime_evidence_helpers.py b/tests/runtime_evidence_helpers.py index f5aa483f3..1ab245a15 100644 --- a/tests/runtime_evidence_helpers.py +++ b/tests/runtime_evidence_helpers.py @@ -2,6 +2,33 @@ from src.agent_runtime.journal import mark_dispatch, record_action +def server_authorized_executor(executor): + """Give standalone dispatcher fixtures their explicit server grants. + + These existing suites exercise handlers, policy, confinement and approvals. + Their fixture grants cover the declared native tool registry, independently + of the proposed block. Request-authority denial tests use the raw dispatcher. + """ + from functools import wraps + from inspect import signature + from src.agent_runtime.authority import OperationGrant, RequestAuthority + from src.tool_policy import known_tool_names + from src.turn_contract import canonical_tool + call_signature = signature(executor) + @wraps(executor) + async def execute(*args, **kwargs): + bound = call_signature.bind(*args, **kwargs) + parameters = bound.arguments + kwargs.setdefault("request_authority", RequestAuthority( + "standalone-test-request", str(parameters.get("owner") or "").strip().casefold(), + str(parameters.get("session_id") or ""), str(parameters.get("workspace") or ""), + tuple(OperationGrant(name) for name in sorted( + {canonical_tool(n) for n in known_tool_names()} | {"list_dir", "find_files"})), + )) + return await executor(*args, **kwargs) + return execute + + def authoritative_executor(function): @record_action async def execute(block, *args, **kwargs): diff --git a/tests/test_agent_external_tool_schemas.py b/tests/test_agent_external_tool_schemas.py index f6aa8c46e..f6688a870 100644 --- a/tests/test_agent_external_tool_schemas.py +++ b/tests/test_agent_external_tool_schemas.py @@ -333,6 +333,7 @@ def test_external_tool_images_are_threaded_as_multimodal_evidence(): def test_declared_external_call_reaches_scoped_bridge(monkeypatch): + from src.agent_runtime.authority import OperationGrant, RequestAuthority bridge_calls = [] round_no = 0 monkeypatch.setattr(agent_loop, "get_setting", lambda key, default=None: default) @@ -367,6 +368,8 @@ def test_declared_external_call_reaches_scoped_bridge(monkeypatch): [{"role": "user", "content": "Perform the declared operation."}], max_rounds=2, owner="pewds", + request_authority=RequestAuthority("declared-fixture", "pewds", "", "", + (OperationGrant("inspect_state"),)), relevant_tools={"inspect_state"}, forced_tools={"inspect_state"}, fallbacks=[], @@ -389,6 +392,7 @@ def test_declared_external_call_reaches_scoped_bridge(monkeypatch): def test_known_native_tool_reaches_scoped_bridge_without_redeclared_schema(monkeypatch): + from src.agent_runtime.authority import create_request_authority bridge_calls = [] round_no = 0 monkeypatch.setattr(agent_loop, "get_setting", lambda key, default=None: default) @@ -435,6 +439,7 @@ def test_known_native_tool_reaches_scoped_bridge_without_redeclared_schema(monke [{"role": "user", "content": "Search email for Project Alpha."}], max_rounds=2, owner="public-user", + request_authority=create_request_authority("Search email for Project Alpha.", owner="public-user"), relevant_tools={"search_emails"}, forced_tools={"search_emails"}, fallbacks=[], diff --git a/tests/test_ask_user_tool.py b/tests/test_ask_user_tool.py index 4094701bb..facaf2d01 100644 --- a/tests/test_ask_user_tool.py +++ b/tests/test_ask_user_tool.py @@ -11,6 +11,9 @@ from src.agent_tools import ToolBlock, TOOL_TAGS # noqa: E402 (import first to from src.tool_execution import NO_TOOL_SECURITY_CONTEXT, execute_tool_block from src.tool_index import ALWAYS_AVAILABLE, BUILTIN_TOOL_DESCRIPTIONS from src.tool_security import is_public_blocked_tool +from tests.runtime_evidence_helpers import server_authorized_executor + +execute_tool_block = server_authorized_executor(execute_tool_block) def _run(content): diff --git a/tests/test_client_tool_routing.py b/tests/test_client_tool_routing.py index 83fe271d1..fe050d7b7 100644 --- a/tests/test_client_tool_routing.py +++ b/tests/test_client_tool_routing.py @@ -14,6 +14,7 @@ from routes.chat_routes import _agent_turn_cwd import routes.chat_routes as chat_routes from src import tool_execution as _te from src.agent_loop import _is_explicit_local_network_request +from tests.runtime_evidence_helpers import server_authorized_executor # Hold module-object references (not just from-imported names): other test # modules re-import src.tool_execution via sys.modules pops, so string-target @@ -28,7 +29,7 @@ _ROUTED_BRIDGE_TOOLS = _te._ROUTED_BRIDGE_TOOLS async def _execute_tool_block_for_unit_tests(*args, **kwargs): """Use the explicit non-security-context test mode for dispatch tests.""" kwargs.setdefault("security_context", _te.NO_TOOL_SECURITY_CONTEXT) - return await _te.execute_tool_block(*args, **kwargs) + return await server_authorized_executor(_te.execute_tool_block)(*args, **kwargs) execute_tool_block = _execute_tool_block_for_unit_tests diff --git a/tests/test_edit_file.py b/tests/test_edit_file.py index b3fb202e1..6f94a3961 100644 --- a/tests/test_edit_file.py +++ b/tests/test_edit_file.py @@ -4,6 +4,14 @@ import os import tempfile import pytest +from tests.runtime_evidence_helpers import server_authorized_executor + + +@pytest.fixture(autouse=True) +def standalone_dispatch_authority(monkeypatch): + from src import tool_execution + monkeypatch.setattr(tool_execution, "execute_tool_block", + server_authorized_executor(tool_execution.execute_tool_block)) from src import tool_security from src.tool_security import ( diff --git a/tests/test_execution_bridge.py b/tests/test_execution_bridge.py index 9c72e0569..bcfc33b96 100644 --- a/tests/test_execution_bridge.py +++ b/tests/test_execution_bridge.py @@ -1,6 +1,7 @@ import asyncio import logging import src.tool_execution as tool_execution +from tests.runtime_evidence_helpers import server_authorized_executor from src.tool_execution import ( AgentExecutionBridge, @@ -11,6 +12,9 @@ from src.tool_execution import ( ) +execute_tool_block = server_authorized_executor(execute_tool_block) + + class Block: tool_type = "host_shell" diff --git a/tests/test_external_context_tool_gate.py b/tests/test_external_context_tool_gate.py index cc8d7f4f0..473cba897 100644 --- a/tests/test_external_context_tool_gate.py +++ b/tests/test_external_context_tool_gate.py @@ -7,6 +7,14 @@ from collections import namedtuple from pathlib import Path import pytest +from tests.runtime_evidence_helpers import server_authorized_executor + + +@pytest.fixture(autouse=True) +def standalone_dispatch_authority(monkeypatch): + from src import tool_execution + monkeypatch.setattr(tool_execution, "execute_tool_block", + server_authorized_executor(tool_execution.execute_tool_block)) from tests.helpers.document_source import document_source from tests.helpers.js_modules import email_library_paths diff --git a/tests/test_image_creation_routing.py b/tests/test_image_creation_routing.py index 9a2189775..a571e960b 100644 --- a/tests/test_image_creation_routing.py +++ b/tests/test_image_creation_routing.py @@ -125,6 +125,8 @@ async def test_generation_dispatch_uses_owner_aware_backend(monkeypatch, exit_co from types import SimpleNamespace from src import ai_interaction, tool_execution from src.agent_runtime.journal import ActionJournal, bind_journal + from tests.runtime_evidence_helpers import server_authorized_executor + execute = server_authorized_executor(tool_execution.execute_tool_block) calls = [] async def generate(content, **kwargs): calls.append((content, kwargs)) @@ -140,9 +142,9 @@ async def test_generation_dispatch_uses_owner_aware_backend(monkeypatch, exit_co block = SimpleNamespace(tool_type='generate_image', content='{"prompt":"A city"}') journal = ActionJournal() with bind_journal(journal): - _, denied = await tool_execution.execute_tool_block(block, owner='pewds', session_id='fixture', + _, denied = await execute(block, owner='pewds', session_id='fixture', disabled_tools={'generate_image'}, security_context=tool_execution.NO_TOOL_SECURITY_CONTEXT) - _, result = await tool_execution.execute_tool_block(block, owner='pewds', session_id='fixture', + _, result = await execute(block, owner='pewds', session_id='fixture', security_context=tool_execution.NO_TOOL_SECURITY_CONTEXT) assert denied['exit_code'] != 0 assert journal.actions[0].execution_id is None diff --git a/tests/test_request_authority.py b/tests/test_request_authority.py new file mode 100644 index 000000000..2f25449ca --- /dev/null +++ b/tests/test_request_authority.py @@ -0,0 +1,584 @@ +"""Request grants are independent of tool offerings and model proposals.""" +import asyncio +from contextlib import nullcontext +from dataclasses import replace +import json +from types import SimpleNamespace +from unittest.mock import AsyncMock, MagicMock + +import pytest + +from src.agent_runtime.authority import ( + MISSING_AUTHORITY, ExactOperation, OperationGrant, RequestAuthority, + active_request_authority, bind_request_authority, create_request_authority, + restore_background_authority, restore_task_authority, save_background_authority, + seal_task_authority, task_operation, with_request_authority, + is_internal_tool_request, require_user_approval_request, + request_authority_for_http, +) +from src.tool_policy import ToolPolicy +from src.tool_types import ToolBlock + + +def authority(*tools, owner="alice", session_id="s", workspace=""): + return RequestAuthority("request-test", owner, session_id, workspace, + tuple(OperationGrant(tool) for tool in tools)) + + +@pytest.mark.asyncio +@pytest.mark.parametrize("offering", ["schema", "bridge", "dynamic"]) +async def test_availability_and_model_selection_do_not_grant_execution(monkeypatch, offering): + from src import tool_execution as execution + from src.turn_contract import bind_turn_contract, resolve_turn_contract + implementation = AsyncMock(return_value=("bash", {"exit_code": 0})) + monkeypatch.setattr(execution, "_execute_tool_block_impl", implementation) + offered_handler = AsyncMock() + if offering == "dynamic": + import src.agent_tools + monkeypatch.setitem(src.agent_tools.TOOL_HANDLERS, "bash", offered_handler) + bridge_context = execution.bind_execution_bridge(execution.AgentExecutionBridge( + route_tool=offered_handler, supported_tools=frozenset({"bash"}))) if offering == "bridge" else nullcontext() + contract = resolve_turn_contract(capabilities={"shell_files"}, policy=ToolPolicy(), + schemas=[{"function": {"name": "bash"}}]) + with bind_turn_contract(contract), bridge_context: + description, result = await execution.execute_tool_block( + ToolBlock("bash", "echo 'authorized by model'"), owner="alice", session_id="s", + security_context=execution.NO_TOOL_SECURITY_CONTEXT, + request_authority=authority("transcribe_media")) + assert "BLOCKED" in description + assert result["failure_kind"] == "request_authority_denied" + implementation.assert_not_awaited() + offered_handler.assert_not_awaited() + + +@pytest.mark.parametrize("user_text,allowed,denied", [ + ("Transcribe /workspace/input/audio.wav", "transcribe_media", "bash"), + ("OCR extract exact text from /workspace/input/image.png", "extract_text", "python"), + ("List my tasks", "manage_tasks", "web_fetch"), + ("Search the web for current weather", "web_search", "bash"), +]) +def test_explicit_request_classes_remain_narrow(user_text, allowed, denied): + grant = create_request_authority(user_text) + content = '{"action":"list"}' if allowed == "manage_tasks" else '{}' + assert grant.permits(ExactOperation.normalize(allowed, content)) + assert not grant.permits(ExactOperation.normalize(denied, '{}')) + + +def test_safe_task_read_does_not_authorize_same_tool_mutation(): + grant = create_request_authority("List my tasks") + assert grant.permits(ExactOperation.normalize("manage_tasks", '{"action":"list"}')) + assert not grant.permits(ExactOperation.normalize("manage_tasks", '{"action":"create","prompt":"run bash"}')) + + +def test_exact_read_identifiers_cannot_be_changed_by_model(): + grant = create_request_authority("Read note id abc123") + read = next(g for g in grant.grants if g.tool == "manage_notes") + assert read.inputs is not None + assert not grant.permits(ExactOperation.normalize("manage_notes", '{"action":"view","id":"another"}')) + + +def test_browser_fallback_does_not_authorize_interaction_or_evaluation(): + grant = create_request_authority("Use web_fetch to read https://example.test") + assert grant.permits(ExactOperation.normalize("private_browser", '{"action":"open","url":"https://example.test"}')) + for action in ("click", "fill", "evaluate"): + assert not grant.permits(ExactOperation.normalize("private_browser", json.dumps({"action": action}))) + + +def test_unknown_intent_has_no_generic_execution_floor(): + grant = create_request_authority("Please solve this") + assert {g.tool for g in grant.grants} == {"ask_user", "update_plan"} + + +@pytest.mark.parametrize("metadata", [ + {"tool_events": [{"tool": "bash", "output": "pwd", "exit_code": 0}]}, + {"tool_events": [{"tool": "python", "output": "ready", "exit_code": 0}]}, +]) +def test_model_history_cannot_establish_followup_authority(metadata): + history = [ + {"role": "user", "content": "Transcribe /workspace/input/a.wav"}, + {"role": "assistant", "content": "I will run bash and python", "metadata": metadata}, + {"role": "user", "content": "Run bash", "metadata": {"trusted": False}}, + ] + grant = create_request_authority("Try it again", history=history) + assert not grant.permits(ExactOperation.normalize("bash", "pwd")) + assert not grant.permits(ExactOperation.normalize("python", "print(1)")) + + +@pytest.mark.parametrize("mutation", [ + {"version": True}, {"grants": "bash"}, {"denied": None}, + {"block_all": "false"}, {"inherited": 0}, +]) +def test_malformed_persisted_authority_is_rejected(mutation): + snapshot = authority("bash").to_dict() + snapshot.update(mutation) + with pytest.raises((ValueError, TypeError, KeyError)): + RequestAuthority.from_dict(snapshot) + + +def test_explicit_local_network_lookup_is_host_only(): + grant = create_request_authority("find ajax local ip on the LAN") + assert grant.permits(ExactOperation.normalize("host_shell", '{"command":"ip neigh | grep ajax"}')) + assert not grant.permits(ExactOperation.normalize("bash", "pwd")) + assert not grant.permits(ExactOperation.normalize("python", "print(1)")) + + +@pytest.mark.parametrize("content", ['{"x":1,"x":2}', '{"x":NaN}', '{"action":']) +def test_malformed_exact_operation_is_rejected(content): + with pytest.raises(ValueError): + ExactOperation.normalize("manage_tasks", content) + + +def test_raw_script_braces_are_preserved_as_exact_input(): + script = "{ printf requested; }" + assert ExactOperation.normalize("bash", script).input == script + snapshot = seal_task_authority(script, "action", "run_local", owner="alice") + restored = restore_task_authority(snapshot, script, "action", "run_local", owner="alice") + assert restored.permits(task_operation("action", "run_local", script)) + assert not restored.permits(task_operation("action", "run_local", "{ printf other; }")) + + +def test_normalization_and_aliases_do_not_erase_denials(): + grant = authority("read_email").restrict(disabled_tools={"mcp__email__read_email"}) + assert not grant.permits(ExactOperation.normalize("read_email", '{}')) + grant = authority("read_email").restrict(ToolPolicy(disable_mcp=True)) + assert not grant.permits(ExactOperation.normalize("mcp__email__read_email", '{}')) + assert ExactOperation.normalize("manage_tasks", '{ "action": "list" }').input == '{"action":"list"}' + + +@pytest.mark.asyncio +@pytest.mark.parametrize("state", [MISSING_AUTHORITY, None, {}, "authorized"]) +async def test_missing_and_malformed_dispatch_authority_fail_closed(monkeypatch, state): + from src import tool_execution as execution + implementation = AsyncMock() + monkeypatch.setattr(execution, "_execute_tool_block_impl", implementation) + _, result = await execution.execute_tool_block(ToolBlock("bash", "pwd"), + security_context=execution.NO_TOOL_SECURITY_CONTEXT, request_authority=state) + assert result["blocked"] is True + implementation.assert_not_awaited() + + +@pytest.mark.asyncio +async def test_dispatch_checks_grants_and_current_disabled_policy(monkeypatch): + from src import tool_execution as execution + implementation = AsyncMock(return_value=("bash", {"exit_code": 0})) + monkeypatch.setattr(execution, "_execute_tool_block_impl", implementation) + for disabled in (set(), {"bash"}): + _, result = await execution.execute_tool_block(ToolBlock("bash", "pwd"), + owner="alice", session_id="s", disabled_tools=disabled, + security_context=execution.NO_TOOL_SECURITY_CONTEXT, + request_authority=authority("bash")) + assert result["exit_code"] == (1 if disabled else 0) + assert implementation.await_count == 1 + + +@pytest.mark.asyncio +async def test_nested_stream_intersects_and_restores_parent_on_close(): + seen = [] + @with_request_authority + async def child(messages, request_authority=MISSING_AUTHORITY, owner="alice", session_id="s"): + seen.append(active_request_authority()) + yield "child" + parent = authority("transcribe_media") + with bind_request_authority(parent): + stream = child([{"role": "user", "content": "Run bash"}], + request_authority=authority("bash", "transcribe_media")) + assert await anext(stream) == "child" + assert not seen[0].permits(ExactOperation.normalize("bash", "pwd")) + assert seen[0].permits(ExactOperation.normalize("transcribe_media", '{}')) + await stream.aclose() + assert active_request_authority() is parent + assert active_request_authority() is None + + +@pytest.mark.asyncio +async def test_retries_and_provider_changes_do_not_recreate_authority(): + seen = [] + @with_request_authority + async def run(messages, request_authority=MISSING_AUTHORITY, owner="alice", session_id="s"): + for proposed in ("transcribe_media", "bash", "python"): + seen.append(active_request_authority()) + yield active_request_authority().permits(ExactOperation.normalize(proposed, '{}')) + assert [x async for x in run([{"role": "user", "content": "Transcribe /workspace/input/a.wav"}])] == [True, False, False] + assert all(value is seen[0] for value in seen) + + +def test_task_snapshot_caps_model_payload_and_rejects_changed_or_missing_state(): + parent = authority("manage_tasks") + with bind_request_authority(parent): + snapshot = seal_task_authority("Run bash in the workspace", "llm", None, owner="alice") + restored = restore_task_authority(snapshot, "Run bash in the workspace", "llm", None, owner="alice") + assert not restored.permits(ExactOperation.normalize("bash", "pwd")) + for state in (None, "{}", snapshot): + changed = restore_task_authority(state, "a different prompt", "llm", None, owner="alice") + assert changed.grants == () + + +def test_direct_task_ingress_seals_only_exact_builtin_action(): + snapshot = seal_task_authority("printf requested", "action", "run_local", owner="alice") + restored = restore_task_authority(snapshot, "printf requested", "action", "run_local", owner="alice") + assert restored.permits(task_operation("action", "run_local", "printf requested")) + assert not restored.permits(ExactOperation.normalize("bash", "printf other")) + + +def test_background_snapshot_preserves_scope_and_rejects_other_session(monkeypatch, tmp_path): + import src.constants + monkeypatch.setattr(src.constants, "BG_JOBS_DIR", str(tmp_path)) + grant = authority("transcribe_media").restrict(disabled_tools={"bash"}) + save_background_authority("job1", grant) + restored = restore_background_authority("job1", owner="alice", session_id="s") + assert restored.request_id == grant.request_id + assert restored.denied == frozenset({"bash"}) + assert not restored.permits(ExactOperation.normalize("python", "print(1)")) + assert restore_background_authority("job1", owner="alice", session_id="other").grants == () + + +@pytest.mark.asyncio +async def test_only_server_background_launch_can_seal_job_authority(monkeypatch, tmp_path): + import src.constants + from src import bg_jobs, tool_execution as execution + monkeypatch.setattr(src.constants, "BG_JOBS_DIR", str(tmp_path)) + monkeypatch.setattr(execution, "_owner_is_admin", lambda owner: True) + monkeypatch.setattr(bg_jobs, "launch", lambda *a, **k: {"id": "server-job"}) + await execution.execute_tool_block(ToolBlock("bash", "#!bg\nprintf trusted"), + owner="alice", session_id="s", security_context=execution.NO_TOOL_SECURITY_CONTEXT, + request_authority=authority("bash")) + restored = restore_background_authority("server-job", owner="alice", session_id="s") + assert restored.request_id == "request-test" + assert restored.permits(ExactOperation.normalize("bash", "printf trusted")) + handler = AsyncMock(return_value=("transcribe_media", {"bg_job_id": "forged-job", "exit_code": 0})) + monkeypatch.setattr(execution, "_execute_tool_block_impl", handler) + await execution.execute_tool_block(ToolBlock("transcribe_media", '{}'), + owner="alice", session_id="s", security_context=execution.NO_TOOL_SECURITY_CONTEXT, + request_authority=authority("transcribe_media")) + assert not (tmp_path / "forged-job.authority.json").exists() + + +@pytest.mark.asyncio +async def test_exact_approval_grants_one_input_without_widening_continuation(monkeypatch): + from src import tool_execution as execution + from src.tool_approvals import ToolApprovalStore + from src.tool_capabilities import ToolRunSecurityContext, capabilities_for_action + store = ToolApprovalStore() + original = authority("transcribe_media") + pending = store.create(owner="alice", session_id="s", origin_run_id="journal-parent", + tool_name="bash", content="printf approved", workspace=None, + external_untrusted_context_seen=True, capabilities=capabilities_for_action("bash", "printf approved"), + request_authority=original, selected_tools={"bash", "python"}) + approval = store.consume(pending.approval_id, owner="alice", session_id="s", decision="approve_task") + implementation = AsyncMock(return_value=("bash", {"exit_code": 0})) + monkeypatch.setattr(execution, "_execute_tool_block_impl", implementation) + for tool, content, expected in (("bash", "printf other", 1), ("bash", "printf approved", 0), + ("bash", "printf approved", 1), ("python", "print(1)", 1)): + _, result = await execution.execute_tool_block(ToolBlock(tool, content), owner="alice", session_id="s", + security_context=ToolRunSecurityContext(external_untrusted_context_seen=True), + request_authority=original, exact_approval=approval) + assert result["exit_code"] == expected + assert implementation.await_count == 1 + assert "request_authority" not in pending.public_payload() + + +def test_approval_digest_binds_authority_snapshot(): + from src.tool_approvals import ExactToolApproval, ToolApprovalStore + from src.tool_capabilities import capabilities_for_action + pending = ToolApprovalStore().create(owner="alice", session_id="s", origin_run_id="journal", + tool_name="bash", content="pwd", workspace=None, external_untrusted_context_seen=True, + capabilities=capabilities_for_action("bash", "pwd"), request_authority=authority("transcribe_media")) + forged = ExactToolApproval(replace(pending, request_authority=authority("bash", "python"))) + assert not forged.matches(owner="alice", session_id="s", workspace=None, tool_name="bash", content="pwd") + + +@pytest.mark.asyncio +async def test_nested_approval_cannot_cross_parent_class_ceiling(monkeypatch): + from src import tool_execution as execution + from src.tool_approvals import ToolApprovalStore + from src.tool_capabilities import ToolRunSecurityContext, capabilities_for_action + store = ToolApprovalStore() + pending = store.create(owner="alice", session_id="s", origin_run_id="parent", + tool_name="bash", content="pwd", workspace=None, external_untrusted_context_seen=True, + capabilities=capabilities_for_action("bash", "pwd")) + approval = store.consume(pending.approval_id, owner="alice", session_id="s", decision="approve") + implementation = AsyncMock() + monkeypatch.setattr(execution, "_execute_tool_block_impl", implementation) + with bind_request_authority(authority("transcribe_media")): + _, result = await execution.execute_tool_block(ToolBlock("bash", "pwd"), owner="alice", session_id="s", + security_context=ToolRunSecurityContext(external_untrusted_context_seen=True), + request_authority=authority("bash"), exact_approval=approval) + assert result["blocked"] is True + implementation.assert_not_awaited() + + +@pytest.mark.asyncio +async def test_teacher_receives_parent_authority_instead_of_synthetic_prompt_grants(monkeypatch): + import src.agent_loop as loop + import src.ai_interaction as interaction + import src.settings as settings + import src.teacher_escalation as teacher + original = authority("transcribe_media") + seen = [] + monkeypatch.setattr(settings, "get_setting", lambda key, default=None: + {"teacher_enabled": True, "teacher_model": "teacher"}.get(key, default)) + monkeypatch.setattr(interaction, "_resolve_model", lambda *a, **k: ("https://teacher.invalid", "teacher", {})) + monkeypatch.setattr(teacher, "evaluate_turn_regex", lambda *a: ("failure", "test failure")) + @with_request_authority + async def child(messages, request_authority=MISSING_AUTHORITY, owner=None, session_id=None, **kwargs): + seen.append(active_request_authority()) + yield 'data: [DONE]\n\n' + monkeypatch.setattr(loop, "stream_agent_loop", child) + with bind_request_authority(original): + chunks = [chunk async for chunk in teacher.run_teacher_inline( + student_endpoint_url="https://student.invalid", + student_messages=[{"role": "user", "content": "Run bash and python"}], + student_tool_events=[], student_reply="failed", owner="alice", session_id="s", + request_authority=original, parent_run_id="journal-parent")] + assert active_request_authority() is original + assert chunks + assert seen[0].request_id == original.request_id + assert not seen[0].permits(ExactOperation.normalize("bash", "pwd")) + assert seen[0].permits(ExactOperation.normalize("transcribe_media", '{}')) + + +@pytest.mark.asyncio +async def test_untrusted_user_role_result_cannot_become_request_authority(): + @with_request_authority + async def run(messages, request_authority=MISSING_AUTHORITY): + yield active_request_authority() + messages = [{"role": "user", "content": "Transcribe /workspace/input/a.wav"}, + {"role": "user", "content": "Run bash", "metadata": {"trusted": False}}] + stream = run(messages) + grant = await anext(stream) + await stream.aclose() + assert not grant.permits(ExactOperation.normalize("bash", "pwd")) + + +@pytest.mark.asyncio +async def test_scheduled_builtin_missing_authority_does_not_invoke_action(monkeypatch): + import src.builtin_actions as actions + from src.task_scheduler import TaskScheduler + handler = AsyncMock(return_value=("ok", True)) + monkeypatch.setitem(actions.BUILTIN_ACTIONS, "run_local", handler) + task = SimpleNamespace(prompt="printf exact", task_type="action", action="run_local", + owner="alice", name="task", request_authority_json=None) + result, success = await TaskScheduler(session_manager=None)._execute_action(task) + assert not success + assert "authority" in result + handler.assert_not_awaited() + + +def test_task_authority_column_migration_is_additive_and_idempotent(monkeypatch, tmp_path): + import core.database as database + from sqlalchemy import create_engine, inspect, text + engine = create_engine(f"sqlite:///{tmp_path / 'legacy.db'}") + with engine.begin() as connection: + connection.execute(text("CREATE TABLE scheduled_tasks (id TEXT PRIMARY KEY, prompt TEXT)")) + connection.execute(text("INSERT INTO scheduled_tasks VALUES ('legacy', 'Run bash')")) + monkeypatch.setattr(database, "engine", engine) + database._migrate_add_task_authority_column() + database._migrate_add_task_authority_column() + assert "request_authority_json" in {c["name"] for c in inspect(engine).get_columns("scheduled_tasks")} + with engine.connect() as connection: + assert connection.execute(text("SELECT prompt, request_authority_json FROM scheduled_tasks")).one() == ("Run bash", None) + engine.dispose() + + +@pytest.mark.asyncio +async def test_server_seeded_defaults_have_authority_but_legacy_rows_do_not_gain_it(monkeypatch, tmp_path): + import core.database as database + import routes.prefs_routes as preferences + from sqlalchemy import create_engine + from sqlalchemy.orm import sessionmaker + from src.task_scheduler import TaskScheduler + engine = create_engine(f"sqlite:///{tmp_path / 'tasks.db'}") + database.Base.metadata.create_all(engine) + sessions = sessionmaker(bind=engine) + monkeypatch.setattr(database, "SessionLocal", sessions) + monkeypatch.setattr(preferences, "_load_for_user", lambda owner: {}) + scheduler = TaskScheduler(session_manager=None) + monkeypatch.setattr(scheduler, "ensure_assistant_defaults", AsyncMock()) + with sessions() as db: + db.add(database.ScheduledTask(id="legacy", owner="alice", name="Legacy housekeeping", + task_type="action", action="tidy_sessions")) + db.commit() + await scheduler.ensure_defaults("alice") + with sessions() as db: + tasks = db.query(database.ScheduledTask).all() + assert len(tasks) > 1 + for task in tasks: + grant = restore_task_authority(task.request_authority_json, task.prompt, + task.task_type, task.action, owner=task.owner) + assert grant.permits(task_operation(task.task_type, task.action, task.prompt)) == (task.id != "legacy") + engine.dispose() + + +@pytest.mark.asyncio +async def test_dispatch_binds_explicit_authority_for_nested_handler(monkeypatch): + from src import tool_execution as execution + seen = [] + async def handler(*args, **kwargs): + seen.append(active_request_authority()) + with bind_request_authority(authority("bash", "transcribe_media")) as child: + assert not child.permits(ExactOperation.normalize("bash", "pwd")) + return "transcribe_media", {"exit_code": 0} + monkeypatch.setattr(execution, "_execute_tool_block_impl", handler) + await execution.execute_tool_block(ToolBlock("transcribe_media", '{}'), owner="alice", session_id="s", + security_context=execution.NO_TOOL_SECURITY_CONTEXT, request_authority=authority("transcribe_media")) + assert seen + assert active_request_authority() is None + + +def test_tool_http_task_payload_is_not_new_user_authority(): + from core.middleware import INTERNAL_TOOL_HEADER, INTERNAL_TOOL_TOKEN + from routes.task.task_routes import _seal_request_task_authority + request = SimpleNamespace(headers={INTERNAL_TOOL_HEADER: INTERNAL_TOOL_TOKEN}) + snapshot = _seal_request_task_authority(request, "Run bash in workspace", "llm", None, "alice") + restored = restore_task_authority(snapshot, "Run bash in workspace", "llm", None, owner="alice") + assert not restored.permits(ExactOperation.normalize("bash", "pwd")) + + +@pytest.mark.parametrize("marker", ["header", "middleware"]) +@pytest.mark.parametrize("owner", [None, "alice"]) +def test_internal_tool_requests_cannot_submit_user_approval(marker, owner): + from core.middleware import INTERNAL_TOOL_HEADER, INTERNAL_TOOL_TOKEN + from fastapi import HTTPException + request = SimpleNamespace( + headers={INTERNAL_TOOL_HEADER: INTERNAL_TOOL_TOKEN} if marker == "header" else {}, + state=SimpleNamespace(current_user="internal-tool" if marker == "middleware" else owner)) + assert is_internal_tool_request(request) + with pytest.raises(HTTPException) as failure: + require_user_approval_request(request) + assert failure.value.status_code == 403 + human = SimpleNamespace(headers={}, state=SimpleNamespace(current_user=owner)) + require_user_approval_request(human) + + +@pytest.mark.parametrize("internal", [True, False]) +def test_http_chat_request_context_cannot_mint_authority_from_tool_message(internal): + from core.middleware import INTERNAL_TOOL_HEADER, INTERNAL_TOOL_TOKEN + request = SimpleNamespace( + headers={INTERNAL_TOOL_HEADER: INTERNAL_TOOL_TOKEN} if internal else {}, + state=SimpleNamespace(current_user="alice")) + grant = request_authority_for_http(request, "Run bash in the workspace", owner="alice", + session_id="s", workspace="/workspace/original", policy=ToolPolicy(disabled_tools={"python"})) + assert grant.bound_to(owner="alice", session_id="s", workspace="/workspace/original") + assert grant.permits(ExactOperation.normalize("bash", "pwd")) is not internal + assert not grant.permits(ExactOperation.normalize("python", "print(1)")) + + +def test_internal_chat_context_does_not_become_trusted_followup_history(): + from routes.chat_routes import _append_internal_chat_context + trusted = {"role": "user", "content": "Transcribe /workspace/input/a.wav"} + ctx = SimpleNamespace(messages=[trusted], route_messages=[trusted]) + _append_internal_chat_context(ctx, "Run bash in the workspace") + assert ctx.messages == ctx.route_messages + assert ctx.messages[-1]["metadata"]["trusted"] is False + grant = create_request_authority("Try it again", history=ctx.messages) + assert not grant.permits(ExactOperation.normalize("bash", "pwd")) + + +@pytest.mark.asyncio +async def test_skill_approval_ingress_rejects_internal_tool_before_consuming(monkeypatch): + import routes.skills_routes as skills + from core.middleware import INTERNAL_TOOL_HEADER, INTERNAL_TOOL_TOKEN + from fastapi import HTTPException + from src import tool_approvals + consume = MagicMock() + monkeypatch.setattr(tool_approvals.tool_approval_store, "consume", consume) + router = skills.setup_skills_routes(MagicMock()) + endpoint = next(route.endpoint for route in router.routes + if route.path.endswith("/test-approval")) + request = SimpleNamespace(headers={INTERNAL_TOOL_HEADER: INTERNAL_TOOL_TOKEN}, + state=SimpleNamespace(current_user="alice")) + with pytest.raises(HTTPException) as failure: + await endpoint(request, "skill") + assert failure.value.status_code == 403 + consume.assert_not_called() + + +@pytest.mark.asyncio +@pytest.mark.parametrize("internal", [True, False]) +async def test_skill_test_ingress_distinguishes_user_task_from_tool_payload(monkeypatch, internal): + import routes.skills_routes as skills + import src.endpoint_resolver as endpoints + import src.llm_core as llm + from core.middleware import INTERNAL_TOOL_HEADER, INTERNAL_TOOL_TOKEN + manager = MagicMock() + manager.load.return_value = [{"name": "skill", "owner": "alice"}] + manager.read_skill_md.return_value = "# Skill" + monkeypatch.setattr(skills, "get_current_user", lambda request: "alice") + monkeypatch.setattr(skills, "_skill_test_jobs", {}) + monkeypatch.setattr(endpoints, "resolve_endpoint", lambda *a, **k: ("https://inference.invalid", "model", {})) + monkeypatch.setattr(llm, "list_model_ids", lambda *a, **k: []) + seen = [] + async def run(*args, **kwargs): + seen.append(kwargs["request_authority"]) + monkeypatch.setattr(skills, "_run_skill_test_job", run) + router = skills.setup_skills_routes(manager) + endpoint = next(route.endpoint for route in router.routes if route.path.endswith("/{skill_id}/test")) + request = SimpleNamespace( + headers={INTERNAL_TOOL_HEADER: INTERNAL_TOOL_TOKEN} if internal else {}, + state=SimpleNamespace(current_user="alice"), + json=AsyncMock(return_value={"task": "Run bash in the workspace"})) + result = await endpoint(request, "skill") + await asyncio.sleep(0) + assert result["status"] == "running" + assert len(seen) == 1 + assert seen[0].permits(ExactOperation.normalize("bash", "pwd")) is not internal + + +@pytest.mark.asyncio +async def test_scheduled_override_cannot_replace_stored_authority(monkeypatch): + from src import agent_loop + from src.task_scheduler import TaskScheduler + seen = [] + async def fake_stream(*args, **kwargs): + seen.append(kwargs) + yield 'data: {"delta":"done"}\n\n' + yield 'data: [DONE]\n\n' + monkeypatch.setattr(agent_loop, "stream_agent_loop", fake_stream) + task = SimpleNamespace(prompt="Transcribe /workspace/input/a.wav", task_type="llm", action=None, + owner="alice", name="test", max_steps=1) + with bind_request_authority(authority("transcribe_media", workspace="/workspace/original")): + task.request_authority_json = seal_task_authority(task.prompt, "llm", None, owner="alice") + await TaskScheduler(session_manager=None)._run_agent_loop("https://inference.invalid", "test", task, "s", + override_user_message="Run bash and python") + grant = seen[0]["request_authority"] + assert grant.permits(ExactOperation.normalize("transcribe_media", '{}')) + assert not grant.permits(ExactOperation.normalize("bash", "pwd")) + assert grant.workspace == "/workspace/original" + assert seen[0]["workspace"] == grant.workspace + + +@pytest.mark.asyncio +async def test_loop_retains_denial_before_offered_inventory_reconciliation(monkeypatch): + from src import agent_loop as loop, tool_execution as execution + from src.turn_contract import resolve_turn_contract + implementation = AsyncMock(return_value=("bash", {"exit_code": 0, "output": "unexpected"})) + monkeypatch.setattr(execution, "_execute_tool_block_impl", implementation) + monkeypatch.setattr(loop, "get_setting", lambda key, default=None: default) + monkeypatch.setattr(loop, "get_mcp_manager", lambda: None) + monkeypatch.setattr(loop, "estimate_tokens", lambda *a, **k: 10) + async def fake_stream(*args, **kwargs): + yield 'data: ' + json.dumps({"delta": '```bash\npwd\n```'}) + '\n\n' + monkeypatch.setattr(loop, "stream_llm_with_fallback", fake_stream) + contract = resolve_turn_contract(capabilities={"shell_files"}, selected_tools={"bash"}, + schemas=[{"function": {"name": "bash"}}], policy=ToolPolicy()) + chunks = [chunk async for chunk in loop.stream_agent_loop( + "https://inference.invalid", "test", [{"role": "user", "content": "Run bash pwd"}], + owner="alice", session_id="s", max_rounds=1, turn_contract=contract, disabled_tools={"bash"})] + implementation.assert_not_awaited() + assert chunks[-1] == 'data: [DONE]\n\n' + assert active_request_authority() is None + + +@pytest.mark.asyncio +async def test_authority_denial_cannot_create_completion_receipt(monkeypatch): + from src import tool_execution as execution + from src.agent_runtime.journal import ActionJournal, bind_journal + journal = ActionJournal() + with bind_journal(journal): + await execution.execute_tool_block(ToolBlock("bash", "pwd"), owner="alice", session_id="s", + security_context=execution.NO_TOOL_SECURITY_CONTEXT, + request_authority=authority("transcribe_media")) + assert len(journal.actions) == 1 + assert journal.actions[0].execution_id is None + assert journal.actions[0].outcome["authoritative"] is False + assert journal.actions[0].outcome["blocked"] is True diff --git a/tests/test_review_regressions.py b/tests/test_review_regressions.py index a2df0ba65..a05c74e10 100644 --- a/tests/test_review_regressions.py +++ b/tests/test_review_regressions.py @@ -14,9 +14,10 @@ from src.preset_manager import PresetManager async def _execute_without_run_context(execute_tool_block, *args, **kwargs): from src.tool_execution import NO_TOOL_SECURITY_CONTEXT + from tests.runtime_evidence_helpers import server_authorized_executor kwargs.setdefault("security_context", NO_TOOL_SECURITY_CONTEXT) - return await execute_tool_block(*args, **kwargs) + return await server_authorized_executor(execute_tool_block)(*args, **kwargs) class _FakeColumn: diff --git a/tests/test_runtime_evidence_contract.py b/tests/test_runtime_evidence_contract.py index 45774e40a..c7760e723 100644 --- a/tests/test_runtime_evidence_contract.py +++ b/tests/test_runtime_evidence_contract.py @@ -6,6 +6,14 @@ import json import os import pytest +from tests.runtime_evidence_helpers import server_authorized_executor + + +@pytest.fixture(autouse=True) +def standalone_dispatch_authority(monkeypatch): + from src import tool_execution + monkeypatch.setattr(tool_execution, "execute_tool_block", + server_authorized_executor(tool_execution.execute_tool_block)) from src.agent_evidence import CompletionRequirements, EvidenceLedger, EvidenceKind from src.agent_runtime.completion import completion_answer, with_completion_gate diff --git a/tests/test_task_cookbook_admin_gate.py b/tests/test_task_cookbook_admin_gate.py index d7e72f9ef..09a0b48c0 100644 --- a/tests/test_task_cookbook_admin_gate.py +++ b/tests/test_task_cookbook_admin_gate.py @@ -88,7 +88,7 @@ def builtin_action_info(monkeypatch): def _req(user): - return SimpleNamespace(state=SimpleNamespace(current_user=user)) + return SimpleNamespace(state=SimpleNamespace(current_user=user), headers={}) def _endpoint(method, path): diff --git a/tests/test_task_scheduler_cancel.py b/tests/test_task_scheduler_cancel.py index d0f080533..73b81d3cf 100644 --- a/tests/test_task_scheduler_cancel.py +++ b/tests/test_task_scheduler_cancel.py @@ -25,6 +25,7 @@ def _setup_db(tmp_path, monkeypatch): next_run = Column(DateTime) last_run = Column(DateTime) prompt = Column(Text, default='') + request_authority_json = Column(Text) trigger_type = Column(String, default='schedule') schedule = Column(String, default='daily') scheduled_time = Column(String, default='08:00') @@ -187,9 +188,11 @@ def test_running_task_cancel_keeps_event_loop_responsive_during_database_lock(tm from src.task_scheduler import TaskScheduler from src.builtin_actions import BUILTIN_ACTIONS monkeypatch.setenv('BACKGROUND_TASK_FOREGROUND_GATE', 'false') + from src.agent_runtime.authority import seal_task_authority with session_local() as db: db.add(ScheduledTask(id='running-task', owner='alice', name='Fixture action', - task_type='action', action='fixture_wait', status='active')) + task_type='action', action='fixture_wait', status='active', + request_authority_json=seal_task_authority('', 'action', 'fixture_wait', owner='alice'))) db.commit() started = asyncio.Event() async def external_action(**kwargs): diff --git a/tests/test_tool_approvals.py b/tests/test_tool_approvals.py index b75127b7d..aefa7931b 100644 --- a/tests/test_tool_approvals.py +++ b/tests/test_tool_approvals.py @@ -4,6 +4,14 @@ import time from collections import namedtuple import pytest +from tests.runtime_evidence_helpers import server_authorized_executor + + +@pytest.fixture(autouse=True) +def standalone_dispatch_authority(monkeypatch): + from src import tool_execution + monkeypatch.setattr(tool_execution, "execute_tool_block", + server_authorized_executor(tool_execution.execute_tool_block)) from src.tool_approvals import ToolApprovalStore, document_content_digest from src.tool_capabilities import ToolRunSecurityContext, capabilities_for_action diff --git a/tests/test_tool_path_confinement.py b/tests/test_tool_path_confinement.py index d8e1400fc..8c3e60414 100644 --- a/tests/test_tool_path_confinement.py +++ b/tests/test_tool_path_confinement.py @@ -18,6 +18,14 @@ from types import SimpleNamespace from unittest.mock import patch import pytest +from tests.runtime_evidence_helpers import server_authorized_executor + + +@pytest.fixture(autouse=True) +def standalone_dispatch_authority(monkeypatch): + from src import tool_execution + monkeypatch.setattr(tool_execution, "execute_tool_block", + server_authorized_executor(tool_execution.execute_tool_block)) def _make_block(tool_type, content): diff --git a/tests/test_tool_policy.py b/tests/test_tool_policy.py index 166ede246..40d875cb2 100644 --- a/tests/test_tool_policy.py +++ b/tests/test_tool_policy.py @@ -17,6 +17,9 @@ from src.tool_policy import ( web_search_enabled_for_turn, ) from src.turn_contract import requested_capabilities +from tests.runtime_evidence_helpers import server_authorized_executor + +execute_tool_block = server_authorized_executor(execute_tool_block) def _collect(gen): diff --git a/tests/test_turn_contract.py b/tests/test_turn_contract.py index 74c03ff67..1ff2d14a5 100644 --- a/tests/test_turn_contract.py +++ b/tests/test_turn_contract.py @@ -1806,6 +1806,8 @@ def test_candidate_schema_changes_cannot_mutate_contract(): async def test_dispatcher_enforces_bound_contract_without_external_mutations(monkeypatch): from src import tool_execution, tool_implementations from src.tool_execution import NO_TOOL_SECURITY_CONTEXT, execute_tool_block + from tests.runtime_evidence_helpers import server_authorized_executor + execute_tool_block = server_authorized_executor(execute_tool_block) handler = AsyncMock(return_value={"events": [], "exit_code": 0}) monkeypatch.setattr(tool_implementations, "do_manage_calendar", handler) diff --git a/tests/test_turn_contract_integration.py b/tests/test_turn_contract_integration.py index caa9bae42..642ad96c1 100644 --- a/tests/test_turn_contract_integration.py +++ b/tests/test_turn_contract_integration.py @@ -4,6 +4,14 @@ import json from unittest.mock import AsyncMock, Mock import pytest +from tests.runtime_evidence_helpers import server_authorized_executor + + +@pytest.fixture(autouse=True) +def standalone_dispatch_authority(monkeypatch): + from src import tool_execution + monkeypatch.setattr(tool_execution, "execute_tool_block", + server_authorized_executor(tool_execution.execute_tool_block)) from src.tool_policy import ToolPolicy from src.tool_schemas import FUNCTION_TOOL_SCHEMAS diff --git a/tests/test_update_plan_tool.py b/tests/test_update_plan_tool.py index 842be246f..9973d6bbd 100644 --- a/tests/test_update_plan_tool.py +++ b/tests/test_update_plan_tool.py @@ -11,6 +11,9 @@ from src.agent_tools import ToolBlock, TOOL_TAGS # import first to avoid circul from src.tool_execution import NO_TOOL_SECURITY_CONTEXT, execute_tool_block from src.tool_index import ALWAYS_AVAILABLE, BUILTIN_TOOL_DESCRIPTIONS from src.tool_security import is_public_blocked_tool +from tests.runtime_evidence_helpers import server_authorized_executor + +execute_tool_block = server_authorized_executor(execute_tool_block) def _run(content): diff --git a/tests/test_weather_search_recovery.py b/tests/test_weather_search_recovery.py index 043120d20..e974b7787 100644 --- a/tests/test_weather_search_recovery.py +++ b/tests/test_weather_search_recovery.py @@ -10,6 +10,9 @@ from src.tool_policy import WEB_TOOL_NAMES, ToolPolicy from src.tool_schemas import FUNCTION_TOOL_SCHEMAS from src.tool_types import ToolBlock from src.tool_execution import NO_TOOL_SECURITY_CONTEXT, execute_tool_block +from tests.runtime_evidence_helpers import server_authorized_executor + +execute_tool_block = server_authorized_executor(execute_tool_block) from src.turn_contract import FAMILY_TOOLS diff --git a/tests/test_workspace_confine.py b/tests/test_workspace_confine.py index 25ca7c192..d033b6481 100644 --- a/tests/test_workspace_confine.py +++ b/tests/test_workspace_confine.py @@ -17,6 +17,7 @@ import tempfile from types import SimpleNamespace import pytest +from tests.runtime_evidence_helpers import server_authorized_executor from src.tool_execution import ( NO_TOOL_SECURITY_CONTEXT, @@ -31,6 +32,9 @@ from src.tool_execution import ( ) +_execute_tool_block = server_authorized_executor(_execute_tool_block) + + async def execute_tool_block(*args, **kwargs): kwargs.setdefault("security_context", NO_TOOL_SECURITY_CONTEXT) return await _execute_tool_block(*args, **kwargs) diff --git a/website/configuration-reference.md b/website/configuration-reference.md index fc2a3bdeb..69b3bea60 100644 --- a/website/configuration-reference.md +++ b/website/configuration-reference.md @@ -222,7 +222,7 @@ Listed for completeness. Setting one of these on a real install is either a no-o | `ODYSSEUS_QA_TEACHER_TIMEOUT` | `'120'` | `scripts/odysseus_conversation_qa.py:372` | Timeout in seconds for that call. Clamped to 15-120. | | `ODYSSEUS_RUNTIME_REVISION` | `''` | `routes/chat_helpers.py:198` (+1 more) | Revision string stamped into each captured SFT trace record, so a trace can be tied back to the build that produced it. | | `ODYSSEUS_SFT_DISABLE_WORKSPACE_TOOLS` | `'1'` | `src/agent_loop.py:7407` | On by default. Keeps synthetic personal-assistant fixtures out of workspace mode; set 0, false, no or off to let them through. | -| `ODYSSEUS_SFT_FORCE_UTC_TIMEZONE` | `'0'` | `routes/chat_routes.py:2070` | Truthy forces `sft_` accounts to UTC for deterministic batch generation. Interactive accounts still follow the browser timezone. | +| `ODYSSEUS_SFT_FORCE_UTC_TIMEZONE` | `'0'` | `routes/chat_routes.py:2080` | Truthy forces `sft_` accounts to UTC for deterministic batch generation. Interactive accounts still follow the browser timezone. | | `ODYSSEUS_SFT_TRACE_CAPTURE` | `'1'` | `routes/chat_helpers.py:161` (+1 more) | On by default, but only for owners whose name starts with `sft_`. Set 0, false, no or off to stop writing training traces. | | `ODYSSEUS_SFT_TRACE_DIR` | *unset* | `routes/chat_helpers.py:195` (+2 more) | Directory the SFT trace JSONL files are written to. Defaults to `sft_traces` under the data directory. | | `ODYSSEUS_SKIP_RUN_HINT` | *unset* | `setup.py:284` | Any non-empty value suppresses the `start the server with` hint at the end of setup. `start-macos.sh` sets it because it starts the server itself. |