Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
49 commits
Select commit Hold shift + click to select a range
3c99e0f
refactor(harness): split openenv.core.harness into a package
splusq Aug 28, 2026
01bb0ce
feat(harness): RFC 005 foundation types for agentic harnesses
splusq Aug 28, 2026
f3d90c2
feat(harness): HarnessEnvironment, subprocess helper, and MCP tool br…
splusq Aug 28, 2026
385fded
feat(server): production /harness WebSocket route and mode wiring
splusq Aug 28, 2026
f39c662
Merge branch 'main' into rfc-005/pr1-harness-package-split
splusq Aug 31, 2026
b6e6984
Merge branch 'main' into rfc-005/pr2-harness-foundation-types
splusq Aug 31, 2026
025191d
Merge branch 'main' into rfc-005/pr3-harness-environment-runtime
splusq Aug 31, 2026
58dd957
Merge branch 'main' into rfc-005/pr4-harness-production-route
splusq Aug 31, 2026
671f5c0
Merge branch 'main' into rfc-005/pr1-harness-package-split
splusq Sep 2, 2026
933e113
Merge branch 'main' into rfc-005/pr2-harness-foundation-types
splusq Sep 2, 2026
668ab2e
fix(harness): address Bugbot review on the environment runtime
splusq Sep 2, 2026
ec68a5d
Merge branch 'rfc-005/pr3-harness-environment-runtime' into rfc-005/p…
splusq Sep 2, 2026
dc3cf59
fix(server): bound and terminate every production /harness turn
splusq Sep 2, 2026
7cda7a3
fix: refresh harness package split
burtenshaw Sep 17, 2026
c9e4316
fix: refresh harness foundation types
burtenshaw Sep 17, 2026
cf760ba
fix: refresh harness runtime
burtenshaw Sep 17, 2026
38d92c3
fix: refresh harness production route
burtenshaw Sep 17, 2026
1a6caa5
fix: wake harness readers on eof
burtenshaw Sep 17, 2026
545ba9d
fix: propagate harness eof cleanup
burtenshaw Sep 17, 2026
01b09e5
fix: refresh latest core state changes
burtenshaw Sep 17, 2026
5c11106
fix: refresh foundation base
burtenshaw Sep 17, 2026
8736543
fix: refresh runtime base
burtenshaw Sep 17, 2026
ebb4491
fix: preserve state schemas and harness mode
burtenshaw Sep 17, 2026
40fbe50
fix: preserve live harness sessions
burtenshaw Sep 17, 2026
a81f0a2
chore: refresh latest main
burtenshaw Sep 17, 2026
4de0d41
chore: refresh foundation base
burtenshaw Sep 17, 2026
7401fe0
chore: refresh runtime base
burtenshaw Sep 17, 2026
1111156
chore: refresh production base
burtenshaw Sep 17, 2026
2e622ad
Merge remote-tracking branch 'origin/main' into rfc-005/pr3-harness-e…
splusq Sep 17, 2026
4ad1fb2
fix(harness): apply transforms and clean up exited processes
splusq Sep 17, 2026
e317462
fix(harness): clean up cancelled reset and startup
splusq Sep 17, 2026
26a5043
Merge branch 'main' into rfc-005/pr3-harness-environment-runtime
splusq Sep 23, 2026
d42d67f
Merge branch 'main' into rfc-005/pr4-harness-production-route
splusq Sep 23, 2026
7715b12
fix: reject production turns when harness process has died
splusq Sep 23, 2026
a2bf6f6
Merge branch 'main' into rfc-005/pr3-harness-environment-runtime
splusq Sep 24, 2026
35df865
Merge branch 'main' into rfc-005/pr4-harness-production-route
splusq Sep 24, 2026
7225191
Merge main into RFC 005 harness runtime; resolve package split conflicts
splusq Sep 25, 2026
11f0de3
Merge main into RFC 005 production harness branch
splusq Sep 25, 2026
15f5376
fix(harness): distinguish recoverable protocol errors
splusq Sep 25, 2026
0cc7d91
fix(harness): clean up cancelled conversational turns
splusq Sep 25, 2026
87a7772
Merge main after RFC 005 foundation types; preserve runtime exports
splusq Sep 28, 2026
8d4bcf8
Merge main after RFC 005 foundation types landed
splusq Sep 28, 2026
bd7f0f2
fix(harness): clean up when rubric reset fails
splusq Sep 28, 2026
0d87b5c
Merge remote-tracking branch 'origin/main' into codex/resolve-pr-1100
splusq Sep 28, 2026
d945683
Merge branch 'main' into rfc-005/pr3-harness-environment-runtime
splusq Sep 28, 2026
079a6d5
Merge branch 'main' into rfc-005/pr3-harness-environment-runtime
burtenshaw Sep 29, 2026
388e2a1
fix: sync harness runtime
burtenshaw Sep 29, 2026
e4fd153
fix: harden harness sessions
burtenshaw Sep 29, 2026
eecdd9f
fix: finish harness cleanup
burtenshaw Sep 29, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
227 changes: 226 additions & 1 deletion src/openenv/core/env_server/http_server.py
Original file line number Diff line number Diff line change
Expand Up @@ -640,6 +640,212 @@ def concurrency_config(self) -> ConcurrencyConfig:
"""Return the concurrency configuration."""
return self._concurrency_config

