Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
23 changes: 20 additions & 3 deletions examples/worker-lifecycle-state-smoke.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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,
Expand Down Expand Up @@ -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", [])}

Expand Down
45 changes: 41 additions & 4 deletions loopx/control_plane/agents/directory.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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."""

Expand All @@ -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),
Expand Down Expand Up @@ -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.

Expand All @@ -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 = [
Expand All @@ -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
Expand Down
Loading
Loading