diff --git a/docs/reference/goal-chat-continuation.md b/docs/reference/goal-chat-continuation.md index 00569ce34..7f9ab29ce 100644 --- a/docs/reference/goal-chat-continuation.md +++ b/docs/reference/goal-chat-continuation.md @@ -100,7 +100,10 @@ messages; images use the ordinary conversation after pausing. The Chat hard timeout remains in force; this is not an unattended daemon. - To roll back, pause/close the Chat service before installing an older build. Disabling mode or deleting a binding does not cancel already admitted children; - use their own execution/recovery controls and retain their evidence. + stop one with `loopx delegation stop --execute` (or `stop_delegation`), read + its receipt, and retain its evidence. Only `settled` proves the worker + acknowledged and released its locks and its native host exited; stopped work + needs a new operation id. The ordinary native command path also remains available without delegation: @@ -143,6 +146,8 @@ For a disposable mixed-team setup, use the 协调员保留只读沙箱,成员权限来自各自执行绑定,不继承管家的扩大权限。 成员通过验收与协调员报告、整个 Goal 验收分别显示;本模式不直接完成报告 Todo 或整个 Goal。额度是含历史用量的总量,正在执行的请求可能超额,成员另行计量。 -回滚旧版本前先暂停或关闭 Chat 服务;退出或撤销绑定不自动取消已启动的成员。 +回滚旧版本前先暂停或关闭 Chat 服务;退出或撤销绑定不自动取消已启动的成员, +用 `loopx delegation stop --execute`(或 `stop_delegation`)停止单个成员并阅读回执: +只有 `settled` 证明 worker 已确认并释放锁且原生 host 已退出;已停止的工作需要新的 operation id。 此按钮目前限本机 managed Codex Goal 对话,不宣称 Lark、挂接会话或其他主力 驱动等价。可用下方示例准备一次隔离的本地 DSH+云端 Ark 协作。 diff --git a/docs/reference/local-delegation.md b/docs/reference/local-delegation.md index 52b8ec5b6..5462a62ed 100644 --- a/docs/reference/local-delegation.md +++ b/docs/reference/local-delegation.md @@ -157,6 +157,7 @@ delegate start --binding-id independent-review --operation-id review-round-1 \ --brief-file request.json --execute delegate read --operation-id review-round-1 delegate wait --operation-id review-round-1 +delegate stop --operation-id review-round-1 --execute ``` Inspection uses the bound worker workspace as its actual safety scan root. If @@ -221,6 +222,68 @@ own authorized peers supplies `--parent-request-id` on start. CLI and MCP share grant validation, detached execution, wait/readback and recovery rather than maintaining separate rules. +`stop --execute` ends one member's bounded work and returns a receipt that +states what was proven. The request is written beside the execution record +(`.stop.json`), never into it, so a worker that is still holding the +operation cannot overwrite it. A worker on this machine receives `SIGTERM` for +its whole process group, which ends its Turn child; the native host runs in +its own process group, and its supervisor terminates that group once the Turn +child is gone. The worker acknowledges from under its own lock, marks the +record `stopped` and releases its hard task lease. When nobody holds the operation, the requester +acknowledges itself. A worker on another machine is never signalled; it finds +the request at its next checkpoint or at its next record write, which is +refused. The receipt `phase` is `settled` only when an acknowledgement exists, +the operation lock is free, the member's Turn lane holder record shows it +released by the stopped worker (the lane is read, never taken), the native +host the Turn launched has exited together with every process in its group, +and a required hard task lease was actually released. The obligation is read +from the operation record, not from the stop sidecar: the acknowledgement is +persisted before the lease is released, so a crash in between must not turn +"not yet written" into "nothing was owed". A release that failed is +retried under the stop's own lock on the next read, so it never becomes a +`settled` receipt that leaves the member's Todo blocked until the lease TTL; +while it is unproven the stop stays `acknowledged` with +`required_lease_release_unproven`. If that host cannot be attributed, or its +supervisor never finished cleaning up, the stop stays `acknowledged` and a +later `stop` rereads it. On a platform without process groups the launched +host cannot be proven drained at all, so `stop --execute` fails with an +actionable error naming that boundary rather than leaving a receipt no read +can settle. `unknown` +means the holder vanished before acknowledging, and `noop` means the +work was already accepted, rejected or stopped. `requested` or `acknowledged` +means it is still winding down: call `stop` again. A grace timeout never turns +into a receipt. Stopped work is not resumed; `resume` refuses it and a new +scope needs a new operation id. The Turn journal keeps its `in_progress` entry +for inspection, and the record is never rewritten as a completion. A stopped +member's Todo stays open, so the coordinator decides what happens next. The +member's Todo completion and reply publication commit under the same lock a +stop takes, so a stop written first means neither effect lands, and a stop +written after both leaves their acceptance intact. + +中文:`stop --execute` 结束一个成员的有界工作,并返回一份只陈述已证明事实的 +回执。停止请求写在执行记录旁边的 `.stop.json`,从不写进记录本身, +因此仍持有该 operation 的 worker 无法覆盖它。本机 worker 会收到整个进程组的 +`SIGTERM`,其 Turn 子进程随之结束;原生 host 在自己的进程组中运行,Turn 子进程 +退出后由其 supervisor 终止整个 host 进程组。worker 在自己的锁下确认,把记录标为 +`stopped` 并释放硬任务租约。没有持有者时由请求方自行确认。另一台机器上的 +worker 不会被发信号,它在下一个检查点或下一次写记录时发现请求,写入被拒绝。 +只有存在确认、operation 锁已释放、成员 Turn lane 的持有者记录显示已被停止的 +worker 释放(只读 lane,从不获取)、该 Turn 启动的原生 host 及其进程组内所有进程 +都已退出,且必需的硬任务租约确实释放成功时,`phase` 才是 `settled`。该义务取自 +操作记录而非 stop sidecar:确认会先于释放落盘,因此两者之间发生崩溃时,不能把 +「尚未写入」当成「本就不需要释放」。释放失败会在下一次读取时于 stop 自己的锁下重试, +因此不会产生一份「已结算」却让成员 Todo 被租约阻塞到 TTL 的回执;在释放得到证明前,停止保持 `acknowledged`,原因为 +`required_lease_release_unproven`。host 无法归属或其 supervisor 未完成清理时, +停止保持 `acknowledged`,之后再次调用 `stop` 会重新读取。在没有进程组的平台上, +启动过的 host 根本无法被证明已收尾,因此 `stop --execute` 会以指明该平台边界的 +可操作错误失败,而不是留下一份任何读取都无法结算的回执。`unknown` 表示持有者在确认前消失;`noop` 表示工作已 accepted、 +rejected 或 stopped;`requested`/`acknowledged` 表示仍在收尾,再次调用 `stop`。 +宽限期超时永远不会变成回执。已停止的工作不能 `resume`,新范围需要新的 +operation id。Turn journal 保留 `in_progress` 条目供检查,记录不会被改写成完成; +成员的 Todo 仍然打开,由协调者决定下一步。成员的 Todo 完成与回执发布在 stop +所取的同一把锁下提交,因此先写入停止则两个效果都不会落地,后写入停止则其验收结果 +保持不变。 + This entrypoint does not create Agents, grant bindings or wake an idle Codex conversation. The existing host/LoopX continuation policy owns the next lead turn. The conversation remains persistent independently of whether autonomous @@ -583,7 +646,8 @@ operation and observed artifact hashes in its existing inbox. Pending and delive receipts stay distinct from application; a retry after an uncertain response reuses the exact message and operation id. **Pause coordinator** stays in the panel and reports its actual scope. Dispatched members continue independently; -this entrypoint cannot stop the whole team. Ordinary polling does not read artifact +this entrypoint cannot stop the whole team. Stop one member explicitly with +`delegation stop --execute` or `stop_delegation` and read its receipt. Ordinary polling does not read artifact bodies or run preflight. Closing the panel changes no work state. This local operator entrypoint does not grant a Lark audience access. @@ -593,7 +657,8 @@ operator entrypoint does not grant a Lark audience access. 可展开查看,返回列表保留位置与键盘焦点。协调员运行时,可把执行标识、看到的 产物哈希和反馈投递到原收件箱;等待投递、已交付和已应用不能混为一谈。不确定响应 后重试同一消息和标识,避免重复投递。面板内的「暂停协调员」显示实际反馈,但不会 -停止已派发成员,也不宣称整个团队停止。暂停时仍可检查证据;读取不启动模型。 +停止已派发成员,也不宣称整个团队停止。要停止某个成员,显式使用 +`delegation stop --execute` 或 `stop_delegation` 并阅读其回执。暂停时仍可检查证据;读取不启动模型。 Screenshots use isolated synthetic research data, not a live-model qualification: [desktop evidence](../assets/personal-workspace/team-evidence-desktop.png), @@ -623,6 +688,10 @@ unchanged and cannot launch workers. With it, the Agent can: response is normal. `read_delegation` reads the durable original operation. 4. If `recovery_required` is true, call `resume_delegation` with that same id. This cannot retarget the work or silently create a replacement Turn. +5. Call `stop_delegation(operation_id)` to end one member. Read its `phase`: + `settled` is the only receipt that the worker acknowledged and released its + locks and that the native host and its process group exited; `unknown` means the holder vanished first; `noop` means the work had + already ended. Stopped work cannot be resumed; use a new operation id. Configure the member's host to expose its own identity-bound collaboration tools. It reads `DELEGATION.json`, independently calls `assess_request`, and @@ -663,6 +732,7 @@ concurrent executions still use the same kernel lock and original Turn journal. | Requesting MCP conversation closes | The detached bounded worker continues; another connection reads the original operation. | | Duplicate start/resume while work runs | Operation identity, task lock and Turn journal prevent another concurrent execution. | | Worker process or machine stops | Reconnect with the same operator configuration and credentials, then resume the original Turn. | +| Member stopped on request | The worker acknowledges under its lock, its Turn child is ended and the host supervisor terminates the host group, its lease is released; `settled` needs that acknowledgement, free locks and an exited host group, `unknown` means the holder vanished first. The record is `stopped`; resume refuses it. | | Ark is computing without local tools | The already-started cloud turn can continue. It is not dependent on the local conversation. | | Ark requests a local tool while the host is absent | It waits for the local tool result. Recovery observes the original input/session and executes only previously unstarted tool calls. | | Tool execution or send acknowledgement is uncertain | Do not repeat the effect. Preserve the receipt/session for explicit reconciliation. | @@ -678,7 +748,8 @@ or default executor change; those existing configuration surfaces are untouched. To disable new admission, remove the caller's grants or remove `--execution-config` from the host. A stopped Goal refuses new starts/resumes; existing completed results remain readable. Disabling does not kill work already -running. Retain receipts, stop or reconcile owned workers, and confirm cloud +running; `delegation stop --execute` ends one member and returns a receipt. +Retain receipts, stop or reconcile owned workers, and confirm cloud resource cleanup before deleting a disposable runtime. The optional adapter's cleanup command never grants task completion. diff --git a/loopx/cli_commands/delegation.py b/loopx/cli_commands/delegation.py index 9d2fb72a9..e422e68c6 100644 --- a/loopx/cli_commands/delegation.py +++ b/loopx/cli_commands/delegation.py @@ -22,7 +22,7 @@ def register_delegation( "delegation", help="Launch and recover authorized peer work; returns JSON." ) add_format(parser) - parser.add_argument("delegation_action", choices=("list", "operations", "inspect", "start", "read", "wait", "resume", "adopt")) + parser.add_argument("delegation_action", choices=("list", "operations", "inspect", "start", "read", "wait", "resume", "adopt", "stop")) parser.add_argument("--goal-id", required=True) parser.add_argument("--agent-id", required=True, help="Calling registered Agent, not the worker.") parser.add_argument("--execution-config", type=Path, required=True, @@ -34,7 +34,7 @@ def register_delegation( parser.add_argument("--parent-request-id", help="For start: the request received by this coordinator.") parser.add_argument("--limit", type=int, help="For operations: page size, 1–50 (default 20).") parser.add_argument("--cursor", help="For operations: next_cursor returned by the previous page.") - parser.add_argument("--execute", action="store_true", help="Required for start/resume/adopt; grants no additional authority.") + parser.add_argument("--execute", action="store_true", help="Required for start/resume/adopt/stop; grants no additional authority.") def handle_delegation( @@ -45,10 +45,10 @@ def handle_delegation( action = args.delegation_action try: - if action in {"start", "resume", "adopt"} and not args.execute: + if action in {"start", "resume", "adopt", "stop"} and not args.execute: raise ValueError(f"delegation {action} requires --execute") - if action not in {"start", "resume", "adopt"} and args.execute: - raise ValueError("--execute is only valid for start/resume/adopt") + if action not in {"start", "resume", "adopt", "stop"} and args.execute: + raise ValueError("--execute is only valid for start/resume/adopt/stop") if action not in {"list", "operations", "inspect"} and not args.operation_id: raise ValueError(f"delegation {action} requires --operation-id") if action in {"list", "operations", "inspect"} and args.operation_id: @@ -86,6 +86,8 @@ def handle_delegation( result = service.read(args.operation_id) elif action == "wait": result = service.wait(args.operation_id) + elif action == "stop": + result = service.stop(args.operation_id, execute=args.execute) else: result = service.resume(args.operation_id) payload = {"ok": True, **result} diff --git a/loopx/collaboration_mcp.py b/loopx/collaboration_mcp.py index 426c0916b..bcbce71d5 100644 --- a/loopx/collaboration_mcp.py +++ b/loopx/collaboration_mcp.py @@ -14,6 +14,7 @@ import hashlib import json import os +import signal import stat import subprocess import sys @@ -26,7 +27,10 @@ if TYPE_CHECKING: from mcp.server.fastmcp import FastMCP -from .file_lock import exclusive_file_lock, LockAcquisitionPolicy, LockAcquireTimeoutError +from .file_lock import ( + exclusive_file_lock, lock_holder_host_label, lock_holder_liveness, + LOCK_HOLDER_FOREIGN_HOST, LOCK_HOLDER_LIVE, LockAcquisitionPolicy, LockAcquireTimeoutError, +) from .control_plane.effect_runtime import ( effect_runtime_request_scope, effect_runtime_result, EffectRuntimeRemoteError, ) @@ -38,7 +42,19 @@ turn_journal_path, ) from .control_plane.turn_driver.host_binding import turn_host_arg_option +from .control_plane.turn_driver.host_process_transport import ( + HOST_PROCESS_DRAINING, HOST_PROCESS_RECORD_ENV, HOST_PROCESS_UNSUPPORTED_PLATFORM, + host_process_drain, +) +from .control_plane.turn_driver.lane_fence import ( + TURN_LANE_ABSENT, TURN_LANE_DEAD, TURN_LANE_LIVE, TURN_LANE_RELEASED, + turn_lane_liveness, turn_lane_target, +) +from .control_plane.work_items.task_lease import release_task_lease from .control_plane.collaboration.inbox import _hash, _read, _write, _root, _receipt +from .control_plane.collaboration.delegation_inventory import ( + DELEGATION_HOST_PROCESS_SUFFIX, DELEGATION_STOP_RECEIPT_SUFFIX, +) from .control_plane.collaboration.peers import return_result from .control_plane.collaboration.inbox import acknowledge, _entry, normalize_request from .control_plane.collaboration.goal_instance_scope import ( @@ -325,6 +341,84 @@ def consume_peer_result(request_id: str) -> dict: ) +DELEGATION_STOP_SCHEMA_VERSION = "loopx_delegation_stop_v0" +# Observations that no worker may reopen; a stop against one is a no-op receipt. +DELEGATION_TERMINAL_STATUSES = frozenset({"accepted", "rejected", "stopped"}) +DELEGATION_STOP_OPEN_PHASES = frozenset({"requested", "acknowledged"}) +DELEGATION_STOP_TERMINAL_PHASES = frozenset({"settled", "unknown"}) +# How long a signalled same-host worker may take to acknowledge before SIGKILL. +DELEGATION_STOP_GRACE_SECONDS = 10.0 +DELEGATION_STOPPED_MESSAGE = "delegation operation was stopped; start a new operation id" + + +class DelegationStopRequested(BaseException): + """A stop reached the worker that owns this operation; it must acknowledge, not finish. + + A ``BaseException`` like ``KeyboardInterrupt``: a termination request must + not be swallowed by an ``except Exception`` and turned into further work. + """ + + def __init__(self, source: str) -> None: + super().__init__(source) + self.source = source + + +class DelegationFenced(DelegationStopRequested): + """A stop this process never acknowledged fences its execution-record write. + + The write is refused before it happens: a late-returning or other-host + worker records no Turn result, completes no Todo and publishes nothing. + """ + + def __init__(self) -> None: + super().__init__("fenced") + + +class _WorkerStopSignal: + """Turn SIGTERM into a stop request only when a stop was written for this operation. + + Without a stop receipt the signal keeps its default meaning, so a shutdown + still leaves the operation recoverable by ``resume`` instead of stopping it. + Later signals are absorbed while the acknowledgement is written. + """ + + def __init__(self, stop_path: Path) -> None: + self.stop_path = stop_path + self.armed = True + + def __call__(self, signum: int, frame: object) -> None: + if not self.armed: + return + try: + requested = self.stop_path.exists() + except OSError: + requested = False + if not requested: + signal.signal(signum, signal.SIG_DFL) + os.kill(os.getpid(), signum) + return + self.armed = False + raise DelegationStopRequested("SIGTERM") + + def disarm(self) -> None: + self.armed = False + + +def install_worker_stop_signal(stop_path: Path) -> _WorkerStopSignal | None: + """Install the detached worker's SIGTERM handler; ``None`` where signals are unsupported.""" + + if not hasattr(signal, "SIGTERM"): + return None + handler = _WorkerStopSignal(stop_path) + try: + signal.signal(signal.SIGTERM, handler) + except (ValueError, OSError): + # Not the main thread, or a platform without handler support: the + # worker still honours stop files at every checkpoint and fenced write. + return None + return handler + + class Delegations: """Host IO for bound peer work; typed grants and observations stay in TS. @@ -336,6 +430,7 @@ def __init__(self, root: Path, registry: Path, goal_id: str, agent_id: str, conf self.root, self.registry = root.resolve(), registry.resolve() self.goal_id, self.agent_id, self.config = goal_id, agent_id, config.resolve() self._goal_ref_lock = Lock() + self._stop_signal: _WorkerStopSignal | None = None try: self.goal_ref = capture_collaboration_goal_ref( self.registry, @@ -377,6 +472,27 @@ def directory(self) -> dict: def path(self, operation_id: str) -> Path: return _root(self.root) / "executions" / _hash([self.goal_id, self.agent_id]) / (_hash(operation_id) + ".json") + @staticmethod + def _stop_path(path: Path) -> Path: + """The stop receipt sits beside its execution record and is never merged into it.""" + return path.with_name(path.stem + DELEGATION_STOP_RECEIPT_SUFFIX) + + @staticmethod + def _host_process_record(path: Path) -> Path: + """Where this operation's Turn names the native Host it launched, for drain readback.""" + return path.with_name(path.stem + DELEGATION_HOST_PROCESS_SUFFIX) + + @staticmethod + def _dispatch_lock(path: Path) -> Path: + return path.with_suffix(".dispatch") + + @staticmethod + def _read_stop(path: Path) -> dict | None: + stop_path = Delegations._stop_path(path) + if not stop_path.exists(): + return None + return _read(stop_path) + def operations(self, *, limit: int = 20, cursor: str | None = None) -> dict: from .control_plane.collaboration.delegation_inventory import read_delegation_inventory @@ -557,15 +673,22 @@ def _spawn(self, operation_id: str) -> None: def resume(self, operation_id: str) -> dict: path = self.path(operation_id) + if self._read_stop(path) is not None: + raise ValueError(DELEGATION_STOPPED_MESSAGE) try: with exclusive_file_lock( path, policy=LockAcquisitionPolicy.SINGLE_FLIGHT ): row = _read(path) binding = self._bound(row) + if row["status"] == "stopped": + raise ValueError(DELEGATION_STOPPED_MESSAGE) if row["status"] == "rejected": - self._recover_validated_settlement(path, row, binding) - should_spawn = row["status"] not in {"accepted", "rejected"} + try: + self._recover_validated_settlement(path, row, binding) + except DelegationFenced: + raise ValueError(DELEGATION_STOPPED_MESSAGE) from None + should_spawn = row["status"] not in DELEGATION_TERMINAL_STATUSES if should_spawn: self.binding( row["identity"]["binding"]["id"], require_active=True @@ -582,7 +705,8 @@ def wait(self, operation_id: str) -> dict: """Observe for at most 15 seconds; waiting neither starts nor resumes work.""" for _ in range(5): result = self.read(operation_id) - if result["status"] in {"accepted", "rejected"} or result["recovery_required"]: + if (result["status"] in DELEGATION_TERMINAL_STATUSES or result["recovery_required"] + or result.get("stop", {}).get("phase") in DELEGATION_STOP_TERMINAL_PHASES): return result time.sleep(3) return self.read(operation_id) @@ -669,7 +793,7 @@ def _recover_validated_settlement( "error": None, } row.pop("error", None) - _write(path, row) + self._fenced_write(path, row) return True def adopt_result(self, operation_id: str, consumer_operation_id: str) -> dict: @@ -692,11 +816,15 @@ def _read_current(self, operation_id: str) -> dict: active = False except LockAcquireTimeoutError: active = True + # Resume refuses an operation with a stop receipt, so it never needs recovery. + stop = self._read_stop(path) result = {"operation_id": operation_id, "request_id": row["identity"]["request_id"], "agent_id": binding["agent_id"], "todo_id": binding["todo_id"], "status": row["status"], "worker_active": active, - "recovery_required": not active and row["status"] not in {"accepted", "rejected"} - and time.time() - row.get("created_at", 0) > 15} + "recovery_required": not active and row["status"] not in DELEGATION_TERMINAL_STATUSES + and stop is None and time.time() - row.get("created_at", 0) > 15} + if stop is not None: + result["stop"] = {"stop_id": stop["stop_id"], "phase": stop["phase"]} if row["status"] == "accepted": # A saved receipt cannot hide an amended task, verifier or output. artifacts = self._accepted(binding) @@ -707,19 +835,398 @@ def _read_current(self, operation_id: str) -> dict: result["error"] = row["error"] return result - def _observe(self, path: Path, row: dict, status: str, **facts) -> None: + def _observe(self, path: Path, row: dict, status: str, *, already_locked: bool = False, **facts) -> None: decision = effect_runtime_result("collaboration.delegation.observe", { "from": row["status"], "to": status, **facts, }) row.update(status=decision["status"]) + if already_locked: + self._fenced_write_locked(path, row) + else: + self._fenced_write(path, row) + + def _fenced_write(self, path: Path, row: dict) -> None: + """Write the execution record only while no unacknowledged stop fences this process. + + The stop receipt is re-read under the dispatch lock on every write, so a + worker that returns after a stop it never saw writes nothing at all. + """ + + with exclusive_file_lock(self._dispatch_lock(path)): + self._fenced_write_locked(path, row) + + def _fenced_write_locked(self, path: Path, row: dict) -> None: + """The same write for a caller that already holds the dispatch lock. + + The lock is a kernel file lock, so it is not reentrant: a caller that + widened its critical section to cover a whole effect group must use this + entry point rather than nesting ``_fenced_write``. + """ + + stop = self._read_stop(path) + if stop is not None and not self._acknowledged_here(stop): + raise DelegationFenced() _write(path, row) - def _cli(self, binding: dict, *args: str, timeout: int = 60) -> dict: + @staticmethod + def _acknowledged_here(stop: dict) -> bool: + ack = stop.get("ack") + return (isinstance(ack, dict) and ack.get("pid") == os.getpid() + and ack.get("host") == lock_holder_host_label()) + + def _raise_if_stop_requested(self, path: Path) -> None: + """Worker checkpoint: leave before the next host launch or Todo effect.""" + if self._read_stop(path) is not None: + raise DelegationStopRequested("stop_file") + + @staticmethod + def _worker_identity() -> dict: + return { + "pid": os.getpid(), + "pgid": os.getpgid(0) if hasattr(os, "getpgid") else None, + "host": lock_holder_host_label(), + } + + def _lane_target(self, binding: dict) -> Path: + return turn_lane_target(runtime_root=self.root, goal_id=self.goal_id, + plan={"turn_envelope": {"agent_id": binding["agent_id"]}}) + + def _operation_lock_free(self, path: Path) -> bool: + """Probe this operation's own kernel lock; only ever called once its stop receipt exists. + + Unlike the Turn lane, this lock admits nothing but this operation, and a + probe holding it for an instant refuses no legitimate acquisition once + the receipt is written: ``resume``, its only single-flight acquirer, + refuses a stopped operation before it touches the lock; ``execute``, + adoption and the requester acknowledgement wait through brief holders + with the mutation policy; and a status read already makes this same + instant observation. Before the receipt exists a ``resume`` is still + legitimate, so the holder is then read from its record instead. + """ + + try: + with exclusive_file_lock(path, policy=LockAcquisitionPolicy.SINGLE_FLIGHT): + return True + except LockAcquireTimeoutError: + return False + + @staticmethod + def _recorded_worker(row: dict, stop: dict) -> dict | None: + worker = row.get("worker") + if not isinstance(worker, dict): + worker = stop.get("worker") + return worker if isinstance(worker, dict) else None + + def _worker_lane_released(self, row: dict, stop: dict, binding: dict) -> tuple[bool, str]: + """Say whether the stopped worker's Turn has let go of the member's lane, read-only. + + This never takes the lane lock: a probe holding it for an instant would + refuse a legitimate Turn of the same member racing that instant with + ``turn_lane_in_flight``. The lane's last holder record decides instead. + Released, dead or absent is released. A live holder on this machine is + released only when it sits outside the recorded worker's process group, + because the worker's run-once child runs in that group; a holder that + cannot be attributed, another host's holder and an unreadable record + prove nothing, so the typed decision keeps the stop open. + """ + + lane = turn_lane_liveness(self._lane_target(binding)) + state = lane["state"] + if state in {TURN_LANE_RELEASED, TURN_LANE_DEAD, TURN_LANE_ABSENT}: + return True, state + worker = self._recorded_worker(row, stop) + if (state != TURN_LANE_LIVE or worker is None or not hasattr(os, "getpgid") + or worker.get("host") != lock_holder_host_label() + or not isinstance(worker.get("pgid"), int)): + return False, state + try: + return os.getpgid(lane["holder"]["pid"]) != worker["pgid"], state + except ProcessLookupError: + return True, TURN_LANE_DEAD # the holder exited between the two reads + except OSError: + return False, state + + def _turn_journal_status(self, row: dict, binding: dict) -> str | None: + turn_key = row.get("turn_key") or self._matching_turn_key(row, binding) + if not turn_key: + return None + journal = load_turn_journal(turn_journal_path(self.root, goal_id=self.goal_id, turn_key=turn_key)) + status = journal.get("status") if isinstance(journal, dict) else None + return str(status) if status else None + + def _new_stop_record(self, row: dict, *, requested_by: str, worker: dict | None) -> dict: + requested_at = time.time() + return { + "schema_version": DELEGATION_STOP_SCHEMA_VERSION, + "stop_id": _hash([row["identity"]["operation_id"], requested_by, requested_at])[:32], + "operation_id": row["identity"]["operation_id"], + "request_id": row["identity"]["request_id"], + "phase": "requested", + "reason": "awaiting_acknowledgement", + "requested_by": requested_by, + "requested_at": requested_at, + "requested_status": row["status"], + "worker": worker, + "ack": None, + "lease": None, + "settled": None, + } + + def _stop_receipt(self, row: dict, binding: dict, stop: dict | None) -> dict: + receipt = { + "operation_id": row["identity"]["operation_id"], + "request_id": row["identity"]["request_id"], + "agent_id": binding["agent_id"], "todo_id": binding["todo_id"], + "status": row["status"], + } + if stop is None: + # Nothing was written: a terminal observation cannot be stopped, and + # repeating the request returns exactly this receipt again. + receipt.update(phase="noop", reason="delegation already " + row["status"], stop=None) + else: + receipt.update(phase=stop["phase"], reason=stop.get("reason"), stop=stop) + return receipt + + def stop(self, operation_id: str, *, execute: bool) -> dict: + """Stop one bounded member and return a receipt that says what was proven. + + ``requested`` is written beside the execution record, never into it. + When no worker holds the operation, this caller takes the lock, marks + the record stopped and releases the hard lease itself. A same-host + holder is signalled by process group and given a bounded grace to + acknowledge; another host's holder is left to find the request at its + next checkpoint or fenced write. ``settled`` and ``unknown`` come from + the typed decision over lock facts and the drain of the native Host the + Turn launched, whose TS supervisor alone terminates it; elapsed time + proves nothing. + """ + + require_operation_id(operation_id) + if not execute: + raise ValueError("delegation stop requires execute") + path = self.path(operation_id) + if not path.exists(): + raise ValueError("unknown delegation operation; start_delegation returns the operation_id to stop") + with exclusive_file_lock(self._dispatch_lock(path)): + row = _read(path) + binding = self._bound(row) + stop = self._read_stop(path) + if stop is None: + if row["status"] in DELEGATION_TERMINAL_STATUSES: + return self._stop_receipt(row, binding, None) + stop = self._new_stop_record(row, requested_by=self.agent_id, + worker=self._lock_holder_worker(path, row)) + _write(self._stop_path(path), stop) + if stop["phase"] in DELEGATION_STOP_OPEN_PHASES and stop.get("ack") is None: + if stop.get("worker") is None: + # No worker was named when the request was written, so whoever owns + # the operation acknowledges it: this caller once the lock is free. + # The wait rides out a status read's instant hold, which must not + # be mistaken for a holder that vanished. + try: + with exclusive_file_lock(path): + self._acknowledge_stop(path, _read(path), binding, source="requester") + except LockAcquireTimeoutError: + pass # an unnamed holder meets the request at its next checkpoint or write + else: + # A Host left behind by an earlier worker may still be terminating. + self._await_host_drain(path, time.monotonic() + DELEGATION_STOP_GRACE_SECONDS) + else: + # Only the named worker acknowledges. If it vanishes first, the typed + # decision reports unknown instead of a requester settlement. + self._signal_worker(path, stop) + return self._settle_stop(path) + + def _lock_holder_worker(self, path: Path, row: dict) -> dict | None: + """Name the recorded worker while it is the operation lock's unreleased holder. + + Read from the holder record, never the kernel lock: this runs before the + stop receipt exists, when a probe could refuse a legitimate ``resume``. + Only the worker identity the execution record names can become a signal + target, so a status reader's instant holder record is never taken for it. + """ + + state, holder = lock_holder_liveness(path) + recorded = row.get("worker") + if state not in {LOCK_HOLDER_LIVE, LOCK_HOLDER_FOREIGN_HOST} or not isinstance(recorded, dict): + return None + if holder.get("pid") != recorded.get("pid") or holder.get("host") != recorded.get("host"): + return None + return {key: recorded.get(key) for key in ("pid", "pgid", "host")} + + def _signal_worker(self, path: Path, stop: dict) -> None: + """Terminate a same-host holder's process group; never signal across hosts.""" + + worker = stop.get("worker") + if (not isinstance(worker, dict) or worker.get("host") != lock_holder_host_label() + or not hasattr(os, "killpg") or not isinstance(worker.get("pid"), int)): + return + pid, pgid = worker["pid"], worker.get("pgid") or worker["pid"] + if pgid == os.getpgid(0): + raise ValueError("delegation stop refuses to signal its own process group") + try: + if os.getpgid(pid) != pgid: + return # the pid was reused by an unrelated process + os.killpg(pgid, signal.SIGTERM) + except ProcessLookupError: + return + # Poll the release facts the decision needs; the deadline only bounds this + # call, and a Host still draining when it passes leaves the stop open. + deadline = time.monotonic() + DELEGATION_STOP_GRACE_SECONDS + while time.monotonic() < deadline: + if self._operation_lock_free(path): + self._await_host_drain(path, deadline) + return + time.sleep(0.2) + try: + os.killpg(pgid, signal.SIGKILL) + except ProcessLookupError: + return + deadline = time.monotonic() + 2.0 + while time.monotonic() < deadline and not self._operation_lock_free(path): + time.sleep(0.1) + # The Host supervisor sits outside the worker's group and cleans up on its own. + self._await_host_drain(path, time.monotonic() + DELEGATION_STOP_GRACE_SECONDS) + + def _await_host_drain(self, path: Path, deadline: float) -> None: + """Wait, never kill: the TS Host supervisor owns terminating its process group.""" + while (host_process_drain(self._host_process_record(path)) == HOST_PROCESS_DRAINING + and time.monotonic() < deadline): + time.sleep(0.1) + + def _acknowledge_stop(self, path: Path, row: dict, binding: dict, *, source: str) -> None: + """Acknowledge from under the operation lock: mark stopped, then release the hard lease. + + Only the operation-lock holder calls this. The record is transitioned as + it is on disk, so state that a fenced write refused stays unwritten; the + lease is released from what this process acquired, which may be newer + than the record. A missing stop, one already acknowledged or finished, + and a record that already reached a terminal observation stay untouched. + """ + + if self._stop_signal is not None: + self._stop_signal.disarm() + with exclusive_file_lock(self._dispatch_lock(path)): + stop = self._read_stop(path) + current = _read(path) + if (stop is None or stop.get("ack") is not None + or stop["phase"] not in DELEGATION_STOP_OPEN_PHASES + or current["status"] in DELEGATION_TERMINAL_STATUSES): + return + observed = current["status"] + transition = effect_runtime_result("collaboration.delegation.observe", { + "from": observed, "to": "stopped", + }) + # This process holds the operation lock; its lane and Host drain are read later. + phase = effect_runtime_result("collaboration.delegation.stop", { + "phase": stop["phase"], "acknowledged": True, + "operation_lock_free": False, "worker_lane_released": False, + "host_process": HOST_PROCESS_DRAINING, + }) + current["status"] = transition["status"] + _write(path, current) + stop.update(phase=phase["phase"], reason=phase["reason"], ack={ + "pid": os.getpid(), "host": lock_holder_host_label(), "at": time.time(), + "source": source, "observed_status": observed, + "turn_key": (current.get("turn_key") or row.get("turn_key") + or self._matching_turn_key(current, binding)), + }) + _write(self._stop_path(path), stop) + try: + self._clear_delegation_bootstrap(row, binding) + except (OSError, ValueError): + pass # the bootstrap is host input; its state never blocks the receipt + lease = self._release_delegation_lease(row, binding) + with exclusive_file_lock(self._dispatch_lock(path)): + stop = self._read_stop(path) or stop + stop["lease"] = lease + _write(self._stop_path(path), stop) + + def _release_delegation_lease(self, row: dict, binding: dict) -> dict: + lease = row.get("task_lease") + if not isinstance(lease, dict) or lease.get("required") is not True: + return {"required": False, "released": None} + try: + result = release_task_lease( + runtime_root=self.root, goal_id=self.goal_id, todo_id=binding["todo_id"], + owner=binding["agent_id"], idempotency_key=str(lease["idempotency_key"]), + expected_version=lease.get("version"), registry_path=self.registry, + ) + except (ValueError, OSError, RuntimeError, EffectRuntimeRemoteError) as exc: + return {"required": True, "released": False, "error": str(exc)[:180]} + return {"required": True, "released": result.get("released") is True, + "missing": result.get("missing") is True} + + def _settle_stop(self, path: Path) -> dict: + with exclusive_file_lock(self._dispatch_lock(path)): + row = _read(path) + binding = self._bound(row) + stop = self._read_stop(path) + if stop is None: + return self._stop_receipt(row, binding, None) + if stop["phase"] not in DELEGATION_STOP_OPEN_PHASES: + return self._stop_receipt(row, binding, stop) + # A launched Host on a platform that cannot prove its group exited + # has no converging stop: say so plainly rather than leaving the + # caller with an acknowledged receipt it can never settle. + unsupported = host_process_drain(self._host_process_record(path)) == HOST_PROCESS_UNSUPPORTED_PLATFORM + if unsupported: + raise ValueError( + "delegation stop cannot prove the launched Host drained on this platform: " + "process groups are unavailable, so the Host supervisor is best-effort. " + "Stop the member's Host through its own supervisor and re-read the receipt." + ) + facts = {"operation_lock_free": self._operation_lock_free(path)} + facts["worker_lane_released"], lane_state = self._worker_lane_released(row, stop, binding) + # Read last: a Host seen drained after its worker and lane let go stays drained. + facts["host_process"] = host_process_drain(self._host_process_record(path)) + # The obligation comes from the operation record, not from the stop + # sidecar. The acknowledgement writes its ACK before it releases the + # lease, so a process loss in between leaves the sidecar with no + # `lease` field at all — and reading that as "nothing was owed" + # settles a stop whose member still holds an active hard lease, with + # resume already refused and the Todo blocked until the TTL. + # + # The canonical `row.task_lease.required` is the source of truth, so + # the obligation survives the crash. A release is retried here, under + # the same lock that guards the record: every later read is another + # attempt rather than one failure becoming permanent. + owed = isinstance(row.get("task_lease"), dict) and row["task_lease"].get("required") is True + lease = stop.get("lease") if isinstance(stop.get("lease"), dict) else {} + if owed and lease.get("released") is not True: + retried = self._release_delegation_lease(row, binding) + if retried != lease: + stop["lease"] = retried + _write(self._stop_path(path), stop) + lease = retried + if owed: + facts["lease_released"] = lease.get("released") is True + decision = effect_runtime_result("collaboration.delegation.stop", { + "phase": stop["phase"], "acknowledged": stop.get("ack") is not None, + "timed_out": time.time() - stop["requested_at"] > DELEGATION_STOP_GRACE_SECONDS, + **facts, + }) + if decision["phase"] != stop["phase"] or decision.get("reason") != stop.get("reason"): + stop.update(phase=decision["phase"], reason=decision.get("reason")) + if decision["phase"] in DELEGATION_STOP_TERMINAL_PHASES: + stop["settled"] = { + "at": time.time(), **facts, "lane_state": lane_state, + "turn_journal_status": self._turn_journal_status(row, binding), + } + _write(self._stop_path(path), stop) + return self._stop_receipt(row, binding, stop) + + def _cli(self, binding: dict, *args: str, timeout: int = 60, host_record: Path | None = None) -> dict: + environment = _pinned_release_environment() + environment.pop(HOST_PROCESS_RECORD_ENV, None) + if host_record is not None: + # The Turn's Host transport names the process group its supervisor owns. + environment[HOST_PROCESS_RECORD_ENV] = str(host_record) completed = subprocess.run([*_python_module_command("loopx.cli"), "--registry", str(self.registry), "--runtime-root", str(self.root), "--format", "json", *args, ], cwd=binding["workspace"], capture_output=True, text=True, encoding="utf-8", - timeout=timeout, env=_pinned_release_environment()) + timeout=timeout, env=environment) try: value = json.loads(completed.stdout) except ValueError as exc: @@ -769,21 +1276,34 @@ def execute(self, operation_id: str) -> None: # The existing bounded mutation policy still excludes concurrent workers. with exclusive_file_lock(path): row = _read(path) - if row["status"] in {"accepted", "rejected"}: + if row["status"] in DELEGATION_TERMINAL_STATUSES: return - binding = self._bound(row, require_active=True) - row.pop("error", None) - _write(path, row) - # Different request ids cannot run the same assigned task concurrently. - task_lock = _root(self.root) / "execution-slots" / _hash([self.goal_id, binding["todo_id"]]) - with exclusive_file_lock(task_lock, policy=LockAcquisitionPolicy.SINGLE_FLIGHT): - try: - self._execute(path, row, binding) - except (ValueError, KeyError, subprocess.TimeoutExpired, EffectRuntimeRemoteError) as exc: - row["error"] = str(exc)[:180] if isinstance(exc, (ValueError, EffectRuntimeRemoteError)) else type(exc).__name__ - _write(path, row) - if row["status"] == "prepared": - self._observe(path, row, "rejected") + binding = self._bound(row) + if self._read_stop(path) is not None: + # The stop arrived before any worker owned the operation: this + # holder acknowledges it from under the lock and launches nothing. + self._acknowledge_stop(path, row, binding, source="worker_entry") + return + try: + binding = self._bound(row, require_active=True) + row.pop("error", None) + row["worker"] = self._worker_identity() + self._fenced_write(path, row) + # Different request ids cannot run the same assigned task concurrently. + task_lock = _root(self.root) / "execution-slots" / _hash([self.goal_id, binding["todo_id"]]) + with exclusive_file_lock(task_lock, policy=LockAcquisitionPolicy.SINGLE_FLIGHT): + try: + self._execute(path, row, binding) + except (ValueError, KeyError, subprocess.TimeoutExpired, EffectRuntimeRemoteError) as exc: + row["error"] = str(exc)[:180] if isinstance(exc, (ValueError, EffectRuntimeRemoteError)) else type(exc).__name__ + self._fenced_write(path, row) + if row["status"] == "prepared": + self._observe(path, row, "rejected") + except DelegationStopRequested as stop: + # SIGTERM, a checkpoint or a fenced write: the host child is already + # gone (its run-once exits with this process's exception), the + # bootstrap was cleared, and only the acknowledgement remains. + self._acknowledge_stop(path, row, binding, source=stop.source) def _execution_arguments(self, binding: dict, operation_id: str) -> list[str]: """Exactly the same profile, workspace and validation arguments for preview/run.""" @@ -848,7 +1368,7 @@ def _record_turn_result( if publish: self._observe(path, row, "turn_returned") else: - _write(path, row) + self._fenced_write(path, row) def _receiver_adopted(self, row: dict, binding: dict) -> bool: request_id = row["identity"]["request_id"] @@ -950,7 +1470,7 @@ def _acquire_delegation_lease( goal_id=self.goal_id, ): row["task_lease"] = {"required": False, "handoff_mode": "legacy"} - _write(path, row) + self._fenced_write(path, row) return row["task_lease"] handoff_mode = show_goal_handoff_mode( registry_path=self.registry, @@ -962,7 +1482,7 @@ def _acquire_delegation_lease( "required": False, "handoff_mode": handoff_mode, } - _write(path, row) + self._fenced_write(path, row) return row["task_lease"] lease_key = self._turn_instance_id(row) result = self._cli( @@ -1004,7 +1524,7 @@ def _acquire_delegation_lease( "idempotency_key": lease_key, "version": lease["version"], } - _write(path, row) + self._fenced_write(path, row) return row["task_lease"] def _complete_delegated_todo(self, row: dict, binding: dict) -> None: @@ -1049,6 +1569,7 @@ def _execute(self, path: Path, row: dict, binding: dict) -> None: common = ["--goal-id", self.goal_id, "--agent-id", binding["agent_id"]] execution = self._execution_arguments(binding, row["identity"]["operation_id"]) try: + self._raise_if_stop_requested(path) if row["status"] == "prepared": acceptance = delegation_validation.capture(self, binding) if acceptance["plan"]["state"] != "ready" or not acceptance["files_current"]: @@ -1058,6 +1579,7 @@ def _execute(self, path: Path, row: dict, binding: dict) -> None: self._acquire_delegation_lease(path, row, binding) self._observe(path, row, "running") if row["status"] == "running": + self._raise_if_stop_requested(path) turn_key = self._matching_turn_key(row, binding) selector = ( ["--resume-turn-key", turn_key] @@ -1070,7 +1592,8 @@ def _execute(self, path: Path, row: dict, binding: dict) -> None: ] ) result = self._cli(binding, "turn", "run-once", *common, *selector, *execution, - "--execute", timeout=binding["timeout_seconds"] + 60) + "--execute", timeout=binding["timeout_seconds"] + 60, + host_record=self._host_process_record(path)) self._record_turn_result(path, row, result) finally: # The compatibility bootstrap is private host input. Keeping it @@ -1100,8 +1623,10 @@ def _execute(self, path: Path, row: dict, binding: dict) -> None: ) if not isinstance(row.get("task_lease"), dict): self._acquire_delegation_lease(path, row, binding) + self._raise_if_stop_requested(path) self._complete_delegated_todo(row, binding) todo_completed_for_settlement = True + self._raise_if_stop_requested(path) result = self._cli( binding, "turn", @@ -1112,6 +1637,7 @@ def _execute(self, path: Path, row: dict, binding: dict) -> None: *execution, "--execute", timeout=binding["timeout_seconds"] + 60, + host_record=self._host_process_record(path), ) self._record_turn_result(path, row, result, publish=False) if result.get("status") != "committed" or result.get("result_kind") != "validated_progress": @@ -1128,34 +1654,43 @@ def _execute(self, path: Path, row: dict, binding: dict) -> None: return self._bound(row, require_active=True) # revocation or rebinding while the model ran delegation_results.require_dependencies(self, binding, delegation_results.operation_brief(self, row)) - if not todo_completed_for_settlement: - self._complete_delegated_todo(row, binding) - row["artifacts"] = self._accepted(binding) - if not (_root(self.root) / "replies" / request_id / "conclusion.json").exists(): - return_result( - self.root, - self.goal_id, - binding["agent_id"], - request_id, - json.dumps( - { - "todo_id": binding["todo_id"], - "status": "accepted", - "artifacts": [ - {k: v for k, v in item.items() if k != "text"} - for item in row["artifacts"] - ], - } - ), - registry=self.registry, - caller_goal_ref=self._caller_goal_ref(), - ) - self._observe(path, row, "accepted", canonical_done=True, acceptance_ready=True, artifacts_current=True) + # Both effects commit inside the dispatch lock that a stop also + # takes, so the two sides linearize: either a stop is written first + # and neither effect runs, or both effects commit first and the stop + # that follows reports a record that already reached its terminal + # observation. Committing them outside the lock let a stop settle + # for a member whose Todo and reply had already landed. + with exclusive_file_lock(self._dispatch_lock(path)): + self._raise_if_stop_requested(path) + if not todo_completed_for_settlement: + self._complete_delegated_todo(row, binding) + row["artifacts"] = self._accepted(binding) + if not (_root(self.root) / "replies" / request_id / "conclusion.json").exists(): + return_result( + self.root, + self.goal_id, + binding["agent_id"], + request_id, + json.dumps( + { + "todo_id": binding["todo_id"], + "status": "accepted", + "artifacts": [ + {k: v for k, v in item.items() if k != "text"} + for item in row["artifacts"] + ], + } + ), + registry=self.registry, + caller_goal_ref=self._caller_goal_ref(), + ) + self._observe(path, row, "accepted", already_locked=True, + canonical_done=True, acceptance_ready=True, artifacts_current=True) except (ValueError, KeyError, subprocess.TimeoutExpired, EffectRuntimeRemoteError) as exc: # Retain uncertain execution for explicit same-operation recovery. # No fresh Turn is ever created because its client timed out. row["error"] = str(exc)[:180] if isinstance(exc, (ValueError, EffectRuntimeRemoteError)) else type(exc).__name__ - _write(path, row) + self._fenced_write(path, row) def register_delegation_tools(server, delegations: Delegations) -> None: @@ -1222,6 +1757,20 @@ def resume_delegation(operation_id: str) -> dict: """Reconnect an interrupted original execution; never launch a replacement Turn.""" return delegations.resume(operation_id) + @server.tool() + async def stop_delegation(operation_id: str) -> dict: + """Stop one original operation and return what was proven, not what was hoped. + + settled: the worker acknowledged, released the operation and let go of its + Turn lane, and the native host and its process group exited. + acknowledged/requested: still winding down; call again. unknown: + the named worker vanished before acknowledging; inspect its Turn and task + lease before reusing the task. noop: already accepted/rejected/stopped. + Stopped work is not resumed; a new scope needs a new operation id. Elapsed + time is never a receipt. + """ + return await asyncio.to_thread(delegations.stop, operation_id, execute=True) + def main(): parser = argparse.ArgumentParser(description=__doc__) @@ -1243,10 +1792,14 @@ def main(): if args.delegation_action == "validate": service._validate(service._bound(_read(service.path(args.operation_id)))) else: + service._stop_signal = install_worker_stop_signal( + service._stop_path(service.path(args.operation_id))) try: service.execute(args.operation_id) except LockAcquireTimeoutError: pass # Another worker still owns the operation after the bounded wait. + except DelegationStopRequested: + pass # Stopped before owning the operation; the holder acknowledges. return if args.workspace is None: parser.error("--workspace is required when serving MCP") diff --git a/loopx/control_plane/collaboration/delegation.ts b/loopx/control_plane/collaboration/delegation.ts index a09f5afcd..134604526 100644 --- a/loopx/control_plane/collaboration/delegation.ts +++ b/loopx/control_plane/collaboration/delegation.ts @@ -81,7 +81,7 @@ export function selectDelegationBinding(params: JsonObject): JsonObject { return binding; } -type Observation = "prepared" | "running" | "turn_returned" | "accepted" | "rejected"; +type Observation = "prepared" | "running" | "turn_returned" | "accepted" | "rejected" | "stopped"; function boundedReason(value: unknown, fallback: string): string { if (typeof value !== "string") return fallback; @@ -278,8 +278,8 @@ export function delegationPreflight(params: JsonObject): JsonObject { }; } const transitions: Record = { - prepared: ["running", "rejected"], running: ["turn_returned", "rejected"], - turn_returned: ["accepted", "rejected"], accepted: [], rejected: [], + prepared: ["running", "rejected", "stopped"], running: ["turn_returned", "rejected", "stopped"], + turn_returned: ["accepted", "rejected", "stopped"], accepted: [], rejected: [], stopped: [], }; /** Page only the caller's existing journal. A cursor is not a fleet snapshot. */ @@ -342,6 +342,70 @@ export function transitionDelegationObservation(params: JsonObject): JsonObject return {status: to}; } +type StopPhase = "requested" | "acknowledged" | "settled" | "unknown"; +const openStopPhases: readonly StopPhase[] = ["requested", "acknowledged"]; +/** What the host read back about the native Host process the operation launched. */ +type HostProcessDrain = "not_launched" | "drained" | "draining" | "unattributable"; +const hostProcessDrains: readonly HostProcessDrain[] = ["not_launched", "drained", "draining", "unattributable"]; + +/** Advance one stop request from host release facts; a receipt is never inferred from time. + * + * ``settled`` needs the acknowledgement of a process that held the operation + * lock, that lock free again, the member's Turn lane released by the stopped + * worker's process group, and the native Host the operation launched drained + * together with its process group. A worker and its lane can let go while the + * Host supervisor is still terminating the Host, so their release proves + * nothing about the Host. The host reads the lane from its holder record and + * never takes it, so a legitimate Turn is not refused, and a holder it cannot + * attribute is not released. A Host drain that cannot be attributed keeps an + * acknowledged stop open so a later read with the same identity can still + * settle it. Everything released without an acknowledgement means the named + * holder vanished before recording what it observed, which is ``unknown`` + * rather than a fake settlement. A grace timeout on its own moves nothing: a + * worker still holding a lock still runs. + */ +export function decideDelegationStop(params: JsonObject): JsonObject { + const phase = params.phase as StopPhase; + requireThat(openStopPhases.includes(phase), "delegation stop decision requires an open stop phase"); + requireThat(typeof params.acknowledged === "boolean", "delegation stop acknowledgement fact required"); + requireThat(typeof params.operation_lock_free === "boolean" && typeof params.worker_lane_released === "boolean", + "delegation stop release facts required"); + requireThat(hostProcessDrains.includes(params.host_process as HostProcessDrain), + "delegation stop host process drain fact required"); + // A required hard lease is released by the stop itself. Its release is part + // of what the receipt promises: a member whose lease is still held can block + // its Todo until the lease TTL, which is not a safe stop and is not something + // the owner can act on. `undefined` means the operation held no required + // lease. + requireThat(params.lease_released === undefined || typeof params.lease_released === "boolean", + "delegation stop lease release fact must be boolean"); + requireThat(params.timed_out === undefined || typeof params.timed_out === "boolean", + "delegation stop timeout fact must be boolean"); + requireThat(phase !== "acknowledged" || params.acknowledged === true, + "an acknowledged stop cannot lose its acknowledgement"); + const operationFree = params.operation_lock_free === true; + const workerReleased = operationFree && params.worker_lane_released === true; + const host = params.host_process as HostProcessDrain; + const hostDrained = host === "drained" || host === "not_launched"; + const leaseReleased = params.lease_released !== false; + const pending = !operationFree ? "operation_lock_still_held" + : params.worker_lane_released !== true ? "worker_lane_release_unproven" + : host === "draining" ? "host_process_still_running" + : host !== "drained" && host !== "not_launched" ? "host_process_drain_unproven" + : "required_lease_release_unproven"; + if (params.acknowledged === true) { + if (workerReleased && hostDrained && leaseReleased) { + return {phase: "settled", terminal: true, reason: "acknowledged_worker_and_host_released"}; + } + return {phase: "acknowledged", terminal: false, reason: pending}; + } + if (workerReleased && host !== "draining" && leaseReleased) { + return {phase: "unknown", terminal: true, reason: "holder_gone_without_acknowledgement"}; + } + return {phase: "requested", terminal: false, reason: operationFree ? pending + : params.timed_out === true ? "holder_still_running_after_grace" : "awaiting_acknowledgement"}; +} + /** Repair only a false terminal observation after the exact Turn validated. * * This does not retry model work. The host boundary must prove that the diff --git a/loopx/control_plane/collaboration/delegation_context.py b/loopx/control_plane/collaboration/delegation_context.py index 881ac31c9..6b799f643 100644 --- a/loopx/control_plane/collaboration/delegation_context.py +++ b/loopx/control_plane/collaboration/delegation_context.py @@ -152,6 +152,7 @@ def project_delegation_context( "turn_returned", "accepted", "rejected", + "stopped", "unavailable", ) if statuses[key] diff --git a/loopx/control_plane/collaboration/delegation_inventory.py b/loopx/control_plane/collaboration/delegation_inventory.py index 66bd66a9f..4fb2fc3a8 100644 --- a/loopx/control_plane/collaboration/delegation_inventory.py +++ b/loopx/control_plane/collaboration/delegation_inventory.py @@ -12,6 +12,12 @@ if TYPE_CHECKING: from ...collaboration_mcp import Delegations +# Files kept beside an execution record and never merged into it: the stop +# receipt, and the record naming the native Host its Turn launched. +DELEGATION_STOP_RECEIPT_SUFFIX = ".stop.json" +DELEGATION_HOST_PROCESS_SUFFIX = ".host.json" +DELEGATION_RECORD_SIDECAR_SUFFIXES = (DELEGATION_STOP_RECEIPT_SUFFIX, DELEGATION_HOST_PROCESS_SUFFIX) + def read_delegation_inventory(service: Delegations, *, limit: int = 20, cursor: str | None = None) -> dict: @@ -24,7 +30,7 @@ def addresses(): try: entries = directory.iterdir() for path in entries: - if path.suffix != ".json": + if path.suffix != ".json" or path.name.endswith(DELEGATION_RECORD_SIDECAR_SUFFIXES): continue if not BARE_SHA256_PATTERN.fullmatch(path.stem): raise ValueError("unexpected delegation record address; reconcile inventory storage") diff --git a/loopx/control_plane/effect_runtime_handlers.ts b/loopx/control_plane/effect_runtime_handlers.ts index cf04abcb7..7f46d2169 100644 --- a/loopx/control_plane/effect_runtime_handlers.ts +++ b/loopx/control_plane/effect_runtime_handlers.ts @@ -17,7 +17,7 @@ import {evaluateUserCompletion} from "./todos/user_completion.ts"; import {projectTodoSuccession} from "./todos/succession.ts"; import {projectLegacyTodoWorkCounts} from "./todos/summary_lanes.ts"; import {sealProjectionEnvelope} from "./projection_envelope.ts"; -import {recordDelegationAdoption, delegationInventoryItem, delegationInventoryQuery, delegationPreflight, delegationTurnPlanDecision, delegationValidationPlan, recoverValidatedDelegationSettlement, selectDelegationBinding, transitionDelegationObservation} from "./collaboration/delegation.ts"; +import {recordDelegationAdoption, decideDelegationStop, delegationInventoryItem, delegationInventoryQuery, delegationPreflight, delegationTurnPlanDecision, delegationValidationPlan, recoverValidatedDelegationSettlement, selectDelegationBinding, transitionDelegationObservation} from "./collaboration/delegation.ts"; import {resolveConversationTrigger} from "./collaboration/conversation_trigger.ts"; import {admitGoalDraft} from "./collaboration/goal_draft.ts"; import {planChatMode} from "./collaboration/chat_mode.ts"; @@ -754,6 +754,7 @@ export function createEffectRuntimeHandlers( ["chat.turn.accept", planChatTurnAcceptance], ["collaboration.delegation.observe", transitionDelegationObservation], ["collaboration.delegation.recover_validated_settlement", recoverValidatedDelegationSettlement], + ["collaboration.delegation.stop", decideDelegationStop], ["collaboration.delegation.adoption", recordDelegationAdoption], [ "collaboration.request.normalize", diff --git a/loopx/control_plane/subagent_context.ts b/loopx/control_plane/subagent_context.ts index 56f86e81f..c21fadd14 100644 --- a/loopx/control_plane/subagent_context.ts +++ b/loopx/control_plane/subagent_context.ts @@ -174,7 +174,7 @@ function boundedDelegationContext(value: unknown): JsonObject | null { const operationReceipts: JsonObject = {}; if (rawReceipts) { for (const key of ["observed", "prepared", "running", "turn_returned", "accepted", - "rejected", "unavailable", "recovery_required"]) { + "rejected", "stopped", "unavailable", "recovery_required"]) { if (Number.isInteger(rawReceipts[key]) && Number(rawReceipts[key]) >= 0) { operationReceipts[key] = Math.min(Number(rawReceipts[key]), 10_000); } diff --git a/loopx/control_plane/turn_driver/host_process.ts b/loopx/control_plane/turn_driver/host_process.ts index 685a5c39f..0a752a32b 100644 --- a/loopx/control_plane/turn_driver/host_process.ts +++ b/loopx/control_plane/turn_driver/host_process.ts @@ -22,6 +22,9 @@ export interface HostProcessResult { group_signal_sent: boolean; } export type HostProcessOutput = {kind: "stdout" | "stderr"; text: string}; +/** The group this owner will clean up, reported once the Host is spawned. + * ``process_group`` is null where cleanup is tree best effort (Windows). */ +export type HostProcessSpawned = {kind: "spawned"; pid: number; process_group: number | null}; export const HOST_PROCESS_TERMINATE_GRACE_MS = 300; /** Restrict transport size separately from the caller's public result budget. */ @@ -47,7 +50,8 @@ function signalGroup(child: ChildProcessWithoutNullStreams, signal: NodeJS.Signa } export async function runHostProcess(request: HostProcessRequest, - output: (item: HostProcessOutput) => Promise, signal?: AbortSignal): Promise { + output: (item: HostProcessOutput) => Promise, signal?: AbortSignal, + spawned?: (item: HostProcessSpawned) => Promise): 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}; @@ -116,6 +120,12 @@ export async function runHostProcess(request: HostProcessRequest, }; const reads = Promise.all([read("stdout"), read("stderr")]); child.stdin.on("error", () => {}); // A Host may close stdin before consuming it. + if (child.pid && spawned) { + // A caller that cannot record the owned group cancels rather than run unaccounted. + try { await spawned({kind: "spawned", pid: child.pid, + process_group: process.platform === "win32" ? null : child.pid}); } + catch { complete = false; stop("cancelled"); } + } child.stdin.end(request.input); try { await exited; diff --git a/loopx/control_plane/turn_driver/host_process_bridge.ts b/loopx/control_plane/turn_driver/host_process_bridge.ts index b1c261c8b..d34038cfa 100644 --- a/loopx/control_plane/turn_driver/host_process_bridge.ts +++ b/loopx/control_plane/turn_driver/host_process_bridge.ts @@ -20,7 +20,7 @@ process.stdin.on("data", (chunk: Buffer) => { 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)); } + try { await emit(await runHostProcess(decodeHostProcessRequest(JSON.parse(line)), emit, owner.signal, emit)); } 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 index 4f51148fa..95405308c 100644 --- a/loopx/control_plane/turn_driver/host_process_transport.py +++ b/loopx/control_plane/turn_driver/host_process_transport.py @@ -3,12 +3,15 @@ from __future__ import annotations import json +import os import subprocess import sys +import tempfile from collections.abc import Callable, Sequence from pathlib import Path from typing import Any +from ...file_lock import lock_holder_host_label from ..effect_runtime import _node_executable # Keep Python's Windows executable/batch launcher compatibility. No timeout, @@ -16,6 +19,89 @@ _WINDOWS_COMMAND_RELAY = "import subprocess,sys;sys.exit(subprocess.call(sys.argv[1:]))" +# A launching owner that must later prove its Host drained names a record path +# here. The transport consumes it: the Host never inherits it, so a nested +# LoopX run inside the Host cannot overwrite its parent's record. +HOST_PROCESS_RECORD_ENV = "LOOPX_HOST_PROCESS_RECORD" +HOST_PROCESS_RECORD_SCHEMA_VERSION = "loopx_host_process_record_v0" +# Drain facts read back from a record; the caller's typed decision interprets them. +HOST_PROCESS_NOT_LAUNCHED = "not_launched" +HOST_PROCESS_DRAINED = "drained" +HOST_PROCESS_DRAINING = "draining" +HOST_PROCESS_UNATTRIBUTABLE = "unattributable" +HOST_PROCESS_UNSUPPORTED_PLATFORM = "unsupported_platform" + + +def host_process_drain_supported() -> bool: + """Whether this platform can prove an owned Host group has exited. + + The TS supervisor cleans up a process group on POSIX and a process tree + best-effort on Windows. Only the group gives the stop a fact it can prove, + so the caller must say so rather than reporting a settlement it cannot + support. + """ + + return hasattr(os, "killpg") + + +def _write_host_process_record(path: Path, record: dict[str, Any]) -> None: + path.parent.mkdir(parents=True, exist_ok=True) + descriptor, temporary = tempfile.mkstemp(prefix=f".{path.name}.", suffix=".tmp", dir=path.parent) + try: + with os.fdopen(descriptor, "w", encoding="utf-8") as handle: + json.dump(record, handle, separators=(",", ":")) + os.replace(temporary, path) + finally: + Path(temporary).unlink(missing_ok=True) + + +def _process_group_present(pgid: int) -> bool: + try: + os.killpg(pgid, 0) + except ProcessLookupError: + return False + except PermissionError: + return True # it exists; it is simply not ours to signal + return True + + +def host_process_drain(record_path: Path) -> str: + """Read whether the Host a record names, and its process group, have exited. + + Read-only: nothing is signalled, and the TS-owned supervisor keeps cleanup. + No record means no Host was launched under it. ``draining`` while the + supervising bridge or the Host's group still has a member. A record from + another machine, one without a reported group whose supervisor is gone, + or a platform without process groups proves nothing: ``unattributable``. + """ + + try: + record = json.loads(record_path.read_text(encoding="utf-8")) + except FileNotFoundError: + return HOST_PROCESS_NOT_LAUNCHED + except (OSError, ValueError): + return HOST_PROCESS_UNATTRIBUTABLE + if (not isinstance(record, dict) or record.get("schema_version") != HOST_PROCESS_RECORD_SCHEMA_VERSION + or record.get("host") != lock_holder_host_label()): + return HOST_PROCESS_UNATTRIBUTABLE + if not hasattr(os, "killpg"): + # A launched Host on a platform without process groups is never proven + # drained. This is a platform boundary, not an attribution failure. + return HOST_PROCESS_UNSUPPORTED_PLATFORM + bridge, group = record.get("bridge_pid"), record.get("process_group") + if record.get("phase") != "finished": + if not isinstance(bridge, int) or bridge <= 1: + return HOST_PROCESS_UNATTRIBUTABLE + if _process_group_present(bridge): # the bridge leads its own session + return HOST_PROCESS_DRAINING + if group is None: + # The supervisor left before reporting a group it may already have spawned. + return HOST_PROCESS_UNATTRIBUTABLE if record.get("phase") == "launching" else HOST_PROCESS_DRAINED + if not isinstance(group, int) or group <= 1: + return HOST_PROCESS_UNATTRIBUTABLE + return HOST_PROCESS_DRAINING if _process_group_present(group) else HOST_PROCESS_DRAINED + + class HostOutputLines: """Frame LF records without retaining raw trajectories or an unbounded line.""" @@ -76,6 +162,9 @@ def run_host_process( "stdout_limit_bytes": stdout_limit_bytes, } bridge = Path(__file__).with_name("host_process_bridge.ts") + environment = os.environ.copy() + record_value = environment.pop(HOST_PROCESS_RECORD_ENV, "") + record_path = Path(record_value) if record_value else None with subprocess.Popen( [ _node_executable(), @@ -90,10 +179,19 @@ def run_host_process( encoding="utf-8", errors="strict", start_new_session=True, + env=environment, ) as proc: assert proc.stdin is not None and proc.stdout is not None result = None + record = None try: + if record_path is not None: + # Written before the request, so no Host exists that the record + # does not name; a record that cannot be written launches nothing. + record = {"schema_version": HOST_PROCESS_RECORD_SCHEMA_VERSION, + "host": lock_holder_host_label(), "owner_pid": os.getpid(), + "bridge_pid": proc.pid, "phase": "launching", "process_group": None} + _write_host_process_record(record_path, record) proc.stdin.write( json.dumps(request, ensure_ascii=False, separators=(",", ":")) + "\n" ) @@ -101,7 +199,13 @@ def run_host_process( for line in proc.stdout: event = json.loads(line) kind = event.get("kind") - if kind in {"stdout", "stderr"} and isinstance(event.get("text"), str): + if kind == "spawned" and isinstance(event.get("pid"), int): + if record is not None: + group = event.get("process_group") + record.update(phase="spawned", host_pid=event["pid"], + process_group=group if isinstance(group, int) else None) + _write_host_process_record(record_path, record) + elif 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"]) @@ -124,6 +228,10 @@ def run_host_process( except subprocess.TimeoutExpired: proc.kill() proc.wait() + if record is not None and result is not None and proc.returncode == 0: + # The supervisor returned only after cleaning its group; the group is still re-read. + record["phase"] = "finished" + _write_host_process_record(record_path, record) 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 index 7da98714e..558725565 100644 --- a/tests/control_plane/test_host_process.py +++ b/tests/control_plane/test_host_process.py @@ -2,6 +2,7 @@ from __future__ import annotations +import json import os import signal import subprocess @@ -174,3 +175,65 @@ def test_windows_transport_relay_preserves_argv_and_stdin(tmp_path: Path) -> Non ) assert result["ok"] is True assert result["value"] == {"args": values, "input": {}} + + +@pytest.mark.skipif(os.name == "nt", reason="POSIX process-group drain readback") +def test_host_process_record_names_the_owned_group_and_is_not_inherited(tmp_path: Path, monkeypatch) -> None: + from loopx.control_plane.turn_driver.host_process_transport import ( + HOST_PROCESS_RECORD_ENV, host_process_drain, + ) + + record_path = tmp_path / "op.host.json" + assert host_process_drain(record_path) == "not_launched" + monkeypatch.setenv(HOST_PROCESS_RECORD_ENV, str(record_path)) + host = ("import json,os,sys;print(json.dumps({'env': os.environ.get(%r), 'pid': os.getpid()," + " 'pgid': os.getpgid(0)}))" % HOST_PROCESS_RECORD_ENV) + result = _run_host({}, argv=[sys.executable, "-c", host], project=tmp_path, timeout_seconds=5) + assert result["ok"] is True + # A nested LoopX run inside the Host cannot overwrite its parent's record. + assert result["value"]["env"] is None + record = json.loads(record_path.read_text()) + assert record["phase"] == "finished" + assert record["host_pid"] == result["value"]["pid"] == record["process_group"] == result["value"]["pgid"] + assert host_process_drain(record_path) == "drained" + + +@pytest.mark.skipif(os.name == "nt", reason="POSIX process-group drain readback") +def test_host_process_drain_reads_live_groups_and_refuses_unattributable_records(tmp_path: Path) -> None: + from loopx.control_plane.turn_driver.host_process_transport import ( + HOST_PROCESS_RECORD_SCHEMA_VERSION, host_process_drain, + ) + from loopx.file_lock import lock_holder_host_label + + path = tmp_path / "op.host.json" + live = subprocess.Popen([sys.executable, "-c", "import time;time.sleep(60)"], start_new_session=True) + gone = subprocess.Popen([sys.executable, "-c", "pass"], start_new_session=True) + gone.wait(timeout=10) + + def drain(**fields): + path.write_text(json.dumps({"schema_version": HOST_PROCESS_RECORD_SCHEMA_VERSION, + "host": lock_holder_host_label(), "owner_pid": 1, **fields})) + return host_process_drain(path) + + try: + # A live supervisor or a live Host group is still draining. + assert drain(phase="launching", bridge_pid=live.pid, process_group=None) == "draining" + assert drain(phase="spawned", bridge_pid=gone.pid, process_group=live.pid) == "draining" + assert drain(phase="finished", bridge_pid=gone.pid, process_group=live.pid) == "draining" + assert drain(phase="spawned", bridge_pid=gone.pid, process_group=gone.pid) == "drained" + # A supervisor gone before it reported a group may have spawned one anyway. + assert drain(phase="launching", bridge_pid=gone.pid, process_group=None) == "unattributable" + assert drain(phase="finished", bridge_pid=gone.pid, process_group=None) == "drained" + for fields in ({"phase": "spawned", "bridge_pid": None, "process_group": gone.pid}, + {"phase": "spawned", "bridge_pid": gone.pid, "process_group": "1"}, + {"phase": "spawned", "bridge_pid": gone.pid, "process_group": 1}, + {"phase": "spawned", "bridge_pid": gone.pid, "process_group": gone.pid, + "host": "another-machine"}, + {"phase": "spawned", "bridge_pid": gone.pid, "process_group": gone.pid, + "schema_version": "other"}): + assert drain(**fields) == "unattributable", fields + path.write_text("{not json") + assert host_process_drain(path) == "unattributable" + finally: + live.kill() + live.wait(timeout=10) diff --git a/tests/control_plane_ts/delegation.test.ts b/tests/control_plane_ts/delegation.test.ts index 7ae9fe4d4..953358667 100644 --- a/tests/control_plane_ts/delegation.test.ts +++ b/tests/control_plane_ts/delegation.test.ts @@ -1,6 +1,6 @@ import test from "node:test"; import assert from "node:assert/strict"; -import {recordDelegationAdoption, delegationInventoryItem, delegationInventoryQuery, delegationPreflight, delegationTurnPlanDecision, delegationValidationPlan, recoverValidatedDelegationSettlement, selectDelegationBinding, transitionDelegationObservation} from "../../loopx/control_plane/collaboration/delegation.ts"; +import {recordDelegationAdoption, decideDelegationStop, delegationInventoryItem, delegationInventoryQuery, delegationPreflight, delegationTurnPlanDecision, delegationValidationPlan, recoverValidatedDelegationSettlement, selectDelegationBinding, transitionDelegationObservation} from "../../loopx/control_plane/collaboration/delegation.ts"; import {canonicalAuthoritySha256} from "../../loopx/control_plane/coordination/authority_store_codec.ts"; import {projectTurnSelectionRejection} from "../../loopx/control_plane/turn_driver/selection_rejection.ts"; @@ -102,6 +102,92 @@ test("message receipt and model return do not imply accepted work", () => { canonical_done: true, acceptance_ready: true, artifacts_current: true}), {status: "accepted"}); }); +test("stopped is terminal and reachable only from open observations", () => { + for (const from of ["prepared", "running", "turn_returned"]) + assert.deepEqual(transitionDelegationObservation({from, to: "stopped"}), {status: "stopped"}); + assert.deepEqual(transitionDelegationObservation({from: "stopped", to: "stopped"}), {status: "stopped"}); + for (const from of ["accepted", "rejected"]) + assert.throws(() => transitionDelegationObservation({from, to: "stopped"}), /transition/); + for (const to of ["running", "turn_returned", "accepted", "rejected"]) + assert.throws(() => transitionDelegationObservation({from: "stopped", to}), /transition/); + const observation = {operation_id: "op-1", request_id: "req", agent_id: "reviewer", todo_id: "todo_review", + status: "stopped", worker_active: false, recovery_required: false}; + assert.equal(delegationInventoryItem({record: {record_id: "a".repeat(64), operation_id: "op-1"}, + observation}).status, "stopped"); +}); + +test("a stop settles only on an acknowledgement plus released holders; time alone proves nothing", () => { + const open = {phase: "requested", acknowledged: false, operation_lock_free: false, worker_lane_released: false, + host_process: "drained"}; + assert.deepEqual(decideDelegationStop(open), {phase: "requested", terminal: false, reason: "awaiting_acknowledgement"}); + assert.deepEqual(decideDelegationStop({...open, timed_out: true}), + {phase: "requested", terminal: false, reason: "holder_still_running_after_grace"}); + // A lane release without a free operation lock is not a vanished holder. + assert.deepEqual(decideDelegationStop({...open, worker_lane_released: true, timed_out: true}), + {phase: "requested", terminal: false, reason: "holder_still_running_after_grace"}); + // A free operation lock with an unattributed lane holder proves nothing yet. + assert.deepEqual(decideDelegationStop({...open, operation_lock_free: true, timed_out: true}), + {phase: "requested", terminal: false, reason: "worker_lane_release_unproven"}); + assert.deepEqual(decideDelegationStop({...open, operation_lock_free: true, worker_lane_released: true}), + {phase: "unknown", terminal: true, reason: "holder_gone_without_acknowledgement"}); + const acked = {phase: "acknowledged", acknowledged: true, operation_lock_free: false, worker_lane_released: false, + host_process: "drained"}; + assert.deepEqual(decideDelegationStop(acked), {phase: "acknowledged", terminal: false, reason: "operation_lock_still_held"}); + assert.deepEqual(decideDelegationStop({...acked, worker_lane_released: true}), + {phase: "acknowledged", terminal: false, reason: "operation_lock_still_held"}); + assert.deepEqual(decideDelegationStop({...acked, operation_lock_free: true}), + {phase: "acknowledged", terminal: false, reason: "worker_lane_release_unproven"}); + assert.deepEqual(decideDelegationStop({...acked, phase: "requested", operation_lock_free: true, worker_lane_released: true}), + {phase: "settled", terminal: true, reason: "acknowledged_worker_and_host_released"}); + assert.deepEqual(decideDelegationStop({...acked, operation_lock_free: true, worker_lane_released: true, timed_out: true}), + {phase: "settled", terminal: true, reason: "acknowledged_worker_and_host_released"}); + for (const patch of [{phase: "settled"}, {phase: "unknown"}, {phase: "noop"}, {acknowledged: "yes"}, + {host_process: undefined}, {host_process: "exited"}, {host_process: true}, + {operation_lock_free: 1}, {worker_lane_released: undefined}, {lane_lock_free: true, worker_lane_released: undefined}, + {timed_out: "later"}, {phase: "acknowledged", acknowledged: false}]) + assert.throws(() => decideDelegationStop({...open, ...patch})); +}); + +test("a released worker and lane never settle a stop while the native Host still drains", () => { + const released = {phase: "acknowledged", acknowledged: true, operation_lock_free: true, worker_lane_released: true}; + assert.deepEqual(decideDelegationStop({...released, host_process: "draining", timed_out: true}), + {phase: "acknowledged", terminal: false, reason: "host_process_still_running"}); + // Without an attributable drain the stop stays open for a later same-identity read. + assert.deepEqual(decideDelegationStop({...released, host_process: "unattributable"}), + {phase: "acknowledged", terminal: false, reason: "host_process_drain_unproven"}); + for (const host_process of ["drained", "not_launched"]) + assert.deepEqual(decideDelegationStop({...released, host_process}), + {phase: "settled", terminal: true, reason: "acknowledged_worker_and_host_released"}); + // A required lease the stop could not release is not a settlement: the + // member's Todo can stay blocked by it until the lease TTL. + for (const host_process of ["drained", "not_launched"]) { + assert.deepEqual(decideDelegationStop({...released, host_process, lease_released: false}), + {phase: "acknowledged", terminal: false, reason: "required_lease_release_unproven"}); + // An operation that held no required lease omits the fact, which settles. + assert.deepEqual(decideDelegationStop({...released, host_process}), + {phase: "settled", terminal: true, reason: "acknowledged_worker_and_host_released"}); + } + assert.throws(() => decideDelegationStop({...released, host_process: "drained", lease_released: "yes"}), + /lease release fact/); + // A vanished holder whose required lease is still held is not terminal either. + const unacked = {phase: "requested", acknowledged: false, operation_lock_free: true, + worker_lane_released: true, host_process: "drained"}; + assert.deepEqual(decideDelegationStop({...unacked, lease_released: false}), + {phase: "requested", terminal: false, reason: "required_lease_release_unproven"}); + assert.deepEqual(decideDelegationStop(unacked), + {phase: "unknown", terminal: true, reason: "holder_gone_without_acknowledgement"}); + // A held lock still dominates a drained Host. + assert.deepEqual(decideDelegationStop({...released, operation_lock_free: false, host_process: "drained"}), + {phase: "acknowledged", terminal: false, reason: "operation_lock_still_held"}); + // A vanished holder is unknown only once its Host is no longer seen running. + const vanished = {...released, phase: "requested", acknowledged: false}; + assert.deepEqual(decideDelegationStop({...vanished, host_process: "draining"}), + {phase: "requested", terminal: false, reason: "host_process_still_running"}); + for (const host_process of ["drained", "not_launched", "unattributable"]) + assert.deepEqual(decideDelegationStop({...vanished, host_process}), + {phase: "unknown", terminal: true, reason: "holder_gone_without_acknowledgement"}); +}); + test("a false rejection can reopen only for exact validated settlement recovery", () => { const evidence = { from: "rejected", diff --git a/tests/control_plane_ts/host_process.test.ts b/tests/control_plane_ts/host_process.test.ts index 23b235a8e..7eaaaf541 100644 --- a/tests/control_plane_ts/host_process.test.ts +++ b/tests/control_plane_ts/host_process.test.ts @@ -1,4 +1,5 @@ import assert from "node:assert/strict"; +import {existsSync} from "node:fs"; import {mkdtemp, readFile, rm} from "node:fs/promises"; import {join} from "node:path"; import {tmpdir} from "node:os"; @@ -78,3 +79,26 @@ test("output consumer failure cancels execution rather than leaving an orphan", async () => { throw new Error("consumer left"); }); assert.equal(result.outcome, "cancelled"); assert.equal(result.output_complete, false); }); + +test("the spawned Host group is reported once before input, and an unrecorded group never runs", {skip: process.platform === "win32"}, async t => { + const root = await mkdtemp(join(tmpdir(), "loopx-host-spawned-")); + t.after(() => rm(root, {recursive: true, force: true})); + const seen: unknown[] = []; + let stdout = ""; + const result = await runHostProcess(request(`process.stdout.write(String(process.pid)+' '+String(require('child_process').execSync('ps -o pgid= -p '+process.pid)).trim())`), + async item => { stdout += item.text; }, undefined, async item => { seen.push(item); }); + const [pid, pgid] = stdout.split(" ").map(Number); + assert.equal(result.outcome, "exited"); + assert.deepEqual(seen, [{kind: "spawned", pid, process_group: pid}]); assert.equal(pgid, pid); + // A caller that cannot record the owned group runs nothing unaccounted for, and + // the armed host proves it really started, so this is not an unspawned process. + const script = (marker: string) => `require('fs').writeFileSync(${JSON.stringify(marker)},'') + process.stdin.on('data',()=>{});setInterval(()=>{},1000)`; + const recorded = join(root, "recorded"); + const refused = await runHostProcess(request(script(recorded)), async () => {}, undefined, async () => { + const until = Date.now() + 2000; + while (!existsSync(recorded) && Date.now() < until) await delay(5); + assert.ok(existsSync(recorded), "the Host never started"); + throw new Error("record unavailable"); }); + assert.equal(refused.outcome, "cancelled"); assert.equal(refused.output_complete, false); +}); diff --git a/tests/test_delegation_cli.py b/tests/test_delegation_cli.py index ee751375a..95355fbbf 100644 --- a/tests/test_delegation_cli.py +++ b/tests/test_delegation_cli.py @@ -84,6 +84,34 @@ def test_attached_cli_disconnect_retry_and_verified_return(service): assert "artifacts" not in inventory["items"][0] +def test_cli_stop_settles_a_running_member_and_refuses_resume(service): + root, runner = service + (root / "hold").touch() + source = root / "brief.json" + source.write_text(json.dumps(brief())) + status, started = cli(runner, "start", "--binding-id", "analysis", "--operation-id", "cli-stop", + "--brief-file", str(source), "--execute") + assert status == 0, started + deadline = time.monotonic() + 45 + while not (root / "host-started").exists() and time.monotonic() < deadline: + time.sleep(0.1) + assert (root / "host-started").exists() + status, refused = cli(runner, "stop", "--operation-id", "cli-stop") + assert status == 1 and "--execute" in refused["error"] + assert not runner._stop_path(runner.path("cli-stop")).exists() + status, stopped = cli(runner, "stop", "--operation-id", "cli-stop", "--execute") + assert status == 0 and stopped["phase"] == "settled" and stopped["status"] == "stopped", stopped + assert stopped["stop"]["ack"]["source"] == "SIGTERM" + status, again = cli(runner, "stop", "--operation-id", "cli-stop", "--execute") + assert status == 0 and again == stopped + status, resumed = cli(runner, "resume", "--operation-id", "cli-stop", "--execute") + assert status == 1 and "start a new operation id" in resumed["error"] + status, observed = cli(runner, "read", "--operation-id", "cli-stop") + assert status == 0 and observed["status"] == "stopped" and observed["stop"]["phase"] == "settled" + assert (root / "analyst" / "initial" / "host-invocations").read_text() == "1" + assert not demo.canonical_tasks(root)["todo_analyst-initial"]["done"] + + def test_cli_invalid_inputs_do_not_launch_work(service): root, runner = service bad = root / "bad.json" diff --git a/tests/test_delegation_inventory.py b/tests/test_delegation_inventory.py index 8486eb4a5..9ac402e4c 100644 --- a/tests/test_delegation_inventory.py +++ b/tests/test_delegation_inventory.py @@ -56,6 +56,15 @@ def test_corruption_and_stopped_worker_do_not_hide_healthy_sibling(service, monk assert by_id["stopped"]["recovery_required"] assert len(page["items"]) == 4 and not page["page_readback_complete"] assert sum(row["status"] == "unavailable" for row in page["items"]) == 2 + # A stop receipt beside its record is not another record and reads back as stopped. + assert runner.stop("healthy", execute=True)["phase"] == "settled" + assert runner._stop_path(runner.path("healthy")).exists() + # Nor is the record naming the native Host an operation's Turn launched. + runner._host_process_record(runner.path("healthy")).write_text("{}") + page = runner.operations() + by_id = {row["operation_id"]: row for row in page["items"] if row["operation_id"]} + assert len(page["items"]) == 4 and by_id["healthy"]["status"] == "stopped" + assert not by_id["healthy"]["recovery_required"] def test_unknown_requester_and_unreadable_source_are_not_empty_inventory(service, monkeypatch): diff --git a/tests/test_local_delegation.py b/tests/test_local_delegation.py index d29b13fd7..51733ca48 100644 --- a/tests/test_local_delegation.py +++ b/tests/test_local_delegation.py @@ -1,6 +1,7 @@ """Production delegation/Turn/TS completion with an explicit fixture model host.""" import json import asyncio +import os from pathlib import Path import subprocess import sys @@ -16,13 +17,16 @@ sys.path.insert(0, str(Path(__file__).resolve().parents[1] / "examples" / "managed-research-team")) import research_team as demo # noqa: E402 from test_managed_research_scenario import fixture # noqa: E402 -from loopx.collaboration_mcp import Delegations # noqa: E402 +from loopx import collaboration_mcp as delegation_module # noqa: E402 +from loopx.collaboration_mcp import DelegationFenced, DelegationStopRequested, Delegations # noqa: E402 from loopx.control_plane.collaboration.peers import returns # noqa: E402 from loopx.control_plane.collaboration.inbox import _read # noqa: E402 -from loopx.file_lock import exclusive_file_lock # noqa: E402 +from loopx.control_plane.turn_driver.lane_fence import turn_lane_liveness, turn_lane_singleflight # noqa: E402 +from loopx.file_lock import exclusive_file_lock, try_exclusive_file_lock # noqa: E402 +from loopx.file_lock import lock_holder_host_label # noqa: E402 -HOST = '''import json, sys, time +HOST = '''import json, os, sys, time from pathlib import Path from loopx.control_plane.turn_driver.host_candidate import build_result from loopx.control_plane.collaboration.inbox import acknowledge @@ -34,7 +38,13 @@ actor = envelope['agent_id'] counter = workspace / 'host-invocations' counter.write_text(str(int(counter.read_text()) + 1 if counter.exists() else 1)) +if (root / 'ignore-term').exists(): + import signal, subprocess + signal.signal(signal.SIGTERM, signal.SIG_IGN) + child = subprocess.Popen([sys.executable, '-c', 'import signal, time; signal.signal(signal.SIGTERM, signal.SIG_IGN); time.sleep(600)']) + (root / 'host-child-pid').write_text(str(child.pid)) if (root / 'hold').exists(): + (root / 'host-pid').write_text(str(os.getpid())) (root / 'host-started').touch() while not (root / 'release').exists(): time.sleep(0.1) delegation = json.loads((workspace / 'DELEGATION.json').read_text()) @@ -169,6 +179,7 @@ async def disconnect_requester(): assert not (root / "host-started").exists() inventory = await session.call_tool("list_delegations", {}) assert not inventory.isError and json.loads(inventory.content[0].text)["items"] == [] + assert "stop_delegation" in {tool.name for tool in (await session.list_tools()).tools} result = await session.call_tool("start_delegation", { "binding_id": "analysis", "operation_id": "analysis-1", "brief": brief()}) assert not result.isError @@ -197,6 +208,13 @@ async def disconnect_requester(): assert len(returned) == 1 assert returned[0]["decision"] == "adopt" assert wait(reconnected)["artifacts"] == result["artifacts"] + # Accepted work cannot be stopped: nothing is written and the receipt repeats exactly. + noop = reconnected.stop("analysis-1", execute=True) + assert noop["phase"] == "noop" and noop["status"] == "accepted" and noop["stop"] is None + assert reconnected.stop("analysis-1", execute=True) == noop + assert not reconnected._stop_path(reconnected.path("analysis-1")).exists() + assert wait(reconnected)["artifacts"] == result["artifacts"] + assert "stop" not in reconnected.read("analysis-1") changed_brief = {**brief(), "purpose": "Changed instruction"} with pytest.raises(ValueError, match="identity conflict"): reconnected.start("analysis", "analysis-1", changed_brief) @@ -265,3 +283,544 @@ def read_on_publish(path, row, status, **facts): assert len(terminal_reads) == 1 assert not demo.canonical_tasks(root)["todo_analyst-initial"]["done"] assert returns(runner.root, runner.goal_id, "lead")["items"] == [] + + +def until(predicate, timeout=45): + deadline = time.monotonic() + timeout + while time.monotonic() < deadline: + if predicate(): + return True + time.sleep(0.1) + return predicate() + + +def process_gone(pid): + """Platform-valid: a zombie has exited; a process ``ps`` cannot see is gone.""" + try: + os.kill(pid, 0) + except ProcessLookupError: + return True + except PermissionError: + return False + state = subprocess.run(["ps", "-o", "stat=", "-p", str(pid)], + capture_output=True, text=True).stdout.strip() + return not state or state.startswith("Z") + + +def start_held_worker(service, operation="analysis-stop"): + root, runner = service + (root / "hold").touch() + runner.start("analysis", operation, brief()) + assert until(lambda: (root / "host-started").exists()), _read(runner.path(operation)) + return int((root / "host-pid").read_text()) + + +def test_stop_while_executing_is_acknowledged_by_the_worker_and_settles(service, monkeypatch): + """The detached worker acknowledges SIGTERM under its own lock; time proves nothing.""" + root, runner = service + host_pid = start_held_worker(service) + path = runner.path("analysis-stop") + before = _read(path) + assert before["status"] == "running" and before["worker"]["pid"] == before["worker"]["pgid"] + # The real run-once child holds the member's lane from inside the worker's group, + # which is what lets settlement attribute the lane without ever taking it. + binding = runner.binding("analysis") + lane = turn_lane_liveness(runner._lane_target(binding)) + assert lane["state"] == "live" and lane["holder"]["pid"] != before["worker"]["pid"] + assert os.getpgid(lane["holder"]["pid"]) == before["worker"]["pgid"] + receipt = runner.stop("analysis-stop", execute=True) + host_gone = process_gone(host_pid) # the instant settled returns, not after a wait + assert receipt["phase"] == "settled" and receipt["status"] == "stopped", receipt + assert host_gone + stop = receipt["stop"] + assert stop["requested_by"] == "lead" and stop["requested_status"] == "running" + assert stop["worker"]["pid"] == before["worker"]["pid"] + assert stop["ack"]["pid"] == before["worker"]["pid"] and stop["ack"]["source"] == "SIGTERM" + assert stop["ack"]["observed_status"] == "running" and stop["ack"]["turn_key"] + assert stop["settled"]["operation_lock_free"] and stop["settled"]["worker_lane_released"] + assert stop["settled"]["lane_state"] in {"dead", "released"} + assert stop["settled"]["host_process"] == "drained" + assert stop["settled"]["turn_journal_status"] == "in_progress" + assert stop["lease"] == {"required": False, "released": None} + # The acknowledged record is final: nobody writes it again, the Todo stays open, + # the worker exits and the member's Turn lane can be taken. + frozen = path.read_bytes() + assert until(lambda: process_gone(before["worker"]["pid"]), timeout=20) + with try_exclusive_file_lock(runner._lane_target(binding)) as held: + assert held is not None + assert not demo.canonical_tasks(root)["todo_analyst-initial"]["done"] + assert not (root / "analyst" / "initial" / "DELEGATION.json").exists() + assert path.read_bytes() == frozen + assert runner.stop("analysis-stop", execute=True) == receipt + observed = runner.read("analysis-stop") + assert observed["status"] == "stopped" and not observed["recovery_required"] + assert observed["stop"] == {"stop_id": stop["stop_id"], "phase": "settled"} + assert runner.wait("analysis-stop")["status"] == "stopped" + monkeypatch.setattr(runner, "_spawn", lambda _: pytest.fail("stopped work must not respawn")) + with pytest.raises(ValueError, match="start a new operation id"): + runner.resume("analysis-stop") + assert path.read_bytes() == frozen + page = runner.operations() + assert page["page_readback_complete"] and page["items"][0]["status"] == "stopped" + assert returns(runner.root, runner.goal_id, "lead")["items"] == [] + + +def test_stop_without_a_holder_is_acknowledged_by_the_requester(service, monkeypatch): + root, runner = service + monkeypatch.setattr(runner, "_spawn", lambda _: None) + runner.start("analysis", "analysis-idle", brief()) + receipt = runner.stop("analysis-idle", execute=True) + assert receipt["phase"] == "settled" and receipt["status"] == "stopped" + assert receipt["stop"]["worker"] is None and receipt["stop"]["requested_status"] == "prepared" + assert receipt["stop"]["ack"]["pid"] == os.getpid() and receipt["stop"]["ack"]["source"] == "requester" + assert receipt["stop"]["settled"]["turn_journal_status"] is None + assert receipt["stop"]["settled"]["lane_state"] in {"absent", "released"} + assert receipt["stop"]["settled"]["host_process"] == "not_launched" + frozen = runner.path("analysis-idle").read_bytes() + with pytest.raises(ValueError, match="start a new operation id"): + runner.resume("analysis-idle") + runner.execute("analysis-idle") # a late worker finds terminal work and launches nothing + assert runner.path("analysis-idle").read_bytes() == frozen + assert not (root / "host-started").exists() + assert runner.stop("analysis-idle", execute=True) == receipt + with pytest.raises(ValueError, match="requires execute"): + runner.stop("analysis-idle", execute=False) + with pytest.raises(ValueError, match="unknown delegation operation"): + runner.stop("never-started", execute=True) + + +def test_worker_killed_before_acknowledging_is_unknown_not_settled(service, monkeypatch): + """A vanished named holder never becomes a settlement; the stop still fences resume.""" + import signal + + root, runner = service + host_pid = start_held_worker(service) + path = runner.path("analysis-stop") + worker = _read(path)["worker"] + killed = [] + + def kill_without_grace(target, stop): + if killed: + return # later calls find nothing left to signal + killed.append(stop["worker"]) + assert stop["worker"] == worker and worker["pgid"] != os.getpgid(0) + os.killpg(worker["pgid"], signal.SIGKILL) + assert until(lambda: runner._operation_lock_free(target), timeout=20) + + monkeypatch.setattr(runner, "_signal_worker", kill_without_grace) + first = runner.stop("analysis-stop", execute=True) + assert first["stop"]["ack"] is None and first["phase"] in {"requested", "unknown"}, first + # The operation lock is free now, yet the requester never acknowledges for a + # named worker: the outcome converges on unknown once its Turn child is reaped. + assert until(lambda: runner.stop("analysis-stop", execute=True)["phase"] == "unknown", timeout=20) + receipt = runner.stop("analysis-stop", execute=True) + assert receipt["status"] == "running" and receipt["stop"]["ack"] is None, receipt + assert receipt["stop"]["settled"]["operation_lock_free"] and receipt["stop"]["settled"]["worker_lane_released"] + assert receipt["stop"]["settled"]["turn_journal_status"] == "in_progress" + assert receipt["stop"]["lease"] is None and len(killed) == 1 + assert until(lambda: process_gone(host_pid), timeout=20) + assert runner.stop("analysis-stop", execute=True) == receipt + with pytest.raises(ValueError, match="start a new operation id"): + runner.resume("analysis-stop") + observed = runner.read("analysis-stop") + assert observed["stop"]["phase"] == "unknown" and not observed["recovery_required"] + assert runner.wait("analysis-stop")["stop"]["phase"] == "unknown" + assert not demo.canonical_tasks(root)["todo_analyst-initial"]["done"] + + +def test_sigterm_without_a_stop_keeps_the_operation_recoverable(service, monkeypatch): + """A shutdown signal is not a stop: the worker dies as before and resume stays available.""" + import signal + + root, runner = service + start_held_worker(service, "analysis-term") + path = runner.path("analysis-term") + worker = _read(path)["worker"] + target = runner._lane_target(runner.binding("analysis")) + turn_child = turn_lane_liveness(target)["holder"]["pid"] + try: + os.kill(worker["pid"], signal.SIGTERM) + assert until(lambda: runner._operation_lock_free(path), timeout=20) + # Default termination, exactly as before: only the worker died, and its + # orphaned Turn child still holds the member's lane. + lane = turn_lane_liveness(target) + assert lane["state"] == "live" and lane["holder"]["pid"] == turn_child + assert not runner._stop_path(path).exists() + assert _read(path)["status"] == "running" and "stop" not in runner.read("analysis-term") + spawned = [] + monkeypatch.setattr(runner, "_spawn", spawned.append) + runner.resume("analysis-term") + assert spawned == ["analysis-term"] + finally: + try: + os.killpg(worker["pgid"], signal.SIGKILL) # the orphaned Turn child + except ProcessLookupError: + pass + + +def test_stop_settlement_never_refuses_a_concurrent_turn_on_the_member_lane(service, monkeypatch): + """Settling reads the lane holder record; a real Turn racing it is always admitted.""" + import threading + from loopx import file_lock + + _, runner = service + monkeypatch.setattr(runner, "_spawn", lambda _: None) + runner.start("analysis", "analysis-race", brief()) + path = runner.path("analysis-race") + binding = runner.binding("analysis") + target = runner._lane_target(binding) + lane_lock = target.with_name(target.name + ".lock") + opened = [] + real_open = file_lock._open_lock_descriptor + + def recording_open(lock_path, **kwargs): + opened.append((threading.get_ident(), Path(lock_path))) + return real_open(lock_path, **kwargs) + + monkeypatch.setattr(file_lock, "_open_lock_descriptor", recording_open) + lane = {"runtime_root": runner.root, "goal_id": runner.goal_id, + "plan": {"turn_envelope": {"agent_id": binding["agent_id"]}}} + admitted, refused, phases, settlers = [], [], [], [] + done = Event() + + def settle_continuously(): + settlers.append(threading.get_ident()) + while not done.is_set(): + phases.append(runner.stop("analysis-race", execute=True)["phase"]) + + # An unnamed holder keeps the operation open, so every stop call settles again + # and reads the member's lane while other Turns of that member take it. + with exclusive_file_lock(path), ThreadPoolExecutor(max_workers=1) as pool: + future = pool.submit(settle_continuously) + try: + assert until(lambda: len(phases) >= 2, timeout=30) + for _ in range(200): + with turn_lane_singleflight(**lane) as held: + (admitted if held is not None else refused).append(held) + assert until(lambda: len(phases) >= 4, timeout=30) + finally: + done.set() + future.result(timeout=30) + assert refused == [] and len(admitted) == 200 + assert set(phases) == {"requested"} + assert [lock for ident, lock in opened if ident in settlers and lock == lane_lock] == [] + # Once the holder is gone the requester acknowledges; settling still never takes the lane. + opened.clear() + receipt = runner.stop("analysis-race", execute=True) + assert receipt["phase"] == "settled" and receipt["stop"]["settled"]["lane_state"] == "released" + assert lane_lock not in {lock for _, lock in opened} + + +def test_fenced_write_after_another_process_stop_writes_nothing(service, monkeypatch): + root, runner = service + monkeypatch.setattr(runner, "_spawn", lambda _: None) + runner.start("analysis", "analysis-fenced", brief()) + path = runner.path("analysis-fenced") + row = _read(path) + foreign = runner._new_stop_record(row, requested_by="other-host-lead", worker=None) + foreign.update(phase="acknowledged", ack={"pid": 1, "host": "elsewhere", "at": 0.0, + "source": "requester", "observed_status": "running", + "turn_key": None}) + from loopx.control_plane.collaboration.inbox import _write + + _write(runner._stop_path(path), foreign) + frozen = path.read_bytes() + with pytest.raises(DelegationFenced): + runner._fenced_write(path, {**row, "status": "running"}) + with pytest.raises(DelegationFenced): + runner._observe(path, dict(row), "running") + with pytest.raises(DelegationFenced): + runner._record_turn_result(path, {**row, "status": "running"}, + {"status": "committed", "result_kind": "validated_progress"}) + runner.execute("analysis-fenced") # the foreign acknowledgement stands; nothing is rewritten + assert path.read_bytes() == frozen + assert _read(runner._stop_path(path)) == foreign + assert not (root / "host-started").exists() + assert not demo.canonical_tasks(root)["todo_analyst-initial"]["done"] + # A stop this process may acknowledge is taken from under the lock at entry. + runner.start("analysis", "analysis-entry", brief()) + entry = runner.path("analysis-entry") + _write(runner._stop_path(entry), runner._new_stop_record(_read(entry), requested_by="lead", worker=None)) + runner.execute("analysis-entry") + assert _read(entry)["status"] == "stopped" + acknowledged = _read(runner._stop_path(entry)) + assert acknowledged["phase"] == "acknowledged" and acknowledged["ack"]["source"] == "worker_entry" + assert runner.stop("analysis-entry", execute=True)["phase"] == "settled" + + +def test_stop_settles_only_after_the_owned_host_and_its_descendants_exit(service): + """A host and its same-group child that ignore SIGTERM keep the stop open until they exit. + + The postcondition is checked at the instant ``settled`` returns, not after a wait. + """ + root, runner = service + (root / "ignore-term").touch() + host_pid = start_held_worker(service) + assert until(lambda: (root / "host-child-pid").exists()) + child_pid = int((root / "host-child-pid").read_text()) + assert not process_gone(host_pid) and not process_gone(child_pid) + receipt = runner.stop("analysis-stop", execute=True) + host_gone, child_gone = process_gone(host_pid), process_gone(child_pid) + assert receipt["phase"] == "settled", receipt + assert host_gone and child_gone, (receipt, host_gone, child_gone) + settled = receipt["stop"]["settled"] + assert settled["host_process"] == "drained", settled + # Rereading a settled stop neither reopens it nor admits or completes anything. + frozen = runner.path("analysis-stop").read_bytes() + assert runner.stop("analysis-stop", execute=True) == receipt + assert runner.read("analysis-stop")["stop"]["phase"] == "settled" + assert runner.path("analysis-stop").read_bytes() == frozen + assert int((root / "analyst" / "initial" / "host-invocations").read_text()) == 1 + assert not demo.canonical_tasks(root)["todo_analyst-initial"]["done"] + assert returns(runner.root, runner.goal_id, "lead")["items"] == [] + + +def test_interrupted_host_cleanup_keeps_the_stop_open_until_a_reread_sees_it_drained(service, monkeypatch): + """A Host supervisor that never finishes cleaning up leaves no settlement to claim. + + Rereads with the same identity stay acknowledged without admitting or completing + anything, and settle once the Host's group is observed gone. + """ + import signal + from loopx import collaboration_mcp + + root, runner = service + monkeypatch.setattr(collaboration_mcp, "DELEGATION_STOP_GRACE_SECONDS", 2.0) + host_pid = start_held_worker(service) + path = runner.path("analysis-stop") + record = json.loads(runner._host_process_record(path).read_text()) + assert record["phase"] == "spawned" and record["host_pid"] == host_pid == record["process_group"] + bridge = record["bridge_pid"] + os.kill(bridge, signal.SIGSTOP) # the supervisor cannot run its cleanup + try: + first = runner.stop("analysis-stop", execute=True) + assert first["phase"] == "acknowledged" and first["status"] == "stopped", first + assert first["reason"] == "host_process_still_running" and not process_gone(host_pid) + os.kill(bridge, signal.SIGKILL) # and now never will: the Host is orphaned + assert until(lambda: process_gone(bridge), timeout=20) + frozen = path.read_bytes() + for _ in range(3): + again = runner.stop("analysis-stop", execute=True) + assert again["phase"] == "acknowledged" and again["reason"] == "host_process_still_running", again + assert again["stop"]["stop_id"] == first["stop"]["stop_id"] + assert again["stop"]["ack"] == first["stop"]["ack"] and again["stop"]["settled"] is None + assert runner.read("analysis-stop")["stop"]["phase"] == "acknowledged" + assert not process_gone(host_pid) and path.read_bytes() == frozen + with pytest.raises(ValueError, match="start a new operation id"): + runner.resume("analysis-stop") + finally: + for target, sig in ((bridge, signal.SIGCONT), (host_pid, signal.SIGKILL)): + try: + os.killpg(target, sig) if target == host_pid else os.kill(target, sig) + except ProcessLookupError: + pass + assert until(lambda: process_gone(host_pid), timeout=20) + receipt = runner.stop("analysis-stop", execute=True) + assert receipt["phase"] == "settled" and receipt["stop"]["settled"]["host_process"] == "drained", receipt + assert receipt["stop"]["stop_id"] == first["stop"]["stop_id"] + assert runner.stop("analysis-stop", execute=True) == receipt + assert path.read_bytes() == frozen + assert int((root / "analyst" / "initial" / "host-invocations").read_text()) == 1 + assert not demo.canonical_tasks(root)["todo_analyst-initial"]["done"] + assert returns(runner.root, runner.goal_id, "lead")["items"] == [] + + +def test_a_stop_written_before_the_effects_linearizes_against_them(service, monkeypatch): + """Todo completion and reply publication commit under the stop's own lock. + + The pre-check alone was not a boundary: the worker could pass it, a stop + could be persisted, and both external effects still committed before the + fenced record write, leaving `settled`/`stopped` for a member whose work had + landed. Holding the dispatch lock across both effects makes the two sides + linearize in either order. + """ + from loopx.control_plane.collaboration.inbox import _write as write_inbox + + root, runner = service + runner.start("analysis", "analysis-race", brief()) + path = runner.path("analysis-race") + row = _read(path) + # Reproduce the interleaving: the stop exists before the worker reaches its + # effects. A same-process request is what the worker then acknowledges. + write_inbox(runner._stop_path(path), runner._new_stop_record( + row, requested_by=runner.agent_id, worker=None)) + + # The checkpoint after the model returns refuses to run the effects. + with pytest.raises(DelegationStopRequested): + with exclusive_file_lock(runner._dispatch_lock(path)): + runner._raise_if_stop_requested(path) + + runner.execute("analysis-race") + assert _read(path)["status"] == "stopped" + # Neither effect committed for a stop that was written first. + assert not demo.canonical_tasks(root)["todo_analyst-initial"]["done"] + assert not (root / "runtime" / "replies" / "analysis-race" / "conclusion.json").exists() + assert not (root / "host-started").exists() + receipt = runner.stop("analysis-race", execute=True) + assert receipt["phase"] == "settled" and receipt["status"] == "stopped" + + +def test_a_failed_required_lease_release_keeps_the_stop_open_and_retries(service, monkeypatch): + """A required lease the stop could not release is not a settlement. + + The member's Todo can stay blocked by that lease until its TTL, so reporting + `settled` would be a terminal claim the owner cannot act on. The release is + retried under the stop's own lock on the next read instead of being attempted + once and forgotten. + """ + from loopx import collaboration_mcp as delegation + from loopx.control_plane.collaboration.inbox import _write as write_inbox + + root, runner = service + monkeypatch.setattr(runner, "_spawn", lambda _: None) + runner.start("analysis", "analysis-lease", brief()) + path = runner.path("analysis-lease") + row = _read(path) + # A bounded member task holds a required lease; make the record say so. + row["task_lease"] = {"required": True, "idempotency_key": "lease-1", "version": 1} + row["status"] = "stopped" + runner._fenced_write(path, row) + write_inbox(runner._stop_path(path), { + **runner._new_stop_record(row, requested_by=runner.agent_id, worker=None), + "phase": "acknowledged", + "ack": {"pid": os.getpid(), "host": lock_holder_host_label(), "at": time.time(), + "source": "requester", "observed_status": "stopped", "turn_key": None}, + # The acknowledgement could not release it, which is what the receipt records. + "lease": {"required": True, "released": False, "error": "authority unavailable"}, + }) + + attempts = [] + outcomes = [RuntimeError("authority temporarily unavailable"), {"released": True}] + + def flaky_release(**kwargs): + attempts.append(kwargs) + outcome = outcomes[min(len(attempts) - 1, len(outcomes) - 1)] + if isinstance(outcome, Exception): + raise outcome + return outcome + + monkeypatch.setattr(delegation, "release_task_lease", flaky_release) + + # The first settle read tries the release, fails, and must not settle. + receipt = runner.stop("analysis-lease", execute=True) + assert attempts, "the stop never attempted the required release" + assert receipt["phase"] == "acknowledged", receipt + assert receipt["stop"]["reason"] == "required_lease_release_unproven" + assert receipt["stop"]["lease"]["released"] is not True + + # The next read retries the release and only then settles. + settled = runner.stop("analysis-lease", execute=True) + assert len(attempts) >= 2, "the failed release was never retried" + assert settled["phase"] == "settled", settled + assert settled["stop"]["lease"]["released"] is True + assert settled["stop"]["settled"]["lease_released"] is True + + # Once settled the receipt is stable, and a released lease is not re-attempted. + before = len(attempts) + assert runner.stop("analysis-lease", execute=True) == settled + assert len(attempts) == before + + +def test_a_launched_host_on_a_platform_without_process_groups_fails_fast(service, monkeypatch): + """A stop that cannot prove its Host drained says so instead of never settling. + + Windows cleanup is process-tree best effort, so no fact proves the Host's + descendants exited. An acknowledged receipt the caller can never settle is + worse than an actionable refusal that names the platform boundary. + """ + from loopx.control_plane.turn_driver import host_process_transport + + root, runner = service + monkeypatch.setattr(runner, "_spawn", lambda _: None) + runner.start("analysis", "analysis-platform", brief()) + path = runner.path("analysis-platform") + record = runner._host_process_record(path) + record.parent.mkdir(parents=True, exist_ok=True) + record.write_text(json.dumps({ + "schema_version": host_process_transport.HOST_PROCESS_RECORD_SCHEMA_VERSION, + "host": lock_holder_host_label(), "phase": "finished", + "bridge_pid": os.getpid(), "process_group": os.getpid(), + })) + + # The launched Host is real, but this platform cannot prove it drained. + assert host_process_transport.host_process_drain(record) == host_process_transport.HOST_PROCESS_DRAINED + monkeypatch.delattr(host_process_transport.os, "killpg", raising=False) + assert host_process_transport.host_process_drain(record) == host_process_transport.HOST_PROCESS_UNSUPPORTED_PLATFORM + + with pytest.raises(ValueError, match="cannot prove the launched Host drained"): + runner.stop("analysis-platform", execute=True) + + +def test_a_crash_between_the_ack_and_the_lease_result_keeps_the_stop_open(service, monkeypatch): + """The lease obligation survives process loss after the acknowledgement. + + `_acknowledge_stop` writes the ACK before it releases the lease, so a crash + in between leaves the sidecar with no `lease` field. Reading that as "nothing + was owed" settles a stop whose member still holds an active hard lease, with + resume already refused and the Todo blocked until the TTL. The obligation + comes from the operation record, which the crash cannot lose. + """ + from loopx.control_plane.collaboration.inbox import _write as write_inbox + + root, runner = service + monkeypatch.setattr(runner, "_spawn", lambda _: None) + runner.start("analysis", "analysis-crash", brief()) + path = runner.path("analysis-crash") + row = _read(path) + # The member holds a real required lease, as a bounded task does. + row["task_lease"] = {"required": True, "idempotency_key": "lease-crash", "version": 1} + right_after_ack = {**row, "status": "stopped"} + runner._fenced_write(path, right_after_ack) + + # The crash window: ACK persisted, no lease result written yet. + write_inbox(runner._stop_path(path), { + **runner._new_stop_record(right_after_ack, requested_by=runner.agent_id, worker=None), + "phase": "acknowledged", + "ack": {"pid": os.getpid(), "host": lock_holder_host_label(), "at": time.time(), + "source": "requester", "observed_status": "stopped", "turn_key": None}, + }) + sidecar = _read(runner._stop_path(path)) + assert "lease" not in sidecar or sidecar.get("lease") is None + assert _read(path)["task_lease"]["required"] is True + + # A release that succeeds on this read lets the stop settle, and the receipt + # says the lease really is gone. + monkeypatch.setattr(delegation_module, "release_task_lease", + lambda **kw: {"released": True}) + settled = runner.stop("analysis-crash", execute=True) + assert settled["phase"] == "settled", settled + assert settled["stop"]["lease"]["released"] is True + assert settled["stop"]["settled"]["lease_released"] is True + + +def test_a_crash_between_the_ack_and_the_lease_result_never_settles_unreleased(service, monkeypatch): + """The same window, with the release still failing, must not report `settled`. + + `settled` tells the owner the member is safely stopped. Claiming it while a + required lease is provably still active is exactly the terminal distortion + the crash window used to produce. + """ + from loopx.control_plane.collaboration.inbox import _write as write_inbox + + root, runner = service + monkeypatch.setattr(runner, "_spawn", lambda _: None) + runner.start("analysis", "analysis-crash-open", brief()) + path = runner.path("analysis-crash-open") + row = _read(path) + row["task_lease"] = {"required": True, "idempotency_key": "lease-open", "version": 1} + right_after_ack = {**row, "status": "stopped"} + runner._fenced_write(path, right_after_ack) + write_inbox(runner._stop_path(path), { + **runner._new_stop_record(right_after_ack, requested_by=runner.agent_id, worker=None), + "phase": "acknowledged", + "ack": {"pid": os.getpid(), "host": lock_holder_host_label(), "at": time.time(), + "source": "requester", "observed_status": "stopped", "turn_key": None}, + }) + + def still_held(**kwargs): + raise RuntimeError("authority unavailable") + + monkeypatch.setattr(delegation_module, "release_task_lease", still_held) + receipt = runner.stop("analysis-crash-open", execute=True) + assert receipt["phase"] == "acknowledged", receipt + assert receipt["stop"]["reason"] == "required_lease_release_unproven"