def _factory_produces_harness_env(self) -> bool:
"""Return whether the env factory produces a HarnessEnvironment."""
import inspect

# Lazy import: openenv.core.harness imports env_server modules, so a
# top-level import here would be circular.
from ..harness.environment import HarnessEnvironment

if inspect.isclass(self._env_factory):
return issubclass(self._env_factory, HarnessEnvironment)
_temp_env = self._env_factory()
try:
return isinstance(_temp_env, HarnessEnvironment)
finally:
_temp_env.close()

def _register_harness_route(self, app: FastAPI) -> None:
"""
Register the production `/harness` WebSocket route (RFC 005).

Each connection gets its own environment session: connecting resets
the environment (which starts the harness process and injects tools),
and each `{"type": "message", "content": ...}` frame runs one
conversational turn, streamed back as `HarnessEvent` JSON frames
ending with a `turn_complete` event. Malformed client frames receive a
recoverable `protocol_error` response without starting a turn; the
connection remains usable. Terminal failures use `error`.
"""
# Lazy import to avoid a circular import with openenv.core.harness.
from ..harness.adapter import HarnessNotRunningError
from ..harness.events import (
HarnessClientMessage,
HarnessEvent,
HarnessEventType,
HarnessProtocolError,
)

@app.websocket("/harness")
async def harness_websocket_endpoint(websocket: WebSocket):
await websocket.accept()

session_id = None
session_env = None

async def send_error(message: str, code: WSErrorCode) -> None:
error_response = WSErrorResponse(
data={"message": message, "code": code}
)
await websocket.send_text(error_response.model_dump_json())

async def send_protocol_error(message: str, code: WSErrorCode) -> None:
error_response = HarnessProtocolError(
data={"message": message, "code": code}
)
await websocket.send_text(error_response.model_dump_json())

async def send_harness_error(message: str) -> None:
"""Emit a terminal ERROR event in the harness event stream."""
error_event = HarnessEvent(
type=HarnessEventType.ERROR,
data={"message": message, "recoverable": False},
)
await websocket.send_text(error_event.model_dump_json())
Comment thread
splusq marked this conversation as resolved.

try:
session_id, session_env = await self._create_session()
# Protect the live harness from idle reaping and HTTP session
# close, including startup and turns that emit no events.
self._session_websocket_attachments.add(session_id)

async with AsyncExitStack() as stack:
mcp_session_factory = getattr(session_env, "mcp_session", None)
if callable(mcp_session_factory):
mcp_session_cm = cast(
AsyncContextManager[Any], mcp_session_factory()
)
await stack.enter_async_context(mcp_session_cm)

# Starts the harness process and injects environment tools
await session_env.reset_async()
await websocket.send_text(
json.dumps(
{
"type": "session_started",
"data": {
"session_id": session_id,
"harness": session_env.adapter.config.name,
},
}
)
)

while True:
raw_message = await websocket.receive_text()

try:
message_dict = json.loads(raw_message)
except json.JSONDecodeError as e:
await send_protocol_error(
f"Invalid JSON: {e}", WSErrorCode.INVALID_JSON
)
continue
try:
client_message = HarnessClientMessage(**message_dict)
except (ValidationError, TypeError) as e:
await send_protocol_error(
f"Invalid message: {e}",
WSErrorCode.VALIDATION_ERROR,
)
continue

self._update_session_activity(session_id, increment_step=True)

async def stream_turn(content: str) -> bool:
"""Stream one turn; True if it ended with TURN_COMPLETE."""
saw_terminal = False
adapter = session_env.adapter
if not await adapter.is_alive():
raise HarnessNotRunningError(
"harness process is not running"
)
async for event in adapter.send_message_streaming(content):
await websocket.send_text(event.model_dump_json())
# Record progress throughout the turn.
self._update_session_activity(session_id)
Comment thread
cursor[bot] marked this conversation as resolved.
saw_terminal = (
event.type is HarnessEventType.TURN_COMPLETE
)
return saw_terminal

