From 64005a14f7955884c6132e3d39be00d7121bd33a Mon Sep 17 00:00:00 2001 From: hallerite Date: Thu, 23 Jul 2026 18:48:21 +0200 Subject: [PATCH 1/7] feat(v1): run RLM over persistent ACP --- .../references/REFERENCE.md | 7 +- tests/v1/fixtures/pyproject.toml | 1 + tests/v1/test_e2e.py | 3 + verifiers/v1/acp/__init__.py | 79 ++++- verifiers/v1/acp/_runner.py | 312 +++++++++++++++--- verifiers/v1/harnesses/rlm/harness.py | 66 ++-- 6 files changed, 394 insertions(+), 74 deletions(-) diff --git a/skills/evaluate-environments/references/REFERENCE.md b/skills/evaluate-environments/references/REFERENCE.md index 7e25bb1f49..5ff58fa076 100644 --- a/skills/evaluate-environments/references/REFERENCE.md +++ b/skills/evaluate-environments/references/REFERENCE.md @@ -286,11 +286,14 @@ Installs the Codex CLI into the runtime and runs `codex exec`. #### `RLMHarnessConfig` — `id: "rlm"` -Installs the rlm CLI and runs it. Knobs map onto `RLM_*` env vars; base `HarnessConfig.env` passes any other `RLM_*` var through verbatim. +Installs rlm-harness and runs its ACP agent. MCP tools become pre-imported +IPython skills while the model-facing tool surface remains `ipython`. Knobs map +onto `RLM_*` env vars; base `HarnessConfig.env` passes any other `RLM_*` var +through verbatim. | Field | Type | Default | Notes | | --- | --- | --- | --- | -| `version` | `str` | `"main"` | Git ref (branch/tag/commit) of rlm to install. | +| `version` | `str` | `"56218f33796ecbe465445bc43948886354fde196"` | Git ref (branch/tag/commit) of rlm-harness to install. | | `max_depth` | `int` | `0` | Recursion depth rlm may spawn sub-harnesses to (`RLM_MAX_DEPTH`). | | `skills` | `list["edit" \| "search"]` | `[]` | Built-in rlm skills to enable (`RLM_SKILLS`). Empty enables none. | | `summarize_at_tokens` | `int \| (int, int) \| None` | `None` | Auto-compaction threshold (`RLM_SUMMARIZE_AT_TOKENS`): compact once context grows past this many tokens. An int is fixed; a `(lo, hi)` pair draws a per-group threshold (seeded by task index). `None` disables. Ints must be positive. | diff --git a/tests/v1/fixtures/pyproject.toml b/tests/v1/fixtures/pyproject.toml index 8b26b2dfaf..908dba4e62 100644 --- a/tests/v1/fixtures/pyproject.toml +++ b/tests/v1/fixtures/pyproject.toml @@ -15,6 +15,7 @@ build-backend = "hatchling.build" [tool.hatch.build] include = [ "counter_tool_v1.py", + "echo_acp_resume_v1.py", "echo_tool_v1.py", "tool_response_image_v1.py", "pyproject.toml", diff --git a/tests/v1/test_e2e.py b/tests/v1/test_e2e.py index 25d3033714..c854361077 100644 --- a/tests/v1/test_e2e.py +++ b/tests/v1/test_e2e.py @@ -54,6 +54,7 @@ def _pair(a: str, b: str, id: str, *extra_marks): # retain MCP access after resuming. Cover every harness in the local container runtime, # plus one remote placement for the sandbox/tunnel boundary. ACP_RESUME_PLACEMENTS = [ + _pair("rlm", "docker", "rlm-acp-in-docker"), _pair("kimi-code", "docker", "kimi-code-acp-in-docker"), _pair("pi", "docker", "pi-acp-in-docker"), _pair("pool", "docker", "pool-acp-in-docker"), @@ -205,6 +206,8 @@ async def test_acp_resume_with_tool(run_v1, harness, harness_runtime, tmp_path): assert segments[1]["terminated"] is False assert "tool" in segments[1]["roles"] assert segments[1]["tool_outputs"] + if harness == "rlm": + assert "turns_since_last_compaction" in trace.metrics @pytest.mark.e2e diff --git a/verifiers/v1/acp/__init__.py b/verifiers/v1/acp/__init__.py index 0c17b7e21c..ad579a87c7 100644 --- a/verifiers/v1/acp/__init__.py +++ b/verifiers/v1/acp/__init__.py @@ -2,7 +2,7 @@ import json import secrets -from pathlib import Path +from pathlib import Path, PurePosixPath from verifiers.v1.dialects.chat import message_to_wire from verifiers.v1.harness import Harness @@ -33,6 +33,7 @@ async def run( mcp_urls: dict[str, str] | None = None, system_prompt: str | None = None, session_path: str | None = None, + sidecar_path: str | None = None, ) -> ProgramResult: if prompt is None: raise ValueError("ACP requires a prompt") @@ -51,6 +52,24 @@ async def run( program = await runtime.prepare_uv_script( ACP_SOURCE, {**env, "UV_FROZEN": "false"} ) + sidecar_log = None + if sidecar_path is not None: + sidecar_dir = self._sidecar_dir(sidecar_path) + sidecar_log = f"{sidecar_dir}/acp.log" + exists = await runtime.run(["test", "-S", sidecar_path], {}) + if exists.exit_code != 0: + created = await runtime.run( + ["mkdir", "-p", "-m", "700", sidecar_dir], {} + ) + if created.exit_code != 0: + raise RuntimeError( + f"ACP sidecar directory failed: {created.stderr.strip()}" + ) + await runtime.run_background( + [*program, "serve", sidecar_path], + env, + sidecar_log, + ) directory = f".vf-acp-{secrets.token_hex(8)}" created = await runtime.run(["mkdir", "-m", "700", directory], {}) if created.exit_code != 0: @@ -58,7 +77,63 @@ async def run( path = f"{directory}/config.json" try: await runtime.write(path, json.dumps(config).encode()) - result = await runtime.run_program([*program, path], env) + command = ( + [*program, "request", path, sidecar_path] + if sidecar_path is not None + else [*program, "once", path] + ) + result = await runtime.run_program(command, env) + if sidecar_log is not None and result.exit_code != 0: + log = await runtime.run(["tail", "-c", "4000", sidecar_log], {}) + if log.exit_code == 0 and log.stdout: + result = ProgramResult( + exit_code=result.exit_code, + stdout=result.stdout, + stderr=( + f"{result.stderr.rstrip()}\n\nACP sidecar log:\n" + f"{log.stdout.rstrip()}" + ).lstrip(), + ) return result finally: await run_shielded(runtime.run(["rm", "-rf", directory], {})) + + async def close( + self, runtime: Runtime, sidecar_path: str, *, remove: bool = True + ) -> None: + """Stop a persistent ACP session, optionally keeping its artifacts.""" + sidecar_dir = self._sidecar_dir(sidecar_path) + exists = await runtime.run(["test", "-S", sidecar_path], {}) + failure = "" + try: + if exists.exit_code == 0: + program = await runtime.prepare_uv_script( + ACP_SOURCE, {"UV_FROZEN": "false"} + ) + result = await runtime.run([*program, "shutdown", sidecar_path], {}) + if result.exit_code != 0: + log = await runtime.run( + ["tail", "-c", "4000", f"{sidecar_dir}/acp.log"], {} + ) + failure = ( + result.stderr.strip() + or result.stdout.strip() + or "ACP sidecar shutdown failed" + ) + if log.exit_code == 0 and log.stdout: + failure = ( + f"{failure}\n\nACP sidecar log:\n{log.stdout.rstrip()}" + ) + finally: + if remove: + await run_shielded(runtime.run(["rm", "-rf", sidecar_dir], {})) + if failure: + raise RuntimeError(failure) + + @staticmethod + def _sidecar_dir(sidecar_path: str) -> str: + path = PurePosixPath(sidecar_path) + parent = str(path.parent) + if path.is_absolute() or ".." in path.parts or parent in ("", ".", "/"): + raise ValueError("ACP sidecar must live in a private subdirectory") + return parent diff --git a/verifiers/v1/acp/_runner.py b/verifiers/v1/acp/_runner.py index 5b794e610f..ae472ccf9a 100644 --- a/verifiers/v1/acp/_runner.py +++ b/verifiers/v1/acp/_runner.py @@ -2,12 +2,14 @@ # requires-python = ">=3.10,<3.15" # dependencies = ["agent-client-protocol==0.11.0"] # /// -"""Run one harness segment through an ACP agent.""" +"""Run harness segments through an ACP agent.""" import asyncio import json import os import sys +import traceback +from contextlib import AsyncExitStack, suppress from pathlib import Path from typing import Any @@ -30,12 +32,18 @@ TextContentBlock, ) +MAX_PACKET_BYTES = 128 * 1024 * 1024 + class VerifiersClient(Client): def __init__(self) -> None: self.visible_reply = "" self.message_id: str | None = None + def reset(self) -> None: + self.visible_reply = "" + self.message_id = None + async def session_update(self, session_id: str, update: Any, **kwargs: Any) -> None: if not isinstance(update, AgentMessageChunk) or not isinstance( update.content, TextContentBlock @@ -105,7 +113,62 @@ def content_blocks(messages: list[dict], supports_images: bool) -> list: return blocks -async def run_client(config: dict) -> None: +def mcp_servers(config: dict) -> list[HttpMcpServer]: + return [ + HttpMcpServer(type="http", name=name, url=url, headers=[]) + for name, url in config["mcp_urls"].items() + ] + + +def segment_messages(config: dict, is_new: bool) -> list[dict]: + messages = config["messages"] + if not is_new: + last_assistant = max( + ( + index + for index, message in enumerate(messages) + if message.get("role") == "assistant" + ), + default=-1, + ) + messages = messages[last_assistant + 1 :] + if is_new and config["system_prompt"]: + messages = [ + {"role": "system", "content": config["system_prompt"]}, + *messages, + ] + return messages + + +async def prompt( + client: VerifiersClient, + connection: Any, + capabilities: Any, + session_id: str, + config: dict, + *, + is_new: bool, +) -> str: + prompt_capabilities = capabilities and capabilities.prompt_capabilities + supports_images = bool(prompt_capabilities and prompt_capabilities.image) + blocks = content_blocks(segment_messages(config, is_new), supports_images) + if not blocks: + raise ValueError("ACP prompt has no content") + client.reset() + try: + response = await connection.prompt(session_id=session_id, prompt=blocks) + except RequestError as error: + detail = error.data.get("details") if isinstance(error.data, dict) else None + raise RuntimeError(detail or str(error)) from error + if not client.visible_reply.strip(): + raise RuntimeError( + "ACP agent produced no visible reply " + f"(stop_reason={response.stop_reason!r})" + ) + return client.visible_reply + + +async def run_once(config: dict) -> str: client = VerifiersClient() command = config["command"] async with spawn_agent_process( @@ -120,72 +183,231 @@ async def run_client(config: dict) -> None: client_capabilities=ClientCapabilities(), ) capabilities = initialized.agent_capabilities - prompt_capabilities = capabilities and capabilities.prompt_capabilities - supports_images = bool(prompt_capabilities and prompt_capabilities.image) - mcp_servers = [ - HttpMcpServer(type="http", name=name, url=url, headers=[]) - for name, url in config["mcp_urls"].items() - ] session_path = Path(config["session_path"]) if config["session_path"] else None is_new = session_path is None or not session_path.exists() + servers = mcp_servers(config) if is_new: - session = await connection.new_session( - cwd=os.getcwd(), mcp_servers=mcp_servers - ) + session = await connection.new_session(cwd=os.getcwd(), mcp_servers=servers) session_id = session.session_id else: session_id = session_path.read_text().strip() session_capabilities = capabilities and capabilities.session_capabilities if session_capabilities and session_capabilities.resume is not None: await connection.resume_session( - cwd=os.getcwd(), session_id=session_id, mcp_servers=mcp_servers + cwd=os.getcwd(), session_id=session_id, mcp_servers=servers ) elif capabilities and capabilities.load_session: await connection.load_session( - cwd=os.getcwd(), session_id=session_id, mcp_servers=mcp_servers + cwd=os.getcwd(), session_id=session_id, mcp_servers=servers ) else: raise RuntimeError("ACP agent does not support resuming sessions") - messages = config["messages"] - if not is_new: - last_assistant = max( - ( - index - for index, message in enumerate(messages) - if message.get("role") == "assistant" - ), - default=-1, - ) - messages = messages[last_assistant + 1 :] - if is_new and config["system_prompt"]: - messages = [ - {"role": "system", "content": config["system_prompt"]}, - *messages, - ] - prompt = content_blocks(messages, supports_images) - if not prompt: - raise ValueError("ACP prompt has no content") - client.visible_reply = "" - client.message_id = None - try: - await connection.prompt(session_id=session_id, prompt=prompt) - except RequestError as error: - detail = error.data.get("details") if isinstance(error.data, dict) else None - raise RuntimeError(detail or str(error)) from error - if not client.visible_reply.strip(): - raise RuntimeError("ACP agent produced no visible reply") - sys.stdout.write(client.visible_reply) + reply = await prompt( + client, + connection, + capabilities, + session_id, + config, + is_new=is_new, + ) if session_path and is_new: session_path.parent.mkdir(parents=True, exist_ok=True) session_path.write_text(session_id) + return reply -async def main() -> None: - path = Path(sys.argv[1]) +class PersistentSession: + """One live ACP process, connection, and session shared by several segments.""" + + def __init__(self) -> None: + self.stack = AsyncExitStack() + self.client = VerifiersClient() + self.connection: Any = None + self.capabilities: Any = None + self.session_id: str | None = None + self.command: list[str] | None = None + self.server_urls: dict[str, str] | None = None + self.system_prompt: str | None = None + self.is_new = True + + async def start(self, config: dict) -> None: + command = config["command"] + try: + self.connection, _process = await self.stack.enter_async_context( + spawn_agent_process( + self.client, + command[0], + *command[1:], + env=os.environ.copy(), + transport_kwargs={"stderr": None}, + ) + ) + initialized = await self.connection.initialize( + protocol_version=PROTOCOL_VERSION, + client_capabilities=ClientCapabilities(), + ) + self.capabilities = initialized.agent_capabilities + session = await self.connection.new_session( + cwd=os.getcwd(), mcp_servers=mcp_servers(config) + ) + except BaseException: + await self.stack.aclose() + raise + self.session_id = session.session_id + self.command = command + self.server_urls = config["mcp_urls"] + self.system_prompt = config["system_prompt"] + + async def run(self, config: dict) -> str: + if self.connection is None: + await self.start(config) + elif ( + config["command"] != self.command + or config["mcp_urls"] != self.server_urls + or config["system_prompt"] != self.system_prompt + ): + raise RuntimeError("ACP sidecar session configuration changed") + assert self.session_id is not None + reply = await prompt( + self.client, + self.connection, + self.capabilities, + self.session_id, + config, + is_new=self.is_new, + ) + self.is_new = False + return reply + + async def close(self) -> None: + if self.connection is not None and self.session_id is not None: + session_capabilities = ( + self.capabilities and self.capabilities.session_capabilities + ) + if session_capabilities and session_capabilities.close is not None: + with suppress(Exception): + await self.connection.close_session(session_id=self.session_id) + await self.stack.aclose() + self.connection = None + + +async def read_packet(reader: asyncio.StreamReader) -> dict: + size = int.from_bytes(await reader.readexactly(8), "big") + if size > MAX_PACKET_BYTES: + raise ValueError(f"ACP sidecar packet is too large: {size} bytes") + return json.loads((await reader.readexactly(size)).decode()) + + +async def write_packet(writer: asyncio.StreamWriter, value: dict) -> None: + data = json.dumps(value, ensure_ascii=False).encode() + writer.write(len(data).to_bytes(8, "big")) + writer.write(data) + await writer.drain() + + +async def serve_sidecar(socket_path: str) -> None: + path = Path(socket_path) + path.unlink(missing_ok=True) + session = PersistentSession() + lock = asyncio.Lock() + shutdown = asyncio.Event() + + async def handle( + reader: asyncio.StreamReader, writer: asyncio.StreamWriter + ) -> None: + try: + request = await read_packet(reader) + async with lock: + operation = request.get("operation") + if operation == "prompt": + reply = await session.run(request["config"]) + response = {"ok": True, "reply": reply} + elif operation == "shutdown": + try: + await session.close() + response = {"ok": True} + finally: + shutdown.set() + else: + raise ValueError(f"unknown ACP sidecar operation: {operation!r}") + except Exception as error: + traceback.print_exc() + response = { + "ok": False, + "error": f"{type(error).__name__}: {error}", + } + try: + await write_packet(writer, response) + finally: + writer.close() + await writer.wait_closed() + + server = await asyncio.start_unix_server(handle, path=socket_path) + os.chmod(path, 0o600) + try: + async with server: + await shutdown.wait() + finally: + server.close() + await server.wait_closed() + await session.close() + path.unlink(missing_ok=True) + + +async def connect( + socket_path: str, + wait_seconds: float = 60, +) -> tuple[asyncio.StreamReader, asyncio.StreamWriter]: + loop = asyncio.get_running_loop() + deadline = loop.time() + wait_seconds + while True: + try: + return await asyncio.open_unix_connection(socket_path) + except (FileNotFoundError, ConnectionRefusedError): + if loop.time() >= deadline: + raise RuntimeError("timed out waiting for ACP sidecar") + await asyncio.sleep(0.1) + + +async def request_sidecar( + socket_path: str, request: dict, wait_seconds: float = 60 +) -> dict: + reader, writer = await connect(socket_path, wait_seconds) + try: + await write_packet(writer, request) + response = await read_packet(reader) + finally: + writer.close() + await writer.wait_closed() + if not response.get("ok"): + raise RuntimeError(response.get("error") or "ACP sidecar request failed") + return response + + +def read_config(path_value: str) -> dict: + path = Path(path_value) config = json.loads(path.read_text()) path.unlink() - await run_client(config) + return config + + +async def main() -> None: + operation = sys.argv[1] + if operation == "once": + sys.stdout.write(await run_once(read_config(sys.argv[2]))) + elif operation == "serve": + await serve_sidecar(sys.argv[2]) + elif operation == "request": + response = await request_sidecar( + sys.argv[3], + {"operation": "prompt", "config": read_config(sys.argv[2])}, + ) + sys.stdout.write(response["reply"]) + elif operation == "shutdown": + await request_sidecar(sys.argv[2], {"operation": "shutdown"}, wait_seconds=2) + else: + raise ValueError(f"unknown ACP runner operation: {operation!r}") if __name__ == "__main__": diff --git a/verifiers/v1/harnesses/rlm/harness.py b/verifiers/v1/harnesses/rlm/harness.py index 5569cb430a..68a57b12d3 100644 --- a/verifiers/v1/harnesses/rlm/harness.py +++ b/verifiers/v1/harnesses/rlm/harness.py @@ -1,4 +1,4 @@ -"""RLM exposes `RLM_MCP_CONFIG` tools as pre-imported IPython skills.""" +"""RLM over ACP, with MCP tools exposed as pre-imported IPython skills.""" import json import logging @@ -8,30 +8,30 @@ from pydantic import model_validator -from verifiers.v1.harness import Harness, HarnessConfig +from verifiers.v1.acp import ACP from verifiers.v1.clients import ModelContext from verifiers.v1.decorators import metric -from verifiers.v1.dialects.chat import message_to_wire +from verifiers.v1.harness import Harness, HarnessConfig from verifiers.v1.runtimes import ProgramResult, Runtime -from verifiers.v1.trace import Trace from verifiers.v1.task import TaskData +from verifiers.v1.trace import Trace logger = logging.getLogger(__name__) BuiltinSkill = Literal["edit", "search"] -RLM_REPO = "github.com/PrimeIntellect-ai/rlm.git" -# rlm writes its session under $RLM_HOME/sessions//; point it at a workdir- -# relative dir so it stays in the runtime (and is cleaned up with the workdir). -RLM_HOME = ".rlm" +RLM_REPO = "github.com/PrimeIntellect-ai/rlm-harness.git" +RLM_VERSION = "56218f33796ecbe465445bc43948886354fde196" RLM_DIR = "/tmp/vf-rlm" RLM_BIN = f"{RLM_DIR}/bin/rlm" SKILLS_DIR = "/task/rlm-skills" +RLM_STATE_DIR = ".vf-rlm" +RLM_ACP = ACP() class RLMHarnessConfig(HarnessConfig): - version: str = "main" - """Git ref (branch, tag, or commit) of rlm to install.""" + version: str = RLM_VERSION + """Git ref (branch, tag, or commit) of rlm-harness to install.""" max_depth: int = 0 """Recursion depth rlm may spawn sub-harnesses to (RLM_MAX_DEPTH).""" builtin_skills: list[BuiltinSkill] = [] @@ -92,10 +92,11 @@ async def setup(self, runtime: Runtime) -> None: logger.info("rlm: ensuring rlm is installed (version=%s)", self.config.version) ensure = shlex.quote(f"[ -x {RLM_BIN} ] || ({install})") guarded = f"mkdir -p {RLM_DIR} && flock {RLM_DIR}/install.lock sh -c {ensure}" - env = {**self.config.resolved_env, "RLM_HOME": RLM_HOME} + env = self.config.resolved_env result = await runtime.run(["sh", "-c", guarded], env) if result.exit_code != 0: raise RuntimeError(f"rlm install failed: {result.stderr.strip()[-500:]}") + await RLM_ACP.setup(self, runtime) def summarize_threshold(self, task_idx: int | None) -> str: """The `RLM_SUMMARIZE_AT_TOKENS` value: a range draws per-group (seeded by task index — @@ -122,34 +123,34 @@ async def launch( system_prompt, prompt = self.resolve_prompt(data) if prompt is None: raise ValueError("RLM requires a prompt") - if not isinstance(prompt, str): - prompt = json.dumps( - [message_to_wire(message) for message in prompt], ensure_ascii=False - ) env = { **self.config.resolved_env, "RLM_BASE_URL": endpoint, "RLM_API_KEY": secret, "RLM_MODEL": ctx.model, "RLM_MAX_DEPTH": str(self.config.max_depth), - "RLM_HOME": RLM_HOME, + "RLM_HOME": self._home(trace), "RLM_SUMMARIZE_AT_TOKENS": self.summarize_threshold(data.idx), } if system_prompt is not None: env["RLM_APPEND_TO_SYSTEM_PROMPT"] = system_prompt if self.config.builtin_skills: env["RLM_SKILLS"] = ",".join(self.config.builtin_skills) - if mcp_urls: - env["RLM_MCP_CONFIG"] = json.dumps( - {"mcpServers": {name: {"url": url} for name, url in mcp_urls.items()}} - ) - # RLM has no interactive mode; resumed segments explicitly replay the transcript. - return await runtime.run_program([RLM_BIN, "--", prompt], env) + return await RLM_ACP.run( + runtime, + env, + [RLM_BIN, "--acp"], + prompt, + mcp_urls=mcp_urls, + sidecar_path=self._sidecar(trace), + ) @metric - async def rlm(self, runtime: Runtime) -> dict[str, float]: - # Stateless continuation creates one session per segment; report the latest. - latest = f'cat "$(ls -t {RLM_HOME}/sessions/*/meta.json | head -1)"' + async def rlm(self, trace: Trace, runtime: Runtime) -> dict[str, float]: + # Closing finalizes RLM's meta.json; keep the state until cleanup. + await RLM_ACP.close(runtime, self._sidecar(trace), remove=False) + home = shlex.quote(self._home(trace)) + latest = f'cat "$(ls -t {home}/sessions/*/meta.json | head -1)"' result = await runtime.run(["sh", "-c", latest], {}) if result.exit_code != 0 or not result.stdout.strip(): return {} @@ -162,3 +163,18 @@ async def rlm(self, runtime: Runtime) -> dict[str, float]: for key, value in meta.get("metrics", {}).items() if isinstance(value, (int, float)) and not isinstance(value, bool) } + + async def cleanup(self, trace: Trace, runtime: Runtime) -> None: + await RLM_ACP.close(runtime, self._sidecar(trace)) + + @staticmethod + def _state_dir(trace: Trace) -> str: + return f"{RLM_STATE_DIR}/{trace.id}" + + @classmethod + def _home(cls, trace: Trace) -> str: + return f"{cls._state_dir(trace)}/home" + + @classmethod + def _sidecar(cls, trace: Trace) -> str: + return f"{cls._state_dir(trace)}/acp.sock" From e5a51da04f63905699e56553593ce6867394fe65 Mon Sep 17 00:00:00 2001 From: hallerite Date: Thu, 23 Jul 2026 19:21:31 +0200 Subject: [PATCH 2/7] fix(v1): recover persistent ACP sidecars --- verifiers/v1/acp/__init__.py | 9 ++++-- verifiers/v1/acp/_runner.py | 53 ++++++++++++++++++++++++++---------- 2 files changed, 46 insertions(+), 16 deletions(-) diff --git a/verifiers/v1/acp/__init__.py b/verifiers/v1/acp/__init__.py index ad579a87c7..31c1911a04 100644 --- a/verifiers/v1/acp/__init__.py +++ b/verifiers/v1/acp/__init__.py @@ -56,8 +56,13 @@ async def run( if sidecar_path is not None: sidecar_dir = self._sidecar_dir(sidecar_path) sidecar_log = f"{sidecar_dir}/acp.log" - exists = await runtime.run(["test", "-S", sidecar_path], {}) - if exists.exit_code != 0: + probe = await runtime.run([*program, "probe", sidecar_path], {}) + if probe.exit_code != 0: + removed = await runtime.run(["rm", "-f", sidecar_path], {}) + if removed.exit_code != 0: + raise RuntimeError( + f"stale ACP sidecar cleanup failed: {removed.stderr.strip()}" + ) created = await runtime.run( ["mkdir", "-p", "-m", "700", sidecar_dir], {} ) diff --git a/verifiers/v1/acp/_runner.py b/verifiers/v1/acp/_runner.py index ae472ccf9a..a4646695b0 100644 --- a/verifiers/v1/acp/_runner.py +++ b/verifiers/v1/acp/_runner.py @@ -252,7 +252,18 @@ async def start(self, config: dict) -> None: cwd=os.getcwd(), mcp_servers=mcp_servers(config) ) except BaseException: - await self.stack.aclose() + try: + await self.stack.aclose() + except BaseException: + pass + self.stack = AsyncExitStack() + self.connection = None + self.capabilities = None + self.session_id = None + self.command = None + self.server_urls = None + self.system_prompt = None + self.is_new = True raise self.session_id = session.session_id self.command = command @@ -318,19 +329,24 @@ async def handle( ) -> None: try: request = await read_packet(reader) - async with lock: - operation = request.get("operation") - if operation == "prompt": - reply = await session.run(request["config"]) - response = {"ok": True, "reply": reply} - elif operation == "shutdown": - try: - await session.close() - response = {"ok": True} - finally: - shutdown.set() - else: - raise ValueError(f"unknown ACP sidecar operation: {operation!r}") + operation = request.get("operation") + if operation == "ping": + response = {"ok": True} + else: + async with lock: + if operation == "prompt": + reply = await session.run(request["config"]) + response = {"ok": True, "reply": reply} + elif operation == "shutdown": + try: + await session.close() + response = {"ok": True} + finally: + shutdown.set() + else: + raise ValueError( + f"unknown ACP sidecar operation: {operation!r}" + ) except Exception as error: traceback.print_exc() response = { @@ -406,6 +422,15 @@ async def main() -> None: sys.stdout.write(response["reply"]) elif operation == "shutdown": await request_sidecar(sys.argv[2], {"operation": "shutdown"}, wait_seconds=2) + elif operation == "probe": + await asyncio.wait_for( + request_sidecar( + sys.argv[2], + {"operation": "ping"}, + wait_seconds=0, + ), + timeout=2, + ) else: raise ValueError(f"unknown ACP runner operation: {operation!r}") From 1aaad1131f81ea071b299fb58c373f4f73416f17 Mon Sep 17 00:00:00 2001 From: hallerite Date: Thu, 23 Jul 2026 19:31:37 +0200 Subject: [PATCH 3/7] fix(v1): serialize ACP sidecar startup --- verifiers/v1/acp/__init__.py | 66 ++++++++++++++++++++++++++---------- verifiers/v1/acp/_runner.py | 5 +-- 2 files changed, 52 insertions(+), 19 deletions(-) diff --git a/verifiers/v1/acp/__init__.py b/verifiers/v1/acp/__init__.py index 31c1911a04..1629f966cd 100644 --- a/verifiers/v1/acp/__init__.py +++ b/verifiers/v1/acp/__init__.py @@ -1,8 +1,10 @@ """Public Agent Client Protocol support for harness programs.""" +import asyncio import json import secrets from pathlib import Path, PurePosixPath +from weakref import WeakValueDictionary from verifiers.v1.dialects.chat import message_to_wire from verifiers.v1.harness import Harness @@ -18,6 +20,19 @@ class ACP: """Run an ACP agent.""" + def __init__(self) -> None: + self._sidecar_locks: WeakValueDictionary[tuple[int, str], asyncio.Lock] = ( + WeakValueDictionary() + ) + + def _sidecar_lock(self, runtime: Runtime, sidecar_path: str) -> asyncio.Lock: + key = (id(runtime), sidecar_path) + lock = self._sidecar_locks.get(key) + if lock is None: + lock = asyncio.Lock() + self._sidecar_locks[key] = lock + return lock + async def setup(self, harness: Harness, runtime: Runtime) -> None: await runtime.prepare_uv_script( ACP_SOURCE, {**harness.config.resolved_env, "UV_FROZEN": "false"} @@ -56,25 +71,42 @@ async def run( if sidecar_path is not None: sidecar_dir = self._sidecar_dir(sidecar_path) sidecar_log = f"{sidecar_dir}/acp.log" - probe = await runtime.run([*program, "probe", sidecar_path], {}) - if probe.exit_code != 0: - removed = await runtime.run(["rm", "-f", sidecar_path], {}) - if removed.exit_code != 0: - raise RuntimeError( - f"stale ACP sidecar cleanup failed: {removed.stderr.strip()}" + async with self._sidecar_lock(runtime, sidecar_path): + probe = await runtime.run([*program, "probe", sidecar_path], {}) + if probe.exit_code != 0: + removed = await runtime.run(["rm", "-f", sidecar_path], {}) + if removed.exit_code != 0: + raise RuntimeError( + "stale ACP sidecar cleanup failed: " + f"{removed.stderr.strip()}" + ) + created = await runtime.run( + ["mkdir", "-p", "-m", "700", sidecar_dir], {} ) - created = await runtime.run( - ["mkdir", "-p", "-m", "700", sidecar_dir], {} - ) - if created.exit_code != 0: - raise RuntimeError( - f"ACP sidecar directory failed: {created.stderr.strip()}" + if created.exit_code != 0: + raise RuntimeError( + f"ACP sidecar directory failed: {created.stderr.strip()}" + ) + await runtime.run_background( + [*program, "serve", sidecar_path], + env, + sidecar_log, ) - await runtime.run_background( - [*program, "serve", sidecar_path], - env, - sidecar_log, - ) + ready = await runtime.run( + [*program, "probe", sidecar_path, "60"], {} + ) + if ready.exit_code != 0: + log = await runtime.run(["tail", "-c", "4000", sidecar_log], {}) + detail = ( + ready.stderr.strip() + or ready.stdout.strip() + or "sidecar did not become ready" + ) + if log.exit_code == 0 and log.stdout: + detail = ( + f"{detail}\n\nACP sidecar log:\n{log.stdout.rstrip()}" + ) + raise RuntimeError(f"ACP sidecar failed to start: {detail}") directory = f".vf-acp-{secrets.token_hex(8)}" created = await runtime.run(["mkdir", "-m", "700", directory], {}) if created.exit_code != 0: diff --git a/verifiers/v1/acp/_runner.py b/verifiers/v1/acp/_runner.py index a4646695b0..8052957382 100644 --- a/verifiers/v1/acp/_runner.py +++ b/verifiers/v1/acp/_runner.py @@ -423,13 +423,14 @@ async def main() -> None: elif operation == "shutdown": await request_sidecar(sys.argv[2], {"operation": "shutdown"}, wait_seconds=2) elif operation == "probe": + wait_seconds = float(sys.argv[3]) if len(sys.argv) > 3 else 0 await asyncio.wait_for( request_sidecar( sys.argv[2], {"operation": "ping"}, - wait_seconds=0, + wait_seconds=wait_seconds, ), - timeout=2, + timeout=max(2, wait_seconds + 1), ) else: raise ValueError(f"unknown ACP runner operation: {operation!r}") From b26cc43ccdf125d86908830aa6ee59f135ddcec5 Mon Sep 17 00:00:00 2001 From: hallerite Date: Thu, 23 Jul 2026 19:37:50 +0200 Subject: [PATCH 4/7] fix(v1): bound ACP sidecar lifecycle --- verifiers/v1/acp/__init__.py | 10 ++++++- verifiers/v1/acp/_runner.py | 52 ++++++++++++++++++++++++------------ 2 files changed, 44 insertions(+), 18 deletions(-) diff --git a/verifiers/v1/acp/__init__.py b/verifiers/v1/acp/__init__.py index 1629f966cd..ce08cc2b71 100644 --- a/verifiers/v1/acp/__init__.py +++ b/verifiers/v1/acp/__init__.py @@ -13,6 +13,7 @@ from verifiers.v1.utils.aio import run_shielded ACP_SOURCE = (Path(__file__).resolve().parent / "_runner.py").read_text() +PROBE_UNAVAILABLE_EXIT_CODE = 75 __all__ = ["ACP"] @@ -73,7 +74,7 @@ async def run( sidecar_log = f"{sidecar_dir}/acp.log" async with self._sidecar_lock(runtime, sidecar_path): probe = await runtime.run([*program, "probe", sidecar_path], {}) - if probe.exit_code != 0: + if probe.exit_code == PROBE_UNAVAILABLE_EXIT_CODE: removed = await runtime.run(["rm", "-f", sidecar_path], {}) if removed.exit_code != 0: raise RuntimeError( @@ -107,6 +108,13 @@ async def run( f"{detail}\n\nACP sidecar log:\n{log.stdout.rstrip()}" ) raise RuntimeError(f"ACP sidecar failed to start: {detail}") + elif probe.exit_code != 0: + detail = ( + probe.stderr.strip() + or probe.stdout.strip() + or "sidecar did not respond" + ) + raise RuntimeError(f"ACP sidecar probe failed: {detail}") directory = f".vf-acp-{secrets.token_hex(8)}" created = await runtime.run(["mkdir", "-m", "700", directory], {}) if created.exit_code != 0: diff --git a/verifiers/v1/acp/_runner.py b/verifiers/v1/acp/_runner.py index 8052957382..adbe499019 100644 --- a/verifiers/v1/acp/_runner.py +++ b/verifiers/v1/acp/_runner.py @@ -33,6 +33,7 @@ ) MAX_PACKET_BYTES = 128 * 1024 * 1024 +PROBE_UNAVAILABLE_EXIT_CODE = 75 class VerifiersClient(Client): @@ -332,17 +333,17 @@ async def handle( operation = request.get("operation") if operation == "ping": response = {"ok": True} + elif operation == "shutdown": + try: + await session.close() + response = {"ok": True} + finally: + shutdown.set() else: async with lock: if operation == "prompt": reply = await session.run(request["config"]) response = {"ok": True, "reply": reply} - elif operation == "shutdown": - try: - await session.close() - response = {"ok": True} - finally: - shutdown.set() else: raise ValueError( f"unknown ACP sidecar operation: {operation!r}" @@ -387,12 +388,19 @@ async def connect( async def request_sidecar( - socket_path: str, request: dict, wait_seconds: float = 60 + socket_path: str, + request: dict, + wait_seconds: float = 60, + response_seconds: float | None = None, ) -> dict: reader, writer = await connect(socket_path, wait_seconds) try: await write_packet(writer, request) - response = await read_packet(reader) + response = ( + await read_packet(reader) + if response_seconds is None + else await asyncio.wait_for(read_packet(reader), response_seconds) + ) finally: writer.close() await writer.wait_closed() @@ -421,17 +429,27 @@ async def main() -> None: ) sys.stdout.write(response["reply"]) elif operation == "shutdown": - await request_sidecar(sys.argv[2], {"operation": "shutdown"}, wait_seconds=2) + await request_sidecar( + sys.argv[2], + {"operation": "shutdown"}, + wait_seconds=2, + response_seconds=5, + ) elif operation == "probe": wait_seconds = float(sys.argv[3]) if len(sys.argv) > 3 else 0 - await asyncio.wait_for( - request_sidecar( - sys.argv[2], - {"operation": "ping"}, - wait_seconds=wait_seconds, - ), - timeout=max(2, wait_seconds + 1), - ) + try: + await asyncio.wait_for( + request_sidecar( + sys.argv[2], + {"operation": "ping"}, + wait_seconds=wait_seconds, + ), + timeout=max(2, wait_seconds + 1), + ) + except RuntimeError as error: + if str(error) == "timed out waiting for ACP sidecar": + raise SystemExit(PROBE_UNAVAILABLE_EXIT_CODE) from None + raise else: raise ValueError(f"unknown ACP runner operation: {operation!r}") From 54864d07d962221131e8f35d292b812229bf142c Mon Sep 17 00:00:00 2001 From: hallerite Date: Thu, 23 Jul 2026 19:42:22 +0200 Subject: [PATCH 5/7] fix(v1): keep RLM metrics best-effort --- verifiers/v1/harnesses/rlm/harness.py | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/verifiers/v1/harnesses/rlm/harness.py b/verifiers/v1/harnesses/rlm/harness.py index 68a57b12d3..c29f2a828f 100644 --- a/verifiers/v1/harnesses/rlm/harness.py +++ b/verifiers/v1/harnesses/rlm/harness.py @@ -148,7 +148,10 @@ async def launch( @metric async def rlm(self, trace: Trace, runtime: Runtime) -> dict[str, float]: # Closing finalizes RLM's meta.json; keep the state until cleanup. - await RLM_ACP.close(runtime, self._sidecar(trace), remove=False) + try: + await RLM_ACP.close(runtime, self._sidecar(trace), remove=False) + except Exception: + logger.warning("rlm: sidecar shutdown failed before metrics", exc_info=True) home = shlex.quote(self._home(trace)) latest = f'cat "$(ls -t {home}/sessions/*/meta.json | head -1)"' result = await runtime.run(["sh", "-c", latest], {}) From d323b4602711b66ebb24917ed8646734b4265eb8 Mon Sep 17 00:00:00 2001 From: hallerite Date: Thu, 23 Jul 2026 19:49:45 +0200 Subject: [PATCH 6/7] fix(v1): harden ACP sidecar lifecycle --- verifiers/v1/acp/_runner.py | 111 +++++++++++++++++++++++++----------- 1 file changed, 77 insertions(+), 34 deletions(-) diff --git a/verifiers/v1/acp/_runner.py b/verifiers/v1/acp/_runner.py index adbe499019..0591e5e316 100644 --- a/verifiers/v1/acp/_runner.py +++ b/verifiers/v1/acp/_runner.py @@ -222,8 +222,11 @@ class PersistentSession: """One live ACP process, connection, and session shared by several segments.""" def __init__(self) -> None: - self.stack = AsyncExitStack() self.client = VerifiersClient() + self._reset() + + def _reset(self) -> None: + self.stack = AsyncExitStack() self.connection: Any = None self.capabilities: Any = None self.session_id: str | None = None @@ -257,19 +260,13 @@ async def start(self, config: dict) -> None: await self.stack.aclose() except BaseException: pass - self.stack = AsyncExitStack() - self.connection = None - self.capabilities = None - self.session_id = None - self.command = None - self.server_urls = None - self.system_prompt = None - self.is_new = True + self._reset() raise self.session_id = session.session_id self.command = command self.server_urls = config["mcp_urls"] self.system_prompt = config["system_prompt"] + self.is_new = True async def run(self, config: dict) -> str: if self.connection is None: @@ -293,15 +290,17 @@ async def run(self, config: dict) -> str: return reply async def close(self) -> None: - if self.connection is not None and self.session_id is not None: - session_capabilities = ( - self.capabilities and self.capabilities.session_capabilities - ) - if session_capabilities and session_capabilities.close is not None: - with suppress(Exception): - await self.connection.close_session(session_id=self.session_id) - await self.stack.aclose() - self.connection = None + try: + if self.connection is not None and self.session_id is not None: + session_capabilities = ( + self.capabilities and self.capabilities.session_capabilities + ) + if session_capabilities and session_capabilities.close is not None: + with suppress(Exception): + await self.connection.close_session(session_id=self.session_id) + await self.stack.aclose() + finally: + self._reset() async def read_packet(reader: asyncio.StreamReader) -> dict: @@ -323,31 +322,69 @@ async def serve_sidecar(socket_path: str) -> None: path.unlink(missing_ok=True) session = PersistentSession() lock = asyncio.Lock() + stop_lock = asyncio.Lock() shutdown = asyncio.Event() + active_prompt: asyncio.Task[str] | None = None + + async def run_prompt(config: dict) -> str: + nonlocal active_prompt + async with lock: + if shutdown.is_set(): + raise RuntimeError("ACP sidecar is shutting down") + active_prompt = asyncio.create_task(session.run(config)) + try: + return await active_prompt + finally: + active_prompt = None + + async def stop_session() -> None: + shutdown.set() + async with stop_lock: + task = active_prompt + if task is not None and not task.done(): + task.cancel() + await asyncio.gather(task, return_exceptions=True) + # `run_prompt` releases this after its cancelled session request unwinds. + # Holding it for close prevents a waiting prompt from racing a restart. + async with lock: + await session.close() async def handle( reader: asyncio.StreamReader, writer: asyncio.StreamWriter ) -> None: + response: dict | None = None try: request = await read_packet(reader) operation = request.get("operation") if operation == "ping": response = {"ok": True} elif operation == "shutdown": - try: - await session.close() - response = {"ok": True} - finally: - shutdown.set() + await stop_session() + response = {"ok": True} + elif operation == "prompt": + prompt_task = asyncio.create_task(run_prompt(request["config"])) + disconnect_task = asyncio.create_task(reader.read()) + done, _ = await asyncio.wait( + (prompt_task, disconnect_task), + return_when=asyncio.FIRST_COMPLETED, + ) + if prompt_task in done: + disconnect_task.cancel() + await asyncio.gather(disconnect_task, return_exceptions=True) + response = {"ok": True, "reply": await prompt_task} + else: + # The short-lived request process was cancelled or timed out. + # Stop both the prompt and its agent process so neither can keep + # consuming tokens without a client. + prompt_task.cancel() + await asyncio.gather(prompt_task, return_exceptions=True) + await stop_session() + return else: - async with lock: - if operation == "prompt": - reply = await session.run(request["config"]) - response = {"ok": True, "reply": reply} - else: - raise ValueError( - f"unknown ACP sidecar operation: {operation!r}" - ) + raise ValueError(f"unknown ACP sidecar operation: {operation!r}") + except asyncio.CancelledError: + if not shutdown.is_set(): + raise except Exception as error: traceback.print_exc() response = { @@ -355,10 +392,14 @@ async def handle( "error": f"{type(error).__name__}: {error}", } try: - await write_packet(writer, response) + if response is not None: + await write_packet(writer, response) + except (BrokenPipeError, ConnectionResetError): + pass finally: writer.close() - await writer.wait_closed() + with suppress(BrokenPipeError, ConnectionResetError): + await writer.wait_closed() server = await asyncio.start_unix_server(handle, path=socket_path) os.chmod(path, 0o600) @@ -368,7 +409,7 @@ async def handle( finally: server.close() await server.wait_closed() - await session.close() + await stop_session() path.unlink(missing_ok=True) @@ -450,6 +491,8 @@ async def main() -> None: if str(error) == "timed out waiting for ACP sidecar": raise SystemExit(PROBE_UNAVAILABLE_EXIT_CODE) from None raise + except TimeoutError: + raise SystemExit(PROBE_UNAVAILABLE_EXIT_CODE) from None else: raise ValueError(f"unknown ACP runner operation: {operation!r}") From b447e180db3a0dfb3a379e6cf4391dbf2e070f4e Mon Sep 17 00:00:00 2001 From: hallerite Date: Thu, 23 Jul 2026 20:07:43 +0200 Subject: [PATCH 7/7] fix(v1): preserve ACP sidecar ownership --- verifiers/v1/acp/__init__.py | 15 +++++++++------ verifiers/v1/acp/_runner.py | 2 -- 2 files changed, 9 insertions(+), 8 deletions(-) diff --git a/verifiers/v1/acp/__init__.py b/verifiers/v1/acp/__init__.py index ce08cc2b71..a46386c03a 100644 --- a/verifiers/v1/acp/__init__.py +++ b/verifiers/v1/acp/__init__.py @@ -4,7 +4,7 @@ import json import secrets from pathlib import Path, PurePosixPath -from weakref import WeakValueDictionary +from weakref import WeakKeyDictionary from verifiers.v1.dialects.chat import message_to_wire from verifiers.v1.harness import Harness @@ -22,16 +22,19 @@ class ACP: """Run an ACP agent.""" def __init__(self) -> None: - self._sidecar_locks: WeakValueDictionary[tuple[int, str], asyncio.Lock] = ( - WeakValueDictionary() + self._sidecar_locks: WeakKeyDictionary[Runtime, dict[str, asyncio.Lock]] = ( + WeakKeyDictionary() ) def _sidecar_lock(self, runtime: Runtime, sidecar_path: str) -> asyncio.Lock: - key = (id(runtime), sidecar_path) - lock = self._sidecar_locks.get(key) + locks = self._sidecar_locks.get(runtime) + if locks is None: + locks = {} + self._sidecar_locks[runtime] = locks + lock = locks.get(sidecar_path) if lock is None: lock = asyncio.Lock() - self._sidecar_locks[key] = lock + locks[sidecar_path] = lock return lock async def setup(self, harness: Harness, runtime: Runtime) -> None: diff --git a/verifiers/v1/acp/_runner.py b/verifiers/v1/acp/_runner.py index 0591e5e316..428328dd18 100644 --- a/verifiers/v1/acp/_runner.py +++ b/verifiers/v1/acp/_runner.py @@ -491,8 +491,6 @@ async def main() -> None: if str(error) == "timed out waiting for ACP sidecar": raise SystemExit(PROBE_UNAVAILABLE_EXIT_CODE) from None raise - except TimeoutError: - raise SystemExit(PROBE_UNAVAILABLE_EXIT_CODE) from None else: raise ValueError(f"unknown ACP runner operation: {operation!r}")