From 6485ebe4fb5a45c8847ffa3038c8f08dff32f4c1 Mon Sep 17 00:00:00 2001 From: song <22676124+songoow@users.noreply.github.com> Date: Tue, 29 Sep 2026 06:23:12 -0400 Subject: [PATCH 1/9] feat(turn-driver): read Turn lane liveness without taking the lane Add `lock_holder_liveness` to the file-lock owner and `turn_lane_liveness` on top of it. Both classify a lane's last executing Turn from its holder record alone: `released` on a clean exit, `dead` when the record names this machine and the pid is gone, `foreign_host` when the pid cannot be checked here, `unreadable` when a lock file carries no parseable record, and `live` only when a same-host pid is still alive. The probe never touches the kernel lock. A probe that acquired it for an instant would refuse a real `run-once --execute` racing that instant with `turn_lane_in_flight` for nothing; the new test drives the real fence wrapper concurrently with a continuous probe and proves the Turn is admitted exactly once. The holder host label is now single-sourced so the writer and the readers cannot disagree on what "this machine" is. Co-Authored-By: Claude Fable 5.1 Signed-off-by: song <22676124+songoow@users.noreply.github.com> --- loopx/control_plane/turn_driver/lane_fence.py | 67 +++++--- loopx/file_lock.py | 91 +++++++++-- tests/test_turn_lane_fence.py | 144 +++++++++++++++++- 3 files changed, 274 insertions(+), 28 deletions(-) diff --git a/loopx/control_plane/turn_driver/lane_fence.py b/loopx/control_plane/turn_driver/lane_fence.py index 5ac99e8f7..9f8ef860e 100644 --- a/loopx/control_plane/turn_driver/lane_fence.py +++ b/loopx/control_plane/turn_driver/lane_fence.py @@ -17,12 +17,20 @@ from contextlib import contextmanager from functools import wraps import hashlib -import json from pathlib import Path import re from typing import Any -from ...file_lock import lock_holder_path, try_exclusive_file_lock +from ...file_lock import ( + LOCK_HOLDER_ABSENT, + LOCK_HOLDER_DEAD, + LOCK_HOLDER_FOREIGN_HOST, + LOCK_HOLDER_LIVE, + LOCK_HOLDER_RELEASED, + LOCK_HOLDER_UNREADABLE, + lock_holder_liveness, + try_exclusive_file_lock, +) # Typed refusal for a lane whose single executor is already busy. The reason is # a fact about this lane, so a caller can retry it unchanged once it clears. @@ -103,6 +111,18 @@ def turn_lane_singleflight( yield lock_path +def _public_holder(record: Mapping[str, Any]) -> dict[str, Any]: + projection: dict[str, Any] = {} + for field in TURN_LANE_HOLDER_TEXT_FIELDS: + value = record.get(field) + if isinstance(value, str) and value: + projection[field] = value + pid = record.get("pid") + if isinstance(pid, int): + projection["pid"] = pid + return projection + + def turn_lane_holder_readback(target: Path) -> dict[str, Any]: """Return the public-safe identity of the Turn holding one lane, else ``{}``. @@ -113,21 +133,34 @@ def turn_lane_holder_readback(target: Path) -> dict[str, Any]: hosts share one runtime root. """ - try: - record = json.loads(lock_holder_path(target).read_text(encoding="utf-8")) - except (OSError, ValueError): - return {} - if not isinstance(record, Mapping): - return {} - projection: dict[str, Any] = {} - for field in TURN_LANE_HOLDER_TEXT_FIELDS: - value = record.get(field) - if isinstance(value, str) and value: - projection[field] = value - pid = record.get("pid") - if isinstance(pid, int): - projection["pid"] = pid - return projection + _state, record = lock_holder_liveness(target) + return _public_holder(record) + + +# Lane liveness vocabulary: the lock owner's holder states, named here so a +# projection can switch on them without learning the lock record format. +TURN_LANE_LIVE = LOCK_HOLDER_LIVE +TURN_LANE_RELEASED = LOCK_HOLDER_RELEASED +TURN_LANE_DEAD = LOCK_HOLDER_DEAD +TURN_LANE_FOREIGN_HOST = LOCK_HOLDER_FOREIGN_HOST +TURN_LANE_UNREADABLE = LOCK_HOLDER_UNREADABLE +TURN_LANE_ABSENT = LOCK_HOLDER_ABSENT + + +def turn_lane_liveness(target: Path) -> dict[str, Any]: + """Say whether one lane's last executing Turn is still running, read-only. + + The answer comes from the holder record alone: ``released_at`` for a clean + exit, the machine name for whether the pid can be checked here, and pid + liveness for a holder that never released. This never takes the lane lock, + not even for an instant: a probe that did would refuse a real Turn racing + the same instant with ``turn_lane_in_flight`` for no reason. ``live`` is + the only state that is evidence of execution; ``foreign_host`` and + ``unreadable`` are unknowns a consumer must fail closed on. + """ + + state, record = lock_holder_liveness(target) + return {"state": state, "holder": _public_holder(record)} def turn_lane_in_flight_record( diff --git a/loopx/file_lock.py b/loopx/file_lock.py index 4e25c05f9..7ffa39682 100644 --- a/loopx/file_lock.py +++ b/loopx/file_lock.py @@ -175,7 +175,7 @@ def _identity( # hosts sharing one runtime root can both read the holder, so the record # names its own machine and a reader never has to guess which host a pid # belongs to. The name is a sanitized label, not a path or a secret. - "host": _safe_label(socket.gethostname(), fallback="unknown"), + "host": lock_holder_host_label(), "agent_id": _safe_label( agent_id or os.environ.get("LOOPX_AGENT_ID"), fallback="unknown", @@ -268,14 +268,8 @@ def _mark_released( pass -def _read_holder_record(lock_path: Path) -> dict[str, object]: - try: - payload = json.loads(lock_path.read_text(encoding="utf-8")) - except (OSError, json.JSONDecodeError): - return {} - if not isinstance(payload, dict): - return {} - allowed = { +_HOLDER_RECORD_FIELDS = frozenset( + { "schema_version", "lock_id", "policy", @@ -286,7 +280,84 @@ def _read_holder_record(lock_path: Path) -> dict[str, object]: "acquired_at", "released_at", } - return {key: payload[key] for key in allowed if key in payload} +) + + +def _filter_holder_record(payload: object) -> dict[str, object]: + if not isinstance(payload, dict): + return {} + return {key: payload[key] for key in _HOLDER_RECORD_FIELDS if key in payload} + + +def _read_holder_record(lock_path: Path) -> dict[str, object]: + try: + payload = json.loads(lock_path.read_text(encoding="utf-8")) + except (OSError, json.JSONDecodeError): + return {} + return _filter_holder_record(payload) + + +def lock_holder_host_label() -> str: + """The machine label a holder record carries; readers compare against it.""" + + return _safe_label(socket.gethostname(), fallback="unknown") + + +# Liveness of a lock's last holder, read from its record alone. The kernel lock +# is never probed: a probe would hold the lock for an instant, and a real +# single-flight acquisition racing that instant would be refused for nothing. +LOCK_HOLDER_LIVE = "live" +LOCK_HOLDER_RELEASED = "released" +LOCK_HOLDER_DEAD = "dead" +LOCK_HOLDER_FOREIGN_HOST = "foreign_host" +LOCK_HOLDER_UNREADABLE = "unreadable" +LOCK_HOLDER_ABSENT = "absent" +LOCK_HOLDER_LIVENESS_STATES = ( + LOCK_HOLDER_LIVE, + LOCK_HOLDER_RELEASED, + LOCK_HOLDER_DEAD, + LOCK_HOLDER_FOREIGN_HOST, + LOCK_HOLDER_UNREADABLE, + LOCK_HOLDER_ABSENT, +) + + +def lock_holder_liveness(path: Path) -> tuple[str, dict[str, object]]: + """Classify the last holder of one lock without touching the kernel lock. + + Returns the liveness state and the filtered holder record. ``released`` + means the holder wrote ``released_at`` on a clean exit; ``dead`` means the + record names this machine and the pid is gone, which is what a crashed or + killed holder leaves behind; ``foreign_host`` means the pid cannot be + checked from here; ``unreadable`` means a lock file exists but carries no + parseable record, for example mid-acquisition. Only ``live`` is evidence of + a running holder, and even that is pid liveness, not the kernel lock: a + reused pid can keep a crashed holder looking alive until the next holder + overwrites the record. + """ + + holder_path = lock_holder_path(path) + try: + text = holder_path.read_text(encoding="utf-8") + except FileNotFoundError: + return LOCK_HOLDER_ABSENT, {} + except OSError: + return LOCK_HOLDER_UNREADABLE, {} + try: + record = _filter_holder_record(json.loads(text)) + except ValueError: + return LOCK_HOLDER_UNREADABLE, {} + if not record: + return LOCK_HOLDER_UNREADABLE, {} + released_at = record.get("released_at") + if isinstance(released_at, str) and released_at: + return LOCK_HOLDER_RELEASED, record + if record.get("host") != lock_holder_host_label(): + return LOCK_HOLDER_FOREIGN_HOST, record + pid = record.get("pid") + if isinstance(pid, bool) or not isinstance(pid, int): + return LOCK_HOLDER_UNREADABLE, record + return (LOCK_HOLDER_LIVE if process_is_alive(pid) else LOCK_HOLDER_DEAD), record def _operator_action(holder: dict[str, object], *, retry_mode: str) -> dict[str, object]: diff --git a/tests/test_turn_lane_fence.py b/tests/test_turn_lane_fence.py index 0341a6666..2f5179f4b 100644 --- a/tests/test_turn_lane_fence.py +++ b/tests/test_turn_lane_fence.py @@ -8,16 +8,29 @@ from __future__ import annotations +import json +import os import socket +import threading +import time from pathlib import Path from loopx.control_plane.turn_driver.executor import run_loopx_turn_once -from loopx.file_lock import _safe_label +from loopx.file_lock import _safe_label, lock_holder_path +from loopx.control_plane.turn_driver import lane_fence from loopx.control_plane.turn_driver.lane_fence import ( REMEDY_WAIT_FOR_IN_FLIGHT_TURN, + single_executor_per_turn_lane, + TURN_LANE_ABSENT, + TURN_LANE_DEAD, + TURN_LANE_FOREIGN_HOST, TURN_LANE_IN_FLIGHT, + TURN_LANE_LIVE, TURN_LANE_OPERATION, + TURN_LANE_RELEASED, + TURN_LANE_UNREADABLE, turn_lane_holder_readback, + turn_lane_liveness, turn_lane_singleflight, turn_lane_target, ) @@ -128,3 +141,132 @@ def test_the_holder_readback_stays_public_safe(tmp_path: Path) -> None: # The private lock identity and the runtime path never leave the process. assert str(tmp_path) not in str(holder) assert turn_lane_holder_readback(tmp_path / "absent.lane") == {} + + +def _lane(tmp_path: Path) -> Path: + return turn_lane_target( + runtime_root=tmp_path / "runtime", goal_id=GOAL_ID, plan=_plan() + ) + + +def _rewrite_holder(target: Path, **changes: object) -> None: + """Edit the holder record the way a crash or another machine would leave it.""" + + holder_path = lock_holder_path(target) + record = json.loads(holder_path.read_text(encoding="utf-8")) + record.pop("released_at", None) + record.update(changes) + holder_path.write_text(json.dumps(record), encoding="utf-8") + + +def test_liveness_follows_the_lane_from_absent_to_live_to_released( + tmp_path: Path, +) -> None: + target = _lane(tmp_path) + assert turn_lane_liveness(target) == {"state": TURN_LANE_ABSENT, "holder": {}} + + with turn_lane_singleflight( + runtime_root=tmp_path / "runtime", goal_id=GOAL_ID, plan=_plan() + ) as held: + assert held is not None + live = turn_lane_liveness(target) + + assert live["state"] == TURN_LANE_LIVE + # The holder is the same public-safe readback a refusal names. + assert live["holder"] == turn_lane_holder_readback(target) | {"pid": os.getpid()} + assert live["holder"]["pid"] == os.getpid() + assert str(tmp_path) not in json.dumps(live) + # A clean exit is a release, whatever the pid does afterwards. + assert turn_lane_liveness(target)["state"] == TURN_LANE_RELEASED + + +def test_liveness_fails_closed_on_dead_foreign_and_unreadable_holders( + tmp_path: Path, +) -> None: + target = _lane(tmp_path) + with turn_lane_singleflight( + runtime_root=tmp_path / "runtime", goal_id=GOAL_ID, plan=_plan() + ): + pass + + # A killed Turn never writes released_at; its pid is gone on this machine. + dead_pid = os.getpid() + while True: + dead_pid += 1 + try: + os.kill(dead_pid, 0) + except ProcessLookupError: + break + except OSError: + continue + if dead_pid > os.getpid() + 100_000: + raise AssertionError("no free pid found near this process") + _rewrite_holder(target, pid=dead_pid) + assert turn_lane_liveness(target)["state"] == TURN_LANE_DEAD + + # A holder on another machine cannot be pid-checked here, even if that pid + # happens to be alive on this one. + _rewrite_holder(target, pid=os.getpid(), host="another-machine") + foreign = turn_lane_liveness(target) + assert foreign["state"] == TURN_LANE_FOREIGN_HOST + assert foreign["holder"]["host"] == "another-machine" + + # A lock file with no parseable record is mid-acquisition or corrupt: not + # absent, and not evidence of anything. + lock_holder_path(target).write_text("", encoding="utf-8") + assert turn_lane_liveness(target) == {"state": TURN_LANE_UNREADABLE, "holder": {}} + lock_holder_path(target).write_text("{}", encoding="utf-8") + assert turn_lane_liveness(target)["state"] == TURN_LANE_UNREADABLE + + +def test_the_liveness_probe_never_refuses_a_concurrent_executing_turn( + tmp_path: Path, monkeypatch +) -> None: + """The probe reads a record; it never takes the lane, not even for an instant.""" + + fence_calls: list[str] = [] + real_fence = lane_fence.try_exclusive_file_lock + + def counting_fence(*args, **kwargs): + fence_calls.append(str(kwargs.get("operation"))) + return real_fence(*args, **kwargs) + + monkeypatch.setattr(lane_fence, "try_exclusive_file_lock", counting_fence) + target = _lane(tmp_path) + turn_lane_liveness(target) + turn_lane_holder_readback(target) + assert fence_calls == [] + + # The executing entry is the real fence wrapper run-once --execute goes + # through; only the Turn body is a stand-in that holds the lane a moment. + @single_executor_per_turn_lane( + lambda plan, record, **kwargs: {**record, "effects": kwargs["effects"]} + ) + def executing_turn(plan, *, runtime_root, goal_id, execute): + time.sleep(0.3) + return {"status": "committed", "held": turn_lane_liveness(target)["state"]} + + observed: set[str] = set() + stop = threading.Event() + + def probe() -> None: + while not stop.is_set(): + observed.add(turn_lane_liveness(target)["state"]) + + prober = threading.Thread(target=probe, daemon=True) + prober.start() + try: + payload = executing_turn( + _plan(), runtime_root=tmp_path / "runtime", goal_id=GOAL_ID, execute=True + ) + finally: + stop.set() + prober.join(timeout=5) + + # The Turn took the fence exactly once and was never told the lane was busy. + assert fence_calls == [TURN_LANE_OPERATION] + assert payload == {"status": "committed", "held": TURN_LANE_LIVE} + assert payload.get("reason") != TURN_LANE_IN_FLIGHT + assert observed <= {TURN_LANE_ABSENT, TURN_LANE_LIVE, TURN_LANE_RELEASED, TURN_LANE_UNREADABLE} + assert TURN_LANE_LIVE in observed + assert turn_lane_liveness(target)["state"] == TURN_LANE_RELEASED From 692f4f6bf14aef86792d6558b542436c9511d10c Mon Sep 17 00:00:00 2001 From: song <22676124+songoow@users.noreply.github.com> Date: Tue, 29 Sep 2026 06:28:55 -0400 Subject: [PATCH 2/9] fix(agents): derive worker executing state from lane and lease facts `_agent_state` returned `executing` whenever the current Todo had been updated in the last eight hours. No execution fact backed it, so a dead peer with a fresh Todo read as running and peer activation admitted it. Add `agents/execution_facts.py`, a read-only collector over facts LoopX already keeps: Turn lane liveness per Goal and agent (a delegated member executes inside its own lane), the delegation worker's operation lock, and the task lease on a claimed Todo (canonical head after cutover, local lease files before). Status collection attaches the map to the payload as `agent_execution_facts` and passes it to the projection, so the peer directory re-projects the same facts and drops `lease_state_not_projected` when they are present; each agent row carries the evidence as `execution`. Derivation order is now blocked, monitoring/waiting, `executing` only when the lane is live or a delegation worker holds its lock, then the new `unknown` when the holder is on a foreign host or unreadable, or an active lease has expired with nothing live, then bound/launchable/addressable/ registered. `EXECUTING_ACTIVITY_THRESHOLD_HOURS` is gone; Todo timestamps remain only in `last_activity_at` and `stale_claim_hint`. Default behavior change: open work updated within eight hours no longer reads as `executing`; without facts it is `bound` or `launchable`. The lifecycle tests and the worker-lifecycle smoke encode the new rule, with facts fixtures built from the real runtime layout. Co-Authored-By: Claude Fable 5.1 Signed-off-by: song <22676124+songoow@users.noreply.github.com> --- examples/worker-lifecycle-state-smoke.py | 23 +- loopx/control_plane/agents/directory.py | 23 +- loopx/control_plane/agents/execution_facts.py | 217 +++++++++++++++++ .../agents/management_projection.py | 130 ++++++++-- loopx/control_plane/status/collection.py | 10 + .../test_agent_lifecycle_state.py | 225 ++++++++++++++++-- 6 files changed, 581 insertions(+), 47 deletions(-) create mode 100644 loopx/control_plane/agents/execution_facts.py diff --git a/examples/worker-lifecycle-state-smoke.py b/examples/worker-lifecycle-state-smoke.py index c74f3d196..7cc48d28b 100644 --- a/examples/worker-lifecycle-state-smoke.py +++ b/examples/worker-lifecycle-state-smoke.py @@ -2,7 +2,7 @@ """Smoke test for worker lifecycle state projection. Verifies that the agent management projection correctly derives lifecycle -states from existing facts (registry, todo, session binding, activity). +states from existing facts (registry, todo, session binding, execution facts). Run from the repository root: uv run --extra test python examples/worker-lifecycle-state-smoke.py @@ -19,9 +19,14 @@ def _recent_activity() -> str: - """Activity timestamp within the activity threshold (8 hours).""" + """A fresh Todo update: activity, which on its own never means execution.""" return (datetime.now(timezone.utc) - timedelta(hours=1)).isoformat() + +# Execution facts are what make a worker `executing`: here the worker's Turn +# lane is live. The projection reads them; it never infers them from age. +EXECUTION_FACTS = {"worker-executing": {"lane": "live"}} + from loopx.control_plane.agents.management_projection import ( # noqa: E402 WORKER_LIFECYCLE_STATE_ADDRESSABLE, WORKER_LIFECYCLE_STATE_BLOCKED, @@ -112,7 +117,19 @@ def build_status_payload() -> dict: def main() -> int: payload = build_status_payload() - projection = build_agent_management_projection(payload) + projection = build_agent_management_projection( + payload, execution_facts=EXECUTION_FACTS + ) + without_facts = build_agent_management_projection(payload) + timestamp_only = { + a["agent_id"]: a["state"] for a in without_facts.get("agents", []) + }.get("worker-executing") + if timestamp_only == WORKER_LIFECYCLE_STATE_EXECUTING: + print( + "FAIL: a fresh Todo timestamp without execution facts must not " + "read as executing" + ) + return 1 agents = {a["agent_id"]: a for a in projection.get("agents", [])} diff --git a/loopx/control_plane/agents/directory.py b/loopx/control_plane/agents/directory.py index 80379e654..ecca13aa1 100644 --- a/loopx/control_plane/agents/directory.py +++ b/loopx/control_plane/agents/directory.py @@ -15,8 +15,9 @@ - it reports no presence, because no presence provider is registered, and it says so in `presence_coverage` instead of leaving a reader to guess between "not running" and "this machine cannot see it"; -- it does not project a lease epoch, which the current projection does not own, - and it names that gap as a limitation. +- it projects lease state only when the status payload carries execution + facts; a payload without them gets `lease_state_not_projected` instead of + a guess. """ from __future__ import annotations @@ -120,9 +121,12 @@ def _work_block(agent_row: Mapping[str, Any]) -> dict[str, Any] | None: if claimed_by and _compact(todo.get("updated_at"), limit=60) else CLAIM_AGE_UNKNOWN ) + lease = _as_mapping(_as_mapping(agent_row.get("execution")).get("lease")) work: dict[str, Any] = { "todo_id": normalize_todo_id(todo_id) or todo_id, "todo_status": _compact(todo.get("status"), limit=40) or "unknown", + "lease_status": _compact(lease.get("status"), limit=40), + "lease_expired": lease.get("expired") if isinstance(lease.get("expired"), bool) else None, "task_class": _compact(todo.get("task_class"), limit=60), "action_kind": _compact(todo.get("action_kind"), limit=60), "priority": _compact(todo.get("priority"), limit=20), @@ -174,6 +178,7 @@ def build_peer_agent_directory( goal_id: str | None = None, caller_agent_id: str | None = None, available_capabilities: Any = None, + execution_facts: Any = None, ) -> dict[str, Any]: """Return a bounded `peer_agent_directory_v0` packet for one Goal. @@ -182,13 +187,22 @@ def build_peer_agent_directory( against the registry rather than asserted by the caller; when it is absent the packet records that the caller identity was not supplied instead of inventing one. + + `execution_facts` defaults to the `agent_execution_facts` status collection + attached to the payload, so this re-projection reads the same lane, worker + and lease facts the management projection did. """ payload = status_payload if isinstance(status_payload, Mapping) else {} resolved_goal = _compact(goal_id or payload.get("goal_filter"), limit=120) caller = _compact(caller_agent_id, limit=120) + if execution_facts is None: + execution_facts = payload.get("agent_execution_facts") + facts = execution_facts if isinstance(execution_facts, Mapping) else None projection = build_agent_management_projection( - dict(payload), available_capabilities=available_capabilities + dict(payload), + available_capabilities=available_capabilities, + execution_facts=dict(facts) if facts is not None else None, ) agent_rows = [row for row in _as_list(projection.get("agents")) if isinstance(row, Mapping)] registered_agent_ids = [ @@ -200,8 +214,9 @@ def build_peer_agent_directory( limitations = [ LIMITATION_PRESENCE_PROVIDER_UNAVAILABLE, LIMITATION_PRESENCE_IS_ADVISORY, - LIMITATION_LEASE_STATE_NOT_PROJECTED, ] + if facts is None: + limitations.append(LIMITATION_LEASE_STATE_NOT_PROJECTED) gaps: list[dict[str, Any]] = [] if caller and caller not in registered_agent_ids: # An unregistered caller gets a scope gap, never a listing it has no diff --git a/loopx/control_plane/agents/execution_facts.py b/loopx/control_plane/agents/execution_facts.py new file mode 100644 index 000000000..a03f102f7 --- /dev/null +++ b/loopx/control_plane/agents/execution_facts.py @@ -0,0 +1,217 @@ +"""Execution facts behind the worker lifecycle projection: read, never lock. + +`executing` has to be backed by something that runs, not by a fresh Todo +timestamp. This module gathers the facts LoopX already keeps about running +work, keyed by agent id, so the agent management projection can derive +`executing` and `unknown` from them instead of from activity age: + +- the Turn lane holder record, one per Goal and agent, read through + `turn_lane_liveness`; a delegated member executes inside its own lane too; +- the delegation worker's operation lock, whose holder is the detached worker + process executing a delegated Turn; +- the task lease on a claimed Todo, active, expired or released. + +It writes nothing, takes no lock and keeps no state of its own: every input is +owned elsewhere and this is a bounded read model over them. Absence of a fact +is reported as absence, never as "not running". +""" + +from __future__ import annotations + +from collections.abc import Iterable, Mapping +from pathlib import Path +from typing import Any + +from ...file_lock import LOCK_HOLDER_LIVE, lock_holder_liveness +from ..collaboration.inbox import _read as read_manager_context_record +from ..collaboration.inbox import _root as manager_context_root +from ..coordination.local_authority import ( + LocalCoordinationAuthorityUnavailable, + read_canonical_todos_if_promoted, +) +from ..runtime.time import now_utc +from ..todos.contract import normalize_todo_claimed_by +from ..turn_driver.lane_fence import ( + TURN_LANE_ABSENT, + TURN_LANE_DEAD, + TURN_LANE_FOREIGN_HOST, + TURN_LANE_LIVE, + TURN_LANE_RELEASED, + TURN_LANE_UNREADABLE, + turn_lane_liveness, + turn_lane_target, +) +from ..work_items.local_lease_record import TaskLeaseError, read_lease +from ..work_items.task_lease import lease_expires_at, task_lease_dir +from .management_projection import projected_agent_goals + +# One agent can hold one lane per Goal; the row reports the strongest fact. +# Unknowns outrank a plain "not running": a reader must fail closed on them. +_LANE_PRECEDENCE = ( + TURN_LANE_LIVE, + TURN_LANE_FOREIGN_HOST, + TURN_LANE_UNREADABLE, + TURN_LANE_DEAD, + TURN_LANE_RELEASED, + TURN_LANE_ABSENT, +) +LEASE_STATUS_ACTIVE = "active" +LEASE_STATUS_UNAVAILABLE = "unavailable" +_LANE_HOLDER_FIELDS = ("host", "pid", "acquired_at") + + +def _lane_rank(state: str) -> int: + return _LANE_PRECEDENCE.index(state) if state in _LANE_PRECEDENCE else len(_LANE_PRECEDENCE) + + +def _lease_rank(lease: Mapping[str, Any]) -> int: + """Fresher evidence first: an unexpired claim outranks an expired one.""" + + status = lease.get("status") + if status == LEASE_STATUS_ACTIVE: + return 0 if lease.get("expired") is False else 1 + if status == LEASE_STATUS_UNAVAILABLE: + return 3 + return 2 + + +def _lane_fact(runtime_root: Path, *, goal_id: str, agent_id: str) -> dict[str, Any]: + target = turn_lane_target( + runtime_root=runtime_root, + goal_id=goal_id, + plan={"turn_envelope": {"agent_id": agent_id}}, + ) + liveness = turn_lane_liveness(target) + fact: dict[str, Any] = {"lane": liveness["state"]} + holder = { + key: liveness["holder"][key] + for key in _LANE_HOLDER_FIELDS + if liveness["holder"].get(key) not in (None, "") + } + if holder and liveness["state"] != TURN_LANE_ABSENT: + fact["lane_holder"] = holder + return fact + + +def _goal_leases(runtime_root: Path, goal_id: str) -> list[dict[str, Any]] | None: + """Leases at the canonical head after cutover, else the local lease files. + + ``None`` means the lease authority could not be read, which is reported as + ``unavailable`` rather than as "no lease". + """ + + try: + canonical = read_canonical_todos_if_promoted( + runtime_root=runtime_root, goal_id=goal_id, include_leases=True + ) + except LocalCoordinationAuthorityUnavailable: + return None + if canonical is not None: + return [dict(lease) for lease in canonical.get("leases") or []] + lease_dir = task_lease_dir(runtime_root=runtime_root, goal_id=goal_id) + if not lease_dir.is_dir(): + return [] + leases: list[dict[str, Any]] = [] + for path in sorted(lease_dir.glob("todo_*.json")): + try: + lease = read_lease(path) + except TaskLeaseError: + # One corrupt peer lease must not hide the healthy ones. + continue + if lease: + leases.append(lease) + return leases + + +def _lease_fact(lease: Mapping[str, Any], *, at: Any) -> dict[str, Any]: + status = str(lease.get("status") or "").strip().lower() or "unknown" + fact: dict[str, Any] = {"status": status} + if status == LEASE_STATUS_ACTIVE: + expires_at = lease_expires_at(dict(lease)) + fact["expired"] = not (expires_at and expires_at > at) + return fact + + +def _delegation_worker_agents(runtime_root: Path) -> set[str]: + """Agents whose delegation worker process still holds its operation lock.""" + + executions = manager_context_root(runtime_root) / "executions" + if not executions.is_dir(): + return set() + active: set[str] = set() + for path in sorted(executions.glob("*/*.json")): + if path.is_symlink() or not path.is_file(): + continue + state, _record = lock_holder_liveness(path) + if state != LOCK_HOLDER_LIVE: + continue + try: + row = read_manager_context_record(path) + raw_agent = row["identity"]["binding"]["agent_id"] + except (OSError, ValueError, KeyError, TypeError): + continue + agent_id = normalize_todo_claimed_by(raw_agent) or str(raw_agent or "").strip() + if agent_id: + active.add(agent_id) + return active + + +def _merge_lane(row: dict[str, Any], fact: Mapping[str, Any]) -> None: + if "lane" not in row or _lane_rank(fact["lane"]) < _lane_rank(row["lane"]): + row.pop("lane_holder", None) + row.update(fact) + + +def _merge_lease(row: dict[str, Any], fact: Mapping[str, Any]) -> None: + if "lease" not in row or _lease_rank(fact) < _lease_rank(row["lease"]): + row["lease"] = dict(fact) + + +def collect_agent_execution_facts( + *, runtime_root: Path | str | None, status_payload: Mapping[str, Any] +) -> dict[str, dict[str, Any]]: + """Return ``agent_id -> {lane, lane_holder?, delegation_worker_active, lease?}``. + + Agents and Goals are the ones the management projection would row: the + Goal's registered agents plus every Todo claimant. Without a runtime root + there are no facts to read and the result is empty, which the projection + treats as "facts not collected", not as "nobody is running". + """ + + if runtime_root is None: + return {} + root = Path(runtime_root) + if not root.is_dir(): + return {} + facts: dict[str, dict[str, Any]] = {} + at = now_utc() + for goal_id, agents in projected_agent_goals(status_payload).items(): + leases = _goal_leases(root, goal_id) if agents else [] + for agent_id, spellings in agents.items(): + row = facts.setdefault(agent_id, {"delegation_worker_active": False}) + for spelling in spellings: + _merge_lane(row, _lane_fact(root, goal_id=goal_id, agent_id=spelling)) + if leases is None: + _merge_lease(row, {"status": LEASE_STATUS_UNAVAILABLE}) + for lease in leases or []: + owner = normalize_todo_claimed_by(lease.get("owner")) + if owner in facts: + _merge_lease(facts[owner], _lease_fact(lease, at=at)) + for agent_id in _delegation_worker_agents(root): + facts.setdefault(agent_id, {"lane": TURN_LANE_ABSENT})["delegation_worker_active"] = True + return facts + + +def execution_fact_rows( + facts: Mapping[str, Mapping[str, Any]] | None, agent_ids: Iterable[str] +) -> dict[str, dict[str, Any]]: + """Narrow a facts map to the named agents, for agent-lane compaction.""" + + if not isinstance(facts, Mapping): + return {} + wanted = set(agent_ids) + return { + agent_id: dict(row) + for agent_id, row in facts.items() + if agent_id in wanted and isinstance(row, Mapping) + } diff --git a/loopx/control_plane/agents/management_projection.py b/loopx/control_plane/agents/management_projection.py index b8f9e1d64..021d055b1 100644 --- a/loopx/control_plane/agents/management_projection.py +++ b/loopx/control_plane/agents/management_projection.py @@ -29,7 +29,6 @@ MAX_SESSION_BINDING_CANDIDATES = 3 MAX_WORKSPACE_SCOPES = 4 STALE_CLAIM_THRESHOLD_HOURS = 36 -EXECUTING_ACTIVITY_THRESHOLD_HOURS = 8 MATERIAL_LIFECYCLE_CAPABILITY = "material_lifecycle" _TODO_GROUP_LIST_KEYS = tuple( @@ -189,6 +188,42 @@ def _registered_agents(status_payload: dict[str, Any]) -> dict[str, dict[str, An return rows +def projected_agent_goals( + status_payload: dict[str, Any], +) -> dict[str, dict[str, list[str]]]: + """Return ``goal_id -> {agent_id -> [spellings]}`` for the agents this view rows. + + The execution facts collector reads one Turn lane per Goal and agent, so it + needs the same agent set this projection rows: the Goal's registered agents + and every Todo claimant. Registered ids keep their registry spelling next to + the normalized one because the Turn envelope names the lane with the former + while Todo claims carry the latter. + """ + + goals: dict[str, dict[str, list[str]]] = {} + + def add(goal_id: str | None, raw_agent: Any) -> None: + agent_id = normalize_todo_claimed_by(raw_agent) + if not goal_id or not agent_id: + return + spellings = goals.setdefault(goal_id, {}).setdefault(agent_id, [agent_id]) + raw = str(raw_agent or "").strip() + if raw and raw not in spellings: + spellings.append(raw) + + run_history = _as_dict(status_payload.get("run_history")) + for goal in _as_list(run_history.get("goals")): + if not isinstance(goal, dict): + continue + goal_id = _compact(goal.get("id"), limit=180) + for raw_agent in _as_list(_as_dict(goal.get("coordination")).get("registered_agents")): + add(goal_id, raw_agent) + for todo in _iter_status_todos(status_payload): + if not _is_done(todo): + add(_compact(todo.get("goal_id"), limit=180), _todo_agent_id(todo)) + return goals + + def _agent_material_frontiers( status_payload: dict[str, Any], ) -> dict[tuple[str, str], dict[str, Any]]: @@ -474,23 +509,29 @@ def _agent_state( *, current: dict[str, Any] | None = None, has_session_binding: bool = False, - last_activity_at: str | None = None, + execution: dict[str, Any] | None = None, ) -> str: """Derive the worker lifecycle state from existing facts only. The state is a projection over registry membership, todo claims, - session bindings, and activity timestamps. It does not introduce a - second source of truth: every input is already owned by another - contract (registry, todo, session binding, or run history). + session bindings and execution facts. It does not introduce a second + source of truth: every input is already owned by another contract + (registry, todo, session binding, Turn lane, delegation lock, lease). State priority (highest first): 1. blocked — current todo is blocked or a blocker 2. monitoring / waiting — monitor-only or non-open current work - 3. executing — current open work updated within the activity threshold - 4. bound — has session binding and active todo - 5. launchable — has active todo, no session binding - 6. addressable — has session binding but no active todo - 7. registered — registered in registry, no binding or todo + 3. executing — the Turn lane is live or a delegation worker holds its lock + 4. unknown — the lane holder cannot be checked here (foreign host or + unreadable record), or an active lease has expired while + nothing is live; consumers fail closed on it + 5. bound — has session binding and active todo + 6. launchable — has active todo, no session binding + 7. addressable — has session binding but no active todo + 8. registered — registered in registry, no binding or todo + + A Todo timestamp never makes a worker `executing`: activity age stays in + `last_activity_at` and `stale_claim_hint`, where it is a hint, not liveness. """ open_todos = [todo for todo in todos if not _is_done(todo)] @@ -508,13 +549,15 @@ def _agent_state( if _todo_status(current) != "open": return "waiting" - # Activity describes the selected work, not updates to unrelated todos. - if current and not _is_done(current) and last_activity_at: - parsed = parse_timestamp(last_activity_at) - if parsed: - age_hours = (now_utc() - parsed).total_seconds() / 3600 - if 0 <= age_hours <= EXECUTING_ACTIVITY_THRESHOLD_HOURS: - return WORKER_LIFECYCLE_STATE_EXECUTING + facts = execution or {} + lane = str(facts.get("lane") or "") + if lane == EXECUTION_LANE_LIVE or facts.get("delegation_worker_active") is True: + return WORKER_LIFECYCLE_STATE_EXECUTING + lease = _as_dict(facts.get("lease")) + if lane in EXECUTION_LANE_UNKNOWN_STATES or ( + lease.get("status") == "active" and lease.get("expired") is True + ): + return WORKER_LIFECYCLE_STATE_UNKNOWN # Bound: has session binding and active todo. if current and not _is_done(current) and has_session_binding: @@ -535,13 +578,48 @@ def _agent_state( # Worker lifecycle state vocabulary for R2 small-team execution qualification. # These states are derived from existing facts only; they do not introduce a # second source of truth. The projection reads registry membership, todo -# claims, session bindings, and activity timestamps — nothing else. +# claims, session bindings and execution facts (Turn lane liveness, delegation +# worker locks, task leases) -- nothing else. WORKER_LIFECYCLE_STATE_REGISTERED = "registered" WORKER_LIFECYCLE_STATE_ADDRESSABLE = "addressable" WORKER_LIFECYCLE_STATE_BOUND = "bound" WORKER_LIFECYCLE_STATE_LAUNCHABLE = "launchable" WORKER_LIFECYCLE_STATE_EXECUTING = "executing" WORKER_LIFECYCLE_STATE_BLOCKED = "blocked" +WORKER_LIFECYCLE_STATE_UNKNOWN = "unknown" + +# Execution-fact lane states this projection switches on. `live` is the only +# evidence of execution; the two unknowns are lanes whose holder this machine +# cannot vouch for, so the row must not read as idle either. +EXECUTION_LANE_LIVE = "live" +EXECUTION_LANE_UNKNOWN_STATES = frozenset({"foreign_host", "unreadable"}) + + +def _execution_row(facts: Any) -> dict[str, Any] | None: + """Project one agent's execution facts as read-only evidence for its state.""" + + if not isinstance(facts, dict): + return None + row: dict[str, Any] = {} + lane = _compact(facts.get("lane"), limit=40) + if lane: + row["lane"] = lane + holder = _as_dict(facts.get("lane_holder")) + if holder: + row["lane_holder"] = { + key: holder[key] + for key in ("host", "pid", "acquired_at") + if holder.get(key) not in (None, "") + } + if facts.get("delegation_worker_active") is True: + row["delegation_worker_active"] = True + lease = _as_dict(facts.get("lease")) + lease_status = _compact(lease.get("status"), limit=40) + if lease_status: + row["lease"] = {"status": lease_status} + if isinstance(lease.get("expired"), bool): + row["lease"]["expired"] = lease["expired"] + return row or None def _last_activity(todos: list[dict[str, Any]]) -> str | None: @@ -594,12 +672,18 @@ def build_agent_management_projection( status_payload: dict[str, Any], *, available_capabilities: Any = None, + execution_facts: Any = None, ) -> dict[str, Any]: """Build a read-only agent management view from a status payload. This is a projection over existing LoopX status/todo/history state. It does not allocate tasks, dispatch agents, reclaim stale claims, or expose write actions. + + ``execution_facts`` maps agent id to the facts `collect_agent_execution_facts` + reads (Turn lane liveness, delegation worker lock, task lease). Without them + no row can be `executing` or `unknown`: the projection then says what the + durable state proves and nothing more. """ rows_by_agent = _registered_agents(status_payload) @@ -619,6 +703,7 @@ def build_agent_management_projection( # keep a bounded candidate summary plus the full count, so no consumer can # read the surviving row as the only route to that peer. session_binding_candidates = _collect_session_binding_candidates(status_payload) + facts_by_agent = execution_facts if isinstance(execution_facts, dict) else {} seen_todos: set[tuple[str, str, str, str]] = set() for todo in _iter_status_todos(status_payload): @@ -667,11 +752,12 @@ def build_agent_management_projection( if ref not in handoff_refs: handoff_refs.append(ref) last_activity = _last_activity(todos) + execution = _execution_row(facts_by_agent.get(agent_id)) agent_state = _agent_state( all_todos, current=current, has_session_binding=agent_id in session_binding_candidates, - last_activity_at=_last_activity([current]) if current else None, + execution=execution, ) agent_row: dict[str, Any] = { "agent_id": agent_id, @@ -680,6 +766,7 @@ def build_agent_management_projection( "current_todo": _todo_row(current) if current else None, "next_action": _safe_next_action(current), "last_activity_at": last_activity, + "execution": execution, "evidence_refs": evidence_refs[:MAX_REFS], "handoff_refs": handoff_refs[:MAX_REFS], "goal_ids": _as_list(raw_row.get("_goal_ids"))[:MAX_REFS], @@ -739,6 +826,11 @@ def build_agent_management_projection( } if material_frontiers: source_summary["material_frontier_count"] = len(material_frontiers) + if isinstance(execution_facts, dict): + source_summary["execution_fact_source"] = ( + "turn lane holder records, delegation worker locks, task leases" + ) + source_summary["execution_fact_agent_count"] = len(facts_by_agent) projection: dict[str, Any] = { "schema_version": AGENT_MANAGEMENT_PROJECTION_SCHEMA_VERSION, diff --git a/loopx/control_plane/status/collection.py b/loopx/control_plane/status/collection.py index 3bca96593..78f1efc77 100644 --- a/loopx/control_plane/status/collection.py +++ b/loopx/control_plane/status/collection.py @@ -8,6 +8,7 @@ from pathlib import Path from typing import Any, Callable +from ..agents.execution_facts import collect_agent_execution_facts from ..coordination.local_authority import CanonicalTodoSnapshot from ..goals.acceptance_observation import attach_goal_acceptance_observations from ..goals.artifact_lifecycle import attach_goal_artifact_lifecycle_projections @@ -230,9 +231,18 @@ def collect_status( "registry_revision": registry_activation_revision(registry), } payload["runtime_projection_routes"] = runtime_projection_route_health + # Lane liveness, delegation worker locks and leases are what make a worker + # `executing` or `unknown`; they travel with the payload so re-projections + # such as the peer directory read the same facts this projection did. + execution_facts = collect_agent_execution_facts( + runtime_root=runtime_root, status_payload=payload + ) + if execution_facts: + payload["agent_execution_facts"] = execution_facts agent_management_projection = context.build_agent_management_projection( payload, available_capabilities=available_capabilities, + execution_facts=execution_facts, ) if agent_management_projection.get("agents"): payload["agent_management_projection"] = agent_management_projection diff --git a/tests/control_plane/test_agent_lifecycle_state.py b/tests/control_plane/test_agent_lifecycle_state.py index 2f82553c3..7be140b07 100644 --- a/tests/control_plane/test_agent_lifecycle_state.py +++ b/tests/control_plane/test_agent_lifecycle_state.py @@ -1,30 +1,48 @@ -"""Exercise worker states through the public projection and real peer admission.""" +"""Exercise worker states through the public projection and real peer admission. + +`executing` and `unknown` come from execution facts (Turn lane liveness, +delegation worker locks, task leases), never from a Todo timestamp. The facts +fixtures here are real: a lane held in-process, a holder record edited the +way a crash or another machine leaves it, a lease file, a delegation row lock. +""" from datetime import datetime, timedelta, timezone +import json +import os +from pathlib import Path import pytest from loopx.control_plane.agents import management_projection as projection +from loopx.control_plane.agents.execution_facts import collect_agent_execution_facts +from loopx.control_plane.collaboration.inbox import _root as manager_context_root from loopx.control_plane.quota.task_orchestration import apply_task_orchestration_contract +from loopx.control_plane.turn_driver.lane_fence import turn_lane_singleflight, turn_lane_target +from loopx.file_lock import exclusive_file_lock, lock_holder_path NOW = datetime(2026, 9, 18, 12, tzinfo=timezone.utc) +GOAL = "test-goal" def build_projection(monkeypatch, *, age=None, binding=False, status="open", - task_class="advancement_task", has_todo=True, extra_todos=()): + task_class="advancement_task", has_todo=True, extra_todos=(), + facts=None, runtime_root=None): monkeypatch.setattr(projection, "now_utc", lambda: NOW) - todo = {"todo_id": "todo_peer", "goal_id": "test-goal", "role": "agent", + todo = {"todo_id": "todo_peer", "goal_id": GOAL, "role": "agent", "claimed_by": "peer", "status": status, "task_class": task_class, "action_kind": "inspect", "text": "Inspect the public contract."} if age is not None: todo["updated_at"] = (NOW - timedelta(hours=age)).isoformat() todos = ([todo] if has_todo else []) + list(extra_todos) - payload = {"goal_filter": "test-goal", "run_history": {"goals": [{ - "id": "test-goal", "coordination": { + payload = {"goal_filter": GOAL, "run_history": {"goals": [{ + "id": GOAL, "coordination": { "registered_agents": ["peer"], "thread_agent_bindings": [{"agent_id": "peer", "thread_id": "thread-peer", "host_surface": "codex-app"}] if binding else [], }}]}, "todo_index": {"items": todos}} - return projection.build_agent_management_projection(payload), todo + if runtime_root is not None: + facts = collect_agent_execution_facts(runtime_root=runtime_root, status_payload=payload) + execution_facts = {"peer": facts} if facts is not None and runtime_root is None else facts + return projection.build_agent_management_projection(payload, execution_facts=execution_facts), todo @pytest.mark.parametrize("kwargs,expected", [ @@ -34,9 +52,11 @@ def build_projection(monkeypatch, *, age=None, binding=False, status="open", ({"status": "done", "binding": True}, "addressable"), ({}, "launchable"), ({"binding": True}, "bound"), - ({"age": 1}, "executing"), - ({"age": 8}, "executing"), - ({"age": 8.01}, "launchable"), + # A fresh Todo update is activity, not execution: nothing runs behind it. + ({"age": 0}, "launchable"), + ({"age": 1}, "launchable"), + ({"age": 1, "binding": True}, "bound"), + ({"age": 8}, "launchable"), ({"age": 9, "binding": True}, "bound"), ({"age": -1}, "launchable"), ({"age": 48}, "launchable"), @@ -45,6 +65,23 @@ def build_projection(monkeypatch, *, age=None, binding=False, status="open", ({"task_class": "continuous_monitor", "age": 1}, "monitoring"), ({"task_class": "continuous_monitor", "age": 48}, "monitoring"), ({"status": "deferred", "age": 1}, "waiting"), + # Execution facts decide `executing` and `unknown`, whatever the Todo age. + ({"age": 30, "facts": {"lane": "live"}}, "executing"), + ({"age": 30, "facts": {"lane": "absent", "delegation_worker_active": True}}, "executing"), + ({"age": 0, "facts": {"lane": "dead"}}, "launchable"), + ({"age": 0, "binding": True, "facts": {"lane": "dead"}}, "bound"), + ({"age": 0, "facts": {"lane": "released"}}, "launchable"), + ({"age": 0, "facts": {"lane": "foreign_host"}}, "unknown"), + ({"age": 0, "facts": {"lane": "unreadable"}}, "unknown"), + ({"age": 0, "facts": {"lane": "absent", "lease": {"status": "active", "expired": True}}}, "unknown"), + ({"age": 0, "facts": {"lane": "absent", "lease": {"status": "active", "expired": False}}}, "launchable"), + ({"age": 0, "facts": {"lane": "absent", "lease": {"status": "released"}}}, "launchable"), + ({"age": 0, "facts": {"lane": "live", "lease": {"status": "active", "expired": True}}}, "executing"), + ({"age": 0, "facts": {"lane": "foreign_host", "delegation_worker_active": True}}, "executing"), + ({"has_todo": False, "facts": {"lane": "live"}}, "executing"), + ({"has_todo": False, "facts": {"lane": "foreign_host"}}, "unknown"), + ({"status": "blocked", "facts": {"lane": "live"}}, "blocked"), + ({"task_class": "continuous_monitor", "facts": {"lane": "live"}}, "monitoring"), ]) def test_projected_state(monkeypatch, kwargs, expected): packet, _ = build_projection(monkeypatch, **kwargs) @@ -52,11 +89,27 @@ def test_projected_state(monkeypatch, kwargs, expected): assert row["state"] == expected assert "lifecycle_state" not in row assert ("session_binding_candidates" in row) == kwargs.get("binding", False) + assert ("execution" in row) == ("facts" in kwargs) assert packet["truth_contract"]["projection_is_writable"] is False + assert not hasattr(projection, "EXECUTING_ACTIVITY_THRESHOLD_HOURS") + + +def test_execution_facts_are_projected_as_evidence_not_authority(monkeypatch): + facts = {"lane": "foreign_host", "lane_holder": {"host": "elsewhere", "pid": 7, "acquired_at": "2026-09-18T11:00:00Z"}, + "lease": {"status": "active", "expired": True}, "delegation_worker_active": False} + packet, _ = build_projection(monkeypatch, age=0, facts=facts) + row = packet["agents"][0] + assert row["state"] == "unknown" + assert row["execution"] == {"lane": "foreign_host", "lane_holder": facts["lane_holder"], + "lease": {"status": "active", "expired": True}} + assert packet["source_summary"]["execution_fact_agent_count"] == 1 + assert "execution_fact_source" in packet["source_summary"] + without, _ = build_projection(monkeypatch, age=0) + assert "execution_fact_agent_count" not in without["source_summary"] def test_unrelated_blocked_activity_does_not_change_current_work(monkeypatch): - other = {"todo_id": "todo_blocked", "goal_id": "test-goal", "role": "agent", + other = {"todo_id": "todo_blocked", "goal_id": GOAL, "role": "agent", "claimed_by": "peer", "status": "blocked", "task_class": "blocker", "updated_at": NOW.isoformat()} packet, _ = build_projection(monkeypatch, extra_todos=[other]) @@ -66,19 +119,149 @@ def test_unrelated_blocked_activity_does_not_change_current_work(monkeypatch): assert row["state"] == "launchable" -@pytest.mark.parametrize("age,binding,capability,resume_ready,reason", [ - (1, True, True, True, None), - (1, False, True, True, None), - (9, True, True, True, None), - (9, False, True, True, None), - (48, True, True, True, "peer_runtime_stale"), - (None, True, True, True, "peer_runtime_stale"), - (1, True, False, True, "peer_agent_activation_unavailable"), - (1, True, True, False, "peer_lane_not_resume_ready"), +# --- facts collected from the real runtime layout --------------------------- + +def _lane(runtime_root: Path) -> Path: + return turn_lane_target(runtime_root=runtime_root, goal_id=GOAL, + plan={"turn_envelope": {"agent_id": "peer"}}) + + +def _rewrite_holder(target: Path, **changes) -> None: + holder_path = lock_holder_path(target) + record = json.loads(holder_path.read_text(encoding="utf-8")) + record.pop("released_at", None) + record.update(changes) + holder_path.write_text(json.dumps(record), encoding="utf-8") + + +def _dead_pid() -> int: + pid = os.getpid() + while True: + pid += 1 + try: + os.kill(pid, 0) + except ProcessLookupError: + return pid + except OSError: + continue + + +def _write_lease(runtime_root: Path, *, expires_in_hours: float, status: str = "active") -> None: + now = datetime.now(timezone.utc) + path = runtime_root / "goals" / GOAL / "task-leases" / "todo_peer.json" + path.parent.mkdir(parents=True, exist_ok=True) + path.write_text(json.dumps({ + "schema_version": "task_lease_v0", "goal_id": GOAL, "todo_id": "todo_peer", "owner": "peer", + "idempotency_key": "execution-peer", "version": 1, "lease_epoch": 1, "status": status, + "write_scopes": ["src/**"], "acquire_ttl_seconds": 600, "acquired_at": now.isoformat(), + "updated_at": now.isoformat(), "expires_at": (now + timedelta(hours=expires_in_hours)).isoformat(), + }), encoding="utf-8") + + +def test_a_lane_held_in_process_makes_a_stale_todo_executing(monkeypatch, tmp_path): + runtime_root = tmp_path / "runtime" + plan = {"turn_envelope": {"agent_id": "peer"}} + with turn_lane_singleflight(runtime_root=runtime_root, goal_id=GOAL, plan=plan) as held: + assert held is not None + packet, _ = build_projection(monkeypatch, age=30, runtime_root=runtime_root) + row = packet["agents"][0] + assert row["state"] == "executing" + assert row["execution"]["lane"] == "live" + assert row["execution"]["lane_holder"]["pid"] == os.getpid() + assert "stale_claim_hint" not in row + # Once the Turn settles cleanly the same stale Todo is only launchable. + after, _ = build_projection(monkeypatch, age=30, runtime_root=runtime_root) + assert after["agents"][0]["state"] == "launchable" + assert after["agents"][0]["execution"]["lane"] == "released" + + +def test_a_dead_holder_behind_a_fresh_todo_is_not_executing(monkeypatch, tmp_path): + runtime_root = tmp_path / "runtime" + plan = {"turn_envelope": {"agent_id": "peer"}} + with turn_lane_singleflight(runtime_root=runtime_root, goal_id=GOAL, plan=plan): + pass + _rewrite_holder(_lane(runtime_root), pid=_dead_pid()) + packet, _ = build_projection(monkeypatch, age=0, runtime_root=runtime_root) + assert packet["agents"][0]["state"] == "launchable" + assert packet["agents"][0]["execution"]["lane"] == "dead" + bound, _ = build_projection(monkeypatch, age=0, binding=True, runtime_root=runtime_root) + assert bound["agents"][0]["state"] == "bound" + + +def test_a_foreign_host_holder_is_unknown(monkeypatch, tmp_path): + runtime_root = tmp_path / "runtime" + plan = {"turn_envelope": {"agent_id": "peer"}} + with turn_lane_singleflight(runtime_root=runtime_root, goal_id=GOAL, plan=plan): + pass + _rewrite_holder(_lane(runtime_root), host="another-machine") + packet, _ = build_projection(monkeypatch, age=0, binding=True, runtime_root=runtime_root) + row = packet["agents"][0] + assert row["state"] == "unknown" + assert row["execution"]["lane"] == "foreign_host" + assert row["execution"]["lane_holder"]["host"] == "another-machine" + + +def test_an_expired_active_lease_with_nothing_live_is_unknown(monkeypatch, tmp_path): + runtime_root = tmp_path / "runtime" + _write_lease(runtime_root, expires_in_hours=-1) + packet, _ = build_projection(monkeypatch, age=0, runtime_root=runtime_root) + row = packet["agents"][0] + assert row["state"] == "unknown" + assert row["execution"] == {"lane": "absent", "lease": {"status": "active", "expired": True}} + _write_lease(runtime_root, expires_in_hours=1) + fresh, _ = build_projection(monkeypatch, age=0, runtime_root=runtime_root) + assert fresh["agents"][0]["state"] == "launchable" + assert fresh["agents"][0]["execution"]["lease"] == {"status": "active", "expired": False} + _write_lease(runtime_root, expires_in_hours=-1, status="released") + released, _ = build_projection(monkeypatch, age=0, runtime_root=runtime_root) + assert released["agents"][0]["state"] == "launchable" + assert released["agents"][0]["execution"]["lease"] == {"status": "released"} + + +def test_a_live_delegation_worker_lock_is_executing(monkeypatch, tmp_path): + runtime_root = tmp_path / "runtime" + row_path = manager_context_root(runtime_root) / "executions" / ("a" * 64) / ("b" * 64 + ".json") + row_path.parent.mkdir(parents=True) + row_path.write_text(json.dumps({"identity": {"binding": {"id": "b1", "agent_id": "peer", "todo_id": "todo_peer"}, + "request_id": "r1", "operation_id": "op-1"}, + "status": "running"}), encoding="utf-8") + with exclusive_file_lock(row_path): + packet, _ = build_projection(monkeypatch, age=30, runtime_root=runtime_root) + assert packet["agents"][0]["state"] == "executing" + assert packet["agents"][0]["execution"] == {"lane": "absent", "delegation_worker_active": True} + # The worker exited: its released lock is no evidence of execution. + settled, _ = build_projection(monkeypatch, age=30, runtime_root=runtime_root) + assert settled["agents"][0]["state"] == "launchable" + assert settled["agents"][0]["execution"] == {"lane": "absent"} + + +def test_no_runtime_root_means_no_facts_not_no_execution(monkeypatch): + payload = {"goal_filter": GOAL, "run_history": {"goals": [{"id": GOAL, "coordination": {"registered_agents": ["peer"]}}]}} + assert collect_agent_execution_facts(runtime_root=None, status_payload=payload) == {} + packet, _ = build_projection(monkeypatch, age=0) + assert "execution" not in packet["agents"][0] + + +# --- real projection to real peer admission ---------------------------------- + +@pytest.mark.parametrize("age,binding,facts,capability,resume_ready,reason", [ + (1, True, None, True, True, None), + (1, False, None, True, True, None), + (9, True, None, True, True, None), + (9, False, None, True, True, None), + (30, False, {"lane": "live"}, True, True, None), + (48, True, None, True, True, "peer_runtime_stale"), + (None, True, None, True, True, "peer_runtime_stale"), + (1, True, None, False, True, "peer_agent_activation_unavailable"), + (1, True, None, True, False, "peer_lane_not_resume_ready"), + # Unknown liveness fails closed: no activation on a peer this machine cannot vouch for. + (1, True, {"lane": "foreign_host"}, True, True, "peer_runtime_not_active"), + (1, False, {"lane": "unreadable"}, True, True, "peer_runtime_not_active"), + (1, True, {"lane": "absent", "lease": {"status": "active", "expired": True}}, True, True, "peer_runtime_not_active"), ]) -def test_real_projection_to_peer_admission(monkeypatch, age, binding, capability, +def test_real_projection_to_peer_admission(monkeypatch, age, binding, facts, capability, resume_ready, reason): - packet, todo = build_projection(monkeypatch, age=age, binding=binding) + packet, todo = build_projection(monkeypatch, age=age, binding=binding, facts=facts) todo.update(resume_when="todo_done:todo_dependency", resume_ready=resume_ready) summary = {"items": [todo]} contract, lane = apply_task_orchestration_contract( From 442c8d12a670a82e8905cdfb58f4c851b9e3cd2f Mon Sep 17 00:00:00 2001 From: song <22676124+songoow@users.noreply.github.com> Date: Tue, 29 Sep 2026 10:06:35 -0400 Subject: [PATCH 3/9] fix(agents): carry execution facts per row and fail closed on unknown peers Execution facts now travel on each agent row as `execution`, collected for the Goal's registered agents and open-Todo claimants together with the leases on those open Todos, so the peer directory re-projects the same facts instead of a side payload. Status names the collected sources as a typed `execution_facts` summary. Peer activation admission is an allow-list: `executing` requires a live Turn lane or delegation worker, and `unknown` (holder on another host, unreadable record, expired lease with nothing live) or any unlisted state is not active. Co-Authored-By: Claude Opus 5.5 (1M context) Signed-off-by: song <22676124+songoow@users.noreply.github.com> --- loopx/control_plane/agents/directory.py | 30 ++- loopx/control_plane/agents/execution_facts.py | 205 ++++++++++-------- .../agents/management_projection.py | 53 +++-- .../control_plane/quota/peer_orchestration.ts | 4 + loopx/control_plane/status/collection.py | 6 +- .../test_agent_lifecycle_state.py | 66 ++++-- .../peer_orchestration.test.ts | 15 ++ 7 files changed, 239 insertions(+), 140 deletions(-) diff --git a/loopx/control_plane/agents/directory.py b/loopx/control_plane/agents/directory.py index ecca13aa1..49a8099d7 100644 --- a/loopx/control_plane/agents/directory.py +++ b/loopx/control_plane/agents/directory.py @@ -104,6 +104,28 @@ def _publishable_route(route: dict[str, Any]) -> tuple[dict[str, Any], int]: return published, withheld +def _projected_execution_facts(payload: Mapping[str, Any]) -> dict[str, Any] | None: + """Recover the execution facts status collection projected onto agent rows. + + ``None`` when the payload's projection says no facts were collected, so an + older or fact-less payload keeps `lease_state_not_projected` instead of + reading as "leases projected, none found". + """ + + projection = _as_mapping(payload.get("agent_management_projection")) + summary = _as_mapping(_as_mapping(projection.get("source_summary")).get("execution_facts")) + if summary.get("collected") is not True: + return None + facts: dict[str, Any] = {} + for row in _as_list(projection.get("agents")): + row = _as_mapping(row) + agent_id = str(row.get("agent_id") or "").strip() + execution = row.get("execution") + if agent_id and isinstance(execution, Mapping): + facts[agent_id] = dict(execution) + return facts + + def _work_block(agent_row: Mapping[str, Any]) -> dict[str, Any] | None: """Project one Agent's bounded work facts, or nothing when it holds none.""" @@ -188,16 +210,16 @@ def build_peer_agent_directory( the packet records that the caller identity was not supplied instead of inventing one. - `execution_facts` defaults to the `agent_execution_facts` status collection - attached to the payload, so this re-projection reads the same lane, worker - and lease facts the management projection did. + `execution_facts` defaults to the facts the payload's management projection + already carries on its rows, so this re-projection reads the same lane, + worker and lease facts status collection did instead of reading them again. """ payload = status_payload if isinstance(status_payload, Mapping) else {} resolved_goal = _compact(goal_id or payload.get("goal_filter"), limit=120) caller = _compact(caller_agent_id, limit=120) if execution_facts is None: - execution_facts = payload.get("agent_execution_facts") + execution_facts = _projected_execution_facts(payload) facts = execution_facts if isinstance(execution_facts, Mapping) else None projection = build_agent_management_projection( dict(payload), diff --git a/loopx/control_plane/agents/execution_facts.py b/loopx/control_plane/agents/execution_facts.py index a03f102f7..8af2b3bbb 100644 --- a/loopx/control_plane/agents/execution_facts.py +++ b/loopx/control_plane/agents/execution_facts.py @@ -1,34 +1,36 @@ """Execution facts behind the worker lifecycle projection: read, never lock. `executing` has to be backed by something that runs, not by a fresh Todo -timestamp. This module gathers the facts LoopX already keeps about running -work, keyed by agent id, so the agent management projection can derive -`executing` and `unknown` from them instead of from activity age: +timestamp. This module gathers, per agent, the facts LoopX already keeps about +running work, so the agent management projection can derive `executing` and +`unknown` from them instead of from activity age: - the Turn lane holder record, one per Goal and agent, read through `turn_lane_liveness`; a delegated member executes inside its own lane too; -- the delegation worker's operation lock, whose holder is the detached worker - process executing a delegated Turn; -- the task lease on a claimed Todo, active, expired or released. +- the delegation worker: a delegation journal row whose operation lock and + whose execution slot for the delegated Todo are both held by a live process; +- the task leases on the Goal's open Todos, read from the canonical head after + cutover and from the local lease files before it. It writes nothing, takes no lock and keeps no state of its own: every input is -owned elsewhere and this is a bounded read model over them. Absence of a fact -is reported as absence, never as "not running". +owned elsewhere and this is a bounded read model over them. A fact that cannot +be read is reported as such (`unreadable`, `unavailable`), never as "not +running". """ from __future__ import annotations from collections.abc import Iterable, Mapping +from datetime import datetime from pathlib import Path +import re from typing import Any from ...file_lock import LOCK_HOLDER_LIVE, lock_holder_liveness +from ..collaboration.inbox import _hash as manager_context_hash from ..collaboration.inbox import _read as read_manager_context_record from ..collaboration.inbox import _root as manager_context_root -from ..coordination.local_authority import ( - LocalCoordinationAuthorityUnavailable, - read_canonical_todos_if_promoted, -) +from ..coordination.local_authority import read_canonical_todos_if_promoted from ..runtime.time import now_utc from ..todos.contract import normalize_todo_claimed_by from ..turn_driver.lane_fence import ( @@ -42,7 +44,7 @@ turn_lane_target, ) from ..work_items.local_lease_record import TaskLeaseError, read_lease -from ..work_items.task_lease import lease_expires_at, task_lease_dir +from ..work_items.task_lease import lease_expires_at, normalize_goal_id, task_lease_path from .management_projection import projected_agent_goals # One agent can hold one lane per Goal; the row reports the strongest fact. @@ -58,6 +60,7 @@ LEASE_STATUS_ACTIVE = "active" LEASE_STATUS_UNAVAILABLE = "unavailable" _LANE_HOLDER_FIELDS = ("host", "pid", "acquired_at") +_JOURNAL_ADDRESS = re.compile(r"[a-f0-9]{64}") def _lane_rank(state: str) -> int: @@ -65,7 +68,7 @@ def _lane_rank(state: str) -> int: def _lease_rank(lease: Mapping[str, Any]) -> int: - """Fresher evidence first: an unexpired claim outranks an expired one.""" + """Fresher evidence first: an unexpired lease outranks an expired one.""" status = lease.get("status") if status == LEASE_STATUS_ACTIVE: @@ -88,76 +91,93 @@ def _lane_fact(runtime_root: Path, *, goal_id: str, agent_id: str) -> dict[str, for key in _LANE_HOLDER_FIELDS if liveness["holder"].get(key) not in (None, "") } - if holder and liveness["state"] != TURN_LANE_ABSENT: + if holder: fact["lane_holder"] = holder return fact -def _goal_leases(runtime_root: Path, goal_id: str) -> list[dict[str, Any]] | None: - """Leases at the canonical head after cutover, else the local lease files. +def _delegation_worker_agents( + runtime_root: Path, *, goal_id: str, requesters: Iterable[str] +) -> set[str]: + """Delegated members whose worker is executing their Todo right now. - ``None`` means the lease authority could not be read, which is reported as - ``unavailable`` rather than as "no lease". + A journal row lives under its Goal and requester, so only this Goal's + requesters are read. The row's operation lock is also taken briefly by + readers and by result adoption, so a live operation holder alone is not a + worker; the execution slot for the delegated Todo is taken only by the + worker while it runs, which makes a live slot holder the execution fact. + """ + + store = manager_context_root(runtime_root) + active: set[str] = set() + for requester in dict.fromkeys(requesters): + journal = store / "executions" / manager_context_hash([goal_id, requester]) + try: + rows = sorted(journal.glob("*.json")) if journal.is_dir() else [] + except OSError: + continue + for path in rows: + if not _JOURNAL_ADDRESS.fullmatch(path.stem) or path.is_symlink() or not path.is_file(): + continue + if lock_holder_liveness(path)[0] != LOCK_HOLDER_LIVE: + continue + try: + binding = read_manager_context_record(path)["identity"]["binding"] + member, todo_id = binding["agent_id"], binding["todo_id"] + slot = store / "execution-slots" / manager_context_hash([goal_id, todo_id]) + except (OSError, ValueError, KeyError, TypeError): + continue + agent_id = normalize_todo_claimed_by(member) + if agent_id and lock_holder_liveness(slot)[0] == LOCK_HOLDER_LIVE: + active.add(agent_id) + return active + + +def _open_todo_leases( + runtime_root: Path, *, goal_id: str, todo_ids: set[str] +) -> list[Mapping[str, Any]] | None: + """Leases on the Goal's open Todos; ``None`` when the lease authority is unreadable. + + After cutover the canonical head is the only lease source; its failure is + reported, never replaced by the local files it superseded. The failure + classes are the ones the ownership observation already isolates. """ try: canonical = read_canonical_todos_if_promoted( runtime_root=runtime_root, goal_id=goal_id, include_leases=True ) - except LocalCoordinationAuthorityUnavailable: + except (OSError, RuntimeError, ValueError): return None if canonical is not None: - return [dict(lease) for lease in canonical.get("leases") or []] - lease_dir = task_lease_dir(runtime_root=runtime_root, goal_id=goal_id) - if not lease_dir.is_dir(): - return [] - leases: list[dict[str, Any]] = [] - for path in sorted(lease_dir.glob("todo_*.json")): - try: - lease = read_lease(path) - except TaskLeaseError: - # One corrupt peer lease must not hide the healthy ones. - continue - if lease: - leases.append(lease) - return leases - - -def _lease_fact(lease: Mapping[str, Any], *, at: Any) -> dict[str, Any]: + leases = [lease for lease in canonical.get("leases") or [] if isinstance(lease, Mapping)] + else: + leases = [] + for todo_id in sorted(todo_ids): + try: + lease = read_lease( + task_lease_path(runtime_root=runtime_root, goal_id=goal_id, todo_id=todo_id) + ) + except (TaskLeaseError, OSError): + # One corrupt peer lease must not hide the healthy ones. + continue + if lease: + leases.append(lease) + return [lease for lease in leases if lease.get("todo_id") in todo_ids] + + +def _lease_fact(lease: Mapping[str, Any], *, at: datetime) -> dict[str, Any]: status = str(lease.get("status") or "").strip().lower() or "unknown" fact: dict[str, Any] = {"status": status} if status == LEASE_STATUS_ACTIVE: expires_at = lease_expires_at(dict(lease)) - fact["expired"] = not (expires_at and expires_at > at) + # An active lease whose expiry cannot be read is not proven current. + fact["expired"] = not (expires_at is not None and expires_at > at) return fact -def _delegation_worker_agents(runtime_root: Path) -> set[str]: - """Agents whose delegation worker process still holds its operation lock.""" - - executions = manager_context_root(runtime_root) / "executions" - if not executions.is_dir(): - return set() - active: set[str] = set() - for path in sorted(executions.glob("*/*.json")): - if path.is_symlink() or not path.is_file(): - continue - state, _record = lock_holder_liveness(path) - if state != LOCK_HOLDER_LIVE: - continue - try: - row = read_manager_context_record(path) - raw_agent = row["identity"]["binding"]["agent_id"] - except (OSError, ValueError, KeyError, TypeError): - continue - agent_id = normalize_todo_claimed_by(raw_agent) or str(raw_agent or "").strip() - if agent_id: - active.add(agent_id) - return active - - def _merge_lane(row: dict[str, Any], fact: Mapping[str, Any]) -> None: - if "lane" not in row or _lane_rank(fact["lane"]) < _lane_rank(row["lane"]): + if _lane_rank(fact["lane"]) < _lane_rank(row["lane"]): row.pop("lane_holder", None) row.update(fact) @@ -169,49 +189,46 @@ def _merge_lease(row: dict[str, Any], fact: Mapping[str, Any]) -> None: def collect_agent_execution_facts( *, runtime_root: Path | str | None, status_payload: Mapping[str, Any] -) -> dict[str, dict[str, Any]]: +) -> dict[str, dict[str, Any]] | None: """Return ``agent_id -> {lane, lane_holder?, delegation_worker_active, lease?}``. - Agents and Goals are the ones the management projection would row: the - Goal's registered agents plus every Todo claimant. Without a runtime root - there are no facts to read and the result is empty, which the projection - treats as "facts not collected", not as "nobody is running". + Agents and Goals are the ones the management projection rows: each Goal's + registered agents plus its open-Todo claimants. ``None`` means no runtime + root could be read, which the projection reports as facts not collected, + never as "nobody is running". """ if runtime_root is None: - return {} + return None root = Path(runtime_root) if not root.is_dir(): - return {} + return None facts: dict[str, dict[str, Any]] = {} at = now_utc() - for goal_id, agents in projected_agent_goals(status_payload).items(): - leases = _goal_leases(root, goal_id) if agents else [] - for agent_id, spellings in agents.items(): - row = facts.setdefault(agent_id, {"delegation_worker_active": False}) - for spelling in spellings: + for raw_goal_id, agents in projected_agent_goals(dict(status_payload)).items(): + try: + goal_id = normalize_goal_id(raw_goal_id) + except TaskLeaseError: + continue + workers = _delegation_worker_agents( + root, + goal_id=goal_id, + requesters=(spelling for work in agents.values() for spelling in work["spellings"]), + ) + todo_ids = {todo_id for work in agents.values() for todo_id in work["open_todo_ids"]} + leases = _open_todo_leases(root, goal_id=goal_id, todo_ids=todo_ids) if todo_ids else [] + for agent_id, work in agents.items(): + row = facts.setdefault( + agent_id, {"lane": TURN_LANE_ABSENT, "delegation_worker_active": False} + ) + for spelling in work["spellings"]: _merge_lane(row, _lane_fact(root, goal_id=goal_id, agent_id=spelling)) - if leases is None: + if agent_id in workers: + row["delegation_worker_active"] = True + if leases is None and work["open_todo_ids"]: _merge_lease(row, {"status": LEASE_STATUS_UNAVAILABLE}) for lease in leases or []: owner = normalize_todo_claimed_by(lease.get("owner")) - if owner in facts: + if owner in agents: _merge_lease(facts[owner], _lease_fact(lease, at=at)) - for agent_id in _delegation_worker_agents(root): - facts.setdefault(agent_id, {"lane": TURN_LANE_ABSENT})["delegation_worker_active"] = True return facts - - -def execution_fact_rows( - facts: Mapping[str, Mapping[str, Any]] | None, agent_ids: Iterable[str] -) -> dict[str, dict[str, Any]]: - """Narrow a facts map to the named agents, for agent-lane compaction.""" - - if not isinstance(facts, Mapping): - return {} - wanted = set(agent_ids) - return { - agent_id: dict(row) - for agent_id, row in facts.items() - if agent_id in wanted and isinstance(row, Mapping) - } diff --git a/loopx/control_plane/agents/management_projection.py b/loopx/control_plane/agents/management_projection.py index 021d055b1..725a4ec3a 100644 --- a/loopx/control_plane/agents/management_projection.py +++ b/loopx/control_plane/agents/management_projection.py @@ -190,37 +190,45 @@ def _registered_agents(status_payload: dict[str, Any]) -> dict[str, dict[str, An def projected_agent_goals( status_payload: dict[str, Any], -) -> dict[str, dict[str, list[str]]]: - """Return ``goal_id -> {agent_id -> [spellings]}`` for the agents this view rows. - - The execution facts collector reads one Turn lane per Goal and agent, so it - needs the same agent set this projection rows: the Goal's registered agents - and every Todo claimant. Registered ids keep their registry spelling next to - the normalized one because the Turn envelope names the lane with the former - while Todo claims carry the latter. +) -> dict[str, dict[str, dict[str, list[str]]]]: + """Return ``goal_id -> agent_id -> {spellings, open_todo_ids}`` for this view's agents. + + The execution facts collector reads one Turn lane per Goal and agent and the + leases on open Todos, so it needs the same agents this projection rows: each + Goal's registered agents and its open-Todo claimants. Registered ids keep + their registry spelling next to the normalized one, because the Turn + envelope names a lane with the former while Todo claims carry the latter. """ - goals: dict[str, dict[str, list[str]]] = {} + goals: dict[str, dict[str, dict[str, list[str]]]] = {} - def add(goal_id: str | None, raw_agent: Any) -> None: + def add(goal_id: Any, raw_agent: Any, *, todo_id: str | None = None) -> None: + goal = str(goal_id or "").strip() agent_id = normalize_todo_claimed_by(raw_agent) - if not goal_id or not agent_id: + if not goal or not agent_id: return - spellings = goals.setdefault(goal_id, {}).setdefault(agent_id, [agent_id]) + work = goals.setdefault(goal, {}).setdefault( + agent_id, {"spellings": [agent_id], "open_todo_ids": []} + ) raw = str(raw_agent or "").strip() - if raw and raw not in spellings: - spellings.append(raw) + if raw and raw not in work["spellings"]: + work["spellings"].append(raw) + if todo_id and todo_id not in work["open_todo_ids"]: + work["open_todo_ids"].append(todo_id) run_history = _as_dict(status_payload.get("run_history")) for goal in _as_list(run_history.get("goals")): if not isinstance(goal, dict): continue - goal_id = _compact(goal.get("id"), limit=180) for raw_agent in _as_list(_as_dict(goal.get("coordination")).get("registered_agents")): - add(goal_id, raw_agent) + add(goal.get("id"), raw_agent) for todo in _iter_status_todos(status_payload): if not _is_done(todo): - add(_compact(todo.get("goal_id"), limit=180), _todo_agent_id(todo)) + add( + todo.get("goal_id"), + _todo_agent_id(todo), + todo_id=normalize_todo_id(todo.get("todo_id")), + ) return goals @@ -593,6 +601,8 @@ def _agent_state( # cannot vouch for, so the row must not read as idle either. EXECUTION_LANE_LIVE = "live" EXECUTION_LANE_UNKNOWN_STATES = frozenset({"foreign_host", "unreadable"}) +# What `source_summary.execution_facts.sources` names when facts were collected. +EXECUTION_FACT_SOURCES = ("turn_lane_holder", "delegation_worker_lock", "task_lease") def _execution_row(facts: Any) -> dict[str, Any] | None: @@ -827,10 +837,11 @@ def build_agent_management_projection( if material_frontiers: source_summary["material_frontier_count"] = len(material_frontiers) if isinstance(execution_facts, dict): - source_summary["execution_fact_source"] = ( - "turn lane holder records, delegation worker locks, task leases" - ) - source_summary["execution_fact_agent_count"] = len(facts_by_agent) + source_summary["execution_facts"] = { + "collected": True, + "sources": list(EXECUTION_FACT_SOURCES), + "agent_count": len(facts_by_agent), + } projection: dict[str, Any] = { "schema_version": AGENT_MANAGEMENT_PROJECTION_SCHEMA_VERSION, diff --git a/loopx/control_plane/quota/peer_orchestration.ts b/loopx/control_plane/quota/peer_orchestration.ts index 073130c52..5c7d65c25 100644 --- a/loopx/control_plane/quota/peer_orchestration.ts +++ b/loopx/control_plane/quota/peer_orchestration.ts @@ -3,6 +3,10 @@ import type { JsonObject } from "../effect_program.ts"; import { jsonObject, requireJsonObject } from "../runtime_decode.ts"; type ActivationState = "ready" | "blocked"; +// Admission is an allow-list. `executing` is backed by a live Turn lane or +// delegation worker; `bound`/`launchable` are durable work without a process. +// `unknown` (holder on another host, unreadable record, expired lease with +// nothing live) and every unlisted or missing state fail closed as not active. const activeStates = new Set(["running", "monitoring", "executing", "bound", "launchable"]); const rows = (value: unknown): JsonObject[] => Array.isArray(value) ? value.flatMap(item => { const row = jsonObject(item); return row ? [row] : []; }) : []; diff --git a/loopx/control_plane/status/collection.py b/loopx/control_plane/status/collection.py index 78f1efc77..f0359a6c8 100644 --- a/loopx/control_plane/status/collection.py +++ b/loopx/control_plane/status/collection.py @@ -232,13 +232,11 @@ def collect_status( } payload["runtime_projection_routes"] = runtime_projection_route_health # Lane liveness, delegation worker locks and leases are what make a worker - # `executing` or `unknown`; they travel with the payload so re-projections - # such as the peer directory read the same facts this projection did. + # `executing` or `unknown`. Each row carries its facts as `execution`, so a + # re-projection such as the peer directory reads them from there. execution_facts = collect_agent_execution_facts( runtime_root=runtime_root, status_payload=payload ) - if execution_facts: - payload["agent_execution_facts"] = execution_facts agent_management_projection = context.build_agent_management_projection( payload, available_capabilities=available_capabilities, diff --git a/tests/control_plane/test_agent_lifecycle_state.py b/tests/control_plane/test_agent_lifecycle_state.py index 7be140b07..28cb974ce 100644 --- a/tests/control_plane/test_agent_lifecycle_state.py +++ b/tests/control_plane/test_agent_lifecycle_state.py @@ -14,6 +14,7 @@ from loopx.control_plane.agents import management_projection as projection from loopx.control_plane.agents.execution_facts import collect_agent_execution_facts +from loopx.control_plane.collaboration.inbox import _hash as manager_context_hash from loopx.control_plane.collaboration.inbox import _root as manager_context_root from loopx.control_plane.quota.task_orchestration import apply_task_orchestration_contract from loopx.control_plane.turn_driver.lane_fence import turn_lane_singleflight, turn_lane_target @@ -25,7 +26,7 @@ def build_projection(monkeypatch, *, age=None, binding=False, status="open", task_class="advancement_task", has_todo=True, extra_todos=(), - facts=None, runtime_root=None): + facts=None, runtime_root=None, registered=("peer",)): monkeypatch.setattr(projection, "now_utc", lambda: NOW) todo = {"todo_id": "todo_peer", "goal_id": GOAL, "role": "agent", "claimed_by": "peer", "status": status, "task_class": task_class, @@ -35,7 +36,7 @@ def build_projection(monkeypatch, *, age=None, binding=False, status="open", todos = ([todo] if has_todo else []) + list(extra_todos) payload = {"goal_filter": GOAL, "run_history": {"goals": [{ "id": GOAL, "coordination": { - "registered_agents": ["peer"], + "registered_agents": list(registered), "thread_agent_bindings": [{"agent_id": "peer", "thread_id": "thread-peer", "host_surface": "codex-app"}] if binding else [], }}]}, "todo_index": {"items": todos}} @@ -102,10 +103,13 @@ def test_execution_facts_are_projected_as_evidence_not_authority(monkeypatch): assert row["state"] == "unknown" assert row["execution"] == {"lane": "foreign_host", "lane_holder": facts["lane_holder"], "lease": {"status": "active", "expired": True}} - assert packet["source_summary"]["execution_fact_agent_count"] == 1 - assert "execution_fact_source" in packet["source_summary"] + assert packet["source_summary"]["execution_facts"] == { + "collected": True, + "sources": ["turn_lane_holder", "delegation_worker_lock", "task_lease"], + "agent_count": 1, + } without, _ = build_projection(monkeypatch, age=0) - assert "execution_fact_agent_count" not in without["source_summary"] + assert "execution_facts" not in without["source_summary"] def test_unrelated_blocked_activity_does_not_change_current_work(monkeypatch): @@ -218,26 +222,54 @@ def test_an_expired_active_lease_with_nothing_live_is_unknown(monkeypatch, tmp_p assert released["agents"][0]["execution"]["lease"] == {"status": "released"} -def test_a_live_delegation_worker_lock_is_executing(monkeypatch, tmp_path): - runtime_root = tmp_path / "runtime" - row_path = manager_context_root(runtime_root) / "executions" / ("a" * 64) / ("b" * 64 + ".json") - row_path.parent.mkdir(parents=True) +def _delegation_row(runtime_root: Path, *, goal: str, requester: str) -> tuple[Path, Path]: + """A journal row where the delegation service keeps it, plus its execution slot.""" + store = manager_context_root(runtime_root) + row_path = store / "executions" / manager_context_hash([goal, requester]) / ("b" * 64 + ".json") + row_path.parent.mkdir(parents=True, exist_ok=True) row_path.write_text(json.dumps({"identity": {"binding": {"id": "b1", "agent_id": "peer", "todo_id": "todo_peer"}, "request_id": "r1", "operation_id": "op-1"}, "status": "running"}), encoding="utf-8") + return row_path, store / "execution-slots" / manager_context_hash([goal, "todo_peer"]) + + +def _peer_row(packet): + return next(row for row in packet["agents"] if row["agent_id"] == "peer") + + +def test_a_live_delegation_worker_is_executing(monkeypatch, tmp_path): + runtime_root = tmp_path / "runtime" + registered = ("coordinator", "peer") + row_path, slot = _delegation_row(runtime_root, goal=GOAL, requester="coordinator") + with exclusive_file_lock(row_path), exclusive_file_lock(slot): + packet, _ = build_projection(monkeypatch, age=30, runtime_root=runtime_root, registered=registered) + assert _peer_row(packet)["state"] == "executing" + assert _peer_row(packet)["execution"] == {"lane": "absent", "delegation_worker_active": True} + # The worker exited: its released locks are no evidence of execution. + settled, _ = build_projection(monkeypatch, age=30, runtime_root=runtime_root, registered=registered) + assert _peer_row(settled)["state"] == "launchable" + assert _peer_row(settled)["execution"] == {"lane": "absent"} + + +def test_an_operation_lock_without_its_execution_slot_is_not_a_worker(monkeypatch, tmp_path): + """Readers and result adoption also hold the row lock briefly; only the worker holds the slot.""" + runtime_root = tmp_path / "runtime" + registered = ("coordinator", "peer") + row_path, _slot = _delegation_row(runtime_root, goal=GOAL, requester="coordinator") with exclusive_file_lock(row_path): - packet, _ = build_projection(monkeypatch, age=30, runtime_root=runtime_root) - assert packet["agents"][0]["state"] == "executing" - assert packet["agents"][0]["execution"] == {"lane": "absent", "delegation_worker_active": True} - # The worker exited: its released lock is no evidence of execution. - settled, _ = build_projection(monkeypatch, age=30, runtime_root=runtime_root) - assert settled["agents"][0]["state"] == "launchable" - assert settled["agents"][0]["execution"] == {"lane": "absent"} + packet, _ = build_projection(monkeypatch, age=0, runtime_root=runtime_root, registered=registered) + assert _peer_row(packet)["state"] == "launchable" + # A worker executing under another Goal's journal is not this Goal's fact. + other_row, other_slot = _delegation_row(runtime_root, goal="other-goal", requester="coordinator") + with exclusive_file_lock(other_row), exclusive_file_lock(other_slot): + scoped, _ = build_projection(monkeypatch, age=0, runtime_root=runtime_root, registered=registered) + assert _peer_row(scoped)["state"] == "launchable" def test_no_runtime_root_means_no_facts_not_no_execution(monkeypatch): payload = {"goal_filter": GOAL, "run_history": {"goals": [{"id": GOAL, "coordination": {"registered_agents": ["peer"]}}]}} - assert collect_agent_execution_facts(runtime_root=None, status_payload=payload) == {} + assert collect_agent_execution_facts(runtime_root=None, status_payload=payload) is None + assert collect_agent_execution_facts(runtime_root="/nonexistent-loopx-runtime", status_payload=payload) is None packet, _ = build_projection(monkeypatch, age=0) assert "execution" not in packet["agents"][0] diff --git a/tests/control_plane_ts/peer_orchestration.test.ts b/tests/control_plane_ts/peer_orchestration.test.ts index 2d8ac8437..e524cbacc 100644 --- a/tests/control_plane_ts/peer_orchestration.test.ts +++ b/tests/control_plane_ts/peer_orchestration.test.ts @@ -23,6 +23,9 @@ test("runtime, activation and dependency gates remain independently enforced", ( { available_capabilities: input.available_capabilities, agents: [], extra: {}, reasons: ["peer_liveness_unavailable"] }, { available_capabilities: input.available_capabilities, agents: [{ agent_id: "worker", state: "running", stale_claim_hint: true }], extra: {}, reasons: ["peer_runtime_stale"] }, { available_capabilities: input.available_capabilities, agents: [{ agent_id: "worker", state: "dormant" }], extra: {}, reasons: ["peer_runtime_not_active"] }, + // Liveness this machine cannot vouch for is not activation evidence. + { available_capabilities: input.available_capabilities, agents: [{ agent_id: "worker", state: "unknown" }], extra: {}, reasons: ["peer_runtime_not_active"] }, + { available_capabilities: input.available_capabilities, agents: [{ agent_id: "worker" }], extra: {}, reasons: ["peer_runtime_not_active"] }, { available_capabilities: input.available_capabilities, agents: input.agents, extra: { resume_when: "todo_done:dependency", resume_ready: false }, reasons: ["peer_lane_not_resume_ready"] }, ]; for (const row of cases) { @@ -33,6 +36,18 @@ test("runtime, activation and dependency gates remain independently enforced", ( } }); +test("execution-backed and durable-work states are the only admitted states", () => { + for (const state of ["executing", "bound", "launchable", "monitoring", "running"]) { + const result = projectPeerOrchestration({ ...input, agents: [{ agent_id: "worker", state }], items: [task("task")] })!; + assert.equal(result.execution_state, "ready", state); + } + for (const state of ["unknown", "registered", "addressable", "blocked", "waiting", "stale", "scope_wait"]) { + const result = projectPeerOrchestration({ ...input, agents: [{ agent_id: "worker", state }], items: [task("task")] })!; + assert.equal(result.execution_state, "blocked", state); + assert.deepEqual((result.blocked_peer_lanes as any[])[0].reason_codes, ["peer_runtime_not_active"]); + } +}); + test("closed, deferred, unregistered, self and non-advancement rows never grant activation", () => { const items = [task("done", { done: true }), task("blocked", { status: "blocked" }), task("deferred", { status: "deferred" }), task("foreign", { claimed_by: "foreign" }), From b5a611cd1cf61e6370c697584e9c9aba3b14145a Mon Sep 17 00:00:00 2001 From: song <22676124+songoow@users.noreply.github.com> Date: Tue, 29 Sep 2026 10:17:02 -0400 Subject: [PATCH 4/9] chore(census): follow the moved status registry read Collecting execution facts in collect_status moved its registry read. Regenerated with scripts/generate_project_registry_io_manifest.py; the site and its classification are unchanged. Co-Authored-By: Claude Opus 5.5 (1M context) Signed-off-by: song <22676124+songoow@users.noreply.github.com> --- loopx/semantics/project_registry_io_manifest_v1.json | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/loopx/semantics/project_registry_io_manifest_v1.json b/loopx/semantics/project_registry_io_manifest_v1.json index 1d92d2924..0378be2f0 100644 --- a/loopx/semantics/project_registry_io_manifest_v1.json +++ b/loopx/semantics/project_registry_io_manifest_v1.json @@ -1527,7 +1527,7 @@ }, { "site": "loopx/control_plane/status/collection.py::.collect_status::codec_read:load_registry#1", - "line": 98, + "line": 99, "column": 16, "kind": "codec_read", "api": "load_registry", From b4c547fb0f25faeb3fc7825b0d0ba5a9160ae8e1 Mon Sep 17 00:00:00 2001 From: song <22676124+songoow@users.noreply.github.com> Date: Tue, 29 Sep 2026 12:17:38 -0400 Subject: [PATCH 5/9] fix(agents): keep the execution facts summary within the status output budget The agent-facing CLI differential rejected status JSON that grew by 270 characters and 12 lines against a 153-character, 5-line allowance. The growth was the constant source summary, not the per-agent facts: its source list and agent count had no reader, while the peer directory only needs to tell "facts collected, none found" from "no facts collected". The summary is now one `execution_facts_collected` flag, and each row still carries its own `execution` evidence. The budget limits are unchanged. Co-Authored-By: Claude Opus 5.5 (1M context) Signed-off-by: song <22676124+songoow@users.noreply.github.com> --- loopx/control_plane/agents/directory.py | 4 ++-- loopx/control_plane/agents/management_projection.py | 10 +++------- tests/control_plane/test_agent_lifecycle_state.py | 8 ++------ 3 files changed, 7 insertions(+), 15 deletions(-) diff --git a/loopx/control_plane/agents/directory.py b/loopx/control_plane/agents/directory.py index 49a8099d7..1f522ae10 100644 --- a/loopx/control_plane/agents/directory.py +++ b/loopx/control_plane/agents/directory.py @@ -113,8 +113,8 @@ def _projected_execution_facts(payload: Mapping[str, Any]) -> dict[str, Any] | N """ projection = _as_mapping(payload.get("agent_management_projection")) - summary = _as_mapping(_as_mapping(projection.get("source_summary")).get("execution_facts")) - if summary.get("collected") is not True: + summary = _as_mapping(projection.get("source_summary")) + if summary.get("execution_facts_collected") is not True: return None facts: dict[str, Any] = {} for row in _as_list(projection.get("agents")): diff --git a/loopx/control_plane/agents/management_projection.py b/loopx/control_plane/agents/management_projection.py index 725a4ec3a..809921a89 100644 --- a/loopx/control_plane/agents/management_projection.py +++ b/loopx/control_plane/agents/management_projection.py @@ -601,8 +601,6 @@ def _agent_state( # cannot vouch for, so the row must not read as idle either. EXECUTION_LANE_LIVE = "live" EXECUTION_LANE_UNKNOWN_STATES = frozenset({"foreign_host", "unreadable"}) -# What `source_summary.execution_facts.sources` names when facts were collected. -EXECUTION_FACT_SOURCES = ("turn_lane_holder", "delegation_worker_lock", "task_lease") def _execution_row(facts: Any) -> dict[str, Any] | None: @@ -837,11 +835,9 @@ def build_agent_management_projection( if material_frontiers: source_summary["material_frontier_count"] = len(material_frontiers) if isinstance(execution_facts, dict): - source_summary["execution_facts"] = { - "collected": True, - "sources": list(EXECUTION_FACT_SOURCES), - "agent_count": len(facts_by_agent), - } + # One flag: readers only need to tell "facts collected, none found" + # from "no facts collected"; each row carries its own evidence. + source_summary["execution_facts_collected"] = True projection: dict[str, Any] = { "schema_version": AGENT_MANAGEMENT_PROJECTION_SCHEMA_VERSION, diff --git a/tests/control_plane/test_agent_lifecycle_state.py b/tests/control_plane/test_agent_lifecycle_state.py index 28cb974ce..f31bc744b 100644 --- a/tests/control_plane/test_agent_lifecycle_state.py +++ b/tests/control_plane/test_agent_lifecycle_state.py @@ -103,13 +103,9 @@ def test_execution_facts_are_projected_as_evidence_not_authority(monkeypatch): assert row["state"] == "unknown" assert row["execution"] == {"lane": "foreign_host", "lane_holder": facts["lane_holder"], "lease": {"status": "active", "expired": True}} - assert packet["source_summary"]["execution_facts"] == { - "collected": True, - "sources": ["turn_lane_holder", "delegation_worker_lock", "task_lease"], - "agent_count": 1, - } + assert packet["source_summary"]["execution_facts_collected"] is True without, _ = build_projection(monkeypatch, age=0) - assert "execution_facts" not in without["source_summary"] + assert "execution_facts_collected" not in without["source_summary"] def test_unrelated_blocked_activity_does_not_change_current_work(monkeypatch): From 3bbe2dfdb31b5fdb117d1b04c002b40771f90e5b Mon Sep 17 00:00:00 2001 From: song <22676124+songoow@users.noreply.github.com> Date: Tue, 29 Sep 2026 12:31:10 -0400 Subject: [PATCH 6/9] chore(census): follow registry reads moved on main Regenerated with scripts/generate_project_registry_io_manifest.py after rebasing onto main. Site ids and classifications are unchanged. Co-Authored-By: Claude Opus 5.5 (1M context) Signed-off-by: song <22676124+songoow@users.noreply.github.com> --- loopx/semantics/project_registry_io_manifest_v1.json | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/loopx/semantics/project_registry_io_manifest_v1.json b/loopx/semantics/project_registry_io_manifest_v1.json index 0378be2f0..a6d5733f8 100644 --- a/loopx/semantics/project_registry_io_manifest_v1.json +++ b/loopx/semantics/project_registry_io_manifest_v1.json @@ -1799,7 +1799,7 @@ }, { "site": "loopx/history.py::.collect_history::codec_read:load_registry#1", - "line": 341, + "line": 342, "column": 20, "kind": "codec_read", "api": "load_registry", @@ -1807,7 +1807,7 @@ }, { "site": "loopx/history.py::.inspect_index_duplicates::codec_read:load_registry#1", - "line": 597, + "line": 598, "column": 16, "kind": "codec_read", "api": "load_registry", @@ -1815,7 +1815,7 @@ }, { "site": "loopx/history.py::.rebuild_index_artifact_collisions::codec_read:load_registry#1", - "line": 811, + "line": 812, "column": 16, "kind": "codec_read", "api": "load_registry", @@ -1823,7 +1823,7 @@ }, { "site": "loopx/history.py::.repair_index_duplicates::codec_read:load_registry#1", - "line": 701, + "line": 702, "column": 16, "kind": "codec_read", "api": "load_registry", From a40909c1573c39d85b54b8c156e0d8b40861c8e7 Mon Sep 17 00:00:00 2001 From: song Date: Wed, 30 Sep 2026 10:47:04 +0800 Subject: [PATCH 7/9] fix(agents): attribute a delegation slot only to its own in-flight operation A live operation lock and a live execution slot for the same Todo were two independent facts: a reader holding an old, accepted operation of a re-bound Todo was projected as executing while another operation's worker held the slot. The worker holds its operation lock and, nested inside it, the slot in one process. A slot now attributes to at most one operation: in flight per the delegation transitions, held by the slot holder's own host and pid, taken before the slot, and the latest such operation; a tie names nobody. The read stays lock-free and grants nothing. Negative cases: a two-process historical reader, a reader that locked first, settled operations under the slot holder, released and crashed records, a reused pid, another agent's live lane, and a lease left by the previous claimant. Signed-off-by: song --- loopx/control_plane/agents/execution_facts.py | 58 ++++-- .../test_agent_lifecycle_state.py | 181 +++++++++++++++++- 2 files changed, 220 insertions(+), 19 deletions(-) diff --git a/loopx/control_plane/agents/execution_facts.py b/loopx/control_plane/agents/execution_facts.py index 8af2b3bbb..a6757e19d 100644 --- a/loopx/control_plane/agents/execution_facts.py +++ b/loopx/control_plane/agents/execution_facts.py @@ -7,8 +7,8 @@ - the Turn lane holder record, one per Goal and agent, read through `turn_lane_liveness`; a delegated member executes inside its own lane too; -- the delegation worker: a delegation journal row whose operation lock and - whose execution slot for the delegated Todo are both held by a live process; +- the delegation worker: an in-flight delegation journal row whose operation + lock and the delegated Todo's execution slot are held by one live process; - the task leases on the Goal's open Todos, read from the canonical head after cutover and from the local lease files before it. @@ -31,7 +31,7 @@ from ..collaboration.inbox import _read as read_manager_context_record from ..collaboration.inbox import _root as manager_context_root from ..coordination.local_authority import read_canonical_todos_if_promoted -from ..runtime.time import now_utc +from ..runtime.time import now_utc, parse_timestamp from ..todos.contract import normalize_todo_claimed_by from ..turn_driver.lane_fence import ( TURN_LANE_ABSENT, @@ -61,6 +61,10 @@ LEASE_STATUS_UNAVAILABLE = "unavailable" _LANE_HOLDER_FIELDS = ("host", "pid", "acquired_at") _JOURNAL_ADDRESS = re.compile(r"[a-f0-9]{64}") +# Observations that still have a transition in the delegation owner +# (`loopx/control_plane/collaboration/delegation.ts`); `accepted` and `rejected` +# are final, and an unlisted status is not evidence of a running worker. +_DELEGATION_IN_FLIGHT_STATUSES = frozenset({"prepared", "running", "turn_returned"}) def _lane_rank(state: str) -> int: @@ -96,20 +100,29 @@ def _lane_fact(runtime_root: Path, *, goal_id: str, agent_id: str) -> dict[str, return fact +def _same_process(holder: Mapping[str, Any], other: Mapping[str, Any]) -> bool: + return bool(holder.get("host")) and all(holder.get(key) == other.get(key) for key in ("host", "pid")) + + def _delegation_worker_agents( runtime_root: Path, *, goal_id: str, requesters: Iterable[str] ) -> set[str]: """Delegated members whose worker is executing their Todo right now. A journal row lives under its Goal and requester, so only this Goal's - requesters are read. The row's operation lock is also taken briefly by - readers and by result adoption, so a live operation holder alone is not a - worker; the execution slot for the delegated Todo is taken only by the - worker while it runs, which makes a live slot holder the execution fact. + requesters are read. The worker holds the row's operation lock and, nested + inside it, the execution slot of the delegated Todo, in one process. Two + live holders are two facts until they are the same process: readers and + result adoption hold an operation lock briefly, so a re-bound Todo can show + a live reader on an old operation while another operation's worker holds + the slot. A slot therefore attributes to at most one operation: in flight, + held by the slot's own process, and the latest such operation to be taken + before the slot. An operation with no further transition (see + ``delegation.ts``) is never a worker, whoever holds its lock. """ store = manager_context_root(runtime_root) - active: set[str] = set() + candidates: dict[Path, list[tuple[datetime, str, Mapping[str, Any]]]] = {} for requester in dict.fromkeys(requesters): journal = store / "executions" / manager_context_hash([goal_id, requester]) try: @@ -119,17 +132,38 @@ def _delegation_worker_agents( for path in rows: if not _JOURNAL_ADDRESS.fullmatch(path.stem) or path.is_symlink() or not path.is_file(): continue - if lock_holder_liveness(path)[0] != LOCK_HOLDER_LIVE: + liveness, holder = lock_holder_liveness(path) + acquired_at = parse_timestamp(holder.get("acquired_at")) + if liveness != LOCK_HOLDER_LIVE or acquired_at is None: continue try: - binding = read_manager_context_record(path)["identity"]["binding"] + row = read_manager_context_record(path) + binding = row["identity"]["binding"] member, todo_id = binding["agent_id"], binding["todo_id"] slot = store / "execution-slots" / manager_context_hash([goal_id, todo_id]) except (OSError, ValueError, KeyError, TypeError): continue agent_id = normalize_todo_claimed_by(member) - if agent_id and lock_holder_liveness(slot)[0] == LOCK_HOLDER_LIVE: - active.add(agent_id) + if agent_id and row.get("status") in _DELEGATION_IN_FLIGHT_STATUSES: + candidates.setdefault(slot, []).append((acquired_at, agent_id, holder)) + active: set[str] = set() + for slot, operations in candidates.items(): + liveness, slot_holder = lock_holder_liveness(slot) + slot_acquired_at = parse_timestamp(slot_holder.get("acquired_at")) + if liveness != LOCK_HOLDER_LIVE or slot_acquired_at is None: + continue + owned = [ + (acquired_at, agent_id) + for acquired_at, agent_id, holder in operations + if _same_process(holder, slot_holder) and acquired_at <= slot_acquired_at + ] + if not owned: + continue + latest = max(acquired_at for acquired_at, _agent in owned) + agents = {agent_id for acquired_at, agent_id in owned if acquired_at == latest} + # A tie between members cannot name the worker; fail closed on it. + if len(agents) == 1: + active |= agents return active diff --git a/tests/control_plane/test_agent_lifecycle_state.py b/tests/control_plane/test_agent_lifecycle_state.py index f31bc744b..61abbed95 100644 --- a/tests/control_plane/test_agent_lifecycle_state.py +++ b/tests/control_plane/test_agent_lifecycle_state.py @@ -5,10 +5,13 @@ fixtures here are real: a lane held in-process, a holder record edited the way a crash or another machine leaves it, a lease file, a delegation row lock. """ +from contextlib import contextmanager from datetime import datetime, timedelta, timezone import json import os from pathlib import Path +import subprocess +import sys import pytest @@ -146,12 +149,13 @@ def _dead_pid() -> int: continue -def _write_lease(runtime_root: Path, *, expires_in_hours: float, status: str = "active") -> None: +def _write_lease(runtime_root: Path, *, expires_in_hours: float, status: str = "active", + owner: str = "peer") -> None: now = datetime.now(timezone.utc) path = runtime_root / "goals" / GOAL / "task-leases" / "todo_peer.json" path.parent.mkdir(parents=True, exist_ok=True) path.write_text(json.dumps({ - "schema_version": "task_lease_v0", "goal_id": GOAL, "todo_id": "todo_peer", "owner": "peer", + "schema_version": "task_lease_v0", "goal_id": GOAL, "todo_id": "todo_peer", "owner": owner, "idempotency_key": "execution-peer", "version": 1, "lease_epoch": 1, "status": status, "write_scopes": ["src/**"], "acquire_ttl_seconds": 600, "acquired_at": now.isoformat(), "updated_at": now.isoformat(), "expires_at": (now + timedelta(hours=expires_in_hours)).isoformat(), @@ -218,17 +222,42 @@ def test_an_expired_active_lease_with_nothing_live_is_unknown(monkeypatch, tmp_p assert released["agents"][0]["execution"]["lease"] == {"status": "released"} -def _delegation_row(runtime_root: Path, *, goal: str, requester: str) -> tuple[Path, Path]: +def _delegation_row(runtime_root: Path, *, goal: str, requester: str, agent: str = "peer", + operation: str = "op-1", status: str = "running") -> tuple[Path, Path]: """A journal row where the delegation service keeps it, plus its execution slot.""" store = manager_context_root(runtime_root) - row_path = store / "executions" / manager_context_hash([goal, requester]) / ("b" * 64 + ".json") + row_path = (store / "executions" / manager_context_hash([goal, requester]) + / (manager_context_hash(operation) + ".json")) row_path.parent.mkdir(parents=True, exist_ok=True) - row_path.write_text(json.dumps({"identity": {"binding": {"id": "b1", "agent_id": "peer", "todo_id": "todo_peer"}, - "request_id": "r1", "operation_id": "op-1"}, - "status": "running"}), encoding="utf-8") + row_path.write_text(json.dumps({"identity": {"binding": {"id": "b1", "agent_id": agent, "todo_id": "todo_peer"}, + "request_id": "r1", "operation_id": operation}, + "status": status}), encoding="utf-8") return row_path, store / "execution-slots" / manager_context_hash([goal, "todo_peer"]) +_WORKER = """ +import sys +from pathlib import Path +from loopx.file_lock import exclusive_file_lock +with exclusive_file_lock(Path(sys.argv[1])), exclusive_file_lock(Path(sys.argv[2])): + print("held", flush=True) + sys.stdin.readline() +""" + + +@contextmanager +def _worker_process(row_path: Path, slot: Path): + """A separate process holding one operation lock and the Todo's slot, as a worker does.""" + worker = subprocess.Popen([sys.executable, "-c", _WORKER, str(row_path), str(slot)], + stdin=subprocess.PIPE, stdout=subprocess.PIPE, text=True) + try: + assert worker.stdout.readline().strip() == "held" + yield worker.pid + finally: + worker.stdin.close() + worker.wait(timeout=30) + + def _peer_row(packet): return next(row for row in packet["agents"] if row["agent_id"] == "peer") @@ -262,6 +291,144 @@ def test_an_operation_lock_without_its_execution_slot_is_not_a_worker(monkeypatc assert _peer_row(scoped)["state"] == "launchable" +REUSED = ("coordinator", "peer", "old-peer") + + +def test_a_historical_result_reader_does_not_borrow_the_current_workers_slot(monkeypatch, tmp_path): + """The Todo was re-bound: old-peer's accepted operation is only being read while peer's worker runs.""" + runtime_root = tmp_path / "runtime" + old_row, slot = _delegation_row(runtime_root, goal=GOAL, requester="coordinator", agent="old-peer", + operation="op-old", status="accepted") + new_row, _ = _delegation_row(runtime_root, goal=GOAL, requester="coordinator", operation="op-new") + with _worker_process(new_row, slot) as worker_pid, exclusive_file_lock(old_row): + assert worker_pid != os.getpid() + packet, _ = build_projection(monkeypatch, age=30, runtime_root=runtime_root, registered=REUSED) + facts = collect_agent_execution_facts(runtime_root=runtime_root, status_payload=_payload(REUSED)) + by_agent = {row["agent_id"]: row for row in packet["agents"]} + assert by_agent["peer"]["state"] == "executing" + assert by_agent["old-peer"]["state"] == "registered" + assert facts["old-peer"]["delegation_worker_active"] is False + + +def test_an_operation_held_by_another_process_than_the_slot_is_not_that_workers(monkeypatch, tmp_path): + """Two live holders are two facts: only the process holding the slot is the worker.""" + runtime_root = tmp_path / "runtime" + # Still in flight, so only the holder identity can tell the two operations apart. + old_row, slot = _delegation_row(runtime_root, goal=GOAL, requester="old-requester", agent="old-peer", + operation="op-old", status="running") + new_row, _ = _delegation_row(runtime_root, goal=GOAL, requester="coordinator", operation="op-new") + registered = ("coordinator", "old-requester", "peer", "old-peer") + with _worker_process(new_row, slot), exclusive_file_lock(old_row): + packet, _ = build_projection(monkeypatch, age=30, runtime_root=runtime_root, registered=registered) + by_agent = {row["agent_id"]: row for row in packet["agents"]} + assert by_agent["peer"]["state"] == "executing" + assert by_agent["old-peer"]["state"] == "registered" + + +def test_a_reader_that_locked_first_is_not_the_worker_of_an_unseen_operation(monkeypatch, tmp_path): + """The reader took the old operation before the worker took the slot; order alone cannot tell them apart.""" + runtime_root = tmp_path / "runtime" + old_row, slot = _delegation_row(runtime_root, goal=GOAL, requester="coordinator", agent="old-peer", + operation="op-old", status="running") + # The worker's own row is under a requester this Goal does not project. + new_row, _ = _delegation_row(runtime_root, goal=GOAL, requester="unlisted", operation="op-new") + with exclusive_file_lock(old_row), _worker_process(new_row, slot): + packet, _ = build_projection(monkeypatch, age=30, runtime_root=runtime_root, registered=REUSED) + by_agent = {row["agent_id"]: row for row in packet["agents"]} + assert by_agent["old-peer"]["state"] == "registered" + assert by_agent["peer"]["state"] == "launchable" + + +def test_a_crashed_operation_whose_pid_was_reused_does_not_share_the_slot(monkeypatch, tmp_path): + """A slot names one worker: the latest in-flight operation its process took before the slot.""" + runtime_root = tmp_path / "runtime" + old_row, slot = _delegation_row(runtime_root, goal=GOAL, requester="coordinator", agent="old-peer", + operation="op-old", status="running") + new_row, _ = _delegation_row(runtime_root, goal=GOAL, requester="coordinator", operation="op-new") + with exclusive_file_lock(old_row): + pass + with _worker_process(new_row, slot) as worker_pid: + # The old worker crashed mid-run and its pid now belongs to the new worker. + _rewrite_holder(old_row, pid=worker_pid) + packet, _ = build_projection(monkeypatch, age=30, runtime_root=runtime_root, registered=REUSED) + by_agent = {row["agent_id"]: row for row in packet["agents"]} + assert by_agent["peer"]["state"] == "executing" + assert by_agent["old-peer"]["state"] == "registered" + # When the records cannot say which operation came last, neither is named. + new_acquired = json.loads(lock_holder_path(new_row).read_text(encoding="utf-8"))["acquired_at"] + _rewrite_holder(old_row, pid=worker_pid, acquired_at=new_acquired) + tied, _ = build_projection(monkeypatch, age=30, runtime_root=runtime_root, registered=REUSED) + # The slot is taken inside its operation: a lock the process took after the slot is not its operation. + slot_acquired = datetime.fromisoformat( + json.loads(lock_holder_path(slot).read_text(encoding="utf-8"))["acquired_at"].replace("Z", "+00:00")) + _rewrite_holder(old_row, pid=worker_pid, acquired_at=(slot_acquired + timedelta(seconds=1)).isoformat()) + later, _ = build_projection(monkeypatch, age=30, runtime_root=runtime_root, registered=REUSED) + assert {row["agent_id"]: row["state"] for row in later["agents"]} == { + "coordinator": "registered", "peer": "executing", "old-peer": "registered"} + assert {row["agent_id"]: row["state"] for row in tied["agents"]} == { + "coordinator": "registered", "peer": "launchable", "old-peer": "registered"} + + +def test_a_settled_operation_is_no_worker_even_under_the_slot_holder(monkeypatch, tmp_path): + """An accepted or rejected operation has no further transition: its holder is a reader.""" + runtime_root = tmp_path / "runtime" + for status in ("accepted", "rejected"): + row_path, slot = _delegation_row(runtime_root, goal=GOAL, requester="coordinator", status=status) + with exclusive_file_lock(row_path), exclusive_file_lock(slot): + packet, _ = build_projection(monkeypatch, age=30, runtime_root=runtime_root, + registered=("coordinator", "peer")) + assert _peer_row(packet)["state"] == "launchable", status + + +def test_a_released_or_stale_worker_record_is_no_worker(monkeypatch, tmp_path): + runtime_root = tmp_path / "runtime" + registered = ("coordinator", "peer") + row_path, slot = _delegation_row(runtime_root, goal=GOAL, requester="coordinator") + with exclusive_file_lock(row_path), exclusive_file_lock(slot): + pass + released, _ = build_projection(monkeypatch, age=0, runtime_root=runtime_root, registered=registered) + assert _peer_row(released)["state"] == "launchable" + # A crashed worker leaves both records naming its pid; a dead pid executes nothing. + dead = _dead_pid() + _rewrite_holder(row_path, pid=dead) + _rewrite_holder(slot, pid=dead) + crashed, _ = build_projection(monkeypatch, age=0, runtime_root=runtime_root, registered=registered) + assert _peer_row(crashed)["state"] == "launchable" + + +def test_one_agents_live_lane_is_not_another_agents_execution(monkeypatch, tmp_path): + runtime_root = tmp_path / "runtime" + plan = {"turn_envelope": {"agent_id": "old-peer"}} + with turn_lane_singleflight(runtime_root=runtime_root, goal_id=GOAL, plan=plan) as held: + assert held is not None + packet, _ = build_projection(monkeypatch, age=0, runtime_root=runtime_root, registered=REUSED) + by_agent = {row["agent_id"]: row for row in packet["agents"]} + assert by_agent["old-peer"]["state"] == "executing" + assert by_agent["peer"]["state"] == "launchable" + assert by_agent["peer"]["execution"]["lane"] == "absent" + + +def test_a_lease_left_by_the_previous_claimant_stays_with_its_owner(monkeypatch, tmp_path): + """A re-bound Todo's stale lease is the old owner's unverifiable fact, never the new claimant's.""" + runtime_root = tmp_path / "runtime" + _write_lease(runtime_root, expires_in_hours=-1, owner="old-peer") + packet, _ = build_projection(monkeypatch, age=0, runtime_root=runtime_root, registered=REUSED) + by_agent = {row["agent_id"]: row for row in packet["agents"]} + assert by_agent["peer"]["state"] == "launchable" + assert "lease" not in by_agent["peer"]["execution"] + assert by_agent["old-peer"]["state"] == "unknown" + _write_lease(runtime_root, expires_in_hours=-1, owner="old-peer", status="released") + released, _ = build_projection(monkeypatch, age=0, runtime_root=runtime_root, registered=REUSED) + assert {row["agent_id"]: row["state"] for row in released["agents"]}["old-peer"] == "registered" + + +def _payload(registered): + todo = {"todo_id": "todo_peer", "goal_id": GOAL, "role": "agent", "claimed_by": "peer", "status": "open"} + return {"goal_filter": GOAL, "run_history": {"goals": [{ + "id": GOAL, "coordination": {"registered_agents": list(registered)}}]}, + "todo_index": {"items": [todo]}} + + def test_no_runtime_root_means_no_facts_not_no_execution(monkeypatch): payload = {"goal_filter": GOAL, "run_history": {"goals": [{"id": GOAL, "coordination": {"registered_agents": ["peer"]}}]}} assert collect_agent_execution_facts(runtime_root=None, status_payload=payload) is None From 6da1f3bf752666ca134d43044bba3aa1bd0f48d6 Mon Sep 17 00:00:00 2001 From: song Date: Wed, 30 Sep 2026 11:05:22 +0800 Subject: [PATCH 8/9] fix(agents): borrow the journal address shape from the digest owner Main now pins one owner for the whole-value SHA-256 shape. The execution facts reader checked delegation journal filenames with its own pattern; it now uses BARE_SHA256_PATTERN as delegation_inventory does and is pinned as a consumer. Signed-off-by: song --- loopx/control_plane/agents/execution_facts.py | 5 ++--- tests/architecture/test_content_digest_single_owner.py | 1 + 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/loopx/control_plane/agents/execution_facts.py b/loopx/control_plane/agents/execution_facts.py index a6757e19d..c1e7d26c1 100644 --- a/loopx/control_plane/agents/execution_facts.py +++ b/loopx/control_plane/agents/execution_facts.py @@ -23,13 +23,13 @@ from collections.abc import Iterable, Mapping from datetime import datetime from pathlib import Path -import re from typing import Any from ...file_lock import LOCK_HOLDER_LIVE, lock_holder_liveness from ..collaboration.inbox import _hash as manager_context_hash from ..collaboration.inbox import _read as read_manager_context_record from ..collaboration.inbox import _root as manager_context_root +from ..content_digest import BARE_SHA256_PATTERN from ..coordination.local_authority import read_canonical_todos_if_promoted from ..runtime.time import now_utc, parse_timestamp from ..todos.contract import normalize_todo_claimed_by @@ -60,7 +60,6 @@ LEASE_STATUS_ACTIVE = "active" LEASE_STATUS_UNAVAILABLE = "unavailable" _LANE_HOLDER_FIELDS = ("host", "pid", "acquired_at") -_JOURNAL_ADDRESS = re.compile(r"[a-f0-9]{64}") # Observations that still have a transition in the delegation owner # (`loopx/control_plane/collaboration/delegation.ts`); `accepted` and `rejected` # are final, and an unlisted status is not evidence of a running worker. @@ -130,7 +129,7 @@ def _delegation_worker_agents( except OSError: continue for path in rows: - if not _JOURNAL_ADDRESS.fullmatch(path.stem) or path.is_symlink() or not path.is_file(): + if not BARE_SHA256_PATTERN.fullmatch(path.stem) or path.is_symlink() or not path.is_file(): continue liveness, holder = lock_holder_liveness(path) acquired_at = parse_timestamp(holder.get("acquired_at")) diff --git a/tests/architecture/test_content_digest_single_owner.py b/tests/architecture/test_content_digest_single_owner.py index 8304b7e08..45298959a 100644 --- a/tests/architecture/test_content_digest_single_owner.py +++ b/tests/architecture/test_content_digest_single_owner.py @@ -175,6 +175,7 @@ "loopx.capabilities.progress_review.receipt", "loopx.chat_action_normalization", "loopx.configuration_transaction", + "loopx.control_plane.agents.execution_facts", "loopx.control_plane.collaboration.delegation_inventory", "loopx.control_plane.collaboration.inbox", "loopx.control_plane.collaboration.peers", From 4a71c3b4e63320edebf8cfcc041564b38f1da993 Mon Sep 17 00:00:00 2001 From: song Date: Wed, 30 Sep 2026 12:58:19 +0800 Subject: [PATCH 9/9] test(agents): keep the delegation slot fixtures in one place Move the projection payload helper next to the other fixtures and drop an assertion the worker-process helper already guarantees. Signed-off-by: song --- .../test_agent_lifecycle_state.py | 19 +++++++++---------- 1 file changed, 9 insertions(+), 10 deletions(-) diff --git a/tests/control_plane/test_agent_lifecycle_state.py b/tests/control_plane/test_agent_lifecycle_state.py index 61abbed95..1983b925a 100644 --- a/tests/control_plane/test_agent_lifecycle_state.py +++ b/tests/control_plane/test_agent_lifecycle_state.py @@ -291,7 +291,14 @@ def test_an_operation_lock_without_its_execution_slot_is_not_a_worker(monkeypatc assert _peer_row(scoped)["state"] == "launchable" -REUSED = ("coordinator", "peer", "old-peer") +def _payload(registered): + todo = {"todo_id": "todo_peer", "goal_id": GOAL, "role": "agent", "claimed_by": "peer", "status": "open"} + return {"goal_filter": GOAL, "run_history": {"goals": [{ + "id": GOAL, "coordination": {"registered_agents": list(registered)}}]}, + "todo_index": {"items": [todo]}} + + +REUSED = ("coordinator", "peer", "old-peer") # a Todo re-bound from old-peer to peer def test_a_historical_result_reader_does_not_borrow_the_current_workers_slot(monkeypatch, tmp_path): @@ -300,8 +307,7 @@ def test_a_historical_result_reader_does_not_borrow_the_current_workers_slot(mon old_row, slot = _delegation_row(runtime_root, goal=GOAL, requester="coordinator", agent="old-peer", operation="op-old", status="accepted") new_row, _ = _delegation_row(runtime_root, goal=GOAL, requester="coordinator", operation="op-new") - with _worker_process(new_row, slot) as worker_pid, exclusive_file_lock(old_row): - assert worker_pid != os.getpid() + with _worker_process(new_row, slot), exclusive_file_lock(old_row): packet, _ = build_projection(monkeypatch, age=30, runtime_root=runtime_root, registered=REUSED) facts = collect_agent_execution_facts(runtime_root=runtime_root, status_payload=_payload(REUSED)) by_agent = {row["agent_id"]: row for row in packet["agents"]} @@ -422,13 +428,6 @@ def test_a_lease_left_by_the_previous_claimant_stays_with_its_owner(monkeypatch, assert {row["agent_id"]: row["state"] for row in released["agents"]}["old-peer"] == "registered" -def _payload(registered): - todo = {"todo_id": "todo_peer", "goal_id": GOAL, "role": "agent", "claimed_by": "peer", "status": "open"} - return {"goal_filter": GOAL, "run_history": {"goals": [{ - "id": GOAL, "coordination": {"registered_agents": list(registered)}}]}, - "todo_index": {"items": [todo]}} - - def test_no_runtime_root_means_no_facts_not_no_execution(monkeypatch): payload = {"goal_filter": GOAL, "run_history": {"goals": [{"id": GOAL, "coordination": {"registered_agents": ["peer"]}}]}} assert collect_agent_execution_facts(runtime_root=None, status_payload=payload) is None