# Bound the turn in wall-clock time, matching what
# simulation mode does in HarnessEnvironment._run_turn.
# Without this a hung harness holds the session open
# forever, and the server sits at capacity.
turn_timeout_s = session_env.adapter.config.session_timeout_s
try:
completed = await asyncio.wait_for(
stream_turn(client_message.content),
turn_timeout_s,
)
Comment thread
splusq marked this conversation as resolved.
except asyncio.TimeoutError:
await send_harness_error(
f"harness turn exceeded {turn_timeout_s} seconds"
)
break
except HarnessNotRunningError:
await send_harness_error("harness process is not running")
break
except Exception:
# Harness state after a crash is undefined; end
# the session so a reconnect gets a fresh one.
# Adapter exceptions can contain credentials or
# subprocess output; do not expose them to clients.
await send_harness_error("harness turn failed")
break

if not completed:
# send_message() raises HarnessError here; the
# socket equivalent is to say so and end the
# session, rather than leaving a client that
# blocks on the terminal event waiting forever.
await send_harness_error(
"harness event stream ended without a "
"TURN_COMPLETE event"
)
break
Comment thread
cursor[bot] marked this conversation as resolved.

except WebSocketDisconnect:
pass
except SessionCapacityError as e:
await send_error(str(e), WSErrorCode.CAPACITY_REACHED)
except EnvironmentFactoryError as e:
await send_error(str(e), WSErrorCode.FACTORY_ERROR)
except Exception:
try:
await send_error(
"harness session failed", WSErrorCode.SESSION_ERROR
)
except (RuntimeError, WebSocketDisconnect):
pass
finally:
if session_id:
# Release ownership without an await so cancellation cannot
# leave a session permanently exempt from idle reaping.
self._session_websocket_attachments.discard(session_id)
cleanup = asyncio.create_task(self._destroy_session(session_id))
try:
await asyncio.shield(cleanup)
except asyncio.CancelledError:
# ASGI cancellation must not orphan a running harness.
# Finish teardown before propagating cancellation, even
# when the request's cancel scope cancels us repeatedly.
while not cleanup.done():
try:
await asyncio.shield(cleanup)
except asyncio.CancelledError:
pass
cleanup.result()
raise
try:
await websocket.close()
except (RuntimeError, WebSocketDisconnect):
# TestClient raises RuntimeError, real ASGI servers raise
# WebSocketDisconnect when the client is already gone.
pass

def register_routes(
self, app: FastAPI, mode: ServerMode | str = ServerMode.SIMULATION
) -> None:
Expand Down Expand Up @@ -1290,6 +1496,11 @@ async def mcp_websocket_endpoint(websocket: WebSocket):
except (RuntimeError, WebSocketDisconnect):
pass

# In production mode, a harness environment is exposed directly to
# clients via a streaming WebSocket (RFC 005).
if mode == ServerMode.PRODUCTION and self._factory_produces_harness_env():
self._register_harness_route(app)

# Register simulation control routes only in simulation mode
if mode == ServerMode.SIMULATION:

