diff --git a/environments/compact/compact/harness.py b/environments/compact/compact/harness.py index 8fdf309f9..5b173d559 100644 --- a/environments/compact/compact/harness.py +++ b/environments/compact/compact/harness.py @@ -51,4 +51,7 @@ async def launch( {"mcpServers": {name: {"url": url} for name, url in mcp_urls.items()}} ) program = await runtime.prepare_uv_script(PROGRAM_SOURCE, self.config.env) - return await runtime.run_program([*program, trace.task.data.prompt], env) + prompt = trace.task.data.prompt + if not isinstance(prompt, str): + raise ValueError("compact harness requires a string task prompt") + return await runtime.run_program(program, env, stdin=prompt.encode("utf-8")) diff --git a/environments/compact/compact/program.py b/environments/compact/compact/program.py index c9d3736da..ef7d7970c 100644 --- a/environments/compact/compact/program.py +++ b/environments/compact/compact/program.py @@ -99,7 +99,7 @@ async def call_mcp(dispatch: dict, name: str, arguments: dict) -> str: async def main() -> None: - task = sys.argv[1] + task = sys.stdin.read() config = json.loads(os.environ.get("MCP_CONFIG", "{}")) notes: str | None = None # the durable memory carried across turns tool_output: str | None = None # the last tool result, kept for exactly one turn diff --git a/verifiers/v1/harnesses/bash/harness.py b/verifiers/v1/harnesses/bash/harness.py index 071017800..d7643c9b0 100644 --- a/verifiers/v1/harnesses/bash/harness.py +++ b/verifiers/v1/harnesses/bash/harness.py @@ -68,19 +68,13 @@ async def launch( f"--base-url={endpoint}", f"--api-key={secret}", f"--model={ctx.model}", - f"--system-prompt={system_prompt}", + "--payload-stdin", ] if self.config.edit: args.append("--edit") if self.config.search: - # Resolve the key and keep it OUT of the program env: it's handed to the program over - # argv (--serper-key), so popping it here stops the agent's `bash` subprocesses from - # inheriting it via $SERPER_API_KEY / /proc/self/environ. Prefer a key set in the harness - # env (--harness.env / forward_env); fall back to the host env only when the key is - # *absent* (None), not present-but-empty — a rollout setting SERPER_API_KEY="" is - # deliberately masking the host secret, so honor that (the check below then fails loudly - # rather than leaking the host key). The pop is scoped to search=true, so an unrelated - # key forwarded for the agent's own bash-side use is left untouched. + # Keep the search key out of the program environment so agent-spawned Bash commands + # cannot inherit it. It is bounded secret/config data rather than task payload. serper_key = env.pop("SERPER_API_KEY", None) if serper_key is None: serper_key = os.environ.get("SERPER_API_KEY") @@ -90,30 +84,23 @@ async def launch( "(the host env or --harness.env)" ) args += ["--search", f"--serper-key={serper_key}"] - if mcp_urls: - # The program connects to the tool servers over HTTP; hand it a standard - # `mcpServers` URL config (the `mcp` client itself comes from the uv deps). - args.append( - "--mcp-config=" - + json.dumps( - { - "mcpServers": { - name: {"url": url} for name, url in mcp_urls.items() - } - } - ) - ) - if isinstance(prompt, str): - args.append(f"--prompt={prompt}") - elif prompt is not None: - # Base64 images can exceed exec limits, so hand Messages off through a file. - path = f".vf-initial-messages-{trace.id}.json" - await runtime.write( - path, - json.dumps([message_to_wire(m) for m in prompt]).encode(), - ) - args.append(f"--initial-messages-file={path}") + payload = { + "system_prompt": system_prompt, + "prompt": prompt if isinstance(prompt, str) else None, + "initial_messages": ( + [message_to_wire(message) for message in prompt] + if prompt is not None and not isinstance(prompt, str) + else [] + ), + "mcp_config": { + "mcpServers": {name: {"url": url} for name, url in mcp_urls.items()} + }, + } program = await runtime.prepare_uv_script( PROGRAM_SOURCE, self.config.resolved_env ) - return await runtime.run_program([*program, *args], env) + return await runtime.run_program( + [*program, *args], + env, + stdin=json.dumps(payload).encode("utf-8"), + ) diff --git a/verifiers/v1/harnesses/bash/program.py b/verifiers/v1/harnesses/bash/program.py index ff17096ad..10ab7a138 100644 --- a/verifiers/v1/harnesses/bash/program.py +++ b/verifiers/v1/harnesses/bash/program.py @@ -8,6 +8,7 @@ import asyncio import json import subprocess +import sys from contextlib import AsyncExitStack, asynccontextmanager, suppress from pathlib import Path @@ -136,7 +137,11 @@ def run_search(query: str, api_key: str, num_results: int = 5) -> str: def run_bash(command: str) -> str: try: result = subprocess.run( - ["bash", "-c", command], capture_output=True, text=True, timeout=3600 + ["bash"], + input=command, + capture_output=True, + text=True, + timeout=3600, ) return result.stdout + result.stderr except Exception as e: @@ -300,7 +305,9 @@ def parse_args() -> argparse.Namespace: parser.add_argument("--api-key", required=True) parser.add_argument("--model", required=True) parser.add_argument("--system-prompt", default="") + parser.add_argument("--payload-stdin", action="store_true") parser.add_argument("--prompt", default="") + parser.add_argument("--prompt-file", default="") parser.add_argument("--initial-messages-file", default="") parser.add_argument("--mcp-config", default="") parser.add_argument("--edit", action="store_true") @@ -312,13 +319,26 @@ def parse_args() -> argparse.Namespace: async def main() -> None: args = parse_args() initial = [] + prompt = args.prompt + system_prompt = args.system_prompt + mcp_config = args.mcp_config + if args.payload_stdin: + payload = json.load(sys.stdin) + system_prompt = payload.get("system_prompt") or "" + prompt = payload.get("prompt") or "" + initial = payload.get("initial_messages") or [] + mcp_config = json.dumps(payload.get("mcp_config") or {}) + if args.prompt_file: + path = Path(args.prompt_file) + prompt = path.read_text(encoding="utf-8") + path.unlink() if args.initial_messages_file: path = Path(args.initial_messages_file) payload = path.read_bytes() path.unlink() initial = json.loads(payload) client = AsyncOpenAI(base_url=args.base_url, api_key=args.api_key) - config = json.loads(args.mcp_config or "{}") + config = json.loads(mcp_config or "{}") tools = [BASH_TOOL] reserved = {"bash"} if args.edit: @@ -333,15 +353,11 @@ async def main() -> None: else ([], {}, {}) ) tools += mcp_tools - messages = ( - [{"role": "system", "content": args.system_prompt}] - if args.system_prompt - else [] - ) + messages = [{"role": "system", "content": system_prompt}] if system_prompt else [] if initial: messages.extend(initial) - elif args.prompt: - messages.append({"role": "user", "content": args.prompt}) + elif prompt: + messages.append({"role": "user", "content": prompt}) while True: message = await chat(client, args.model, messages, tools) messages.append(message.model_dump(exclude_none=True)) diff --git a/verifiers/v1/harnesses/claude_code/harness.py b/verifiers/v1/harnesses/claude_code/harness.py index aab2a2af5..418e12413 100644 --- a/verifiers/v1/harnesses/claude_code/harness.py +++ b/verifiers/v1/harnesses/claude_code/harness.py @@ -79,7 +79,9 @@ async def launch( ctx.model, ] if system_prompt: - argv += ["--append-system-prompt", system_prompt] + system_prompt_path = f".vf-claude-system-{trace.id}.txt" + await runtime.write(system_prompt_path, system_prompt.encode("utf-8")) + argv += ["--append-system-prompt-file", system_prompt_path] argv += [ arg for tool in self.config.disabled_tools or [] @@ -97,6 +99,7 @@ async def launch( mcp_path, "--strict-mcp-config", "--", - instruction or "", ] - return await runtime.run_program(argv, env) + return await runtime.run_program( + argv, env, stdin=(instruction or "").encode("utf-8") + ) diff --git a/verifiers/v1/harnesses/codex/harness.py b/verifiers/v1/harnesses/codex/harness.py index 6d183d0a3..3665cd247 100644 --- a/verifiers/v1/harnesses/codex/harness.py +++ b/verifiers/v1/harnesses/codex/harness.py @@ -119,6 +119,8 @@ async def launch( image_args += ["-i", path] image_index += 1 prompt = "\n\n".join(texts) + if prompt is None: + raise ValueError("Codex requires a task prompt (it has no user simulator)") # codex authenticates to the interception server with the session secret (its provider # api key) and posts Responses calls to `{endpoint}/responses`. env = {**self.config.resolved_env, KEY_VAR: secret} @@ -163,10 +165,10 @@ async def launch( *tool_config, *image_args, "--", - prompt, + "-", ] try: - return await runtime.run_program(argv, env) + return await runtime.run_program(argv, env, stdin=prompt.encode("utf-8")) finally: if image_args: try: diff --git a/verifiers/v1/harnesses/null/harness.py b/verifiers/v1/harnesses/null/harness.py index 3af929208..b4e200d9d 100644 --- a/verifiers/v1/harnesses/null/harness.py +++ b/verifiers/v1/harnesses/null/harness.py @@ -38,33 +38,23 @@ async def launch( f"--base-url={endpoint}", f"--api-key={secret}", f"--model={ctx.model}", + "--payload-stdin", ] - if system_prompt: - args.append(f"--system-prompt={system_prompt}") - if mcp_urls: - # The program connects to the tool servers over HTTP; hand it a standard - # `mcpServers` URL config (the `mcp` client itself comes from the uv deps). - args.append( - "--mcp-config=" - + json.dumps( - { - "mcpServers": { - name: {"url": url} for name, url in mcp_urls.items() - } - } - ) - ) - if isinstance(prompt, str): - args.append(f"--prompt={prompt}") - elif prompt is not None: - # Base64 images can exceed exec limits, so hand Messages off through a file. - path = f".vf-initial-messages-{trace.id}.json" - await runtime.write( - path, - json.dumps([message_to_wire(m) for m in prompt]).encode(), - ) - args.append(f"--initial-messages-file={path}") + payload = { + "system_prompt": system_prompt, + "prompt": prompt if isinstance(prompt, str) else None, + "initial_messages": ( + [message_to_wire(message) for message in prompt] + if prompt is not None and not isinstance(prompt, str) + else [] + ), + "mcp_config": { + "mcpServers": {name: {"url": url} for name, url in mcp_urls.items()} + }, + } program = await runtime.prepare_uv_script( PROGRAM_SOURCE, self.config.resolved_env ) - return await runtime.run_program([*program, *args], env) + return await runtime.run_program( + [*program, *args], env, stdin=json.dumps(payload).encode("utf-8") + ) diff --git a/verifiers/v1/harnesses/null/program.py b/verifiers/v1/harnesses/null/program.py index 633ede6a4..87fdbaa53 100644 --- a/verifiers/v1/harnesses/null/program.py +++ b/verifiers/v1/harnesses/null/program.py @@ -7,6 +7,7 @@ import argparse import asyncio import json +import sys from contextlib import AsyncExitStack, asynccontextmanager, suppress from pathlib import Path @@ -144,6 +145,7 @@ def parse_args() -> argparse.Namespace: parser.add_argument("--api-key", required=True) parser.add_argument("--model", required=True) parser.add_argument("--system-prompt", default="") + parser.add_argument("--payload-stdin", action="store_true") parser.add_argument("--prompt", default="") parser.add_argument("--initial-messages-file", default="") parser.add_argument("--mcp-config", default="") @@ -153,28 +155,33 @@ def parse_args() -> argparse.Namespace: async def main() -> None: args = parse_args() initial = [] + prompt = args.prompt + system_prompt = args.system_prompt + mcp_config = args.mcp_config + if args.payload_stdin: + payload = json.load(sys.stdin) + system_prompt = payload.get("system_prompt") or "" + prompt = payload.get("prompt") or "" + initial = payload.get("initial_messages") or [] + mcp_config = json.dumps(payload.get("mcp_config") or {}) if args.initial_messages_file: path = Path(args.initial_messages_file) payload = path.read_bytes() path.unlink() initial = json.loads(payload) client = AsyncOpenAI(base_url=args.base_url, api_key=args.api_key) - config = json.loads(args.mcp_config or "{}") + config = json.loads(mcp_config or "{}") if config.get("mcpServers"): # Bound only tool enumeration; each session is opened and closed within this task. async with asyncio.timeout(60): tools, dispatch, servers = await connect_mcp(config) else: tools, dispatch, servers = [], {}, {} - messages = ( - [{"role": "system", "content": args.system_prompt}] - if args.system_prompt - else [] - ) + messages = [{"role": "system", "content": system_prompt}] if system_prompt else [] if initial: messages.extend(initial) - elif args.prompt: - messages.append({"role": "user", "content": args.prompt}) + elif prompt: + messages.append({"role": "user", "content": prompt}) while True: message = await chat(client, args.model, messages, tools) messages.append(message.model_dump(exclude_none=True)) diff --git a/verifiers/v1/harnesses/pi/harness.py b/verifiers/v1/harnesses/pi/harness.py index 0dff3f160..c30813463 100644 --- a/verifiers/v1/harnesses/pi/harness.py +++ b/verifiers/v1/harnesses/pi/harness.py @@ -234,7 +234,12 @@ async def launch( if self.config.disabled_tools else [] ) - system_args = ["--append-system-prompt", system_prompt] if system_prompt else [] + system_args: list[str] = [] + if system_prompt: + system_prompt_path = f"{agent_dir}/system-prompt.txt" + await runtime.write(system_prompt_path, system_prompt.encode("utf-8")) + # Pi interprets an existing path as the system-prompt file contents. + system_args = ["--append-system-prompt", system_prompt_path] argv = [ "sh", "-c", diff --git a/verifiers/v1/harnesses/terminus_2/harness.py b/verifiers/v1/harnesses/terminus_2/harness.py index 7ac800635..46ebb6f59 100644 --- a/verifiers/v1/harnesses/terminus_2/harness.py +++ b/verifiers/v1/harnesses/terminus_2/harness.py @@ -1,3 +1,4 @@ +import json import logging from pathlib import Path @@ -58,13 +59,15 @@ async def launch( f"--base-url={endpoint}", f"--api-key={secret}", f"--model={ctx.model}", - f"--system-prompt={system_prompt or ''}", - f"--task={prompt}", + "--payload-stdin", ] + payload = json.dumps( + {"system_prompt": system_prompt or "", "task": prompt} + ).encode("utf-8") try: source = PROGRAM_SOURCE.replace("{version}", self.config.version) program = await runtime.prepare_uv_script(source, self.config.resolved_env) - return await runtime.run_program([*program, *args], env) + return await runtime.run_program([*program, *args], env, stdin=payload) finally: # Harbor normally destroys its whole sandbox; this adapter borrows the # Verifiers runtime, so clean up Terminus's detached tmux server ourselves. diff --git a/verifiers/v1/harnesses/terminus_2/program.py b/verifiers/v1/harnesses/terminus_2/program.py index 9eee3cf7b..b3b6f001a 100644 --- a/verifiers/v1/harnesses/terminus_2/program.py +++ b/verifiers/v1/harnesses/terminus_2/program.py @@ -5,8 +5,10 @@ import argparse import asyncio +import json import os import subprocess +import sys from pathlib import Path, PurePosixPath from harbor.agents.terminus_2 import Terminus2 @@ -29,8 +31,8 @@ async def exec( ) -> ExecResult: _ = user result = subprocess.run( - command, - shell=True, + ["sh"], + input=command, cwd=cwd, env={**os.environ, **(env or {})}, capture_output=True, @@ -50,13 +52,18 @@ def parse_args() -> argparse.Namespace: parser.add_argument("--api-key", required=True) parser.add_argument("--model", required=True) parser.add_argument("--system-prompt", default="") - parser.add_argument("--task", required=True) + parser.add_argument("--task", default="") + parser.add_argument("--payload-stdin", action="store_true") return parser.parse_args() async def main() -> None: args = parse_args() model, system_prompt, task = args.model, args.system_prompt, args.task + if args.payload_stdin: + payload = json.load(sys.stdin) + system_prompt = payload.get("system_prompt") or "" + task = payload.get("task") or "" logs_dir = Path(os.environ["TMUX_TMPDIR"]) logs_dir.mkdir(mode=0o700, exist_ok=True) EnvironmentPaths.agent_dir = PurePosixPath(logs_dir) diff --git a/verifiers/v1/runtimes/base.py b/verifiers/v1/runtimes/base.py index b8a31fcdd..93324f685 100644 --- a/verifiers/v1/runtimes/base.py +++ b/verifiers/v1/runtimes/base.py @@ -145,13 +145,56 @@ def cleanup(self) -> None: async def run(self, argv: list[str], env: dict[str, str]) -> ProgramResult: pass - async def run_program(self, argv: list[str], env: dict[str, str]) -> ProgramResult: - """Run the harness's MAIN program — the rollout itself (a possibly long-lived, stateful, - agentic run) — as opposed to the short idempotent infra ops (write / mv / install / - provisioning) that go through `run`. No framework layer may replay this argv: doing so - against the rollout's persistent trace would fork a duplicate branch. Provider SDKs may - still retry individual safe transport operations underneath `run`.""" - return await self.run(argv, env) + async def run_program( + self, + argv: list[str], + env: dict[str, str], + *, + stdin: bytes | None = None, + ) -> ProgramResult: + """Run the harness's MAIN program once, optionally feeding binary stdin. + + Runtimes with a native stdin channel override ``_run_program_stdin``; the portable + fallback stages the bytes in the runtime workspace and redirects them through a bounded + shell wrapper. Either path keeps unbounded payloads out of argv and the environment. No + framework layer may replay the program: doing so against the rollout's persistent trace + would fork a duplicate branch. + """ + self._validate_exec_payload(argv, env) + if stdin is None: + return await self.run(argv, env) + return await self._run_program_stdin(argv, env, stdin) + + async def _run_program_stdin( + self, argv: list[str], env: dict[str, str], stdin: bytes + ) -> ProgramResult: + """File-backed stdin fallback for runtimes whose execution API has no stdin channel.""" + path = f".vf-stdin-{uuid.uuid4().hex}.bin" + await self.write(path, stdin) + command = ( + 'chmod 600 "$0"; status=0; "$@" < "$0" || status=$?; ' + 'rm -f -- "$0"; exit "$status"' + ) + return await self.run(["sh", "-c", command, path, *argv], env) + + @staticmethod + def _validate_exec_payload(argv: list[str], env: dict[str, str]) -> None: + """Fail clearly before exec when one argv/env entry exceeds Linux MAX_ARG_STRLEN.""" + limit = 128 * 1024 + entries = [(f"argv[{i}]", value) for i, value in enumerate(argv)] + entries += [(f"env[{key!r}]", f"{key}={value}") for key, value in env.items()] + for name, value in entries: + if "\0" in value: + raise ValueError( + f"{name} contains a NUL byte and cannot be passed to exec" + ) + size = len(value.encode("utf-8")) + 1 + if size > limit: + raise ValueError( + f"{name} is {size:,} bytes, exceeding the portable 128 KiB exec " + "argument limit; pass unbounded content through run_program(stdin=...) " + "or a runtime file" + ) async def run_background( self, argv: list[str], env: dict[str, str], log: str diff --git a/verifiers/v1/runtimes/docker.py b/verifiers/v1/runtimes/docker.py index f99919db8..6f14e656e 100644 --- a/verifiers/v1/runtimes/docker.py +++ b/verifiers/v1/runtimes/docker.py @@ -135,6 +135,32 @@ async def run(self, argv: list[str], env: dict[str, str]) -> ProgramResult: "exec", *env_args, "--workdir", self.config.workdir, self._container, *argv ) + async def _run_program_stdin( + self, argv: list[str], env: dict[str, str], stdin: bytes + ) -> ProgramResult: + env_args = [ + arg for key, value in env.items() for arg in ("--env", f"{key}={value}") + ] + proc = await asyncio.create_subprocess_exec( + "docker", + "exec", + "-i", + *env_args, + "--workdir", + self.config.workdir, + self._container, + *argv, + stdin=asyncio.subprocess.PIPE, + stdout=asyncio.subprocess.PIPE, + stderr=asyncio.subprocess.PIPE, + ) + stdout, stderr = await proc.communicate(input=stdin) + return ProgramResult( + exit_code=proc.returncode or 0, + stdout=stdout.decode(errors="replace"), + stderr=stderr.decode(errors="replace"), + ) + async def run_background( self, argv: list[str], env: dict[str, str], log: str ) -> None: diff --git a/verifiers/v1/runtimes/subprocess.py b/verifiers/v1/runtimes/subprocess.py index daf648362..fcf5bf2fb 100644 --- a/verifiers/v1/runtimes/subprocess.py +++ b/verifiers/v1/runtimes/subprocess.py @@ -73,6 +73,32 @@ async def run(self, argv: list[str], env: dict[str, str]) -> ProgramResult: stderr=stderr.decode(errors="replace"), ) + async def _run_program_stdin( + self, argv: list[str], env: dict[str, str], stdin: bytes + ) -> ProgramResult: + full_env = {k: v for k, v in os.environ.items() if "API_KEY" not in k.upper()} + full_env.update(env) + proc = await asyncio.create_subprocess_exec( + *argv, + env=full_env, + cwd=self.workdir, + stdin=asyncio.subprocess.PIPE, + stdout=asyncio.subprocess.PIPE, + stderr=asyncio.subprocess.PIPE, + start_new_session=True, + ) + try: + stdout, stderr = await proc.communicate(input=stdin) + finally: + if proc.returncode is None: + with contextlib.suppress(ProcessLookupError, PermissionError): + os.killpg(os.getpgid(proc.pid), signal.SIGKILL) + return ProgramResult( + exit_code=proc.returncode or 0, + stdout=stdout.decode(errors="replace"), + stderr=stderr.decode(errors="replace"), + ) + async def run_background( self, argv: list[str], env: dict[str, str], log: str ) -> None: