Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
16 changes: 12 additions & 4 deletions environments/compact/compact/harness.py
Original file line number Diff line number Diff line change
Expand Up @@ -11,9 +11,10 @@
import json
from pathlib import Path

from verifiers.v1.harness import Harness, HarnessConfig
from verifiers.v1.clients import ModelContext
from verifiers.v1.harness import Harness, HarnessConfig
from verifiers.v1.runtimes import ProgramResult, Runtime
from verifiers.v1.task import TaskData
from verifiers.v1.trace import Trace

PROGRAM_SOURCE = (Path(__file__).resolve().parent / "program.py").read_text()
Expand All @@ -28,7 +29,9 @@ class CompactingHarness(Harness[CompactingHarnessConfig]):
SUPPORTS_MCP = True

async def setup(self, runtime: Runtime) -> None:
await runtime.prepare_uv_script(PROGRAM_SOURCE, self.config.env)
await runtime.prepare_uv_script(
PROGRAM_SOURCE, {**self.config.resolved_env, "UV_FROZEN": "false"}
)

async def launch(
self,
Expand All @@ -38,8 +41,11 @@ async def launch(
endpoint: str,
secret: str,
mcp_urls: dict[str, str],
data: TaskData,
) -> ProgramResult:
_, prompt = self.resolve_prompt(data)
env = {
**self.config.resolved_env,
"OPENAI_BASE_URL": endpoint,
"OPENAI_API_KEY": secret,
"OPENAI_MODEL": ctx.model,
Expand All @@ -50,5 +56,7 @@ async def launch(
env["MCP_CONFIG"] = json.dumps(
{"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)
program = await runtime.prepare_uv_script(
PROGRAM_SOURCE, {**self.config.resolved_env, "UV_FROZEN": "false"}
)
return await runtime.run_program([*program, prompt], env)
2 changes: 2 additions & 0 deletions verifiers/v1/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@

from pydantic_config import BaseConfig

from verifiers.v1.acp import ACP
from verifiers.v1.clients import (
BaseClientConfig,
Client,
Expand Down Expand Up @@ -238,6 +239,7 @@
"BaseConfig",
"Harness",
"HarnessConfig",
"ACP",
"ModelContext",
"Runtime",
"RuntimeConfig",
Expand Down
65 changes: 65 additions & 0 deletions verifiers/v1/acp/__init__.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,65 @@
"""Public Agent Client Protocol support for harness programs."""

import json
import secrets
from dataclasses import replace
from pathlib import Path

from verifiers.v1.dialects.chat import message_to_wire
from verifiers.v1.harness import Harness
from verifiers.v1.runtimes import ProgramResult, Runtime
from verifiers.v1.types import Messages
from verifiers.v1.utils.aio import run_shielded

ACP_SOURCE = (Path(__file__).resolve().parent / "_runner.py").read_text()

__all__ = ["ACP"]


class ACP:
"""Run an ACP agent."""

async def setup(self, harness: Harness, runtime: Runtime) -> None:
await runtime.prepare_uv_script(
ACP_SOURCE, {**harness.config.resolved_env, "UV_FROZEN": "false"}
)

async def run(
self,
runtime: Runtime,
env: dict[str, str],
command: list[str],
prompt: str | Messages | None,
*,
mcp_urls: dict[str, str] | None = None,
system_prompt: str | None = None,
session_path: str | None = None,
) -> ProgramResult:
if prompt is None:
raise ValueError("ACP requires a prompt")
messages = (
[{"role": "user", "content": prompt}]
if isinstance(prompt, str)
else [message_to_wire(message) for message in prompt]
)
config = {
"command": command,
"messages": messages,
"mcp_urls": mcp_urls or {},
"system_prompt": system_prompt or "",
"session_path": session_path,
}
program = await runtime.prepare_uv_script(
ACP_SOURCE, {**env, "UV_FROZEN": "false"}
)
directory = f".vf-acp-{secrets.token_hex(8)}"
created = await runtime.run(["mkdir", "-m", "700", directory], {})
if created.exit_code != 0:
raise RuntimeError(f"ACP config directory failed: {created.stderr.strip()}")
path = f"{directory}/config.json"
try:
await runtime.write(path, json.dumps(config).encode())
result = await runtime.run_program([*program, path], env)
return replace(result, visible_reply=result.stdout.strip())
finally:
await run_shielded(runtime.run(["rm", "-rf", directory], {}))
213 changes: 213 additions & 0 deletions verifiers/v1/acp/_runner.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,213 @@
# /// script
# requires-python = ">=3.10,<3.15"
# dependencies = ["agent-client-protocol==0.11.0"]
# ///
"""Run one harness segment through an ACP agent."""

import asyncio
import json
import os
import sys
from pathlib import Path
from typing import Any

from acp import (
PROTOCOL_VERSION,
Client,
RequestError,
image_block,
spawn_agent_process,
text_block,
)
from acp.schema import (
AgentMessageChunk,
AllowedOutcome,
ClientCapabilities,
DeniedOutcome,
FileSystemCapabilities,
HttpMcpServer,
PermissionOption,
ReadTextFileResponse,
RequestPermissionResponse,
TextContentBlock,
WriteTextFileResponse,
)


class VerifiersClient(Client):
def __init__(self) -> None:
self.visible_reply = ""
self.message_id: str | None = 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
):
return
if update.message_id is not None and update.message_id != self.message_id:
self.visible_reply = ""
self.message_id = update.message_id
self.visible_reply += update.content.text

async def read_text_file(
self,
session_id: str,
path: str,
line: int | None = None,
limit: int | None = None,
**kwargs: Any,
) -> ReadTextFileResponse:
lines = Path(path).read_text().splitlines(keepends=True)
start = (line or 1) - 1
return ReadTextFileResponse(
content="".join(lines[start : start + limit if limit is not None else None])
)

async def write_text_file(
self,
session_id: str,
path: str,
content: str,
**kwargs: Any,
) -> WriteTextFileResponse:
Path(path).write_text(content)
return WriteTextFileResponse()

async def request_permission(
self,
session_id: str,
tool_call: Any,
options: list[PermissionOption],
**kwargs: Any,
) -> RequestPermissionResponse:
option = next(
(item for item in options if item.kind in ("allow_once", "allow_always")),
None,
)
outcome = (
AllowedOutcome(outcome="selected", option_id=option.option_id)
if option
else DeniedOutcome(outcome="cancelled")
)
return RequestPermissionResponse(outcome=outcome)


def content_blocks(messages: list[dict], supports_images: bool) -> list:
blocks = []
transcript = len(messages) != 1 or messages[0].get("role") != "user"
for message in messages:
if blocks:
blocks.append(text_block("\n\n"))
if transcript:
blocks.append(text_block(f"[{message.get('role', 'message')}]\n"))
content = message.get("content") or ""
parts = (
[{"type": "text", "text": content}] if isinstance(content, str) else content
)
for part in parts:
if part["type"] == "text":
blocks.append(text_block(part["text"]))
continue
if not supports_images:
raise ValueError("ACP agent does not support image prompts")
url = part["image_url"]["url"]
metadata, separator, data = url.partition(",")
media_type, *parameters = metadata.removeprefix("data:").split(";")
if (
not separator
or not metadata.startswith("data:image/")
or not any(value.lower() == "base64" for value in parameters)
):
raise ValueError("ACP image prompts require base64 data:image URLs")
blocks.append(image_block(data, media_type))
metadata = {
key: value
for key, value in message.items()
if key not in ("role", "content") and value
}
if metadata:
blocks.append(text_block("\n" + json.dumps(metadata, ensure_ascii=False)))
return blocks


async def run_client(config: dict) -> None:
client = VerifiersClient()
command = config["command"]
async with spawn_agent_process(
client,
command[0],
*command[1:],
env=os.environ.copy(),
transport_kwargs={"stderr": None},
) as (connection, _process):
initialized = await connection.initialize(
protocol_version=PROTOCOL_VERSION,
client_capabilities=ClientCapabilities(
fs=FileSystemCapabilities(read_text_file=True, write_text_file=True)
),
)
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()
if is_new:
session = await connection.new_session(
cwd=os.getcwd(), mcp_servers=mcp_servers
)
session_id = session.session_id
else:
if not capabilities or not capabilities.load_session:
raise RuntimeError("ACP agent does not support loading sessions")
session_id = session_path.read_text().strip()
await connection.load_session(
cwd=os.getcwd(), session_id=session_id, mcp_servers=mcp_servers
)

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)
if session_path and is_new:
session_path.parent.mkdir(parents=True, exist_ok=True)
session_path.write_text(session_id)


async def main() -> None:
path = Path(sys.argv[1])
config = json.loads(path.read_text())
path.unlink()
await run_client(config)


if __name__ == "__main__":
asyncio.run(main())
13 changes: 9 additions & 4 deletions verifiers/v1/agent.py
Original file line number Diff line number Diff line change
Expand Up @@ -193,12 +193,17 @@ async def _turn(self, message: str | Messages | None) -> Reply:
elif message is not None:
messages = _as_messages(message)
self._started = True
turns_before = self.trace.num_turns
await self._run.step(messages)
if self.trace.num_turns > turns_before:
result = await self._run.step(messages)
if result is not None:
# The segment answered — even if a limit or @stop then ended the
# exchange, that surfaces as the NEXT turn's stopped reply.
return Reply(text=self.trace.last_reply)
return Reply(
text=(
result.visible_reply
if result.visible_reply is not None
else self.trace.last_reply
)
)
self._over = True
return Reply(text="", stopped=True)

Expand Down
7 changes: 3 additions & 4 deletions verifiers/v1/harness.py
Original file line number Diff line number Diff line change
Expand Up @@ -111,7 +111,7 @@ async def run(
mcp_urls: dict[str, str],
data: TaskData,
messages: Messages | None = None,
) -> None:
) -> ProgramResult:
"""Run ONE segment of the exchange: the program from launch (or, with
`messages`, the user's next turn(s) via `resume`) until it yields — a segment
ends when the program exits. The rollout loop owns the exchange across
Expand All @@ -125,14 +125,13 @@ async def run(
result = await self.resume(
ctx, trace, runtime, endpoint, secret, mcp_urls, data, messages
)
if trace.stop_condition is not None:
return # a @stop refused a turn mid-rollout; the harness's exit is expected
if result.exit_code != 0:
if trace.stop_condition is None and result.exit_code != 0:
# The real cause is at the END of a traceback, so keep the tail.
detail = (result.stderr or result.stdout).strip()[-2000:] or "<no output>"
raise HarnessError(
f"harness {self.config.id!r} exited {result.exit_code}: {detail}"
)
return result

async def score(self, trace: Trace, runtime: Runtime) -> None:
"""Run this harness's `@metric` methods over the finished trace, recording
Expand Down
Loading
Loading