Expand Down Expand Up @@ -1781,6 +1992,8 @@ def create_app(
show_default_tab: bool = True,
title_override: Optional[str] = None,
state_cls: Type[State] = State,
*,
mode: Optional[ServerMode | str] = None,
) -> FastAPI:
"""
Create a FastAPI application with or without web interface.
Expand Down Expand Up @@ -1822,6 +2035,9 @@ def create_app(
state_cls (`Type[State]`, *optional*, defaults to `State`):
The `State` subclass this environment reports, used for the `/state`
response model and the `state` entry of `/schema`.
mode (`ServerMode` or `str`, *optional*):
Server mode. When `None`, resolved from the `OPENENV_MODE`
environment variable, defaulting to simulation.

Returns:
`FastAPI` application instance with or without web interface and README integration.
Expand Down Expand Up @@ -1851,6 +2067,7 @@ def create_app(
custom_tab_primary=custom_tab_primary,
show_default_tab=show_default_tab,
title_override=title_override,
mode=mode,
)
else:
# Use standard FastAPI app without web interface
Expand All @@ -1862,6 +2079,7 @@ def create_app(
concurrency_config,
env_name=env_name,
state_cls=state_cls,
mode=mode,
)


Expand All @@ -1873,6 +2091,8 @@ def create_fastapi_app(
concurrency_config: Optional[ConcurrencyConfig] = None,
env_name: Optional[str] = None,
state_cls: Type[State] = State,
*,
mode: Optional[ServerMode | str] = None,
) -> FastAPI:
"""
Create a FastAPI application with comprehensive documentation.
Expand All @@ -1895,6 +2115,9 @@ def create_fastapi_app(
state_cls (`Type[State]`, *optional*, defaults to `State`):
The `State` subclass this environment reports, used for the `/state`
response model and the `state` entry of `/schema`.
mode (`ServerMode` or `str`, *optional*):
Server mode. When `None`, resolved from the `OPENENV_MODE`
environment variable, defaulting to simulation.

Returns:
`FastAPI` application instance.
Expand Down Expand Up @@ -1975,5 +2198,7 @@ def create_fastapi_app(
env_name=env_name,
state_cls=state_cls,
)
server.register_routes(app)
if mode is None:
mode = os.environ.get("OPENENV_MODE", ServerMode.SIMULATION.value)
server.register_routes(app, mode=mode)
return app
6 changes: 6 additions & 0 deletions src/openenv/core/env_server/web_interface.py
Original file line number Diff line number Diff line change
Expand Up @@ -435,6 +435,8 @@ def create_web_interface_app(
show_default_tab: bool = True,
title_override: Optional[str] = None,
state_cls: Type[State] = State,
*,
mode: Optional[Any] = None,
) -> FastAPI:
"""
Create a FastAPI application with web interface for the given environment.
Expand Down Expand Up @@ -467,6 +469,9 @@ def create_web_interface_app(
title instead of the default ``"OpenEnv Agentic Environment: {name}"``.
state_cls: The State subclass this environment reports. Used for the /state
response model and the state entry of /schema. Defaults to State.
mode: Server mode (``ServerMode`` or string). When ``None``, resolved
from the ``OPENENV_MODE`` environment variable, defaulting to
simulation.

Returns:
FastAPI application instance with web interface
Expand All @@ -482,6 +487,7 @@ def create_web_interface_app(
concurrency_config,
env_name=env_name,
state_cls=state_cls,
mode=mode,
)

# Load environment metadata
Expand Down
17 changes: 13 additions & 4 deletions src/openenv/core/harness/__init__.py
Original file line number Diff line number Diff line change
@@ -1,16 +1,17 @@
# SPDX-License-Identifier: BSD-3-Clause

"""Harness helpers for training, evaluation, and wrapping external agents.
"""Harness integration helpers for training, evaluation, and wrapping agents.

This package hosts two complementary layers:

1. **Trainer-side rollout API** (``openenv.core.harness.rollout``): a harness
drives an entire episode in one call (``run_white_box``/``run_black_box``)
against a resource session. Used by ``openenv collect`` and the training
tutorials.
2. **Turn-based agentic harness API** (RFC 005): the types describing an
external harness such as OpenClaw or Claude Code, where each ``step()`` is
one conversational turn. See
2. **Turn-based agentic harness API** (RFC 005): an external harness such as
OpenClaw or Claude Code runs inside the environment container, and each
``step()`` is one conversational turn. See
[`~openenv.core.harness.environment.HarnessEnvironment`] and
[`~openenv.core.harness.adapter.AgenticHarnessAdapter`].

Both layers are importable from ``openenv.core.harness``.
Expand All @@ -23,14 +24,17 @@
HarnessStartupError,
HarnessTurnTimeoutError,
)
from .bridge import build_bridge_server, HarnessMCPBridge
from .config import HarnessConfig, HarnessTransport
from .environment import HarnessAction, HarnessEnvironment
from .events import (
events_to_metadata,
HarnessClientMessage,
HarnessEvent,
HarnessEventType,
HarnessResponse,
)
from .process import HarnessProcess
from .rollout import ( # noqa: F401 (_resolve_env_reward: private back-compat re-export)
_resolve_env_reward,
build_harness_rollout_func,
Expand Down Expand Up @@ -80,16 +84,21 @@
"build_harness_rollout_func",
# Turn-based agentic harness API (RFC 005)
"AgenticHarnessAdapter",
"HarnessAction",
"HarnessClientMessage",
"HarnessConfig",
"HarnessEnvironment",
"HarnessError",
"HarnessEvent",
"HarnessEventType",
"HarnessMCPBridge",
"HarnessNotRunningError",
"HarnessProcess",
"HarnessResponse",
"HarnessStartupError",
"HarnessTransport",
"HarnessTurnTimeoutError",
"build_bridge_server",
"events_to_metadata",
"resolve_tool_conflicts",
]
Loading
Loading