fix(harness): bound MCP bridge shutdown and retain live servers - #1099
burtenshaw merged 37 commits into
Conversation
Moves the trainer-side rollout API out of the package __init__ and into `openenv.core.harness.rollout`, leaving __init__ as a re-export shim. No behavior change: every name previously importable from `openenv.core.harness` still is, and is the same object. The module was ~730 lines living directly in __init__ with a docstring noting it sat outside the stable surface "while RFC 005 is still under review". Splitting it now makes room for the RFC 005 turn-based agentic harness layer to land in sibling modules instead of growing the __init__ further. Also re-exports the private `_resolve_env_reward`, which tests/scripts/test_browsergym_harness_eval_examples.py imports from the package root, and points `collect.py` at `.rollout` directly rather than importing from its own package. Consumers left untouched and verified: `openenv collect`, pi_env, opencode_env, browsergym_env, reasoning_gym_env, openspiel_env. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Adds the type layer for wrapping an external agentic harness (Claude Code, OpenClaw, Codex) as an OpenEnv environment. No runtime behavior yet — this PR is types plus their unit tests. - `config.py`: `HarnessConfig` / `HarnessTransport`. `session_timeout_s` is documented as bounding ONE conversational turn, per the RFC's temporal-semantics section (the field comment in the RFC is ambiguous; flagging for reviewer sign-off). - `events.py`: `HarnessEventType` / `HarnessEvent` / `HarnessResponse`, plus `events_to_metadata()`, the sanctioned JSON-safe path for putting events into `Observation.metadata` so they survive wire serialization. - `adapter.py`: `AgenticHarnessAdapter` ABC and its error hierarchy. - `tools.py`: `resolve_tool_conflicts()` for the RFC's tool-name collision rules (`env_` prefixing, error on ambiguity). Two deliberate deviations from the RFC text, both because the RFC is stale against the code: 1. The RFC's `ToolDefinition` does not exist; the type is `Tool` (`env_server/mcp_types.py`), reused here rather than duplicated. Same for `RESERVED_TOOL_NAMES`, which `resolve_tool_conflicts` re-checks as defense in depth. 2. `send_message()` is concrete rather than abstract. Streaming is the single abstract turn primitive and `send_message()` drains it, which removes duplication from every concrete adapter and makes the terminal TURN_COMPLETE event an enforced contract instead of a convention. The ABC is named `AgenticHarnessAdapter` to avoid colliding with the rollout layer's existing `HarnessAdapter`. Worth discussing whether to rename the rollout classes instead and reclaim the RFC's plain names -- see the PR description. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
…idge Makes the RFC 005 types runnable: an environment that owns a harness subprocess, hands it the environment's MCP tools, and turns each step() into one conversational turn. - `environment.py`: `HarnessAction` + `HarnessEnvironment(MCPEnvironment)`. reset() stops any live harness, enumerates and conflict-resolves the env tools, starts the bridge, injects, then starts the harness -- injection strictly before start, per the RFC. step() runs one turn; MCP actions keep their normal routing. Rubrics run after the turn completes, outside the harness's control loop, preserving RFC 004's reward boundary. - `process.py`: `HarnessProcess`, a loop-agnostic Popen + reader-thread helper (readiness gating, stderr-tail diagnostics, idempotent stop with SIGTERM -> SIGKILL escalation over the process group). - `bridge.py`: `HarnessMCPBridge`, serving the env's FastMCP tool surface over loopback HTTP for the harness to consume. Three decisions worth reviewer attention: 1. `HarnessEnvironment` subclasses `MCPEnvironment` and substitutes an empty internal FastMCP when `mcp=None`. `MCPEnvironment` requires `mcp_server` positionally, so the RFC's optional-mcp constructor cannot be written literally; this keeps reserved-name validation, tool enumeration and mcp_session() integration for free. 2. Popen + threads rather than asyncio subprocess transports, because the same instance must work across event loops: the sync facade spins a fresh loop per call (run_async_safely) while the server keeps one long-lived loop. Asyncio subprocess transports are bound to their creating loop. 3. The bridge is a separate loopback server rather than a reuse of the env server's /mcp endpoint. Reusing /mcp would put the orchestration routes (/reset, /step, /state) on the same origin the harness can reach, violating the RFC's security boundary; it would also hand the harness a *different* env instance, since WS /mcp creates its own session. Keeping it separate makes the boundary structural rather than filter-based. Turn timeouts and harness crashes become terminal observations (done=True, metadata.error_type) rather than exceptions, so a training loop scores the episode and moves on instead of unwinding. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Four findings from the automated review, all confirmed against the code before fixing. 1. Renamed tools were unreachable (High, reported on huggingface#1100). Conflict resolution renames a colliding env tool before injection (read_file -> env_read_file), but the bridge served the source FastMCP unchanged, so the harness was handed a name that did not resolve. Adds `build_bridge_server()`, which serves a renamed view built with FastMCP's own `Tool.from_tool(tool, name=...)`, and returns the source server untouched when there is nothing to rename. The new test fails against the old code with `['add', 'read_file'] != ['add', 'env_read_file']`, which is the bug exactly. 2. Reset skipped adapter cleanup (Medium). `reset_async` only stopped the adapter when `is_alive()` was true, but a harness that died on its own reports False while still holding an unwaited process, open pipes and live reader threads; the next `start()` then overwrote that state and leaked it. `stop()` is contractually idempotent, so it is now called unconditionally, and also on the `start()` failure path. 3. Subprocess I/O lacked an explicit encoding (Medium). `text=True` alone decodes with the locale encoding, which is frequently ASCII in a container, while harness output is routinely not. Worse, `UnicodeDecodeError` is a `ValueError`, which the reader thread caught and exited on -- so one non-ASCII byte silently stopped stdout pumping and the turn hung until its timeout. Now `encoding="utf-8"` with `errors="replace"`, and the reader's handler is narrowed to the pipe-closed case it was meant for. 4. Timeouts were reported as crashes (Low). `HarnessTurnTimeoutError` subclasses `HarnessError`, so an adapter raising the dedicated timeout exception was labelled `harness_crashed`. It is now caught first and mapped to `turn_timeout`. Also from finding 1's root cause: mode-specific tools registered with `tool(mode=...)` live on the environment, not the FastMCP server, so the bridge cannot serve them either. They are now excluded from injection with a warning rather than advertised and then failing on call. Supporting them needs a decision about what a mode means inside a harness turn, which is left as a follow-up. Note this only ever manifested when the env's `_mode` matched the tool's -- the test sets it explicitly, since otherwise the assertion would pass vacuously. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Two findings from the automated review on huggingface#1100. The third (renamed tools unreachable through the bridge) was the same root cause as a finding on huggingface#1099 and is fixed there. 1. Production turns ignored session_timeout_s (Medium). The /harness handler streamed `send_message_streaming` with no bound, while simulation mode wraps the same call in `asyncio.wait_for` inside `HarnessEnvironment._run_turn`. A hung harness therefore held its session open indefinitely, and since HarnessEnvironment is SUPPORTS_CONCURRENT_SESSIONS=False with the idle reaper off by default, the server stayed pinned at capacity. The turn is now bounded by the adapter's `session_timeout_s`, matching simulation semantics. 2. A stream that ended without TURN_COMPLETE hung the client (Medium). `send_message()` raises HarnessError in that case, but the socket loop silently went back to waiting for the next client frame, so a client blocking on the terminal event waited forever. The handler now detects it and emits a terminal ERROR event before ending the session. Both paths end the session rather than continuing, so a reconnect gets a fresh harness -- consistent with how a mid-stream crash was already handled. Adds `send_harness_error()` since three paths now emit the same terminal ERROR frame. Both new tests assert `server.active_sessions == 0` afterwards, which is the property finding 1 was really about. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
caiotheodoro
left a comment
There was a problem hiding this comment.
src/openenv/core/harness/process.py:221 — bug: read_line(timeout_s=None) doesn't do what the docstring promises ("None blocks until a line or EOF"). The reader thread's for line in stream: sink(line) (line 253) just returns when the pipe closes — no sentinel goes into the queue — so self._stdout_queue.get(timeout=None) (line 236) blocks forever past real EOF.
Repro: start a process that prints one line and exits, drain that line, then call read_line(timeout_s=None) again.
first read_line -> hello
is_running() -> False
calling read_line(timeout_s=None) on an exited process...
BUG CONFIRMED: read_line(timeout_s=None) hung for 5s past EOF instead of returning None
STILL RUNNING AFTER 12s
Wrapping the call in asyncio.wait_for lets the caller move on, but the to_thread worker stays parked in queue.get() — it isn't one of the daemon _reader_threads that stop() joins, so it outlives stop() and (being a non-daemon default-executor thread) blocks the interpreter from exiting.
Nothing in this PR calls read_line with timeout_s=None yet, but it's public API on a class future harness adapters will drive, and the contract as documented is false. Have the reader push a sentinel (or close the queue) on stream EOF and have read_line return None for it instead of relying on a bare queue.get().
|
The docs for this PR live here. All of your documentation changes will be reflected on that endpoint. The docs are available until 30 days after the last update. |
|
@caiotheodoro, your EOF review is addressed by 1a6caa57, which is already included in this PR. The stdout reader now enqueues an EOF sentinel from its Real-subprocess regression tests cover:
These cases passed in the latest 79-test harness run. The full repository test run also passed: 2,552 passed, 120 skipped. |
|
@caiotheodoro can you please review again? |
…1100) * refactor(harness): split openenv.core.harness into a package Moves the trainer-side rollout API out of the package __init__ and into `openenv.core.harness.rollout`, leaving __init__ as a re-export shim. No behavior change: every name previously importable from `openenv.core.harness` still is, and is the same object. The module was ~730 lines living directly in __init__ with a docstring noting it sat outside the stable surface "while RFC 005 is still under review". Splitting it now makes room for the RFC 005 turn-based agentic harness layer to land in sibling modules instead of growing the __init__ further. Also re-exports the private `_resolve_env_reward`, which tests/scripts/test_browsergym_harness_eval_examples.py imports from the package root, and points `collect.py` at `.rollout` directly rather than importing from its own package. Consumers left untouched and verified: `openenv collect`, pi_env, opencode_env, browsergym_env, reasoning_gym_env, openspiel_env. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> * feat(harness): RFC 005 foundation types for agentic harnesses Adds the type layer for wrapping an external agentic harness (Claude Code, OpenClaw, Codex) as an OpenEnv environment. No runtime behavior yet — this PR is types plus their unit tests. - `config.py`: `HarnessConfig` / `HarnessTransport`. `session_timeout_s` is documented as bounding ONE conversational turn, per the RFC's temporal-semantics section (the field comment in the RFC is ambiguous; flagging for reviewer sign-off). - `events.py`: `HarnessEventType` / `HarnessEvent` / `HarnessResponse`, plus `events_to_metadata()`, the sanctioned JSON-safe path for putting events into `Observation.metadata` so they survive wire serialization. - `adapter.py`: `AgenticHarnessAdapter` ABC and its error hierarchy. - `tools.py`: `resolve_tool_conflicts()` for the RFC's tool-name collision rules (`env_` prefixing, error on ambiguity). Two deliberate deviations from the RFC text, both because the RFC is stale against the code: 1. The RFC's `ToolDefinition` does not exist; the type is `Tool` (`env_server/mcp_types.py`), reused here rather than duplicated. Same for `RESERVED_TOOL_NAMES`, which `resolve_tool_conflicts` re-checks as defense in depth. 2. `send_message()` is concrete rather than abstract. Streaming is the single abstract turn primitive and `send_message()` drains it, which removes duplication from every concrete adapter and makes the terminal TURN_COMPLETE event an enforced contract instead of a convention. The ABC is named `AgenticHarnessAdapter` to avoid colliding with the rollout layer's existing `HarnessAdapter`. Worth discussing whether to rename the rollout classes instead and reclaim the RFC's plain names -- see the PR description. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> * feat(harness): HarnessEnvironment, subprocess helper, and MCP tool bridge Makes the RFC 005 types runnable: an environment that owns a harness subprocess, hands it the environment's MCP tools, and turns each step() into one conversational turn. - `environment.py`: `HarnessAction` + `HarnessEnvironment(MCPEnvironment)`. reset() stops any live harness, enumerates and conflict-resolves the env tools, starts the bridge, injects, then starts the harness -- injection strictly before start, per the RFC. step() runs one turn; MCP actions keep their normal routing. Rubrics run after the turn completes, outside the harness's control loop, preserving RFC 004's reward boundary. - `process.py`: `HarnessProcess`, a loop-agnostic Popen + reader-thread helper (readiness gating, stderr-tail diagnostics, idempotent stop with SIGTERM -> SIGKILL escalation over the process group). - `bridge.py`: `HarnessMCPBridge`, serving the env's FastMCP tool surface over loopback HTTP for the harness to consume. Three decisions worth reviewer attention: 1. `HarnessEnvironment` subclasses `MCPEnvironment` and substitutes an empty internal FastMCP when `mcp=None`. `MCPEnvironment` requires `mcp_server` positionally, so the RFC's optional-mcp constructor cannot be written literally; this keeps reserved-name validation, tool enumeration and mcp_session() integration for free. 2. Popen + threads rather than asyncio subprocess transports, because the same instance must work across event loops: the sync facade spins a fresh loop per call (run_async_safely) while the server keeps one long-lived loop. Asyncio subprocess transports are bound to their creating loop. 3. The bridge is a separate loopback server rather than a reuse of the env server's /mcp endpoint. Reusing /mcp would put the orchestration routes (/reset, /step, /state) on the same origin the harness can reach, violating the RFC's security boundary; it would also hand the harness a *different* env instance, since WS /mcp creates its own session. Keeping it separate makes the boundary structural rather than filter-based. Turn timeouts and harness crashes become terminal observations (done=True, metadata.error_type) rather than exceptions, so a training loop scores the episode and moves on instead of unwinding. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> * feat(server): production /harness WebSocket route and mode wiring Exposes a harness environment directly to clients in production mode, and gives deployments a way to actually select that mode. `/harness` is registered only when mode is PRODUCTION *and* the env factory produces a `HarnessEnvironment`. Connecting opens a session and resets the env (starting the harness and injecting tools); each `{"type": "message", "content": ...}` frame runs one turn, streamed back as HarnessEvent frames terminated by turn_complete. Malformed frames get a WSErrorResponse without dropping the connection; an adapter crash streams a terminal error event and ends the session, since harness state after a crash is undefined. The handler is modelled on the existing WS /mcp handler rather than the RFC's pseudocode: it goes through `_create_session()` so capacity limits, the AsyncExitStack, per-session executors and the idle reaper all apply. The RFC sketch calls the factory directly and would bypass all of it. Session activity is touched per streamed event so a long turn is not reaped as idle. Mode wiring: `create_app` / `create_fastapi_app` (and the web-interface factory) take a keyword-only `mode`, resolved from `OPENENV_MODE` when omitted. Previously `create_fastapi_app` hardcoded `register_routes(app)`, so nothing outside tests could ever select production. Default is unchanged (simulation), mirroring the existing `OPENENV_CLIENT_MODE` convention on the client side. Two notes for review: - Harness-env detection uses a lazy import of `HarnessEnvironment` inside the method: `openenv.core.harness` imports env_server modules, so a top-level import would be circular. Detection may instantiate the factory once to probe it, which is safe because constructing a HarnessEnvironment starts nothing -- asserted by a test. - `websocket.close()` on an already-gone client raises RuntimeError under TestClient but WebSocketDisconnect under a real ASGI server; both are now caught. Found by the end-to-end test, which drives a real uvicorn server. The pre-existing /mcp and /ws handlers have the same latent gap and are deliberately left untouched here. Includes the end-to-end test for the whole stack: a real subprocess harness that reads the injected MCP config, calls an env tool over the live bridge, and streams turns -- exercised through both the simulation step() API and a real uvicorn server's /harness socket. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> * fix(harness): address Bugbot review on the environment runtime Four findings from the automated review, all confirmed against the code before fixing. 1. Renamed tools were unreachable (High, reported on #1100). Conflict resolution renames a colliding env tool before injection (read_file -> env_read_file), but the bridge served the source FastMCP unchanged, so the harness was handed a name that did not resolve. Adds `build_bridge_server()`, which serves a renamed view built with FastMCP's own `Tool.from_tool(tool, name=...)`, and returns the source server untouched when there is nothing to rename. The new test fails against the old code with `['add', 'read_file'] != ['add', 'env_read_file']`, which is the bug exactly. 2. Reset skipped adapter cleanup (Medium). `reset_async` only stopped the adapter when `is_alive()` was true, but a harness that died on its own reports False while still holding an unwaited process, open pipes and live reader threads; the next `start()` then overwrote that state and leaked it. `stop()` is contractually idempotent, so it is now called unconditionally, and also on the `start()` failure path. 3. Subprocess I/O lacked an explicit encoding (Medium). `text=True` alone decodes with the locale encoding, which is frequently ASCII in a container, while harness output is routinely not. Worse, `UnicodeDecodeError` is a `ValueError`, which the reader thread caught and exited on -- so one non-ASCII byte silently stopped stdout pumping and the turn hung until its timeout. Now `encoding="utf-8"` with `errors="replace"`, and the reader's handler is narrowed to the pipe-closed case it was meant for. 4. Timeouts were reported as crashes (Low). `HarnessTurnTimeoutError` subclasses `HarnessError`, so an adapter raising the dedicated timeout exception was labelled `harness_crashed`. It is now caught first and mapped to `turn_timeout`. Also from finding 1's root cause: mode-specific tools registered with `tool(mode=...)` live on the environment, not the FastMCP server, so the bridge cannot serve them either. They are now excluded from injection with a warning rather than advertised and then failing on call. Supporting them needs a decision about what a mode means inside a harness turn, which is left as a follow-up. Note this only ever manifested when the env's `_mode` matched the tool's -- the test sets it explicitly, since otherwise the assertion would pass vacuously. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> * fix(server): bound and terminate every production /harness turn Two findings from the automated review on #1100. The third (renamed tools unreachable through the bridge) was the same root cause as a finding on #1099 and is fixed there. 1. Production turns ignored session_timeout_s (Medium). The /harness handler streamed `send_message_streaming` with no bound, while simulation mode wraps the same call in `asyncio.wait_for` inside `HarnessEnvironment._run_turn`. A hung harness therefore held its session open indefinitely, and since HarnessEnvironment is SUPPORTS_CONCURRENT_SESSIONS=False with the idle reaper off by default, the server stayed pinned at capacity. The turn is now bounded by the adapter's `session_timeout_s`, matching simulation semantics. 2. A stream that ended without TURN_COMPLETE hung the client (Medium). `send_message()` raises HarnessError in that case, but the socket loop silently went back to waiting for the next client frame, so a client blocking on the terminal event waited forever. The handler now detects it and emits a terminal ERROR event before ending the session. Both paths end the session rather than continuing, so a reconnect gets a fresh harness -- consistent with how a mid-stream crash was already handled. Adds `send_harness_error()` since three paths now emit the same terminal ERROR frame. Both new tests assert `server.active_sessions == 0` afterwards, which is the property finding 1 was really about. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> * fix: wake harness readers on eof * fix: preserve live harness sessions * fix(harness): apply transforms and clean up exited processes * fix(harness): clean up cancelled reset and startup * fix: reject production turns when harness process has died * fix(harness): distinguish recoverable protocol errors * fix(harness): clean up cancelled conversational turns * fix(harness): clean up when rubric reset fails * fix: harden harness sessions * fix: finish harness cleanup --------- Co-authored-by: Claude Opus 5 <noreply@anthropic.com> Co-authored-by: burtenshaw <ben.burtenshaw@gmail.com>
There was a problem hiding this comment.
Cursor Bugbot has reviewed your changes using default effort and found 1 potential issue.
❌ Bugbot Autofix is OFF. To automatically fix reported issues with cloud agents, enable autofix in the Cursor dashboard.
Reviewed by Cursor Bugbot for commit 24e345a. Configure here.

An active HTTP stream could outlast
HarnessMCPBridge.stop()and leave a live server thread untracked. Shutdown now bounds Uvicorn's graceful wait, requests forced exit if needed, and only clears the server reference after the thread exits. If the deadline expires,stop()raisesHarnessError, retains ownership for a retry, andstart()rejects the still-stopping server.The RFC 005 runtime is already on
mainthrough #1100. This PR preserves those newer runtime fixes and now adds only the bridge shutdown fix and regression tests.Test plan
smolagentsis unavailable. Repository-wide formatting reports 56 existing files; these are unchanged by this fix.Note
Medium Risk
Changes harness MCP bridge and episode reset/close lifecycle; incorrect behavior could leak threads or block resets, but scope is localized and covered by new regression tests.
Overview
MCP bridge shutdown is now time-bounded and fails loudly when the Uvicorn thread does not exit.
HarnessMCPBridge.stop()caps graceful shutdown, escalates toforce_exit, and raisesHarnessErroron timeout while keeping the server/thread handles so callers can retrystop();start()refuses to spin up a replacement while a stop is still in flight.HarnessEnvironment aligns with that contract:
reset_asyncstops the existing bridge synchronously (propagating shutdown errors), skips episode cleanup that would spawn a second bridge after a failed stop, andclose()only drops_bridgeafter a successful stop so a laterclose()can retry.Regression tests cover active HTTP streams during shutdown, retained servers after deadline expiry, and reset/close behavior when bridge stop fails or is cancelled.
Reviewed by Cursor Bugbot for commit 6e5c36c. Bugbot is set up for automated code reviews on this repo. Configure here.