From cf0ad84547e27af00ce0990e1cf504d26d67efea Mon Sep 17 00:00:00 2001 From: Mika Senghaas Date: Wed, 26 Aug 2026 19:22:59 +0000 Subject: [PATCH 01/14] add bash context compaction Keep long Bash harness runs alive when the provider rejects an overlong\nprompt. The harness requests a nano-rlm-style checkpoint summary and\ncontinues from a fresh branch.\n\nMake compaction the default while preserving clean context-length\ntruncation when users disable it. --- verifiers/v1/harnesses/bash/harness.py | 6 ++ verifiers/v1/harnesses/bash/program.py | 126 +++++++++++++++++++++++-- verifiers/v1/interception/server.py | 21 ++++- 3 files changed, 142 insertions(+), 11 deletions(-) diff --git a/verifiers/v1/harnesses/bash/harness.py b/verifiers/v1/harnesses/bash/harness.py index 92273025cc..10f48888d5 100644 --- a/verifiers/v1/harnesses/bash/harness.py +++ b/verifiers/v1/harnesses/bash/harness.py @@ -28,6 +28,10 @@ class BashHarnessConfig(HarnessConfig): + compaction: bool = True + """Compact the conversation at the model context limit. The harness asks the model for a + handoff summary, then continues from a fresh context that contains the summary.""" + edit: bool = True """Offer the local `edit` tool (single-occurrence string replacement in a file) alongside `bash`. On by default; set `--env.agent.harness.edit false` for a bash-only agent.""" @@ -77,6 +81,8 @@ async def launch( ] if tool_interception_url: args.append(f"--tool-interception-url={tool_interception_url}") + if self.config.compaction: + args.append("--compaction") if self.config.edit: args.append("--edit") if self.config.search: diff --git a/verifiers/v1/harnesses/bash/program.py b/verifiers/v1/harnesses/bash/program.py index 7a31c6cad5..0e6e16c45a 100644 --- a/verifiers/v1/harnesses/bash/program.py +++ b/verifiers/v1/harnesses/bash/program.py @@ -12,13 +12,40 @@ from pathlib import Path import httpx -from openai import AsyncOpenAI +from openai import AsyncOpenAI, BadRequestError from tenacity import AsyncRetrying, stop_after_attempt, wait_exponential_jitter SERPER_URL = "https://google.serper.dev/search" MCP_CALL_ATTEMPTS = 6 MCP_TIMEOUT = 600.0 +CONTEXT_COMPACTION_HEADER = "X-Verifiers-Context-Compaction" + +CHECKPOINT_COMPACTION_PROMPT = ( + "You are performing a CONTEXT CHECKPOINT COMPACTION. " + "Create a handoff summary for another LLM that will resume the task.\n" + "\n" + "Include:\n" + "- Current progress and key decisions made\n" + "- Important context, constraints, or user preferences\n" + "- What remains to be done (clear next steps)\n" + "- Any critical data, examples, or references needed to continue\n" + "\n" + "Be concise, structured, and focused on helping the next LLM " + "seamlessly continue the work." +) + +POST_COMPACTION_FRAMING = ( + "Another language model started to solve this problem and produced " + "a summary of its thinking process. Use this to build on the work " + "that has already been done and avoid duplicating work. Here is " + "the summary produced by the other language model, use the " + "information in this summary to assist with your own analysis:" +) + +COMPACTED_TOOL_RESULT = ( + "[tool output dropped because it exceeded the model context limit]" +) BASH_TOOL = { @@ -178,14 +205,84 @@ def run_edit(path: str, old_str: str, new_str: str) -> str: async def chat( - client: AsyncOpenAI, model: str, messages: list[dict], tools: list[dict] + client: AsyncOpenAI, + model: str, + messages: list[dict], + tools: list[dict], + *, + allow_compaction: bool = False, + tool_choice: str | None = None, ): - completion = await client.chat.completions.create( - model=model, messages=messages, tools=tools or None - ) + kwargs = {"model": model, "messages": messages, "tools": tools or None} + if allow_compaction: + kwargs["extra_headers"] = {CONTEXT_COMPACTION_HEADER: "1"} + if tools and tool_choice is not None: + kwargs["tool_choice"] = tool_choice + completion = await client.chat.completions.create(**kwargs) return completion.choices[0].message +def is_context_length_error(error: BadRequestError) -> bool: + details = f"{error} {error.body or ''}".casefold() + return "context_length" in details or "context length" in details + + +def drop_latest_tool_result(messages: list[dict]) -> bool: + """Replace one tool result so the checkpoint request can fit in context.""" + for index in range(len(messages) - 1, -1, -1): + message = messages[index] + if message.get("role") != "tool": + continue + if message.get("content") == COMPACTED_TOOL_RESULT: + continue + messages[index] = {**message, "content": COMPACTED_TOOL_RESULT} + return True + return False + + +def can_compact(messages: list[dict]) -> bool: + return any( + message.get("role") == "tool" + and message.get("content") != COMPACTED_TOOL_RESULT + for message in messages + ) + + +async def compact( + client: AsyncOpenAI, + model: str, + messages: list[dict], + tools: list[dict], +) -> list[dict]: + """Create a handoff summary after removing only the tool output needed to fit it.""" + system_messages = [ + message for message in messages if message.get("role") == "system" + ] + while drop_latest_tool_result(messages): + checkpoint_messages = [ + *messages, + {"role": "user", "content": CHECKPOINT_COMPACTION_PROMPT}, + ] + try: + summary = await chat( + client, + model, + checkpoint_messages, + tools, + allow_compaction=can_compact(messages), + tool_choice="none", + ) + except BadRequestError as error: + if not is_context_length_error(error): + raise + if not can_compact(messages): + raise + continue + framed = POST_COMPACTION_FRAMING + "\n\n" + (summary.content or "") + return [*system_messages, {"role": "user", "content": framed}] + raise RuntimeError("context compaction could not make the checkpoint prompt fit") + + @asynccontextmanager async def mcp_session(spec: dict): """One fresh streamable-HTTP session to an MCP server, opened and closed within the caller's @@ -328,6 +425,7 @@ def parse_args() -> argparse.Namespace: parser.add_argument("--initial-messages-file", default="") parser.add_argument("--mcp-config", default="") parser.add_argument("--tool-interception-url", default="") + parser.add_argument("--compaction", action="store_true") parser.add_argument("--edit", action="store_true") parser.add_argument("--search", action="store_true") parser.add_argument("--serper-key", default="") @@ -373,7 +471,23 @@ async def main() -> None: elif args.prompt: messages.append({"role": "user", "content": args.prompt}) while True: - message = await chat(client, args.model, messages, tools) + try: + message = await chat( + client, + args.model, + messages, + tools, + allow_compaction=args.compaction and can_compact(messages), + ) + except BadRequestError as error: + if ( + not args.compaction + or not can_compact(messages) + or not is_context_length_error(error) + ): + raise + messages = await compact(client, args.model, messages, tools) + continue messages.append(message.model_dump(exclude_none=True)) if not message.tool_calls: break diff --git a/verifiers/v1/interception/server.py b/verifiers/v1/interception/server.py index 9b2ba7be7b..09c1dd42b9 100644 --- a/verifiers/v1/interception/server.py +++ b/verifiers/v1/interception/server.py @@ -81,6 +81,8 @@ # Attempt counter the stainless-generated SDKs (OpenAI, Anthropic) send on every request: # 0 on the first attempt, incremented on each retry of the same request. RETRY_COUNT_HEADER = "x-stainless-retry-count" +CONTEXT_COMPACTION_HEADER = "X-Verifiers-Context-Compaction" +"""Internal opt-in for a harness that can recover from an overlong prompt.""" IDEMPOTENCY_KEY_HEADER = "Idempotency-Key" IDEMPOTENCY_CACHE_TTL_SECONDS = 600 IDEMPOTENCY_CACHE_MAX_COMPLETED = 64 @@ -493,6 +495,12 @@ async def handle_request( body = dialect.apply_overrides(body, session.ctx.model, session.ctx.sampling) streaming = dialect.streaming(body) upstream_headers = dict(request.headers) + context_compaction = request.headers.get(CONTEXT_COMPACTION_HEADER) == "1" + upstream_headers = { + name: value + for name, value in upstream_headers.items() + if name.lower() != CONTEXT_COMPACTION_HEADER.lower() + } logger.debug( "intercept %s: id=%s stream=%s", request.path, @@ -715,14 +723,17 @@ async def sample() -> web.Response: status=400, ) except OverlongPromptError as e: - # An overlong prompt is a budget limit, not a crash: end the rollout - # cleanly as a truncation — refuse the call to halt the harness (same - # shape as `refused` above). error = e - session.trace.stop("context_length") + if not context_compaction: + session.trace.stop("context_length") logger.debug("prompt too long: id=%s", session.trace.id) + message = ( + "context_length" + if context_compaction + else "rollout stopped: context_length" + ) return web.json_response( - dialect.error_body("rollout stopped: context_length"), + dialect.error_body(message), status=400, ) except RolloutError as e: From b60b651d22287028440d593c59612eb34caffc97 Mon Sep 17 00:00:00 2001 From: Mika Senghaas Date: Wed, 26 Aug 2026 20:19:25 +0000 Subject: [PATCH 02/14] recover from context overflow --- pyproject.toml | 1 + tests/v1/conftest.py | 9 +- tests/v1/fixtures/context_compaction_v1.py | 86 +++++++++++++++++ tests/v1/test_e2e.py | 64 +++++++++++++ verifiers/v1/harnesses/bash/harness.py | 11 ++- verifiers/v1/harnesses/bash/program.py | 102 ++++++++++++++------- verifiers/v1/harnesses/rlm/harness.py | 4 +- verifiers/v1/interception/server.py | 23 +---- 8 files changed, 236 insertions(+), 64 deletions(-) create mode 100644 tests/v1/fixtures/context_compaction_v1.py diff --git a/pyproject.toml b/pyproject.toml index 056ebad809..6781bba70d 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -290,6 +290,7 @@ markers = [ "bash: v1 e2e cases on the bash harness", "browser_use: v1 e2e cases on the browser_use harness", "rlm: v1 e2e cases on the rlm harness", + "compaction: v1 context-compaction E2E cases", "kimi_code: v1 e2e cases on the kimi-code harness", "pi: v1 e2e cases on the pi harness", "pool: v1 e2e cases on the pool harness", diff --git a/tests/v1/conftest.py b/tests/v1/conftest.py index 0b5ba36de9..63e1549624 100644 --- a/tests/v1/conftest.py +++ b/tests/v1/conftest.py @@ -96,7 +96,7 @@ def pytest_configure(config) -> None: def pytest_collection_modifyitems(config, items) -> None: """Skip the live-model tests (marked `e2e`) when no model endpoint is configured, so the rest of the suite (e.g. config parsing) still runs in a keyless environment.""" - if os.environ.get("PRIME_API_KEY"): + if os.environ.get("PRIME_API_KEY") or os.environ.get("VF_COMPACTION_E2E_BASE_URL"): return skip = pytest.mark.skip(reason="needs PRIME_API_KEY") for item in items: @@ -124,7 +124,7 @@ def _eval_config( harness: str | HarnessConfig | None = "null", n: int = 1, num_tasks: int = 1, - max_tokens: int = 2048, + max_tokens: int | None = 2048, max_turns: int | None = 4, rollout_timeout: float = 180, taskset_overrides: dict | None = None, @@ -132,6 +132,8 @@ def _eval_config( env: dict | None = None, pool: dict | None = None, reasoning_effort: str | None = None, + model: str = CI_MODEL, + client: dict | None = None, server: bool = False, ) -> EvalConfig: """Build the smallest `EvalConfig` that still exercises the path, shared by the in-process @@ -187,7 +189,8 @@ def _eval_config( serve=({"pool": pool} if pool else {}) if server else None, output_dir=output_dir.parent, run={"dir": output_dir.name}, - model=CI_MODEL, + model=model, + client=client or {"type": "eval"}, ) diff --git a/tests/v1/fixtures/context_compaction_v1.py b/tests/v1/fixtures/context_compaction_v1.py new file mode 100644 index 0000000000..79ba6bfb62 --- /dev/null +++ b/tests/v1/fixtures/context_compaction_v1.py @@ -0,0 +1,86 @@ +"""Context-compaction E2E scenarios for harness agent loops.""" + +from typing import Literal + +from pydantic import Field + +import verifiers.v1 as vf + + +class OverflowToolsetConfig(vf.ToolsetConfig): + payload_chars: int = Field(65_536, gt=0) + + +class OverflowToolset(vf.Toolset[OverflowToolsetConfig]): + TOOL_PREFIX = "overflow" + + def __init__(self, config: OverflowToolsetConfig): + super().__init__(config) + self.called = False + + @vf.tool + def overflow_context(self) -> str: + """Return a payload that is intentionally larger than the model context.""" + if self.called: + return "The overflow already occurred. Answer `recovered` now." + self.called = True + block = "0123456789abcdef " + repeats = self.config.payload_chars // len(block) + 1 + return (block * repeats)[: self.config.payload_chars] + + +class ContextCompactionTaskConfig(vf.TaskConfig): + scenario: Literal["decode", "tool_result"] = "decode" + payload_chars: int = Field(65_536, gt=0) + tools: OverflowToolsetConfig = OverflowToolsetConfig() + + +class ContextCompactionTask( + vf.Task[vf.TaskData, vf.State, ContextCompactionTaskConfig] +): + @classmethod + def toolsets(cls, config: ContextCompactionTaskConfig) -> list[vf.Toolset]: + if config.scenario != "tool_result": + return [] + tool_config = config.tools.model_copy( + update={"payload_chars": config.payload_chars} + ) + return [OverflowToolset(tool_config)] + + @vf.reward + async def compacted(self, trace: vf.Trace) -> float: + return float(trace.num_branches > 1) + + +class ContextCompactionConfig(vf.TasksetConfig): + task: ContextCompactionTaskConfig = ContextCompactionTaskConfig() + + +class ContextCompactionTaskset( + vf.Taskset[ContextCompactionTask, ContextCompactionConfig] +): + def load(self) -> list[ContextCompactionTask]: + if self.config.task.scenario == "decode": + prompt = ( + "Write `x ` repeatedly. Do not use tools and do not stop. " + "Continue until the model context ends the decode." + ) + else: + prompt = ( + "Call the `overflow_context` tool exactly once, then answer `recovered`. " + "In an RLM IPython session, call it with " + "`result = await overflow_overflow_context(); print(result)`." + ) + return [ + ContextCompactionTask( + vf.TaskData(idx=0, prompt=prompt), + self.config.task, + ) + ] + + +__all__ = ["ContextCompactionTaskset"] + + +if __name__ == "__main__": + OverflowToolset(OverflowToolsetConfig()).run() diff --git a/tests/v1/test_e2e.py b/tests/v1/test_e2e.py index cef9bf7e16..9e73157f67 100644 --- a/tests/v1/test_e2e.py +++ b/tests/v1/test_e2e.py @@ -5,8 +5,12 @@ with distinct networking — instead of fanning the full cross product. prime/modal rows are local-only (their marks are excluded in CI).""" +import os + import pytest +from verifiers.v1.utils.loaders import harness_config_type + mark = pytest.mark @@ -156,6 +160,66 @@ async def test_single_turn(run_v1, harness, harness_runtime, tmp_path): assert call.time.duration > 0 +@pytest.mark.e2e +@pytest.mark.compaction +@pytest.mark.docker +@pytest.mark.parametrize("scenario", ["decode", "tool_result"]) +@pytest.mark.parametrize("harness_id", ["bash", "rlm"]) +async def test_context_compaction_matrix(run_v1, scenario, harness_id, tmp_path): + """Both in-house loops recover when decoding or tool output fills context.""" + base_url = os.environ.get("VF_COMPACTION_E2E_BASE_URL") + model = os.environ.get("VF_COMPACTION_E2E_MODEL") + context_window = int(os.environ.get("VF_COMPACTION_E2E_CONTEXT_WINDOW", "4096")) + if not base_url or not model: + pytest.skip("needs a local compaction E2E model") + + harness = { + "id": harness_id, + "summarize_at_tokens": context_window, + **({"compaction": True} if harness_id == "bash" else {}), + } + sampling = ( + {"extra_body": {"ignore_eos": True}} + if scenario == "decode" + else {"temperature": 0.0} + ) + (trace,) = await run_v1( + "context-compaction-v1", + harness=harness_config_type(harness_id).model_validate(harness), + runtime={"type": "docker"}, + env={"agent": {"sampling": sampling, "max_output_tokens": None}}, + client={ + "type": "eval", + "base_url": base_url, + "api_key_var": "VF_COMPACTION_E2E_API_KEY", + }, + model=model, + max_tokens=None if scenario == "decode" else 512, + max_turns=8, + rollout_timeout=600, + taskset_overrides={ + "task": { + "scenario": scenario, + "payload_chars": context_window * 12, + } + }, + output_dir=tmp_path / f"{scenario}-{harness_id}", + ) + + assert trace.ok, trace.errors + assert trace.num_branches > 1 + assert trace.rewards["compacted"].score == 1.0 + if scenario == "decode": + assert any(call.finish_reason == "length" for call in trace.calls) + else: + assert any( + call.error is not None and call.error.type == "OverlongPromptError" + for call in trace.calls + ) + if harness_id == "rlm": + assert trace.metrics["num_compactions"] >= 1 + + @pytest.mark.e2e @pytest.mark.browser_use @pytest.mark.docker diff --git a/verifiers/v1/harnesses/bash/harness.py b/verifiers/v1/harnesses/bash/harness.py index 10f48888d5..f97e30b2b0 100644 --- a/verifiers/v1/harnesses/bash/harness.py +++ b/verifiers/v1/harnesses/bash/harness.py @@ -2,6 +2,8 @@ import os from pathlib import Path +from pydantic import PositiveInt + from verifiers.v1.clients import ModelContext from verifiers.v1.configs.harness import HarnessConfig from verifiers.v1.dialects.chat import message_to_wire @@ -29,8 +31,11 @@ class BashHarnessConfig(HarnessConfig): compaction: bool = True - """Compact the conversation at the model context limit. The harness asks the model for a - handoff summary, then continues from a fresh context that contains the summary.""" + """Recover from context exhaustion with a handoff summary and a fresh branch.""" + + summarize_at_tokens: PositiveInt | None = None + """Compact proactively when the estimated active context reaches this threshold. When unset, + compaction still recovers from provider context errors.""" edit: bool = True """Offer the local `edit` tool (single-occurrence string replacement in a file) alongside @@ -83,6 +88,8 @@ async def launch( args.append(f"--tool-interception-url={tool_interception_url}") if self.config.compaction: args.append("--compaction") + if self.config.summarize_at_tokens is not None: + args.append(f"--summarize-at-tokens={self.config.summarize_at_tokens}") if self.config.edit: args.append("--edit") if self.config.search: diff --git a/verifiers/v1/harnesses/bash/program.py b/verifiers/v1/harnesses/bash/program.py index 0e6e16c45a..1ed6e16afa 100644 --- a/verifiers/v1/harnesses/bash/program.py +++ b/verifiers/v1/harnesses/bash/program.py @@ -19,7 +19,6 @@ MCP_CALL_ATTEMPTS = 6 MCP_TIMEOUT = 600.0 -CONTEXT_COMPACTION_HEADER = "X-Verifiers-Context-Compaction" CHECKPOINT_COMPACTION_PROMPT = ( "You are performing a CONTEXT CHECKPOINT COMPACTION. " @@ -210,21 +209,28 @@ async def chat( messages: list[dict], tools: list[dict], *, - allow_compaction: bool = False, tool_choice: str | None = None, ): kwargs = {"model": model, "messages": messages, "tools": tools or None} - if allow_compaction: - kwargs["extra_headers"] = {CONTEXT_COMPACTION_HEADER: "1"} if tools and tool_choice is not None: kwargs["tool_choice"] = tool_choice - completion = await client.chat.completions.create(**kwargs) - return completion.choices[0].message + return await client.chat.completions.create(**kwargs) def is_context_length_error(error: BadRequestError) -> bool: details = f"{error} {error.body or ''}".casefold() - return "context_length" in details or "context length" in details + return any( + marker in details + for marker in ( + "request entity too large", + "context_length", + "context length", + "context window", + "prompt is too long", + "too many tokens", + "token limit exceeded", + ) + ) def drop_latest_tool_result(messages: list[dict]) -> bool: @@ -240,12 +246,15 @@ def drop_latest_tool_result(messages: list[dict]) -> bool: return False -def can_compact(messages: list[dict]) -> bool: - return any( - message.get("role") == "tool" - and message.get("content") != COMPACTED_TOOL_RESULT - for message in messages - ) +def estimated_tokens(value: str) -> int: + return (len(value) + 3) // 4 + + +def context_tokens(completion) -> int: + usage = completion.usage + if usage is None: + return 0 + return (usage.prompt_tokens or 0) + (usage.completion_tokens or 0) async def compact( @@ -254,33 +263,32 @@ async def compact( messages: list[dict], tools: list[dict], ) -> list[dict]: - """Create a handoff summary after removing only the tool output needed to fit it.""" + """Create a handoff summary, removing tool output only when the checkpoint overflows.""" system_messages = [ message for message in messages if message.get("role") == "system" ] - while drop_latest_tool_result(messages): + while True: checkpoint_messages = [ *messages, {"role": "user", "content": CHECKPOINT_COMPACTION_PROMPT}, ] try: - summary = await chat( + completion = await chat( client, model, checkpoint_messages, tools, - allow_compaction=can_compact(messages), tool_choice="none", ) except BadRequestError as error: if not is_context_length_error(error): raise - if not can_compact(messages): + if not drop_latest_tool_result(messages): raise continue + summary = completion.choices[0].message framed = POST_COMPACTION_FRAMING + "\n\n" + (summary.content or "") return [*system_messages, {"role": "user", "content": framed}] - raise RuntimeError("context compaction could not make the checkpoint prompt fit") @asynccontextmanager @@ -426,6 +434,7 @@ def parse_args() -> argparse.Namespace: parser.add_argument("--mcp-config", default="") parser.add_argument("--tool-interception-url", default="") parser.add_argument("--compaction", action="store_true") + parser.add_argument("--summarize-at-tokens", type=int) parser.add_argument("--edit", action="store_true") parser.add_argument("--search", action="store_true") parser.add_argument("--serper-key", default="") @@ -471,26 +480,37 @@ async def main() -> None: elif args.prompt: messages.append({"role": "user", "content": args.prompt}) while True: - try: - message = await chat( - client, - args.model, - messages, - tools, - allow_compaction=args.compaction and can_compact(messages), - ) - except BadRequestError as error: + recovered_overflow = False + while True: + try: + completion = await chat(client, args.model, messages, tools) + except BadRequestError as error: + if ( + not args.compaction + or recovered_overflow + or not is_context_length_error(error) + ): + raise + messages = await compact(client, args.model, messages, tools) + recovered_overflow = True + continue + choice = completion.choices[0] if ( - not args.compaction - or not can_compact(messages) - or not is_context_length_error(error) + args.compaction + and args.summarize_at_tokens is not None + and choice.finish_reason == "length" + and context_tokens(completion) >= args.summarize_at_tokens + and not recovered_overflow ): - raise - messages = await compact(client, args.model, messages, tools) - continue + messages = await compact(client, args.model, messages, tools) + recovered_overflow = True + continue + break + message = choice.message messages.append(message.model_dump(exclude_none=True)) if not message.tool_calls: break + tool_result_tokens = 0 for call in message.tool_calls: name = call.function.name tool_message = { @@ -509,7 +529,11 @@ async def main() -> None: tool_message, ) if decision["action"] == "rewrite": - messages.append(decision["message"]) + rewritten = decision["message"] + messages.append(rewritten) + tool_result_tokens += estimated_tokens( + str(rewritten.get("content", "")) + ) continue try: tool_args = json.loads(call.function.arguments or "{}") @@ -554,6 +578,14 @@ async def main() -> None: if decision["action"] == "rewrite": tool_message = decision["message"] messages.append(tool_message) + tool_result_tokens += estimated_tokens(str(tool_message["content"])) + if ( + args.compaction + and args.summarize_at_tokens is not None + and context_tokens(completion) + tool_result_tokens + >= args.summarize_at_tokens + ): + messages = await compact(client, args.model, messages, tools) if tool_client is not None: await tool_client.aclose() diff --git a/verifiers/v1/harnesses/rlm/harness.py b/verifiers/v1/harnesses/rlm/harness.py index 639d556c2e..9c9407cbb4 100644 --- a/verifiers/v1/harnesses/rlm/harness.py +++ b/verifiers/v1/harnesses/rlm/harness.py @@ -35,9 +35,7 @@ class _SessionSnapshot(BaseModel): class RLMHarnessConfig(HarnessConfig): - version: str = Field( - default="d4ce3e10e63b359f4f3d432d58a77471e9e21fe7", min_length=1 - ) + version: str = Field(default="9f64353", min_length=1) """Git ref (branch, tag, or commit) of nano-rlm to install.""" max_depth: int = 0 """Recursion depth RLM may spawn sub-harnesses to.""" diff --git a/verifiers/v1/interception/server.py b/verifiers/v1/interception/server.py index 09c1dd42b9..c485d9cb0b 100644 --- a/verifiers/v1/interception/server.py +++ b/verifiers/v1/interception/server.py @@ -81,8 +81,6 @@ # Attempt counter the stainless-generated SDKs (OpenAI, Anthropic) send on every request: # 0 on the first attempt, incremented on each retry of the same request. RETRY_COUNT_HEADER = "x-stainless-retry-count" -CONTEXT_COMPACTION_HEADER = "X-Verifiers-Context-Compaction" -"""Internal opt-in for a harness that can recover from an overlong prompt.""" IDEMPOTENCY_KEY_HEADER = "Idempotency-Key" IDEMPOTENCY_CACHE_TTL_SECONDS = 600 IDEMPOTENCY_CACHE_MAX_COMPLETED = 64 @@ -495,12 +493,6 @@ async def handle_request( body = dialect.apply_overrides(body, session.ctx.model, session.ctx.sampling) streaming = dialect.streaming(body) upstream_headers = dict(request.headers) - context_compaction = request.headers.get(CONTEXT_COMPACTION_HEADER) == "1" - upstream_headers = { - name: value - for name, value in upstream_headers.items() - if name.lower() != CONTEXT_COMPACTION_HEADER.lower() - } logger.debug( "intercept %s: id=%s stream=%s", request.path, @@ -724,16 +716,9 @@ async def sample() -> web.Response: ) except OverlongPromptError as e: error = e - if not context_compaction: - session.trace.stop("context_length") logger.debug("prompt too long: id=%s", session.trace.id) - message = ( - "context_length" - if context_compaction - else "rollout stopped: context_length" - ) return web.json_response( - dialect.error_body(message), + dialect.error_body("context_length"), status=400, ) except RolloutError as e: @@ -816,10 +801,9 @@ async def _stream( ) except OverlongPromptError as e: error = e - session.trace.stop("context_length") logger.debug("prompt too long: id=%s", session.trace.id) return web.json_response( - dialect.error_body("rollout stopped: context_length"), status=400 + dialect.error_body("context_length"), status=400 ) except RolloutError as e: error = e @@ -1020,10 +1004,7 @@ async def _stream( await resp.write_eof() return resp except OverlongPromptError as e: - # A streamed terminal provider failure is discovered only after its response body - # was relayed. Context exhaustion remains a clean truncation like earlier failures. error = e - session.trace.stop("context_length") logger.debug("prompt too long: id=%s", session.trace.id) return resp except RolloutError as e: From c1a13eee625e24d00e4a41cda856d3d46aaa362a Mon Sep 17 00:00:00 2001 From: Mika Senghaas Date: Wed, 26 Aug 2026 21:54:21 +0000 Subject: [PATCH 03/14] simplify compaction flow --- pyproject.toml | 1 - tests/v1/conftest.py | 9 +- tests/v1/fixtures/context_compaction_v1.py | 86 --------- tests/v1/test_context.py | 70 ++++++++ tests/v1/test_e2e.py | 64 ------- verifiers/v1/clients/context.py | 50 ++++++ verifiers/v1/configs/harness.py | 27 ++- verifiers/v1/harnesses/bash/harness.py | 22 ++- verifiers/v1/harnesses/bash/program.py | 193 +++++++++++---------- verifiers/v1/harnesses/rlm/harness.py | 52 +++--- 10 files changed, 283 insertions(+), 291 deletions(-) delete mode 100644 tests/v1/fixtures/context_compaction_v1.py create mode 100644 tests/v1/test_context.py create mode 100644 verifiers/v1/clients/context.py diff --git a/pyproject.toml b/pyproject.toml index 6781bba70d..056ebad809 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -290,7 +290,6 @@ markers = [ "bash: v1 e2e cases on the bash harness", "browser_use: v1 e2e cases on the browser_use harness", "rlm: v1 e2e cases on the rlm harness", - "compaction: v1 context-compaction E2E cases", "kimi_code: v1 e2e cases on the kimi-code harness", "pi: v1 e2e cases on the pi harness", "pool: v1 e2e cases on the pool harness", diff --git a/tests/v1/conftest.py b/tests/v1/conftest.py index 63e1549624..0b5ba36de9 100644 --- a/tests/v1/conftest.py +++ b/tests/v1/conftest.py @@ -96,7 +96,7 @@ def pytest_configure(config) -> None: def pytest_collection_modifyitems(config, items) -> None: """Skip the live-model tests (marked `e2e`) when no model endpoint is configured, so the rest of the suite (e.g. config parsing) still runs in a keyless environment.""" - if os.environ.get("PRIME_API_KEY") or os.environ.get("VF_COMPACTION_E2E_BASE_URL"): + if os.environ.get("PRIME_API_KEY"): return skip = pytest.mark.skip(reason="needs PRIME_API_KEY") for item in items: @@ -124,7 +124,7 @@ def _eval_config( harness: str | HarnessConfig | None = "null", n: int = 1, num_tasks: int = 1, - max_tokens: int | None = 2048, + max_tokens: int = 2048, max_turns: int | None = 4, rollout_timeout: float = 180, taskset_overrides: dict | None = None, @@ -132,8 +132,6 @@ def _eval_config( env: dict | None = None, pool: dict | None = None, reasoning_effort: str | None = None, - model: str = CI_MODEL, - client: dict | None = None, server: bool = False, ) -> EvalConfig: """Build the smallest `EvalConfig` that still exercises the path, shared by the in-process @@ -189,8 +187,7 @@ def _eval_config( serve=({"pool": pool} if pool else {}) if server else None, output_dir=output_dir.parent, run={"dir": output_dir.name}, - model=model, - client=client or {"type": "eval"}, + model=CI_MODEL, ) diff --git a/tests/v1/fixtures/context_compaction_v1.py b/tests/v1/fixtures/context_compaction_v1.py deleted file mode 100644 index 79ba6bfb62..0000000000 --- a/tests/v1/fixtures/context_compaction_v1.py +++ /dev/null @@ -1,86 +0,0 @@ -"""Context-compaction E2E scenarios for harness agent loops.""" - -from typing import Literal - -from pydantic import Field - -import verifiers.v1 as vf - - -class OverflowToolsetConfig(vf.ToolsetConfig): - payload_chars: int = Field(65_536, gt=0) - - -class OverflowToolset(vf.Toolset[OverflowToolsetConfig]): - TOOL_PREFIX = "overflow" - - def __init__(self, config: OverflowToolsetConfig): - super().__init__(config) - self.called = False - - @vf.tool - def overflow_context(self) -> str: - """Return a payload that is intentionally larger than the model context.""" - if self.called: - return "The overflow already occurred. Answer `recovered` now." - self.called = True - block = "0123456789abcdef " - repeats = self.config.payload_chars // len(block) + 1 - return (block * repeats)[: self.config.payload_chars] - - -class ContextCompactionTaskConfig(vf.TaskConfig): - scenario: Literal["decode", "tool_result"] = "decode" - payload_chars: int = Field(65_536, gt=0) - tools: OverflowToolsetConfig = OverflowToolsetConfig() - - -class ContextCompactionTask( - vf.Task[vf.TaskData, vf.State, ContextCompactionTaskConfig] -): - @classmethod - def toolsets(cls, config: ContextCompactionTaskConfig) -> list[vf.Toolset]: - if config.scenario != "tool_result": - return [] - tool_config = config.tools.model_copy( - update={"payload_chars": config.payload_chars} - ) - return [OverflowToolset(tool_config)] - - @vf.reward - async def compacted(self, trace: vf.Trace) -> float: - return float(trace.num_branches > 1) - - -class ContextCompactionConfig(vf.TasksetConfig): - task: ContextCompactionTaskConfig = ContextCompactionTaskConfig() - - -class ContextCompactionTaskset( - vf.Taskset[ContextCompactionTask, ContextCompactionConfig] -): - def load(self) -> list[ContextCompactionTask]: - if self.config.task.scenario == "decode": - prompt = ( - "Write `x ` repeatedly. Do not use tools and do not stop. " - "Continue until the model context ends the decode." - ) - else: - prompt = ( - "Call the `overflow_context` tool exactly once, then answer `recovered`. " - "In an RLM IPython session, call it with " - "`result = await overflow_overflow_context(); print(result)`." - ) - return [ - ContextCompactionTask( - vf.TaskData(idx=0, prompt=prompt), - self.config.task, - ) - ] - - -__all__ = ["ContextCompactionTaskset"] - - -if __name__ == "__main__": - OverflowToolset(OverflowToolsetConfig()).run() diff --git a/tests/v1/test_context.py b/tests/v1/test_context.py new file mode 100644 index 0000000000..b827c4ae1f --- /dev/null +++ b/tests/v1/test_context.py @@ -0,0 +1,70 @@ +import httpx +from openai import BadRequestError + +from verifiers.v1.clients.context import ( + compaction_threshold, + model_context_window, +) +from verifiers.v1.configs.harness import CompactionConfig +from verifiers.v1.harnesses.bash.harness import BashHarnessConfig +from verifiers.v1.harnesses.bash.program import context_error +from verifiers.v1.harnesses.rlm.harness import RLMHarnessConfig + + +def test_model_context_window_reads_vllm_extension() -> None: + payload = { + "data": [ + {"id": "other", "max_model_len": 1}, + {"id": "target", "max_model_len": 32_768}, + ] + } + + assert model_context_window(payload, "target") == 32_768 + + +def test_model_context_window_accepts_common_provider_extensions() -> None: + payload = {"data": [{"id": "target", "context_length": 128_000}]} + + assert model_context_window(payload, "target") == 128_000 + + +def test_model_context_window_is_unknown_for_standard_model_card() -> None: + payload = { + "data": [ + { + "id": "target", + "object": "model", + "created": 1, + "owned_by": "provider", + } + ] + } + + assert model_context_window(payload, "target") is None + + +def test_compaction_threshold_reserves_ten_percent() -> None: + assert compaction_threshold(32_768) == 29_491 + + +def test_compaction_is_disabled_by_default_for_both_harnesses() -> None: + assert BashHarnessConfig().compaction is None + assert RLMHarnessConfig().compaction is None + + +def test_compaction_config_has_shared_automatic_default() -> None: + assert CompactionConfig().summarize_at_tokens is None + + +def test_threshold_is_learned_from_provider_error() -> None: + response = httpx.Response( + 400, + request=httpx.Request("POST", "http://provider/v1/chat/completions"), + ) + error = BadRequestError( + "maximum context length is 32,768 tokens", + response=response, + body={"error": {"message": "maximum context length is 32,768 tokens"}}, + ) + + assert context_error(error) == (True, 29_491) diff --git a/tests/v1/test_e2e.py b/tests/v1/test_e2e.py index 9e73157f67..cef9bf7e16 100644 --- a/tests/v1/test_e2e.py +++ b/tests/v1/test_e2e.py @@ -5,12 +5,8 @@ with distinct networking — instead of fanning the full cross product. prime/modal rows are local-only (their marks are excluded in CI).""" -import os - import pytest -from verifiers.v1.utils.loaders import harness_config_type - mark = pytest.mark @@ -160,66 +156,6 @@ async def test_single_turn(run_v1, harness, harness_runtime, tmp_path): assert call.time.duration > 0 -@pytest.mark.e2e -@pytest.mark.compaction -@pytest.mark.docker -@pytest.mark.parametrize("scenario", ["decode", "tool_result"]) -@pytest.mark.parametrize("harness_id", ["bash", "rlm"]) -async def test_context_compaction_matrix(run_v1, scenario, harness_id, tmp_path): - """Both in-house loops recover when decoding or tool output fills context.""" - base_url = os.environ.get("VF_COMPACTION_E2E_BASE_URL") - model = os.environ.get("VF_COMPACTION_E2E_MODEL") - context_window = int(os.environ.get("VF_COMPACTION_E2E_CONTEXT_WINDOW", "4096")) - if not base_url or not model: - pytest.skip("needs a local compaction E2E model") - - harness = { - "id": harness_id, - "summarize_at_tokens": context_window, - **({"compaction": True} if harness_id == "bash" else {}), - } - sampling = ( - {"extra_body": {"ignore_eos": True}} - if scenario == "decode" - else {"temperature": 0.0} - ) - (trace,) = await run_v1( - "context-compaction-v1", - harness=harness_config_type(harness_id).model_validate(harness), - runtime={"type": "docker"}, - env={"agent": {"sampling": sampling, "max_output_tokens": None}}, - client={ - "type": "eval", - "base_url": base_url, - "api_key_var": "VF_COMPACTION_E2E_API_KEY", - }, - model=model, - max_tokens=None if scenario == "decode" else 512, - max_turns=8, - rollout_timeout=600, - taskset_overrides={ - "task": { - "scenario": scenario, - "payload_chars": context_window * 12, - } - }, - output_dir=tmp_path / f"{scenario}-{harness_id}", - ) - - assert trace.ok, trace.errors - assert trace.num_branches > 1 - assert trace.rewards["compacted"].score == 1.0 - if scenario == "decode": - assert any(call.finish_reason == "length" for call in trace.calls) - else: - assert any( - call.error is not None and call.error.type == "OverlongPromptError" - for call in trace.calls - ) - if harness_id == "rlm": - assert trace.metrics["num_compactions"] >= 1 - - @pytest.mark.e2e @pytest.mark.browser_use @pytest.mark.docker diff --git a/verifiers/v1/clients/context.py b/verifiers/v1/clients/context.py new file mode 100644 index 0000000000..0920a76775 --- /dev/null +++ b/verifiers/v1/clients/context.py @@ -0,0 +1,50 @@ +"""Model context-window discovery for OpenAI-compatible endpoints.""" + +from collections.abc import Mapping +from typing import Any, cast + +from openai import APIError + +from verifiers.v1.clients.base import build_async_openai +from verifiers.v1.clients.client import ModelContext + +CONTEXT_WINDOW_FIELDS = ( + "max_model_len", + "context_length", + "context_window", + "max_context_length", +) +_context_window_cache: dict[tuple[str, str], int | None] = {} + + +def model_context_window(payload: Mapping[str, Any], model: str) -> int | None: + """Read a provider context-window extension from one model card.""" + for card in payload.get("data") or []: + if not isinstance(card, Mapping) or card.get("id") != model: + continue + for field in CONTEXT_WINDOW_FIELDS: + value = card.get(field) + if isinstance(value, int) and not isinstance(value, bool) and value > 0: + return value + break + return None + + +def compaction_threshold(context_window: int) -> int: + """Reserve ten percent of the model context for checkpointing.""" + return max(1, context_window * 9 // 10) + + +async def resolve_compaction_threshold(ctx: ModelContext) -> int | None: + """Discover a model's proactive compaction threshold when advertised.""" + key = (ctx.client.model_dump_json(), ctx.model) + if key not in _context_window_cache: + try: + async with build_async_openai(ctx.client) as client: + payload = await client.get("/models", cast_to=cast(Any, dict[str, Any])) + _context_window_cache[key] = model_context_window(payload, ctx.model) + except APIError: + _context_window_cache[key] = None + + window = _context_window_cache[key] + return compaction_threshold(window) if window is not None else None diff --git a/verifiers/v1/configs/harness.py b/verifiers/v1/configs/harness.py index 068a8169f1..844018c8fc 100644 --- a/verifiers/v1/configs/harness.py +++ b/verifiers/v1/configs/harness.py @@ -3,14 +3,39 @@ from __future__ import annotations import os +import random from pathlib import Path -from pydantic import ConfigDict, Field, FiniteFloat +from pydantic import ConfigDict, Field, FiniteFloat, PositiveInt, model_validator from pydantic_config import BaseConfig from verifiers.v1.types import ID +class CompactionConfig(BaseConfig): + """Optional context compaction policy for in-house agent loops.""" + + summarize_at_tokens: PositiveInt | tuple[PositiveInt, PositiveInt] | None = None + """Compact at this token count. A pair draws a task-seeded threshold. When unset, use + 90% of the model context window when the provider advertises it.""" + + @model_validator(mode="after") + def validate_range(self) -> CompactionConfig: + value = self.summarize_at_tokens + if isinstance(value, tuple) and value[0] > value[1]: + raise ValueError( + "`summarize_at_tokens` range must be (lo, hi) with lo <= hi." + ) + return self + + def summarize_threshold(self, task_idx: int | None) -> int | None: + value = self.summarize_at_tokens + if isinstance(value, tuple): + lo, hi = value + return random.Random(task_idx or 0).randint(lo, hi) + return value + + class HarnessConfig(BaseConfig): id: ID = "bash" """Installed harness package, set through the seat's diff --git a/verifiers/v1/harnesses/bash/harness.py b/verifiers/v1/harnesses/bash/harness.py index f97e30b2b0..920bf9160b 100644 --- a/verifiers/v1/harnesses/bash/harness.py +++ b/verifiers/v1/harnesses/bash/harness.py @@ -2,10 +2,9 @@ import os from pathlib import Path -from pydantic import PositiveInt - from verifiers.v1.clients import ModelContext -from verifiers.v1.configs.harness import HarnessConfig +from verifiers.v1.clients.context import resolve_compaction_threshold +from verifiers.v1.configs.harness import CompactionConfig, HarnessConfig from verifiers.v1.dialects.chat import message_to_wire from verifiers.v1.harness import Harness from verifiers.v1.runtimes import ProgramResult, Runtime @@ -30,12 +29,8 @@ class BashHarnessConfig(HarnessConfig): - compaction: bool = True - """Recover from context exhaustion with a handoff summary and a fresh branch.""" - - summarize_at_tokens: PositiveInt | None = None - """Compact proactively when the estimated active context reaches this threshold. When unset, - compaction still recovers from provider context errors.""" + compaction: CompactionConfig | None = None + """Context compaction policy. Set an empty config to use automatic thresholds.""" edit: bool = True """Offer the local `edit` tool (single-occurrence string replacement in a file) alongside @@ -86,10 +81,13 @@ async def launch( ] if tool_interception_url: args.append(f"--tool-interception-url={tool_interception_url}") - if self.config.compaction: + if self.config.compaction is not None: args.append("--compaction") - if self.config.summarize_at_tokens is not None: - args.append(f"--summarize-at-tokens={self.config.summarize_at_tokens}") + threshold = self.config.compaction.summarize_threshold(data.idx) + if threshold is None: + threshold = await resolve_compaction_threshold(ctx) + if threshold is not None: + args.append(f"--summarize-at-tokens={threshold}") if self.config.edit: args.append("--edit") if self.config.search: diff --git a/verifiers/v1/harnesses/bash/program.py b/verifiers/v1/harnesses/bash/program.py index 1ed6e16afa..407ef1c6b0 100644 --- a/verifiers/v1/harnesses/bash/program.py +++ b/verifiers/v1/harnesses/bash/program.py @@ -7,6 +7,7 @@ import argparse import asyncio import json +import re import subprocess from contextlib import AsyncExitStack, asynccontextmanager, suppress from pathlib import Path @@ -20,30 +21,34 @@ MCP_CALL_ATTEMPTS = 6 MCP_TIMEOUT = 600.0 -CHECKPOINT_COMPACTION_PROMPT = ( - "You are performing a CONTEXT CHECKPOINT COMPACTION. " - "Create a handoff summary for another LLM that will resume the task.\n" - "\n" - "Include:\n" - "- Current progress and key decisions made\n" - "- Important context, constraints, or user preferences\n" - "- What remains to be done (clear next steps)\n" - "- Any critical data, examples, or references needed to continue\n" - "\n" - "Be concise, structured, and focused on helping the next LLM " - "seamlessly continue the work." -) - -POST_COMPACTION_FRAMING = ( - "Another language model started to solve this problem and produced " - "a summary of its thinking process. Use this to build on the work " - "that has already been done and avoid duplicating work. Here is " - "the summary produced by the other language model, use the " - "information in this summary to assist with your own analysis:" -) - -COMPACTED_TOOL_RESULT = ( - "[tool output dropped because it exceeded the model context limit]" +CHECKPOINT_COMPACTION_PROMPT = """You are performing a CONTEXT CHECKPOINT COMPACTION. Create a handoff summary for another LLM that will resume the task. + +Include: +- Current progress and key decisions made +- Important context, constraints, or user preferences +- What remains to be done (clear next steps) +- Any critical data, examples, or references needed to continue + +Be concise, structured, and focused on helping the next LLM seamlessly continue the work.""" + +POST_COMPACTION_FRAMING = """Another language model started to solve this problem and produced \ +a summary of its thinking process. Use this to build on the work \ +that has already been done and avoid duplicating work. Here is \ +the summary produced by the other language model, use the \ +information in this summary to assist with your own analysis:""" + +COMPACTED_TOOL_RESULT = "[tool output dropped because it exceeded the context limit]" + +CONTEXT_WINDOW_PATTERNS = ( + re.compile( + r"(?:maximum|max(?:imum)?)[^.\n]{0,40}(?:context length|context window)" + r"[^\d]{0,20}([\d,]+)", + re.IGNORECASE, + ), + re.compile( + r"[\"']?(?:max_model_len|context_length)[\"']?\s*[:=]\s*([\d,]+)", + re.IGNORECASE, + ), ) @@ -217,10 +222,10 @@ async def chat( return await client.chat.completions.create(**kwargs) -def is_context_length_error(error: BadRequestError) -> bool: - details = f"{error} {error.body or ''}".casefold() - return any( - marker in details +def context_error(error: BadRequestError) -> tuple[bool, int | None]: + details = f"{error} {error.body or ''}" + overflow = any( + marker in details.casefold() for marker in ( "request entity too large", "context_length", @@ -231,6 +236,12 @@ def is_context_length_error(error: BadRequestError) -> bool: "token limit exceeded", ) ) + for pattern in CONTEXT_WINDOW_PATTERNS: + match = pattern.search(details) + if match: + context_window = int(match.group(1).replace(",", "")) + return overflow, max(1, context_window * 9 // 10) + return overflow, None def drop_latest_tool_result(messages: list[dict]) -> bool: @@ -257,38 +268,62 @@ def context_tokens(completion) -> int: return (usage.prompt_tokens or 0) + (usage.completion_tokens or 0) -async def compact( - client: AsyncOpenAI, - model: str, - messages: list[dict], - tools: list[dict], -) -> list[dict]: - """Create a handoff summary, removing tool output only when the checkpoint overflows.""" - system_messages = [ - message for message in messages if message.get("role") == "system" - ] - while True: - checkpoint_messages = [ - *messages, - {"role": "user", "content": CHECKPOINT_COMPACTION_PROMPT}, - ] +class Compactor: + """Compact once and retry once when a model turn exhausts its context.""" + + def __init__(self, client, model, tools, enabled, threshold): + self.client = client + self.model = model + self.tools = tools + self.enabled = enabled + self.threshold = threshold + + def reached(self, completion, extra_tokens: int = 0) -> bool: + return ( + self.enabled + and self.threshold is not None + and context_tokens(completion) + extra_tokens >= self.threshold + ) + + async def complete(self, messages: list[dict]): try: - completion = await chat( - client, - model, - checkpoint_messages, - tools, - tool_choice="none", - ) + completion = await chat(self.client, self.model, messages, self.tools) except BadRequestError as error: - if not is_context_length_error(error): + overflow, threshold = context_error(error) + if not self.enabled or not overflow: raise - if not drop_latest_tool_result(messages): - raise - continue - summary = completion.choices[0].message - framed = POST_COMPACTION_FRAMING + "\n\n" + (summary.content or "") - return [*system_messages, {"role": "user", "content": framed}] + if self.threshold is None: + self.threshold = threshold + else: + choice = completion.choices[0] + if choice.finish_reason != "length" or not self.reached(completion): + return completion, messages + + messages = await self.compact(messages) + completion = await chat(self.client, self.model, messages, self.tools) + return completion, messages + + async def compact(self, messages: list[dict]) -> list[dict]: + system = [message for message in messages if message.get("role") == "system"] + while True: + checkpoint = [ + *messages, + {"role": "user", "content": CHECKPOINT_COMPACTION_PROMPT}, + ] + try: + completion = await chat( + self.client, + self.model, + checkpoint, + self.tools, + tool_choice="none", + ) + summary = completion.choices[0].message.content or "" + framed = POST_COMPACTION_FRAMING + "\n\n" + summary + return [*system, {"role": "user", "content": framed}] + except BadRequestError as error: + if not context_error(error)[0] or not drop_latest_tool_result(messages): + raise @asynccontextmanager @@ -479,33 +514,16 @@ async def main() -> None: messages.extend(initial) elif args.prompt: messages.append({"role": "user", "content": args.prompt}) + compactor = Compactor( + client, + args.model, + tools, + args.compaction, + args.summarize_at_tokens, + ) while True: - recovered_overflow = False - while True: - try: - completion = await chat(client, args.model, messages, tools) - except BadRequestError as error: - if ( - not args.compaction - or recovered_overflow - or not is_context_length_error(error) - ): - raise - messages = await compact(client, args.model, messages, tools) - recovered_overflow = True - continue - choice = completion.choices[0] - if ( - args.compaction - and args.summarize_at_tokens is not None - and choice.finish_reason == "length" - and context_tokens(completion) >= args.summarize_at_tokens - and not recovered_overflow - ): - messages = await compact(client, args.model, messages, tools) - recovered_overflow = True - continue - break + completion, messages = await compactor.complete(messages) + choice = completion.choices[0] message = choice.message messages.append(message.model_dump(exclude_none=True)) if not message.tool_calls: @@ -579,13 +597,8 @@ async def main() -> None: tool_message = decision["message"] messages.append(tool_message) tool_result_tokens += estimated_tokens(str(tool_message["content"])) - if ( - args.compaction - and args.summarize_at_tokens is not None - and context_tokens(completion) + tool_result_tokens - >= args.summarize_at_tokens - ): - messages = await compact(client, args.model, messages, tools) + if compactor.reached(completion, tool_result_tokens): + messages = await compactor.compact(messages) if tool_client is not None: await tool_client.aclose() diff --git a/verifiers/v1/harnesses/rlm/harness.py b/verifiers/v1/harnesses/rlm/harness.py index 9c9407cbb4..e92b6cfc6d 100644 --- a/verifiers/v1/harnesses/rlm/harness.py +++ b/verifiers/v1/harnesses/rlm/harness.py @@ -1,15 +1,15 @@ """RLM over ACP, with MCP tools exposed as pre-imported IPython skills.""" import logging -import random import shlex from typing import Literal -from pydantic import BaseModel, ConfigDict, Field, PositiveInt, model_validator +from pydantic import BaseModel, ConfigDict, Field, model_validator from verifiers.v1.acp import ACPConfig, ACPHarness, ACPTurn, JsonObject from verifiers.v1.clients import ModelContext -from verifiers.v1.configs.harness import HarnessConfig +from verifiers.v1.clients.context import resolve_compaction_threshold +from verifiers.v1.configs.harness import CompactionConfig, HarnessConfig from verifiers.v1.runtimes import Runtime from verifiers.v1.task import TaskData from verifiers.v1.trace import Trace @@ -35,27 +35,15 @@ class _SessionSnapshot(BaseModel): class RLMHarnessConfig(HarnessConfig): - version: str = Field(default="9f64353", min_length=1) + version: str = Field(default="1e45450", min_length=1) """Git ref (branch, tag, or commit) of nano-rlm to install.""" max_depth: int = 0 """Recursion depth RLM may spawn sub-harnesses to.""" builtin_skills: list[BuiltinSkill] = Field(default_factory=list) """Built-in rlm skills to enable (RLM_SKILLS), e.g. `["edit"]`; empty enables none. The tool set is fixed (ipython); the base `skills` field takes SKILL.md paths.""" - summarize_at_tokens: PositiveInt | tuple[PositiveInt, PositiveInt] | None = None - """Auto-compaction threshold (RLM_SUMMARIZE_AT_TOKENS): compact the context once it grows - past this many tokens. An int is a fixed threshold; a `(lo, hi)` pair draws a per-group - threshold (seeded by the task index, so a task's rollouts share one draw and tasks vary). - `None` disables auto-compaction; ints must be positive.""" - - @model_validator(mode="after") - def validate_range(self) -> "RLMHarnessConfig": - value = self.summarize_at_tokens - if isinstance(value, tuple) and value[0] > value[1]: - raise ValueError( - "`summarize_at_tokens` range must be (lo, hi) with lo <= hi." - ) - return self + compaction: CompactionConfig | None = None + """Context compaction policy. Set an empty config to use automatic thresholds.""" @model_validator(mode="after") def reject_disabled_tools(self) -> "RLMHarnessConfig": @@ -95,16 +83,6 @@ async def setup(self, runtime: Runtime) -> None: raise RuntimeError(f"rlm install failed: {result.stderr.strip()[-500:]}") await super().setup(runtime) - def summarize_threshold(self, task_idx: int | None) -> int | None: - """Resolve a fixed or per-task compaction threshold.""" - value = self.config.summarize_at_tokens - if value is None: - return None - if isinstance(value, tuple): - lo, hi = value - return random.Random(task_idx or 0).randint(lo, hi) - return value - def _runtime_metadata( self, ctx: ModelContext, @@ -112,8 +90,8 @@ def _runtime_metadata( runtime: Runtime, endpoint: str, secret: str, - data: TaskData, system_prompt: str | None, + compaction: dict[str, int | None] | None, ) -> JsonObject: payload = { "session_id": trace.id, @@ -124,7 +102,7 @@ def _runtime_metadata( }, "policy": { "max_depth": self.config.max_depth, - "summarize_at_tokens": self.summarize_threshold(data.idx), + "compaction": compaction, "max_concurrent_subagents": max(4, self.config.max_depth), }, "system_prompt_path": None, @@ -146,12 +124,24 @@ async def prepare_acp( data: TaskData, ) -> ACPConfig: system_prompt, prompt = self.resolve_prompt(data) + compaction = None + if self.config.compaction is not None: + summarize_at_tokens = self.config.compaction.summarize_threshold(data.idx) + if summarize_at_tokens is None: + summarize_at_tokens = await resolve_compaction_threshold(ctx) + compaction = {"summarize_at_tokens": summarize_at_tokens} return ACPConfig( env={**self.config.resolved_env, "RLM_HOME": self._home(trace)}, command=[RLM_BIN, "--acp"], prompt=prompt, session_meta=self._runtime_metadata( - ctx, trace, runtime, endpoint, secret, data, system_prompt + ctx, + trace, + runtime, + endpoint, + secret, + system_prompt, + compaction, ), ) From 4bb99de312eb6185d47bdfd61af84e50c30094a9 Mon Sep 17 00:00:00 2001 From: Mika Senghaas Date: Wed, 26 Aug 2026 21:54:25 +0000 Subject: [PATCH 04/14] propagate context limit errors --- verifiers/v1/interception/server.py | 4 ---- 1 file changed, 4 deletions(-) diff --git a/verifiers/v1/interception/server.py b/verifiers/v1/interception/server.py index c485d9cb0b..5a824ade0a 100644 --- a/verifiers/v1/interception/server.py +++ b/verifiers/v1/interception/server.py @@ -1003,10 +1003,6 @@ async def _stream( await resp.write(event) await resp.write_eof() return resp - except OverlongPromptError as e: - error = e - logger.debug("prompt too long: id=%s", session.trace.id) - return resp except RolloutError as e: # A streamed terminal provider failure is discovered only after the # response body has been relayed. Keep it off the graph and preserve From 14721235db3bd8ea0b1c38d49d5a50b3ae6522f9 Mon Sep 17 00:00:00 2001 From: Mika Senghaas Date: Wed, 26 Aug 2026 22:12:47 +0000 Subject: [PATCH 05/14] move compaction config into each harness Review: compaction is a per-loop policy, so each in-house harness owns its CompactionConfig and threshold discovery instead of sharing a top-level config class and a clients module. The duplicated resolver caches per (base_url, model) and queries the upstream /models card host-side, since the interception server serves no /models route for agent-side discovery. Co-Authored-By: Claude Fable 5 --- tests/v1/test_context.py | 70 -------------------------- verifiers/v1/clients/context.py | 50 ------------------ verifiers/v1/configs/harness.py | 27 +--------- verifiers/v1/harnesses/bash/harness.py | 68 ++++++++++++++++++++++++- verifiers/v1/harnesses/rlm/harness.py | 69 +++++++++++++++++++++++-- 5 files changed, 132 insertions(+), 152 deletions(-) delete mode 100644 tests/v1/test_context.py delete mode 100644 verifiers/v1/clients/context.py diff --git a/tests/v1/test_context.py b/tests/v1/test_context.py deleted file mode 100644 index b827c4ae1f..0000000000 --- a/tests/v1/test_context.py +++ /dev/null @@ -1,70 +0,0 @@ -import httpx -from openai import BadRequestError - -from verifiers.v1.clients.context import ( - compaction_threshold, - model_context_window, -) -from verifiers.v1.configs.harness import CompactionConfig -from verifiers.v1.harnesses.bash.harness import BashHarnessConfig -from verifiers.v1.harnesses.bash.program import context_error -from verifiers.v1.harnesses.rlm.harness import RLMHarnessConfig - - -def test_model_context_window_reads_vllm_extension() -> None: - payload = { - "data": [ - {"id": "other", "max_model_len": 1}, - {"id": "target", "max_model_len": 32_768}, - ] - } - - assert model_context_window(payload, "target") == 32_768 - - -def test_model_context_window_accepts_common_provider_extensions() -> None: - payload = {"data": [{"id": "target", "context_length": 128_000}]} - - assert model_context_window(payload, "target") == 128_000 - - -def test_model_context_window_is_unknown_for_standard_model_card() -> None: - payload = { - "data": [ - { - "id": "target", - "object": "model", - "created": 1, - "owned_by": "provider", - } - ] - } - - assert model_context_window(payload, "target") is None - - -def test_compaction_threshold_reserves_ten_percent() -> None: - assert compaction_threshold(32_768) == 29_491 - - -def test_compaction_is_disabled_by_default_for_both_harnesses() -> None: - assert BashHarnessConfig().compaction is None - assert RLMHarnessConfig().compaction is None - - -def test_compaction_config_has_shared_automatic_default() -> None: - assert CompactionConfig().summarize_at_tokens is None - - -def test_threshold_is_learned_from_provider_error() -> None: - response = httpx.Response( - 400, - request=httpx.Request("POST", "http://provider/v1/chat/completions"), - ) - error = BadRequestError( - "maximum context length is 32,768 tokens", - response=response, - body={"error": {"message": "maximum context length is 32,768 tokens"}}, - ) - - assert context_error(error) == (True, 29_491) diff --git a/verifiers/v1/clients/context.py b/verifiers/v1/clients/context.py deleted file mode 100644 index 0920a76775..0000000000 --- a/verifiers/v1/clients/context.py +++ /dev/null @@ -1,50 +0,0 @@ -"""Model context-window discovery for OpenAI-compatible endpoints.""" - -from collections.abc import Mapping -from typing import Any, cast - -from openai import APIError - -from verifiers.v1.clients.base import build_async_openai -from verifiers.v1.clients.client import ModelContext - -CONTEXT_WINDOW_FIELDS = ( - "max_model_len", - "context_length", - "context_window", - "max_context_length", -) -_context_window_cache: dict[tuple[str, str], int | None] = {} - - -def model_context_window(payload: Mapping[str, Any], model: str) -> int | None: - """Read a provider context-window extension from one model card.""" - for card in payload.get("data") or []: - if not isinstance(card, Mapping) or card.get("id") != model: - continue - for field in CONTEXT_WINDOW_FIELDS: - value = card.get(field) - if isinstance(value, int) and not isinstance(value, bool) and value > 0: - return value - break - return None - - -def compaction_threshold(context_window: int) -> int: - """Reserve ten percent of the model context for checkpointing.""" - return max(1, context_window * 9 // 10) - - -async def resolve_compaction_threshold(ctx: ModelContext) -> int | None: - """Discover a model's proactive compaction threshold when advertised.""" - key = (ctx.client.model_dump_json(), ctx.model) - if key not in _context_window_cache: - try: - async with build_async_openai(ctx.client) as client: - payload = await client.get("/models", cast_to=cast(Any, dict[str, Any])) - _context_window_cache[key] = model_context_window(payload, ctx.model) - except APIError: - _context_window_cache[key] = None - - window = _context_window_cache[key] - return compaction_threshold(window) if window is not None else None diff --git a/verifiers/v1/configs/harness.py b/verifiers/v1/configs/harness.py index 844018c8fc..068a8169f1 100644 --- a/verifiers/v1/configs/harness.py +++ b/verifiers/v1/configs/harness.py @@ -3,39 +3,14 @@ from __future__ import annotations import os -import random from pathlib import Path -from pydantic import ConfigDict, Field, FiniteFloat, PositiveInt, model_validator +from pydantic import ConfigDict, Field, FiniteFloat from pydantic_config import BaseConfig from verifiers.v1.types import ID -class CompactionConfig(BaseConfig): - """Optional context compaction policy for in-house agent loops.""" - - summarize_at_tokens: PositiveInt | tuple[PositiveInt, PositiveInt] | None = None - """Compact at this token count. A pair draws a task-seeded threshold. When unset, use - 90% of the model context window when the provider advertises it.""" - - @model_validator(mode="after") - def validate_range(self) -> CompactionConfig: - value = self.summarize_at_tokens - if isinstance(value, tuple) and value[0] > value[1]: - raise ValueError( - "`summarize_at_tokens` range must be (lo, hi) with lo <= hi." - ) - return self - - def summarize_threshold(self, task_idx: int | None) -> int | None: - value = self.summarize_at_tokens - if isinstance(value, tuple): - lo, hi = value - return random.Random(task_idx or 0).randint(lo, hi) - return value - - class HarnessConfig(BaseConfig): id: ID = "bash" """Installed harness package, set through the seat's diff --git a/verifiers/v1/harnesses/bash/harness.py b/verifiers/v1/harnesses/bash/harness.py index 920bf9160b..00309d2fbe 100644 --- a/verifiers/v1/harnesses/bash/harness.py +++ b/verifiers/v1/harnesses/bash/harness.py @@ -1,10 +1,16 @@ import json import os +import random from pathlib import Path +from typing import Any, cast + +from openai import APIError +from pydantic import PositiveInt, model_validator +from pydantic_config import BaseConfig from verifiers.v1.clients import ModelContext -from verifiers.v1.clients.context import resolve_compaction_threshold -from verifiers.v1.configs.harness import CompactionConfig, HarnessConfig +from verifiers.v1.clients.base import build_async_openai +from verifiers.v1.configs.harness import HarnessConfig from verifiers.v1.dialects.chat import message_to_wire from verifiers.v1.harness import Harness from verifiers.v1.runtimes import ProgramResult, Runtime @@ -28,6 +34,64 @@ ) +CONTEXT_WINDOW_FIELDS = ( + "max_model_len", + "context_length", + "context_window", + "max_context_length", +) +_context_window_cache: dict[tuple[str, str], int | None] = {} + + +class CompactionConfig(BaseConfig): + """Context compaction policy for the bash agent loop.""" + + summarize_at_tokens: PositiveInt | tuple[PositiveInt, PositiveInt] | None = None + """Compact at this token count. A pair draws a task-seeded threshold. When unset, use + 90% of the model context window when the provider advertises it.""" + + @model_validator(mode="after") + def validate_range(self) -> "CompactionConfig": + value = self.summarize_at_tokens + if isinstance(value, tuple) and value[0] > value[1]: + raise ValueError( + "`summarize_at_tokens` range must be (lo, hi) with lo <= hi." + ) + return self + + def summarize_threshold(self, task_idx: int | None) -> int | None: + value = self.summarize_at_tokens + if isinstance(value, tuple): + lo, hi = value + return random.Random(task_idx or 0).randint(lo, hi) + return value + + +async def resolve_compaction_threshold(ctx: ModelContext) -> int | None: + """90% of the model context window, when the provider's model card advertises one.""" + key = (ctx.client.base_url, ctx.model) + if key not in _context_window_cache: + window = None + try: + async with build_async_openai(ctx.client) as client: + payload = await client.get("/models", cast_to=cast(Any, dict[str, Any])) + except APIError: + payload = {} + for card in payload.get("data") or []: + if not isinstance(card, dict) or card.get("id") != ctx.model: + continue + for field in CONTEXT_WINDOW_FIELDS: + value = card.get(field) + if isinstance(value, int) and not isinstance(value, bool) and value > 0: + window = value + break + break + _context_window_cache[key] = window + + window = _context_window_cache[key] + return max(1, window * 9 // 10) if window is not None else None + + class BashHarnessConfig(HarnessConfig): compaction: CompactionConfig | None = None """Context compaction policy. Set an empty config to use automatic thresholds.""" diff --git a/verifiers/v1/harnesses/rlm/harness.py b/verifiers/v1/harnesses/rlm/harness.py index e92b6cfc6d..6949cc0512 100644 --- a/verifiers/v1/harnesses/rlm/harness.py +++ b/verifiers/v1/harnesses/rlm/harness.py @@ -1,15 +1,18 @@ """RLM over ACP, with MCP tools exposed as pre-imported IPython skills.""" import logging +import random import shlex -from typing import Literal +from typing import Any, Literal, cast -from pydantic import BaseModel, ConfigDict, Field, model_validator +from openai import APIError +from pydantic import BaseModel, ConfigDict, Field, PositiveInt, model_validator +from pydantic_config import BaseConfig from verifiers.v1.acp import ACPConfig, ACPHarness, ACPTurn, JsonObject from verifiers.v1.clients import ModelContext -from verifiers.v1.clients.context import resolve_compaction_threshold -from verifiers.v1.configs.harness import CompactionConfig, HarnessConfig +from verifiers.v1.clients.base import build_async_openai +from verifiers.v1.configs.harness import HarnessConfig from verifiers.v1.runtimes import Runtime from verifiers.v1.task import TaskData from verifiers.v1.trace import Trace @@ -34,6 +37,64 @@ class _SessionSnapshot(BaseModel): metrics: dict[str, int | float] +CONTEXT_WINDOW_FIELDS = ( + "max_model_len", + "context_length", + "context_window", + "max_context_length", +) +_context_window_cache: dict[tuple[str, str], int | None] = {} + + +class CompactionConfig(BaseConfig): + """Context compaction policy for the RLM agent loop.""" + + summarize_at_tokens: PositiveInt | tuple[PositiveInt, PositiveInt] | None = None + """Compact at this token count. A pair draws a task-seeded threshold. When unset, use + 90% of the model context window when the provider advertises it.""" + + @model_validator(mode="after") + def validate_range(self) -> "CompactionConfig": + value = self.summarize_at_tokens + if isinstance(value, tuple) and value[0] > value[1]: + raise ValueError( + "`summarize_at_tokens` range must be (lo, hi) with lo <= hi." + ) + return self + + def summarize_threshold(self, task_idx: int | None) -> int | None: + value = self.summarize_at_tokens + if isinstance(value, tuple): + lo, hi = value + return random.Random(task_idx or 0).randint(lo, hi) + return value + + +async def resolve_compaction_threshold(ctx: ModelContext) -> int | None: + """90% of the model context window, when the provider's model card advertises one.""" + key = (ctx.client.base_url, ctx.model) + if key not in _context_window_cache: + window = None + try: + async with build_async_openai(ctx.client) as client: + payload = await client.get("/models", cast_to=cast(Any, dict[str, Any])) + except APIError: + payload = {} + for card in payload.get("data") or []: + if not isinstance(card, dict) or card.get("id") != ctx.model: + continue + for field in CONTEXT_WINDOW_FIELDS: + value = card.get(field) + if isinstance(value, int) and not isinstance(value, bool) and value > 0: + window = value + break + break + _context_window_cache[key] = window + + window = _context_window_cache[key] + return max(1, window * 9 // 10) if window is not None else None + + class RLMHarnessConfig(HarnessConfig): version: str = Field(default="1e45450", min_length=1) """Git ref (branch, tag, or commit) of nano-rlm to install.""" From 43a0b6d5223e9b6b9c0b1b8f656a944a3084011b Mon Sep 17 00:00:00 2001 From: Mika Senghaas Date: Wed, 26 Aug 2026 22:12:47 +0000 Subject: [PATCH 06/14] relay overlong prompt errors untouched Review: interception rewrote overlong prompt failures into a synthesized "context_length" body. Let them fall through to the generic RolloutError relay instead: the harness receives the provider's original message (from which it can learn the context window) with the error's 400 status. A rollout whose final call overflowed without recovery now fails like any other provider error; "context_length" is gone as a stop condition, so drop it from Trace.is_truncated. Co-Authored-By: Claude Fable 5 --- verifiers/v1/errors.py | 12 ++++++------ verifiers/v1/interception/server.py | 14 -------------- verifiers/v1/trace.py | 1 - 3 files changed, 6 insertions(+), 21 deletions(-) diff --git a/verifiers/v1/errors.py b/verifiers/v1/errors.py index 98a1df523d..75dd11cca5 100644 --- a/verifiers/v1/errors.py +++ b/verifiers/v1/errors.py @@ -43,10 +43,10 @@ def __init__(self, message: str = "", *, status_code: int = 502) -> None: class OverlongPromptError(ProviderError): - """The prompt exceeded the model's context window — a budget limit, ended as a clean - truncation rather than recorded as an error. Defaults to a 400 (what the interception - server surfaces for it — deterministic, so an SDK never retries it); `model_error` - keeps the provider's real status when the failure carried one.""" + """The prompt exceeded the model's context window. Relayed to the harness like any other + provider error so it can compact and retry. Defaults to a 400 (deterministic, so an SDK + never retries it); `model_error` keeps the provider's real status when the failure + carried one.""" def __init__(self, message: str = "", *, status_code: int = 400) -> None: super().__init__(message, status_code=status_code) @@ -130,8 +130,8 @@ def _provider_status(e: OpenAIError | str) -> int: def model_error( e: OpenAIError | str, *, status_code: int | None = None ) -> ProviderError: - """Map a provider failure to our error type: an overlong prompt (a budget limit the interception - server turns into a clean truncation) is told apart from any other provider call failure, which + """Map a provider failure to our error type: an overlong prompt (which a harness may compact + and recover from) is told apart from any other provider call failure, which becomes a plain `ProviderError`. `status_code` is the HTTP status surfaced to the harness (whose SDK then retries 5xx/429/timeout and not 4xx); derived from an SDK error when not given. Accepts an SDK error (the renderer) or the provider's raw error body (the httpx proxy).""" diff --git a/verifiers/v1/interception/server.py b/verifiers/v1/interception/server.py index 5a824ade0a..fcb5cc9b41 100644 --- a/verifiers/v1/interception/server.py +++ b/verifiers/v1/interception/server.py @@ -46,7 +46,6 @@ is_sse_done_event, ) from verifiers.v1.errors import ( - OverlongPromptError, ProviderError, RolloutError, TaskError, @@ -714,13 +713,6 @@ async def sample() -> web.Response: dialect.error_body(f"rollout stopped: {stopped}"), status=400, ) - except OverlongPromptError as e: - error = e - logger.debug("prompt too long: id=%s", session.trace.id) - return web.json_response( - dialect.error_body("context_length"), - status=400, - ) except RolloutError as e: # Stash the real cause; the rollout re-raises it after the harness returns. # Relay the provider's status so the harness SDK retries 5xx/429 and not 4xx. @@ -799,12 +791,6 @@ async def _stream( headers=request.headers, session_id=session.trace.id, ) - except OverlongPromptError as e: - error = e - logger.debug("prompt too long: id=%s", session.trace.id) - return web.json_response( - dialect.error_body("context_length"), status=400 - ) except RolloutError as e: error = e session.error = e diff --git a/verifiers/v1/trace.py b/verifiers/v1/trace.py index 10d21a3948..8e1bafdbcd 100644 --- a/verifiers/v1/trace.py +++ b/verifiers/v1/trace.py @@ -522,7 +522,6 @@ def is_truncated(self) -> bool: "max_input_tokens", "max_output_tokens", "max_total_tokens", - "context_length", ): return True last = next((c for c in reversed(self.calls) if c.error is None), None) From c5b1b2e9809a7e5b927f40be255069850fcefeab Mon Sep 17 00:00:00 2001 From: Mika Senghaas Date: Wed, 26 Aug 2026 22:27:42 +0000 Subject: [PATCH 07/14] discover compaction thresholds harness-side Review: threshold discovery belongs to the agent loops, not the vf host. The bash program reads the provider's /models card itself (like nano-rlm's engine already does) and the vf harnesses only forward the explicitly configured threshold. The RLM policy crosses ACP flat (compaction toggle + summarize_at_tokens) to match nano-rlm's flattened ExecutionPolicy; pin bumped to f452d52. Behind the interception server (which serves no /models route) discovery yields nothing and both loops fall back to reactive compaction, learning the threshold from the relayed provider error. Co-Authored-By: Claude Fable 5 --- verifiers/v1/harnesses/bash/harness.py | 39 ---------------- verifiers/v1/harnesses/bash/program.py | 28 +++++++++++- verifiers/v1/harnesses/rlm/harness.py | 62 ++++---------------------- 3 files changed, 36 insertions(+), 93 deletions(-) diff --git a/verifiers/v1/harnesses/bash/harness.py b/verifiers/v1/harnesses/bash/harness.py index 00309d2fbe..07c33e424e 100644 --- a/verifiers/v1/harnesses/bash/harness.py +++ b/verifiers/v1/harnesses/bash/harness.py @@ -2,14 +2,11 @@ import os import random from pathlib import Path -from typing import Any, cast -from openai import APIError from pydantic import PositiveInt, model_validator from pydantic_config import BaseConfig from verifiers.v1.clients import ModelContext -from verifiers.v1.clients.base import build_async_openai from verifiers.v1.configs.harness import HarnessConfig from verifiers.v1.dialects.chat import message_to_wire from verifiers.v1.harness import Harness @@ -34,15 +31,6 @@ ) -CONTEXT_WINDOW_FIELDS = ( - "max_model_len", - "context_length", - "context_window", - "max_context_length", -) -_context_window_cache: dict[tuple[str, str], int | None] = {} - - class CompactionConfig(BaseConfig): """Context compaction policy for the bash agent loop.""" @@ -67,31 +55,6 @@ def summarize_threshold(self, task_idx: int | None) -> int | None: return value -async def resolve_compaction_threshold(ctx: ModelContext) -> int | None: - """90% of the model context window, when the provider's model card advertises one.""" - key = (ctx.client.base_url, ctx.model) - if key not in _context_window_cache: - window = None - try: - async with build_async_openai(ctx.client) as client: - payload = await client.get("/models", cast_to=cast(Any, dict[str, Any])) - except APIError: - payload = {} - for card in payload.get("data") or []: - if not isinstance(card, dict) or card.get("id") != ctx.model: - continue - for field in CONTEXT_WINDOW_FIELDS: - value = card.get(field) - if isinstance(value, int) and not isinstance(value, bool) and value > 0: - window = value - break - break - _context_window_cache[key] = window - - window = _context_window_cache[key] - return max(1, window * 9 // 10) if window is not None else None - - class BashHarnessConfig(HarnessConfig): compaction: CompactionConfig | None = None """Context compaction policy. Set an empty config to use automatic thresholds.""" @@ -148,8 +111,6 @@ async def launch( if self.config.compaction is not None: args.append("--compaction") threshold = self.config.compaction.summarize_threshold(data.idx) - if threshold is None: - threshold = await resolve_compaction_threshold(ctx) if threshold is not None: args.append(f"--summarize-at-tokens={threshold}") if self.config.edit: diff --git a/verifiers/v1/harnesses/bash/program.py b/verifiers/v1/harnesses/bash/program.py index 407ef1c6b0..8a02584833 100644 --- a/verifiers/v1/harnesses/bash/program.py +++ b/verifiers/v1/harnesses/bash/program.py @@ -13,7 +13,7 @@ from pathlib import Path import httpx -from openai import AsyncOpenAI, BadRequestError +from openai import APIError, AsyncOpenAI, BadRequestError from tenacity import AsyncRetrying, stop_after_attempt, wait_exponential_jitter SERPER_URL = "https://google.serper.dev/search" @@ -51,6 +51,30 @@ ), ) +CONTEXT_WINDOW_FIELDS = ( + "max_model_len", + "context_length", + "context_window", + "max_context_length", +) + + +async def discover_threshold(client: AsyncOpenAI, model: str) -> int | None: + """90% of the model context window, when the provider's model card advertises one.""" + try: + payload = await client.get("/models", cast_to=dict) + except APIError: + return None + for card in payload.get("data") or []: + if not isinstance(card, dict) or card.get("id") != model: + continue + for field in CONTEXT_WINDOW_FIELDS: + value = card.get(field) + if isinstance(value, int) and not isinstance(value, bool) and value > 0: + return max(1, value * 9 // 10) + break + return None + BASH_TOOL = { "type": "function", @@ -521,6 +545,8 @@ async def main() -> None: args.compaction, args.summarize_at_tokens, ) + if compactor.enabled and compactor.threshold is None: + compactor.threshold = await discover_threshold(client, args.model) while True: completion, messages = await compactor.complete(messages) choice = completion.choices[0] diff --git a/verifiers/v1/harnesses/rlm/harness.py b/verifiers/v1/harnesses/rlm/harness.py index 6949cc0512..f9ffa14374 100644 --- a/verifiers/v1/harnesses/rlm/harness.py +++ b/verifiers/v1/harnesses/rlm/harness.py @@ -3,15 +3,13 @@ import logging import random import shlex -from typing import Any, Literal, cast +from typing import Literal -from openai import APIError from pydantic import BaseModel, ConfigDict, Field, PositiveInt, model_validator from pydantic_config import BaseConfig from verifiers.v1.acp import ACPConfig, ACPHarness, ACPTurn, JsonObject from verifiers.v1.clients import ModelContext -from verifiers.v1.clients.base import build_async_openai from verifiers.v1.configs.harness import HarnessConfig from verifiers.v1.runtimes import Runtime from verifiers.v1.task import TaskData @@ -37,15 +35,6 @@ class _SessionSnapshot(BaseModel): metrics: dict[str, int | float] -CONTEXT_WINDOW_FIELDS = ( - "max_model_len", - "context_length", - "context_window", - "max_context_length", -) -_context_window_cache: dict[tuple[str, str], int | None] = {} - - class CompactionConfig(BaseConfig): """Context compaction policy for the RLM agent loop.""" @@ -70,33 +59,8 @@ def summarize_threshold(self, task_idx: int | None) -> int | None: return value -async def resolve_compaction_threshold(ctx: ModelContext) -> int | None: - """90% of the model context window, when the provider's model card advertises one.""" - key = (ctx.client.base_url, ctx.model) - if key not in _context_window_cache: - window = None - try: - async with build_async_openai(ctx.client) as client: - payload = await client.get("/models", cast_to=cast(Any, dict[str, Any])) - except APIError: - payload = {} - for card in payload.get("data") or []: - if not isinstance(card, dict) or card.get("id") != ctx.model: - continue - for field in CONTEXT_WINDOW_FIELDS: - value = card.get(field) - if isinstance(value, int) and not isinstance(value, bool) and value > 0: - window = value - break - break - _context_window_cache[key] = window - - window = _context_window_cache[key] - return max(1, window * 9 // 10) if window is not None else None - - class RLMHarnessConfig(HarnessConfig): - version: str = Field(default="1e45450", min_length=1) + version: str = Field(default="f452d52", min_length=1) """Git ref (branch, tag, or commit) of nano-rlm to install.""" max_depth: int = 0 """Recursion depth RLM may spawn sub-harnesses to.""" @@ -152,8 +116,9 @@ def _runtime_metadata( endpoint: str, secret: str, system_prompt: str | None, - compaction: dict[str, int | None] | None, + data: TaskData, ) -> JsonObject: + compaction = self.config.compaction payload = { "session_id": trace.id, "model": ctx.model, @@ -163,7 +128,10 @@ def _runtime_metadata( }, "policy": { "max_depth": self.config.max_depth, - "compaction": compaction, + "compaction": compaction is not None, + "summarize_at_tokens": ( + compaction.summarize_threshold(data.idx) if compaction else None + ), "max_concurrent_subagents": max(4, self.config.max_depth), }, "system_prompt_path": None, @@ -185,24 +153,12 @@ async def prepare_acp( data: TaskData, ) -> ACPConfig: system_prompt, prompt = self.resolve_prompt(data) - compaction = None - if self.config.compaction is not None: - summarize_at_tokens = self.config.compaction.summarize_threshold(data.idx) - if summarize_at_tokens is None: - summarize_at_tokens = await resolve_compaction_threshold(ctx) - compaction = {"summarize_at_tokens": summarize_at_tokens} return ACPConfig( env={**self.config.resolved_env, "RLM_HOME": self._home(trace)}, command=[RLM_BIN, "--acp"], prompt=prompt, session_meta=self._runtime_metadata( - ctx, - trace, - runtime, - endpoint, - secret, - system_prompt, - compaction, + ctx, trace, runtime, endpoint, secret, system_prompt, data ), ) From ff011aa79d306d57683c217c109816518a3ce61c Mon Sep 17 00:00:00 2001 From: Mika Senghaas Date: Wed, 26 Aug 2026 22:27:43 +0000 Subject: [PATCH 08/14] remove unused overlong prompt error Review: nothing catches or type-checks OverlongPromptError since interception relays it like any other provider failure. Raise a plain ProviderError with a deterministic 400 at the two sites that produced it (renderer pre-flight overflow, Responses context_length_exceeded) and drop the phrase-sniffing from model_error. Co-Authored-By: Claude Fable 5 --- verifiers/v1/clients/train.py | 11 ++++---- verifiers/v1/dialects/responses.py | 8 ++---- verifiers/v1/errors.py | 43 +++--------------------------- 3 files changed, 12 insertions(+), 50 deletions(-) diff --git a/verifiers/v1/clients/train.py b/verifiers/v1/clients/train.py index 528b8e12a5..6ffca79f70 100644 --- a/verifiers/v1/clients/train.py +++ b/verifiers/v1/clients/train.py @@ -10,8 +10,7 @@ from typing import Any, ClassVar, TypeVar from openai import OpenAIError -from renderers import OverlongPromptError as RendererOverlongPromptError -from renderers import RenderedTokens, Renderer, RendererConfig +from renderers import OverlongPromptError, RenderedTokens, Renderer, RendererConfig from renderers.base import ToolCallParseStatus, is_multimodal from verifiers.v1.clients.base import build_async_openai @@ -19,7 +18,7 @@ from verifiers.v1.configs.client import TrainClientConfig from verifiers.v1.dialects import FINISH_REASONS, ChatDialect, Dialect, parse_tools from verifiers.v1.dialects.chat import message_to_wire -from verifiers.v1.errors import OverlongPromptError, model_error +from verifiers.v1.errors import ProviderError, model_error from verifiers.v1.graph import PendingTurn from verifiers.v1.types import ( AssistantMessage, @@ -436,8 +435,10 @@ def bridge(): if session_id else None, ) - except RendererOverlongPromptError as e: - raise OverlongPromptError(str(e)) from e + except OverlongPromptError as e: + # The renderer's pre-flight overflow never reached the provider: a + # deterministic 400, so the harness SDK never retries it. + raise ProviderError(str(e), status_code=400) from e except OpenAIError as e: raise model_error(e) from e response = response_from_generate( diff --git a/verifiers/v1/dialects/responses.py b/verifiers/v1/dialects/responses.py index fdbcba4967..c558318e8c 100644 --- a/verifiers/v1/dialects/responses.py +++ b/verifiers/v1/dialects/responses.py @@ -28,7 +28,7 @@ parse_sse_event, provider_allowed_domains, ) -from verifiers.v1.errors import OverlongPromptError, model_error +from verifiers.v1.errors import model_error from verifiers.v1.types import ( AssistantMessage, ContentPart, @@ -322,15 +322,11 @@ def response_from_wire(response: OpenAIResponse) -> Response: code = error.get("code") if isinstance(error, dict) else None message = error.get("message") if isinstance(error, dict) else None detail = ": ".join(str(value) for value in (status, code, message) if value) - if code == "context_length_exceeded": - raise OverlongPromptError( - f"upstream Responses request did not complete: {detail}" - ) status_code = ( 429 if code in ("rate_limit_exceeded", "rate_limit_error") else 400 - if code == "invalid_prompt" + if code in ("invalid_prompt", "context_length_exceeded") else 502 ) raise model_error( diff --git a/verifiers/v1/errors.py b/verifiers/v1/errors.py index 75dd11cca5..e4766cee98 100644 --- a/verifiers/v1/errors.py +++ b/verifiers/v1/errors.py @@ -42,16 +42,6 @@ def __init__(self, message: str = "", *, status_code: int = 502) -> None: self.status_code = status_code -class OverlongPromptError(ProviderError): - """The prompt exceeded the model's context window. Relayed to the harness like any other - provider error so it can compact and retry. Defaults to a 400 (deterministic, so an SDK - never retries it); `model_error` keeps the provider's real status when the failure - carried one.""" - - def __init__(self, message: str = "", *, status_code: int = 400) -> None: - super().__init__(message, status_code=status_code) - - class HarnessError(RolloutError): """The harness failed to install or launch, or its agent process exited unsuccessfully.""" @@ -99,20 +89,6 @@ async def boundary(error_cls: type[RolloutError], what: str) -> AsyncIterator[No raise error_cls(f"{what}: {type(e).__name__}: {e}") from e -_CONTEXT_LENGTH_PHRASES = ( - "this model's maximum context length is", - "is longer than the model's context length", - "is longer than the maximum model length", - "exceeds the model's context length", - "exceed the configured limit", - "exceeds the configured limit", - "exceeded model", - "prompt_too_long", - "context length", - "maximum model length", -) - - def _provider_status(e: OpenAIError | str) -> int: """The HTTP status to surface for an SDK error: the provider's own for an HTTP status error, a retryable 5xx for a transport/timeout fault, else 502.""" @@ -130,23 +106,12 @@ def _provider_status(e: OpenAIError | str) -> int: def model_error( e: OpenAIError | str, *, status_code: int | None = None ) -> ProviderError: - """Map a provider failure to our error type: an overlong prompt (which a harness may compact - and recover from) is told apart from any other provider call failure, which - becomes a plain `ProviderError`. `status_code` is the HTTP status surfaced to the harness (whose - SDK then retries 5xx/429/timeout and not 4xx); derived from an SDK error when not given. Accepts - an SDK error (the renderer) or the provider's raw error body (the httpx proxy).""" - from openai import APIStatusError - + """Map a provider failure to a `ProviderError`. `status_code` is the HTTP status surfaced to + the harness (whose SDK then retries 5xx/429/timeout and not 4xx); derived from an SDK error + when not given. Accepts an SDK error (the renderer) or the provider's raw error body (the + httpx proxy).""" # Some SDK errors stringify empty; fall back to the type so the message is never blank. text = str(e) or (type(e).__name__ if isinstance(e, BaseException) else "") - if any(phrase in text.casefold() for phrase in _CONTEXT_LENGTH_PHRASES): - # Keep the provider's real status when the failure carried one; else the class - # default (the 400 the interception server surfaces for overlong prompts). - if status_code is None and isinstance(e, APIStatusError): - status_code = e.status_code - return OverlongPromptError( - text, **({} if status_code is None else {"status_code": status_code}) - ) return ProviderError( text, status_code=status_code if status_code is not None else _provider_status(e), From 70e8723a24bdc33143c7aa95797d6478d30ea37e Mon Sep 17 00:00:00 2001 From: Mika Senghaas Date: Wed, 26 Aug 2026 22:42:56 +0000 Subject: [PATCH 09/14] serve models through interception Add GET /v1/models to the interception server, relaying the upstream listing through the session's shared client (cached per endpoint, one upstream fetch per process). Agent loops can now discover their compaction threshold in training and eval alike. The path is dialect-universal: OpenAI and Anthropic SDKs both list models at /v1/models, so one route serves every dialect and only the auth carrier differs (tried per dialect). The response schema stays the upstream's; only OpenAI-compatible engines advertise a context window (vLLM's max_model_len), others fall back to reactive compaction. Also parse the listing as a parameterized mapping in the bash program - the OpenAI SDK cannot construct a bare dict, so discovery raised ValueError instead of returning None. Pin nano-rlm da9f4a4 with the same fix. Co-Authored-By: Claude Fable 5 --- verifiers/v1/clients/client.py | 6 ++++++ verifiers/v1/clients/eval.py | 24 +++++++++++++++++++++++ verifiers/v1/clients/train.py | 16 ++++++++++++++- verifiers/v1/harnesses/bash/program.py | 4 +++- verifiers/v1/harnesses/rlm/harness.py | 2 +- verifiers/v1/interception/server.py | 27 ++++++++++++++++++++++++++ 6 files changed, 76 insertions(+), 3 deletions(-) diff --git a/verifiers/v1/clients/client.py b/verifiers/v1/clients/client.py index 5b4bab1a0b..9b218ba936 100644 --- a/verifiers/v1/clients/client.py +++ b/verifiers/v1/clients/client.py @@ -66,6 +66,12 @@ async def relay_aux( as native JSON and return the provider JSON. Only the relay (eval) client supports it.""" raise NotImplementedError(f"{type(self).__name__} does not relay aux routes") + async def models(self, dialect: Dialect) -> dict: + """The upstream `GET /v1/models` listing, cached per client. Agent loops read a + provider context-window extension (e.g. vLLM's `max_model_len`) from it to set their + compaction threshold; the body is relayed as the provider sent it.""" + raise NotImplementedError(f"{type(self).__name__} does not list models") + async def close(self) -> None: pass diff --git a/verifiers/v1/clients/eval.py b/verifiers/v1/clients/eval.py index 74ac1c8d52..7645a6d515 100644 --- a/verifiers/v1/clients/eval.py +++ b/verifiers/v1/clients/eval.py @@ -1,5 +1,6 @@ """The eval client: proxies harness-native request to the provider.""" +import asyncio import re from collections.abc import Mapping @@ -65,6 +66,8 @@ def __init__(self, config: BaseClientConfig) -> None: # the dialect's provider authentication is applied. self.headers = dict(config.headers or {}) self.client = httpx.AsyncClient(timeout=DEFAULT_TIMEOUT, limits=DEFAULT_LIMITS) + self._models: dict | None = None + self._models_lock = asyncio.Lock() async def get_response( self, @@ -159,6 +162,27 @@ async def _request( f"upstream {response.status_code}: {text}", status_code=response.status_code ) + async def models(self, dialect: Dialect) -> dict: + # One fetch serves every rollout on this endpoint (the lock stops a launch stampede). + async with self._models_lock: + if self._models is None: + headers = httpx.Headers(self.headers) + headers.update(dialect.auth_headers(self.api_key)) + try: + response = await self.client.get( + join_url(self.base_url, "/v1/models"), headers=headers + ) + response.raise_for_status() + except httpx.HTTPStatusError as e: + raise model_error( + f"upstream {e.response.status_code}: {e.response.text}", + status_code=e.response.status_code, + ) from e + except httpx.HTTPError as e: + raise model_error(str(e), status_code=503) from e + self._models = from_json(response.content) + return self._models + async def relay( self, dialect: Dialect, diff --git a/verifiers/v1/clients/train.py b/verifiers/v1/clients/train.py index 6ffca79f70..b7a6ecb757 100644 --- a/verifiers/v1/clients/train.py +++ b/verifiers/v1/clients/train.py @@ -7,7 +7,7 @@ from collections.abc import AsyncIterator, Callable, Mapping from contextlib import asynccontextmanager from dataclasses import dataclass, field -from typing import Any, ClassVar, TypeVar +from typing import Any, ClassVar, TypeVar, cast from openai import OpenAIError from renderers import OverlongPromptError, RenderedTokens, Renderer, RendererConfig @@ -312,6 +312,8 @@ class TrainClient(Client): def __init__(self, config: TrainClientConfig) -> None: self.config = config self.client = build_async_openai(config) + self._models: dict | None = None + self._models_lock = asyncio.Lock() # The per-request model is only known at call time; a config that pins the renderer # model can warm now, which is every training run (prime-rl always pins it). if config.renderer_model_name is not None: @@ -321,6 +323,18 @@ def __init__(self, config: TrainClientConfig) -> None: multiplex=config.multiplex, ).warm() + async def models(self, dialect: Dialect) -> dict: + # One fetch serves every rollout on this engine (the lock stops a launch stampede). + async with self._models_lock: + if self._models is None: + try: + self._models = await self.client.get( + "/models", cast_to=cast(Any, dict[str, Any]) + ) + except OpenAIError as e: + raise model_error(e) from e + return self._models + async def get_response( self, dialect: Dialect, diff --git a/verifiers/v1/harnesses/bash/program.py b/verifiers/v1/harnesses/bash/program.py index 8a02584833..e66078dd89 100644 --- a/verifiers/v1/harnesses/bash/program.py +++ b/verifiers/v1/harnesses/bash/program.py @@ -11,6 +11,7 @@ import subprocess from contextlib import AsyncExitStack, asynccontextmanager, suppress from pathlib import Path +from typing import Any import httpx from openai import APIError, AsyncOpenAI, BadRequestError @@ -62,7 +63,8 @@ async def discover_threshold(client: AsyncOpenAI, model: str) -> int | None: """90% of the model context window, when the provider's model card advertises one.""" try: - payload = await client.get("/models", cast_to=dict) + # The SDK needs a parameterized mapping type to parse into (bare `dict` fails). + payload = await client.get("/models", cast_to=dict[str, Any]) except APIError: return None for card in payload.get("data") or []: diff --git a/verifiers/v1/harnesses/rlm/harness.py b/verifiers/v1/harnesses/rlm/harness.py index f9ffa14374..91883eaa23 100644 --- a/verifiers/v1/harnesses/rlm/harness.py +++ b/verifiers/v1/harnesses/rlm/harness.py @@ -60,7 +60,7 @@ def summarize_threshold(self, task_idx: int | None) -> int | None: class RLMHarnessConfig(HarnessConfig): - version: str = Field(default="f452d52", min_length=1) + version: str = Field(default="da9f4a4", min_length=1) """Git ref (branch, tag, or commit) of nano-rlm to install.""" max_depth: int = 0 """Recursion depth RLM may spawn sub-harnesses to.""" diff --git a/verifiers/v1/interception/server.py b/verifiers/v1/interception/server.py index fcb5cc9b41..d7b8698910 100644 --- a/verifiers/v1/interception/server.py +++ b/verifiers/v1/interception/server.py @@ -326,6 +326,9 @@ async def start(self) -> None: app.router.add_post(route, self._handler_for(dialect)) for aux in dialect.aux_routes: app.router.add_post(aux, self._aux_handler_for(dialect, aux)) + # One models route serves every dialect: OpenAI and Anthropic SDKs both list + # models at `GET /v1/models` (the response schema is the upstream's). + app.router.add_get("/v1/models", self.handle_models) # Tool servers use a state-only capability; the model bearer cannot reach these. app.router.add_get("/state", self.handle_state_get) app.router.add_put("/state", self.handle_state_put) @@ -1053,6 +1056,30 @@ async def handle_aux( return web.json_response(dialect.error_body(str(e)), status=502) return web.json_response(result) + async def handle_models(self, request: web.Request) -> web.Response: + """`GET /v1/models`: relay the upstream model listing so agent loops can read a + provider context-window extension (e.g. vLLM's `max_model_len`). The path is shared + by every dialect; only the auth carrier differs, so the bearer is tried per dialect. + Never recorded on the trace, and a failure never fails the rollout.""" + for dialect in DIALECTS: + session = self.sessions.get(dialect.secret(request.headers)) + if session is not None: + break + else: + return web.json_response({"error": "unauthorized"}, status=401) + session.adopt(asyncio.current_task()) + logger.debug("intercept models: id=%s", session.trace.id) + try: + payload = await session.client.models(dialect) + except NotImplementedError as e: + return web.json_response(dialect.error_body(str(e)), status=404) + except RolloutError as e: + logger.warning("models call failed: id=%s %s", session.trace.id, e) + return web.json_response( + dialect.error_body(str(e)), status=getattr(e, "status_code", 502) + ) + return web.json_response(payload) + def _session_for( self, request: web.Request, *, allow_service: bool = False ) -> RolloutSession | None: From 5b48a268a6830b49d2182ecacb138a530d3f8e3b Mon Sep 17 00:00:00 2001 From: Mika Senghaas Date: Thu, 27 Aug 2026 02:51:43 +0000 Subject: [PATCH 10/14] drop ranged compaction thresholds summarize_at_tokens is a plain token count now - the (lo, hi) task-seeded draw is no longer used. Co-Authored-By: Claude Fable 5 --- verifiers/v1/harnesses/bash/harness.py | 27 +++++-------------------- verifiers/v1/harnesses/rlm/harness.py | 28 +++++--------------------- 2 files changed, 10 insertions(+), 45 deletions(-) diff --git a/verifiers/v1/harnesses/bash/harness.py b/verifiers/v1/harnesses/bash/harness.py index 07c33e424e..cc5a876742 100644 --- a/verifiers/v1/harnesses/bash/harness.py +++ b/verifiers/v1/harnesses/bash/harness.py @@ -1,9 +1,8 @@ import json import os -import random from pathlib import Path -from pydantic import PositiveInt, model_validator +from pydantic import PositiveInt from pydantic_config import BaseConfig from verifiers.v1.clients import ModelContext @@ -34,25 +33,9 @@ class CompactionConfig(BaseConfig): """Context compaction policy for the bash agent loop.""" - summarize_at_tokens: PositiveInt | tuple[PositiveInt, PositiveInt] | None = None - """Compact at this token count. A pair draws a task-seeded threshold. When unset, use - 90% of the model context window when the provider advertises it.""" - - @model_validator(mode="after") - def validate_range(self) -> "CompactionConfig": - value = self.summarize_at_tokens - if isinstance(value, tuple) and value[0] > value[1]: - raise ValueError( - "`summarize_at_tokens` range must be (lo, hi) with lo <= hi." - ) - return self - - def summarize_threshold(self, task_idx: int | None) -> int | None: - value = self.summarize_at_tokens - if isinstance(value, tuple): - lo, hi = value - return random.Random(task_idx or 0).randint(lo, hi) - return value + summarize_at_tokens: PositiveInt | None = None + """Compact at this token count. When unset, use 90% of the model context window when + the provider advertises it.""" class BashHarnessConfig(HarnessConfig): @@ -110,7 +93,7 @@ async def launch( args.append(f"--tool-interception-url={tool_interception_url}") if self.config.compaction is not None: args.append("--compaction") - threshold = self.config.compaction.summarize_threshold(data.idx) + threshold = self.config.compaction.summarize_at_tokens if threshold is not None: args.append(f"--summarize-at-tokens={threshold}") if self.config.edit: diff --git a/verifiers/v1/harnesses/rlm/harness.py b/verifiers/v1/harnesses/rlm/harness.py index 91883eaa23..79aabdf546 100644 --- a/verifiers/v1/harnesses/rlm/harness.py +++ b/verifiers/v1/harnesses/rlm/harness.py @@ -1,7 +1,6 @@ """RLM over ACP, with MCP tools exposed as pre-imported IPython skills.""" import logging -import random import shlex from typing import Literal @@ -38,25 +37,9 @@ class _SessionSnapshot(BaseModel): class CompactionConfig(BaseConfig): """Context compaction policy for the RLM agent loop.""" - summarize_at_tokens: PositiveInt | tuple[PositiveInt, PositiveInt] | None = None - """Compact at this token count. A pair draws a task-seeded threshold. When unset, use - 90% of the model context window when the provider advertises it.""" - - @model_validator(mode="after") - def validate_range(self) -> "CompactionConfig": - value = self.summarize_at_tokens - if isinstance(value, tuple) and value[0] > value[1]: - raise ValueError( - "`summarize_at_tokens` range must be (lo, hi) with lo <= hi." - ) - return self - - def summarize_threshold(self, task_idx: int | None) -> int | None: - value = self.summarize_at_tokens - if isinstance(value, tuple): - lo, hi = value - return random.Random(task_idx or 0).randint(lo, hi) - return value + summarize_at_tokens: PositiveInt | None = None + """Compact at this token count. When unset, use 90% of the model context window when + the provider advertises it.""" class RLMHarnessConfig(HarnessConfig): @@ -116,7 +99,6 @@ def _runtime_metadata( endpoint: str, secret: str, system_prompt: str | None, - data: TaskData, ) -> JsonObject: compaction = self.config.compaction payload = { @@ -130,7 +112,7 @@ def _runtime_metadata( "max_depth": self.config.max_depth, "compaction": compaction is not None, "summarize_at_tokens": ( - compaction.summarize_threshold(data.idx) if compaction else None + compaction.summarize_at_tokens if compaction else None ), "max_concurrent_subagents": max(4, self.config.max_depth), }, @@ -158,7 +140,7 @@ async def prepare_acp( command=[RLM_BIN, "--acp"], prompt=prompt, session_meta=self._runtime_metadata( - ctx, trace, runtime, endpoint, secret, system_prompt, data + ctx, trace, runtime, endpoint, secret, system_prompt ), ) From f63d13c6bd6cb39565e461fe0b3775c937120097 Mon Sep 17 00:00:00 2001 From: Mika Senghaas Date: Thu, 27 Aug 2026 02:55:55 +0000 Subject: [PATCH 11/14] relay the models listing statelessly Review: the models listing is a proxy concern, not a client one. handle_models relays GET /v1/models straight from the session's endpoint config (the matched dialect contributes bearer extraction and provider auth), the upstream body and status pass through verbatim, and the Client.models method and its per-endpoint cache go away - every request hits the provider, like any other relay. Also order compaction last in BashHarnessConfig. Co-Authored-By: Claude Fable 5 --- verifiers/v1/clients/client.py | 6 ------ verifiers/v1/clients/eval.py | 24 --------------------- verifiers/v1/clients/train.py | 16 +------------- verifiers/v1/harnesses/bash/harness.py | 6 +++--- verifiers/v1/interception/server.py | 29 +++++++++++++++++--------- 5 files changed, 23 insertions(+), 58 deletions(-) diff --git a/verifiers/v1/clients/client.py b/verifiers/v1/clients/client.py index 9b218ba936..5b4bab1a0b 100644 --- a/verifiers/v1/clients/client.py +++ b/verifiers/v1/clients/client.py @@ -66,12 +66,6 @@ async def relay_aux( as native JSON and return the provider JSON. Only the relay (eval) client supports it.""" raise NotImplementedError(f"{type(self).__name__} does not relay aux routes") - async def models(self, dialect: Dialect) -> dict: - """The upstream `GET /v1/models` listing, cached per client. Agent loops read a - provider context-window extension (e.g. vLLM's `max_model_len`) from it to set their - compaction threshold; the body is relayed as the provider sent it.""" - raise NotImplementedError(f"{type(self).__name__} does not list models") - async def close(self) -> None: pass diff --git a/verifiers/v1/clients/eval.py b/verifiers/v1/clients/eval.py index 7645a6d515..74ac1c8d52 100644 --- a/verifiers/v1/clients/eval.py +++ b/verifiers/v1/clients/eval.py @@ -1,6 +1,5 @@ """The eval client: proxies harness-native request to the provider.""" -import asyncio import re from collections.abc import Mapping @@ -66,8 +65,6 @@ def __init__(self, config: BaseClientConfig) -> None: # the dialect's provider authentication is applied. self.headers = dict(config.headers or {}) self.client = httpx.AsyncClient(timeout=DEFAULT_TIMEOUT, limits=DEFAULT_LIMITS) - self._models: dict | None = None - self._models_lock = asyncio.Lock() async def get_response( self, @@ -162,27 +159,6 @@ async def _request( f"upstream {response.status_code}: {text}", status_code=response.status_code ) - async def models(self, dialect: Dialect) -> dict: - # One fetch serves every rollout on this endpoint (the lock stops a launch stampede). - async with self._models_lock: - if self._models is None: - headers = httpx.Headers(self.headers) - headers.update(dialect.auth_headers(self.api_key)) - try: - response = await self.client.get( - join_url(self.base_url, "/v1/models"), headers=headers - ) - response.raise_for_status() - except httpx.HTTPStatusError as e: - raise model_error( - f"upstream {e.response.status_code}: {e.response.text}", - status_code=e.response.status_code, - ) from e - except httpx.HTTPError as e: - raise model_error(str(e), status_code=503) from e - self._models = from_json(response.content) - return self._models - async def relay( self, dialect: Dialect, diff --git a/verifiers/v1/clients/train.py b/verifiers/v1/clients/train.py index b7a6ecb757..6ffca79f70 100644 --- a/verifiers/v1/clients/train.py +++ b/verifiers/v1/clients/train.py @@ -7,7 +7,7 @@ from collections.abc import AsyncIterator, Callable, Mapping from contextlib import asynccontextmanager from dataclasses import dataclass, field -from typing import Any, ClassVar, TypeVar, cast +from typing import Any, ClassVar, TypeVar from openai import OpenAIError from renderers import OverlongPromptError, RenderedTokens, Renderer, RendererConfig @@ -312,8 +312,6 @@ class TrainClient(Client): def __init__(self, config: TrainClientConfig) -> None: self.config = config self.client = build_async_openai(config) - self._models: dict | None = None - self._models_lock = asyncio.Lock() # The per-request model is only known at call time; a config that pins the renderer # model can warm now, which is every training run (prime-rl always pins it). if config.renderer_model_name is not None: @@ -323,18 +321,6 @@ def __init__(self, config: TrainClientConfig) -> None: multiplex=config.multiplex, ).warm() - async def models(self, dialect: Dialect) -> dict: - # One fetch serves every rollout on this engine (the lock stops a launch stampede). - async with self._models_lock: - if self._models is None: - try: - self._models = await self.client.get( - "/models", cast_to=cast(Any, dict[str, Any]) - ) - except OpenAIError as e: - raise model_error(e) from e - return self._models - async def get_response( self, dialect: Dialect, diff --git a/verifiers/v1/harnesses/bash/harness.py b/verifiers/v1/harnesses/bash/harness.py index cc5a876742..4100dcec47 100644 --- a/verifiers/v1/harnesses/bash/harness.py +++ b/verifiers/v1/harnesses/bash/harness.py @@ -39,9 +39,6 @@ class CompactionConfig(BaseConfig): class BashHarnessConfig(HarnessConfig): - compaction: CompactionConfig | None = None - """Context compaction policy. Set an empty config to use automatic thresholds.""" - edit: bool = True """Offer the local `edit` tool (single-occurrence string replacement in a file) alongside `bash`. On by default; set `--env.agent.harness.edit false` for a bash-only agent.""" @@ -51,6 +48,9 @@ class BashHarnessConfig(HarnessConfig): eval environment; the key is handed to the program over argv (like the interception secret) so the agent's `bash` subprocesses don't inherit it.""" + compaction: CompactionConfig | None = None + """Context compaction policy. Set an empty config to use automatic thresholds.""" + class BashHarness(Harness[BashHarnessConfig]): APPENDS_SYSTEM_PROMPT = True diff --git a/verifiers/v1/interception/server.py b/verifiers/v1/interception/server.py index d7b8698910..8731cbdadf 100644 --- a/verifiers/v1/interception/server.py +++ b/verifiers/v1/interception/server.py @@ -33,13 +33,15 @@ from tempfile import SpooledTemporaryFile from typing import Literal +import httpx from aiohttp import web from pydantic import ValidationError from pydantic_core import PydanticSerializationError, from_json, to_json from verifiers.v1 import graph from verifiers.v1.clients import Client, resolve_client -from verifiers.v1.configs.client import BaseClientConfig +from verifiers.v1.clients.base import DEFAULT_TIMEOUT, join_url +from verifiers.v1.configs.client import BaseClientConfig, resolve_api_key from verifiers.v1.dialects import DIALECTS, Dialect from verifiers.v1.dialects.base import ( PROVIDER_CAPABILITY_POLICY_CODE, @@ -1060,7 +1062,8 @@ async def handle_models(self, request: web.Request) -> web.Response: """`GET /v1/models`: relay the upstream model listing so agent loops can read a provider context-window extension (e.g. vLLM's `max_model_len`). The path is shared by every dialect; only the auth carrier differs, so the bearer is tried per dialect. - Never recorded on the trace, and a failure never fails the rollout.""" + A pure relay from the session's endpoint config — never recorded on the trace, and a + failure never fails the rollout.""" for dialect in DIALECTS: session = self.sessions.get(dialect.secret(request.headers)) if session is not None: @@ -1069,16 +1072,22 @@ async def handle_models(self, request: web.Request) -> web.Response: return web.json_response({"error": "unauthorized"}, status=401) session.adopt(asyncio.current_task()) logger.debug("intercept models: id=%s", session.trace.id) + config = session.ctx.client + headers = dict(config.headers or {}) + headers.update(dialect.auth_headers(resolve_api_key(config))) try: - payload = await session.client.models(dialect) - except NotImplementedError as e: - return web.json_response(dialect.error_body(str(e)), status=404) - except RolloutError as e: + async with httpx.AsyncClient(timeout=DEFAULT_TIMEOUT) as client: + upstream = await client.get( + join_url(config.base_url, "/v1/models"), headers=headers + ) + except httpx.HTTPError as e: logger.warning("models call failed: id=%s %s", session.trace.id, e) - return web.json_response( - dialect.error_body(str(e)), status=getattr(e, "status_code", 502) - ) - return web.json_response(payload) + return web.json_response(dialect.error_body(str(e)), status=502) + return web.Response( + body=upstream.content, + status=upstream.status_code, + content_type="application/json", + ) def _session_for( self, request: web.Request, *, allow_service: bool = False From dffe17d3482687440166047aa74d75c8b261cc6c Mon Sep 17 00:00:00 2001 From: Mika Senghaas Date: Thu, 27 Aug 2026 03:13:14 +0000 Subject: [PATCH 12/14] resample checkpoint summaries The model can ignore tool_choice="none" on the checkpoint turn and reply with a tool call and no text (observed on ~6% of compactions with deepseek-v4-flash on Terminal-Bench 2) - the rebuilt branch then starts with no context on the work done so far. Resample the checkpoint up to three times until it yields a text summary; pin nano-rlm ac8fdb0 with the same fix. Co-Authored-By: Claude Fable 5 --- verifiers/v1/harnesses/bash/program.py | 24 ++++++++++++++++-------- verifiers/v1/harnesses/rlm/harness.py | 2 +- 2 files changed, 17 insertions(+), 9 deletions(-) diff --git a/verifiers/v1/harnesses/bash/program.py b/verifiers/v1/harnesses/bash/program.py index e66078dd89..09fa54442b 100644 --- a/verifiers/v1/harnesses/bash/program.py +++ b/verifiers/v1/harnesses/bash/program.py @@ -21,6 +21,7 @@ MCP_CALL_ATTEMPTS = 6 MCP_TIMEOUT = 600.0 +CHECKPOINT_ATTEMPTS = 3 CHECKPOINT_COMPACTION_PROMPT = """You are performing a CONTEXT CHECKPOINT COMPACTION. Create a handoff summary for another LLM that will resume the task. @@ -337,14 +338,21 @@ async def compact(self, messages: list[dict]) -> list[dict]: {"role": "user", "content": CHECKPOINT_COMPACTION_PROMPT}, ] try: - completion = await chat( - self.client, - self.model, - checkpoint, - self.tools, - tool_choice="none", - ) - summary = completion.choices[0].message.content or "" + # The model can ignore `tool_choice="none"` and reply with a tool call + # and no text; resample so the next branch never starts context-free. + summary = "" + for _ in range(CHECKPOINT_ATTEMPTS): + completion = await chat( + self.client, + self.model, + checkpoint, + self.tools, + tool_choice="none", + ) + message = completion.choices[0].message + if not message.tool_calls and (message.content or "").strip(): + summary = message.content + break framed = POST_COMPACTION_FRAMING + "\n\n" + summary return [*system, {"role": "user", "content": framed}] except BadRequestError as error: diff --git a/verifiers/v1/harnesses/rlm/harness.py b/verifiers/v1/harnesses/rlm/harness.py index 79aabdf546..8edd82ef12 100644 --- a/verifiers/v1/harnesses/rlm/harness.py +++ b/verifiers/v1/harnesses/rlm/harness.py @@ -43,7 +43,7 @@ class CompactionConfig(BaseConfig): class RLMHarnessConfig(HarnessConfig): - version: str = Field(default="da9f4a4", min_length=1) + version: str = Field(default="ac8fdb0", min_length=1) """Git ref (branch, tag, or commit) of nano-rlm to install.""" max_depth: int = 0 """Recursion depth RLM may spawn sub-harnesses to.""" From db96ac11a32ed15a8f7767b129f32cf8dc288dd8 Mon Sep 17 00:00:00 2001 From: Mika Senghaas Date: Thu, 27 Aug 2026 03:21:32 +0000 Subject: [PATCH 13/14] forbid tool calls in the checkpoint prompt The failed checkpoints were not disobedience: the model understood the request and chose to run one more state-gathering tool call before summarizing, which the loop never grants it. Say explicitly that the summary must come from the conversation as it stands; pin nano-rlm b1b4140 with the same wording. Co-Authored-By: Claude Fable 5 --- verifiers/v1/harnesses/bash/program.py | 4 +++- verifiers/v1/harnesses/rlm/harness.py | 2 +- 2 files changed, 4 insertions(+), 2 deletions(-) diff --git a/verifiers/v1/harnesses/bash/program.py b/verifiers/v1/harnesses/bash/program.py index 09fa54442b..1e49eb0be6 100644 --- a/verifiers/v1/harnesses/bash/program.py +++ b/verifiers/v1/harnesses/bash/program.py @@ -31,7 +31,9 @@ - What remains to be done (clear next steps) - Any critical data, examples, or references needed to continue -Be concise, structured, and focused on helping the next LLM seamlessly continue the work.""" +Be concise, structured, and focused on helping the next LLM seamlessly continue the work. + +Reply with the summary as plain text. Do not call any tools - summarize from the conversation as it stands.""" POST_COMPACTION_FRAMING = """Another language model started to solve this problem and produced \ a summary of its thinking process. Use this to build on the work \ diff --git a/verifiers/v1/harnesses/rlm/harness.py b/verifiers/v1/harnesses/rlm/harness.py index 8edd82ef12..037178622f 100644 --- a/verifiers/v1/harnesses/rlm/harness.py +++ b/verifiers/v1/harnesses/rlm/harness.py @@ -43,7 +43,7 @@ class CompactionConfig(BaseConfig): class RLMHarnessConfig(HarnessConfig): - version: str = Field(default="ac8fdb0", min_length=1) + version: str = Field(default="b1b4140", min_length=1) """Git ref (branch, tag, or commit) of nano-rlm to install.""" max_depth: int = 0 """Recursion depth RLM may spawn sub-harnesses to.""" From daec3696455be7039da8b06b7ef7f81047c1a6f5 Mon Sep 17 00:00:00 2001 From: Mika Senghaas Date: Thu, 27 Aug 2026 17:32:41 +0000 Subject: [PATCH 14/14] simplify checkpoint to one attempt One checkpoint call per compaction - the hardened prompt is the guard against tool-call replies. Also name estimated_tokens' input for what it counts. Pin nano-rlm f5c14aa with the same changes. Co-Authored-By: Claude Fable 5 --- verifiers/v1/harnesses/bash/program.py | 29 ++++++++++---------------- verifiers/v1/harnesses/rlm/harness.py | 2 +- 2 files changed, 12 insertions(+), 19 deletions(-) diff --git a/verifiers/v1/harnesses/bash/program.py b/verifiers/v1/harnesses/bash/program.py index 1e49eb0be6..c53a375cf0 100644 --- a/verifiers/v1/harnesses/bash/program.py +++ b/verifiers/v1/harnesses/bash/program.py @@ -21,7 +21,6 @@ MCP_CALL_ATTEMPTS = 6 MCP_TIMEOUT = 600.0 -CHECKPOINT_ATTEMPTS = 3 CHECKPOINT_COMPACTION_PROMPT = """You are performing a CONTEXT CHECKPOINT COMPACTION. Create a handoff summary for another LLM that will resume the task. @@ -286,8 +285,9 @@ def drop_latest_tool_result(messages: list[dict]) -> bool: return False -def estimated_tokens(value: str) -> int: - return (len(value) + 3) // 4 +def estimated_tokens(chars: str) -> int: + """Rough token count at four characters per token.""" + return (len(chars) + 3) // 4 def context_tokens(completion) -> int: @@ -340,21 +340,14 @@ async def compact(self, messages: list[dict]) -> list[dict]: {"role": "user", "content": CHECKPOINT_COMPACTION_PROMPT}, ] try: - # The model can ignore `tool_choice="none"` and reply with a tool call - # and no text; resample so the next branch never starts context-free. - summary = "" - for _ in range(CHECKPOINT_ATTEMPTS): - completion = await chat( - self.client, - self.model, - checkpoint, - self.tools, - tool_choice="none", - ) - message = completion.choices[0].message - if not message.tool_calls and (message.content or "").strip(): - summary = message.content - break + completion = await chat( + self.client, + self.model, + checkpoint, + self.tools, + tool_choice="none", + ) + summary = completion.choices[0].message.content or "" framed = POST_COMPACTION_FRAMING + "\n\n" + summary return [*system, {"role": "user", "content": framed}] except BadRequestError as error: diff --git a/verifiers/v1/harnesses/rlm/harness.py b/verifiers/v1/harnesses/rlm/harness.py index 037178622f..a503828542 100644 --- a/verifiers/v1/harnesses/rlm/harness.py +++ b/verifiers/v1/harnesses/rlm/harness.py @@ -43,7 +43,7 @@ class CompactionConfig(BaseConfig): class RLMHarnessConfig(HarnessConfig): - version: str = Field(default="b1b4140", min_length=1) + version: str = Field(default="f5c14aa", min_length=1) """Git ref (branch, tag, or commit) of nano-rlm to install.""" max_depth: int = 0 """Recursion depth RLM may spawn sub-harnesses to."""