diff --git a/anton/cloud_turn/__init__.py b/anton/cloud_turn/__init__.py new file mode 100644 index 00000000..653c54ec --- /dev/null +++ b/anton/cloud_turn/__init__.py @@ -0,0 +1,13 @@ +"""Cloud turn: run one full anton turn inside a sandbox pod. + +`python -m anton.cloud_turn` reads a :class:`TurnRequestV1` as one JSON line on +stdin, runs the turn to completion against the mounted workspace with a +cloud-safe session, and emits `delta` / `turn_completed` / `turn_failed` JSONL +events on stdout (diagnostics to stderr). It is the headless counterpart of the +desktop CLI host: the same ChatSession, built cloud-safe. +""" + +from anton.cloud_turn.contract import TurnRequestV1 +from anton.cloud_turn.session import build_cloud_chat_session + +__all__ = ["TurnRequestV1", "build_cloud_chat_session"] diff --git a/anton/cloud_turn/__main__.py b/anton/cloud_turn/__main__.py new file mode 100644 index 00000000..fc0e631c --- /dev/null +++ b/anton/cloud_turn/__main__.py @@ -0,0 +1,127 @@ +"""`python -m anton.cloud_turn` - the sandbox-pod turn entrypoint. + +Contract (matches scratchpad-controller + cowork-server): + stdin : ONE newline-terminated TurnRequestV1 JSON line (controller closes stdin) + stdout: JSONL events - `delta` / `turn_completed` / `turn_failed`, nothing else + stderr: diagnostic logs + full tracebacks + exit : 0 (the controller detects the terminal from the event, not the code) + +stdout is isolated at the OS file-descriptor level: the real FD 1 is duplicated +to a private, non-inheritable descriptor used only for protocol events, and +process FD 1 is redirected to stderr. So `print()`, `os.write(1, ...)`, +native-library writes, and child-process stdout can never corrupt the stream. +""" + +from __future__ import annotations + +import asyncio +import contextlib +import inspect +import json +import logging +import os +import sys + +from anton.cloud_turn.contract import TurnRequestV1 +from anton.cloud_turn.session import build_cloud_chat_session + +logger = logging.getLogger(__name__) + +#: Bound the single request line so a malformed/huge stdin can't exhaust memory. +MAX_REQUEST_BYTES = 10 * 1024 * 1024 +#: Keep wire error strings short and log-safe. +MAX_ERROR_MESSAGE_CHARS = 300 + + +def _scrub(exc: Exception) -> str: + """Short, credential-scrubbed error string for the wire. Full traceback + stays on stderr (logged by the caller).""" + from anton.utils.datasources import scrub_credentials + + text = scrub_credentials(f"{type(exc).__name__}: {exc}") + if len(text) > MAX_ERROR_MESSAGE_CHARS: + text = text[: MAX_ERROR_MESSAGE_CHARS - 1] + "…" + return text + + +@contextlib.contextmanager +def _isolated_protocol_stdout(): + """OS-level stdout isolation. Yields ``emit(event: dict)`` writing JSONL to + the saved protocol descriptor; everything else (FD 1) goes to stderr.""" + sys.stdout.flush() + sys.stderr.flush() + stderr_fd = sys.stderr.fileno() + + protocol_fd = os.dup(1) + os.set_inheritable(protocol_fd, False) # children never inherit the protocol channel + os.dup2(stderr_fd, 1) # any write to fd 1 now lands on stderr + saved_sys_stdout = sys.stdout + sys.stdout = sys.stderr + logging.basicConfig(stream=sys.stderr, level=logging.INFO) + + def emit(event: dict) -> None: + data = (json.dumps(event) + "\n").encode("utf-8") + view = memoryview(data) + while view: # os.write may partial-write; loop until fully flushed + n = os.write(protocol_fd, view) + view = view[n:] + + try: + yield emit + finally: + with contextlib.suppress(Exception): + sys.stderr.flush() + sys.stdout = saved_sys_stdout + os.close(protocol_fd) + + +async def _close(session) -> None: + close = getattr(session, "close", None) + if close is None: + return + try: + result = close() + if inspect.isawaitable(result): + await result + except Exception: + logger.warning("cloud session close failed (non-fatal)", exc_info=True) + + +async def stream_turn(raw_line: str, emit, session_builder=None) -> None: + """Parse the request, run one turn, and emit exactly one terminal event. + + Streaming: assistant text is emitted as ``delta`` events as it arrives, then + a bare ``turn_completed``. Any failure (parse or turn) -> one ``turn_failed`` + with a scrubbed error string. + """ + from anton.core.llm.provider import StreamTextDelta + + builder = session_builder or build_cloud_chat_session + session = None + try: + req = TurnRequestV1.from_json(raw_line) + session = builder(req) + async for event in session.turn_stream(req.input): + if isinstance(event, StreamTextDelta): + emit({"kind": "delta", "text": event.text or ""}) + emit({"kind": "turn_completed"}) + except Exception as exc: + # Full traceback -> stderr only; wire carries a short scrubbed string. + logger.exception("cloud turn failed") + emit({"kind": "turn_failed", "error": _scrub(exc)}) + finally: + if session is not None: + await _close(session) + + +def main(argv: list[str] | None = None) -> int: + with _isolated_protocol_stdout() as emit: + # One bounded line (the controller writes a single JSON line + \n, then + # closes stdin). ``readline`` returns on the newline without blocking. + raw_line = sys.stdin.readline(MAX_REQUEST_BYTES + 1) + asyncio.run(stream_turn(raw_line, emit)) + return 0 + + +if __name__ == "__main__": + sys.exit(main()) diff --git a/anton/cloud_turn/contract.py b/anton/cloud_turn/contract.py new file mode 100644 index 00000000..833d05c6 --- /dev/null +++ b/anton/cloud_turn/contract.py @@ -0,0 +1,45 @@ +"""Wire contract between the scratchpad-controller and the pod entrypoint. + +Matches what the controller sends (`scratchpad_controller.anton_turn.request_line`) +and what cowork-server consumes off the reply stream. Intentionally minimal and +data-only: the entrypoint reads ONE newline-terminated JSON line on stdin. + +Events written back on stdout (JSONL) are the three the controller translates: + {"kind": "delta", "text": "..."} - streamed assistant text, one per chunk + {"kind": "turn_completed"} - terminal success (no payload) + {"kind": "turn_failed", "error": "..."} - terminal failure (scrubbed string) +""" + +from __future__ import annotations + +import json +from dataclasses import dataclass, field + + +@dataclass +class TurnRequestV1: + """One turn to run in the pod. Sent as a single JSON line on stdin.""" + + protocol_version: int + conversation_id: str + input: str + #: Mount path the controller passes; the pod uses its own trusted mount and + #: does not act on this value (kept for wire-compatibility). See session.py. + workspace_path: str | None = None + #: Optional model override; None uses the settings default. + model: str | None = None + #: DB-authoritative ordered history ({"role","content"} dicts). The pod never + #: loads its own history; cowork-server owns persistence. + history: list = field(default_factory=list) + + @staticmethod + def from_json(raw: str) -> "TurnRequestV1": + d = json.loads(raw) + return TurnRequestV1( + protocol_version=int(d["protocol_version"]), + conversation_id=str(d["conversation_id"]), + input=str(d["input"]), + workspace_path=d.get("workspace_path"), + model=d.get("model"), + history=d.get("history") or [], + ) diff --git a/anton/cloud_turn/session.py b/anton/cloud_turn/session.py new file mode 100644 index 00000000..ae7b24eb --- /dev/null +++ b/anton/cloud_turn/session.py @@ -0,0 +1,130 @@ +"""Cloud-safe ChatSession builder. + +Assembles a :class:`~anton.core.session.ChatSession` scoped to the tenant's +mounted workspace with the desktop-only, tenant-leaky behaviours OFF. It does +NOT call the desktop ``build_chat_session`` (which loads workspace ``.env``, +uses ``~/.anton`` personal memory, and injects vault creds into ``os.environ``). + +Safety posture (all internal — nothing here is on the wire): + +* Trusted pod-side workspace mount, never taken from the request. +* No dotenv loading (``AntonSettings(_env_file=None)``, shared into Workspace). +* Personal memory / connectors / data-vault / disk history OFF. +* Scratchpad subprocess env built from a non-secret allowlist, so generated + code can't read the provider key (interim until the Plan-5 gateway removes + the key from the pod entirely). +* Only reviewed, headless-safe tools are exposed (scratchpad + artifacts). +""" + +from __future__ import annotations + +import logging +import os +from pathlib import Path +from typing import TYPE_CHECKING + +from anton.cloud_turn.contract import TurnRequestV1 + +if TYPE_CHECKING: + from anton.core.session import ChatSession + +logger = logging.getLogger(__name__) + +#: Trusted mount path — pod-side config, never from the wire request. +DEFAULT_CLOUD_WORKSPACE_PATH = "/workspace" +#: Operator/CI override for the mount path (pod-side env var, not request data). +_WORKSPACE_PATH_ENV = "ANTON_CLOUD_WORKSPACE_PATH" + +#: The only tools exposed in a cloud turn: scratchpad + the workspace-scoped +#: artifact tools. Everything else core registers is dropped. +CLOUD_TOOL_ALLOWLIST = frozenset( + { + "scratchpad", + "create_artifact", + "list_artifacts", + "open_artifact", + "update_artifact", + } +) + + +def resolve_trusted_workspace_path() -> Path: + """Resolve the trusted workspace mount path (never from the wire request). + + Reads :data:`_WORKSPACE_PATH_ENV` or falls back to + :data:`DEFAULT_CLOUD_WORKSPACE_PATH`; rejects relative paths and ``..``, + then canonicalises so downstream containment checks compare a real path. + """ + raw = os.environ.get(_WORKSPACE_PATH_ENV) or DEFAULT_CLOUD_WORKSPACE_PATH + if not os.path.isabs(raw): + raise ValueError( + f"trusted workspace path must be absolute, got {raw!r} " + f"(set {_WORKSPACE_PATH_ENV} to an absolute path)" + ) + if ".." in Path(raw).parts: + raise ValueError(f"trusted workspace path must not contain '..': {raw!r}") + resolved = Path(raw).resolve() + resolved.mkdir(parents=True, exist_ok=True) + if not resolved.is_dir(): + raise ValueError(f"trusted workspace path is not a directory: {resolved}") + return resolved + + +def build_cloud_chat_session(request: TurnRequestV1) -> "ChatSession": + """Assemble a cloud-safe ChatSession for one turn. + + History is DB-authoritative (from the request). The workspace path is the + trusted pod mount, NOT ``request.workspace_path``. + """ + from anton.config.settings import AntonSettings + from anton.core.backends.local import sanitized_scratchpad_runtime_factory + from anton.core.llm.client import LLMClient + from anton.core.session import ChatSession, ChatSessionConfig + from anton.workspace import Workspace + + base = resolve_trusted_workspace_path() + + # `_env_file=None`: never load the AntonSettings .env chain (~/.anton/.env, + # ~/.cowork/.env, /workspace/.env). Same object passed to Workspace so it + # doesn't build a second, dotenv-loading one. + settings = AntonSettings(_env_file=None) + settings.resolve_workspace(str(base)) + if request.model: + settings.planning_model = request.model + # Skills stay in the workspace, never the pod-shared ~/.anton. + settings.skills_root = base / ".anton" / "skills" + + workspace = Workspace(base, settings=settings) + workspace.initialize() + # No apply_env_to_process(): loading workspace .env into the process env + # would expose tenant secrets to cell code. + + llm_client = LLMClient.from_settings(settings) + + config = ChatSessionConfig( + llm_client=llm_client, + settings=settings, + workspace=workspace, + session_id=request.conversation_id, + harness="cloud", + # DB-authoritative history; the pod never loads its own. + initial_history=list(request.history) if request.history else None, + console=None, # headless + cortex=None, # personal memory OFF + episodic=None, + self_awareness=None, + data_vault=None, # connectors OFF + history_store=None, # disk history OFF (DB authoritative) + tools=[], # no host connector/publish tools + tool_allowlist=CLOUD_TOOL_ALLOWLIST, # only reviewed tools survive the build + runtime_factory=sanitized_scratchpad_runtime_factory, # secret-free scratchpad env + web_search_enabled=False, + web_fetch_enabled=False, + ) + + session = ChatSession(config) + logger.info( + "cloud session built conversation=%s workspace=%s tools=%s", + request.conversation_id, base, sorted(CLOUD_TOOL_ALLOWLIST), + ) + return session diff --git a/anton/core/artifacts/store.py b/anton/core/artifacts/store.py index 19cb4672..135d60b1 100644 --- a/anton/core/artifacts/store.py +++ b/anton/core/artifacts/store.py @@ -117,6 +117,13 @@ def ensure_root(self) -> Path: return self._root def folder_for(self, slug: str) -> Path: + # `open`/`update` take the slug straight from tool input. Reject anything + # that resolves outside the root (``..``, absolute, or symlink escape); + # nested-but-contained paths are fine (create never emits them). + candidate = (self._root / slug).resolve() + root = self._root.resolve() + if candidate != root and root not in candidate.parents: + raise ValueError(f"artifact slug escapes the workspace: {slug!r}") return self._root / slug def metadata_path(self, slug: str) -> Path: @@ -385,7 +392,12 @@ def _save(self, artifact: Artifact) -> None: self.readme_path(artifact.slug).write_text(readme, encoding="utf-8") def _load_silent(self, slug: str) -> Artifact | None: - path = self.metadata_path(slug) + try: + path = self.metadata_path(slug) + except ValueError: + # Escaping slug → treat as "no such artifact" so open/update stay graceful. + logger.warning("Rejected out-of-workspace artifact slug %r", slug) + return None if not path.is_file(): return None try: diff --git a/anton/core/backends/local.py b/anton/core/backends/local.py index d73ee85c..3bfbffc6 100644 --- a/anton/core/backends/local.py +++ b/anton/core/backends/local.py @@ -77,6 +77,47 @@ def _utf8_env(base: "os._Environ[str] | dict[str, str]") -> dict[str, str]: return env +# Sanitized scratchpad env contract: in sanitize_env mode the child inherits +# ONLY these non-secret names (and only if the parent has them) — never a +# provider key, ANTON_* credential, gateway token, or datasource secret. +# Python runtime vars (PYTHONUTF8, PYTHONPATH) are set explicitly in start(), +# never inherited, so nothing can redirect imports. +_SCRATCHPAD_ENV_ALLOWLIST = frozenset( + { + # process / tooling + "PATH", + "HOME", + "USER", + "LOGNAME", + "SHELL", + # temp dirs + "TMPDIR", + "TMP", + "TEMP", + # locale / tty + "LANG", + "LC_ALL", + "LC_CTYPE", + "LANGUAGE", + "TZ", + "TERM", + # TLS trust roots (HTTPS verification from cell code) + "SSL_CERT_FILE", + "SSL_CERT_DIR", + "REQUESTS_CA_BUNDLE", + "CURL_CA_BUNDLE", + # Windows: required for subprocess/socket startup + "SYSTEMROOT", + "SYSTEMDRIVE", + } +) + + +def _sanitized_parent_env() -> dict[str, str]: + """Parent env reduced to the non-secret allowlist.""" + return {k: v for k, v in os.environ.items() if k in _SCRATCHPAD_ENV_ALLOWLIST} + + _MAX_OUTPUT = 10_000 @@ -95,6 +136,7 @@ def __init__( coding_base_url: str, cells: list[Cell] | None = None, workspace_path: Path | None = None, + sanitize_env: bool = False, _venvs_base: Path | None = None, ) -> None: super().__init__( @@ -114,6 +156,9 @@ def __init__( # is a "where to put scratchpad venvs" hint; the explicit # arg is "the agent's project, when known". self._explicit_workspace_path: Path | None = workspace_path + # Cloud/headless: build the child env from the non-secret allowlist. + # Default False = desktop behaviour (full inherited env) unchanged. + self._sanitize_env: bool = sanitize_env self._proc: asyncio.subprocess.Process | None = None self._boot_path: str | None = None self._venv_dir: str | None = None @@ -390,59 +435,64 @@ async def start(self) -> None: os.close(fd) self._boot_path = path - # Force UTF-8 mode in the child so its I/O never depends on the host - # code page (ENG-824). - env = _utf8_env(os.environ) + # Force UTF-8 in the child (ENG-824). In sanitized mode the base is the + # non-secret allowlist, not the full parent env. + env = _utf8_env( + _sanitized_parent_env() if self._sanitize_env else os.environ + ) if self._coding_model: env["ANTON_SCRATCHPAD_MODEL"] = self._coding_model if self._coding_provider: env["ANTON_SCRATCHPAD_PROVIDER"] = self._coding_provider - if "ANTHROPIC_API_KEY" not in env and "ANTON_ANTHROPIC_API_KEY" in env: - env["ANTHROPIC_API_KEY"] = env["ANTON_ANTHROPIC_API_KEY"] - if "OPENAI_API_KEY" not in env and "ANTON_OPENAI_API_KEY" in env: - env["OPENAI_API_KEY"] = env["ANTON_OPENAI_API_KEY"] - if "OPENAI_BASE_URL" not in env and "ANTON_OPENAI_BASE_URL" in env: - env["OPENAI_BASE_URL"] = env["ANTON_OPENAI_BASE_URL"] - if ( - "OPENAI_API_KEY" not in env - and "ANTON_MINDS_API_KEY" in env - and self._coding_provider == "openai-compatible" - ): - env["OPENAI_API_KEY"] = env["ANTON_MINDS_API_KEY"] - if ( - "OPENAI_BASE_URL" not in env - and "ANTON_MINDS_URL" in env - and self._coding_provider == "openai-compatible" - ): - # Host-aware (ENG-436): api.mindshub.ai serves /v1, legacy - # mdb.ai serves /api/v1. The previous hardcoded /api/v1 was - # wrong for mindshub. Mirrors config/settings.py + - # cowork-server minds_chat_base_url. - _minds_base = env["ANTON_MINDS_URL"].rstrip("/") - if _minds_base.endswith("/v1"): - env["OPENAI_BASE_URL"] = _minds_base - elif "mdb.ai" in _minds_base: - env["OPENAI_BASE_URL"] = f"{_minds_base}/api/v1" - else: - env["OPENAI_BASE_URL"] = f"{_minds_base}/v1" - if self._coding_api_key: - sdk_key = { - "anthropic": "ANTHROPIC_API_KEY", - "openai": "OPENAI_API_KEY", - "openai-compatible": "OPENAI_API_KEY", - }.get(self._coding_provider, "") - if sdk_key: - env[sdk_key] = self._coding_api_key - if self._coding_provider in ("openai", "openai-compatible"): - base_url = ( - self._coding_base_url - or env.get("ANTON_OPENAI_BASE_URL") - or env.get("OPENAI_BASE_URL") - or "" - ) - if base_url: - env["OPENAI_BASE_URL"] = base_url - env["ANTON_OPENAI_BASE_URL"] = base_url + # Provider-credential propagation is desktop-only; sanitized mode + # withholds every model key (nested get_llm() is off in the cloud pod). + if not self._sanitize_env: + if "ANTHROPIC_API_KEY" not in env and "ANTON_ANTHROPIC_API_KEY" in env: + env["ANTHROPIC_API_KEY"] = env["ANTON_ANTHROPIC_API_KEY"] + if "OPENAI_API_KEY" not in env and "ANTON_OPENAI_API_KEY" in env: + env["OPENAI_API_KEY"] = env["ANTON_OPENAI_API_KEY"] + if "OPENAI_BASE_URL" not in env and "ANTON_OPENAI_BASE_URL" in env: + env["OPENAI_BASE_URL"] = env["ANTON_OPENAI_BASE_URL"] + if ( + "OPENAI_API_KEY" not in env + and "ANTON_MINDS_API_KEY" in env + and self._coding_provider == "openai-compatible" + ): + env["OPENAI_API_KEY"] = env["ANTON_MINDS_API_KEY"] + if ( + "OPENAI_BASE_URL" not in env + and "ANTON_MINDS_URL" in env + and self._coding_provider == "openai-compatible" + ): + # Host-aware (ENG-436): api.mindshub.ai serves /v1, legacy + # mdb.ai serves /api/v1. The previous hardcoded /api/v1 was + # wrong for mindshub. Mirrors config/settings.py + + # cowork-server minds_chat_base_url. + _minds_base = env["ANTON_MINDS_URL"].rstrip("/") + if _minds_base.endswith("/v1"): + env["OPENAI_BASE_URL"] = _minds_base + elif "mdb.ai" in _minds_base: + env["OPENAI_BASE_URL"] = f"{_minds_base}/api/v1" + else: + env["OPENAI_BASE_URL"] = f"{_minds_base}/v1" + if self._coding_api_key: + sdk_key = { + "anthropic": "ANTHROPIC_API_KEY", + "openai": "OPENAI_API_KEY", + "openai-compatible": "OPENAI_API_KEY", + }.get(self._coding_provider, "") + if sdk_key: + env[sdk_key] = self._coding_api_key + if self._coding_provider in ("openai", "openai-compatible"): + base_url = ( + self._coding_base_url + or env.get("ANTON_OPENAI_BASE_URL") + or env.get("OPENAI_BASE_URL") + or "" + ) + if base_url: + env["OPENAI_BASE_URL"] = base_url + env["ANTON_OPENAI_BASE_URL"] = base_url uv = self._find_uv() if uv: env["ANTON_UV_PATH"] = uv @@ -817,3 +867,27 @@ def local_scratchpad_runtime_factory( cells=cells, workspace_path=workspace_path, ) + + +def sanitized_scratchpad_runtime_factory( + *, + name: str, + coding_provider: str, + coding_model: str, + coding_api_key: str, + coding_base_url: str, + cells: list[Cell] | None, + workspace_path: Path | None, +) -> ScratchpadRuntime: + """Cloud/headless factory: like the local factory but with + ``sanitize_env=True`` so the child env carries no secrets.""" + return LocalScratchpadRuntime( + name=name, + coding_provider=coding_provider, + coding_model=coding_model, + coding_api_key=coding_api_key, + coding_base_url=coding_base_url, + cells=cells, + workspace_path=workspace_path, + sanitize_env=True, + ) diff --git a/anton/core/session.py b/anton/core/session.py index 77eece9a..24846acf 100644 --- a/anton/core/session.py +++ b/anton/core/session.py @@ -299,6 +299,10 @@ class ChatSessionConfig: # settings' `router_enabled` (ANTON_ROUTER_ENABLED); hosts pass an # explicit bool to override per session. router_enabled: bool | None = None + # When set, only these tool names survive the build; ``None`` = full desktop + # set. Applied on every ``_build_tools`` call so a lazy rebuild can't leak a + # non-allowlisted tool. + tool_allowlist: frozenset[str] | None = None class ChatSession: @@ -347,6 +351,7 @@ def __init__(self, config: ChatSessionConfig) -> None: self._act_first = config.act_first self._started_at = config.started_at self._extra_tools = config.tools + self._tool_allowlist = config.tool_allowlist self._workspace = config.workspace self._data_vault = config.data_vault self._console = config.console @@ -879,6 +884,20 @@ def _build_tools(self) -> list[dict]: self._build_core_tools() for tool in self._extra_tools: self.tool_registry.register_tool(tool) + # Enforce the allowlist on every build (None = full desktop set). + if self._tool_allowlist is not None: + built = {t.name for t in self.tool_registry.get_tool_defs()} + # Fail loud on a name that matches no built tool (typo / unavailable + # here) rather than silently dropping it. + unknown = set(self._tool_allowlist) - built + if unknown: + raise ValueError( + "tool_allowlist names not registered in this session: " + + ", ".join(sorted(unknown)) + + f" (available: {', '.join(sorted(built))})" + ) + for name in built - set(self._tool_allowlist): + self.tool_registry.unregister_tool(name) return self.tool_registry.dump() def _build_core_tools(self) -> None: diff --git a/anton/workspace.py b/anton/workspace.py index 1626f751..d3258f21 100644 --- a/anton/workspace.py +++ b/anton/workspace.py @@ -27,14 +27,17 @@ class Workspace: """Manages the .anton/ workspace directory and its files.""" - def __init__(self, base: Path) -> None: + def __init__(self, base: Path, settings: AntonSettings | None = None) -> None: self._base = base self._anton_dir = base / ".anton" self._anton_md = self._anton_dir / "anton.md" self._env_file = self._anton_dir / ".env" self._anton_md_last_read: datetime | None = None - settings = AntonSettings() + # Reuse a caller's settings so the cloud pod's dotenv-disabled + # AntonSettings isn't bypassed by a second one built here. Only + # `artifacts_dir` is read; None = desktop behaviour unchanged. + settings = settings or AntonSettings() self._artifacts_dir = self._anton_dir / settings.artifacts_dir @property diff --git a/pyproject.toml b/pyproject.toml index 4901fe87..a394d5b8 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -58,6 +58,7 @@ asyncio_mode = "auto" pythonpath = ["."] markers = [ "stub_only: test requires stub server (skipped when --live is passed)", + "slow: test spawns a real scratchpad venv/subprocess (seconds)", ] [tool.hatch.build.targets.wheel] diff --git a/tests/cloud_turn_fake_entry.py b/tests/cloud_turn_fake_entry.py new file mode 100644 index 00000000..bf7b6ffa --- /dev/null +++ b/tests/cloud_turn_fake_entry.py @@ -0,0 +1,112 @@ +"""Deterministic subprocess entrypoint for cloud-turn E2E / FD tests. + +Run as ``python tests/cloud_turn_fake_entry.py`` — it installs a deterministic +seam, then calls the REAL ``anton.cloud_turn.__main__.main()`` so the full +process boundary (FD isolation, stdin parse, runner lifecycle, scratchpad) runs +for real. NOT collected by pytest (no ``test_`` prefix). + +Modes (env ``CLOUD_TURN_FAKE_MODE``): +* ``model`` — real ChatSession + real scratchpad, but the LLM is a fake provider + scripted by ``CLOUD_TURN_FAKE_SCRIPT`` (JSON list of steps). No network. +* ``stray`` — replace the session builder with one whose turn writes stray + output to stdout/FD 1/logging, then completes (proves FD isolation). +* ``stray_fail`` — like ``stray`` but the turn raises after the stray output. +""" + +from __future__ import annotations + +import json +import logging +import os +import sys + + +def _install_fake_model() -> None: + import anton.core.llm.client as client_mod + from anton.core.llm.client import LLMClient + from anton.core.llm.provider import ( + LLMProvider, + LLMResponse, + ProviderConnectionInfo, + ToolCall, + Usage, + ) + + script = json.loads(os.environ.get("CLOUD_TURN_FAKE_SCRIPT", "[]")) + + class _FakeProvider(LLMProvider): + name = "fake" + + def __init__(self) -> None: + self._i = 0 + + async def complete(self, *, model, system, messages, tools=None, + tool_choice=None, max_tokens=4096, native_web_tools=None): + step = script[min(self._i, len(script) - 1)] if script else {"text": ""} + self._i += 1 + tool_calls = [] + if "tool" in step: + t = step["tool"] + tool_calls = [ToolCall(id=t.get("id", "t1"), name=t["name"], input=t["input"])] + return LLMResponse( + content=step.get("text", ""), + tool_calls=tool_calls, + usage=Usage(context_pressure=0.0), + ) + + def export_connection_info(self): + return ProviderConnectionInfo(provider="fake", api_key="fake") + + def _fake_from_settings(cls, settings): + prov = _FakeProvider() + return LLMClient( + planning_provider=prov, planning_model="fake-model", + coding_provider=prov, coding_model="fake-model", + ) + + client_mod.LLMClient.from_settings = classmethod(_fake_from_settings) + + +def _install_stray_session(fail: bool) -> None: + import anton.cloud_turn.__main__ as entry_mod + + class _StraySession: + def __init__(self) -> None: + self.history = [] + self.closed = False + + async def turn_stream(self, user_input, **kwargs): + # Every channel that must NOT reach the protocol stream: + print("STRAY via print()") # Python stdout + sys.stdout.write("STRAY via sys.stdout.write\n") # Python stdout + os.write(1, b"STRAY via os.write(1)\n") # direct FD 1 (native-style) + logging.getLogger("some.library").warning("STRAY via logging") + if fail: + raise RuntimeError("boom after stray output") + if False: # make this an async generator + yield + + def close(self): + self.closed = True + + entry_mod.build_cloud_chat_session = lambda request: _StraySession() + + +def main() -> int: + mode = os.environ.get("CLOUD_TURN_FAKE_MODE", "model") + if mode == "model": + _install_fake_model() + elif mode == "stray": + _install_stray_session(fail=False) + elif mode == "stray_fail": + _install_stray_session(fail=True) + else: + raise SystemExit(f"unknown CLOUD_TURN_FAKE_MODE={mode!r}") + + from anton.cloud_turn.__main__ import main as real_main + + return real_main([]) + + +if __name__ == "__main__": + sys.exit(main()) diff --git a/tests/test_cloud_turn_entrypoint.py b/tests/test_cloud_turn_entrypoint.py new file mode 100644 index 00000000..b3f30170 --- /dev/null +++ b/tests/test_cloud_turn_entrypoint.py @@ -0,0 +1,92 @@ +"""Entrypoint wire contract: request parsing + streaming event emission. + +Offline: a fake session (scripted stream events) replaces ChatSession, so no LLM +key or scratchpad is needed. Mirrors the controller/cowork contract: +`delta` -> ... -> `turn_completed` | `turn_failed`. +""" + +from __future__ import annotations + +import asyncio +import json + +from anton.cloud_turn.contract import TurnRequestV1 +from anton.cloud_turn.__main__ import stream_turn +from anton.core.llm.provider import StreamTextDelta + + +# ── contract parsing ───────────────────────────────────────────────────────── + +def test_from_json_parses_full_request(): + req = TurnRequestV1.from_json(json.dumps({ + "protocol_version": 1, "conversation_id": "c", "input": "hi", + "workspace_path": "/workspace", "model": "m", + "history": [{"role": "user", "content": "prev"}], + })) + assert req.conversation_id == "c" + assert req.input == "hi" + assert req.history == [{"role": "user", "content": "prev"}] + + +def test_from_json_history_defaults_empty(): + req = TurnRequestV1.from_json('{"protocol_version":1,"conversation_id":"c","input":"hi"}') + assert req.history == [] + assert req.workspace_path is None and req.model is None + + +# ── streaming event emission ───────────────────────────────────────────────── + +class _FakeSession: + def __init__(self, deltas=(), raise_on_stream=None): + self._deltas = list(deltas) + self._raise = raise_on_stream + self.closed = False + + async def turn_stream(self, user_input, **kwargs): + if self._raise: + raise self._raise + for d in self._deltas: + yield StreamTextDelta(text=d) + + def close(self): + self.closed = True + + +def _drive(session, req_json='{"protocol_version":1,"conversation_id":"c","input":"hi"}'): + events = [] + asyncio.run(stream_turn(req_json, events.append, session_builder=lambda r: session)) + return events + + +def test_emits_deltas_then_completed(): + session = _FakeSession(deltas=["he", "llo"]) + events = _drive(session) + assert events == [ + {"kind": "delta", "text": "he"}, + {"kind": "delta", "text": "llo"}, + {"kind": "turn_completed"}, + ] + assert session.closed is True # session closed on success + + +def test_text_only_no_deltas_still_completes(): + events = _drive(_FakeSession(deltas=[])) + assert events == [{"kind": "turn_completed"}] + + +def test_turn_failure_is_terminal_and_scrubbed(): + session = _FakeSession(raise_on_stream=RuntimeError("boom sk-ant-" + "A" * 80)) + events = _drive(session) + assert len(events) == 1 + assert events[0]["kind"] == "turn_failed" + assert "boom" in events[0]["error"] + assert "sk-ant-" + "A" * 80 not in events[0]["error"] # credential scrubbed + assert session.closed is True # closed even on failure + + +def test_bad_request_is_a_single_turn_failed(): + events = [] + asyncio.run(stream_turn("{not valid json", events.append)) + assert len(events) == 1 + assert events[0]["kind"] == "turn_failed" + assert events[0]["error"] # non-empty, scrubbed diff --git a/tests/test_cloud_turn_process.py b/tests/test_cloud_turn_process.py new file mode 100644 index 00000000..27862ad7 --- /dev/null +++ b/tests/test_cloud_turn_process.py @@ -0,0 +1,170 @@ +"""Real-subprocess tests for the cloud-turn process boundary. + +Runs the actual entrypoint (via cloud_turn_fake_entry.py, which calls the real +main()) as a child process, so FD-level stdout isolation, fresh +process/scratchpad state, and the deterministic E2E are proven end to end. +Event kinds match the controller contract: delta / turn_completed / turn_failed. +""" + +from __future__ import annotations + +import json +import os +import subprocess +import sys +from pathlib import Path + +import pytest + +_HARNESS = str(Path(__file__).parent / "cloud_turn_fake_entry.py") + + +def _run_cli(request, *, workspace, mode="model", script=None, timeout=60): + """Run the entrypoint as a subprocess; return (exit_code, events, stdout, stderr). + Parsing stdout as JSONL is itself the assertion that stdout is clean.""" + env = os.environ.copy() + env["ANTON_CLOUD_WORKSPACE_PATH"] = str(workspace) + env["CLOUD_TURN_FAKE_MODE"] = mode + if script is not None: + env["CLOUD_TURN_FAKE_SCRIPT"] = json.dumps(script) + + stdin = request if isinstance(request, str) else json.dumps(request) + proc = subprocess.run( + [sys.executable, _HARNESS], + input=stdin.encode("utf-8"), + stdout=subprocess.PIPE, stderr=subprocess.PIPE, env=env, timeout=timeout, + ) + lines = [ln for ln in proc.stdout.decode("utf-8").splitlines() if ln.strip()] + events = [json.loads(ln) for ln in lines] # raises if stdout isn't clean JSONL + return proc.returncode, events, proc.stdout.decode(), proc.stderr.decode() + + +def _req(**over): + body = {"protocol_version": 1, "conversation_id": "c", "input": "hi"} + body.update(over) + return body + + +# ── FD-level stdout isolation ──────────────────────────────────────────────── + +def test_stray_stdout_never_corrupts_protocol(tmp_path): + """print(), sys.stdout.write, os.write(1, …) and library logging during the + turn must all land on stderr, leaving stdout a clean protocol stream.""" + _, events, stdout, stderr = _run_cli(_req(), workspace=tmp_path, mode="stray") + assert [e["kind"] for e in events] == ["turn_completed"] + assert "STRAY" not in stdout # nothing stray on the protocol channel + assert "STRAY via os.write(1)" in stderr # direct FD-1 write redirected to stderr + assert "STRAY via print()" in stderr + + +def test_stray_then_failure_stays_clean(tmp_path): + _, events, stdout, _ = _run_cli(_req(), workspace=tmp_path, mode="stray_fail") + assert [e["kind"] for e in events] == ["turn_failed"] + assert "STRAY" not in stdout + assert "Traceback" not in stdout # no traceback leaks onto the wire + assert "boom" in events[-1]["error"] + + +def test_malformed_input_clean_protocol(tmp_path): + _, events, stdout, _ = _run_cli("{not valid json", workspace=tmp_path) + assert len(events) == 1 and events[0]["kind"] == "turn_failed" + assert events[0]["error"] + + +# ── deterministic E2E (real CLI, fake model, no network) ───────────────────── + +def test_e2e_text_only_turn(tmp_path): + """Full path: JSON on stdin -> parse -> cloud-safe session -> streaming + delta events -> turn_completed.""" + _, events, *_ = _run_cli( + _req(input="What is 2 + 2?"), + workspace=tmp_path, mode="model", script=[{"text": "The answer is 4."}], + ) + kinds = [e["kind"] for e in events] + assert kinds[-1] == "turn_completed" + text = "".join(e["text"] for e in events if e["kind"] == "delta") + assert text == "The answer is 4." + + +# ── fresh process + scratchpad state per invocation ────────────────────────── +# Tool output is not on the wire (his contract), so we observe via workspace +# files: Turn A sets a variable + writes a file; Turn B (new process, same +# workspace) reports whether the variable survived. Files persist, runtime does not. + +_CELL_A = ( + "X_SENTINEL = 4242\n" + "open('a_done.txt', 'w').write('a-ok')\n" + "print('set')\n" +) +_CELL_B = ( + "open('b_result.txt', 'w').write('has_x=' + str('X_SENTINEL' in dir()))\n" + "print('done')\n" +) + + +def _scratchpad_step(code): + return {"tool": {"name": "scratchpad", "input": { + "action": "exec", "name": "main", "code": code, + "one_line_description": "test cell", + }}} + + +@pytest.mark.slow +def test_fresh_scratchpad_state_across_processes(tmp_path): + codeA, eventsA, *_ = _run_cli( + _req(conversation_id="A", input="set state"), + workspace=tmp_path, mode="model", timeout=180, + script=[_scratchpad_step(_CELL_A), {"text": "did A"}], + ) + assert eventsA[-1]["kind"] == "turn_completed", eventsA + # Turn A's scratchpad ran and its file persists in the workspace. + assert (tmp_path / "a_done.txt").read_text() == "a-ok" + + codeB, eventsB, *_ = _run_cli( + _req(conversation_id="B", input="read state"), + workspace=tmp_path, mode="model", timeout=180, + script=[_scratchpad_step(_CELL_B), {"text": "did B"}], + ) + assert eventsB[-1]["kind"] == "turn_completed", eventsB + # New process => fresh scratchpad namespace: Turn A's variable is gone. + assert (tmp_path / "b_result.txt").read_text() == "has_x=False" + + +# ── session.close() terminates the inner scratchpad (no orphans) ───────────── + +async def _real_cloud_session(tmp_path, monkeypatch): + import anton.core.llm.client as llm_client_mod + from unittest.mock import AsyncMock, MagicMock + + from anton.core.llm.provider import ProviderConnectionInfo + from anton.cloud_turn.contract import TurnRequestV1 + from anton.cloud_turn.session import build_cloud_chat_session + + monkeypatch.setenv("ANTON_CLOUD_WORKSPACE_PATH", str(tmp_path)) + + def _mk(cls, settings): + llm = AsyncMock() + llm.coding_provider = MagicMock() + llm.coding_provider.export_connection_info = MagicMock( + return_value=ProviderConnectionInfo(provider="anthropic", api_key="test")) + llm.coding_model = "m" + llm.planning_provider = MagicMock() + llm.planning_provider.native_web_tools = MagicMock(return_value=set()) + return llm + + monkeypatch.setattr(llm_client_mod.LLMClient, "from_settings", classmethod(_mk)) + req = TurnRequestV1(protocol_version=1, conversation_id="c", input="hi") + return build_cloud_chat_session(req) + + +@pytest.mark.slow +async def test_session_close_terminates_scratchpad(tmp_path, monkeypatch): + session = await _real_cloud_session(tmp_path, monkeypatch) + pad = await session._scratchpads.get_or_create("main") + await pad.execute("x = 1") + proc = pad._proc + assert proc is not None and proc.returncode is None # alive + + await session.close() + assert pad._proc is None # manager released it + assert proc.returncode is not None # OS process terminated diff --git a/tests/test_cloud_turn_session.py b/tests/test_cloud_turn_session.py new file mode 100644 index 00000000..1954c96a --- /dev/null +++ b/tests/test_cloud_turn_session.py @@ -0,0 +1,296 @@ +"""Safety tests for the cloud-safe ChatSession builder. + +Two layers: +* Config-level (fast, offline): the ChatSessionConfig handed to ChatSession has + every cross-tenant hazard off and the right allowlist/factory. +* Enforcement-level (real code paths): drive the real lazy _build_tools() and a + real sanitized scratchpad subprocess so a claimed property is proven, not just + configured. +""" + +from __future__ import annotations + +import pytest + +import anton.core.llm.client as llm_client_mod +import anton.core.session as session_mod +from anton.cloud_turn.contract import TurnRequestV1 +from anton.cloud_turn.session import ( + CLOUD_TOOL_ALLOWLIST, + _WORKSPACE_PATH_ENV, + build_cloud_chat_session, + resolve_trusted_workspace_path, +) +from anton.core.backends.local import sanitized_scratchpad_runtime_factory + + +class _FakeSession: + """Placeholder return value for the mocked ChatSession - config-level tests + only inspect the captured ChatSessionConfig, never the session.""" + + +def _build(tmp_path, monkeypatch, **req_overrides): + captured: dict = {} + + def fake_chat_session(config): + captured["config"] = config + return _FakeSession() + + monkeypatch.setattr(session_mod, "ChatSession", fake_chat_session) + monkeypatch.setattr( + llm_client_mod.LLMClient, "from_settings", + classmethod(lambda cls, settings: object()), + ) + monkeypatch.setenv(_WORKSPACE_PATH_ENV, str(tmp_path)) + + body = dict(protocol_version=1, conversation_id="conv_1", input="hello") + body.update(req_overrides) + session = build_cloud_chat_session(TurnRequestV1(**body)) + return session, captured["config"] + + +# ── config-level safety ────────────────────────────────────────────────────── + +def test_all_cross_tenant_hazards_off(tmp_path, monkeypatch): + _, cfg = _build(tmp_path, monkeypatch) + assert cfg.cortex is None # personal memory OFF + assert cfg.episodic is None + assert cfg.self_awareness is None + assert cfg.data_vault is None # connectors / vault OFF + assert cfg.history_store is None # disk history OFF + assert cfg.console is None # headless + assert cfg.tools == [] # no host connector/publish tools + assert cfg.web_search_enabled is False + assert cfg.web_fetch_enabled is False + + +def test_scratchpad_uses_sanitized_factory_and_is_workspace_bound(tmp_path, monkeypatch): + _, cfg = _build(tmp_path, monkeypatch) + assert cfg.runtime_factory is sanitized_scratchpad_runtime_factory + assert cfg.workspace is not None + assert cfg.harness == "cloud" + assert cfg.session_id == "conv_1" + + +def test_db_history_is_seeded_not_loaded(tmp_path, monkeypatch): + _, cfg = _build( + tmp_path, monkeypatch, + history=[{"role": "user", "content": "prior turn"}], + ) + assert cfg.initial_history == [{"role": "user", "content": "prior turn"}] + + +def test_config_uses_explicit_tool_allowlist(tmp_path, monkeypatch): + _, cfg = _build(tmp_path, monkeypatch) + assert cfg.tool_allowlist == CLOUD_TOOL_ALLOWLIST + assert "launch_backend" not in cfg.tool_allowlist + assert "scratchpad" in cfg.tool_allowlist + + +def test_model_override_applied(tmp_path, monkeypatch): + _, cfg = _build(tmp_path, monkeypatch, model="claude-opus-4-8") + assert cfg.settings.planning_model == "claude-opus-4-8" + + +# ── trusted workspace path (never from the wire) ───────────────────────────── + +def test_request_workspace_path_is_ignored_trusted_mount_used(tmp_path, monkeypatch): + # A request may carry workspace_path (cowork sends "/workspace"), but the pod + # uses its OWN trusted mount, not the wire value. + monkeypatch.setenv(_WORKSPACE_PATH_ENV, str(tmp_path)) + _, cfg = _build(tmp_path, monkeypatch, workspace_path="/etc/evil") + assert cfg.workspace.base == tmp_path.resolve() + + +def test_resolver_uses_env_override(tmp_path, monkeypatch): + monkeypatch.setenv(_WORKSPACE_PATH_ENV, str(tmp_path)) + assert resolve_trusted_workspace_path() == tmp_path.resolve() + + +def test_resolver_rejects_relative_path(monkeypatch): + monkeypatch.setenv(_WORKSPACE_PATH_ENV, "not/absolute") + with pytest.raises(ValueError, match="absolute"): + resolve_trusted_workspace_path() + + +def test_resolver_rejects_parent_traversal(monkeypatch): + monkeypatch.setenv(_WORKSPACE_PATH_ENV, "/workspace/../etc") + with pytest.raises(ValueError, match=r"\.\."): + resolve_trusted_workspace_path() + + +def test_artifact_tools_cannot_escape_workspace(tmp_path): + from anton.core.artifacts.store import ArtifactStore + + store = ArtifactStore(tmp_path / "artifacts") + for bad in ("../../etc", "/etc/passwd", "a/../../b"): + with pytest.raises(ValueError, match="escapes"): + store.folder_for(bad) + assert store.open("../../../etc/passwd") is None + assert store.folder_for("my-report").parent == (tmp_path / "artifacts") + + +# ── dotenv never loaded ────────────────────────────────────────────────────── + +def test_apply_env_to_process_never_called(tmp_path, monkeypatch): + import anton.workspace as ws_mod + calls = {"apply_env": 0, "init": 0} + real_init = ws_mod.Workspace.initialize + real_apply = ws_mod.Workspace.apply_env_to_process + + monkeypatch.setattr(ws_mod.Workspace, "initialize", + lambda self: (calls.__setitem__("init", calls["init"] + 1), real_init(self))[1]) + monkeypatch.setattr(ws_mod.Workspace, "apply_env_to_process", + lambda self: (calls.__setitem__("apply_env", calls["apply_env"] + 1), real_apply(self))[1]) + + _build(tmp_path, monkeypatch) + assert calls["init"] == 1 + assert calls["apply_env"] == 0 # .env NOT loaded into process env + + +def test_cloud_settings_ignore_dotenv_files(tmp_path, monkeypatch): + from anton.config.settings import AntonSettings + + monkeypatch.delenv("ANTON_ANTHROPIC_API_KEY", raising=False) + monkeypatch.delenv("ANTHROPIC_API_KEY", raising=False) + + envf = tmp_path / "sentinel.env" + envf.write_text("ANTON_ANTHROPIC_API_KEY=DOTENV_SENTINEL\n") + assert AntonSettings(_env_file=str(envf)).anthropic_api_key == "DOTENV_SENTINEL" + + monkeypatch.chdir(tmp_path) + (tmp_path / ".env").write_text("ANTON_ANTHROPIC_API_KEY=DOTENV_SENTINEL\n") + _, cfg = _build(tmp_path, monkeypatch) + assert cfg.settings.anthropic_api_key != "DOTENV_SENTINEL" + + +# ── real enforcement paths ─────────────────────────────────────────────────── + +def _mock_llm(): + from unittest.mock import AsyncMock, MagicMock + + from anton.core.llm.provider import ProviderConnectionInfo + + llm = AsyncMock() + llm.coding_provider = MagicMock() + llm.coding_provider.export_connection_info = MagicMock( + return_value=ProviderConnectionInfo(provider="anthropic", api_key="test") + ) + llm.coding_model = "claude-sonnet-4-6" + llm.planning_provider = MagicMock() + llm.planning_provider.native_web_tools = MagicMock(return_value=set()) + return llm + + +def _real_session_with_workspace(tmp_path, **cfg_overrides): + from unittest.mock import MagicMock + + from anton.core.session import ChatSession, ChatSessionConfig + from anton.workspace import Workspace + + ws = Workspace(tmp_path) + ws.initialize() + session = ChatSession( + ChatSessionConfig(llm_client=_mock_llm(), workspace=ws, **cfg_overrides) + ) + session._scratchpads = MagicMock(available_packages=[]) + return session + + +def test_final_tool_set_equals_allowlist_after_real_build(tmp_path, monkeypatch): + """The registry is built LAZILY at turn time, so the allowlist must be + enforced by the real _build_tools() - assert the EXACT final tool set.""" + from unittest.mock import MagicMock + + monkeypatch.setenv(_WORKSPACE_PATH_ENV, str(tmp_path)) + monkeypatch.setattr( + llm_client_mod.LLMClient, "from_settings", + classmethod(lambda cls, settings: _mock_llm()), + ) + req = TurnRequestV1(protocol_version=1, conversation_id="c", input="hi") + session = build_cloud_chat_session(req) # REAL ChatSession + session._scratchpads = MagicMock(available_packages=[]) + + session._build_tools() + names = {t.name for t in session.tool_registry.get_tool_defs()} + assert names == set(CLOUD_TOOL_ALLOWLIST), ( + f"cloud tool set drifted from the allowlist: {sorted(names)}" + ) + + +def test_tool_allowlist_none_preserves_desktop_tools(tmp_path): + session = _real_session_with_workspace(tmp_path) # allowlist defaults to None + assert session._tool_allowlist is None + session._build_tools() + names = {t.name for t in session.tool_registry.get_tool_defs()} + assert {"launch_backend", "select_path", "scratchpad", "create_artifact"} <= names + + +def test_sanitized_env_contract(monkeypatch): + from anton.core.backends.local import _SCRATCHPAD_ENV_ALLOWLIST, _sanitized_parent_env + + required = {"PATH": "/usr/bin:/bin", "HOME": "/home/anton", "TMPDIR": "/tmp/x", + "LANG": "en_US.UTF-8", "SSL_CERT_FILE": "/etc/ssl/cert.pem"} + for k, v in required.items(): + assert k in _SCRATCHPAD_ENV_ALLOWLIST, f"{k} must be in the contract" + monkeypatch.setenv(k, v) + secrets = {"ANTHROPIC_API_KEY": "sk-ant-x", "ANTON_MINDS_API_KEY": "anton-x", + "ANTON_GATEWAY_TOKEN": "gw-x", "DS_POSTGRES_MAIN__PASSWORD": "ds-x", + "OPENAI_API_KEY": "sk-oai-x", "SOME_RANDOM_SECRET": "nope"} + for k, v in secrets.items(): + monkeypatch.setenv(k, v) + + env = _sanitized_parent_env() + for k, v in required.items(): + assert env.get(k) == v + for k in secrets: + assert k not in env + assert set(env) <= set(_SCRATCHPAD_ENV_ALLOWLIST) + + +async def test_sanitized_scratchpad_cannot_read_secret_env(tmp_path, monkeypatch): + """Real subprocess: with the sanitized factory, secret sentinels in the + parent env must be UNREADABLE from cell code.""" + import json + + monkeypatch.setenv("ANTHROPIC_API_KEY", "sk-ant-SENTINEL") + monkeypatch.setenv("ANTON_GATEWAY_TOKEN", "GATEWAY-SENTINEL") + monkeypatch.setenv("DS_POSTGRES_MAIN__PASSWORD", "DATASOURCE-SENTINEL") + + pad = sanitized_scratchpad_runtime_factory( + name="cloud-sanitize", coding_provider="anthropic", coding_model="", + coding_api_key="", coding_base_url="", cells=None, workspace_path=tmp_path, + ) + await pad.start() + try: + cell = await pad.execute( + "import os, json\n" + "keys = ['ANTHROPIC_API_KEY','ANTON_GATEWAY_TOKEN','DS_POSTGRES_MAIN__PASSWORD']\n" + "print(json.dumps({k: os.environ.get(k) for k in keys}))" + ) + assert cell.error is None, cell.error + seen = json.loads(cell.stdout.strip()) + assert all(v is None for v in seen.values()), f"secret leaked: {seen}" + finally: + await pad.close() + + +async def test_desktop_scratchpad_env_unchanged(tmp_path, monkeypatch): + """Regression: sanitize_env=False (desktop default) STILL inherits the + parent env - the sanitizer is strictly opt-in.""" + from anton.core.backends.local import local_scratchpad_runtime_factory + + monkeypatch.setenv("MY_DESKTOP_SENTINEL", "DESKTOP-VISIBLE") + pad = local_scratchpad_runtime_factory( + name="desktop-env", coding_provider="anthropic", coding_model="", + coding_api_key="", coding_base_url="", cells=None, workspace_path=tmp_path, + ) + await pad.start() + try: + cell = await pad.execute( + "import os; print(os.environ.get('MY_DESKTOP_SENTINEL', 'MISSING'))" + ) + assert cell.error is None, cell.error + assert cell.stdout.strip() == "DESKTOP-VISIBLE" + finally: + await pad.close()