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..1f522ae10 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 @@ -103,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(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")): + 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.""" @@ -120,9 +143,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 +200,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 +209,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 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 = _projected_execution_facts(payload) + 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 +236,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..c1e7d26c1 --- /dev/null +++ b/loopx/control_plane/agents/execution_facts.py @@ -0,0 +1,267 @@ +"""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, 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: 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. + +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. 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 +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 +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, 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. +# 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") +# 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: + 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 lease 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: + fact["lane_holder"] = holder + 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 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) + 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: + rows = sorted(journal.glob("*.json")) if journal.is_dir() else [] + except OSError: + continue + for path in rows: + 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")) + if liveness != LOCK_HOLDER_LIVE or acquired_at is None: + continue + try: + 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 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 + + +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 (OSError, RuntimeError, ValueError): + return None + if canonical is not None: + 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)) + # 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 _merge_lane(row: dict[str, Any], fact: Mapping[str, Any]) -> None: + if _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]] | None: + """Return ``agent_id -> {lane, lane_holder?, delegation_worker_active, lease?}``. + + 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 None + root = Path(runtime_root) + if not root.is_dir(): + return None + facts: dict[str, dict[str, Any]] = {} + at = now_utc() + 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 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 agents: + _merge_lease(facts[owner], _lease_fact(lease, at=at)) + return facts diff --git a/loopx/control_plane/agents/management_projection.py b/loopx/control_plane/agents/management_projection.py index b8f9e1d64..809921a89 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,50 @@ 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, 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, dict[str, list[str]]]] = {} + + 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 or not agent_id: + return + 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 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 + for raw_agent in _as_list(_as_dict(goal.get("coordination")).get("registered_agents")): + add(goal.get("id"), raw_agent) + for todo in _iter_status_todos(status_payload): + if not _is_done(todo): + add( + todo.get("goal_id"), + _todo_agent_id(todo), + todo_id=normalize_todo_id(todo.get("todo_id")), + ) + return goals + + def _agent_material_frontiers( status_payload: dict[str, Any], ) -> dict[tuple[str, str], dict[str, Any]]: @@ -474,23 +517,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 +557,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 +586,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 +680,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 +711,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 +760,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 +774,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 +834,10 @@ def build_agent_management_projection( } if material_frontiers: source_summary["material_frontier_count"] = len(material_frontiers) + if isinstance(execution_facts, dict): + # 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/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 3bca96593..f0359a6c8 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,16 @@ 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`. 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 + ) 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/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/loopx/semantics/project_registry_io_manifest_v1.json b/loopx/semantics/project_registry_io_manifest_v1.json index 956f4c8b0..5ff87ea24 100644 --- a/loopx/semantics/project_registry_io_manifest_v1.json +++ b/loopx/semantics/project_registry_io_manifest_v1.json @@ -1535,7 +1535,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", 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", diff --git a/tests/control_plane/test_agent_lifecycle_state.py b/tests/control_plane/test_agent_lifecycle_state.py index 2f82553c3..1983b925a 100644 --- a/tests/control_plane/test_agent_lifecycle_state.py +++ b/tests/control_plane/test_agent_lifecycle_state.py @@ -1,30 +1,52 @@ -"""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 contextlib import contextmanager from datetime import datetime, timedelta, timezone +import json +import os +from pathlib import Path +import subprocess +import sys 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 _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 +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, registered=("peer",)): 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": { - "registered_agents": ["peer"], + payload = {"goal_filter": GOAL, "run_history": {"goals": [{ + "id": GOAL, "coordination": { + "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}} - 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 +56,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 +69,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 +93,26 @@ 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_facts_collected"] is True + without, _ = build_projection(monkeypatch, age=0) + assert "execution_facts_collected" 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 +122,340 @@ 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", + 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": 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(), + }), 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 _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]) + / (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": 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") + + +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=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 _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): + """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), 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"]} + 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 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 + 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] + + +# --- 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( 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" }), 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