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
62 changes: 43 additions & 19 deletions loopx/control_plane/agents/management_projection.py
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,7 @@
_TODO_GROUP_LIST_KEYS = tuple(
dict.fromkeys(
(
"items",
*TODO_SUMMARY_SOURCE_KEYS,
"executable_backlog_items",
"deferred_items",
Expand Down Expand Up @@ -80,7 +81,9 @@ def _todo_status(todo: dict[str, Any]) -> str:


def _is_done(todo: dict[str, Any]) -> bool:
return bool(todo.get("done")) or _todo_status(todo) in {"done", "archive", "archived"}
# A deferred Todo can carry a checked legacy display flag. Explicit state
# wins; the checkbox fallback belongs only to rows without a status.
return _todo_status(todo) in {"done", "archive", "archived"}



Expand Down Expand Up @@ -314,35 +317,42 @@ def _iter_next_action_todo(

def _iter_status_todos(status_payload: dict[str, Any]) -> Iterable[dict[str, Any]]:
queue = _as_dict(status_payload.get("attention_queue"))
todo_index = _as_dict(status_payload.get("todo_index"))
unavailable_goals = set(_as_list(todo_index.get("unavailable_goal_ids")))
unavailable_goals.update(
item.get("goal_id") for item in _as_list(queue.get("items"))
if isinstance(item, dict) and item.get("todo_source") == "unavailable"
)
candidates: list[dict[str, Any]] = []
hints: list[dict[str, Any]] = []
for item in _as_list(queue.get("items")):
if not isinstance(item, dict):
continue
goal_id = _compact(item.get("goal_id"), limit=180)
yield from _iter_next_action_todo(
hints.extend(_iter_next_action_todo(
item,
goal_id=goal_id,
source="attention_queue.agent_lane_next_action",
)
))
group = _as_dict(item.get("agent_todos"))
yield from _iter_todo_group_items(
candidates.extend(_iter_todo_group_items(
group,
goal_id=goal_id,
source="attention_queue.agent_todos",
)
))
project_asset = _as_dict(item.get("project_asset"))
yield from _iter_next_action_todo(
hints.extend(_iter_next_action_todo(
project_asset,
goal_id=goal_id,
source="project_asset.agent_lane_next_action",
)
))
group = _as_dict(project_asset.get("agent_todos"))
yield from _iter_todo_group_items(
candidates.extend(_iter_todo_group_items(
group,
goal_id=goal_id,
source="project_asset.agent_todos",
)
))

todo_index = _as_dict(status_payload.get("todo_index"))
for todo in _as_list(todo_index.get("items")):
if (
isinstance(todo, dict)
Expand All @@ -351,7 +361,25 @@ def _iter_status_todos(status_payload: dict[str, Any]) -> Iterable[dict[str, Any
):
row = dict(todo)
row.setdefault("source", "todo_index")
yield row
candidates.append(row)

# Display text, ordinals and Next Action hints are not Task identity. Keep
# current source rows before their shorter aliases and historical hints;
# both management and its execution-facts reader consume this one view.
selected = {_todo_identity(todo): todo for todo in hints}
seen: set[tuple[str, str, str, str]] = set()
for row in [*candidates, *hints]:
if row.get("goal_id") in unavailable_goals:
continue
identity = _todo_identity(row)
if identity in seen:
continue
seen.add(identity)
hint = selected.get(identity)
if hint and hint.get("selected_by"):
row = dict(row)
row.setdefault("selected_by", hint["selected_by"])
yield row


def _todo_agent_id(todo: dict[str, Any]) -> str | None:
Expand All @@ -362,11 +390,12 @@ def _todo_agent_id(todo: dict[str, Any]) -> str | None:


def _todo_identity(todo: dict[str, Any]) -> tuple[str, str, str, str]:
todo_id = str(todo.get("todo_id") or "").strip()
return (
str(todo.get("goal_id") or ""),
str(todo.get("todo_id") or ""),
str(todo.get("index") or ""),
str(todo.get("text") or todo.get("title") or ""),
todo_id,
"" if todo_id else str(todo.get("index") or ""),
"" if todo_id else str(todo.get("text") or todo.get("title") or ""),
)


Expand Down Expand Up @@ -713,15 +742,10 @@ def build_agent_management_projection(
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):
agent_id = _todo_agent_id(todo)
if not agent_id:
continue
identity = _todo_identity(todo)
if identity in seen_todos:
continue
seen_todos.add(identity)
row = rows_by_agent.setdefault(
agent_id,
{
Expand Down
13 changes: 8 additions & 5 deletions loopx/control_plane/todos/todo_index.py
Original file line number Diff line number Diff line change
Expand Up @@ -230,11 +230,14 @@ def build_todo_index(
existing["latest_event_kind"] = latest_kind
existing["latest_event_at"] = event_item.get("latest_event_at")
existing["latest_event_status"] = event_item.get("latest_event_status")
if event_item.get("status"):
existing["status"] = event_item.get("status")
existing["done"] = bool(event_item.get("done"))
if event_item.get("agent_id"):
existing["agent_id"] = event_item.get("agent_id")
# Audit receipts describe historical operations. They cannot
# replace the current authority's status or mint a claimant.
if existing.get("source") == "rollout_event_log":
if event_item.get("status"):
existing["status"] = event_item.get("status")
existing["done"] = bool(event_item.get("done"))
if event_item.get("agent_id"):
existing["agent_id"] = event_item.get("agent_id")
# The audit sentence describes the newest event for every row
# kind; only event-only rows also carry `title_source`, and the
# text of an attention-queue row stays authoritative because
Expand Down
54 changes: 54 additions & 0 deletions tests/control_plane/test_peer_agent_directory.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,12 @@

from __future__ import annotations

import pytest

from loopx.control_plane.agents.management_projection import (
build_agent_management_projection,
projected_agent_goals,
)
from loopx.control_plane.agents.directory import (
GAP_AUDIENCE_NOT_AUTHORIZED,
LIMITATION_CALLER_IDENTITY_NOT_SUPPLIED,
Expand Down Expand Up @@ -240,3 +246,51 @@ def test_the_published_candidate_list_respects_its_budget() -> None:
"thread-2",
]
assert "withheld_candidate_count" not in route


@pytest.mark.parametrize("status", ["blocked", "deferred", "done"])
def test_canonical_identity_wins_over_stale_display_variants(status: str) -> None:
payload = _status_payload([WORKING_AGENT])
current = dict(payload["todo_index"]["items"][0], status=status, done=status in ("deferred", "done"))
stale = dict(current, status="open", done=False, index=99,
text="Old shorter display", selected_by="current_agent_claimed_todo")
payload["attention_queue"] = {"items": [{"goal_id": GOAL_ID,
"agent_todos": {"items": [current], "first_open_items": [stale]},
"agent_lane_next_action": stale,
"project_asset": {"agent_todos": {"items": [stale]}}}]}
payload["todo_index"] = {"items": [dict(stale, source="attention_queue")]}

management = build_agent_management_projection(payload)
row = next(row for row in management["agents"] if row["agent_id"] == WORKING_AGENT)
work = _rows_by_agent(build_peer_agent_directory(payload))[WORKING_AGENT]["work"]
if status == "done":
assert "current_todo" not in row
assert work is None
assert projected_agent_goals(payload)[GOAL_ID][WORKING_AGENT]["open_todo_ids"] == []
else:
assert row["current_todo"]["status"] == status
assert work["todo_status"] == status


def test_stale_hint_cannot_restore_a_cleared_claim() -> None:
payload = _status_payload([WORKING_AGENT])
current = dict(payload["todo_index"]["items"][0])
current.pop("claimed_by")
stale = dict(current, claimed_by=WORKING_AGENT, agent_id=WORKING_AGENT)
payload["attention_queue"] = {"items": [{"goal_id": GOAL_ID,
"agent_todos": {"items": [current]}, "agent_lane_next_action": stale}]}
payload["todo_index"] = {"items": []}
assert _rows_by_agent(build_peer_agent_directory(payload))[WORKING_AGENT]["work"] is None
assert projected_agent_goals(payload)[GOAL_ID][WORKING_AGENT]["open_todo_ids"] == []


def test_unavailable_source_does_not_resurrect_project_asset_or_hint() -> None:
payload = _status_payload([WORKING_AGENT])
stale = payload["todo_index"]["items"][0]
payload["attention_queue"] = {"items": [{"goal_id": GOAL_ID,
"todo_source": "unavailable", "agent_lane_next_action": stale,
"project_asset": {"agent_todos": {"items": [stale]}}}]}
payload["todo_index"]["items"] = []
packet = build_peer_agent_directory(payload)
assert _rows_by_agent(packet)[WORKING_AGENT]["work"] is None
assert projected_agent_goals(payload)[GOAL_ID][WORKING_AGENT]["open_todo_ids"] == []
125 changes: 125 additions & 0 deletions tests/control_plane/test_retired_todo_event_source.py
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,33 @@
from loopx.control_plane.testing.canary_harness import run_json_cli_result
from loopx.control_plane.todos import completion_validation
from loopx.todos import complete_goal_todo
from loopx.control_plane.todos.todo_index import build_todo_index
from loopx.control_plane.runtime.public_safety import public_safe_compact_text
from tests.control_plane.canonical_authority_fixture import (
initialize_canonical_authority, isolate_sqlite_runtime,
)
from loopx.control_plane.coordination.runtime_shadow import build_todo_runtime_shadow_projection
from loopx.control_plane.effect_runtime import restart_effect_runtime


@pytest.mark.parametrize("status", ["open", "blocked", "deferred", "done"])
def test_rollout_audit_cannot_overwrite_current_todo_state_or_claim(tmp_path, status):
current = {"todo_id": "todo_current", "text": "Current authority work",
"status": status, "done": status in ("deferred", "done"), "claimed_by": None}
events = [{"goal_id": "goal-a", "todo_id": "todo_current", "event_kind": "todo_update",
"status": "open" if status != "open" else "done", "agent_id": "old-actor",
"summary": "Historical update", "recorded_at": "2026-01-01T00:00:00Z"}]
result = build_todo_index(queue={"items": [{"goal_id": "goal-a",
"agent_todos": {"items": [current]}}]}, history={"goals": [{"id": "goal-a"}]},
runtime_root=tmp_path, public_safe_compact_text=public_safe_compact_text,
events_for_goal=lambda goal_id, **kwargs: events)
row = result["items"][0]
assert row["status"] == status
assert row["done"] is current["done"]
assert row.get("agent_id") is None
assert row["latest_event_status"] == events[0]["status"]
assert row["latest_event_summary"] == "Historical update"
assert row["event_count"] == 1

@pytest.mark.parametrize("alias", [None, "state_event_log", "state_events_file", "event_log"])
@pytest.mark.parametrize("contents", ["{broken\n", '{"event_type":"todo_added"}\n'])
Expand Down Expand Up @@ -175,3 +202,101 @@ def test_http_status_preserves_healthy_goal_when_retired_source_appears(tmp_path
finally:
server.terminate()
server.wait(timeout=10)


@pytest.mark.parametrize("provider", ["file", "sqlite"])
def test_real_directory_lifecycle_uses_canonical_state_despite_stale_audit(tmp_path, monkeypatch, provider):
isolate_sqlite_runtime(tmp_path, monkeypatch)
state = tmp_path / "ACTIVE_GOAL_STATE.md"
state.write_text("# Goal\n\n## Agent Todo\n\n")
runtime = tmp_path / "runtime"
registry = tmp_path / "registry.json"
registry.write_text(json.dumps({"schema_version": 1, "common_runtime_root": str(runtime),
"goals": [{"id": "goal-a", "repo": str(tmp_path), "state_file": str(state),
"domain": "software", "status": "active",
"adapter": {"kind": "read_only_project_map_v0", "status": "connected"},
"coordination": {"registered_agents": ["agent-a", "agent-b", "agent-c"]}}]}))
initialize_canonical_authority(runtime, "goal-a",
build_todo_runtime_shadow_projection(goal_id="goal-a", todos=[], handoff_mode="hard_lease"),
state_path=state, provider=provider)

def cli(*args, success=True):
code, payload = run_json_cli_result(*args, "--goal-id", "goal-a", registry_path=registry)
if success:
assert code == 0 and payload.get("ok") is True, payload
return payload

def readback(expected, lease_status=None, *, agent="agent-a", target=None):
target = target or todo_id
before = cli("todo", "list")
current = next(row for row in before["todos"] if row["todo_id"] == target)
assert current["status"] == expected
payload = cli("status")
indexed = next(row for row in payload["todo_index"]["items"] if row["todo_id"] == target)
assert indexed["status"] == expected
directory = cli("agent-directory", "--agent-id", agent)
work = next(row for row in directory["rows"] if row["agent_id"] == agent)["work"]
if expected == "done":
assert work is None
else:
assert work["todo_id"] == target and work["todo_status"] == expected
assert work["claimed_by"] == agent
assert work.get("lease_status") == lease_status
after = cli("todo", "list")
assert before["authority_read"]["provider_revision"] == after["authority_read"]["provider_revision"]
assert before["todos"] == after["todos"]

try:
added = cli("todo", "add", "--role", "agent", "--text", "Long canonical task " + "work " * 140,
"--claimed-by", "agent-a")
todo_id = added["todo_id"]
lease = cli("task-lease", "acquire", "--todo-id", todo_id, "--owner", "agent-a",
"--idempotency-key", "directory-live-lease")
readback("open", "active")
cli("task-lease", "release", "--todo-id", todo_id, "--owner", "agent-a",
"--idempotency-key", "directory-live-lease", "--expected-version", str(lease["lease"]["version"]))
before = cli("todo", "list")
cli("todo", "update", "--todo-id", todo_id, "--agent-id", "agent-a",
"--status", "blocked", "--reason", "Fixture lifecycle transition",
"--update-operation-id", "directory-blocked", "--clear-resume-when",
"--update-expected-provider-revision", before["authority_read"]["provider_revision"])
readback("blocked", "released")
deferred = cli("todo", "add", "--role", "agent", "--text", "Wait for an explicit decision",
"--status", "deferred", "--resume-when", "todo_done:" + todo_id, "--claimed-by", "agent-b")
readback("deferred", agent="agent-b", target=deferred["todo_id"])
finished = cli("todo", "add", "--role", "agent", "--text", "Validate the terminal readback",
"--claimed-by", "agent-c")
terminal_lease = cli("task-lease", "acquire", "--todo-id", finished["todo_id"], "--owner", "agent-c",
"--idempotency-key", "directory-terminal-lease")
cli("todo", "complete", "--todo-id", finished["todo_id"], "--agent-id", "agent-c",
"--evidence", "validation://directory-lifecycle", "--no-follow-up",
"--task-lease-idempotency-key", "directory-terminal-lease", "--task-lease-expected-version",
str(terminal_lease["lease"]["version"]))
readback("done", agent="agent-c", target=finished["todo_id"])
# A readable stale display and rollout history cannot rescue a lost provider.
state.write_text("## Agent Todo\n- [ ] Old work\n"
f" <!-- loopx:todo todo_id={todo_id} status=open claimed_by=agent-a -->\n")
authority = runtime / "authority" / f"{provider}-v0"
unavailable = runtime / "unavailable-provider"
authority.rename(unavailable)
for command in ("status", "agent-directory"):
failure = subprocess.run([sys.executable, "-m", "loopx.cli", "--registry", str(registry),
"--runtime-root", str(runtime), "--format", "json", command, "--goal-id", "goal-a"],
capture_output=True, text=True, timeout=60)
if failure.stdout.strip():
failed = json.loads(failure.stdout)
if command == "status":
assert not failed["ok"] and not failed.get("todo_index", {}).get("items")
else:
assert all(row["work"] is None for row in failed.get("rows", []))
else:
# Some source-loss paths raise the owning unavailable error.
# That failure is not a successful empty or legacy work view.
assert failure.returncode != 0
assert "LocalCoordinationAuthorityUnavailable:" in failure.stderr
unavailable.rename(authority)
readback("blocked", "released")
readback("deferred", agent="agent-b", target=deferred["todo_id"])
readback("done", agent="agent-c", target=finished["todo_id"])
finally:
restart_effect_runtime()
Loading