diff --git a/docs/reference/protocols/agent-management-projection-v0.md b/docs/reference/protocols/agent-management-projection-v0.md index fcf3fcebcc..057acd1a20 100644 --- a/docs/reference/protocols/agent-management-projection-v0.md +++ b/docs/reference/protocols/agent-management-projection-v0.md @@ -91,7 +91,8 @@ Required fields: - `agent_id`; - `agent_model`: `peer_v1`; - `state`: one of `running`, `waiting`, `blocked`, `monitoring`, - `scope_wait`, `stale`, or `unknown`; + `scope_wait`, `stale`, `unknown`, `registered`, `addressable`, `bound`, + `launchable`, `executing`; - `current_todo`: a `todo_row_v0` object or `null`; - `next_action`: compact local-control next action text. Private project refs are allowed; inline credentials are not. Shareable sinks must redact private @@ -111,6 +112,9 @@ Optional fields: - `handoff_refs`; - `handoff_note`; - `material_frontier`; +- `session_binding`: optional `{thread_id, host_surface}` from existing + `run_history.goals[].coordination.thread_agent_bindings`; absence is not + evidence of a stopped host, and a binding does not prove executable capacity; - `stale_claim_hint`; - `blocked_on`; - `recent_events`; @@ -122,6 +126,42 @@ highest-priority blocked maintenance todo as a separate `blocked_on` `todo_row_v0`. The blocker remains visible without changing todo ownership or making the whole peer appear blocked. +### Worker state refinement and compatibility + +The current producer emits `registered`, `addressable`, `bound`, `launchable`, +`executing`, `blocked`, `monitoring`, and `waiting`. +`running`, `unknown`, `scope_wait`, and `stale` remain +accepted legacy vocabulary; this producer does not emit them. There is no +parallel `lifecycle_state` field. + +This changes the default read projection: registered idle rows previously +reported `unknown` or `waiting`, and open advancement work reported `running`. +The new states refine those observations without adding dispatch authority: + +| State | Observed facts | +| --- | --- | +| `registered` | Registered row without current work or session binding | +| `addressable` | Session binding exists, with no current work | +| `bound` | Current open work and a session binding; no recent work update | +| `launchable` | Current open work without a binding or recent work update | +| `executing` | Current open work updated between zero and eight hours ago | +| `blocked` | Current work is blocked or has the blocker task class | +| `monitoring` | Current work is a monitor, regardless of activity age | +| `waiting` | Other non-open current work, such as deferred work | + +Blocked, monitor and waiting classification precedes activity/binding refinement. +The eight-hour activity window is inclusive, rejects future timestamps, and +uses the current Todo only; it is independent of the 36-hour stale-claim +warning. `last_activity_at` still summarizes the displayed Todos. Neither an +`executing` observation nor `launchable` proves a live process, configured +runtime, available capacity, lease, or permission to start a worker. + +Registered-peer orchestration accepts `executing`, `bound`, and `launchable` +where it accepted legacy `running`, and continues to accept `monitoring` and +cached `running` rows. Its existing stale-claim, observed activation-capability, +and dependency-readiness gates still apply. Idle, blocked, waiting and unknown +rows do not gain admission. Other consumers must tolerate the added strings. + ## Todo Row `todo_row_v0` is the dashboard/review-packet representation of an existing diff --git a/examples/control_plane/agent-management-live-status-smoke.py b/examples/control_plane/agent-management-live-status-smoke.py index 5f83e16c66..09bd58e060 100644 --- a/examples/control_plane/agent-management-live-status-smoke.py +++ b/examples/control_plane/agent-management-live-status-smoke.py @@ -254,11 +254,11 @@ def assert_event_only_todo_receipts_are_not_runtime_work() -> None: if isinstance(row, dict) } assert by_agent["agent-reviewer"]["current_todo"]["todo_id"] == "todo_live_handoff", by_agent - assert by_agent["agent-event-only"]["state"] == "unknown", by_agent + assert by_agent["agent-event-only"]["state"] == "registered", by_agent assert "current_todo" not in by_agent["agent-event-only"], by_agent -def assert_advancement_current_todo_beats_standing_monitor() -> None: +def assert_launchable_advancement_beats_standing_monitor() -> None: payload = fixture_status_payload() payload["run_history"]["goals"][0]["coordination"]["registered_agents"].append("agent-side") payload["attention_queue"]["items"].append( @@ -342,7 +342,7 @@ def assert_advancement_current_todo_beats_standing_monitor() -> None: if isinstance(row, dict) } side_row = by_agent["agent-side"] - assert side_row["state"] == "running", side_row + assert side_row["state"] == "launchable", side_row assert side_row["current_todo"]["todo_id"] == "todo_side_advancement", side_row assert side_row["next_action"] == "Continue projected todo todo_side_advancement.", side_row @@ -369,7 +369,7 @@ def assert_advancement_current_todo_beats_standing_monitor() -> None: for row in stale_blocker.get("agents", []) if isinstance(row, dict) and row.get("agent_id") == "agent-side" ) - assert stale_blocker_row["state"] == "running", stale_blocker_row + assert stale_blocker_row["state"] == "launchable", stale_blocker_row assert stale_blocker_row["current_todo"]["todo_id"] == "todo_side_advancement", stale_blocker_row assert stale_blocker_row["blocked_on"]["todo_id"] == "todo_stale_maintenance_blocker", stale_blocker_row @@ -380,7 +380,7 @@ def assert_advancement_current_todo_beats_standing_monitor() -> None: for row in lane_only.get("agents", []) if isinstance(row, dict) and row.get("agent_id") == "agent-side" ) - assert lane_only_row["state"] == "running", lane_only_row + assert lane_only_row["state"] == "launchable", lane_only_row assert lane_only_row["current_todo"]["todo_id"] == "todo_side_advancement", lane_only_row payload["attention_queue"]["items"][-1]["project_asset"]["agent_todos"][ @@ -508,7 +508,7 @@ def main() -> int: args = parse_args() assert_synthetic_projection() assert_event_only_todo_receipts_are_not_runtime_work() - assert_advancement_current_todo_beats_standing_monitor() + assert_launchable_advancement_beats_standing_monitor() result: dict[str, Any] = { "synthetic": "ok", "bundled_example": assert_bundled_public_example(), diff --git a/examples/worker-lifecycle-state-smoke.py b/examples/worker-lifecycle-state-smoke.py new file mode 100644 index 0000000000..c74f3d196e --- /dev/null +++ b/examples/worker-lifecycle-state-smoke.py @@ -0,0 +1,154 @@ +#!/usr/bin/env python3 +"""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). + +Run from the repository root: + uv run --extra test python examples/worker-lifecycle-state-smoke.py +""" + +from __future__ import annotations + +import sys +from datetime import datetime, timedelta, timezone +from pathlib import Path + +# Add the repository root to the path for direct execution +sys.path.insert(0, str(Path(__file__).parent.parent)) + + +def _recent_activity() -> str: + """Activity timestamp within the activity threshold (8 hours).""" + return (datetime.now(timezone.utc) - timedelta(hours=1)).isoformat() + +from loopx.control_plane.agents.management_projection import ( # noqa: E402 + WORKER_LIFECYCLE_STATE_ADDRESSABLE, + WORKER_LIFECYCLE_STATE_BLOCKED, + WORKER_LIFECYCLE_STATE_BOUND, + WORKER_LIFECYCLE_STATE_EXECUTING, + WORKER_LIFECYCLE_STATE_LAUNCHABLE, + WORKER_LIFECYCLE_STATE_REGISTERED, + build_agent_management_projection, +) + + +def build_status_payload() -> dict: + """Build a minimal status payload with known facts.""" + return { + "goal_filter": "smoke-goal", + "run_history": { + "goals": [ + { + "id": "smoke-goal", + "coordination": { + "registered_agents": [ + "worker-registered", + "worker-addressable", + "worker-bound", + "worker-launchable", + "worker-executing", + "worker-blocked", + ], + "thread_agent_bindings": [ + { + "agent_id": "worker-addressable", + "thread_id": "thread-1", + "host_surface": "codex-app", + }, + { + "agent_id": "worker-bound", + "thread_id": "thread-2", + "host_surface": "codex-app", + }, + { + "agent_id": "worker-executing", + "thread_id": "thread-3", + "host_surface": "codex-app", + }, + ], + }, + } + ] + }, + "attention_queue": { + "items": [ + { + "goal_id": "smoke-goal", + "agent_todos": { + "items": [ + { + "todo_id": "todo-bound", + "claimed_by": "worker-bound", + "status": "open", + "updated_at": "2026-09-15T12:00:00+00:00", + }, + { + "todo_id": "todo-launchable", + "claimed_by": "worker-launchable", + "status": "open", + "updated_at": "2026-09-15T12:00:00+00:00", + }, + { + "todo_id": "todo-executing", + "claimed_by": "worker-executing", + "status": "open", + "updated_at": _recent_activity(), + }, + { + "todo_id": "todo-blocked", + "claimed_by": "worker-blocked", + "status": "blocked", + "task_class": "blocker", + "updated_at": "2026-09-17T12:00:00+00:00", + }, + ] + }, + } + ] + }, + } + + +def main() -> int: + payload = build_status_payload() + projection = build_agent_management_projection(payload) + + agents = {a["agent_id"]: a for a in projection.get("agents", [])} + + # Verify each worker's lifecycle state + expected = { + "worker-registered": WORKER_LIFECYCLE_STATE_REGISTERED, + "worker-addressable": WORKER_LIFECYCLE_STATE_ADDRESSABLE, + "worker-bound": WORKER_LIFECYCLE_STATE_BOUND, + "worker-launchable": WORKER_LIFECYCLE_STATE_LAUNCHABLE, + "worker-executing": WORKER_LIFECYCLE_STATE_EXECUTING, + "worker-blocked": WORKER_LIFECYCLE_STATE_BLOCKED, + } + + failures = [] + for agent_id, expected_state in expected.items(): + agent = agents.get(agent_id) + if agent is None: + failures.append(f"{agent_id}: missing from projection") + continue + actual_state = agent.get("state") + if actual_state != expected_state: + failures.append( + f"{agent_id}: expected {expected_state}, got {actual_state}" + ) + else: + print(f" {agent_id}: {actual_state} ✓") + + if failures: + print("\nFAILURES:") + for failure in failures: + print(f" {failure}") + return 1 + + print(f"\nAll {len(expected)} lifecycle states verified.") + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/loopx/control_plane/agents/management_projection.py b/loopx/control_plane/agents/management_projection.py index 3747c0f52c..7d6d496bea 100644 --- a/loopx/control_plane/agents/management_projection.py +++ b/loopx/control_plane/agents/management_projection.py @@ -25,6 +25,7 @@ MAX_REFS = 1 MAX_WORKSPACE_SCOPES = 4 STALE_CLAIM_THRESHOLD_HOURS = 36 +EXECUTING_ACTIVITY_THRESHOLD_HOURS = 8 MATERIAL_LIFECYCLE_CAPABILITY = "material_lifecycle" _TODO_GROUP_LIST_KEYS = tuple( @@ -474,21 +475,79 @@ def _refs_from_todo(todo: dict[str, Any]) -> tuple[list[str], list[str]]: return evidence_refs[:MAX_REFS], handoff_refs[:MAX_REFS] -def _agent_state(todos: list[dict[str, Any]], *, current: dict[str, Any] | None = None) -> str: +def _agent_state( + todos: list[dict[str, Any]], + *, + current: dict[str, Any] | None = None, + has_session_binding: bool = False, + last_activity_at: str | 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). + + 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 + """ open_todos = [todo for todo in todos if not _is_done(todo)] - if not open_todos: - return "waiting" if todos else "unknown" + + # Blocked takes priority: a blocked worker cannot launch or execute. + # This only applies to the current todo, not all open todos, per the + # protocol contract: "blocker remains visible without making the whole + # peer appear blocked." if current and not _is_done(current): if _todo_status(current) == "blocked" or current.get("task_class") == "blocker": - return "blocked" + return WORKER_LIFECYCLE_STATE_BLOCKED + + if current and not _is_done(current): if _is_monitor_todo(current): return "monitoring" - return "running" - if any(_todo_status(todo) == "blocked" or todo.get("task_class") == "blocker" for todo in open_todos): - return "blocked" - if all(_is_monitor_todo(todo) for todo in open_todos): - return "monitoring" - return "running" + 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 + + # Bound: has session binding and active todo. + if current and not _is_done(current) and has_session_binding: + return WORKER_LIFECYCLE_STATE_BOUND + + # Launchable: has active todo, no session binding. + if open_todos: + return WORKER_LIFECYCLE_STATE_LAUNCHABLE + + # Addressable: has session binding but no active todo. + if has_session_binding: + return WORKER_LIFECYCLE_STATE_ADDRESSABLE + + # Registered: in registry, no binding or todo. + return WORKER_LIFECYCLE_STATE_REGISTERED + + +# 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. +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" def _last_activity(todos: list[dict[str, Any]]) -> str | None: @@ -530,6 +589,27 @@ def build_agent_management_projection( else {} ) goal_filter = _compact(status_payload.get("goal_filter"), limit=180) + + # Session bindings come from run_history.coordination.thread_agent_bindings. + # This is the only source for addressable/bound lifecycle states. + session_bindings: dict[str, dict[str, str]] = {} + run_history = _as_dict(status_payload.get("run_history")) + for raw_goal in _as_list(run_history.get("goals")): + if not isinstance(raw_goal, dict): + continue + coordination = _as_dict(raw_goal.get("coordination")) + for raw_binding in _as_list(coordination.get("thread_agent_bindings")): + if not isinstance(raw_binding, dict): + continue + agent_id = _compact(raw_binding.get("agent_id"), limit=120) + thread_id = _compact(raw_binding.get("thread_id"), limit=120) + host_surface = _compact(raw_binding.get("host_surface"), limit=60) + if agent_id and thread_id: + session_bindings[agent_id] = { + "thread_id": thread_id, + "host_surface": host_surface or "unknown", + } + seen_todos: set[tuple[str, str, str, str]] = set() for todo in _iter_status_todos(status_payload): agent_id = _todo_agent_id(todo) @@ -576,17 +656,26 @@ def build_agent_management_projection( for ref in todo_handoffs: if ref not in handoff_refs: handoff_refs.append(ref) + last_activity = _last_activity(todos) + agent_state = _agent_state( + all_todos, + current=current, + has_session_binding=agent_id in session_bindings, + last_activity_at=_last_activity([current]) if current else None, + ) agent_row: dict[str, Any] = { "agent_id": agent_id, "agent_model": raw_row.get("agent_model") or "unregistered", - "state": _agent_state(all_todos, current=current), + "state": agent_state, "current_todo": _todo_row(current) if current else None, "next_action": _safe_next_action(current), - "last_activity_at": _last_activity(todos), + "last_activity_at": last_activity, "evidence_refs": evidence_refs[:MAX_REFS], "handoff_refs": handoff_refs[:MAX_REFS], "goal_ids": _as_list(raw_row.get("_goal_ids"))[:MAX_REFS], } + if agent_id in session_bindings: + agent_row["session_binding"] = session_bindings[agent_id] material_frontier_key = _agent_material_frontier_key( raw_row=raw_row, current=current, diff --git a/loopx/control_plane/quota/task_orchestration.py b/loopx/control_plane/quota/task_orchestration.py index 8e8312c13f..d73b5499d3 100644 --- a/loopx/control_plane/quota/task_orchestration.py +++ b/loopx/control_plane/quota/task_orchestration.py @@ -372,7 +372,12 @@ def _registered_peer_task_orchestration_contract( reason_codes.append("peer_liveness_unavailable") elif peer_state.get("stale_claim_hint"): reason_codes.append("peer_runtime_stale") - elif peer_state.get("state") not in {"running", "monitoring"}: + # New work-state refinements retain the old running admission policy. + # Activity/claim freshness, activation capability and dependency readiness + # remain independent gates; a binding alone never grants activation. + elif peer_state.get("state") not in { + "running", "monitoring", "executing", "bound", "launchable", + }: reason_codes.append("peer_runtime_not_active") if lane.get("resume_when") and lane.get("resume_ready") is not True: reason_codes.append("peer_lane_not_resume_ready") diff --git a/tests/control_plane/test_agent_lifecycle_state.py b/tests/control_plane/test_agent_lifecycle_state.py new file mode 100644 index 0000000000..2d5d4cb039 --- /dev/null +++ b/tests/control_plane/test_agent_lifecycle_state.py @@ -0,0 +1,100 @@ +"""Exercise worker states through the public projection and real peer admission.""" +from datetime import datetime, timedelta, timezone + +import pytest + +from loopx.control_plane.agents import management_projection as projection +from loopx.control_plane.quota.task_orchestration import apply_task_orchestration_contract + +NOW = datetime(2026, 9, 18, 12, tzinfo=timezone.utc) + + +def build_projection(monkeypatch, *, age=None, binding=False, status="open", + task_class="advancement_task", has_todo=True, extra_todos=()): + monkeypatch.setattr(projection, "now_utc", lambda: NOW) + todo = {"todo_id": "todo_peer", "goal_id": "test-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"], + "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 + + +@pytest.mark.parametrize("kwargs,expected", [ + ({"has_todo": False}, "registered"), + ({"has_todo": False, "binding": True}, "addressable"), + ({"status": "done"}, "registered"), + ({"status": "done", "binding": True}, "addressable"), + ({}, "launchable"), + ({"binding": True}, "bound"), + ({"age": 1}, "executing"), + ({"age": 8}, "executing"), + ({"age": 8.01}, "launchable"), + ({"age": 9, "binding": True}, "bound"), + ({"age": -1}, "launchable"), + ({"age": 48}, "launchable"), + ({"status": "blocked", "age": 1, "binding": True}, "blocked"), + ({"task_class": "blocker", "age": 1}, "blocked"), + ({"task_class": "continuous_monitor", "age": 1}, "monitoring"), + ({"task_class": "continuous_monitor", "age": 48}, "monitoring"), + ({"status": "deferred", "age": 1}, "waiting"), +]) +def test_projected_state(monkeypatch, kwargs, expected): + packet, _ = build_projection(monkeypatch, **kwargs) + row = packet["agents"][0] + assert row["state"] == expected + assert "lifecycle_state" not in row + assert ("session_binding" in row) == kwargs.get("binding", False) + assert packet["truth_contract"]["projection_is_writable"] is False + + +def test_unrelated_blocked_activity_does_not_change_current_work(monkeypatch): + other = {"todo_id": "todo_blocked", "goal_id": "test-goal", "role": "agent", + "claimed_by": "peer", "status": "blocked", "task_class": "blocker", + "updated_at": NOW.isoformat()} + packet, _ = build_projection(monkeypatch, extra_todos=[other]) + row = packet["agents"][0] + assert row["current_todo"]["todo_id"] == "todo_peer" + assert row["blocked_on"]["todo_id"] == "todo_blocked" + 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"), +]) +def test_real_projection_to_peer_admission(monkeypatch, age, binding, capability, + resume_ready, reason): + packet, todo = build_projection(monkeypatch, age=age, binding=binding) + todo.update(resume_when="todo_done:todo_dependency", resume_ready=resume_ready) + summary = {"items": [todo]} + contract, lane = apply_task_orchestration_contract( + fallback_work_lane_contract={"lane": "advancement_task"}, + goal_boundary={"peer_task_coordination": { + "enabled": True, "coordinator_agent_id": "coordinator"}}, + agent_identity={"agent_id": "coordinator", "registered_agents": ["coordinator", "peer"]}, + agent_todo_summary=summary, raw_agent_todo_summary=summary, + available_capabilities=["peer_agent_activation"] if capability else [], + agent_management_projection=packet, + ) + assert contract is not None + assert contract["execution_state"] == ("blocked" if reason else "ready") + if reason: + assert contract["eligible_peer_lanes"] == [] + assert contract["blocked_peer_lanes"][0]["reason_codes"] == [reason] + else: + assert contract["eligible_peer_lanes"][0]["todo_id"] == "todo_peer" + assert lane["lane"] == "task_orchestration"