From ad843ffb21570794c6f20a9941852cb191a89f54 Mon Sep 17 00:00:00 2001 From: LoopX Agent <337587101+loopx-agent@users.noreply.github.com> Date: Mon, 5 Oct 2026 01:03:16 +0800 Subject: [PATCH] fix(runtime): make single-turn execution deadlines opt-in Signed-off-by: LoopX Agent <337587101+loopx-agent@users.noreply.github.com> --- benchmark/LHTB/README.md | 2 +- benchmark/LHTB/run.sh | 8 ++++-- benchmark/LHTB/scripts/render_config.py | 2 +- benchmark/runtime/RUNTIME.md | 6 ++-- benchmark/runtime/codex.py | 6 ++-- benchmark/runtime/harbor.py | 4 +-- benchmark/runtime/worker.py | 11 +++++--- .../configs/shared-heartbeat.yaml | 2 +- benchmark/tests/test_native_codex_goal.py | 10 ++++--- benchmark/tests/test_shared_codex_runtime.py | 17 +++++++++++ docs/reference/protocols/loopx-turn-v0.md | 24 ++++++++++++++++ .../benchmark_toolkit/native_codex_goal.py | 28 +++++++++---------- loopx/chat_agent.py | 15 +++++----- loopx/chat_codex_goal.py | 9 ++++-- loopx/cli_commands/turn_dsh_host.py | 2 +- loopx/cli_commands/turn_registration.py | 3 +- loopx/cli_commands/turn_run_once.py | 2 +- loopx/control_plane/turn_driver/codex_cli.py | 2 +- .../turn_driver/codex_operation_host.py | 6 ++-- loopx/control_plane/turn_driver/executor.py | 6 ++-- .../control_plane/turn_driver/host_process.ts | 8 +++--- .../turn_driver/host_process_transport.py | 4 +-- loopx/extensions/process_runtime.py | 10 +++---- scripts/external_scheduler_worker.py | 12 ++++---- tests/control_plane/test_host_process.py | 19 +++++++++++-- .../control_plane/test_leased_host_process.py | 2 +- tests/control_plane_ts/host_process.test.ts | 12 +++++++- tests/extensions/test_process_runtime.py | 15 ++++++++++ tests/test_chat_codex_goal.py | 4 ++- tests/test_codex_operation_host.py | 5 ++-- tests/test_external_scheduler_worker.py | 15 ++++++++++ 31 files changed, 192 insertions(+), 79 deletions(-) diff --git a/benchmark/LHTB/README.md b/benchmark/LHTB/README.md index da6c3075ca..765a59035a 100644 --- a/benchmark/LHTB/README.md +++ b/benchmark/LHTB/README.md @@ -179,7 +179,7 @@ Defaults leave cleanup room between nested layers: | Harbor agent limit | 5400 s | | LoopX scheduler process | 5080 s | | One wake command | 4800 s | -| One Codex exec | 4700 s | +| One Codex exec | Remaining total phase budget, minus cleanup reserve | The scheduler can perform multiple shorter wakes within 5080 seconds. A single long Codex wake can consume most of that budget, which is expected; Harbor's diff --git a/benchmark/LHTB/run.sh b/benchmark/LHTB/run.sh index a1ca51be6b..7596ee06ef 100755 --- a/benchmark/LHTB/run.sh +++ b/benchmark/LHTB/run.sh @@ -38,7 +38,7 @@ REASONING_EFFORT="${REASONING_EFFORT:-max}" CONCURRENCY="${CONCURRENCY:-4}" AGENT_TIMEOUT_SEC="${AGENT_TIMEOUT_SEC:-5400}" LOOPX_SCHEDULER_TIMEOUT_SEC="${LOOPX_SCHEDULER_TIMEOUT_SEC:-5080}" -LOOPX_CODEX_TURN_TIMEOUT_SEC="${LOOPX_CODEX_TURN_TIMEOUT_SEC:-4700}" +LOOPX_CODEX_TURN_TIMEOUT_SEC="${LOOPX_CODEX_TURN_TIMEOUT_SEC:-}" LHTB_MAX_RETRIES="${LHTB_MAX_RETRIES:-2}" RUNNER_RESTARTS="${RUNNER_RESTARTS:-2}" LHTB_MODELONLY_NETWORK="${LHTB_MODELONLY_NETWORK:-lhtb-modelonly}" @@ -121,6 +121,10 @@ job_name="lhtb-${LOOPX_EXECUTION_MODE}-${LOOPX_TASK_ENTRY}-${LOOPX_ITERATION_CON generated_config="$CODE_DIR/.generated/${job_name}.yaml" jobs_dir="$CODE_DIR/runs" +turn_timeout_args=() +if [[ -n "$LOOPX_CODEX_TURN_TIMEOUT_SEC" ]]; then + turn_timeout_args=(--turn-timeout "$LOOPX_CODEX_TURN_TIMEOUT_SEC") +fi "$VENV/bin/python" "$CODE_DIR/scripts/render_config.py" \ --template "$CODE_DIR/configs/heartbeat-generic-cli.yaml" \ --output "$generated_config" \ @@ -135,7 +139,7 @@ jobs_dir="$CODE_DIR/runs" --planning-timeout "$LOOPX_PLANNING_TIMEOUT_SEC" \ --iteration-context "$LOOPX_ITERATION_CONTEXT" \ --validation-command-json "$LOOPX_VALIDATION_COMMAND_JSON" \ - --turn-timeout "$LOOPX_CODEX_TURN_TIMEOUT_SEC" \ + "${turn_timeout_args[@]}" \ --scheduler-timeout "$LOOPX_SCHEDULER_TIMEOUT_SEC" \ "${task_args[@]}" diff --git a/benchmark/LHTB/scripts/render_config.py b/benchmark/LHTB/scripts/render_config.py index b27023536e..668fb12445 100755 --- a/benchmark/LHTB/scripts/render_config.py +++ b/benchmark/LHTB/scripts/render_config.py @@ -28,7 +28,7 @@ def main() -> int: parser.add_argument("--task-entry", choices=TASK_ENTRIES, default="seeded-todo") parser.add_argument("--planning-timeout", type=float, default=300) parser.add_argument("--validation-command-json", default="[]") - parser.add_argument("--turn-timeout", type=float, default=4700) + parser.add_argument("--turn-timeout", type=float, default=None) parser.add_argument("--scheduler-timeout", type=int, default=5080) args = parser.parse_args() diff --git a/benchmark/runtime/RUNTIME.md b/benchmark/runtime/RUNTIME.md index 595bc27925..5728c0ce12 100644 --- a/benchmark/runtime/RUNTIME.md +++ b/benchmark/runtime/RUNTIME.md @@ -21,7 +21,7 @@ agents: iteration_context: fresh reasoning_effort: max codex_sandbox: danger-full-access - turn_timeout_sec: 4700 + turn_timeout_sec: null scheduler_timeout_sec: 5080 replan_after_todos: 3 ``` @@ -111,7 +111,7 @@ kwargs: task_entry: loopx-planned planning_timeout_sec: 300 iteration_context: fresh - turn_timeout_sec: 4700 + turn_timeout_sec: null scheduler_timeout_sec: 5080 ``` @@ -201,3 +201,5 @@ Install the intended Harbor version for the adapter tests. Real qualification also needs installed Codex, the native Harbor backend and independently checked task output. Unit tests establish no score or model-uplift claim. Validate small jobs through each benchmark's native configuration before launching a study. + +By default, worker calls have no independent turn deadline. Harbor derives their available time from the remaining total phase budget, reserving cleanup and settlement time. An explicit `turn_timeout_sec` remains supported as an operator override. diff --git a/benchmark/runtime/codex.py b/benchmark/runtime/codex.py index e323b1d61d..d4cfe5bb72 100644 --- a/benchmark/runtime/codex.py +++ b/benchmark/runtime/codex.py @@ -20,7 +20,7 @@ class Execution: mode: str = "heartbeat" context: str = "fresh" sandbox: str = "danger-full-access" - timeout_seconds: float = 4700 + timeout_seconds: float | None = None validation_command: tuple[str, ...] = () task_entry: str = "seeded-todo" @@ -35,7 +35,9 @@ def __post_init__(self) -> None: raise ValueError("resume requires mode=turn or mode=heartbeat") if self.sandbox not in SANDBOXES: raise ValueError("unsupported Codex sandbox") - if not math.isfinite(self.timeout_seconds) or self.timeout_seconds <= 0: + if self.timeout_seconds is not None and ( + not math.isfinite(self.timeout_seconds) or self.timeout_seconds <= 0 + ): raise ValueError("execution timeout must be finite and positive") if not isinstance(self.validation_command, (list, tuple)): raise ValueError("validation_command must be an argv list") diff --git a/benchmark/runtime/harbor.py b/benchmark/runtime/harbor.py index b2c7270114..498d153c07 100644 --- a/benchmark/runtime/harbor.py +++ b/benchmark/runtime/harbor.py @@ -53,7 +53,7 @@ def __init__( iteration_context="fresh", codex_sandbox="danger-full-access", validation_command=None, - turn_timeout_sec=4700, + turn_timeout_sec=None, scheduler_timeout_sec=5080, replan_after_todos=3, task_entry="seeded-todo", @@ -66,7 +66,7 @@ def __init__( execution_mode, iteration_context, codex_sandbox, - float(turn_timeout_sec), + float(scheduler_timeout_sec) - 160 if turn_timeout_sec is None else float(turn_timeout_sec), validation_command if validation_command is not None else (), task_entry, ) diff --git a/benchmark/runtime/worker.py b/benchmark/runtime/worker.py index bd336e987c..dc5b60acdc 100644 --- a/benchmark/runtime/worker.py +++ b/benchmark/runtime/worker.py @@ -144,8 +144,8 @@ def turn_command( env["MODEL_NAME"], "--codex-sandbox", execution.sandbox, - "--timeout-seconds", - str(execution.timeout_seconds), + *(["--timeout-seconds", str(execution.timeout_seconds)] + if execution.timeout_seconds is not None else []), "--validation-command-json", json.dumps(execution.validation_command), "--available-capability", @@ -235,7 +235,8 @@ def run_once(env: dict[str, str]) -> dict: mode=env.get("LOOPX_EXECUTION_MODE", "heartbeat"), context=env.get("LOOPX_ITERATION_CONTEXT", "fresh"), sandbox=env.get("LOOPX_CODEX_SANDBOX", "danger-full-access"), - timeout_seconds=float(env.get("LOOPX_CODEX_TURN_TIMEOUT_SEC", "4700")), + timeout_seconds=(float(env["LOOPX_CODEX_TURN_TIMEOUT_SEC"]) + if env.get("LOOPX_CODEX_TURN_TIMEOUT_SEC") else None), validation_command=json.loads(env.get("LOOPX_VALIDATION_COMMAND_JSON", "[]")), task_entry=env.get("LOOPX_TASK_ENTRY", "seeded-todo"), ) @@ -272,7 +273,8 @@ def run_once(env: dict[str, str]) -> dict: receipt.update(ok=True, budget_exhausted=True, host_invoked=False) return receipt execution = replace( - execution, timeout_seconds=min(execution.timeout_seconds, remaining - 160) + execution, timeout_seconds=(remaining - 160 if execution.timeout_seconds is None + else min(execution.timeout_seconds, remaining - 160)) ) prepare_codex_home( home, @@ -330,6 +332,7 @@ def run_once(env: dict[str, str]) -> dict: ) as process: allowance = 150 if execution.mode == "turn" and stage == "execute" else 0 timeout = (float(env["LOOPX_PLANNING_TIMEOUT_SEC"]) if stage == "plan" + else None if execution.timeout_seconds is None else execution.timeout_seconds + allowance) process.communicate( input=body, timeout=timeout diff --git a/benchmark/swe-marathon/configs/shared-heartbeat.yaml b/benchmark/swe-marathon/configs/shared-heartbeat.yaml index da9300877f..1409399d62 100644 --- a/benchmark/swe-marathon/configs/shared-heartbeat.yaml +++ b/benchmark/swe-marathon/configs/shared-heartbeat.yaml @@ -9,5 +9,5 @@ agents: task_entry: seeded-todo # Use loopx-planned for the product planning checkpoint. iteration_context: fresh reasoning_effort: high - turn_timeout_sec: 4700 + turn_timeout_sec: null scheduler_timeout_sec: 5080 diff --git a/benchmark/tests/test_native_codex_goal.py b/benchmark/tests/test_native_codex_goal.py index 7750d6343e..36e84d09c5 100644 --- a/benchmark/tests/test_native_codex_goal.py +++ b/benchmark/tests/test_native_codex_goal.py @@ -293,10 +293,11 @@ def test_terminal_event_preserves_failed_turn_status() -> None: assert turn.turn_status == "failed" -def test_goal_runtime_waits_for_automatic_continuation_until_terminal() -> None: +@pytest.mark.parametrize("timeout", [None, 1]) +def test_goal_runtime_waits_for_automatic_continuation_until_terminal(timeout) -> None: transport = ContinuationTransport() - turn = run_native_goal_until_terminal(transport, _config(), timeout_sec=1) + turn = run_native_goal_until_terminal(transport, _config(), timeout_sec=timeout) methods = [method for method, _ in transport.calls] assert methods.count("turn/start") == 1 @@ -492,8 +493,9 @@ def test_process_cwd_can_differ_from_goal_thread_cwd(tmp_path: Path) -> None: assert turn.terminal_event_observed is True +@pytest.mark.parametrize("timeout", [None, 2]) def test_real_stdio_process_waits_until_native_goal_is_terminal( - tmp_path: Path, + tmp_path: Path, timeout, ) -> None: fake_server = tmp_path / "fake-continuing-codex" _write_fake_app_server(fake_server) @@ -512,7 +514,7 @@ def test_real_stdio_process_waits_until_native_goal_is_terminal( process_command=[sys.executable, str(fake_server)], process_env=process_env, response_timeout_sec=2, - goal_timeout_sec=2, + goal_timeout_sec=timeout, ) assert turn.post_goal_status == "complete" diff --git a/benchmark/tests/test_shared_codex_runtime.py b/benchmark/tests/test_shared_codex_runtime.py index 83e6d2bcf9..a28ee8f8ae 100644 --- a/benchmark/tests/test_shared_codex_runtime.py +++ b/benchmark/tests/test_shared_codex_runtime.py @@ -438,3 +438,20 @@ async def unpack(*args, **kwargs): assert staged == original assert marker.read_text() == "successor" assert uploaded == [b"original"] + + +@pytest.mark.parametrize("total", [1800, 64800]) +def test_harbor_default_execution_budget_tracks_total_trial(tmp_path, total): + pytest.importorskip("harbor") + from benchmark.runtime.harbor import BenchmarkCodex + agent = BenchmarkCodex(logs_dir=tmp_path, model_name="fixture", + scheduler_timeout_sec=total) + assert agent.execution.timeout_seconds == total - 160 + assert agent.scheduler_timeout == total + + +def test_default_execution_has_no_independent_turn_deadline(tmp_path): + execution = Execution(mode="turn", validation_command=("true",)) + env = worker_env(tmp_path) | {"LOOPX_CLI": "loopx", "LOOPX_GOAL_ID": "fixture", "LOOPX_AGENT_ID": "worker", "LOOPX_REGISTRY": "registry", "LOOPX_RUNTIME_ROOT": "runtime"} + assert execution.timeout_seconds is None + assert "--timeout-seconds" not in turn_command(env, execution, "wake-default") diff --git a/docs/reference/protocols/loopx-turn-v0.md b/docs/reference/protocols/loopx-turn-v0.md index 7b28e209ce..7e756d4f81 100644 --- a/docs/reference/protocols/loopx-turn-v0.md +++ b/docs/reference/protocols/loopx-turn-v0.md @@ -22,6 +22,30 @@ The protocol is host-neutral. A Codex CLI adapter is the first target, but the driver lifecycle must not depend on Codex-specific session files, transcript formats, or benchmark task schemas. +### Execution deadline defaults + +`turn run-once` now has no default single-turn wall-time limit (previously +120 seconds). The built-in Codex hosts and generic managed host transport wait +for natural completion or cancellation. `--timeout-seconds` remains an explicit +operator-selected deadline for bounded execution or recovery tests. + +The external heartbeat scheduler likewise has no default wake-command deadline +(previously 600 seconds); `--wake-timeout-seconds` opts into one. Quota probes, +provider request/idle liveness checks, validation commands, output budgets, +lease fencing and process-group cleanup retain their own boundaries. Cancellation +still terminates the owned process tree; disabling an execution timer does not +authorize continued effects after cancellation or lease loss. + +Benchmark adapters retain their declared total trial deadline. The shared Harbor +adapter defaults a call to the remaining trial allowance, including its existing +startup/cleanup reserve, instead of imposing an independent 4700-second wake cap. +Explicit per-call settings in existing experiment configurations remain explicit +protocol choices; omit them for natural continuation. + +单轮执行默认不再计时终止:Turn 的原 120 秒上限和外部 heartbeat 的原 600 秒 +wake 上限均改为显式选择。取消、租约失效、输出限制、请求存活检查和进程清理 +仍生效。benchmark 仍遵守整场总预算,默认不另设 4700 秒的单轮上限。 + ## Mental Model LoopX Turn is a four-stage control loop, not another agent runtime: diff --git a/loopx/capabilities/benchmark_toolkit/native_codex_goal.py b/loopx/capabilities/benchmark_toolkit/native_codex_goal.py index 677465fd5d..95d6bead4f 100644 --- a/loopx/capabilities/benchmark_toolkit/native_codex_goal.py +++ b/loopx/capabilities/benchmark_toolkit/native_codex_goal.py @@ -349,7 +349,7 @@ def wait_native_goal_turn( transport: NativeGoalEventTransport, turn: NativeGoalTurn, *, - timeout_sec: float, + timeout_sec: float | None, completed_before: int | None = None, ) -> NativeGoalTurn: """Drain events until one more correlated turn reaches a terminal event. @@ -359,18 +359,18 @@ def wait_native_goal_turn( preserves the single-turn behavior for existing callers. """ - if timeout_sec <= 0: + if timeout_sec is not None and timeout_sec <= 0: raise ValueError("timeout_sec must be positive") - deadline = time.monotonic() + timeout_sec + deadline = None if timeout_sec is None else time.monotonic() + timeout_sec if completed_before is None: if turn.terminal_event_observed: return turn completed_before = turn.turn_completed_count while turn.turn_completed_count <= completed_before: - remaining = deadline - time.monotonic() - if remaining <= 0: + remaining = None if deadline is None else deadline - time.monotonic() + if remaining is not None and remaining <= 0: raise NativeGoalProtocolError("goal_turn_timeout") - event = transport.next_event(timeout_sec=min(0.25, remaining)) + event = transport.next_event(timeout_sec=0.25 if remaining is None else min(0.25, remaining)) if event is not None: observe_native_goal_event(turn, event) return turn @@ -399,7 +399,7 @@ def run_native_goal_turn( transport: NativeGoalEventTransport, config: NativeGoalConfig, *, - timeout_sec: float, + timeout_sec: float | None, ) -> NativeGoalTurn: """Execute the complete native Goal transaction over an admitted transport.""" @@ -413,7 +413,7 @@ def run_native_goal_until_terminal( transport: NativeGoalEventTransport, config: NativeGoalConfig, *, - timeout_sec: float, + timeout_sec: float | None, on_turn_started: Callable[[NativeGoalTurn], None] | None = None, ) -> NativeGoalTurn: """Run one native Goal until its status leaves ``active``. @@ -423,16 +423,16 @@ def run_native_goal_until_terminal( continuation events and reading the Goal status under one total timeout. """ - if timeout_sec <= 0: + if timeout_sec is not None and timeout_sec <= 0: raise ValueError("timeout_sec must be positive") turn = start_native_goal_turn(transport, config) if on_turn_started is not None: on_turn_started(turn) - deadline = time.monotonic() + timeout_sec + deadline = None if timeout_sec is None else time.monotonic() + timeout_sec completed_before = turn.turn_completed_count while True: - remaining = deadline - time.monotonic() - if remaining <= 0: + remaining = None if deadline is None else deadline - time.monotonic() + if remaining is not None and remaining <= 0: raise NativeGoalDeadlineExceeded("goal_timeout_before_terminal") try: wait_native_goal_turn( @@ -639,7 +639,7 @@ def run_native_goal_process( process_env: Mapping[str, str] | None = None, process_cwd: str | None = None, response_timeout_sec: float = 30, - goal_timeout_sec: float = 21_600, + goal_timeout_sec: float | None = 21_600, ) -> NativeGoalTurn: """Spawn a real app-server and execute one complete native Goal turn.""" @@ -661,7 +661,7 @@ def run_native_goal_process_until_terminal( process_env: Mapping[str, str] | None = None, process_cwd: str | None = None, response_timeout_sec: float = 30, - goal_timeout_sec: float = 21_600, + goal_timeout_sec: float | None = 21_600, ) -> NativeGoalTurn: """Spawn app-server and keep the native Goal alive through continuations.""" diff --git a/loopx/chat_agent.py b/loopx/chat_agent.py index 096d33578c..9ae27066e3 100644 --- a/loopx/chat_agent.py +++ b/loopx/chat_agent.py @@ -432,7 +432,7 @@ class CodexChatAgentSession: reasoning_effort: str | None = None response_timeout_sec: float = 30.0 idle_timeout_sec: float = 180.0 - hard_timeout_sec: float = 900.0 + hard_timeout_sec: float | None = 900.0 next_request_id: int = 5 current_turn_id: str = "" model_catalog_compatibility_applied: bool = False @@ -468,7 +468,7 @@ def start( objective: str, response_timeout_sec: float = 30.0, idle_timeout_sec: float = 180.0, - hard_timeout_sec: float = 900.0, + hard_timeout_sec: float | None = 900.0, resume_thread_id: str | None = None, execution_mode: bool = False, isolate_process_tree: bool = False, @@ -988,7 +988,7 @@ def send( last_activity_at = started_at while True: now = time.monotonic() - if now - started_at >= self.hard_timeout_sec: + if self.hard_timeout_sec is not None and now - started_at >= self.hard_timeout_sec: raise self._timeout_error( "hard_timeout", "Codex Chat turn reached its hard time limit." ) @@ -996,15 +996,14 @@ def send( raise self._timeout_error( "idle_timeout", "Codex Chat turn stopped producing activity." ) - deadline = min( - started_at + self.hard_timeout_sec, - last_activity_at + self.idle_timeout_sec, - ) + deadline = last_activity_at + self.idle_timeout_sec + if self.hard_timeout_sec is not None: + deadline = min(deadline, started_at + self.hard_timeout_sec) try: message = self._next_event(deadline=deadline) except CodexChatAgentError: now = time.monotonic() - if now - started_at >= self.hard_timeout_sec: + if self.hard_timeout_sec is not None and now - started_at >= self.hard_timeout_sec: raise self._timeout_error( "hard_timeout", "Codex Chat turn reached its hard time limit.", diff --git a/loopx/chat_codex_goal.py b/loopx/chat_codex_goal.py index f117d0787b..2d64fc2430 100644 --- a/loopx/chat_codex_goal.py +++ b/loopx/chat_codex_goal.py @@ -248,7 +248,8 @@ def _stop_owned_transport_on_failure(self) -> None: self.session.close() def _observe(self, emit: Callable[[str, dict[str, Any]], None]) -> dict[str, Any]: - deadline = time.monotonic() + self.session.hard_timeout_sec + deadline = (None if self.session.hard_timeout_sec is None + else time.monotonic() + self.session.hard_timeout_sec) parts: list[str] = [] completed_messages: list[str] = [] display = VisibleResponseStreamFilter(protected_paths=[self.session.work_dir]) @@ -276,12 +277,14 @@ def _observe(self, emit: Callable[[str, dict[str, Any]], None]) -> dict[str, Any if parts: completed_messages.append(parse_agent_response("".join(parts), protected_paths=[self.session.work_dir])["message"]) return self._response(goal, messages=completed_messages) - if time.monotonic() >= deadline: + if deadline is not None and time.monotonic() >= deadline: raise self.session._timeout_error( "hard_timeout", "Native Goal reached the Chat time limit; resume it in this conversation.", ) - event = self.session._next_event(deadline=deadline) + event = self.session._next_event( + deadline=deadline if deadline is not None else time.monotonic() + self.session.idle_timeout_sec + ) if self.session._check_server_gate(event): continue if _event_thread_id(event) not in {"", self.session.thread_id}: diff --git a/loopx/cli_commands/turn_dsh_host.py b/loopx/cli_commands/turn_dsh_host.py index 7df0809401..522a5b37da 100644 --- a/loopx/cli_commands/turn_dsh_host.py +++ b/loopx/cli_commands/turn_dsh_host.py @@ -43,7 +43,7 @@ def build_dsh_host_runner( "dsh_home": Path(args.dsh_home) if args.dsh_home else None, "cordis": Path(args.dsh_cordis) if args.dsh_cordis else None, "runtime_bin": args.dsh_runtime_bin, - "request_timeout_seconds": max(1.0, args.timeout_seconds - 5.0), + "request_timeout_seconds": None if args.timeout_seconds is None else max(1.0, args.timeout_seconds - 5.0), "dsh_runner": Path(args.dsh_runner) if args.dsh_runner else None, }.items() if value is not None diff --git a/loopx/cli_commands/turn_registration.py b/loopx/cli_commands/turn_registration.py index da9369406f..ff582534b2 100644 --- a/loopx/cli_commands/turn_registration.py +++ b/loopx/cli_commands/turn_registration.py @@ -311,7 +311,8 @@ def register_turn_commands( "run_dsh_turn(...) used instead of the DeepSeek Harness SDK." ), ) - run_once.add_argument("--timeout-seconds", type=float, default=120.0) + run_once.add_argument("--timeout-seconds", type=float, default=None, + help="Optional execution deadline; by default wait for host completion or cancellation.") run_once.add_argument( "--retry-failed-turn", action="store_true", diff --git a/loopx/cli_commands/turn_run_once.py b/loopx/cli_commands/turn_run_once.py index 631f27b6dc..fb606ead94 100644 --- a/loopx/cli_commands/turn_run_once.py +++ b/loopx/cli_commands/turn_run_once.py @@ -777,7 +777,7 @@ def run_built_in_host( "model": args.codex_model, "reasoning_effort": args.codex_reasoning_effort, "mcp_server": args.codex_mcp_server_json, - "timeout_seconds": max(1.0, args.timeout_seconds - 5.0), + "timeout_seconds": None if args.timeout_seconds is None else max(1.0, args.timeout_seconds - 5.0), } if goal_admission is not None: options["goal_admission"] = goal_admission diff --git a/loopx/control_plane/turn_driver/codex_cli.py b/loopx/control_plane/turn_driver/codex_cli.py index 93224f246c..a382a92701 100644 --- a/loopx/control_plane/turn_driver/codex_cli.py +++ b/loopx/control_plane/turn_driver/codex_cli.py @@ -645,7 +645,7 @@ def run_codex_cli_host( model: str | None = None, reasoning_effort: str | None = None, mcp_server: Mapping[str, Any] | None = None, - timeout_seconds: float = 115.0, + timeout_seconds: float | None = None, goal_admission: FirstPartyHostGoalAdmission | None = None, ) -> dict[str, Any]: if request.get("schema_version") != LOOPX_TURN_HOST_REQUEST_SCHEMA_VERSION: diff --git a/loopx/control_plane/turn_driver/codex_operation_host.py b/loopx/control_plane/turn_driver/codex_operation_host.py index 1df534c970..a5ec78fbcc 100644 --- a/loopx/control_plane/turn_driver/codex_operation_host.py +++ b/loopx/control_plane/turn_driver/codex_operation_host.py @@ -242,7 +242,7 @@ def run_codex_operation_host( reasoning_effort: str | None = None, source_route: Mapping[str, str] | None = None, mcp_server: Mapping[str, Any] | None = None, - timeout_seconds: float = 115, + timeout_seconds: float | None = None, goal_admission: FirstPartyHostGoalAdmission | None = None, confirmed_operation_id: str | None = None, ) -> dict[str, Any]: @@ -359,9 +359,9 @@ def run_codex_operation_host( resume_thread_id=binding["session_id"] if binding else None, dynamic_tools=[OPERATION_TOOL], host_config=host_config, - response_timeout_sec=min(timeout_seconds, 30), + response_timeout_sec=30 if timeout_seconds is None else min(timeout_seconds, 30), hard_timeout_sec=timeout_seconds, - idle_timeout_sec=min(timeout_seconds, 180), + idle_timeout_sec=180 if timeout_seconds is None else min(timeout_seconds, 180), ) def store_binding(): diff --git a/loopx/control_plane/turn_driver/executor.py b/loopx/control_plane/turn_driver/executor.py index 4c2903f14a..f70f313efc 100644 --- a/loopx/control_plane/turn_driver/executor.py +++ b/loopx/control_plane/turn_driver/executor.py @@ -662,7 +662,7 @@ def _run_host( *, argv: Sequence[str], project: Path, - timeout_seconds: float, + timeout_seconds: float | None, ) -> dict[str, Any]: stdout: list[str] = [] stderr_chars = 0 @@ -784,7 +784,7 @@ def _host_result_stage( argv: Sequence[str] | None, completion_lifecycle_configured: bool, project: Path, - timeout_seconds: float, + timeout_seconds: float | None, journal: dict[str, Any], persist_journal: JournalPersist, effects: dict[str, bool], @@ -1276,7 +1276,7 @@ def run_loopx_turn_once( project: Path, runtime_root: Path, goal_id: str, - timeout_seconds: float, + timeout_seconds: float | None, execute: bool, retry_failed: bool = False, task_validator: TaskValidator | None = None, diff --git a/loopx/control_plane/turn_driver/host_process.ts b/loopx/control_plane/turn_driver/host_process.ts index deb6fb32c0..87c1098488 100644 --- a/loopx/control_plane/turn_driver/host_process.ts +++ b/loopx/control_plane/turn_driver/host_process.ts @@ -9,7 +9,7 @@ export interface HostProcessRequest { argv: string[]; cwd: string; input: string; - timeout_ms: number; + timeout_ms: number | null; drain_timeout_ms: number; stdout_limit_bytes: number | null; } @@ -36,7 +36,7 @@ export function decodeHostProcessRequest(value: unknown): HostProcessRequest { if (Object.keys(v).length !== fields.length || fields.some(k => !Object.hasOwn(v, k)) || !Array.isArray(v.argv) || !v.argv.length || !v.argv[0] || v.argv.some(x => typeof x !== "string" || x.includes("\0")) || typeof v.cwd !== "string" || !v.cwd || v.cwd.includes("\0") || typeof v.input !== "string" || - typeof v.timeout_ms !== "number" || !Number.isFinite(v.timeout_ms) || v.timeout_ms <= 0 || v.timeout_ms > 2147483647 || + (v.timeout_ms !== null && (typeof v.timeout_ms !== "number" || !Number.isFinite(v.timeout_ms) || v.timeout_ms <= 0 || v.timeout_ms > 2147483647)) || typeof v.drain_timeout_ms !== "number" || !Number.isFinite(v.drain_timeout_ms) || v.drain_timeout_ms < 0 || v.drain_timeout_ms > 30000 || (v.stdout_limit_bytes !== null && (typeof v.stdout_limit_bytes !== "number" || !Number.isSafeInteger(v.stdout_limit_bytes) || v.stdout_limit_bytes < 1))) throw new TypeError("invalid Host request fields"); @@ -100,7 +100,7 @@ export async function runHostProcess(request: HostProcessRequest, }; const abort = () => stop("cancelled"); signal?.addEventListener("abort", abort, {once: true}); - const deadline = setTimeout(() => stop("timeout"), request.timeout_ms); + const deadline = request.timeout_ms === null ? undefined : setTimeout(() => stop("timeout"), request.timeout_ms); let drainTimer: ReturnType | undefined; const exited = new Promise(resolve => { child.once("error", () => { outcome = "spawn_failed"; resolve(); }); @@ -160,7 +160,7 @@ export async function runHostProcess(request: HostProcessRequest, await clean(); return {...base, outcome, output_complete: complete}; } finally { - clearTimeout(deadline); + if (deadline) clearTimeout(deadline); if (drainTimer) clearTimeout(drainTimer); signal?.removeEventListener("abort", abort); child.stdin.destroy(); child.stdout.destroy(); child.stderr.destroy(); diff --git a/loopx/control_plane/turn_driver/host_process_transport.py b/loopx/control_plane/turn_driver/host_process_transport.py index 31b67917fa..1431cc7f09 100644 --- a/loopx/control_plane/turn_driver/host_process_transport.py +++ b/loopx/control_plane/turn_driver/host_process_transport.py @@ -216,7 +216,7 @@ def run_host_process( *, project: Path, input_text: str, - timeout_seconds: float, + timeout_seconds: float | None, stdout_limit_bytes: int | None = None, drain_timeout_seconds: float = 2, on_stdout: Callable[[str], None] | None = None, @@ -236,7 +236,7 @@ def run_host_process( "argv": command, "cwd": str(project), "input": input_text, - "timeout_ms": max(1.0, timeout_seconds) * 1000, + "timeout_ms": None if timeout_seconds is None else max(1.0, timeout_seconds) * 1000, "drain_timeout_ms": drain_timeout_seconds * 1000, "stdout_limit_bytes": stdout_limit_bytes, } diff --git a/loopx/extensions/process_runtime.py b/loopx/extensions/process_runtime.py index a3d9cec3c7..29b0fb2311 100644 --- a/loopx/extensions/process_runtime.py +++ b/loopx/extensions/process_runtime.py @@ -102,7 +102,7 @@ def run_capped_process( argv: Sequence[str], *, stdin: bytes, - timeout_seconds: float, + timeout_seconds: float | None, output_limit_bytes: int, env: Mapping[str, str] | None = None, cwd: str | Path | None = None, @@ -194,16 +194,16 @@ def write_stdin() -> None: for thread in threads: thread.start() - deadline = time.monotonic() + timeout_seconds + deadline = None if timeout_seconds is None else time.monotonic() + timeout_seconds timed_out = False try: while process.poll() is None: - remaining = deadline - time.monotonic() - if remaining <= 0: + remaining = None if deadline is None else deadline - time.monotonic() + if remaining is not None and remaining <= 0: timed_out = True terminate_process_tree(process, termination_grace_seconds) break - if limit_event.wait(timeout=min(0.05, remaining)): + if limit_event.wait(timeout=0.05 if remaining is None else min(0.05, remaining)): terminate_process_tree(process, termination_grace_seconds) break except BaseException: diff --git a/scripts/external_scheduler_worker.py b/scripts/external_scheduler_worker.py index 5aa7ab5a7b..b474f43193 100644 --- a/scripts/external_scheduler_worker.py +++ b/scripts/external_scheduler_worker.py @@ -264,11 +264,11 @@ def _shell_argv(command: str) -> list[str]: return ["/bin/sh", "-c", command] -def _run_wake(command: str, *, timeout_seconds: float) -> CappedProcessResult: +def _run_wake(command: str, *, timeout_seconds: float | None) -> CappedProcessResult: return run_capped_process( _shell_argv(command), stdin=b"", - timeout_seconds=max(0.01, timeout_seconds), + timeout_seconds=None if timeout_seconds is None else max(0.01, timeout_seconds), output_limit_bytes=PROCESS_OUTPUT_LIMIT_BYTES, termination_grace_seconds=10, ) @@ -376,9 +376,7 @@ def run_worker( if args.wake_cmd: wake_result = _run_wake( args.wake_cmd, - timeout_seconds=float( - getattr(args, "wake_timeout_seconds", 600.0) - ), + timeout_seconds=getattr(args, "wake_timeout_seconds", None), ) wake_rc = wake_result.returncode wake_failure_kind = wake_result.failure_kind @@ -552,8 +550,8 @@ def main(argv: Sequence[str] | None = None) -> int: parser.add_argument( "--wake-timeout-seconds", type=float, - default=600.0, - help="Maximum duration of one configured wake command.", + default=None, + help="Optional wake execution deadline; default waits for completion or cancellation.", ) args = parser.parse_args(argv) if args.registry is None: diff --git a/tests/control_plane/test_host_process.py b/tests/control_plane/test_host_process.py index aa977a5973..c88e88ae8e 100644 --- a/tests/control_plane/test_host_process.py +++ b/tests/control_plane/test_host_process.py @@ -60,7 +60,8 @@ def test_real_generic_host_roundtrip_and_stream_budget(tmp_path: Path) -> None: assert overflow["reason"] == "host stdout exceeded the result budget" -def test_callback_failure_waits_for_owned_host_cleanup(tmp_path: Path) -> None: +@pytest.mark.parametrize("deadline", [None, 20]) +def test_callback_failure_waits_for_owned_host_cleanup(tmp_path: Path, deadline) -> None: def reject(_text: str) -> None: raise ValueError("consumer stopped") @@ -74,7 +75,7 @@ def reject(_text: str) -> None: ], project=tmp_path, input_text="", - timeout_seconds=20, + timeout_seconds=deadline, on_stdout=reject, ) assert time.monotonic() - started < 8 @@ -420,3 +421,17 @@ def written(path, supervises, group, supervision): finally: live.kill() live.wait(timeout=10) + + +def test_default_turn_deadline_is_optional_and_real_transport_accepts_none(tmp_path): + from loopx.cli import build_parser + parser = build_parser() + args = parser.parse_args(["turn", "run-once", "--goal-id", "fixture", + "--agent-id", "agent", "--project", str(tmp_path)]) + assert args.timeout_seconds is None + result = run_host_process( + [sys.executable, "-c", "import time;time.sleep(.15);print('finished')"], + project=tmp_path, input_text="", timeout_seconds=None, + ) + assert result["outcome"] == "exited" + assert result["returncode"] == 0 diff --git a/tests/control_plane/test_leased_host_process.py b/tests/control_plane/test_leased_host_process.py index f2dc428fab..7e0ba2235b 100644 --- a/tests/control_plane/test_leased_host_process.py +++ b/tests/control_plane/test_leased_host_process.py @@ -99,7 +99,7 @@ def record(event,**values): context[field] = [sys.executable, "-c", relay, str(marker), str(trace), phase, str(held_deadline), *context[field]] observed = run_host_process([sys.executable, "-c", "import time;time.sleep(35);print('finished')"], - project=tmp_path, input_text="", timeout_seconds=45, delegated_lease=context) + project=tmp_path, input_text="", timeout_seconds=None, delegated_lease=context) evidence = {"original": context["lease"], "observed": observed, "trace": trace.read_text() if fault in {"lost_reply", "late_reply"} and trace.exists() else ""} # Keep the first observation in JUnit even if subsequent canonical readback diff --git a/tests/control_plane_ts/host_process.test.ts b/tests/control_plane_ts/host_process.test.ts index 6567c7c3e8..274225f8ed 100644 --- a/tests/control_plane_ts/host_process.test.ts +++ b/tests/control_plane_ts/host_process.test.ts @@ -23,6 +23,16 @@ test("Host output is streamed with UTF-8 boundaries and stdin EOF", async () => assert.equal(result.outcome, "exited"); assert.equal(result.returncode, 0); assert.equal(result.output_complete, true); }); +test("null execution deadline waits for natural completion", async () => { + const value = decodeHostProcessRequest(request( + "setTimeout(()=>process.stdout.write('finished'),150)", {timeout_ms: null})); + let text = ""; + const result = await runHostProcess(value, async item => { text += item.text; }); + assert.equal(text, "finished"); + assert.equal(result.outcome, "exited"); + assert.equal(result.returncode, 0); +}); + test("invalid requests and absent executables cannot be mistaken for success", async () => { for (const change of [{argv: []}, {argv: [""]}, {timeout_ms: Infinity}, {timeout_ms: 0}, {stdout_limit_bytes: -1}, {drain_timeout_ms: -1}, {unexpected: true}]) { assert.throws(() => decodeHostProcessRequest({...request(""), ...change}), /invalid/); @@ -76,7 +86,7 @@ for (const mode of ["timeout", "abort", "leader_exit", "closed_pipes"] as const) return timer(() => { controller.abort(); }, 3000); // fixture readiness watchdog }); } - const result = await runHostProcess(request(script, {timeout_ms: mode === "timeout" ? 500 : 3000}), async item => { + const result = await runHostProcess(request(script, {timeout_ms: mode === "timeout" ? 500 : mode === "abort" ? null : 3000}), async item => { if (item.text.includes("ready")) { ready = true; if (mode === "abort") controller.abort(); diff --git a/tests/extensions/test_process_runtime.py b/tests/extensions/test_process_runtime.py index 616ba9462e..ab6dcea469 100644 --- a/tests/extensions/test_process_runtime.py +++ b/tests/extensions/test_process_runtime.py @@ -196,3 +196,18 @@ def test_a_spewing_provider_on_stderr_is_its_own_failure() -> None: assert result.failure_kind == "stderr_too_large" assert result.stdout == b"" + + +def test_no_execution_deadline_preserves_completion_and_output_budget(): + result = run_capped_process( + [sys.executable, "-c", "import time;time.sleep(.1);print('done')"], + stdin=b"", timeout_seconds=None, output_limit_bytes=1024, + ) + assert result.returncode == 0 + assert result.stdout.strip() == b"done" + assert result.failure_kind is None + limited = run_capped_process( + [sys.executable, "-c", "import sys,time;print('x'*2000,flush=True);time.sleep(30)"], + stdin=b"", timeout_seconds=None, output_limit_bytes=1024, + ) + assert limited.failure_kind == "response_too_large" diff --git a/tests/test_chat_codex_goal.py b/tests/test_chat_codex_goal.py index 047dc6bc9c..fc5d1ca72c 100644 --- a/tests/test_chat_codex_goal.py +++ b/tests/test_chat_codex_goal.py @@ -129,8 +129,9 @@ def test_goal_role_is_not_permission_to_take_over_another_binding(changes): validate_goal_chat({**session, **changes}, []) +@pytest.mark.parametrize("hard_deadline", [None, 900]) def test_multiple_native_turns_and_stale_goal_notification_are_reconciled( - tmp_path, monkeypatch + tmp_path, monkeypatch, hard_deadline ): host = Host( tmp_path, @@ -147,6 +148,7 @@ def test_multiple_native_turns_and_stale_goal_notification_are_reconciled( end("two"), ], ) + host.session.hard_timeout_sec = hard_deadline events = [] response = host.driver.run( parse_native_goal_command("/goal start --tokens 10000 Analyze revisions"), diff --git a/tests/test_codex_operation_host.py b/tests/test_codex_operation_host.py index 381352cf6a..9397ab7238 100644 --- a/tests/test_codex_operation_host.py +++ b/tests/test_codex_operation_host.py @@ -471,8 +471,9 @@ def test_registered_source_tokens_survive_owned_native_prepare_and_store_readbac assert refused["ok"] is False +@pytest.mark.parametrize("deadline", [None, 5]) def test_owned_app_server_process_authenticates_native_metadata_and_resumes_same_profile( - tmp_path: Path, + tmp_path: Path, deadline, ) -> None: executable = tmp_path / "fake-codex-operation" executable.write_text(FAKE_SERVER) @@ -485,7 +486,7 @@ def test_owned_app_server_process_authenticates_native_metadata_and_resumes_same "model": "test-model", "reasoning_effort": "xhigh", "source_route": {"host_surface": "codex-app", "thread_id": "source-one"}, - "timeout_seconds": 5, + "timeout_seconds": deadline, } first = run_codex_operation_host(_request(), **options) assert first["result_kind"] == "wait" diff --git a/tests/test_external_scheduler_worker.py b/tests/test_external_scheduler_worker.py index 658060fe81..07c9173077 100644 --- a/tests/test_external_scheduler_worker.py +++ b/tests/test_external_scheduler_worker.py @@ -331,3 +331,18 @@ def test_wake_timeout_enters_failed_backoff_and_kills_descendants( state = json.loads(state_file.read_text(encoding="utf-8")) assert state["last_wake_status"] == "wake_failed" assert state["last_wake_failure_kind"] == "timeout" + + +def test_wake_default_has_no_execution_deadline(monkeypatch, tmp_path): + seen = [] + monkeypatch.setattr(worker, "run_worker", lambda args: seen.append(args) or 0) + assert worker.main(["--goal-id", "fixture", "--agent-id", "agent", + "--registry", str(tmp_path / "registry.json")]) == 0 + assert seen[0].wake_timeout_seconds is None + assert seen[0].quota_timeout_seconds == 30.0 + result = worker._run_wake( + shlex.join([sys.executable, "-c", "import time;time.sleep(.1);print('done')"]), + timeout_seconds=None, + ) + assert result.returncode == 0 + assert result.failure_kind is None