diff --git a/docs/architecture/rfcs/STATUS.md b/docs/architecture/rfcs/STATUS.md index 2a1be7490b..40ffc23e4e 100644 --- a/docs/architecture/rfcs/STATUS.md +++ b/docs/architecture/rfcs/STATUS.md @@ -58,7 +58,7 @@ appendix may keep dated history, but no dated log heading may precede it. | [RFC: Research Exploration Control Plane v0](research-exploration-control-plane-v0.md) | Accepted | none | — | | [RFC: Semantic Vocabulary Convergence and Commit-Time Drift Checks (v0)](semantic-vocabulary-convergence-v0.md) | Accepted | none | [5 entries](ledger/semantic-vocabulary-convergence-v0/) | | [RFC: Shared Goal Alignment and Governed Amendment Protocol (v0)](shared-goal-alignment-and-governed-amendment-v0.md) | Accepted | none | [2 entries](ledger/shared-goal-alignment-and-governed-amendment-v0/) | -| [RFC: LoopX Shared Control-Plane Authority and Pluggable State Providers (v0)](shared-goal-authority-state-provider-v0.md) | Accepted | none | [20 entries](ledger/shared-goal-authority-state-provider-v0/) | +| [RFC: LoopX Shared Control-Plane Authority and Pluggable State Providers (v0)](shared-goal-authority-state-provider-v0.md) | Accepted | none | [22 entries](ledger/shared-goal-authority-state-provider-v0/) | | [RFC: Single-Owner Local Daemon (v0)](single-owner-local-daemon-v0.md) | Accepted | none | — | | [RFC: TypeScript Control-Plane Migration Direction v0](typescript-control-plane-migration-v0.md) | Accepted | none | [12 entries](ledger/typescript-control-plane-migration-v0/) | diff --git a/docs/architecture/rfcs/STATUS.zh-CN.md b/docs/architecture/rfcs/STATUS.zh-CN.md index 1ed8986b9f..81e666a420 100644 --- a/docs/architecture/rfcs/STATUS.zh-CN.md +++ b/docs/architecture/rfcs/STATUS.zh-CN.md @@ -55,7 +55,7 @@ | [RFC:研究型探索控制面 v0](research-exploration-control-plane-v0.zh-CN.md) | 已接受 | 无 | — | | [RFC:语义词表收敛与提交期漂移检查(v0)](semantic-vocabulary-convergence-v0.zh-CN.md) | 已接受 | 无 | [5 条](ledger/semantic-vocabulary-convergence-v0/) | | [RFC:共享 Goal 对齐与受治理 Amendment 协议(v0)](shared-goal-alignment-and-governed-amendment-v0.zh-CN.md) | 已接受 | 无 | [2 条](ledger/shared-goal-alignment-and-governed-amendment-v0/) | -| [RFC:LoopX 共享控制面权威与可插拔状态 Provider(v0)](shared-goal-authority-state-provider-v0.zh-CN.md) | 已接受 | 无 | [20 条](ledger/shared-goal-authority-state-provider-v0/) | +| [RFC:LoopX 共享控制面权威与可插拔状态 Provider(v0)](shared-goal-authority-state-provider-v0.zh-CN.md) | 已接受 | 无 | [22 条](ledger/shared-goal-authority-state-provider-v0/) | | [RFC: Single-Owner Local Daemon (v0)](single-owner-local-daemon-v0.md) | 已接受 | none | — | | [RFC:LoopX 控制面 TypeScript 渐进迁移方向 v0](typescript-control-plane-migration-v0.zh-CN.md) | 已接受 | 无 | [12 条](ledger/typescript-control-plane-migration-v0/) | diff --git a/docs/architecture/rfcs/ledger/shared-goal-authority-state-provider-v0/2026-09-27-host-supervision.md b/docs/architecture/rfcs/ledger/shared-goal-authority-state-provider-v0/2026-09-27-host-supervision.md new file mode 100644 index 0000000000..09bba01147 --- /dev/null +++ b/docs/architecture/rfcs/ledger/shared-goal-authority-state-provider-v0/2026-09-27-host-supervision.md @@ -0,0 +1,72 @@ +# Local defaults: managed Host supervision and reconciled delivery plan + +- Baseline: `fd96e5e25`, audited September 27, 2026. +- Outcome: overall roadmap S2/S4/R5, shared authority external-execution closure, + TS replacement-first migration. No provider or capability is introduced. +- This replaces the **remaining delivery estimate**, not the historical evidence, + in the September 24 reconciliation. The recovery slice proposed in + [#5140](https://github.com/loopx-project/loopx/pull/5140) is still open. +- [中文](2026-09-27-host-supervision.zh-CN.md). + +## Count deliveries, not architectural headings + +Complete-source transport/assembly, transaction capture, canonical pagination, +File v1 automatic backup/upgrade and Python prototype retirement are already on +main. #5013, #5063, #5102 and #5105 are not future work. Two promoted Goals prove +those particular cutovers, not every execution, migration or recovery boundary. + +The old “three packages” and #5140's “three PRs afterwards” were too coarse for +execution protection. An actual subprocess reproduction shows the missing +prerequisite: generic Host timeout kills only its leader, while Codex cleanup +returns early after the leader exits. Both can leave descendants doing work. +Deleting a lease or rejecting a later result does not stop that process. + +The current **four newly planned deliveries include this PR**: + +| Delivery | Observable exit and Python retirement | +| --- | --- | +| **1. Managed subprocess supervision (this PR)** | Generic command and Codex CLI share one TS lifecycle through timeout, caller loss, pipe drain and process-group termination. Retire their separate Python termination/thread-reader implementations. This closes the process component, not the lease component below. | +| **2. Authority-bound execution interval** | Connect the existing provider-neutral lease owner to actual execution: current proof before start, bounded renewal, cancellation on expiry/reclaim/revocation, and uncertain-effect recovery. Reclaim must not silently overlap an old executor. Test with real processes and File/SQLite; explicitly qualify attached Hosts without cancellation. Remove replaced Python decisions rather than create a second lease store. | +| **3. Whole-Goal migration and fenced recovery integration** | Adopt #5140 recovery and #5054 source retirement; cover source drain, reviewed cutover, retained command consumers, projection readback and rollback after later writes. Inventory existing callers before adding writers. Delete legacy decisions only when their actual callers have moved. | +| **4. Default onboarding and bounded Python retirement** | New-Goal creation, settings, CLI, packaged frontend and Lark select the qualified local profile consistently. Existing Goals have explicit upgrade, backup and recovery. Remove remaining replaced Python business writers, retaining necessary rendering and Host IO adapters. | + +**Three planned new PRs remain after this one.** This is a scoped delivery plan, +not an unconditional total or proof that lease supervision has shipped. It is +one additional execution slice compared with #5140's proposed estimate; the +reproduction above is the reason, and this PR does not subtract the uncompleted +lease row. If another slice is needed, amend its named row and evidence. + +Separately, existing open PRs are #5140 (recovery/audit), #5054 (old Todo event +retirement and supervisor logging), and #4931 (SQLite receipt-proof encoding). +Thus the integration inventory is **six named PR deliveries for the File route** +(this + three planned + #5140 + #5054), or **seven for the SQLite route** including +#4931. These counts include already implemented open PRs; they do not mean six +or seven new implementations. #5140's SQLite batch proof read complements #4931; +neither small-suite success qualifies D2. New defects discovered by qualification +can still require changes, so there is no justified guaranteed PR total today. + +D1 consumer parity, profile-specific D2 capacity/recovery/soak and D3 cohort +cutover remain acceptance work, not invented PR allocations. The audited #4224 +1 MiB report still fails receipt p95 (269.03 ms / 50 ms) and scan-100 p95 +(801.81 ms / 250 ms); this process change cannot fix or certify those metrics. +PostgreSQL retains its separate authenticated transport, tenant/identity, +cross-host execution, pooling/failover and operations qualification. Local +process cleanup is reusable across providers because it does not read their +physical layouts or create authority. + +## Ownership and neighboring work + +`control_plane/turn_driver/host_process.ts` owns the managed process lifetime; +its private bridge treats the Python owner's control-pipe EOF as cancellation. +Python adapts transient output, Codex sessions and typed results. Existing +`turn run-once` callers adopt this automatically; no new CLI option, configuration +editor, capability registration, frontend or Lark surface is needed. Attached +App sessions and in-process DSH adapters do not pass through this subprocess +owner and are not represented as newly protected. + +#5141 fences Host state by GoalRef, while #5142 preserves effect uncertainty in +Turn error readback. Neither replaces process supervision. Integration must +retain their admission checks before launching and their recovery observations; +this PR changes neither GoalRef authority nor settlement semantics. + +[Operational behavior and limits](../../../../reference/protocols/loopx-turn-v0.md#managed-host-process-lifetime). diff --git a/docs/architecture/rfcs/ledger/shared-goal-authority-state-provider-v0/2026-09-27-host-supervision.zh-CN.md b/docs/architecture/rfcs/ledger/shared-goal-authority-state-provider-v0/2026-09-27-host-supervision.zh-CN.md new file mode 100644 index 0000000000..c4a01d76a0 --- /dev/null +++ b/docs/architecture/rfcs/ledger/shared-goal-authority-state-provider-v0/2026-09-27-host-supervision.zh-CN.md @@ -0,0 +1,59 @@ +# 本地默认切换:受管 Host 进程监督与交付重估 + +- 基线:`fd96e5e25`,2026-09-27 核对。 +- 目标:总路线 S2/S4/R5、shared authority 外部执行闭环、TS 替换式迁移。 + 不新增 provider 或 capability。 +- 本文替换 9 月 24 日清单的**剩余交付估算**,不改写历史证据。 + [#5140](https://github.com/loopx-project/loopx/pull/5140) 的恢复切片仍未合入。 +- [English](2026-09-27-host-supervision.md)。 + +## 数交付,不数架构标题 + +完整来源传输与组装、事务捕获、canonical 分页、File v1 自动备份升级、Python +原型退役已在 main。#5013、#5063、#5102、#5105 不再计入待开发。两个已晋升 Goal +证明的是那两次切换,不代表全部执行、迁移、恢复边界已完成。 + +旧的“三个大包”以及 #5140 的“之后三个 PR”对执行保护估得过粗。真实子进程复现 +发现前置缺口:通用 Host 超时只杀主进程;Codex 清理在主进程退出后直接返回; +两者都可能遗留继续工作的后代。删除租约或拒绝最终结果不能让这些进程停止。 + +当前规划的**四个新增交付包含本 PR**: + +| 交付 | 可观察退出条件与 Python 退役 | +| --- | --- | +| **1. 受管子进程监督(本 PR)** | 通用命令与 Codex CLI 共用 TS 生命周期,覆盖超时、调用方消失、管道排空和进程组终止;删除各自的 Python 终止及读线程实现。完成进程部分,不冒称下行租约部分已完成。 | +| **2. 权威约束的执行区间** | 把现有 provider-neutral 租约 owner 接到实际执行:启动前当前证明、执行中有界续约、到期/回收/撤权取消、不确定效果恢复。回收不能静默重叠旧执行器。真实进程与 File/SQLite 验证;附着式 Host 无取消能力时明确支持边界。替换 Python 决策,不另建租约存储。 | +| **3. 整 Goal 迁移与带 fence 的恢复闭环** | 接入 #5140 恢复与 #5054 来源退役;覆盖来源排空、已审切换、保留命令消费者、投影读回及存在后续写入时的回退。先盘点 caller,再决定是否需要 writer;只有真实 caller 已迁移才删除旧决策。 | +| **4. 默认入口与有界 Python 清理** | 新 Goal、设置、CLI、打包前端、Lark 一致选择合格的本地 profile;旧 Goal 有显式升级、备份、恢复。删除已替代 Python 业务 writer,保留必要渲染与 Host IO adapter。 | + +**本 PR 之后仍规划三个新增 PR。** 这是有具体边界的计划,不是保证总数,也不代表 +租约监督已经交付。相比 #5140 的提案,明确多拆一个进程切片,原因是上面的复现; +不能用本 PR 抵扣未完成的租约行。后续若再拆,必须修改具体行并给出证据。 + +既有在途 PR 另计:#5140 恢复/审计、#5054 旧 Todo 事件退役与 supervisor 日志、 +#4931 SQLite 回执证明编码。因此 **File 路线的已知合入清单为六项**(本次 + 三项 +待开发 + #5140 + #5054),**SQLite 路线加入 #4931 后为七项**。这里包含已经实现 +但没合入的 PR,绝不是还要新写六七个。#5140 的 SQLite 批量证明读取与 #4931 +互补;小规模测试通过不能代替 D2。资格验证仍可能发现需修改的缺陷,因此不能 +承诺无条件总 PR 数。 + +D1 消费者一致性、各 profile 的 D2 容量/恢复/soak、D3 cohort 切换是验收工作, +不能凭空折算 PR。核对的 #4224 1 MiB 报告仍有 receipt p95 269.03ms / 50ms、 +scan-100 p95 801.81ms / 250ms 未通过,本次进程改造不能修复或认证这些指标。 +PostgreSQL 的认证传输、tenant/identity、跨 Host 执行、连接池/故障切换、运维 +资格仍是独立中期路径。本地进程清理不读取 provider 物理布局、不产生权威, +因此各 provider 可以复用。 + +## 归属与邻近工作 + +`control_plane/turn_driver/host_process.ts` 拥有进程生命周期,私有 bridge 将 Python +调用方控制管道 EOF 视为取消。Python 适配瞬时输出、Codex 会话、typed result。 +现有 `turn run-once` 自动接入,无新 CLI 参数、配置编辑器、capability、前端或 +Lark 界面。附着式 App 会话与进程内 DSH adapter 不经过此子进程 owner,不宣称 +它们因此得到保护。 + +#5141 用 GoalRef 约束 Host 状态;#5142 保留 Turn 错误读回中的效果不确定性。 +两者不能替代进程监督。集成时要保留启动前准入和恢复观察,本 PR 不修改 GoalRef +权限或 settlement 语义。 + +[操作行为和限制](../../../../reference/protocols/loopx-turn-v0.md#managed-host-process-lifetime)。 diff --git a/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md b/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md index 5a839341da..e47af23098 100644 --- a/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md +++ b/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md @@ -36,6 +36,13 @@ existing #5054/#4931 and D2/D3 evidence remain separate. This is not a guarantee count of future defect repairs. [Current inventory, rationale and exits](ledger/shared-goal-authority-state-provider-v0/2026-09-27-recovery-audit.md). +Managed Host supervision lands in the same window as a separate delivered slice: +one TS supervisor now owns generic command and Codex CLI process lifetimes, so +authority-bound execution stays open rather than closing here, and the old three +architectural packages remain a pointer instead of a decrementing PR counter. +Its named plan, changed estimate and boundaries are recorded separately. +[Named plan, changed estimate and boundaries](ledger/shared-goal-authority-state-provider-v0/2026-09-27-host-supervision.md). + File retained-state storage now reuses the existing TS checkpoint/delta codec, stacked on #5063's verified read cache and RPC budgets. Original revisions, receipts and full historical projections survive the physical format upgrade. diff --git a/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.zh-CN.md b/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.zh-CN.md index 87ae3eaa3a..c44ebd5d8a 100644 --- a/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.zh-CN.md +++ b/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.zh-CN.md @@ -30,6 +30,11 @@ PR:本次恢复切片、外部执行区间保护、整 Goal 激活/回退集 缺失证据另列,不能保证最终缺陷修复数量。 [当前清单、依据及退出条件](ledger/shared-goal-authority-state-provider-v0/2026-09-27-recovery-audit.zh-CN.md)。 +同一窗口另有已交付的进程监督切片:通用命令与 Codex CLI 的进程生命周期改由 +一个 TS supervisor 承担,因此执行中租约约束仍开放,不由本次关闭;旧“三个 +架构包”仍是指针,不是递减 PR 计数器。该切片的逐项计划、估算变化与边界单列。 +[核对的逐项计划、估算变化与边界](ledger/shared-goal-authority-state-provider-v0/2026-09-27-host-supervision.zh-CN.md)。 + ## Todo 事件路径退役(2026-09-25) PR #5054 将原先的事件 writer 捕获方案改为删除这条实验性 Todo 来源。 diff --git a/docs/architecture/rfcs/typescript-control-plane-migration-v0.md b/docs/architecture/rfcs/typescript-control-plane-migration-v0.md index 33753edc08..a360916cdb 100644 --- a/docs/architecture/rfcs/typescript-control-plane-migration-v0.md +++ b/docs/architecture/rfcs/typescript-control-plane-migration-v0.md @@ -41,6 +41,13 @@ existing #5054/#4931 and D2/D3 evidence remain separate. This is not a guarantee count of future defect repairs. [Current inventory, rationale and exits](ledger/shared-goal-authority-state-provider-v0/2026-09-27-recovery-audit.md). +Managed Host supervision lands in the same window as a separate delivered slice: +one TS supervisor now owns generic command and Codex CLI process lifetimes, so +authority-bound execution stays open rather than closing here, and the old three +architectural packages remain a pointer instead of a decrementing PR counter. +Its named plan, changed estimate and boundaries are recorded separately. +[Named plan, changed estimate and boundaries](ledger/shared-goal-authority-state-provider-v0/2026-09-27-host-supervision.md). + ## Native authority qualification and prototype retirement (2026-09-26) The coverage-only Python coordination executor, head codec, File provider and diff --git a/docs/architecture/rfcs/typescript-control-plane-migration-v0.zh-CN.md b/docs/architecture/rfcs/typescript-control-plane-migration-v0.zh-CN.md index 9cfbf1e0f4..c477ef4d29 100644 --- a/docs/architecture/rfcs/typescript-control-plane-migration-v0.zh-CN.md +++ b/docs/architecture/rfcs/typescript-control-plane-migration-v0.zh-CN.md @@ -36,6 +36,11 @@ PR:本次恢复切片、外部执行区间保护、整 Goal 激活/回退集 缺失证据另列,不能保证最终缺陷修复数量。 [当前清单、依据及退出条件](ledger/shared-goal-authority-state-provider-v0/2026-09-27-recovery-audit.zh-CN.md)。 +同一窗口另有已交付的进程监督切片:通用命令与 Codex CLI 的进程生命周期改由 +一个 TS supervisor 承担,因此执行中租约约束仍开放,不由本次关闭;旧“三个 +架构包”仍是指针,不是递减 PR 计数器。该切片的逐项计划、估算变化与边界单列。 +[核对的逐项计划、估算变化与边界](ledger/shared-goal-authority-state-provider-v0/2026-09-27-host-supervision.zh-CN.md)。 + ## Todo 事件路径退役(2026-09-25) PR #5054 将原先的事件 writer 捕获方案改为删除这条实验性 Todo 来源。 diff --git a/docs/reference/protocols/loopx-turn-v0.md b/docs/reference/protocols/loopx-turn-v0.md index 5225a313cf..9e7b38b2a0 100644 --- a/docs/reference/protocols/loopx-turn-v0.md +++ b/docs/reference/protocols/loopx-turn-v0.md @@ -253,6 +253,48 @@ dedicated typed result channel; passing `trae chat` directly as the adapter is not sufficient. Check the installed CLI's help and pin the qualified command shape because flags and headless behavior may vary by version. +### Managed Host Process Lifetime + +The generic command executor and built-in Codex CLI adapter now share a TS +process supervisor. Existing `turn run-once` commands need no new option. Node +uses the same supported-version discovery as the control-plane runtime; the +private request stream is limited to 8 MiB and is not a durable protocol. + +A Host leader exiting, its output pipes closing and its descendants stopping +are distinct observations. On POSIX, LoopX starts a dedicated process group, +sends TERM and escalates to KILL after 300 ms, **including when the leader has +already exited**. Normal result return also cleans up leftover group members. +Host commands must not use that group to launch intended persistent services. +Windows retains Python command-launch compatibility (including batch entrypoints) +through a transport-only relay, then attempts tree termination before killing +the leader. Windows uses best-effort process-tree cleanup; this delivery does not claim +POSIX-equivalent cancellation or Windows qualification. + +Timeout, output-consumer failure and loss of the owning Python process trigger +cleanup. The control pipe remains open for the job lifetime; EOF cancels work. +After leader exit, output drain is bounded (normally two seconds), rather than +waiting indefinitely for inherited pipes. Generic stdout is capped at its +existing 12,000-byte result budget while streaming. Codex output is consumed +transiently with LF-framed records capped at 1,048,576 characters and a finite +set of failure categories; an oversized record makes diagnostic observation +incomplete. UTF-8 characters split across byte chunks remain intact. Raw Host +output is not written to LoopX state. + +Generic results require complete output and zero exit status. Codex retains its +existing separate typed result-file contract: incomplete diagnostics do not +invent a failure category, and a validated result file remains usable. Timeout +still preserves the observed opaque session for the existing retry path. No +process observation certifies Todo completion, refunds spend or rolls back an +external effect; independent validation and settlement keep their owners. + +This is **process supervision, not execution authority or a sandbox**. It does +not renew provider leases, prevent stale remote side effects, cancel attached +App sessions, or supervise in-process DSH execution. Descendants that escape the +process group and killing the supervisor itself with SIGKILL are outside this +boundary. Caller death can precede cleanup; the local lane lock alone cannot +certify no overlap with a replacement executor. Authority-bound renewal, +revocation and uncertain-effect recovery remain a separate delivery. + ### Repeatable Codex CLI Qualification The repository includes an opt-in end-to-end qualification that creates an diff --git a/loopx/control_plane/turn_driver/codex_cli.py b/loopx/control_plane/turn_driver/codex_cli.py index dcc15f0af5..6070e4b51f 100644 --- a/loopx/control_plane/turn_driver/codex_cli.py +++ b/loopx/control_plane/turn_driver/codex_cli.py @@ -7,10 +7,7 @@ import os import re import shutil -import signal -import subprocess import tempfile -import threading from collections.abc import Mapping from pathlib import Path from typing import Any @@ -29,6 +26,7 @@ ) from .execution_profile import require_supported_reasoning_effort from .host_failure import BuiltInHostError +from .host_process_transport import HostOutputLines, run_host_process from .transaction import LOOPX_TURN_RESULT_SCHEMA_VERSION, TRANSACTION_PHASES @@ -754,22 +752,6 @@ def _select_failure_category(categories: list[str]) -> str | None: ) -def _terminate_process(proc: subprocess.Popen[str]) -> None: - if proc.poll() is not None: - return - try: - os.killpg(proc.pid, signal.SIGTERM) - except (OSError, ProcessLookupError): - proc.terminate() - try: - proc.wait(timeout=3) - except subprocess.TimeoutExpired: - try: - os.killpg(proc.pid, signal.SIGKILL) - except (OSError, ProcessLookupError): - proc.kill() - - def _codex_command( *, codex_bin: str, @@ -934,67 +916,42 @@ def commit() -> None: session_id=session_id, mcp_server=mcp_server, ) - proc = subprocess.Popen( - command, - cwd=project, - stdin=subprocess.PIPE, - stdout=subprocess.PIPE, - stderr=subprocess.PIPE, - text=True, - encoding="utf-8", - errors="replace", - start_new_session=True, - ) observed_session: list[str] = [] - structured_failure_categories: list[str] = [] - diagnostic_failure_categories: list[str] = [] - - def discard_events() -> None: - assert proc.stdout is not None - for line in proc.stdout: - try: - event = json.loads(line) - except json.JSONDecodeError: - continue - if isinstance(event, dict): - candidate = codex_cli_event_session_id(event) - if candidate and not observed_session: - observed_session.append(candidate) - structured, diagnostic = _event_failure_categories(event) - if structured: - structured_failure_categories.append(structured) - if diagnostic: - diagnostic_failure_categories.append(diagnostic) - - reader = threading.Thread(target=discard_events, daemon=True) - - def discard_stderr() -> None: - assert proc.stderr is not None - for line in proc.stderr: - category = _diagnostic_failure_category(line) - if category: - diagnostic_failure_categories.append(category) - - stderr_reader = threading.Thread(target=discard_stderr, daemon=True) - reader.start() - stderr_reader.start() - assert proc.stdin is not None - timed_out = False - try: - proc.stdin.write(_prompt(request)) - proc.stdin.close() - returncode = proc.wait(timeout=max(1.0, timeout_seconds)) - except subprocess.TimeoutExpired: - _terminate_process(proc) - timed_out = True - returncode = proc.returncode - except BaseException: - _terminate_process(proc) - raise - finally: - reader.join(timeout=OUTPUT_DRAIN_TIMEOUT_SECONDS) - stderr_reader.join(timeout=OUTPUT_DRAIN_TIMEOUT_SECONDS) - output_observation_incomplete = reader.is_alive() or stderr_reader.is_alive() + structured_failure_categories: set[str] = set() + diagnostic_failure_categories: set[str] = set() + + def observe_event(line: str) -> None: + try: + event = json.loads(line) + except json.JSONDecodeError: + return + if isinstance(event, dict): + candidate = codex_cli_event_session_id(event) + if candidate and not observed_session: + observed_session.append(candidate) + structured, diagnostic = _event_failure_categories(event) + if structured: + structured_failure_categories.add(structured) + if diagnostic: + diagnostic_failure_categories.add(diagnostic) + + def observe_stderr(line: str) -> None: + category = _diagnostic_failure_category(line) + if category: + diagnostic_failure_categories.add(category) + + events = HostOutputLines(observe_event) + diagnostics = HostOutputLines(observe_stderr) + observed = run_host_process(command, project=project, input_text=_prompt(request), + timeout_seconds=timeout_seconds, drain_timeout_seconds=OUTPUT_DRAIN_TIMEOUT_SECONDS, + on_stdout=events.feed, on_stderr=diagnostics.feed) + events.finish() + diagnostics.finish() + returncode = observed["returncode"] + timed_out = observed["outcome"] == "timeout" + output_observation_incomplete = not (observed["output_complete"] and events.complete and diagnostics.complete) + if observed["outcome"] not in {"exited", "timeout"}: + raise BuiltInHostError("codex_cli_process_" + observed["outcome"]) if timed_out: if observed_session: store_session(observed_session[0]) @@ -1007,8 +964,8 @@ def discard_stderr() -> None: "unknown" if output_observation_incomplete else ( - _select_failure_category(structured_failure_categories) - or _select_failure_category(diagnostic_failure_categories) + _select_failure_category(list(structured_failure_categories)) + or _select_failure_category(list(diagnostic_failure_categories)) or "exit_nonzero" ) ) diff --git a/loopx/control_plane/turn_driver/executor.py b/loopx/control_plane/turn_driver/executor.py index 7fb67b0ef0..57cc53a97f 100644 --- a/loopx/control_plane/turn_driver/executor.py +++ b/loopx/control_plane/turn_driver/executor.py @@ -3,7 +3,6 @@ from __future__ import annotations import json -import subprocess from collections.abc import Callable, Mapping, Sequence from pathlib import Path from typing import Any @@ -33,6 +32,7 @@ ) from .driver import selected_turn_todo from .execution_readback import execution_payload +from .host_process_transport import run_host_process from .host_binding import ( managed_executor_unavailable_payload, ) @@ -657,34 +657,31 @@ def _run_host( project: Path, timeout_seconds: float, ) -> dict[str, Any]: + stdout: list[str] = [] + stderr_chars = 0 + + def count_stderr(text: str) -> None: + nonlocal stderr_chars + stderr_chars += len(text) + try: - completed = subprocess.run( - list(argv), - cwd=project, - input=json.dumps(request, ensure_ascii=False, separators=(",", ":")), - text=True, encoding="utf-8", errors="replace", - capture_output=True, - timeout=max(1.0, timeout_seconds), - check=False, - ) - except (OSError, subprocess.TimeoutExpired) as exc: + observed = run_host_process(argv, project=project, + input_text=json.dumps(request, ensure_ascii=False, separators=(",", ":")), + timeout_seconds=timeout_seconds, stdout_limit_bytes=HOST_RESULT_MAX_BYTES, + on_stdout=stdout.append, on_stderr=count_stderr) + except (OSError, RuntimeError, ValueError) as exc: return {"ok": False, "reason": type(exc).__name__, "returncode": None} - if completed.returncode != 0: - return { - "ok": False, - "reason": "host command returned non-zero", - "returncode": completed.returncode, - "stderr_chars": len(completed.stderr), - } - encoded = completed.stdout.encode("utf-8") - if len(encoded) > HOST_RESULT_MAX_BYTES: - return { - "ok": False, - "reason": "host stdout exceeded the result budget", - "returncode": 0, - } + if observed["outcome"] == "output_limit": + return {"ok": False, "reason": "host stdout exceeded the result budget", "returncode": observed["returncode"]} + if observed["outcome"] != "exited": + return {"ok": False, "reason": "host process " + observed["outcome"], "returncode": observed["returncode"]} + if not observed["output_complete"]: + return {"ok": False, "reason": "host output observation incomplete", "returncode": observed["returncode"]} + if observed["returncode"] != 0: + return {"ok": False, "reason": "host command returned non-zero", + "returncode": observed["returncode"], "stderr_chars": stderr_chars} try: - value = json.loads(completed.stdout) + value = json.loads("".join(stdout)) except json.JSONDecodeError: return { "ok": False, diff --git a/loopx/control_plane/turn_driver/host_process.ts b/loopx/control_plane/turn_driver/host_process.ts new file mode 100644 index 0000000000..685a5c39fa --- /dev/null +++ b/loopx/control_plane/turn_driver/host_process.ts @@ -0,0 +1,132 @@ +/** Own a managed Host process through exit, pipe drain and descendant cleanup. + * This is process supervision, not a task lease or a sandbox. */ +import {spawn, type ChildProcessWithoutNullStreams} from "node:child_process"; +import {setTimeout as delay} from "node:timers/promises"; +import {StringDecoder} from "node:string_decoder"; + +export interface HostProcessRequest { + argv: string[]; + cwd: string; + input: string; + timeout_ms: number; + drain_timeout_ms: number; + stdout_limit_bytes: number | null; +} +export interface HostProcessResult { + kind: "result"; + outcome: "exited" | "timeout" | "cancelled" | "output_limit" | "spawn_failed"; + returncode: number | null; + signal: string | null; + output_complete: boolean; + cleanup_scope: "process_group" | "process_tree_best_effort"; + group_signal_sent: boolean; +} +export type HostProcessOutput = {kind: "stdout" | "stderr"; text: string}; +export const HOST_PROCESS_TERMINATE_GRACE_MS = 300; + +/** Restrict transport size separately from the caller's public result budget. */ +export function decodeHostProcessRequest(value: unknown): HostProcessRequest { + if (!value || typeof value !== "object" || Array.isArray(value)) throw new TypeError("invalid Host request"); + const v = value as Record; + const fields = ["argv", "cwd", "input", "timeout_ms", "drain_timeout_ms", "stdout_limit_bytes"]; + 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 || + 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"); + return v as unknown as HostProcessRequest; +} + +/** A stopped leader does not prove that its process group has stopped. */ +function signalGroup(child: ChildProcessWithoutNullStreams, signal: NodeJS.Signals): boolean { + if (!child.pid) return false; + try { process.kill(-child.pid, signal); return true; } + catch (error) { if ((error as NodeJS.ErrnoException).code === "ESRCH") return false; throw error; } +} + +export async function runHostProcess(request: HostProcessRequest, + output: (item: HostProcessOutput) => Promise, signal?: AbortSignal): Promise { + const base: HostProcessResult = {kind: "result", outcome: "spawn_failed", returncode: null, signal: null, + output_complete: true, cleanup_scope: process.platform === "win32" ? "process_tree_best_effort" : "process_group", + group_signal_sent: false}; + if (signal?.aborted) return {...base, outcome: "cancelled"}; + const child = spawn(request.argv[0], request.argv.slice(1), {cwd: request.cwd, + stdio: ["pipe", "pipe", "pipe"], detached: process.platform !== "win32", windowsHide: true}); + let outcome: HostProcessResult["outcome"] = "exited"; + let complete = true, forcedDrain = false, stdoutBytes = 0; + let cleanup: Promise | undefined; + const clean = () => cleanup ??= (async () => { + if (process.platform === "win32") { + if (!child.pid) return; + // Try the tree while its leader is still discoverable, before the direct + // kill fallback. A dead leader still makes this best effort on Windows. + await new Promise(resolve => { + const killer = spawn("taskkill", ["/pid", String(child.pid), "/T", "/F"], + {stdio: "ignore", windowsHide: true, timeout: 1000}); + killer.once("error", () => resolve()); + killer.once("exit", code => { base.group_signal_sent = code === 0; resolve(); }); + }); + if (child.exitCode === null && child.signalCode === null) child.kill("SIGKILL"); + return; + } + // Always signal the owned group, including after the leader's exit. + const sent = signalGroup(child, "SIGTERM"); + base.group_signal_sent ||= sent; + if (sent) { await delay(HOST_PROCESS_TERMINATE_GRACE_MS); signalGroup(child, "SIGKILL"); } + })(); + const stop = (reason: HostProcessResult["outcome"]) => { + if (outcome === "exited") outcome = reason; + void clean().catch(() => { complete = false; child.kill("SIGKILL"); }); + }; + const abort = () => stop("cancelled"); + signal?.addEventListener("abort", abort, {once: true}); + const deadline = setTimeout(() => stop("timeout"), request.timeout_ms); + let drainTimer: ReturnType | undefined; + const exited = new Promise(resolve => { + child.once("error", () => { outcome = "spawn_failed"; resolve(); }); + child.once("exit", (code, sig) => { + base.returncode = code; base.signal = sig; + // Descendants may hold inherited pipes forever after the leader exits. + drainTimer = setTimeout(() => { + complete = false; forcedDrain = true; + void clean().catch(() => { complete = false; }).finally(() => { + child.stdout.destroy(); child.stderr.destroy(); + }); + }, request.drain_timeout_ms); + resolve(); + }); + }); + const read = async (kind: "stdout" | "stderr") => { + const decoder = new StringDecoder("utf8"); + try { + for await (const chunk of child[kind]) { + const bytes = chunk as Buffer; + if (kind === "stdout" && request.stdout_limit_bytes !== null) { + stdoutBytes += bytes.length; + if (stdoutBytes > request.stdout_limit_bytes) { complete = false; stop("output_limit"); continue; } + } + const text = decoder.write(bytes); + if (text) await output({kind, text}); // Backpressure, not an unbounded output queue. + } + const tail = decoder.end(); + if (tail && !(kind === "stdout" && outcome === "output_limit")) await output({kind, text: tail}); + } catch { complete = false; if (!forcedDrain) stop("cancelled"); } + }; + const reads = Promise.all([read("stdout"), read("stderr")]); + child.stdin.on("error", () => {}); // A Host may close stdin before consuming it. + child.stdin.end(request.input); + try { + await exited; + await reads; + if (drainTimer) clearTimeout(drainTimer); + await clean(); + return {...base, outcome, output_complete: complete}; + } finally { + 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_bridge.ts b/loopx/control_plane/turn_driver/host_process_bridge.ts new file mode 100644 index 0000000000..b1c261c8bf --- /dev/null +++ b/loopx/control_plane/turn_driver/host_process_bridge.ts @@ -0,0 +1,27 @@ +/** Private parent/child transport. EOF means the owning Python process left. */ +import {once} from "node:events"; +import {decodeHostProcessRequest, runHostProcess} from "./host_process.ts"; +const owner = new AbortController(); +process.stdin.on("end", () => owner.abort()); +process.on("SIGTERM", () => owner.abort()); +process.on("SIGINT", () => owner.abort()); +process.stdout.on("error", () => owner.abort()); +let pending = Buffer.alloc(0), accepted = false; +const emit = async (item: unknown) => { + if (owner.signal.aborted && process.stdout.destroyed) throw new Error("owner disconnected"); + if (!process.stdout.write(JSON.stringify(item) + "\n")) await once(process.stdout, "drain"); +}; +process.stdin.on("data", (chunk: Buffer) => { + if (accepted) return; + pending = Buffer.concat([pending, chunk]); + if (pending.length > 8 * 1024 * 1024) { owner.abort(); process.exitCode = 1; process.stdin.destroy(); return; } + const newline = pending.indexOf(10); + if (newline < 0) return; + accepted = true; + const line = pending.subarray(0, newline).toString("utf8"); pending = Buffer.alloc(0); + void (async () => { + try { await emit(await runHostProcess(decodeHostProcessRequest(JSON.parse(line)), emit, owner.signal)); } + catch { process.exitCode = 1; } + finally { process.stdin.destroy(); } + })(); +}); diff --git a/loopx/control_plane/turn_driver/host_process_transport.py b/loopx/control_plane/turn_driver/host_process_transport.py new file mode 100644 index 0000000000..4f51148fa6 --- /dev/null +++ b/loopx/control_plane/turn_driver/host_process_transport.py @@ -0,0 +1,129 @@ +"""Python transport for the TS-owned managed Host process lifecycle.""" + +from __future__ import annotations + +import json +import subprocess +import sys +from collections.abc import Callable, Sequence +from pathlib import Path +from typing import Any + +from ..effect_runtime import _node_executable + +# Keep Python's Windows executable/batch launcher compatibility. No timeout, +# buffering or lifecycle decision lives in this transport-only child. +_WINDOWS_COMMAND_RELAY = "import subprocess,sys;sys.exit(subprocess.call(sys.argv[1:]))" + + +class HostOutputLines: + """Frame LF records without retaining raw trajectories or an unbounded line.""" + + def __init__(self, consume: Callable[[str], None], max_chars: int = 1_048_576): + self.consume = consume + self.max_chars = max_chars + self.pending = "" + self.dropping = False + self.complete = True + + def feed(self, text: str) -> None: + pieces = text.split("\n") + for index, piece in enumerate(pieces): + if not self.dropping: + if len(self.pending) + len(piece) > self.max_chars: + self.complete = False + self.pending = "" + self.dropping = True + else: + self.pending += piece + if index < len(pieces) - 1: + if not self.dropping: + self.consume(self.pending) + self.pending = "" + self.dropping = False + + def finish(self) -> None: + if self.pending and not self.dropping: + self.consume(self.pending) + self.pending = "" + + +def run_host_process( + argv: Sequence[str], + *, + project: Path, + input_text: str, + timeout_seconds: float, + stdout_limit_bytes: int | None = None, + drain_timeout_seconds: float = 2, + on_stdout: Callable[[str], None] | None = None, + on_stderr: Callable[[str], None] | None = None, +) -> dict[str, Any]: + """Keep the control pipe open until exit; EOF cancels the owned process group. + + Callbacks observe transient chunks. They must not persist raw Host output. + TS owns deadlines, byte budgets, termination and the final observation. + """ + command = list(argv) + if sys.platform == "win32": + command = [sys.executable, "-c", _WINDOWS_COMMAND_RELAY, *command] + request = { + "argv": command, + "cwd": str(project), + "input": input_text, + "timeout_ms": max(1.0, timeout_seconds) * 1000, + "drain_timeout_ms": drain_timeout_seconds * 1000, + "stdout_limit_bytes": stdout_limit_bytes, + } + bridge = Path(__file__).with_name("host_process_bridge.ts") + with subprocess.Popen( + [ + _node_executable(), + "--no-warnings", + "--experimental-strip-types", + str(bridge), + ], + stdin=subprocess.PIPE, + stdout=subprocess.PIPE, + stderr=subprocess.DEVNULL, + text=True, + encoding="utf-8", + errors="strict", + start_new_session=True, + ) as proc: + assert proc.stdin is not None and proc.stdout is not None + result = None + try: + proc.stdin.write( + json.dumps(request, ensure_ascii=False, separators=(",", ":")) + "\n" + ) + proc.stdin.flush() + for line in proc.stdout: + event = json.loads(line) + kind = event.get("kind") + if kind in {"stdout", "stderr"} and isinstance(event.get("text"), str): + consume = on_stdout if kind == "stdout" else on_stderr + if consume is not None: + consume(event["text"]) + elif kind == "result" and event.get("outcome") in { + "exited", + "timeout", + "cancelled", + "output_limit", + "spawn_failed", + }: + result = event + else: + raise RuntimeError("Invalid managed Host process observation") + finally: + # Also runs on callback failure / Ctrl-C. Do not kill the supervisor + # before it has had a chance to terminate its owned Host group. + proc.stdin.close() + try: + proc.wait(timeout=5) + except subprocess.TimeoutExpired: + proc.kill() + proc.wait() + if proc.returncode != 0 or result is None: + raise RuntimeError("Managed Host process supervision returned no result") + return result diff --git a/tests/control_plane/test_host_process.py b/tests/control_plane/test_host_process.py new file mode 100644 index 0000000000..7da98714ef --- /dev/null +++ b/tests/control_plane/test_host_process.py @@ -0,0 +1,176 @@ +"""Real managed Host boundaries: no paid model, no external side effects.""" + +from __future__ import annotations + +import os +import signal +import subprocess +import sys +import time +from pathlib import Path + +import pytest + +from loopx.control_plane.turn_driver.executor import _run_host +from loopx.control_plane.turn_driver.host_process_transport import ( + HostOutputLines, + run_host_process, +) + + +def test_host_output_lines_bound_storage_and_use_lf() -> None: + rows: list[str] = [] + lines = HostOutputLines(rows.append, max_chars=20) + lines.feed("one\u2028two\n" + "x" * 100) + assert lines.pending == "" + lines.feed("tail\nlast") + lines.finish() + assert rows == ["one\u2028two", "last"] + assert not lines.complete + + +def test_real_generic_host_roundtrip_and_stream_budget(tmp_path: Path) -> None: + request = {"message": "one private local request"} + result = _run_host( + request, + argv=[ + sys.executable, + "-c", + "import sys,json;print(json.dumps(json.load(sys.stdin)))", + ], + project=tmp_path, + timeout_seconds=5, + ) + assert result == {"ok": True, "value": request, "returncode": 0} + overflow = _run_host( + request, + argv=[ + sys.executable, + "-c", + "import sys,time;sys.stdout.write('x'*1000000);sys.stdout.flush();time.sleep(30)", + ], + project=tmp_path, + timeout_seconds=5, + ) + assert overflow["ok"] is False + assert overflow["reason"] == "host stdout exceeded the result budget" + + +def test_callback_failure_waits_for_owned_host_cleanup(tmp_path: Path) -> None: + def reject(_text: str) -> None: + raise ValueError("consumer stopped") + + started = time.monotonic() + with pytest.raises(ValueError, match="consumer stopped"): + run_host_process( + [ + sys.executable, + "-c", + "import time;print('ready',flush=True);time.sleep(30)", + ], + project=tmp_path, + input_text="", + timeout_seconds=20, + on_stdout=reject, + ) + assert time.monotonic() - started < 8 + + +@pytest.mark.skipif(os.name == "nt", reason="POSIX process-group cancellation contract") +def test_disappearing_python_owner_cancels_real_host(tmp_path: Path) -> None: + marker = tmp_path / "counter" + pid_path = tmp_path / "pid" + host = f""" +import os,time,signal +from pathlib import Path +signal.signal(signal.SIGTERM, signal.SIG_IGN) +Path({str(pid_path)!r}).write_text(str(os.getpid())) +i=0 +while True: + Path({str(marker)!r}).write_text(str(i));i+=1;time.sleep(.02) +""" + launcher = f""" +from pathlib import Path +from loopx.control_plane.turn_driver.host_process_transport import run_host_process +run_host_process({[sys.executable, "-c", host]!r}, project=Path({str(tmp_path)!r}), input_text='', timeout_seconds=30) +""" + owner = subprocess.Popen( + [sys.executable, "-c", launcher], + stdout=subprocess.DEVNULL, + stderr=subprocess.DEVNULL, + ) + try: + deadline = time.monotonic() + 10 + while not marker.exists() and time.monotonic() < deadline: + time.sleep(0.02) + assert marker.exists(), "Host never started" + owner.kill() + owner.wait(timeout=5) + time.sleep(1) + before = marker.read_text() + time.sleep(0.15) + assert marker.read_text() == before, "Host survived loss of its owner" + finally: + if owner.poll() is None: + owner.kill() + owner.wait(timeout=5) + if pid_path.exists(): + try: + os.kill(int(pid_path.read_text()), signal.SIGKILL) + except ProcessLookupError: + pass + + +@pytest.mark.skipif(os.name == "nt", reason="POSIX process-group cancellation contract") +def test_generic_host_timeout_stops_real_descendant(tmp_path: Path) -> None: + marker = tmp_path / "child-work" + pid_path = tmp_path / "child-pid" + child = f""" +import os,time,signal +from pathlib import Path +signal.signal(signal.SIGTERM, signal.SIG_IGN) +Path({str(pid_path)!r}).write_text(str(os.getpid())) +i=0 +while True: + Path({str(marker)!r}).write_text(str(i));i+=1;time.sleep(.02) +""" + host = f"import subprocess,sys,time;subprocess.Popen({[sys.executable, '-c', child]!r});time.sleep(30)" + try: + result = _run_host( + {}, argv=[sys.executable, "-c", host], project=tmp_path, timeout_seconds=1 + ) + assert result["ok"] is False + assert result["reason"] == "host process timeout" + before = marker.read_text() + time.sleep(0.15) + assert marker.read_text() == before + finally: + if pid_path.exists(): + try: + os.kill(int(pid_path.read_text()), signal.SIGKILL) + except ProcessLookupError: + pass + + +def test_windows_transport_relay_preserves_argv_and_stdin(tmp_path: Path) -> None: + from loopx.control_plane.turn_driver.host_process_transport import _WINDOWS_COMMAND_RELAY + + # The relay is tested here on any OS; real .cmd resolution remains a Windows + # integration obligation. All args stay argv entries, not interpolated code. + values = ["two words", "a&b", "%NAME%", 'one"quote', "界"] + result = _run_host( + {}, + argv=[ + sys.executable, + "-c", + _WINDOWS_COMMAND_RELAY, + sys.executable, + "-c", + "import json,sys;print(json.dumps({'args':sys.argv[1:],'input':json.load(sys.stdin)}))", + *values, + ], + project=tmp_path, + timeout_seconds=5, + ) + assert result["ok"] is True + assert result["value"] == {"args": values, "input": {}} diff --git a/tests/control_plane_ts/host_process.test.ts b/tests/control_plane_ts/host_process.test.ts new file mode 100644 index 0000000000..cbea7db0d4 --- /dev/null +++ b/tests/control_plane_ts/host_process.test.ts @@ -0,0 +1,78 @@ +import assert from "node:assert/strict"; +import {mkdtemp, readFile, rm} from "node:fs/promises"; +import {join} from "node:path"; +import {tmpdir} from "node:os"; +import {setTimeout as delay} from "node:timers/promises"; +import test from "node:test"; +import {decodeHostProcessRequest, runHostProcess, type HostProcessRequest} from "../../loopx/control_plane/turn_driver/host_process.ts"; + +const request = (script: string, overrides: Partial = {}): HostProcessRequest => ({ + argv: [process.execPath, "-e", script], cwd: process.cwd(), input: "request\n", + timeout_ms: 3000, drain_timeout_ms: 50, stdout_limit_bytes: 12000, ...overrides, +}); + +test("Host output is streamed with UTF-8 boundaries and stdin EOF", async () => { + let stdout = "", stderr = ""; + const result = await runHostProcess(request(` + let input='';process.stdin.on('data',x=>input+=x);process.stdin.on('end',()=>{ + const b=Buffer.from('界');process.stdout.write(b.subarray(0,1)); + setTimeout(()=>{process.stdout.write(b.subarray(1));process.stderr.write(input)},5) + });`), async item => { if (item.kind === "stdout") stdout += item.text; else stderr += item.text; }); + assert.equal(stdout, "界"); assert.equal(stderr, "request\n"); + assert.equal(result.outcome, "exited"); assert.equal(result.returncode, 0); assert.equal(result.output_complete, true); +}); + +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/); + } + const result = await runHostProcess(request("", {argv: ["/missing/loopx-test-host"]}), async () => {}); + assert.equal(result.outcome, "spawn_failed"); assert.equal(result.returncode, null); +}); + +test("stdout limit cancels a still-running Host before collecting the full stream", async () => { + let received = 0; + const result = await runHostProcess(request(`setInterval(()=>process.stdout.write('x'.repeat(20000)),1)`, + {stdout_limit_bytes: 1000}), async item => { received += item.text.length; }); + assert.equal(result.outcome, "output_limit"); assert.ok(received <= 1000); assert.equal(result.output_complete, false); +}); + +test("abort before start performs no invocation", async () => { + const controller = new AbortController(); controller.abort(); + const result = await runHostProcess(request("throw new Error('must not start')"), async () => assert.fail(), controller.signal); + assert.equal(result.outcome, "cancelled"); assert.equal(result.returncode, null); +}); + +for (const mode of ["timeout", "abort", "leader_exit", "closed_pipes"] as const) { + test(`${mode}: descendants cannot keep working after managed execution returns`, {skip: process.platform === "win32"}, async t => { + const root = await mkdtemp(join(tmpdir(), "loopx-host-group-")); + t.after(() => rm(root, {recursive: true, force: true})); + const marker = join(root, "counter"); + // Ignore TERM so the test proves escalation and does not merely observe a + // cooperative child. Its marker is the semantic oracle, not a PID lookup. + const child = `const fs=require('fs');let n=0;process.on('SIGTERM',()=>{}); + fs.writeFileSync(${JSON.stringify(marker)},String(n)); + setInterval(()=>fs.writeFileSync(${JSON.stringify(marker)},String(++n)),10)`; + const script = `const{spawn}=require('child_process');const fs=require('fs'); + spawn(process.execPath,['-e',${JSON.stringify(child)}],{stdio:${JSON.stringify(mode === "closed_pipes" ? "ignore" : "inherit")}}); + const timer=setInterval(()=>{if(fs.existsSync(${JSON.stringify(marker)})){ + clearInterval(timer);process.stdout.write('ready\\n'); + ${mode === "leader_exit" || mode === "closed_pipes" ? "process.exit(0)" : "setInterval(()=>{},1000)"} + }},5)`; + const controller = new AbortController(); + const result = await runHostProcess(request(script, {timeout_ms: mode === "timeout" ? 500 : 3000}), async item => { + if (mode === "abort" && item.text.includes("ready")) controller.abort(); + }, controller.signal); + assert.equal(result.outcome, mode === "abort" ? "cancelled" : mode === "timeout" ? "timeout" : "exited"); + assert.equal(result.cleanup_scope, "process_group"); assert.equal(result.group_signal_sent, true); + const counter = await readFile(marker, "utf8"); await delay(100); + assert.equal(await readFile(marker, "utf8"), counter, "child kept changing state after return"); + if (mode === "leader_exit") assert.equal(result.output_complete, false); + }); +} + +test("output consumer failure cancels execution rather than leaving an orphan", async () => { + const result = await runHostProcess(request(`setInterval(()=>process.stdout.write('tick\\n'),10)`), + async () => { throw new Error("consumer left"); }); + assert.equal(result.outcome, "cancelled"); assert.equal(result.output_complete, false); +}); diff --git a/tests/test_loopx_turn_codex_cli.py b/tests/test_loopx_turn_codex_cli.py index ac8a80f35f..17a751a49a 100644 --- a/tests/test_loopx_turn_codex_cli.py +++ b/tests/test_loopx_turn_codex_cli.py @@ -1,6 +1,9 @@ from __future__ import annotations import json +import os +import signal +import time import stat import subprocess import sys @@ -163,6 +166,15 @@ def _fake_codex(tmp_path: Path) -> tuple[Path, Path]: "private_material": "must-not-persist" }), flush=True) raise SystemExit(9) +if os.environ.get("FAKE_CODEX_CHILD_MARKER"): + marker = os.environ["FAKE_CODEX_CHILD_MARKER"] + child = subprocess.Popen([sys.executable, "-c", + "import pathlib,signal,time;signal.signal(signal.SIGTERM,signal.SIG_IGN);" + "p=pathlib.Path(" + repr(marker) + ");n=0\\n" + "while True:\\n p.write_text(str(n));n+=1;time.sleep(.01)"]) + pathlib.Path(marker + ".pid").write_text(str(child.pid)) + while not pathlib.Path(marker).exists(): + time.sleep(.01) if os.environ.get("FAKE_CODEX_SLEEP"): time.sleep(float(os.environ["FAKE_CODEX_SLEEP"])) output_path = pathlib.Path(args[args.index("--output-last-message") + 1]) @@ -1089,3 +1101,36 @@ def test_checkpointed_write_approval_is_scoped_and_absent_by_default(): prompt = _prompt(request) assert "only within its active_write_scope" in prompt assert "publish, and production actions retain their gates" in prompt + + +@pytest.mark.skipif(os.name == "nt", reason="POSIX process-group cleanup contract") +@pytest.mark.parametrize("timeout", [False, True]) +def test_codex_cli_reaps_descendants_after_result_or_timeout( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch, timeout: bool, +) -> None: + executable, log_path = _fake_codex(tmp_path) + marker = tmp_path / "child-work" + monkeypatch.setenv("FAKE_CODEX_LOG", str(log_path)) + monkeypatch.setenv("FAKE_CODEX_CHILD_MARKER", str(marker)) + monkeypatch.setattr("loopx.control_plane.turn_driver.codex_cli.OUTPUT_DRAIN_TIMEOUT_SECONDS", .05) + if timeout: + monkeypatch.setenv("FAKE_CODEX_SLEEP", "30") + try: + kwargs = dict(runtime_root=tmp_path / "runtime", project=tmp_path, + codex_bin=str(executable), timeout_seconds=1 if timeout else 5) + if timeout: + with pytest.raises(BuiltInHostError, match="codex_cli_timeout"): + run_codex_cli_host(_request(), **kwargs) + else: + result = run_codex_cli_host(_request(), **kwargs) + assert result["result_kind"] == "validated_progress" + before = marker.read_text() + time.sleep(.15) + assert marker.read_text() == before, "Codex child kept working after adapter returned" + finally: + pid_path = Path(str(marker) + ".pid") + if pid_path.exists(): + try: + os.kill(int(pid_path.read_text()), signal.SIGKILL) + except ProcessLookupError: + pass