3284 lines
143 KiB
Python
3284 lines
143 KiB
Python
from __future__ import annotations
|
|
|
|
import base64
|
|
import contextlib
|
|
import hashlib
|
|
import http.server
|
|
import json
|
|
import os
|
|
import shutil
|
|
import signal
|
|
import socket
|
|
import subprocess
|
|
import tempfile
|
|
import threading
|
|
import time
|
|
import unittest
|
|
from pathlib import Path
|
|
from typing import Any
|
|
from unittest import mock
|
|
|
|
import mmo_codex_home
|
|
import mmo_gateway
|
|
import mmo_mcp
|
|
import mmo_runtime
|
|
import mmo_util
|
|
import root_runner
|
|
import worker_runner
|
|
from common import ROOT, RuntimeSandbox, create_access_lab_session, root_thread_binding
|
|
from mmo_app_server import (
|
|
APP_SERVER_DYNAMIC_TOOL_TIMEOUT_SECONDS,
|
|
APP_SERVER_GOAL_OBJECTIVE_MAX_CHARS,
|
|
APP_SERVER_PROTOCOL_CODEX_VERSION,
|
|
APP_SERVER_PROTOCOL_FILE_COUNT,
|
|
APP_SERVER_PROTOCOL_SHA256,
|
|
AppServerClient,
|
|
AppServerError,
|
|
PersistentThreadHost,
|
|
UnixWebSocket,
|
|
_dynamic_tool_content_items,
|
|
_flat_mcp_dynamic_tool_name,
|
|
bounded_goal_objective,
|
|
completed_turn_presentable_text,
|
|
require_app_server_codex_version,
|
|
resumed_turns,
|
|
retain_partial_evidence,
|
|
validate_server_request_response,
|
|
)
|
|
from mmo_diagnostics import _mcp_handshake
|
|
from mmo_gateway import ensure_gateway, gateway_status, stop_gateway, stop_idle_gateways
|
|
from mmo_profiles import resolve_profile, set_active_profile
|
|
from mmo_runtime import (
|
|
_terminate_and_reap,
|
|
cancel_job,
|
|
cancel_session,
|
|
create_session,
|
|
finish_session,
|
|
list_jobs,
|
|
load_job,
|
|
load_session,
|
|
mark_session_running,
|
|
read_result,
|
|
spawn_job,
|
|
spawn_jobs,
|
|
wait_for_jobs,
|
|
)
|
|
from mmo_snapshot import compile_profile
|
|
from mmo_util import (
|
|
append_jsonl,
|
|
atomic_write_json,
|
|
http_ready,
|
|
process_alive,
|
|
process_start_token,
|
|
read_json,
|
|
read_toml,
|
|
terminate_process_group,
|
|
toml_dumps,
|
|
)
|
|
from mmo_version import (
|
|
SWITCHYARD_BASELINE_VERSION,
|
|
SWITCHYARD_MCP_NAMESPACE_BRIDGE_VERSION,
|
|
)
|
|
from mmo_workspace import (
|
|
_git_root,
|
|
_nul_paths,
|
|
capture_isolated_patch,
|
|
create_isolated_worktree,
|
|
remove_isolated_worktree,
|
|
)
|
|
from worker_runner import (
|
|
_correlate_artifact_evidence,
|
|
_correlate_command_evidence,
|
|
_fail,
|
|
finalize,
|
|
)
|
|
|
|
|
|
class AdvancedRuntimeTests(unittest.TestCase):
|
|
def test_mcp_path_sanitizer_preserves_opaque_result_contract_fields(self) -> None:
|
|
value = {
|
|
"job": {
|
|
"job_id": "job-1",
|
|
"result_path": "/private/supervisor/result.md",
|
|
"events_path": "/private/supervisor/events.jsonl",
|
|
},
|
|
"content": {
|
|
"result_path": "contract-defined destination",
|
|
"nested": {"events_path": "contract-defined evidence label"},
|
|
},
|
|
"results": {"job-1": {"preview": {"stderr_path": "model-authored diagnostic field"}}},
|
|
"records": [{"socket_path": "/private/supervisor/control.sock"}],
|
|
"live": {"content": {"socket_path": "/private/supervisor/live.sock"}},
|
|
}
|
|
|
|
sanitized = mmo_mcp._sanitize_mcp_result(value)
|
|
|
|
self.assertEqual(sanitized["job"], {"job_id": "job-1"})
|
|
self.assertEqual(sanitized["content"], value["content"])
|
|
self.assertEqual(sanitized["results"], value["results"])
|
|
self.assertEqual(sanitized["records"], [{}])
|
|
self.assertEqual(sanitized["live"], {"content": {}})
|
|
|
|
def test_goal_objective_bounding_is_deterministic_without_replacing_full_work(self) -> None:
|
|
exact = "x" * APP_SERVER_GOAL_OBJECTIVE_MAX_CHARS
|
|
self.assertEqual(bounded_goal_objective(exact), exact)
|
|
long = "A precise delegated task.\n" + ("多模型 evidence\n" * 500)
|
|
first = bounded_goal_objective(long)
|
|
self.assertEqual(first, bounded_goal_objective(long))
|
|
self.assertEqual(len(first), APP_SERVER_GOAL_OBJECTIVE_MAX_CHARS)
|
|
self.assertTrue(first.startswith("A precise delegated task."))
|
|
self.assertIn(hashlib.sha256(long.strip().encode("utf-8")).hexdigest(), first)
|
|
with self.assertRaisesRegex(ValueError, "non-empty"):
|
|
bounded_goal_objective(" ")
|
|
|
|
def test_terminal_root_result_uses_last_presentable_item_from_exact_turn(self) -> None:
|
|
with tempfile.TemporaryDirectory() as temporary:
|
|
events = Path(temporary) / "events.jsonl"
|
|
rows = [
|
|
{
|
|
"direction": "received",
|
|
"message": {
|
|
"method": "item/completed",
|
|
"params": {"turnId": "old", "item": {"type": "plan", "text": "old plan"}},
|
|
},
|
|
},
|
|
{
|
|
"direction": "sent",
|
|
"message": {
|
|
"method": "item/completed",
|
|
"params": {
|
|
"turnId": "wanted",
|
|
"item": {"type": "plan", "text": "sent only"},
|
|
},
|
|
},
|
|
},
|
|
{
|
|
"direction": "received",
|
|
"message": {
|
|
"method": "item/completed",
|
|
"params": {
|
|
"turnId": "wanted",
|
|
"item": {"type": "agentMessage", "text": "interim"},
|
|
},
|
|
},
|
|
},
|
|
{
|
|
"direction": "received",
|
|
"message": {
|
|
"method": "item/completed",
|
|
"params": {
|
|
"turnId": "wanted",
|
|
"item": {"type": "plan", "text": "authoritative plan"},
|
|
},
|
|
},
|
|
},
|
|
]
|
|
events.write_text(
|
|
"{malformed\n" + "".join(json.dumps(row) + "\n" for row in rows),
|
|
encoding="utf-8",
|
|
)
|
|
self.assertEqual(
|
|
completed_turn_presentable_text(events, "wanted"),
|
|
"authoritative plan",
|
|
)
|
|
runner = object.__new__(root_runner.RootRunner)
|
|
runner.events_path = events
|
|
self.assertEqual(
|
|
runner._result_from_turns(
|
|
[
|
|
{
|
|
"id": "old",
|
|
"status": "completed",
|
|
"items": [{"type": "agentMessage", "text": "old fallback"}],
|
|
},
|
|
{
|
|
"id": "wanted",
|
|
"status": "completed",
|
|
"items": [{"type": "agentMessage", "text": "stale turn summary"}],
|
|
},
|
|
]
|
|
),
|
|
"authoritative plan",
|
|
)
|
|
self.assertEqual(
|
|
runner._result_from_turns(
|
|
[
|
|
{
|
|
"id": "old",
|
|
"status": "completed",
|
|
"items": [{"type": "agentMessage", "text": "old fallback"}],
|
|
},
|
|
{"id": "empty", "status": "completed", "items": []},
|
|
]
|
|
),
|
|
"",
|
|
)
|
|
|
|
def test_interactive_bootstrap_goal_activates_only_after_first_turn(self) -> None:
|
|
runner = object.__new__(root_runner.RootRunner)
|
|
runner.bootstrap_goal_pending = True
|
|
runner.lifecycle_timeout = 1200.0
|
|
runner.directory = Path("/synthetic/session")
|
|
runner.state = {"goal": {"objective": "persistent session", "status": "paused"}}
|
|
runner.host = mock.Mock()
|
|
runner.host.active_turn_id = "first-accepted-turn"
|
|
runner.host.set_goal.return_value = {
|
|
"objective": "persistent session",
|
|
"status": "active",
|
|
}
|
|
with (
|
|
mock.patch.object(runner, "_update") as update,
|
|
mock.patch.object(
|
|
root_runner,
|
|
"read_session_record",
|
|
return_value={"status": "running"},
|
|
),
|
|
):
|
|
runner._activate_bootstrap_goal()
|
|
runner.host.set_goal.assert_called_once_with(status="active", timeout=1200.0)
|
|
update.assert_called_once_with(
|
|
status="running",
|
|
root_goal_status="active",
|
|
root_goal_bootstrap_pending=False,
|
|
)
|
|
self.assertFalse(runner.bootstrap_goal_pending)
|
|
|
|
runner.bootstrap_goal_pending = True
|
|
runner.host.active_turn_id = None
|
|
runner.host.set_goal.reset_mock()
|
|
runner._activate_bootstrap_goal()
|
|
runner.host.set_goal.assert_not_called()
|
|
|
|
def test_control_argument_validation_is_shared_without_changing_defaults(self) -> None:
|
|
continued = mmo_runtime._normalize_control_arguments(
|
|
"continue",
|
|
{"goal_token_budget": 20},
|
|
execution_mode="goal",
|
|
current_goal_token_budget=10,
|
|
max_goal_token_budget=20,
|
|
allowed_reasoning_efforts=["low", "high"],
|
|
goal_subject="job",
|
|
finalize_default="unused",
|
|
continue_default="Continue retained work.",
|
|
)
|
|
self.assertEqual(continued["input"], "Continue retained work.")
|
|
self.assertEqual(continued["goal_token_budget"], 20)
|
|
|
|
finalized = mmo_runtime._normalize_control_arguments(
|
|
"finalize",
|
|
{},
|
|
execution_mode="turn",
|
|
current_goal_token_budget=0,
|
|
max_goal_token_budget=0,
|
|
allowed_reasoning_efforts=["low"],
|
|
goal_subject="run",
|
|
finalize_default="Finalize retained evidence.",
|
|
)
|
|
self.assertEqual(finalized["input"], "Finalize retained evidence.")
|
|
with self.assertRaisesRegex(ValueError, "goal-mode run"):
|
|
mmo_runtime._normalize_control_arguments(
|
|
"continue",
|
|
{"goal_token_budget": 1},
|
|
execution_mode="turn",
|
|
current_goal_token_budget=0,
|
|
max_goal_token_budget=0,
|
|
allowed_reasoning_efforts=["low"],
|
|
goal_subject="run",
|
|
finalize_default="unused",
|
|
)
|
|
|
|
def test_detached_root_retains_scheduler_capacity(self) -> None:
|
|
usage = mmo_runtime._active_resource_usage(
|
|
sessions=[
|
|
{
|
|
"status": "detached",
|
|
"root_resource_lock_key": "shared-root-route",
|
|
"root_resource_units": 2,
|
|
}
|
|
],
|
|
jobs=[],
|
|
)
|
|
self.assertEqual(usage, {"shared-root-route": 2})
|
|
|
|
def test_goal_lifecycle_is_token_bounded_without_a_task_clock(self) -> None:
|
|
with RuntimeSandbox():
|
|
agent = resolve_profile("adaptive-engineering")["agents"]["implementation_specialist"]
|
|
self.assertEqual(agent["execution_mode"], "goal")
|
|
self.assertLessEqual(agent["goal_token_budget"], agent["max_goal_token_budget"])
|
|
self.assertGreaterEqual(agent["stall_warning_seconds"], 60)
|
|
for obsolete in (
|
|
"execution_policy",
|
|
"hard_wall_timeout_seconds",
|
|
"initial_active_work_seconds",
|
|
"max_active_work_seconds",
|
|
"renewal_quantum_seconds",
|
|
):
|
|
self.assertNotIn(obsolete, agent)
|
|
|
|
def test_resumed_turn_reconciliation_is_exact_not_latest_wins(self) -> None:
|
|
thread = {
|
|
"turns": [
|
|
{"id": "wanted", "status": "completed"},
|
|
{"id": "unrelated", "status": "failed"},
|
|
{"id": "active", "status": "inProgress"},
|
|
]
|
|
}
|
|
terminal, active = resumed_turns(thread, "wanted")
|
|
self.assertEqual(terminal, {"id": "wanted", "status": "completed"})
|
|
self.assertEqual(active, "active")
|
|
self.assertEqual(resumed_turns(thread, "missing"), (None, "active"))
|
|
|
|
def test_pending_turn_start_recovers_only_history_appended_after_prior_turn(self) -> None:
|
|
thread = {
|
|
"turns": [
|
|
{"id": "prior", "status": "completed"},
|
|
{"id": "accepted-before-disconnect", "status": "completed"},
|
|
]
|
|
}
|
|
self.assertEqual(
|
|
resumed_turns(thread, "prior", turn_start_pending=True),
|
|
({"id": "accepted-before-disconnect", "status": "completed"}, None),
|
|
)
|
|
self.assertEqual(
|
|
resumed_turns(thread, "prior"),
|
|
({"id": "prior", "status": "completed"}, None),
|
|
)
|
|
|
|
def test_replaced_worker_host_interrupts_and_settles_stale_turn_before_continuing(
|
|
self,
|
|
) -> None:
|
|
host = PersistentThreadHost(
|
|
state={
|
|
"thread_id": "thread-1",
|
|
"active_turn_id": "stale-turn",
|
|
"last_turn_id": "stale-turn",
|
|
}
|
|
)
|
|
client = AppServerClient(
|
|
socket_path=ROOT / "unused-stale-worker.sock",
|
|
command=["unused"],
|
|
cwd=ROOT,
|
|
env={},
|
|
events_path=ROOT / "unused-stale-worker-events.jsonl",
|
|
stderr_path=ROOT / "unused-stale-worker-stderr.log",
|
|
approval_policy="never",
|
|
)
|
|
|
|
def request(method: str, params: dict[str, Any], **_kwargs: Any) -> dict[str, Any]:
|
|
if method == "thread/resume":
|
|
return {
|
|
"thread": {
|
|
"id": "thread-1",
|
|
"turns": [{"id": "stale-turn", "status": "inProgress"}],
|
|
"status": {"type": "active"},
|
|
"goal": None,
|
|
}
|
|
}
|
|
if method == "turn/interrupt":
|
|
self.assertEqual(params["turnId"], "stale-turn")
|
|
host.on_message(
|
|
{
|
|
"method": "turn/completed",
|
|
"params": {
|
|
"threadId": "thread-1",
|
|
"turn": {"id": "stale-turn", "status": "interrupted"},
|
|
},
|
|
}
|
|
)
|
|
return {}
|
|
raise AssertionError(method)
|
|
|
|
client.request = mock.Mock(side_effect=request) # type: ignore[method-assign]
|
|
host.attach_client(client)
|
|
host.open_thread(
|
|
"resume",
|
|
{"threadId": "thread-1"},
|
|
timeout=5.0,
|
|
prior_turn_id="stale-turn",
|
|
expected_thread_id="thread-1",
|
|
interrupt_stale=True,
|
|
)
|
|
self.assertIsNone(host.active_turn_id)
|
|
self.assertIsNone(host.completed_turn)
|
|
self.assertEqual(
|
|
[call.args[0] for call in client.request.call_args_list],
|
|
["thread/resume", "turn/interrupt"],
|
|
)
|
|
|
|
def test_replaced_worker_host_preserves_completion_that_races_interrupt_error(
|
|
self,
|
|
) -> None:
|
|
host = PersistentThreadHost(
|
|
state={
|
|
"thread_id": "thread-1",
|
|
"active_turn_id": "stale-turn",
|
|
"last_turn_id": "stale-turn",
|
|
}
|
|
)
|
|
client = AppServerClient(
|
|
socket_path=ROOT / "unused-raced-worker.sock",
|
|
command=["unused"],
|
|
cwd=ROOT,
|
|
env={},
|
|
events_path=ROOT / "unused-raced-worker-events.jsonl",
|
|
stderr_path=ROOT / "unused-raced-worker-stderr.log",
|
|
approval_policy="never",
|
|
)
|
|
|
|
def request(method: str, _params: dict[str, Any], **_kwargs: Any) -> dict[str, Any]:
|
|
if method == "thread/resume":
|
|
return {
|
|
"thread": {
|
|
"id": "thread-1",
|
|
"turns": [{"id": "stale-turn", "status": "inProgress"}],
|
|
"status": {"type": "active"},
|
|
"goal": None,
|
|
}
|
|
}
|
|
if method == "turn/interrupt":
|
|
host.on_message(
|
|
{
|
|
"method": "turn/completed",
|
|
"params": {
|
|
"threadId": "thread-1",
|
|
"turn": {"id": "stale-turn", "status": "completed"},
|
|
},
|
|
}
|
|
)
|
|
raise AppServerError("turn already completed")
|
|
raise AssertionError(method)
|
|
|
|
client.request = mock.Mock(side_effect=request) # type: ignore[method-assign]
|
|
host.attach_client(client)
|
|
host.open_thread(
|
|
"resume",
|
|
{"threadId": "thread-1"},
|
|
timeout=5.0,
|
|
prior_turn_id="stale-turn",
|
|
expected_thread_id="thread-1",
|
|
interrupt_stale=True,
|
|
)
|
|
self.assertIsNone(host.active_turn_id)
|
|
self.assertEqual(
|
|
host.completed_turn,
|
|
{"id": "stale-turn", "status": "completed"},
|
|
)
|
|
|
|
def test_replaced_worker_host_does_not_swallow_a_real_interrupt_error(self) -> None:
|
|
host = PersistentThreadHost(
|
|
state={
|
|
"thread_id": "thread-1",
|
|
"active_turn_id": "stale-turn",
|
|
"last_turn_id": "stale-turn",
|
|
}
|
|
)
|
|
client = AppServerClient(
|
|
socket_path=ROOT / "unused-rejected-worker.sock",
|
|
command=["unused"],
|
|
cwd=ROOT,
|
|
env={},
|
|
events_path=ROOT / "unused-rejected-worker-events.jsonl",
|
|
stderr_path=ROOT / "unused-rejected-worker-stderr.log",
|
|
approval_policy="never",
|
|
)
|
|
client.request = mock.Mock( # type: ignore[method-assign]
|
|
side_effect=[
|
|
{
|
|
"thread": {
|
|
"id": "thread-1",
|
|
"turns": [{"id": "stale-turn", "status": "inProgress"}],
|
|
"status": {"type": "active"},
|
|
"goal": None,
|
|
}
|
|
},
|
|
AppServerError("interrupt rejected"),
|
|
]
|
|
)
|
|
host.attach_client(client)
|
|
with self.assertRaisesRegex(AppServerError, "interrupt rejected"):
|
|
host.open_thread(
|
|
"resume",
|
|
{"threadId": "thread-1"},
|
|
timeout=0.01,
|
|
prior_turn_id="stale-turn",
|
|
expected_thread_id="thread-1",
|
|
interrupt_stale=True,
|
|
)
|
|
|
|
def test_persistent_host_keeps_completion_that_precedes_start_response(self) -> None:
|
|
with tempfile.TemporaryDirectory() as temporary:
|
|
root = Path(temporary)
|
|
host = PersistentThreadHost(state={"thread_id": "thread-1"})
|
|
client = AppServerClient(
|
|
socket_path=root / "app-server.sock",
|
|
command=["unused"],
|
|
cwd=root,
|
|
env={},
|
|
events_path=root / "events.jsonl",
|
|
stderr_path=root / "stderr.log",
|
|
approval_policy="never",
|
|
)
|
|
host.attach_client(client)
|
|
|
|
def complete_before_reply(
|
|
method: str, _params: dict[str, Any], *, timeout: float = 60.0
|
|
) -> dict[str, Any]:
|
|
self.assertEqual(method, "turn/start")
|
|
self.assertEqual(timeout, 5.0)
|
|
host.on_message(
|
|
{
|
|
"method": "turn/completed",
|
|
"params": {"turn": {"id": "turn-1", "status": "completed", "items": []}},
|
|
}
|
|
)
|
|
return {"turn": {"id": "turn-1", "status": "inProgress", "items": []}}
|
|
|
|
with mock.patch.object(client, "request", side_effect=complete_before_reply):
|
|
self.assertEqual(
|
|
host.start_turn(
|
|
[{"type": "text", "text": "work"}],
|
|
effort="high",
|
|
timeout=5.0,
|
|
),
|
|
"turn-1",
|
|
)
|
|
self.assertIsNone(host.active_turn_id)
|
|
self.assertTrue(host.turn_event.is_set())
|
|
terminal = host.completed_turn
|
|
self.assertEqual(terminal.get("status") if terminal else None, "completed")
|
|
|
|
def test_persistent_host_isolates_threads_and_persists_terminal_turn_identity(self) -> None:
|
|
changes: list[dict[str, Any]] = []
|
|
host = PersistentThreadHost(
|
|
state={"thread_id": "root-thread"},
|
|
on_state_change=lambda value, _message: changes.append(dict(value)),
|
|
)
|
|
host.on_message(
|
|
{
|
|
"method": "turn/started",
|
|
"params": {
|
|
"threadId": "native-thread",
|
|
"turn": {"id": "native-turn", "status": "inProgress"},
|
|
},
|
|
}
|
|
)
|
|
host.on_message(
|
|
{
|
|
"method": "thread/goal/updated",
|
|
"params": {
|
|
"threadId": "native-thread",
|
|
"goal": {"threadId": "native-thread", "status": "complete"},
|
|
},
|
|
}
|
|
)
|
|
self.assertIsNone(host.active_turn_id)
|
|
self.assertIsNone(host.last_turn_id)
|
|
self.assertEqual(changes, [])
|
|
|
|
host.on_message(
|
|
{
|
|
"method": "turn/started",
|
|
"params": {
|
|
"threadId": "root-thread",
|
|
"turn": {"id": "root-turn", "status": "inProgress"},
|
|
},
|
|
}
|
|
)
|
|
host.on_message(
|
|
{
|
|
"method": "turn/completed",
|
|
"params": {
|
|
"threadId": "root-thread",
|
|
"turn": {"id": "root-turn", "status": "completed", "items": []},
|
|
},
|
|
}
|
|
)
|
|
self.assertIsNone(host.active_turn_id)
|
|
self.assertEqual(host.last_turn_id, "root-turn")
|
|
terminal = host.completed_turn
|
|
self.assertIsNotNone(terminal)
|
|
self.assertEqual(terminal["id"] if terminal else None, "root-turn")
|
|
self.assertEqual(changes[-1]["last_turn_id"], "root-turn")
|
|
|
|
def test_pending_server_requests_are_partitioned_by_exact_thread(self) -> None:
|
|
client = AppServerClient(
|
|
socket_path=ROOT / "unused-app-server.sock",
|
|
command=["unused"],
|
|
cwd=ROOT,
|
|
env={},
|
|
events_path=ROOT / "unused-events.jsonl",
|
|
stderr_path=ROOT / "unused-stderr.log",
|
|
approval_policy="never",
|
|
)
|
|
client._server_requests = {
|
|
"integer:1": {
|
|
"id": 1,
|
|
"method": "item/tool/requestUserInput",
|
|
"params": {"threadId": "root-thread"},
|
|
},
|
|
"string:1": {
|
|
"id": "1",
|
|
"method": "execCommandApproval",
|
|
"params": {"conversationId": "native-thread"},
|
|
},
|
|
}
|
|
self.assertEqual(
|
|
[item["id"] for item in client.pending_server_requests_for_thread("root-thread")],
|
|
[1],
|
|
)
|
|
self.assertEqual(
|
|
[item["id"] for item in client.pending_server_requests_for_thread("native-thread")],
|
|
["1"],
|
|
)
|
|
|
|
def test_switchyard_bridge_flattens_and_routes_only_discovered_mcp_tools(self) -> None:
|
|
client = AppServerClient(
|
|
socket_path=ROOT / "unused-switchyard-bridge.sock",
|
|
command=["unused"],
|
|
cwd=ROOT,
|
|
env={},
|
|
events_path=ROOT / "unused-switchyard-bridge-events.jsonl",
|
|
stderr_path=ROOT / "unused-switchyard-bridge-stderr.log",
|
|
approval_policy="never",
|
|
)
|
|
client.request = mock.Mock( # type: ignore[method-assign]
|
|
side_effect=[
|
|
{
|
|
"data": [
|
|
{
|
|
"name": "mmo_mesh",
|
|
"tools": {
|
|
"agent_list": {
|
|
"name": "agent_list",
|
|
"description": "List granted agents.",
|
|
"inputSchema": {
|
|
"type": "object",
|
|
"additionalProperties": False,
|
|
},
|
|
}
|
|
},
|
|
}
|
|
],
|
|
"nextCursor": "1",
|
|
},
|
|
{
|
|
"data": [
|
|
{
|
|
"name": "ida.server",
|
|
"tools": {
|
|
"function/get": {
|
|
"name": "function/get",
|
|
"inputSchema": {"type": "object"},
|
|
}
|
|
},
|
|
}
|
|
],
|
|
"nextCursor": None,
|
|
},
|
|
]
|
|
)
|
|
|
|
specs = client.install_switchyard_mcp_bridge(timeout=17.0)
|
|
|
|
self.assertEqual(
|
|
[spec["name"] for spec in specs],
|
|
["mmo_mcp__mmo_mesh__agent_list", "mmo_mcp__ida_server__function_get"],
|
|
)
|
|
self.assertEqual(
|
|
client._dynamic_mcp_tools,
|
|
{
|
|
"mmo_mcp__mmo_mesh__agent_list": ("mmo_mesh", "agent_list"),
|
|
"mmo_mcp__ida_server__function_get": ("ida.server", "function/get"),
|
|
},
|
|
)
|
|
self.assertEqual(
|
|
client.request.call_args_list,
|
|
[
|
|
mock.call(
|
|
"mcpServerStatus/list",
|
|
{"detail": "toolsAndAuthOnly", "limit": 100},
|
|
timeout=17.0,
|
|
),
|
|
mock.call(
|
|
"mcpServerStatus/list",
|
|
{"detail": "toolsAndAuthOnly", "limit": 100, "cursor": "1"},
|
|
timeout=17.0,
|
|
),
|
|
],
|
|
)
|
|
|
|
def test_switchyard_bridge_routes_dynamic_calls_through_codex_mcp(self) -> None:
|
|
client = AppServerClient(
|
|
socket_path=ROOT / "unused-switchyard-call.sock",
|
|
command=["unused"],
|
|
cwd=ROOT,
|
|
env={},
|
|
events_path=ROOT / "unused-switchyard-call-events.jsonl",
|
|
stderr_path=ROOT / "unused-switchyard-call-stderr.log",
|
|
approval_policy="never",
|
|
)
|
|
flat_name = "mmo_mcp__mmo_mesh__agent_list"
|
|
client._dynamic_mcp_tools = {flat_name: ("mmo_mesh", "agent_list")}
|
|
client.request = mock.Mock( # type: ignore[method-assign]
|
|
return_value={
|
|
"content": [{"type": "text", "text": "one agent"}],
|
|
"structuredContent": {"count": 1},
|
|
}
|
|
)
|
|
client.respond = mock.Mock() # type: ignore[method-assign]
|
|
|
|
client._run_dynamic_mcp_tool_call(
|
|
61,
|
|
{
|
|
"params": {
|
|
"threadId": "thread-1",
|
|
"turnId": "turn-1",
|
|
"callId": "call-1",
|
|
"namespace": None,
|
|
"tool": flat_name,
|
|
"arguments": {},
|
|
}
|
|
},
|
|
)
|
|
|
|
client.request.assert_called_once_with(
|
|
"mcpServer/tool/call",
|
|
{
|
|
"threadId": "thread-1",
|
|
"server": "mmo_mesh",
|
|
"tool": "agent_list",
|
|
"arguments": {},
|
|
},
|
|
timeout=APP_SERVER_DYNAMIC_TOOL_TIMEOUT_SECONDS,
|
|
)
|
|
client.respond.assert_called_once_with(
|
|
61,
|
|
{
|
|
"contentItems": [
|
|
{"type": "inputText", "text": "one agent"},
|
|
{"type": "inputText", "text": 'structuredContent={"count":1}'},
|
|
],
|
|
"success": True,
|
|
},
|
|
)
|
|
|
|
def test_switchyard_bridge_preserves_supported_mcp_content_and_name_limits(self) -> None:
|
|
ordinary_name = _flat_mcp_dynamic_tool_name("mmo_mesh", "agent_list")
|
|
self.assertEqual(ordinary_name, "mmo_mcp__mmo_mesh__agent_list")
|
|
self.assertFalse(ordinary_name.startswith("mcp__"))
|
|
long_name = _flat_mcp_dynamic_tool_name("server." * 40, "tool/" * 40)
|
|
self.assertLessEqual(len(long_name), 128)
|
|
self.assertRegex(long_name, r"^[a-zA-Z0-9_-]+$")
|
|
self.assertEqual(
|
|
_dynamic_tool_content_items(
|
|
{
|
|
"content": [
|
|
{"type": "image", "mimeType": "image/png", "data": "AAA"},
|
|
{"type": "audio", "mimeType": "audio/wav", "data": "BBB"},
|
|
]
|
|
}
|
|
),
|
|
[
|
|
{"type": "inputImage", "imageUrl": "data:image/png;base64,AAA"},
|
|
{"type": "inputAudio", "audioUrl": "data:audio/wav;base64,BBB"},
|
|
],
|
|
)
|
|
self.assertEqual(
|
|
_dynamic_tool_content_items(
|
|
{
|
|
"content": [{"type": "text", "text": '{"agents":[]}'}],
|
|
"structuredContent": {"agents": []},
|
|
}
|
|
),
|
|
[{"type": "inputText", "text": '{"agents":[]}'}],
|
|
)
|
|
self.assertEqual(
|
|
_dynamic_tool_content_items(
|
|
{
|
|
"content": [{"type": "text", "text": '[{"name":"one"}]'}],
|
|
"structuredContent": {"result": [{"name": "one"}]},
|
|
}
|
|
),
|
|
[{"type": "inputText", "text": '[{"name":"one"}]'}],
|
|
)
|
|
|
|
def test_temporary_switchyard_bridge_must_be_removed_when_baseline_advances(self) -> None:
|
|
self.assertEqual(SWITCHYARD_MCP_NAMESPACE_BRIDGE_VERSION, "0.2.0")
|
|
self.assertEqual(
|
|
SWITCHYARD_BASELINE_VERSION,
|
|
SWITCHYARD_MCP_NAMESPACE_BRIDGE_VERSION,
|
|
"Switchyard baseline changed: qualify native Codex MCP namespaces and remove the "
|
|
"0.2.0 dynamic-tool bridge before updating this guard",
|
|
)
|
|
|
|
def test_switchyard_version_is_read_from_the_configured_binary(self) -> None:
|
|
with mock.patch.object(
|
|
mmo_gateway.subprocess,
|
|
"run",
|
|
return_value=subprocess.CompletedProcess(
|
|
["/opt/bin/switchyard-server", "--version"],
|
|
0,
|
|
stdout="switchyard-server 0.2.0\n",
|
|
stderr="",
|
|
),
|
|
) as run:
|
|
self.assertEqual(
|
|
mmo_gateway.switchyard_version("/opt/bin/switchyard-server"),
|
|
"0.2.0",
|
|
)
|
|
self.assertEqual(run.call_args.args[0], ["/opt/bin/switchyard-server", "--version"])
|
|
|
|
def test_root_runner_popen_failure_terminalizes_the_created_session(self) -> None:
|
|
with RuntimeSandbox() as box:
|
|
session = create_session(profile="codex-harness-team", cwd=box.workspace)
|
|
with (
|
|
mock.patch.object(
|
|
mmo_runtime.subprocess,
|
|
"Popen",
|
|
side_effect=OSError("synthetic root runner launch failure"),
|
|
),
|
|
self.assertRaisesRegex(OSError, "synthetic root runner launch failure") as raised,
|
|
):
|
|
mmo_runtime._start_root_runner(session)
|
|
failed = load_session(session["session_id"])
|
|
self.assertEqual(failed["status"], "failed")
|
|
self.assertEqual(
|
|
getattr(raised.exception, "mmo_session_id", None), session["session_id"]
|
|
)
|
|
self.assertEqual(getattr(raised.exception, "mmo_session_status", None), "failed")
|
|
|
|
def test_attached_tui_fresh_context_becomes_the_canonical_root(self) -> None:
|
|
with RuntimeSandbox() as box:
|
|
session = create_session(
|
|
profile="codex-harness-team",
|
|
cwd=box.workspace,
|
|
session_kind="interactive",
|
|
)
|
|
session = mmo_runtime._start_root_runner(session)
|
|
initial_thread_id = str(session["root_thread_id"])
|
|
client = AppServerClient(
|
|
socket_path=Path(str(session["root_app_server_socket"])),
|
|
command=None,
|
|
cwd=box.workspace,
|
|
env={},
|
|
events_path=box.root / "attached-tui-events.jsonl",
|
|
stderr_path=box.root / "attached-tui-stderr.log",
|
|
approval_policy="never",
|
|
)
|
|
try:
|
|
token = process_start_token(os.getpid())
|
|
self.assertIsInstance(token, str)
|
|
mmo_runtime.update_session(
|
|
session["session_id"],
|
|
status="running",
|
|
root_client_pid=os.getpid(),
|
|
root_client_start_token=token,
|
|
)
|
|
client.start(timeout=5.0)
|
|
response = client.request(
|
|
"thread/start",
|
|
{
|
|
"cwd": str(box.workspace),
|
|
"sandbox": "workspace-write",
|
|
"approvalPolicy": "never",
|
|
"ephemeral": False,
|
|
"historyMode": "paginated",
|
|
},
|
|
timeout=5.0,
|
|
)
|
|
successor_id = str(response["thread"]["id"])
|
|
deadline = time.monotonic() + 5.0
|
|
while time.monotonic() < deadline:
|
|
current = load_session(session["session_id"])
|
|
if current.get("root_thread_id") == successor_id:
|
|
break
|
|
time.sleep(0.02)
|
|
else:
|
|
self.fail("fresh root context was not adopted")
|
|
self.assertEqual(current["root_thread_generation"], 2)
|
|
self.assertEqual(
|
|
[row["thread_id"] for row in current["root_thread_lineage"]],
|
|
[initial_thread_id, successor_id],
|
|
)
|
|
self.assertEqual(
|
|
current["root_thread_lineage"][0]["successor_thread_id"],
|
|
successor_id,
|
|
)
|
|
self.assertIsNone(current["root_thread_transition"])
|
|
self.assertEqual(
|
|
mmo_runtime.resolve_resume_session(initial_thread_id),
|
|
session["session_id"],
|
|
)
|
|
goal_deadline = time.monotonic() + 5.0
|
|
while time.monotonic() < goal_deadline:
|
|
current = load_session(session["session_id"])
|
|
if isinstance(current.get("root_goal_objective"), str):
|
|
break
|
|
time.sleep(0.02)
|
|
else:
|
|
self.fail("fresh root context did not receive its ongoing session goal")
|
|
self.assertTrue(current["root_goal_bootstrap_pending"])
|
|
finally:
|
|
client.close()
|
|
with contextlib.suppress(Exception):
|
|
cancel_session(session["session_id"])
|
|
|
|
def test_recovery_adopts_newer_top_level_thread_but_not_native_child(self) -> None:
|
|
with RuntimeSandbox() as box:
|
|
session = create_session(profile="codex-harness-team", cwd=box.workspace)
|
|
current_id = "00000000-0000-0000-0000-000000000001"
|
|
successor_id = "00000000-0000-0000-0000-000000000002"
|
|
mmo_runtime.update_session(
|
|
session["session_id"],
|
|
status="detached",
|
|
**root_thread_binding(current_id),
|
|
)
|
|
runner = root_runner.RootRunner(session["session_id"])
|
|
|
|
def thread(
|
|
thread_id: str,
|
|
created_at: int,
|
|
*,
|
|
agent_role: str | None = None,
|
|
parent_id: str | None = None,
|
|
) -> dict[str, Any]:
|
|
return {
|
|
"id": thread_id,
|
|
"sessionId": thread_id,
|
|
"cwd": str(box.workspace),
|
|
"createdAt": created_at,
|
|
"ephemeral": False,
|
|
"parentThreadId": parent_id,
|
|
"forkedFromId": None,
|
|
"agentRole": agent_role,
|
|
"agentNickname": None,
|
|
"status": {"type": "idle"},
|
|
"turns": [],
|
|
"goal": None,
|
|
}
|
|
|
|
summaries = [
|
|
thread(current_id, 1),
|
|
thread(successor_id, 2),
|
|
thread(
|
|
"00000000-0000-0000-0000-000000000003",
|
|
3,
|
|
agent_role="repo_scout",
|
|
parent_id=successor_id,
|
|
),
|
|
]
|
|
full_successor = thread(successor_id, 2)
|
|
full_successor["turns"] = [{"id": "successor-terminal-turn", "status": "completed"}]
|
|
|
|
def request(method: str, params: dict[str, Any], **_kwargs: Any) -> dict[str, Any]:
|
|
if method == "thread/list":
|
|
return {"data": summaries, "nextCursor": None}
|
|
if method == "thread/read":
|
|
self.assertEqual(params["threadId"], successor_id)
|
|
self.assertIs(params["includeTurns"], True)
|
|
return {"thread": full_successor}
|
|
raise AssertionError(method)
|
|
|
|
client = mock.Mock()
|
|
client.request.side_effect = request
|
|
runner.last_progress = time.monotonic() - 100
|
|
runner.stall_reported = True
|
|
runner._recover_root_successors(client)
|
|
recovered = load_session(session["session_id"])
|
|
self.assertEqual(recovered["root_thread_id"], successor_id)
|
|
self.assertEqual(recovered["root_thread_generation"], 2)
|
|
self.assertEqual(
|
|
[row["thread_id"] for row in recovered["root_thread_lineage"]],
|
|
[current_id, successor_id],
|
|
)
|
|
self.assertEqual(recovered["root_last_turn_id"], "successor-terminal-turn")
|
|
self.assertEqual(runner.host.last_turn_id, "successor-terminal-turn")
|
|
self.assertFalse(runner.stall_reported)
|
|
self.assertLess(time.monotonic() - runner.last_progress, 1.0)
|
|
|
|
def test_recovery_does_not_skip_an_unresolved_staged_successor(self) -> None:
|
|
with RuntimeSandbox() as box:
|
|
session = create_session(profile="codex-harness-team", cwd=box.workspace)
|
|
current_id = "00000000-0000-0000-0000-000000000011"
|
|
staged_id = "00000000-0000-0000-0000-000000000012"
|
|
later_id = "00000000-0000-0000-0000-000000000013"
|
|
binding = root_thread_binding(current_id)
|
|
binding["root_thread_transition"] = {
|
|
"from_thread_id": current_id,
|
|
"to_thread_id": staged_id,
|
|
"generation": 2,
|
|
"observed_at": "2026-08-20T00:00:00Z",
|
|
"reason": "fresh_context",
|
|
}
|
|
mmo_runtime.update_session(
|
|
session["session_id"],
|
|
status="detached",
|
|
**binding,
|
|
)
|
|
runner = root_runner.RootRunner(session["session_id"])
|
|
client = mock.Mock()
|
|
client.request.return_value = {
|
|
"data": [
|
|
{
|
|
"id": current_id,
|
|
"sessionId": current_id,
|
|
"cwd": str(box.workspace),
|
|
"createdAt": 1,
|
|
"ephemeral": False,
|
|
"parentThreadId": None,
|
|
"forkedFromId": None,
|
|
"agentRole": None,
|
|
"agentNickname": None,
|
|
"status": {"type": "idle"},
|
|
"turns": [],
|
|
"goal": None,
|
|
},
|
|
{
|
|
"id": later_id,
|
|
"sessionId": later_id,
|
|
"cwd": str(box.workspace),
|
|
"createdAt": 3,
|
|
"ephemeral": False,
|
|
"parentThreadId": None,
|
|
"forkedFromId": None,
|
|
"agentRole": None,
|
|
"agentNickname": None,
|
|
"status": {"type": "idle"},
|
|
"turns": [],
|
|
"goal": None,
|
|
},
|
|
],
|
|
"nextCursor": None,
|
|
}
|
|
runner._recover_root_successors(client)
|
|
recovered = load_session(session["session_id"])
|
|
self.assertEqual(recovered["root_thread_id"], current_id)
|
|
self.assertEqual(recovered["root_thread_generation"], 1)
|
|
self.assertEqual(
|
|
recovered["root_thread_transition"]["to_thread_id"],
|
|
staged_id,
|
|
)
|
|
|
|
def test_root_successor_commit_failure_retains_a_recoverable_staged_transition(
|
|
self,
|
|
) -> None:
|
|
with RuntimeSandbox() as box:
|
|
session = create_session(profile="codex-harness-team", cwd=box.workspace)
|
|
current_id = "00000000-0000-0000-0000-000000000031"
|
|
successor_id = "00000000-0000-0000-0000-000000000032"
|
|
mmo_runtime.update_session(
|
|
session["session_id"],
|
|
status="detached",
|
|
**root_thread_binding(current_id),
|
|
)
|
|
runner = root_runner.RootRunner(session["session_id"])
|
|
|
|
def thread(thread_id: str, created_at: int) -> dict[str, Any]:
|
|
return {
|
|
"id": thread_id,
|
|
"sessionId": thread_id,
|
|
"cwd": str(box.workspace),
|
|
"createdAt": created_at,
|
|
"ephemeral": False,
|
|
"parentThreadId": None,
|
|
"forkedFromId": None,
|
|
"agentRole": None,
|
|
"agentNickname": None,
|
|
"status": {"type": "idle"},
|
|
"turns": [],
|
|
"goal": None,
|
|
}
|
|
|
|
successor = thread(successor_id, 2)
|
|
real_publish = root_runner.publish_session_record
|
|
publications = 0
|
|
|
|
def fail_final_publish(
|
|
directory: Path,
|
|
state: dict[str, Any],
|
|
*,
|
|
mirror_run: bool,
|
|
) -> None:
|
|
nonlocal publications
|
|
publications += 1
|
|
if publications == 2:
|
|
raise OSError("synthetic final lineage publication failure")
|
|
real_publish(directory, state, mirror_run=mirror_run)
|
|
|
|
with mock.patch.object(
|
|
root_runner,
|
|
"publish_session_record",
|
|
side_effect=fail_final_publish,
|
|
):
|
|
with self.assertRaisesRegex(OSError, "final lineage publication failure"):
|
|
runner._commit_root_thread(
|
|
successor,
|
|
reason="fresh_context",
|
|
require_attached_client=False,
|
|
)
|
|
|
|
staged = load_session(session["session_id"])
|
|
self.assertEqual(staged["root_thread_id"], current_id)
|
|
self.assertEqual(staged["root_thread_generation"], 1)
|
|
self.assertEqual(staged["root_thread_transition"]["to_thread_id"], successor_id)
|
|
run = mmo_runtime.load_session_run(
|
|
session["session_id"],
|
|
str(staged["current_run_id"]),
|
|
)
|
|
self.assertEqual(run["root_thread_transition"]["to_thread_id"], successor_id)
|
|
|
|
recovering = root_runner.RootRunner(session["session_id"])
|
|
|
|
def request(method: str, params: dict[str, Any], **_kwargs: Any) -> dict[str, Any]:
|
|
if method == "thread/list":
|
|
return {
|
|
"data": [thread(current_id, 1), thread(successor_id, 2)],
|
|
"nextCursor": None,
|
|
}
|
|
if method == "thread/read":
|
|
self.assertEqual(params["threadId"], successor_id)
|
|
return {"thread": successor}
|
|
raise AssertionError(method)
|
|
|
|
client = mock.Mock()
|
|
client.request.side_effect = request
|
|
recovering._recover_root_successors(client)
|
|
recovered = load_session(session["session_id"])
|
|
self.assertEqual(recovered["root_thread_id"], successor_id)
|
|
self.assertEqual(recovered["root_thread_generation"], 2)
|
|
self.assertIsNone(recovered["root_thread_transition"])
|
|
recovered_run = mmo_runtime.load_session_run(
|
|
session["session_id"],
|
|
str(recovered["current_run_id"]),
|
|
)
|
|
self.assertEqual(recovered_run["root_thread_id"], successor_id)
|
|
self.assertEqual(recovered_run["root_thread_generation"], 2)
|
|
self.assertEqual(
|
|
recovered_run["root_thread_lineage"],
|
|
recovered["root_thread_lineage"],
|
|
)
|
|
self.assertIsNone(recovered_run["root_thread_transition"])
|
|
|
|
def test_late_fresh_context_cannot_mutate_terminal_root_lineage(self) -> None:
|
|
with RuntimeSandbox() as box:
|
|
session = create_session(profile="codex-harness-team", cwd=box.workspace)
|
|
current_id = "00000000-0000-0000-0000-000000000021"
|
|
token = process_start_token(os.getpid())
|
|
self.assertIsInstance(token, str)
|
|
mmo_runtime.update_session(
|
|
session["session_id"],
|
|
root_client_pid=os.getpid(),
|
|
root_client_start_token=token,
|
|
**root_thread_binding(current_id),
|
|
)
|
|
runner = root_runner.RootRunner(session["session_id"])
|
|
finish_session(session["session_id"], exit_code=0)
|
|
runner._on_state_change(
|
|
{
|
|
"thread_started": {
|
|
"id": "00000000-0000-0000-0000-000000000022",
|
|
"sessionId": "00000000-0000-0000-0000-000000000022",
|
|
"cwd": str(box.workspace),
|
|
"createdAt": 2,
|
|
"ephemeral": False,
|
|
"parentThreadId": None,
|
|
"forkedFromId": None,
|
|
"agentRole": None,
|
|
"agentNickname": None,
|
|
"status": {"type": "idle"},
|
|
"turns": [],
|
|
"goal": None,
|
|
}
|
|
},
|
|
{"method": "thread/started"},
|
|
)
|
|
terminal = load_session(session["session_id"])
|
|
self.assertEqual(terminal["status"], "completed")
|
|
self.assertEqual(terminal["root_thread_id"], current_id)
|
|
self.assertEqual(terminal["root_thread_generation"], 1)
|
|
|
|
def test_fresh_context_cannot_replace_a_root_with_turn_start_pending(self) -> None:
|
|
with RuntimeSandbox() as box:
|
|
session = create_session(profile="codex-harness-team", cwd=box.workspace)
|
|
current_id = "00000000-0000-0000-0000-000000000041"
|
|
successor_id = "00000000-0000-0000-0000-000000000042"
|
|
token = process_start_token(os.getpid())
|
|
self.assertIsInstance(token, str)
|
|
mmo_runtime.update_session(
|
|
session["session_id"],
|
|
status="running",
|
|
root_client_pid=os.getpid(),
|
|
root_client_start_token=token,
|
|
root_turn_start_pending=True,
|
|
**root_thread_binding(current_id),
|
|
)
|
|
runner = root_runner.RootRunner(session["session_id"])
|
|
adopted = runner._commit_root_thread(
|
|
{
|
|
"id": successor_id,
|
|
"sessionId": successor_id,
|
|
"cwd": str(box.workspace),
|
|
"createdAt": 2,
|
|
"ephemeral": False,
|
|
"parentThreadId": None,
|
|
"forkedFromId": None,
|
|
"agentRole": None,
|
|
"agentNickname": None,
|
|
"status": {"type": "idle"},
|
|
"turns": [],
|
|
"goal": None,
|
|
},
|
|
reason="fresh_context",
|
|
require_attached_client=True,
|
|
)
|
|
self.assertFalse(adopted)
|
|
current = load_session(session["session_id"])
|
|
self.assertEqual(current["root_thread_id"], current_id)
|
|
self.assertEqual(current["root_thread_generation"], 1)
|
|
finish_session(session["session_id"], exit_code=0)
|
|
|
|
def test_resume_reconciles_dead_session_workers_before_returning_control(self) -> None:
|
|
with RuntimeSandbox() as box:
|
|
session = mmo_runtime._start_root_runner(
|
|
create_session(profile="codex-harness-team", cwd=box.workspace)
|
|
)
|
|
job = spawn_job(
|
|
session_id=session["session_id"],
|
|
caller_agent=session["root_agent"],
|
|
caller_job_id=None,
|
|
agent_id="fresh_critic",
|
|
task_kind="review",
|
|
task="Retain evidence while the worker host is lost. FAKE_SLEEP_SECONDS=30",
|
|
mode="read-only",
|
|
)
|
|
metadata_path = mmo_runtime.job_dir(job["job_id"]) / "metadata.json"
|
|
deadline = time.monotonic() + 10.0
|
|
while time.monotonic() < deadline:
|
|
raw = read_json(metadata_path)
|
|
if isinstance(raw.get("active_turn_id"), str) and isinstance(
|
|
raw.get("runner_pid"), int
|
|
):
|
|
break
|
|
time.sleep(0.02)
|
|
else:
|
|
self.fail("worker did not publish its running process identity")
|
|
mmo_runtime.detach_session(session["session_id"])
|
|
os.killpg(int(raw["runner_pid"]), signal.SIGKILL)
|
|
dead_deadline = time.monotonic() + 5.0
|
|
while process_alive(int(raw["runner_pid"])) and time.monotonic() < dead_deadline:
|
|
time.sleep(0.02)
|
|
self.assertFalse(process_alive(int(raw["runner_pid"])))
|
|
self.assertEqual(read_json(metadata_path)["status"], "running")
|
|
|
|
resumed = mmo_runtime.begin_resume_run(session["session_id"])
|
|
self.assertEqual(resumed["status"], "running")
|
|
suspended = load_job(job["job_id"])
|
|
self.assertEqual(suspended["status"], "suspended")
|
|
self.assertTrue(Path(suspended["partial_result_path"]).is_file())
|
|
audit = Path(load_session(session["session_id"])["audit_path"]).read_text(
|
|
encoding="utf-8"
|
|
)
|
|
self.assertIn('"event":"session_jobs_reconciled"', audit)
|
|
self.assertIn(job["job_id"], audit)
|
|
mmo_runtime.stop_session(session["session_id"], grace_seconds=0)
|
|
|
|
def test_detach_does_not_resurrect_a_concurrently_completed_session(self) -> None:
|
|
with RuntimeSandbox() as box:
|
|
session = create_session(profile="codex-harness-team", cwd=box.workspace)
|
|
session = mmo_runtime.update_session(
|
|
session["session_id"],
|
|
root_control_socket_ready=True,
|
|
**root_thread_binding("root-thread"),
|
|
)
|
|
|
|
def complete_during_delivery(*_args: Any, **_kwargs: Any) -> dict[str, Any]:
|
|
finish_session(
|
|
session["session_id"],
|
|
exit_code=0,
|
|
expected_run_id=session["current_run_id"],
|
|
)
|
|
return {"ok": True, "result": {"detached": True}}
|
|
|
|
with (
|
|
mock.patch.object(
|
|
mmo_runtime, "_root_control_socket", return_value=box.root / "control.sock"
|
|
),
|
|
mock.patch.object(
|
|
mmo_runtime, "send_control_request", side_effect=complete_during_delivery
|
|
),
|
|
):
|
|
detached = mmo_runtime.detach_session(session["session_id"])
|
|
self.assertEqual(detached["session"]["status"], "completed")
|
|
self.assertEqual(load_session(session["session_id"])["status"], "completed")
|
|
|
|
def test_resume_rejects_native_capability_inventory_drift(self) -> None:
|
|
with RuntimeSandbox() as box:
|
|
session = create_session(profile="codex-harness-team", cwd=box.workspace)
|
|
mmo_runtime.update_session(
|
|
session["session_id"],
|
|
status="detached",
|
|
native_token_hashes={"repo_scout": "0" * 64},
|
|
**root_thread_binding("00000000-0000-0000-0000-000000000147"),
|
|
)
|
|
try:
|
|
with self.assertRaisesRegex(RuntimeError, "native capabilities disagree"):
|
|
mmo_runtime.begin_resume_run(session["session_id"])
|
|
finally:
|
|
cancel_session(session["session_id"])
|
|
|
|
def test_native_stop_archives_the_thread_and_refresh_does_not_resurrect_it(self) -> None:
|
|
with RuntimeSandbox() as box:
|
|
session = create_session(profile="codex-harness-team", cwd=box.workspace)
|
|
run_ref = "ar_native_test"
|
|
mmo_runtime.update_session(
|
|
session["session_id"],
|
|
root_native_runs={
|
|
run_ref: {
|
|
"agent_run_ref": run_ref,
|
|
"agent": "repo_scout",
|
|
"backend": "native",
|
|
"thread_id": "native-thread",
|
|
"status": {"type": "idle"},
|
|
"control_revision": 0,
|
|
}
|
|
},
|
|
**root_thread_binding("root-thread"),
|
|
)
|
|
runner = root_runner.RootRunner(session["session_id"])
|
|
calls: list[tuple[str, dict[str, Any]]] = []
|
|
|
|
def request(method: str, params: dict[str, Any], **_kwargs: Any) -> dict[str, Any]:
|
|
calls.append((method, dict(params)))
|
|
if method == "thread/read":
|
|
return {
|
|
"thread": {
|
|
"id": "native-thread",
|
|
"status": {"type": "idle"},
|
|
"turns": [],
|
|
}
|
|
}
|
|
if method == "thread/list":
|
|
return {"data": [], "nextCursor": None}
|
|
return {}
|
|
|
|
runner.client = mock.Mock(request=mock.Mock(side_effect=request))
|
|
stopped = runner._execute_control(
|
|
{
|
|
"action": "stop",
|
|
"target_thread_id": "native-thread",
|
|
"expected_revision": 0,
|
|
}
|
|
)
|
|
self.assertTrue(stopped["stopping"])
|
|
self.assertIn(("thread/archive", {"threadId": "native-thread"}), calls)
|
|
self.assertEqual(
|
|
load_session(session["session_id"])["root_native_runs"][run_ref]["status"],
|
|
"stopped",
|
|
)
|
|
refreshed = runner.refresh_native_runs()
|
|
self.assertEqual(len(refreshed), 1)
|
|
self.assertEqual(refreshed[0]["status"], "stopped")
|
|
finish_session(session["session_id"], exit_code=0)
|
|
|
|
def test_native_detach_does_not_hide_a_later_idle_thread(self) -> None:
|
|
with RuntimeSandbox() as box:
|
|
session = create_session(profile="codex-harness-team", cwd=box.workspace)
|
|
run_ref = "ar_native_detached"
|
|
mmo_runtime.update_session(
|
|
session["session_id"],
|
|
root_native_runs={
|
|
run_ref: {
|
|
"agent_run_ref": run_ref,
|
|
"agent": "repo_scout",
|
|
"backend": "native",
|
|
"thread_id": "native-thread",
|
|
"status": "detached",
|
|
"control_revision": 1,
|
|
}
|
|
},
|
|
**root_thread_binding("root-thread"),
|
|
)
|
|
runner = root_runner.RootRunner(session["session_id"])
|
|
runner.client = mock.Mock(
|
|
request=mock.Mock(
|
|
return_value={
|
|
"data": [
|
|
{
|
|
"id": "native-thread",
|
|
"agentRole": "repo_scout",
|
|
"status": {"type": "idle"},
|
|
}
|
|
],
|
|
"nextCursor": None,
|
|
}
|
|
)
|
|
)
|
|
refreshed = runner.refresh_native_runs()
|
|
self.assertEqual(refreshed[0]["status"], {"type": "idle"})
|
|
self.assertEqual(refreshed[0]["control_revision"], 1)
|
|
finish_session(session["session_id"], exit_code=0)
|
|
|
|
def test_late_root_goal_notification_cannot_resurrect_a_terminal_session(self) -> None:
|
|
with RuntimeSandbox() as box:
|
|
session = create_session(profile="codex-harness-team", cwd=box.workspace)
|
|
runner = root_runner.RootRunner(session["session_id"])
|
|
finish_session(session["session_id"], exit_code=0)
|
|
runner._on_state_change(
|
|
{
|
|
"goal": {
|
|
"status": "active",
|
|
"objective": "already complete",
|
|
"tokensUsed": 1,
|
|
}
|
|
},
|
|
{"method": "thread/goal/updated"},
|
|
)
|
|
self.assertEqual(load_session(session["session_id"])["status"], "completed")
|
|
|
|
def test_root_turn_recovery_markers_fail_closed_on_persistence_error(self) -> None:
|
|
with RuntimeSandbox() as box:
|
|
session = create_session(profile="codex-harness-team", cwd=box.workspace)
|
|
runner = root_runner.RootRunner(session["session_id"])
|
|
with mock.patch.object(
|
|
runner,
|
|
"_update",
|
|
side_effect=OSError("synthetic state persistence failure"),
|
|
):
|
|
with self.assertRaisesRegex(OSError, "synthetic state persistence failure"):
|
|
runner._on_state_change(
|
|
{"turn_start_pending": True},
|
|
{"method": "mmo/turn/reset"},
|
|
)
|
|
runner._on_state_change(
|
|
{"last_observability_event": "warning"},
|
|
{"method": "warning"},
|
|
)
|
|
finish_session(session["session_id"], exit_code=0)
|
|
|
|
def test_runner_launch_failure_after_popen_terminates_and_reaps_host(self) -> None:
|
|
with RuntimeSandbox() as box:
|
|
directory = box.root / "synthetic-job"
|
|
directory.mkdir()
|
|
token = "ephemeral-worker-capability"
|
|
metadata = {"mcp_caller_token_hash": hashlib.sha256(token.encode("utf-8")).hexdigest()}
|
|
sleeper = subprocess.Popen(["sleep", "30"], start_new_session=True)
|
|
try:
|
|
with (
|
|
mock.patch.object(
|
|
mmo_runtime, "read_job_record", side_effect=[metadata, RuntimeError("boom")]
|
|
),
|
|
mock.patch.object(mmo_runtime.subprocess, "Popen", return_value=sleeper),
|
|
):
|
|
with self.assertRaisesRegex(RuntimeError, "boom"):
|
|
mmo_runtime._launch_worker_runner(directory, token)
|
|
sleeper.wait(timeout=5)
|
|
self.assertNotIn(sleeper.pid, mmo_runtime._RUNNER_PROCESSES)
|
|
self.assertFalse(process_alive(sleeper.pid))
|
|
finally:
|
|
if sleeper.poll() is None:
|
|
terminate_process_group(sleeper.pid, grace_seconds=1.0)
|
|
sleeper.wait(timeout=3)
|
|
|
|
def test_persistent_runner_is_reaped_after_asynchronous_exit(self) -> None:
|
|
process = subprocess.Popen(["sleep", "0.1"], start_new_session=True)
|
|
mmo_runtime._track_runner(process)
|
|
deadline = time.monotonic() + 5
|
|
while process.pid in mmo_runtime._RUNNER_PROCESSES:
|
|
self.assertLess(time.monotonic(), deadline)
|
|
time.sleep(0.02)
|
|
self.assertEqual(process.returncode, 0)
|
|
|
|
def test_partial_evidence_prefers_recent_events_with_a_bounded_scan(self) -> None:
|
|
with tempfile.TemporaryDirectory() as temporary:
|
|
directory = Path(temporary)
|
|
events = directory / "events.jsonl"
|
|
for index in range(10):
|
|
append_jsonl(
|
|
events,
|
|
{
|
|
"message": {
|
|
"params": {
|
|
"item": {
|
|
"type": "agentMessage",
|
|
"text": f"message-{index}-" + "x" * 120,
|
|
}
|
|
}
|
|
}
|
|
},
|
|
)
|
|
with mock.patch("mmo_app_server.PARTIAL_EVENT_WINDOW_BYTES", 600):
|
|
retained = retain_partial_evidence({}, directory, reason="synthetic interruption")
|
|
text = Path(retained["partial_result_path"]).read_text(encoding="utf-8")
|
|
self.assertTrue(retained["partial_trace_window_truncated"])
|
|
self.assertIn("message-9-", text)
|
|
self.assertNotIn("message-0-", text)
|
|
self.assertIn("complete durable JSONL trace remains available", text)
|
|
|
|
append_jsonl(
|
|
events,
|
|
{
|
|
"message": {
|
|
"params": {"item": {"type": "agentMessage", "text": "new-cycle-evidence"}}
|
|
}
|
|
},
|
|
)
|
|
second = retain_partial_evidence(
|
|
{"result_state": "read"}, directory, reason="second interruption"
|
|
)
|
|
second_text = Path(second["partial_result_path"]).read_text(encoding="utf-8")
|
|
self.assertIn("second interruption", second_text)
|
|
self.assertIn("new-cycle-evidence", second_text)
|
|
self.assertNotEqual(second["partial_result_sha256"], retained["partial_result_sha256"])
|
|
self.assertEqual(second["result_state"], "unread")
|
|
|
|
def test_contract_repair_does_not_create_fresh_grace_after_deadline(self) -> None:
|
|
with tempfile.TemporaryDirectory() as temporary:
|
|
directory = Path(temporary)
|
|
result_path = directory / "result.md"
|
|
result_path.write_text("not json", encoding="utf-8")
|
|
with (
|
|
mock.patch.object(worker_runner, "update") as update,
|
|
mock.patch.object(worker_runner, "_start_turn") as start_turn,
|
|
):
|
|
result = worker_runner._repair_strict_contract(
|
|
directory=directory,
|
|
state={},
|
|
state_lock=threading.RLock(),
|
|
metadata={
|
|
"contract_enforcement": "strict",
|
|
"finalization_grace_seconds": 120,
|
|
},
|
|
output_contract={
|
|
"type": "object",
|
|
"required": ["answer"],
|
|
"properties": {"answer": {"type": "string"}},
|
|
},
|
|
result_text="not json",
|
|
result_path=result_path,
|
|
deadline=time.monotonic() - 1,
|
|
)
|
|
self.assertEqual(result, "not json")
|
|
update.assert_not_called()
|
|
start_turn.assert_not_called()
|
|
|
|
def test_app_server_lifecycle_covers_tool_mcp_startup_policy(self) -> None:
|
|
resolved = {
|
|
"agents": {
|
|
"plain": {"tool_mcp_servers": {}, "can_spawn": [], "controls": {}},
|
|
"caller": {
|
|
"tool_mcp_servers": {},
|
|
"can_spawn": ["tool_user"],
|
|
"controls": {},
|
|
},
|
|
"tool_user": {
|
|
"tool_mcp_servers": {"slow_tool": {"required": True}},
|
|
"can_spawn": [],
|
|
"controls": {},
|
|
"backends": ["mcp"],
|
|
},
|
|
},
|
|
"tool_mcp_servers": {"slow_tool": {"startup_timeout_sec": 1500.5}},
|
|
}
|
|
self.assertEqual(mmo_runtime.app_server_lifecycle_timeout(resolved, "plain"), 1200.0)
|
|
self.assertEqual(
|
|
mmo_runtime.app_server_lifecycle_timeout(resolved, "tool_user"),
|
|
1530.5,
|
|
)
|
|
self.assertEqual(mmo_codex_home._mesh_tool_timeout(resolved, "caller"), 1560.5)
|
|
|
|
def test_app_server_runtime_admission_rejects_unreviewed_codex_version(self) -> None:
|
|
observed = "0.148.0" if APP_SERVER_PROTOCOL_CODEX_VERSION != "0.148.0" else "0.146.0"
|
|
with mock.patch(
|
|
"mmo_app_server.subprocess.run",
|
|
return_value=subprocess.CompletedProcess(
|
|
["codex", "--version"],
|
|
0,
|
|
stdout=f"codex-cli {observed}\n",
|
|
stderr="",
|
|
),
|
|
):
|
|
with self.assertRaisesRegex(AppServerError, "protocol version mismatch"):
|
|
require_app_server_codex_version("codex")
|
|
|
|
def test_resume_rechecks_the_pinned_codex_executable(self) -> None:
|
|
with RuntimeSandbox() as box:
|
|
session = create_session(profile="codex-harness-team", cwd=box.workspace)
|
|
mmo_runtime.update_session(
|
|
session["session_id"],
|
|
status="detached",
|
|
**root_thread_binding("00000000-0000-0000-0000-000000000147"),
|
|
)
|
|
try:
|
|
with mock.patch.object(
|
|
mmo_runtime,
|
|
"require_app_server_codex_version",
|
|
side_effect=AppServerError("synthetic protocol version mismatch"),
|
|
) as gate:
|
|
with self.assertRaisesRegex(AppServerError, "protocol version mismatch"):
|
|
mmo_runtime.begin_resume_run(session["session_id"])
|
|
gate.assert_called_once_with(str(Path(session["codex_binary"]).resolve()))
|
|
finally:
|
|
cancel_session(session["session_id"])
|
|
|
|
def test_app_server_records_sent_events_under_the_wire_lock(self) -> None:
|
|
client = AppServerClient(
|
|
socket_path=ROOT / "unused-app-server.sock",
|
|
command=["unused"],
|
|
cwd=ROOT,
|
|
env={},
|
|
events_path=ROOT / "unused-events.jsonl",
|
|
stderr_path=ROOT / "unused-stderr.log",
|
|
approval_policy="never",
|
|
)
|
|
transport = mock.Mock(closed=False)
|
|
client.transport = transport
|
|
recorded: list[dict[str, Any]] = []
|
|
|
|
def record(_direction: str, message: dict[str, Any]) -> None:
|
|
self.assertTrue(client._write_lock.locked())
|
|
recorded.append(message)
|
|
|
|
with mock.patch.object(client, "_record", side_effect=record):
|
|
client._send({"method": "test", "params": {}})
|
|
|
|
self.assertEqual([item["method"] for item in recorded], ["test"])
|
|
self.assertNotIn("jsonrpc", recorded[0])
|
|
transport.send_text.assert_called_once()
|
|
|
|
def test_app_server_uses_only_the_reviewed_codex_0149_contract(self) -> None:
|
|
self.assertEqual(APP_SERVER_PROTOCOL_CODEX_VERSION, "0.149.0")
|
|
self.assertEqual(APP_SERVER_PROTOCOL_FILE_COUNT, 401)
|
|
self.assertEqual(
|
|
APP_SERVER_PROTOCOL_SHA256,
|
|
"fcfeaf23728b96ab73916a21302eb7a16629e67ee99f7ee47b60fad6b6e5ee1a",
|
|
)
|
|
|
|
def test_app_server_retries_only_the_exact_overload_error(self) -> None:
|
|
client = AppServerClient(
|
|
socket_path=ROOT / "unused-retry.sock",
|
|
command=["unused"],
|
|
cwd=ROOT,
|
|
env={},
|
|
events_path=ROOT / "unused-retry-events.jsonl",
|
|
stderr_path=ROOT / "unused-retry-stderr.log",
|
|
approval_policy="never",
|
|
)
|
|
client.transport = mock.Mock(closed=False)
|
|
sent: list[dict[str, Any]] = []
|
|
|
|
def send(message: dict[str, Any]) -> None:
|
|
sent.append(message)
|
|
response = (
|
|
{
|
|
"id": message["id"],
|
|
"error": {"code": -32001, "message": "Server overloaded; retry later."},
|
|
}
|
|
if len(sent) < 3
|
|
else {"id": message["id"], "result": {"ok": True}}
|
|
)
|
|
with client._condition:
|
|
client._responses[message["id"]] = response
|
|
client._condition.notify_all()
|
|
|
|
with (
|
|
mock.patch.object(client, "_send", side_effect=send),
|
|
mock.patch("mmo_app_server.random.uniform", return_value=1.0),
|
|
mock.patch("mmo_app_server.time.sleep") as sleep,
|
|
):
|
|
self.assertEqual(client.request("thread/read", {}, timeout=5), {"ok": True})
|
|
self.assertEqual([message["id"] for message in sent], [1, 2, 3])
|
|
self.assertTrue(all("jsonrpc" not in message for message in sent))
|
|
self.assertEqual([call.args[0] for call in sleep.call_args_list], [0.1, 0.2])
|
|
|
|
def test_app_server_does_not_retry_other_minus_32001_errors(self) -> None:
|
|
client = AppServerClient(
|
|
socket_path=ROOT / "unused-no-retry.sock",
|
|
command=["unused"],
|
|
cwd=ROOT,
|
|
env={},
|
|
events_path=ROOT / "unused-no-retry-events.jsonl",
|
|
stderr_path=ROOT / "unused-no-retry-stderr.log",
|
|
approval_policy="never",
|
|
)
|
|
client.transport = mock.Mock(closed=False)
|
|
sent: list[dict[str, Any]] = []
|
|
|
|
def send(message: dict[str, Any]) -> None:
|
|
sent.append(message)
|
|
with client._condition:
|
|
client._responses[message["id"]] = {
|
|
"id": message["id"],
|
|
"error": {"code": -32001, "message": "thread not found"},
|
|
}
|
|
client._condition.notify_all()
|
|
|
|
with mock.patch.object(client, "_send", side_effect=send):
|
|
with self.assertRaisesRegex(AppServerError, "thread not found"):
|
|
client.request("thread/read", {}, timeout=5)
|
|
self.assertEqual(len(sent), 1)
|
|
|
|
@staticmethod
|
|
def _socket_pair_websocket() -> tuple[UnixWebSocket, socket.socket]:
|
|
client_socket, peer = socket.socketpair()
|
|
transport = object.__new__(UnixWebSocket)
|
|
transport.path = ROOT / "socket-pair"
|
|
transport._socket = client_socket
|
|
transport._buffer = bytearray()
|
|
transport._send_lock = threading.Lock()
|
|
transport._closed = False
|
|
return transport, peer
|
|
|
|
def test_websocket_handshake_rejects_missing_or_unsolicited_upgrade_fields(self) -> None:
|
|
variants = {
|
|
"missing connection token": "Connection: keep-alive\r\n",
|
|
"unsolicited extension": (
|
|
"Connection: Upgrade\r\nSec-WebSocket-Extensions: permessage-deflate\r\n"
|
|
),
|
|
}
|
|
for label, extra_headers in variants.items():
|
|
with self.subTest(label=label), tempfile.TemporaryDirectory() as temporary:
|
|
path = Path(temporary) / "app.sock"
|
|
listener = socket.socket(socket.AF_UNIX, socket.SOCK_STREAM)
|
|
listener.bind(str(path))
|
|
listener.listen(1)
|
|
|
|
def serve(
|
|
server_socket: socket.socket = listener,
|
|
response_headers: str = extra_headers,
|
|
) -> None:
|
|
connection, _address = server_socket.accept()
|
|
try:
|
|
request = bytearray()
|
|
while b"\r\n\r\n" not in request:
|
|
request.extend(connection.recv(4096))
|
|
headers = {}
|
|
for line in bytes(request).decode("ascii").split("\r\n")[1:]:
|
|
if ":" in line:
|
|
name, value = line.split(":", 1)
|
|
headers[name.strip().casefold()] = value.strip()
|
|
key = headers["sec-websocket-key"]
|
|
digest = base64.b64encode(
|
|
hashlib.sha1( # noqa: S324 - required by RFC 6455
|
|
(key + "258EAFA5-E914-47DA-95CA-C5AB0DC85B11").encode()
|
|
).digest()
|
|
).decode()
|
|
connection.sendall(
|
|
(
|
|
"HTTP/1.1 101 Switching Protocols\r\n"
|
|
"Upgrade: websocket\r\n"
|
|
f"{response_headers}"
|
|
f"Sec-WebSocket-Accept: {digest}\r\n\r\n"
|
|
).encode("ascii")
|
|
)
|
|
finally:
|
|
connection.close()
|
|
|
|
server = threading.Thread(target=serve)
|
|
server.start()
|
|
try:
|
|
with self.assertRaises(AppServerError):
|
|
UnixWebSocket(path, timeout=2)
|
|
finally:
|
|
server.join(timeout=2)
|
|
listener.close()
|
|
self.assertFalse(server.is_alive())
|
|
|
|
def test_websocket_accepts_fragmented_text_with_control_interleaving(self) -> None:
|
|
transport, peer = self._socket_pair_websocket()
|
|
try:
|
|
peer.sendall(b"\x01\x03hel\x89\x01?\x80\x02lo")
|
|
self.assertEqual(transport.receive_text(), "hello")
|
|
pong = peer.recv(64)
|
|
self.assertEqual(pong[0] & 0x0F, 0x0A)
|
|
self.assertTrue(pong[1] & 0x80)
|
|
finally:
|
|
transport.close()
|
|
peer.close()
|
|
|
|
def test_websocket_rejects_nonconforming_server_frames(self) -> None:
|
|
cases = {
|
|
"masked": b"\x81\x80" + b"mask",
|
|
"rsv": b"\xc1\x00",
|
|
"fragmented control": b"\x09\x00",
|
|
"non-minimal length": b"\x81\x7e\x00\x01x",
|
|
"binary": b"\x82\x00",
|
|
}
|
|
for label, frame in cases.items():
|
|
with self.subTest(label=label):
|
|
transport, peer = self._socket_pair_websocket()
|
|
try:
|
|
peer.sendall(frame)
|
|
with self.assertRaises(AppServerError):
|
|
transport.receive_text()
|
|
self.assertTrue(transport.closed)
|
|
finally:
|
|
transport.close()
|
|
peer.close()
|
|
|
|
def test_websocket_close_emits_a_masked_normal_close_frame(self) -> None:
|
|
transport, peer = self._socket_pair_websocket()
|
|
try:
|
|
transport.close()
|
|
frame = peer.recv(64)
|
|
self.assertEqual(frame[0] & 0x0F, 0x08)
|
|
self.assertTrue(frame[1] & 0x80)
|
|
self.assertEqual(frame[1] & 0x7F, 2)
|
|
finally:
|
|
peer.close()
|
|
|
|
def test_app_server_pending_response_shapes_are_method_specific(self) -> None:
|
|
valid: dict[str, dict[str, Any]] = {
|
|
"item/tool/requestUserInput": {"answers": {"question": {"answers": ["yes"]}}},
|
|
"mcpServer/elicitation/request": {"action": "decline"},
|
|
"item/commandExecution/requestApproval": {"decision": "decline"},
|
|
"item/fileChange/requestApproval": {"decision": "cancel"},
|
|
"item/permissions/requestApproval": {
|
|
"permissions": {"network": {"enabled": False}},
|
|
"scope": "turn",
|
|
},
|
|
"applyPatchApproval": {"decision": {"denied": {"rejection": "no"}}},
|
|
"execCommandApproval": {"decision": "abort"},
|
|
"execCommandApproval:mcp": {"decision": "approved_mcp_policy_amendment"},
|
|
}
|
|
for method, response in valid.items():
|
|
method = method.removesuffix(":mcp")
|
|
with self.subTest(method=method):
|
|
validate_server_request_response(method, response)
|
|
|
|
invalid: dict[str, dict[str, Any]] = {
|
|
"item/tool/requestUserInput": {"decision": "decline"},
|
|
"mcpServer/elicitation/request": {"action": "approve"},
|
|
"item/commandExecution/requestApproval": {"decision": "approved"},
|
|
"item/fileChange/requestApproval": {"decision": "approved"},
|
|
"item/permissions/requestApproval": {"permissions": {"network": True}},
|
|
"applyPatchApproval": {"decision": "decline"},
|
|
"execCommandApproval": {"decision": "decline"},
|
|
}
|
|
for method, response in invalid.items():
|
|
with self.subTest(method=method), self.assertRaises(ValueError):
|
|
validate_server_request_response(method, response)
|
|
|
|
def test_app_server_schema_projection_preserves_shape_for_strict_generation(self) -> None:
|
|
contract = {
|
|
"type": "object",
|
|
"properties": {
|
|
"status": {"type": "string", "enum": ["ok", "blocked"]},
|
|
"sources": {
|
|
"type": "array",
|
|
"items": {"type": "string", "format": "uri"},
|
|
"uniqueItems": True,
|
|
},
|
|
"optional_note": {"type": ["string", "null"]},
|
|
},
|
|
"required": ["status", "sources"],
|
|
"additionalProperties": False,
|
|
"allOf": [
|
|
{
|
|
"if": {"properties": {"status": {"const": "blocked"}}},
|
|
"then": {"properties": {"optional_note": {"type": "string"}}},
|
|
}
|
|
],
|
|
}
|
|
projected = worker_runner._codex_output_schema(contract)
|
|
self.assertIsNotNone(projected)
|
|
assert projected is not None
|
|
self.assertEqual(worker_runner._codex_output_schema_errors(projected), [])
|
|
self.assertEqual(projected["required"], ["status", "sources"])
|
|
self.assertNotIn("optional_note", projected["properties"])
|
|
self.assertNotIn("allOf", projected)
|
|
self.assertNotIn("uniqueItems", projected["properties"]["sources"])
|
|
self.assertNotIn("format", projected["properties"]["sources"]["items"])
|
|
self.assertIn("allOf", contract)
|
|
|
|
def test_codex_output_schema_gate_is_conservative_and_recursive(self) -> None:
|
|
compatible = {
|
|
"type": "object",
|
|
"properties": {
|
|
"status": {"type": "string", "enum": ["ok", "blocked"]},
|
|
"details": {
|
|
"type": "object",
|
|
"properties": {"at": {"type": "string", "format": "date-time"}},
|
|
"required": ["at"],
|
|
"additionalProperties": False,
|
|
},
|
|
},
|
|
"required": ["status", "details"],
|
|
"additionalProperties": False,
|
|
}
|
|
self.assertEqual(worker_runner._codex_output_schema_errors(compatible), [])
|
|
|
|
incompatible_cases = {
|
|
"non-object root": {"type": "array", "items": {"type": "string"}},
|
|
"optional field": {
|
|
"type": "object",
|
|
"properties": {"status": {"type": "string"}},
|
|
"required": [],
|
|
"additionalProperties": False,
|
|
},
|
|
"open object": {
|
|
"type": "object",
|
|
"properties": {},
|
|
"required": [],
|
|
"additionalProperties": True,
|
|
},
|
|
"conditional": {
|
|
**compatible,
|
|
"allOf": [{"if": {"properties": {}}, "then": {"properties": {}}}],
|
|
},
|
|
"unique array": {
|
|
"type": "object",
|
|
"properties": {
|
|
"values": {
|
|
"type": "array",
|
|
"items": {"type": "string"},
|
|
"uniqueItems": True,
|
|
}
|
|
},
|
|
"required": ["values"],
|
|
"additionalProperties": False,
|
|
},
|
|
"unsupported format": {
|
|
"type": "object",
|
|
"properties": {"source": {"type": "string", "format": "uri"}},
|
|
"required": ["source"],
|
|
"additionalProperties": False,
|
|
},
|
|
}
|
|
for label, schema in incompatible_cases.items():
|
|
with self.subTest(label=label):
|
|
self.assertTrue(worker_runner._codex_output_schema_errors(schema))
|
|
|
|
def test_v1_gateway_state_is_rejected_in_place(self) -> None:
|
|
with RuntimeSandbox() as box:
|
|
snapshot_hash = "a" * 64
|
|
gateway_hash = "b" * 64
|
|
gateway_directory = box.state / "gateways" / gateway_hash
|
|
gateway_directory.mkdir(parents=True)
|
|
atomic_write_json(
|
|
gateway_directory / "gateway.json",
|
|
{
|
|
"schema_version": 1,
|
|
"gateway_hash": gateway_hash,
|
|
"snapshot_hash": snapshot_hash,
|
|
"pid": 999_999_999,
|
|
"process_start_token": "not-running",
|
|
"status": "stopped",
|
|
},
|
|
)
|
|
state = mmo_gateway._read_gateway_state_by_key(gateway_hash)
|
|
if state is None:
|
|
self.fail("v1 gateway state disappeared")
|
|
self.assertEqual(state.get("status"), "invalid")
|
|
self.assertEqual(read_json(gateway_directory / "gateway.json")["schema_version"], 1)
|
|
self.assertFalse((box.state / "gateways" / snapshot_hash).exists())
|
|
|
|
def test_v1_session_state_is_rejected_in_place(self) -> None:
|
|
with RuntimeSandbox() as box:
|
|
session = create_session(profile="codex-harness-team", cwd=box.workspace)
|
|
path = box.state / "sessions" / session["session_id"] / "session.json"
|
|
current = read_json(path)
|
|
stale = {**current, "schema_version": 1}
|
|
atomic_write_json(path, stale)
|
|
try:
|
|
with self.assertRaisesRegex(ValueError, "unsupported session state schema"):
|
|
load_session(session["session_id"])
|
|
self.assertEqual(read_json(path)["schema_version"], 1)
|
|
finally:
|
|
atomic_write_json(path, current)
|
|
finish_session(session["session_id"], exit_code=0)
|
|
|
|
def test_root_thread_lineage_is_strictly_ordered_and_self_consistent(self) -> None:
|
|
with RuntimeSandbox() as box:
|
|
session = create_session(profile="codex-harness-team", cwd=box.workspace)
|
|
path = box.state / "sessions" / session["session_id"] / "session.json"
|
|
original = read_json(path)
|
|
first_id = "00000000-0000-0000-0000-000000000201"
|
|
second_id = "00000000-0000-0000-0000-000000000202"
|
|
valid = {
|
|
**original,
|
|
"root_thread_id": second_id,
|
|
"root_thread_generation": 2,
|
|
"root_thread_lineage": [
|
|
{
|
|
"generation": 1,
|
|
"thread_id": first_id,
|
|
"codex_session_id": first_id,
|
|
"adopted_at": "2026-08-20T00:00:00Z",
|
|
"reason": "initial",
|
|
"superseded_at": "2026-08-20T00:01:00Z",
|
|
"successor_thread_id": second_id,
|
|
},
|
|
{
|
|
"generation": 2,
|
|
"thread_id": second_id,
|
|
"codex_session_id": second_id,
|
|
"adopted_at": "2026-08-20T00:01:00Z",
|
|
"reason": "fresh_context",
|
|
},
|
|
],
|
|
"root_thread_transition": None,
|
|
}
|
|
try:
|
|
atomic_write_json(path, valid)
|
|
self.assertEqual(load_session(session["session_id"])["root_thread_id"], second_id)
|
|
|
|
malformed_successor = json.loads(json.dumps(valid))
|
|
malformed_successor["root_thread_lineage"][1] = "not-an-entry"
|
|
malformed_successor["root_thread_id"] = second_id
|
|
atomic_write_json(path, malformed_successor)
|
|
with self.assertRaisesRegex(ValueError, "lineage entry is incomplete"):
|
|
load_session(session["session_id"])
|
|
|
|
wrong_link = json.loads(json.dumps(valid))
|
|
wrong_link["root_thread_lineage"][0]["successor_thread_id"] = first_id
|
|
atomic_write_json(path, wrong_link)
|
|
with self.assertRaisesRegex(ValueError, "lineage entry is incomplete"):
|
|
load_session(session["session_id"])
|
|
|
|
duplicate = json.loads(json.dumps(valid))
|
|
duplicate["root_thread_lineage"][1]["thread_id"] = first_id
|
|
duplicate["root_thread_lineage"][0]["successor_thread_id"] = first_id
|
|
duplicate["root_thread_id"] = first_id
|
|
atomic_write_json(path, duplicate)
|
|
with self.assertRaisesRegex(ValueError, "lineage entry is invalid"):
|
|
load_session(session["session_id"])
|
|
|
|
staged = json.loads(json.dumps(valid))
|
|
staged["root_thread_transition"] = {
|
|
"from_thread_id": first_id,
|
|
"to_thread_id": "00000000-0000-0000-0000-000000000203",
|
|
"generation": 3,
|
|
"observed_at": "2026-08-20T00:02:00Z",
|
|
"reason": "fresh_context",
|
|
}
|
|
atomic_write_json(path, staged)
|
|
with self.assertRaisesRegex(ValueError, "transition is inconsistent"):
|
|
load_session(session["session_id"])
|
|
finally:
|
|
atomic_write_json(path, original)
|
|
finish_session(session["session_id"], exit_code=0)
|
|
|
|
def test_durable_state_identity_must_match_its_directory(self) -> None:
|
|
with RuntimeSandbox() as box:
|
|
session = create_session(profile="codex-harness-team", cwd=box.workspace)
|
|
session_path = box.state / "sessions" / session["session_id"] / "session.json"
|
|
original = read_json(session_path)
|
|
atomic_write_json(session_path, {**original, "session_id": "redirected-session"})
|
|
try:
|
|
with self.assertRaisesRegex(ValueError, "session state identity mismatch"):
|
|
load_session(session["session_id"])
|
|
finally:
|
|
atomic_write_json(session_path, original)
|
|
finish_session(session["session_id"], exit_code=0)
|
|
|
|
directory = box.state / "jobs" / "directory-job"
|
|
directory.mkdir(parents=True)
|
|
atomic_write_json(
|
|
directory / "metadata.json",
|
|
{
|
|
"schema_version": mmo_runtime.MMO_SCHEMA_VERSION,
|
|
"package_version": mmo_runtime.package_version(),
|
|
"job_id": "redirected-job",
|
|
},
|
|
)
|
|
self.assertEqual(mmo_runtime.iter_jobs(strict=False), [])
|
|
with self.assertRaisesRegex(RuntimeError, "job state identity mismatch"):
|
|
mmo_runtime.iter_jobs(strict=True)
|
|
|
|
def test_declared_artifact_paths_and_hashes_are_correlated(self) -> None:
|
|
with tempfile.TemporaryDirectory(prefix="mmo-artifact-") as temporary:
|
|
cwd = Path(temporary)
|
|
artifact = cwd / "screens" / "render.png"
|
|
artifact.parent.mkdir()
|
|
artifact.write_bytes(b"render bytes")
|
|
digest = hashlib.sha256(artifact.read_bytes()).hexdigest()
|
|
errors, identities = _correlate_artifact_evidence(
|
|
{"image_artifacts": [{"relative_path": "screens/render.png", "sha256": digest}]},
|
|
cwd,
|
|
)
|
|
self.assertEqual(errors, [])
|
|
self.assertEqual(identities[0]["relative_path"], "screens/render.png")
|
|
self.assertEqual(identities[0]["sha256"], digest)
|
|
self.assertEqual(identities[0]["size"], len(b"render bytes"))
|
|
self.assertEqual(identities[0]["media_type"], "image/png")
|
|
|
|
errors, identities = _correlate_artifact_evidence(
|
|
{
|
|
"artifacts": [
|
|
{"relative_path": "screens/render.png", "sha256": "0" * 64},
|
|
{"relative_path": "../escape.png", "sha256": digest},
|
|
{"relative_path": "missing.png", "sha256": digest},
|
|
]
|
|
},
|
|
cwd,
|
|
)
|
|
self.assertEqual(identities, [])
|
|
self.assertTrue(any("hash does not match" in item for item in errors), errors)
|
|
self.assertTrue(any("escapes delegated cwd" in item for item in errors), errors)
|
|
self.assertTrue(any("regular file" in item for item in errors), errors)
|
|
|
|
def test_declared_commands_require_correlated_exit_codes(self) -> None:
|
|
with tempfile.TemporaryDirectory(prefix="mmo-command-evidence-") as temporary:
|
|
events = Path(temporary) / "events.jsonl"
|
|
events.write_text(
|
|
json.dumps(
|
|
{
|
|
"type": "item.completed",
|
|
"item": {
|
|
"type": "command_execution",
|
|
"command": "python -m unittest -v",
|
|
"exit_code": 1,
|
|
},
|
|
}
|
|
)
|
|
+ "\n",
|
|
encoding="utf-8",
|
|
)
|
|
self.assertEqual(
|
|
_correlate_command_evidence(
|
|
{"commands": [{"command": "python -m unittest", "exit_code": 1}]},
|
|
events,
|
|
),
|
|
[],
|
|
)
|
|
wrong_exit = _correlate_command_evidence(
|
|
{"commands": [{"command": "python -m unittest", "exit_code": 0}]},
|
|
events,
|
|
)
|
|
self.assertTrue(any("exit_code" in error for error in wrong_exit), wrong_exit)
|
|
absent = _correlate_command_evidence(
|
|
{"commands": [{"command": "ruff check", "exit_code": 0}]},
|
|
events,
|
|
)
|
|
self.assertTrue(any("absent" in error for error in absent), absent)
|
|
|
|
def test_git_nul_paths_use_the_filesystem_codec_losslessly(self) -> None:
|
|
raw_path = b"non-utf8-\xff"
|
|
result = subprocess.CompletedProcess(["git"], 0, raw_path + b"\0", b"")
|
|
decoded = _nul_paths(result)
|
|
self.assertEqual(decoded, {os.fsdecode(raw_path)})
|
|
self.assertEqual({os.fsencode(item) for item in decoded}, {raw_path})
|
|
|
|
def test_doctor_mcp_probe_performs_negotiated_two_phase_handshake(self) -> None:
|
|
with RuntimeSandbox() as box:
|
|
box.init_git()
|
|
session = create_access_lab_session(cwd=box.workspace)
|
|
mark_session_running(session["session_id"], os.getpid())
|
|
try:
|
|
result = _mcp_handshake(session)
|
|
self.assertTrue(result["passed"], result)
|
|
self.assertEqual(result["initialize"]["result"]["protocolVersion"], "2025-06-18")
|
|
self.assertIn("agent_status", result["tools"])
|
|
finally:
|
|
finish_session(session["session_id"], exit_code=0)
|
|
|
|
def test_http_readiness_requires_a_success_status(self) -> None:
|
|
class Handler(http.server.BaseHTTPRequestHandler):
|
|
def do_GET(self) -> None:
|
|
if self.path == "/health":
|
|
self.send_response(302)
|
|
self.send_header("Location", "/ready")
|
|
else:
|
|
self.send_response(204)
|
|
self.end_headers()
|
|
|
|
def log_message(self, _format: str, *args: Any) -> None:
|
|
pass
|
|
|
|
server = http.server.ThreadingHTTPServer(("127.0.0.1", 0), Handler)
|
|
thread = threading.Thread(target=server.serve_forever, daemon=True)
|
|
thread.start()
|
|
try:
|
|
self.assertTrue(http_ready(f"http://127.0.0.1:{server.server_port}/ready"))
|
|
self.assertFalse(http_ready(f"http://127.0.0.1:{server.server_port}/health"))
|
|
finally:
|
|
server.shutdown()
|
|
server.server_close()
|
|
thread.join(timeout=2)
|
|
|
|
def test_cancellation_finalizer_and_orphan_group_cleanup_are_authoritative(self) -> None:
|
|
with RuntimeSandbox() as box:
|
|
directory = box.state / "jobs" / "finalizer-test"
|
|
directory.mkdir(parents=True)
|
|
atomic_write_json(
|
|
directory / "metadata.json",
|
|
{
|
|
"schema_version": mmo_runtime.MMO_SCHEMA_VERSION,
|
|
"package_version": mmo_runtime.package_version(),
|
|
"job_id": directory.name,
|
|
"status": "cancelling",
|
|
"cancel_reason": "operator request",
|
|
},
|
|
)
|
|
final, cancelled = finalize(
|
|
directory,
|
|
status="completed",
|
|
exit_code=0,
|
|
error="late completion must not win",
|
|
)
|
|
self.assertTrue(cancelled)
|
|
self.assertEqual(final["status"], "cancelled")
|
|
self.assertEqual(final["cancel_reason"], "operator request")
|
|
self.assertNotIn("error", final)
|
|
failure_usage = {"input_tokens": 17, "event_count": 1}
|
|
self.assertEqual(
|
|
_fail(directory, 1, "cancelled worker failure", usage=failure_usage), 130
|
|
)
|
|
cancelled = read_json(directory / "metadata.json")
|
|
self.assertEqual(cancelled["warning"], "cancelled worker failure")
|
|
self.assertEqual(cancelled["usage"], failure_usage)
|
|
self.assertNotIn("error", cancelled)
|
|
|
|
child_path = box.root / "orphan.pid"
|
|
process = subprocess.Popen(
|
|
[
|
|
"sh",
|
|
"-c",
|
|
f'nohup sleep 60 >/dev/null 2>&1 & echo "$!" > {child_path}',
|
|
],
|
|
start_new_session=True,
|
|
)
|
|
process.wait(timeout=5)
|
|
child_pid = int(child_path.read_text(encoding="utf-8"))
|
|
try:
|
|
self.assertTrue(process_alive(child_pid))
|
|
_terminate_and_reap(process)
|
|
self.assertFalse(process_alive(child_pid))
|
|
finally:
|
|
terminate_process_group(process.pid, grace_seconds=0.2)
|
|
|
|
def test_jsonl_append_serializes_partial_concurrent_writes(self) -> None:
|
|
with RuntimeSandbox() as box:
|
|
path = box.root / "concurrent.jsonl"
|
|
real_write = os.write
|
|
barrier = threading.Barrier(12)
|
|
failures: list[BaseException] = []
|
|
|
|
def partial_write(fd: int, data: bytes) -> int:
|
|
written = real_write(fd, data[:5])
|
|
time.sleep(0.0005)
|
|
return written
|
|
|
|
def writer(index: int) -> None:
|
|
try:
|
|
barrier.wait(timeout=5)
|
|
append_jsonl(path, {"index": index, "payload": str(index) * 80})
|
|
except BaseException as exc:
|
|
failures.append(exc)
|
|
|
|
with mock.patch.object(mmo_util.os, "write", side_effect=partial_write):
|
|
threads = [threading.Thread(target=writer, args=(index,)) for index in range(12)]
|
|
for thread in threads:
|
|
thread.start()
|
|
for thread in threads:
|
|
thread.join(timeout=10)
|
|
self.assertFalse(failures)
|
|
records = [json.loads(line) for line in path.read_text(encoding="utf-8").splitlines()]
|
|
self.assertEqual(len(records), 12)
|
|
self.assertEqual({item["index"] for item in records}, set(range(12)))
|
|
|
|
def test_json_io_rejects_non_standard_and_non_finite_numbers(self) -> None:
|
|
with RuntimeSandbox() as box:
|
|
source = box.root / "external.json"
|
|
for payload in (
|
|
'{"value": NaN}',
|
|
'{"value": Infinity}',
|
|
'{"value": 1e400}',
|
|
'{"value": 1, "value": 2}',
|
|
'{"value": "\\ud800"}',
|
|
):
|
|
source.write_text(payload, encoding="utf-8")
|
|
with self.subTest(payload=payload), self.assertRaises(ValueError):
|
|
read_json(source)
|
|
|
|
target = box.root / "state.json"
|
|
target.write_text('{"preserved": true}\n', encoding="utf-8")
|
|
with self.assertRaises(ValueError):
|
|
atomic_write_json(target, {"value": float("nan")})
|
|
self.assertEqual(target.read_text(encoding="utf-8"), '{"preserved": true}\n')
|
|
|
|
journal = box.root / "events.jsonl"
|
|
with self.assertRaises(ValueError):
|
|
append_jsonl(journal, {"value": float("inf")})
|
|
self.assertFalse(journal.exists())
|
|
|
|
def test_git_root_preserves_legal_trailing_whitespace(self) -> None:
|
|
with RuntimeSandbox() as box:
|
|
repository = box.workspace / "repository-with-trailing-space "
|
|
subprocess.run(["git", "init", "-q", str(repository)], check=True)
|
|
self.assertEqual(_git_root(repository), repository.resolve())
|
|
|
|
def test_isolated_patch_scope_checks_both_sides_of_a_committed_rename(self) -> None:
|
|
with RuntimeSandbox() as box:
|
|
box.init_git()
|
|
directory = box.root / "isolated-job"
|
|
directory.mkdir()
|
|
isolation = create_isolated_worktree(box.workspace.resolve(), directory, ["inside"], [])
|
|
metadata = {**isolation, "write_scope": ["inside"]}
|
|
worktree = Path(str(isolation["worktree_root"]))
|
|
try:
|
|
(worktree / "inside").mkdir()
|
|
subprocess.run(
|
|
["git", "-C", str(worktree), "mv", "README.md", "inside/README.md"],
|
|
check=True,
|
|
)
|
|
subprocess.run(
|
|
["git", "-C", str(worktree), "commit", "-qm", "worker rename"],
|
|
check=True,
|
|
)
|
|
errors, artifacts, patch = capture_isolated_patch(
|
|
metadata, directory, Path(str(isolation["cwd"]))
|
|
)
|
|
self.assertTrue(any("README.md" in item for item in errors), errors)
|
|
self.assertEqual(artifacts, [])
|
|
self.assertIsNone(patch)
|
|
finally:
|
|
remove_isolated_worktree(metadata)
|
|
|
|
def test_gateway_startup_error_preserves_immediate_exit_log(self) -> None:
|
|
with RuntimeSandbox() as box:
|
|
snapshot = compile_profile("access-efficient-escalation-lab")
|
|
executable = box.root / "failing-switchyard"
|
|
executable.write_text(
|
|
"#!/bin/sh\necho deterministic-startup-failure\nexit 23\n",
|
|
encoding="utf-8",
|
|
)
|
|
executable.chmod(0o755)
|
|
|
|
def wait_for_exit(pid: int) -> None:
|
|
deadline = time.monotonic() + 5
|
|
while process_alive(pid) and time.monotonic() < deadline:
|
|
time.sleep(0.01)
|
|
return None
|
|
|
|
with (
|
|
mock.patch.object(mmo_gateway, "_binary", return_value=str(executable)),
|
|
mock.patch.object(mmo_gateway, "switchyard_version", return_value="0.2.0"),
|
|
mock.patch.object(
|
|
mmo_gateway,
|
|
"process_start_token",
|
|
side_effect=wait_for_exit,
|
|
),
|
|
):
|
|
with self.assertRaisesRegex(RuntimeError, "deterministic-startup-failure"):
|
|
ensure_gateway(snapshot["manifest"]["snapshot_hash"])
|
|
|
|
def test_gateway_launch_is_reaped_when_state_persistence_fails(self) -> None:
|
|
with RuntimeSandbox():
|
|
snapshot = compile_profile("adaptive-engineering")
|
|
launched: list[subprocess.Popen[Any]] = []
|
|
real_popen = mmo_gateway.subprocess.Popen
|
|
|
|
def capture_process(*args: Any, **kwargs: Any) -> subprocess.Popen[Any]:
|
|
process = real_popen(*args, **kwargs)
|
|
launched.append(process)
|
|
return process
|
|
|
|
with (
|
|
mock.patch.object(mmo_gateway.subprocess, "Popen", side_effect=capture_process),
|
|
mock.patch.object(mmo_gateway, "switchyard_version", return_value="0.2.0"),
|
|
mock.patch.object(
|
|
mmo_gateway,
|
|
"atomic_write_json",
|
|
side_effect=OSError("synthetic gateway state failure"),
|
|
),
|
|
self.assertRaisesRegex(OSError, "synthetic gateway state failure"),
|
|
):
|
|
ensure_gateway(snapshot["manifest"]["snapshot_hash"])
|
|
|
|
self.assertEqual(len(launched), 1)
|
|
self.assertIsNotNone(launched[0].poll())
|
|
self.assertNotIn(launched[0].pid, mmo_gateway._GATEWAY_PROCESSES)
|
|
|
|
def test_invalid_gateway_state_cannot_redirect_process_control(self) -> None:
|
|
with RuntimeSandbox() as box:
|
|
snapshot = compile_profile("high-confidence-debugging")
|
|
snapshot_hash = snapshot["manifest"]["snapshot_hash"]
|
|
gateway_hash = snapshot["manifest"]["gateway_hash"]
|
|
directory = box.state / "gateways" / gateway_hash
|
|
directory.mkdir(parents=True)
|
|
atomic_write_json(
|
|
directory / "gateway.json",
|
|
{
|
|
"schema_version": 1,
|
|
"gateway_hash": "0" * 64,
|
|
"pid": os.getpid(),
|
|
"process_start_token": process_start_token(os.getpid()),
|
|
},
|
|
)
|
|
status = gateway_status(snapshot_hash)
|
|
self.assertEqual(status["status"], "invalid")
|
|
with self.assertRaisesRegex(RuntimeError, "invalid gateway state"):
|
|
stop_gateway(snapshot_hash)
|
|
self.assertTrue(process_alive(os.getpid()))
|
|
|
|
def test_gateway_state_cannot_redirect_health_or_model_requests(self) -> None:
|
|
with RuntimeSandbox() as box:
|
|
snapshot = compile_profile("adaptive-engineering")
|
|
snapshot_hash = snapshot["manifest"]["snapshot_hash"]
|
|
gateway_hash = snapshot["manifest"]["gateway_hash"]
|
|
directory = box.state / "gateways" / gateway_hash
|
|
directory.mkdir(parents=True)
|
|
cases = {
|
|
"redirected urls": {
|
|
"host": "127.0.0.1",
|
|
"base_url": "https://attacker.invalid/v1",
|
|
"health_url": "https://attacker.invalid/health",
|
|
},
|
|
"self-consistent non-loopback endpoint": {
|
|
"host": "192.0.2.1",
|
|
"base_url": "http://192.0.2.1:42000/v1",
|
|
"health_url": "http://192.0.2.1:42000/health",
|
|
},
|
|
}
|
|
for label, endpoint in cases.items():
|
|
with self.subTest(label):
|
|
atomic_write_json(
|
|
directory / "gateway.json",
|
|
{
|
|
"schema_version": mmo_runtime.MMO_SCHEMA_VERSION,
|
|
"gateway_hash": gateway_hash,
|
|
"snapshot_hash": snapshot_hash,
|
|
"snapshot_hashes": [snapshot_hash],
|
|
"profile_ids": [snapshot["manifest"]["profile_id"]],
|
|
"pid": os.getpid(),
|
|
"process_start_token": process_start_token(os.getpid()),
|
|
"port": 42000,
|
|
**endpoint,
|
|
},
|
|
)
|
|
with mock.patch.object(mmo_gateway, "http_ready") as ready:
|
|
status = gateway_status(snapshot_hash)
|
|
self.assertEqual(status["status"], "invalid")
|
|
self.assertIn("identity", status["error"])
|
|
ready.assert_not_called()
|
|
|
|
def test_atomic_batch_is_all_or_nothing_and_success_has_one_batch_id(self) -> None:
|
|
with RuntimeSandbox() as box:
|
|
session = create_access_lab_session(cwd=box.workspace)
|
|
mark_session_running(session["session_id"], os.getpid())
|
|
try:
|
|
with self.assertRaises(ValueError):
|
|
spawn_jobs(
|
|
[
|
|
{
|
|
"agent": "literal_scout",
|
|
"literal_task": {
|
|
"operation": "summarize_supplied",
|
|
"text": "README.md exact evidence fixture",
|
|
},
|
|
},
|
|
{
|
|
"agent": "literal_scout",
|
|
"literal_task": {"operation": "architecture"},
|
|
},
|
|
],
|
|
session_id=session["session_id"],
|
|
caller_agent=session["root_agent"],
|
|
caller_job_id=None,
|
|
)
|
|
self.assertEqual(list_jobs(session_id=session["session_id"]), [])
|
|
|
|
admitted = spawn_jobs(
|
|
[
|
|
{
|
|
"agent": "literal_scout",
|
|
"literal_task": {
|
|
"operation": "summarize_supplied",
|
|
"text": "README.md exact evidence fixture",
|
|
},
|
|
},
|
|
{
|
|
"agent": "flagship_escalation",
|
|
"task_kind": "analysis",
|
|
"task": "Analyze the repository fixture read-only and return bounded evidence.",
|
|
},
|
|
],
|
|
session_id=session["session_id"],
|
|
caller_agent=session["root_agent"],
|
|
caller_job_id=None,
|
|
)
|
|
self.assertTrue(admitted["atomic"])
|
|
self.assertEqual(admitted["rejected"], [])
|
|
self.assertEqual(len(admitted["accepted"]), 2)
|
|
batch_ids = {item["batch_id"] for item in admitted["accepted"]}
|
|
self.assertEqual(batch_ids, {admitted["batch_id"]})
|
|
job_ids = [item["job_id"] for item in admitted["accepted"]]
|
|
waited = wait_for_jobs(
|
|
job_ids,
|
|
session_id=session["session_id"],
|
|
timeout_seconds=20,
|
|
include_results=False,
|
|
)
|
|
self.assertFalse(waited["unfinished"], waited)
|
|
completed = [load_job(job_id) for job_id in job_ids]
|
|
self.assertTrue(all(item["status"] == "completed" for item in completed), completed)
|
|
finally:
|
|
finish_session(session["session_id"], exit_code=0)
|
|
|
|
def test_atomic_batch_launch_failure_rolls_back_every_member(self) -> None:
|
|
with RuntimeSandbox() as box:
|
|
profile = box.root / "two-spawn-budget"
|
|
shutil.copytree(ROOT / "profiles" / "access-efficient-escalation-lab", profile)
|
|
profile_path = profile / "profile.toml"
|
|
profile_data = read_toml(profile_path)
|
|
profile_path.write_text(toml_dumps(profile_data), encoding="utf-8")
|
|
session = create_access_lab_session(cwd=box.workspace, profile=profile)
|
|
mark_session_running(session["session_id"], os.getpid())
|
|
original_popen = subprocess.Popen
|
|
calls = 0
|
|
|
|
def flaky_popen(*args: Any, **kwargs: Any) -> subprocess.Popen[Any]:
|
|
nonlocal calls
|
|
calls += 1
|
|
if calls == 2:
|
|
raise OSError("synthetic second-runner launch failure")
|
|
return original_popen(*args, **kwargs)
|
|
|
|
try:
|
|
with mock.patch.object(mmo_runtime.subprocess, "Popen", side_effect=flaky_popen):
|
|
with self.assertRaisesRegex(RuntimeError, "atomic batch launch rolled back"):
|
|
spawn_jobs(
|
|
[
|
|
{
|
|
"agent": "literal_scout",
|
|
"literal_task": {
|
|
"operation": "summarize_supplied",
|
|
"text": "README.md exact evidence fixture FAKE_SLEEP=5",
|
|
},
|
|
},
|
|
{
|
|
"agent": "flagship_escalation",
|
|
"task_kind": "analysis",
|
|
"task": "Analyze the fixture read-only and report exact evidence.",
|
|
},
|
|
],
|
|
session_id=session["session_id"],
|
|
caller_agent=session["root_agent"],
|
|
caller_job_id=None,
|
|
)
|
|
rows = list_jobs(session_id=session["session_id"])
|
|
self.assertEqual(len(rows), 2)
|
|
self.assertTrue(all(row["status"] == "failed" for row in rows), rows)
|
|
self.assertTrue(
|
|
all(load_job(row["job_id"])["atomic_batch_rolled_back"] for row in rows)
|
|
)
|
|
self.assertFalse(
|
|
any(row["status"] in mmo_runtime.ACTIVE_JOB_STATUSES for row in rows)
|
|
)
|
|
retry = spawn_jobs(
|
|
[
|
|
{
|
|
"agent": "literal_scout",
|
|
"literal_task": {
|
|
"operation": "summarize_supplied",
|
|
"text": "README.md exact evidence after rollback",
|
|
},
|
|
},
|
|
{
|
|
"agent": "flagship_escalation",
|
|
"task_kind": "analysis",
|
|
"task": "Analyze the fixture after rollback.",
|
|
},
|
|
],
|
|
session_id=session["session_id"],
|
|
caller_agent=session["root_agent"],
|
|
caller_job_id=None,
|
|
)
|
|
self.assertEqual(len(retry["accepted"]), 2)
|
|
wait_for_jobs(
|
|
[item["job_id"] for item in retry["accepted"]],
|
|
session_id=session["session_id"],
|
|
timeout_seconds=20,
|
|
include_results=False,
|
|
)
|
|
finally:
|
|
finish_session(session["session_id"], exit_code=0)
|
|
|
|
def test_finishing_and_cancelling_sessions_close_admission_before_teardown(self) -> None:
|
|
with RuntimeSandbox() as box:
|
|
session = create_access_lab_session(cwd=box.workspace)
|
|
mark_session_running(session["session_id"], os.getpid())
|
|
entered = threading.Event()
|
|
release = threading.Event()
|
|
original_iter_jobs = mmo_runtime.iter_jobs
|
|
|
|
def blocked_iter_jobs() -> list[dict]:
|
|
entered.set()
|
|
self.assertTrue(release.wait(10), "test failed to release finish_session")
|
|
return original_iter_jobs()
|
|
|
|
outcome: list[dict[str, Any]] = []
|
|
with mock.patch.object(mmo_runtime, "iter_jobs", side_effect=blocked_iter_jobs):
|
|
thread = threading.Thread(
|
|
target=lambda: outcome.append(
|
|
finish_session(session["session_id"], exit_code=0)
|
|
),
|
|
daemon=True,
|
|
)
|
|
thread.start()
|
|
self.assertTrue(entered.wait(5), "finish_session did not enter teardown")
|
|
self.assertEqual(load_session(session["session_id"])["status"], "finishing")
|
|
with self.assertRaisesRegex(RuntimeError, "not admitting"):
|
|
spawn_job(
|
|
session_id=session["session_id"],
|
|
caller_agent=session["root_agent"],
|
|
caller_job_id=None,
|
|
agent_id="literal_scout",
|
|
literal_task={
|
|
"operation": "summarize_supplied",
|
|
"text": "README.md after the session began finishing",
|
|
},
|
|
)
|
|
release.set()
|
|
thread.join(10)
|
|
self.assertFalse(thread.is_alive())
|
|
self.assertEqual(outcome[0]["status"], "completed")
|
|
|
|
second = create_access_lab_session(cwd=box.workspace)
|
|
mark_session_running(second["session_id"], os.getpid())
|
|
cancelled = cancel_session(second["session_id"])
|
|
self.assertEqual(cancelled["session"]["status"], "cancelled")
|
|
with self.assertRaisesRegex(RuntimeError, "not admitting"):
|
|
spawn_job(
|
|
session_id=second["session_id"],
|
|
caller_agent=second["root_agent"],
|
|
caller_job_id=None,
|
|
agent_id="literal_scout",
|
|
literal_task={
|
|
"operation": "summarize_supplied",
|
|
"text": "README.md after cancellation became terminal",
|
|
},
|
|
)
|
|
|
|
def test_nested_delegation_depth_ancestor_and_visibility_are_enforced(self) -> None:
|
|
with RuntimeSandbox() as box:
|
|
custom = box.root / "bounded-research-cycle-test"
|
|
shutil.copytree(ROOT / "profiles" / "bounded-research-organization-lab", custom)
|
|
profile_path = custom / "profile.toml"
|
|
profile_data = read_toml(profile_path)
|
|
profile_data["coordination"]["max_depth"] = 3
|
|
profile_data["coordination"]["max_active_agents"] = 4
|
|
profile_data["agents"]["source_scout"]["can_spawn"] = ["research_lead"]
|
|
profile_path.write_text(toml_dumps(profile_data), encoding="utf-8")
|
|
session = create_session(profile=custom, cwd=box.workspace)
|
|
mark_session_running(session["session_id"], os.getpid())
|
|
try:
|
|
parent = spawn_job(
|
|
session_id=session["session_id"],
|
|
caller_agent=session["root_agent"],
|
|
caller_job_id=None,
|
|
agent_id="research_lead",
|
|
task_kind="research_organization",
|
|
task="Organize one bounded supplied-source investigation independently. FAKE_SLEEP=4",
|
|
)
|
|
child = spawn_job(
|
|
session_id=session["session_id"],
|
|
caller_agent="research_lead",
|
|
caller_job_id=parent["job_id"],
|
|
agent_id="source_scout",
|
|
task_kind="research",
|
|
task="Collect bounded source evidence for the parent without making recommendations. FAKE_SLEEP=4",
|
|
)
|
|
with self.assertRaisesRegex(RuntimeError, "ancestor-role repetition"):
|
|
spawn_job(
|
|
session_id=session["session_id"],
|
|
caller_agent="source_scout",
|
|
caller_job_id=child["job_id"],
|
|
agent_id="research_lead",
|
|
task_kind="research_synthesis",
|
|
task="Attempt an ancestor-role cycle that the bounded hierarchy must reject.",
|
|
)
|
|
|
|
sibling = spawn_job(
|
|
session_id=session["session_id"],
|
|
caller_agent=session["root_agent"],
|
|
caller_job_id=None,
|
|
agent_id="source_scout",
|
|
task_kind="source_verification",
|
|
task="Perform a separate bounded source verification for visibility testing.",
|
|
)
|
|
wait_for_jobs(
|
|
[sibling["job_id"]],
|
|
session_id=session["session_id"],
|
|
timeout_seconds=20,
|
|
include_results=False,
|
|
)
|
|
controlled_result = read_result(
|
|
sibling["job_id"],
|
|
session_id=session["session_id"],
|
|
caller_job_id=parent["job_id"],
|
|
caller_agent="research_lead",
|
|
)
|
|
self.assertTrue(controlled_result["content"])
|
|
root_result = read_result(sibling["job_id"], session_id=session["session_id"])
|
|
self.assertTrue(root_result["content"])
|
|
cancel_job(parent["job_id"], session_id=session["session_id"], cascade=True)
|
|
wait_for_jobs(
|
|
[parent["job_id"], child["job_id"]],
|
|
session_id=session["session_id"],
|
|
timeout_seconds=15,
|
|
include_results=False,
|
|
)
|
|
finally:
|
|
finish_session(session["session_id"], exit_code=0)
|
|
|
|
def test_weighted_resource_capacity_is_reserved_for_whole_batch(self) -> None:
|
|
with RuntimeSandbox() as box:
|
|
custom = box.root / "one-unit-opencode"
|
|
shutil.copytree(ROOT / "profiles" / "access-efficient-escalation-lab", custom)
|
|
profile_path = custom / "profile.toml"
|
|
profile_data = read_toml(profile_path)
|
|
profile_data["agents"]["routine_engineer"]["max_active"] = 2
|
|
profile_path.write_text(toml_dumps(profile_data), encoding="utf-8")
|
|
(custom / "catalog.toml").write_text(
|
|
f"""schema_version = {mmo_runtime.MMO_SCHEMA_VERSION}
|
|
|
|
[resources.opencode_go]
|
|
max_active = 1
|
|
lock_key = "provider:opencode-go"
|
|
description = "One test unit"
|
|
""",
|
|
encoding="utf-8",
|
|
)
|
|
session = create_session(profile=custom, cwd=box.workspace)
|
|
mark_session_running(session["session_id"], os.getpid())
|
|
try:
|
|
with self.assertRaisesRegex(RuntimeError, "resource group"):
|
|
spawn_jobs(
|
|
[
|
|
{
|
|
"agent": "routine_engineer",
|
|
"task_kind": "test",
|
|
"task": "Perform one bounded implementation analysis without writing.",
|
|
},
|
|
{
|
|
"agent": "routine_engineer",
|
|
"task_kind": "test",
|
|
"task": "Perform a second bounded implementation analysis without writing.",
|
|
},
|
|
],
|
|
session_id=session["session_id"],
|
|
caller_agent=session["root_agent"],
|
|
caller_job_id=None,
|
|
)
|
|
self.assertEqual(list_jobs(session_id=session["session_id"]), [])
|
|
finally:
|
|
finish_session(session["session_id"], exit_code=0)
|
|
|
|
def test_role_max_active_is_enforced_across_profile_sessions(self) -> None:
|
|
with RuntimeSandbox() as box:
|
|
custom = box.root / "global-role-limit"
|
|
shutil.copytree(ROOT / "profiles" / "access-efficient-escalation-lab", custom)
|
|
profile_path = custom / "profile.toml"
|
|
first = create_session(profile=custom, cwd=box.workspace)
|
|
profile_data = read_toml(profile_path)
|
|
profile_data["description"] += " (second immutable snapshot)"
|
|
profile_path.write_text(toml_dumps(profile_data), encoding="utf-8")
|
|
second = create_session(profile=custom, cwd=box.workspace)
|
|
self.assertEqual(first["profile_version"], second["profile_version"])
|
|
self.assertNotEqual(first["snapshot_hash"], second["snapshot_hash"])
|
|
for session in (first, second):
|
|
mark_session_running(session["session_id"], os.getpid())
|
|
try:
|
|
active = spawn_job(
|
|
session_id=first["session_id"],
|
|
caller_agent=first["root_agent"],
|
|
caller_job_id=None,
|
|
agent_id="routine_engineer",
|
|
task_kind="test",
|
|
task="Remain active while global role admission is checked. FAKE_SLEEP=5",
|
|
)
|
|
with self.assertRaisesRegex(RuntimeError, "active routine_engineer limit"):
|
|
spawn_job(
|
|
session_id=second["session_id"],
|
|
caller_agent=second["root_agent"],
|
|
caller_job_id=None,
|
|
agent_id="routine_engineer",
|
|
task_kind="test",
|
|
task="Attempt a second instance of this globally bounded profile role.",
|
|
)
|
|
cancel_job(active["job_id"], session_id=first["session_id"])
|
|
wait_for_jobs(
|
|
[active["job_id"]],
|
|
session_id=first["session_id"],
|
|
timeout_seconds=15,
|
|
include_results=False,
|
|
)
|
|
finally:
|
|
for session in (first, second):
|
|
for job in list_jobs(session_id=session["session_id"]):
|
|
if job.get("status") in {"queued", "running", "cancelling"}:
|
|
cancel_job(job["job_id"], session_id=session["session_id"])
|
|
finish_session(session["session_id"], exit_code=0)
|
|
|
|
def test_retired_per_spawn_execution_timeout_is_rejected(self) -> None:
|
|
with RuntimeSandbox() as box:
|
|
session = create_access_lab_session(cwd=box.workspace)
|
|
mark_session_running(session["session_id"], os.getpid())
|
|
try:
|
|
with self.assertRaisesRegex(
|
|
ValueError, "unknown job request fields: timeout_seconds"
|
|
):
|
|
spawn_jobs(
|
|
[
|
|
{
|
|
"agent": "flagship_escalation",
|
|
"task_kind": "analysis",
|
|
"task": "Analyze one bounded concern without changing the workspace.",
|
|
"timeout_seconds": 30,
|
|
}
|
|
],
|
|
session_id=session["session_id"],
|
|
caller_agent=session["root_agent"],
|
|
caller_job_id=None,
|
|
)
|
|
self.assertEqual(list_jobs(session_id=session["session_id"]), [])
|
|
finally:
|
|
finish_session(session["session_id"], exit_code=0)
|
|
|
|
def test_cascade_cancellation_closes_child_admission_atomically(self) -> None:
|
|
with RuntimeSandbox() as box:
|
|
session = create_session(profile="bounded-research-organization-lab", cwd=box.workspace)
|
|
mark_session_running(session["session_id"], os.getpid())
|
|
parent = spawn_job(
|
|
session_id=session["session_id"],
|
|
caller_agent=session["root_agent"],
|
|
caller_job_id=None,
|
|
agent_id="research_lead",
|
|
task_kind="research_organization",
|
|
task="Remain active while cancellation and child admission race. FAKE_SLEEP=5",
|
|
)
|
|
entered = threading.Event()
|
|
release = threading.Event()
|
|
original_descendants = mmo_runtime._descendant_job_ids
|
|
cancel_errors: list[BaseException] = []
|
|
spawn_outcomes: list[str] = []
|
|
|
|
def gated_descendants(job_id: str, *, lock_held: bool = False) -> list[str]:
|
|
entered.set()
|
|
if not release.wait(timeout=5):
|
|
raise TimeoutError("test did not release cancellation discovery")
|
|
return original_descendants(job_id, lock_held=lock_held)
|
|
|
|
def cancel_parent() -> None:
|
|
try:
|
|
cancel_job(parent["job_id"], session_id=session["session_id"], cascade=True)
|
|
except BaseException as exc: # captured for assertion in the test thread
|
|
cancel_errors.append(exc)
|
|
|
|
def spawn_child() -> None:
|
|
try:
|
|
spawn_job(
|
|
session_id=session["session_id"],
|
|
caller_agent="research_lead",
|
|
caller_job_id=parent["job_id"],
|
|
agent_id="source_scout",
|
|
task_kind="research",
|
|
task="This bounded child must not cross the parent's cancellation boundary.",
|
|
)
|
|
spawn_outcomes.append("accepted")
|
|
except Exception as exc:
|
|
spawn_outcomes.append(f"{type(exc).__name__}: {exc}")
|
|
|
|
try:
|
|
with mock.patch("mmo_runtime._descendant_job_ids", side_effect=gated_descendants):
|
|
cancel_thread = threading.Thread(target=cancel_parent)
|
|
cancel_thread.start()
|
|
self.assertTrue(entered.wait(timeout=3))
|
|
spawn_thread = threading.Thread(target=spawn_child)
|
|
spawn_thread.start()
|
|
time.sleep(0.2)
|
|
self.assertTrue(spawn_thread.is_alive())
|
|
release.set()
|
|
cancel_thread.join(timeout=5)
|
|
spawn_thread.join(timeout=5)
|
|
self.assertFalse(cancel_errors, cancel_errors)
|
|
self.assertEqual(len(spawn_outcomes), 1)
|
|
self.assertIn("non-admitting agent job", spawn_outcomes[0])
|
|
finally:
|
|
release.set()
|
|
finish_session(session["session_id"], exit_code=0)
|
|
|
|
def test_nested_session_roots_share_absolute_write_scope_leases(self) -> None:
|
|
with RuntimeSandbox() as box:
|
|
nested = box.workspace / "nested"
|
|
nested.mkdir()
|
|
(nested / "fixture.txt").write_text("nested fixture\n", encoding="utf-8")
|
|
box.init_git()
|
|
custom = box.root / "cross-session-scope"
|
|
shutil.copytree(ROOT / "profiles" / "access-efficient-escalation-lab", custom)
|
|
profile_path = custom / "profile.toml"
|
|
profile_data = read_toml(profile_path)
|
|
profile_data["agents"]["routine_engineer"]["max_active"] = 2
|
|
profile_data["coordination"]["max_active_writers"] = 2
|
|
profile_path.write_text(toml_dumps(profile_data), encoding="utf-8")
|
|
outer = create_session(profile=custom, cwd=box.workspace)
|
|
inner = create_session(profile=custom, cwd=nested)
|
|
mark_session_running(outer["session_id"], os.getpid())
|
|
mark_session_running(inner["session_id"], os.getpid())
|
|
try:
|
|
first = spawn_job(
|
|
session_id=outer["session_id"],
|
|
caller_agent=outer["root_agent"],
|
|
caller_job_id=None,
|
|
agent_id="routine_engineer",
|
|
task_kind="implement",
|
|
task="Hold the nested scope while a second session tests admission. FAKE_SLEEP=4",
|
|
mode="workspace-write",
|
|
write_scope_values=["nested"],
|
|
)
|
|
with self.assertRaisesRegex(RuntimeError, "write scope conflict"):
|
|
spawn_job(
|
|
session_id=inner["session_id"],
|
|
caller_agent=inner["root_agent"],
|
|
caller_job_id=None,
|
|
agent_id="routine_engineer",
|
|
task_kind="implement",
|
|
task="Attempt an overlapping write lease from the nested session root.",
|
|
mode="workspace-write",
|
|
write_scope_values=["."],
|
|
)
|
|
cancel_job(first["job_id"], session_id=outer["session_id"])
|
|
finally:
|
|
finish_session(outer["session_id"], exit_code=0)
|
|
finish_session(inner["session_id"], exit_code=0)
|
|
|
|
def test_builtin_only_skips_gateway_and_gateway_is_reused_then_reaped(self) -> None:
|
|
with RuntimeSandbox() as box:
|
|
builtin = create_session(profile="codex-harness-team", cwd=box.workspace)
|
|
try:
|
|
self.assertIsNone(builtin.get("gateway_base_url"))
|
|
self.assertIsNone(builtin.get("gateway_pid"))
|
|
finally:
|
|
finish_session(builtin["session_id"], exit_code=0)
|
|
|
|
first = create_session(profile="adaptive-engineering", cwd=box.workspace)
|
|
custom = box.root / "adaptive-gateway-reuse"
|
|
shutil.copytree(ROOT / "profiles" / "adaptive-engineering", custom)
|
|
profile_path = custom / "profile.toml"
|
|
profile_data = read_toml(profile_path)
|
|
profile_data["id"] = "adaptive-gateway-reuse"
|
|
profile_path.write_text(toml_dumps(profile_data), encoding="utf-8")
|
|
second = create_session(profile=custom, cwd=box.workspace)
|
|
self.assertNotEqual(first["snapshot_hash"], second["snapshot_hash"])
|
|
self.assertEqual(first["gateway_hash"], second["gateway_hash"])
|
|
self.assertEqual(first["gateway_pid"], second["gateway_pid"])
|
|
finish_session(first["session_id"], exit_code=0)
|
|
finish_session(second["session_id"], exit_code=0)
|
|
time.sleep(1.2)
|
|
stopped = stop_idle_gateways()
|
|
self.assertIn(first["snapshot_hash"], stopped)
|
|
self.assertEqual(gateway_status(first["snapshot_hash"])["status"], "stopped")
|
|
self.assertEqual(gateway_status(second["snapshot_hash"])["status"], "stopped")
|
|
|
|
def test_resume_restores_gateway_before_refreshing_pinned_homes(self) -> None:
|
|
with RuntimeSandbox() as box:
|
|
session = mmo_runtime._start_root_runner(
|
|
create_session(profile="adaptive-engineering", cwd=box.workspace)
|
|
)
|
|
snapshot_hash = str(session["snapshot_hash"])
|
|
prior_gateway_pid = int(session["gateway_pid"])
|
|
prior_gateway_base_url = str(session["gateway_base_url"])
|
|
prior_app_server_pid = int(session["root_app_server_pid"])
|
|
mmo_runtime.detach_session(session["session_id"])
|
|
stop_gateway(snapshot_hash)
|
|
self.assertEqual(gateway_status(snapshot_hash)["status"], "stopped")
|
|
real_refresh = mmo_runtime.refresh_session_homes
|
|
|
|
def refresh_after_gateway(*args: Any, **kwargs: Any) -> dict[str, Any]:
|
|
live = gateway_status(snapshot_hash)
|
|
self.assertEqual(live["status"], "running")
|
|
self.assertEqual(kwargs["gateway_base_url"], live["base_url"])
|
|
return real_refresh(*args, **kwargs)
|
|
|
|
with mock.patch.object(
|
|
mmo_runtime,
|
|
"refresh_session_homes",
|
|
side_effect=refresh_after_gateway,
|
|
):
|
|
resumed = mmo_runtime.begin_resume_run(session["session_id"])
|
|
self.assertEqual(resumed["status"], "running")
|
|
self.assertNotEqual(resumed["gateway_pid"], prior_gateway_pid)
|
|
if resumed["gateway_base_url"] != prior_gateway_base_url:
|
|
self.assertNotEqual(resumed["root_app_server_pid"], prior_app_server_pid)
|
|
else:
|
|
self.assertEqual(resumed["root_app_server_pid"], prior_app_server_pid)
|
|
self.assertEqual(gateway_status(snapshot_hash)["status"], "running")
|
|
mmo_runtime.stop_session(session["session_id"], grace_seconds=0)
|
|
|
|
def test_resume_does_not_clear_an_app_server_replacement_published_during_stop(
|
|
self,
|
|
) -> None:
|
|
with RuntimeSandbox() as box:
|
|
session = mmo_runtime._start_root_runner(
|
|
create_session(
|
|
profile="codex-harness-team",
|
|
cwd=box.workspace,
|
|
session_kind="interactive",
|
|
)
|
|
)
|
|
prior_token = str(session["root_app_server_start_token"])
|
|
replacement: dict[str, Any] = {}
|
|
real_terminate = mmo_runtime.terminate_recorded_process_group
|
|
mmo_runtime.detach_session(session["session_id"])
|
|
mmo_runtime.update_session(
|
|
session["session_id"],
|
|
gateway_base_url="http://stale.invalid/v1",
|
|
switchyard_version="0.2.0",
|
|
root_app_server_lifecycle_timeout_seconds=5.0,
|
|
)
|
|
|
|
def terminate_after_replacement(
|
|
data: dict[str, Any], *, prefix: str, grace_seconds: float
|
|
) -> None:
|
|
real_terminate(data, prefix=prefix, grace_seconds=grace_seconds)
|
|
deadline = time.monotonic() + 10.0
|
|
while time.monotonic() < deadline:
|
|
current = mmo_runtime.read_session_record(
|
|
mmo_runtime.session_dir(session["session_id"])
|
|
)
|
|
token = current.get("root_app_server_start_token")
|
|
if (
|
|
isinstance(token, str)
|
|
and token != prior_token
|
|
and mmo_util.process_matches(current.get("root_app_server_pid"), token)
|
|
):
|
|
replacement.update(current)
|
|
return
|
|
time.sleep(0.02)
|
|
self.fail("root controller did not publish its replacement app-server")
|
|
|
|
try:
|
|
with mock.patch.object(
|
|
mmo_runtime,
|
|
"terminate_recorded_process_group",
|
|
side_effect=terminate_after_replacement,
|
|
):
|
|
resumed = mmo_runtime.begin_resume_run(session["session_id"])
|
|
self.assertEqual(resumed["status"], "running")
|
|
self.assertEqual(
|
|
resumed["root_app_server_start_token"],
|
|
replacement["root_app_server_start_token"],
|
|
)
|
|
self.assertEqual(
|
|
resumed["root_app_server_pid"],
|
|
replacement["root_app_server_pid"],
|
|
)
|
|
time.sleep(0.5)
|
|
self.assertEqual(load_session(session["session_id"])["status"], "running")
|
|
finally:
|
|
with contextlib.suppress(Exception):
|
|
mmo_runtime.stop_session(session["session_id"], grace_seconds=0)
|
|
|
|
def test_gateway_idle_cleanup_fails_closed_on_non_object_session_state(self) -> None:
|
|
with RuntimeSandbox() as box:
|
|
settings_path = box.config / "settings.toml"
|
|
settings_path.write_text(
|
|
settings_path.read_text(encoding="utf-8").replace(
|
|
"gateway_idle_timeout_seconds = 1",
|
|
"gateway_idle_timeout_seconds = 0",
|
|
),
|
|
encoding="utf-8",
|
|
)
|
|
session = create_access_lab_session(cwd=box.workspace)
|
|
mark_session_running(session["session_id"], os.getpid())
|
|
path = box.state / "sessions" / session["session_id"] / "session.json"
|
|
original = read_json(path)
|
|
try:
|
|
path.write_text("[]\n", encoding="utf-8")
|
|
self.assertEqual(mmo_runtime.iter_sessions(strict=False), [])
|
|
with self.assertRaisesRegex(RuntimeError, "invalid session state"):
|
|
mmo_runtime.iter_sessions(strict=True)
|
|
with self.assertRaisesRegex(ValueError, "session state root must be an object"):
|
|
load_session(session["session_id"])
|
|
with self.assertRaisesRegex(RuntimeError, "invalid session state"):
|
|
create_session(profile="codex-harness-team", cwd=box.workspace)
|
|
with self.assertRaisesRegex(RuntimeError, "cannot determine gateway idleness"):
|
|
stop_idle_gateways()
|
|
self.assertEqual(gateway_status(session["snapshot_hash"])["status"], "running")
|
|
finally:
|
|
atomic_write_json(path, original)
|
|
finish_session(session["session_id"], exit_code=0)
|
|
|
|
def test_gateway_idle_accounting_includes_detached_root_and_worker_hosts(self) -> None:
|
|
gateway_hash = "a" * 64
|
|
snapshot_hash = "b" * 64
|
|
session = {
|
|
"session_id": "session-1",
|
|
"snapshot_hash": snapshot_hash,
|
|
"gateway_hash": gateway_hash,
|
|
"gateway_base_url": "http://127.0.0.1:42000/v1",
|
|
"status": "suspended",
|
|
"created_at": "2020-01-01T00:00:00+00:00",
|
|
"last_active_at": "2020-01-01T00:00:01+00:00",
|
|
"root_app_server_pid": 101,
|
|
"root_app_server_start_token": "root-live",
|
|
}
|
|
worker = {
|
|
"job_id": "job-1",
|
|
"session_id": "session-1",
|
|
"status": "suspended",
|
|
"created_at": "2020-01-01T00:00:00+00:00",
|
|
"last_progress_at": "2020-01-01T00:00:02+00:00",
|
|
"app_server_pid": 202,
|
|
"app_server_start_token": "worker-live",
|
|
}
|
|
gateway = {
|
|
"gateway_hash": gateway_hash,
|
|
"snapshot_hash": snapshot_hash,
|
|
"last_used_at": "2020-01-01T00:00:00+00:00",
|
|
"status": "running",
|
|
}
|
|
|
|
def process_is_live(pid: Any, token: Any) -> bool:
|
|
return (pid, token) in {(101, "root-live"), (202, "worker-live")}
|
|
|
|
with (
|
|
mock.patch.object(
|
|
mmo_gateway, "load_settings", return_value={"gateway_idle_timeout_seconds": 0}
|
|
),
|
|
mock.patch.object(mmo_gateway, "_gateway_key", return_value=gateway_hash),
|
|
mock.patch.object(mmo_gateway, "iter_session_records", return_value=[session]),
|
|
mock.patch.object(mmo_gateway, "iter_job_records", return_value=[worker]),
|
|
mock.patch.object(mmo_gateway, "list_gateways", return_value=[gateway]),
|
|
mock.patch.object(mmo_gateway, "process_matches", side_effect=process_is_live),
|
|
mock.patch.object(mmo_gateway.time, "time", return_value=1_577_836_810.0),
|
|
mock.patch.object(mmo_gateway, "stop_gateway") as stop,
|
|
):
|
|
self.assertEqual(stop_idle_gateways(), [])
|
|
stop.assert_not_called()
|
|
|
|
session.pop("root_app_server_pid")
|
|
session.pop("root_app_server_start_token")
|
|
worker.pop("app_server_pid")
|
|
worker.pop("app_server_start_token")
|
|
worker["status"] = "completed"
|
|
worker["finished_at"] = "2020-01-01T00:00:03+00:00"
|
|
session["status"] = "starting"
|
|
session["last_active_at"] = "2020-01-01T00:00:09+00:00"
|
|
self.assertEqual(stop_idle_gateways(), [])
|
|
stop.assert_not_called()
|
|
|
|
session["status"] = "paused"
|
|
worker["status"] = "recovering"
|
|
worker["recovery_requested_at"] = "2020-01-01T00:00:09+00:00"
|
|
self.assertEqual(stop_idle_gateways(), [])
|
|
stop.assert_not_called()
|
|
|
|
worker["status"] = "completed"
|
|
self.assertEqual(stop_idle_gateways(), [snapshot_hash])
|
|
stop.assert_called_once_with(snapshot_hash)
|
|
|
|
def test_admission_fails_closed_on_non_object_job_state(self) -> None:
|
|
with RuntimeSandbox() as box:
|
|
session = create_access_lab_session(cwd=box.workspace)
|
|
mark_session_running(session["session_id"], os.getpid())
|
|
corrupt = box.state / "jobs" / "corrupt-job"
|
|
corrupt.mkdir(parents=True)
|
|
(corrupt / "metadata.json").write_text("[]\n", encoding="utf-8")
|
|
try:
|
|
self.assertEqual(mmo_runtime.iter_jobs(strict=False), [])
|
|
with self.assertRaisesRegex(RuntimeError, "invalid job state"):
|
|
spawn_job(
|
|
session_id=session["session_id"],
|
|
caller_agent=session["root_agent"],
|
|
caller_job_id=None,
|
|
agent_id="flagship_escalation",
|
|
task_kind="analysis",
|
|
task="Attempt admission while durable job accounting is unreadable.",
|
|
)
|
|
finally:
|
|
shutil.rmtree(corrupt)
|
|
finish_session(session["session_id"], exit_code=0)
|
|
|
|
def test_worker_runner_terminates_child_when_post_launch_setup_fails(self) -> None:
|
|
with tempfile.TemporaryDirectory(prefix="mmo-runner-failure-") as temporary:
|
|
directory = Path(temporary) / "job"
|
|
directory.mkdir()
|
|
cwd = Path(temporary) / "workspace"
|
|
cwd.mkdir()
|
|
(directory / "prompt.txt").write_text("bounded task\n", encoding="utf-8")
|
|
metadata = {
|
|
"session_id": "session-test",
|
|
"run_id": "run-test",
|
|
"snapshot_hash": "snapshot-test",
|
|
"agent": "worker",
|
|
"cwd": str(cwd),
|
|
"execution_mode": "turn",
|
|
"goal_token_budget": None,
|
|
"max_goal_token_budget": None,
|
|
"stall_warning_seconds": 300,
|
|
"finalization_grace_seconds": 30,
|
|
"app_server_socket_path": str(Path(temporary) / "worker.sock"),
|
|
"agent_run_ref": "ar_test",
|
|
"result_path": str(directory / "result.md"),
|
|
"events_path": str(directory / "events.jsonl"),
|
|
"stderr_path": str(directory / "stderr.log"),
|
|
"sandbox_mode": "read-only",
|
|
"approval_policy": "never",
|
|
"job_id": "job-test",
|
|
"status": "running",
|
|
}
|
|
session = {
|
|
"session_id": "session-test",
|
|
"current_run_id": "run-test",
|
|
"snapshot_hash": "snapshot-test",
|
|
"codex_binary": "/bin/true",
|
|
"homes": {"worker": {"command_flags": []}},
|
|
}
|
|
client = mock.Mock(pid=12345)
|
|
client.start.return_value = {}
|
|
client.request.return_value = {
|
|
"thread": {"id": "thread-test", "path": str(directory / "rollout.jsonl")}
|
|
}
|
|
with (
|
|
mock.patch.dict(os.environ, {"MMO_CALLER_TOKEN": "ephemeral-test-capability"}),
|
|
mock.patch.object(
|
|
worker_runner.sys,
|
|
"argv",
|
|
["worker_runner.py", str(directory)],
|
|
),
|
|
mock.patch.object(worker_runner, "_read_metadata", return_value=metadata),
|
|
mock.patch.object(worker_runner, "session_dir", return_value=directory),
|
|
mock.patch.object(worker_runner, "read_session_record", return_value=session),
|
|
mock.patch.object(worker_runner, "begin_running", return_value=metadata),
|
|
mock.patch.object(worker_runner, "append_audit"),
|
|
mock.patch.object(
|
|
worker_runner,
|
|
"update",
|
|
side_effect=[metadata, RuntimeError("metadata publication failed")],
|
|
),
|
|
mock.patch.object(worker_runner, "session_environment", return_value={}),
|
|
# This test targets post-launch publication cleanup. Exact
|
|
# executable admission is covered separately and must not
|
|
# move the simulated failure ahead of AppServerClient.start.
|
|
mock.patch.object(worker_runner, "require_app_server_codex_version"),
|
|
mock.patch.object(worker_runner, "AppServerClient", return_value=client),
|
|
mock.patch.object(worker_runner, "_fail", return_value=1),
|
|
mock.patch.object(worker_runner.signal, "signal"),
|
|
):
|
|
self.assertEqual(worker_runner.main(), 1)
|
|
client.stop_host.assert_called_once_with()
|
|
|
|
def test_profile_switch_affects_only_new_sessions(self) -> None:
|
|
with RuntimeSandbox() as box:
|
|
set_active_profile("codex-harness-team")
|
|
first = create_session(cwd=box.workspace)
|
|
set_active_profile("access-efficient-escalation-lab")
|
|
second = create_session(cwd=box.workspace)
|
|
try:
|
|
self.assertEqual(first["profile_id"], "codex-harness-team")
|
|
self.assertEqual(second["profile_id"], "access-efficient-escalation-lab")
|
|
self.assertNotEqual(first["snapshot_hash"], second["snapshot_hash"])
|
|
self.assertEqual(
|
|
load_session(first["session_id"])["snapshot_hash"], first["snapshot_hash"]
|
|
)
|
|
finally:
|
|
finish_session(first["session_id"], exit_code=0)
|
|
finish_session(second["session_id"], exit_code=0)
|
|
|
|
|
|
if __name__ == "__main__":
|
|
unittest.main()
|