Skip to content

Commit cdfe146

Browse files
authored
fix(status): preserve canonical peer work facts over audit rows (#5668)
Keep canonical current Todo status and claimant facts ahead of historical audit events and display aliases. Reconcile stable Goal/Todo identity once for management and execution-facts consumers; preserve explicit deferred state over compatibility checkboxes. Validated the exact PR head with 34 focused and 100 adjacent tests, real File/SQLite CLI source recovery, and full native premerge. The current-main source candidate also passes all 134 regression tests. Installed Bot/GUI/Host adoption is separate. Signed-off-by: LoopX Agent <337587101+loopx-agent@users.noreply.github.com>
1 parent d0333d8 commit cdfe146

4 files changed

Lines changed: 230 additions & 24 deletions

File tree

‎loopx/control_plane/agents/management_projection.py‎

Lines changed: 43 additions & 19 deletions
Original file line numberDiff line numberDiff line change
@@ -34,6 +34,7 @@
3434
_TODO_GROUP_LIST_KEYS = tuple(
3535
dict.fromkeys(
3636
(
37+
"items",
3738
*TODO_SUMMARY_SOURCE_KEYS,
3839
"executable_backlog_items",
3940
"deferred_items",
@@ -80,7 +81,9 @@ def _todo_status(todo: dict[str, Any]) -> str:
8081

8182

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

8588

8689

@@ -314,35 +317,42 @@ def _iter_next_action_todo(
314317

315318
def _iter_status_todos(status_payload: dict[str, Any]) -> Iterable[dict[str, Any]]:
316319
queue = _as_dict(status_payload.get("attention_queue"))
320+
todo_index = _as_dict(status_payload.get("todo_index"))
321+
unavailable_goals = set(_as_list(todo_index.get("unavailable_goal_ids")))
322+
unavailable_goals.update(
323+
item.get("goal_id") for item in _as_list(queue.get("items"))
324+
if isinstance(item, dict) and item.get("todo_source") == "unavailable"
325+
)
326+
candidates: list[dict[str, Any]] = []
327+
hints: list[dict[str, Any]] = []
317328
for item in _as_list(queue.get("items")):
318329
if not isinstance(item, dict):
319330
continue
320331
goal_id = _compact(item.get("goal_id"), limit=180)
321-
yield from _iter_next_action_todo(
332+
hints.extend(_iter_next_action_todo(
322333
item,
323334
goal_id=goal_id,
324335
source="attention_queue.agent_lane_next_action",
325-
)
336+
))
326337
group = _as_dict(item.get("agent_todos"))
327-
yield from _iter_todo_group_items(
338+
candidates.extend(_iter_todo_group_items(
328339
group,
329340
goal_id=goal_id,
330341
source="attention_queue.agent_todos",
331-
)
342+
))
332343
project_asset = _as_dict(item.get("project_asset"))
333-
yield from _iter_next_action_todo(
344+
hints.extend(_iter_next_action_todo(
334345
project_asset,
335346
goal_id=goal_id,
336347
source="project_asset.agent_lane_next_action",
337-
)
348+
))
338349
group = _as_dict(project_asset.get("agent_todos"))
339-
yield from _iter_todo_group_items(
350+
candidates.extend(_iter_todo_group_items(
340351
group,
341352
goal_id=goal_id,
342353
source="project_asset.agent_todos",
343-
)
354+
))
344355

345-
todo_index = _as_dict(status_payload.get("todo_index"))
346356
for todo in _as_list(todo_index.get("items")):
347357
if (
348358
isinstance(todo, dict)
@@ -351,7 +361,25 @@ def _iter_status_todos(status_payload: dict[str, Any]) -> Iterable[dict[str, Any
351361
):
352362
row = dict(todo)
353363
row.setdefault("source", "todo_index")
354-
yield row
364+
candidates.append(row)
365+
366+
# Display text, ordinals and Next Action hints are not Task identity. Keep
367+
# current source rows before their shorter aliases and historical hints;
368+
# both management and its execution-facts reader consume this one view.
369+
selected = {_todo_identity(todo): todo for todo in hints}
370+
seen: set[tuple[str, str, str, str]] = set()
371+
for row in [*candidates, *hints]:
372+
if row.get("goal_id") in unavailable_goals:
373+
continue
374+
identity = _todo_identity(row)
375+
if identity in seen:
376+
continue
377+
seen.add(identity)
378+
hint = selected.get(identity)
379+
if hint and hint.get("selected_by"):
380+
row = dict(row)
381+
row.setdefault("selected_by", hint["selected_by"])
382+
yield row
355383

356384

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

363391

364392
def _todo_identity(todo: dict[str, Any]) -> tuple[str, str, str, str]:
393+
todo_id = str(todo.get("todo_id") or "").strip()
365394
return (
366395
str(todo.get("goal_id") or ""),
367-
str(todo.get("todo_id") or ""),
368-
str(todo.get("index") or ""),
369-
str(todo.get("text") or todo.get("title") or ""),
396+
todo_id,
397+
"" if todo_id else str(todo.get("index") or ""),
398+
"" if todo_id else str(todo.get("text") or todo.get("title") or ""),
370399
)
371400

372401

@@ -713,15 +742,10 @@ def build_agent_management_projection(
713742
session_binding_candidates = _collect_session_binding_candidates(status_payload)
714743
facts_by_agent = execution_facts if isinstance(execution_facts, dict) else {}
715744

716-
seen_todos: set[tuple[str, str, str, str]] = set()
717745
for todo in _iter_status_todos(status_payload):
718746
agent_id = _todo_agent_id(todo)
719747
if not agent_id:
720748
continue
721-
identity = _todo_identity(todo)
722-
if identity in seen_todos:
723-
continue
724-
seen_todos.add(identity)
725749
row = rows_by_agent.setdefault(
726750
agent_id,
727751
{

‎loopx/control_plane/todos/todo_index.py‎

Lines changed: 8 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -230,11 +230,14 @@ def build_todo_index(
230230
existing["latest_event_kind"] = latest_kind
231231
existing["latest_event_at"] = event_item.get("latest_event_at")
232232
existing["latest_event_status"] = event_item.get("latest_event_status")
233-
if event_item.get("status"):
234-
existing["status"] = event_item.get("status")
235-
existing["done"] = bool(event_item.get("done"))
236-
if event_item.get("agent_id"):
237-
existing["agent_id"] = event_item.get("agent_id")
233+
# Audit receipts describe historical operations. They cannot
234+
# replace the current authority's status or mint a claimant.
235+
if existing.get("source") == "rollout_event_log":
236+
if event_item.get("status"):
237+
existing["status"] = event_item.get("status")
238+
existing["done"] = bool(event_item.get("done"))
239+
if event_item.get("agent_id"):
240+
existing["agent_id"] = event_item.get("agent_id")
238241
# The audit sentence describes the newest event for every row
239242
# kind; only event-only rows also carry `title_source`, and the
240243
# text of an attention-queue row stays authoritative because

‎tests/control_plane/test_peer_agent_directory.py‎

Lines changed: 54 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,12 @@
22

33
from __future__ import annotations
44

5+
import pytest
6+
7+
from loopx.control_plane.agents.management_projection import (
8+
build_agent_management_projection,
9+
projected_agent_goals,
10+
)
511
from loopx.control_plane.agents.directory import (
612
GAP_AUDIENCE_NOT_AUTHORIZED,
713
LIMITATION_CALLER_IDENTITY_NOT_SUPPLIED,
@@ -240,3 +246,51 @@ def test_the_published_candidate_list_respects_its_budget() -> None:
240246
"thread-2",
241247
]
242248
assert "withheld_candidate_count" not in route
249+
250+
251+
@pytest.mark.parametrize("status", ["blocked", "deferred", "done"])
252+
def test_canonical_identity_wins_over_stale_display_variants(status: str) -> None:
253+
payload = _status_payload([WORKING_AGENT])
254+
current = dict(payload["todo_index"]["items"][0], status=status, done=status in ("deferred", "done"))
255+
stale = dict(current, status="open", done=False, index=99,
256+
text="Old shorter display", selected_by="current_agent_claimed_todo")
257+
payload["attention_queue"] = {"items": [{"goal_id": GOAL_ID,
258+
"agent_todos": {"items": [current], "first_open_items": [stale]},
259+
"agent_lane_next_action": stale,
260+
"project_asset": {"agent_todos": {"items": [stale]}}}]}
261+
payload["todo_index"] = {"items": [dict(stale, source="attention_queue")]}
262+
263+
management = build_agent_management_projection(payload)
264+
row = next(row for row in management["agents"] if row["agent_id"] == WORKING_AGENT)
265+
work = _rows_by_agent(build_peer_agent_directory(payload))[WORKING_AGENT]["work"]
266+
if status == "done":
267+
assert "current_todo" not in row
268+
assert work is None
269+
assert projected_agent_goals(payload)[GOAL_ID][WORKING_AGENT]["open_todo_ids"] == []
270+
else:
271+
assert row["current_todo"]["status"] == status
272+
assert work["todo_status"] == status
273+
274+
275+
def test_stale_hint_cannot_restore_a_cleared_claim() -> None:
276+
payload = _status_payload([WORKING_AGENT])
277+
current = dict(payload["todo_index"]["items"][0])
278+
current.pop("claimed_by")
279+
stale = dict(current, claimed_by=WORKING_AGENT, agent_id=WORKING_AGENT)
280+
payload["attention_queue"] = {"items": [{"goal_id": GOAL_ID,
281+
"agent_todos": {"items": [current]}, "agent_lane_next_action": stale}]}
282+
payload["todo_index"] = {"items": []}
283+
assert _rows_by_agent(build_peer_agent_directory(payload))[WORKING_AGENT]["work"] is None
284+
assert projected_agent_goals(payload)[GOAL_ID][WORKING_AGENT]["open_todo_ids"] == []
285+
286+
287+
def test_unavailable_source_does_not_resurrect_project_asset_or_hint() -> None:
288+
payload = _status_payload([WORKING_AGENT])
289+
stale = payload["todo_index"]["items"][0]
290+
payload["attention_queue"] = {"items": [{"goal_id": GOAL_ID,
291+
"todo_source": "unavailable", "agent_lane_next_action": stale,
292+
"project_asset": {"agent_todos": {"items": [stale]}}}]}
293+
payload["todo_index"]["items"] = []
294+
packet = build_peer_agent_directory(payload)
295+
assert _rows_by_agent(packet)[WORKING_AGENT]["work"] is None
296+
assert projected_agent_goals(payload)[GOAL_ID][WORKING_AGENT]["open_todo_ids"] == []

‎tests/control_plane/test_retired_todo_event_source.py‎

Lines changed: 125 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,33 @@
1616
from loopx.control_plane.testing.canary_harness import run_json_cli_result
1717
from loopx.control_plane.todos import completion_validation
1818
from loopx.todos import complete_goal_todo
19+
from loopx.control_plane.todos.todo_index import build_todo_index
20+
from loopx.control_plane.runtime.public_safety import public_safe_compact_text
21+
from tests.control_plane.canonical_authority_fixture import (
22+
initialize_canonical_authority, isolate_sqlite_runtime,
23+
)
24+
from loopx.control_plane.coordination.runtime_shadow import build_todo_runtime_shadow_projection
25+
from loopx.control_plane.effect_runtime import restart_effect_runtime
26+
27+
28+
@pytest.mark.parametrize("status", ["open", "blocked", "deferred", "done"])
29+
def test_rollout_audit_cannot_overwrite_current_todo_state_or_claim(tmp_path, status):
30+
current = {"todo_id": "todo_current", "text": "Current authority work",
31+
"status": status, "done": status in ("deferred", "done"), "claimed_by": None}
32+
events = [{"goal_id": "goal-a", "todo_id": "todo_current", "event_kind": "todo_update",
33+
"status": "open" if status != "open" else "done", "agent_id": "old-actor",
34+
"summary": "Historical update", "recorded_at": "2026-01-01T00:00:00Z"}]
35+
result = build_todo_index(queue={"items": [{"goal_id": "goal-a",
36+
"agent_todos": {"items": [current]}}]}, history={"goals": [{"id": "goal-a"}]},
37+
runtime_root=tmp_path, public_safe_compact_text=public_safe_compact_text,
38+
events_for_goal=lambda goal_id, **kwargs: events)
39+
row = result["items"][0]
40+
assert row["status"] == status
41+
assert row["done"] is current["done"]
42+
assert row.get("agent_id") is None
43+
assert row["latest_event_status"] == events[0]["status"]
44+
assert row["latest_event_summary"] == "Historical update"
45+
assert row["event_count"] == 1
1946

2047
@pytest.mark.parametrize("alias", [None, "state_event_log", "state_events_file", "event_log"])
2148
@pytest.mark.parametrize("contents", ["{broken\n", '{"event_type":"todo_added"}\n'])
@@ -175,3 +202,101 @@ def test_http_status_preserves_healthy_goal_when_retired_source_appears(tmp_path
175202
finally:
176203
server.terminate()
177204
server.wait(timeout=10)
205+
206+
207+
@pytest.mark.parametrize("provider", ["file", "sqlite"])
208+
def test_real_directory_lifecycle_uses_canonical_state_despite_stale_audit(tmp_path, monkeypatch, provider):
209+
isolate_sqlite_runtime(tmp_path, monkeypatch)
210+
state = tmp_path / "ACTIVE_GOAL_STATE.md"
211+
state.write_text("# Goal\n\n## Agent Todo\n\n")
212+
runtime = tmp_path / "runtime"
213+
registry = tmp_path / "registry.json"
214+
registry.write_text(json.dumps({"schema_version": 1, "common_runtime_root": str(runtime),
215+
"goals": [{"id": "goal-a", "repo": str(tmp_path), "state_file": str(state),
216+
"domain": "software", "status": "active",
217+
"adapter": {"kind": "read_only_project_map_v0", "status": "connected"},
218+
"coordination": {"registered_agents": ["agent-a", "agent-b", "agent-c"]}}]}))
219+
initialize_canonical_authority(runtime, "goal-a",
220+
build_todo_runtime_shadow_projection(goal_id="goal-a", todos=[], handoff_mode="hard_lease"),
221+
state_path=state, provider=provider)
222+
223+
def cli(*args, success=True):
224+
code, payload = run_json_cli_result(*args, "--goal-id", "goal-a", registry_path=registry)
225+
if success:
226+
assert code == 0 and payload.get("ok") is True, payload
227+
return payload
228+
229+
def readback(expected, lease_status=None, *, agent="agent-a", target=None):
230+
target = target or todo_id
231+
before = cli("todo", "list")
232+
current = next(row for row in before["todos"] if row["todo_id"] == target)
233+
assert current["status"] == expected
234+
payload = cli("status")
235+
indexed = next(row for row in payload["todo_index"]["items"] if row["todo_id"] == target)
236+
assert indexed["status"] == expected
237+
directory = cli("agent-directory", "--agent-id", agent)
238+
work = next(row for row in directory["rows"] if row["agent_id"] == agent)["work"]
239+
if expected == "done":
240+
assert work is None
241+
else:
242+
assert work["todo_id"] == target and work["todo_status"] == expected
243+
assert work["claimed_by"] == agent
244+
assert work.get("lease_status") == lease_status
245+
after = cli("todo", "list")
246+
assert before["authority_read"]["provider_revision"] == after["authority_read"]["provider_revision"]
247+
assert before["todos"] == after["todos"]
248+
249+
try:
250+
added = cli("todo", "add", "--role", "agent", "--text", "Long canonical task " + "work " * 140,
251+
"--claimed-by", "agent-a")
252+
todo_id = added["todo_id"]
253+
lease = cli("task-lease", "acquire", "--todo-id", todo_id, "--owner", "agent-a",
254+
"--idempotency-key", "directory-live-lease")
255+
readback("open", "active")
256+
cli("task-lease", "release", "--todo-id", todo_id, "--owner", "agent-a",
257+
"--idempotency-key", "directory-live-lease", "--expected-version", str(lease["lease"]["version"]))
258+
before = cli("todo", "list")
259+
cli("todo", "update", "--todo-id", todo_id, "--agent-id", "agent-a",
260+
"--status", "blocked", "--reason", "Fixture lifecycle transition",
261+
"--update-operation-id", "directory-blocked", "--clear-resume-when",
262+
"--update-expected-provider-revision", before["authority_read"]["provider_revision"])
263+
readback("blocked", "released")
264+
deferred = cli("todo", "add", "--role", "agent", "--text", "Wait for an explicit decision",
265+
"--status", "deferred", "--resume-when", "todo_done:" + todo_id, "--claimed-by", "agent-b")
266+
readback("deferred", agent="agent-b", target=deferred["todo_id"])
267+
finished = cli("todo", "add", "--role", "agent", "--text", "Validate the terminal readback",
268+
"--claimed-by", "agent-c")
269+
terminal_lease = cli("task-lease", "acquire", "--todo-id", finished["todo_id"], "--owner", "agent-c",
270+
"--idempotency-key", "directory-terminal-lease")
271+
cli("todo", "complete", "--todo-id", finished["todo_id"], "--agent-id", "agent-c",
272+
"--evidence", "validation://directory-lifecycle", "--no-follow-up",
273+
"--task-lease-idempotency-key", "directory-terminal-lease", "--task-lease-expected-version",
274+
str(terminal_lease["lease"]["version"]))
275+
readback("done", agent="agent-c", target=finished["todo_id"])
276+
# A readable stale display and rollout history cannot rescue a lost provider.
277+
state.write_text("## Agent Todo\n- [ ] Old work\n"
278+
f" <!-- loopx:todo todo_id={todo_id} status=open claimed_by=agent-a -->\n")
279+
authority = runtime / "authority" / f"{provider}-v0"
280+
unavailable = runtime / "unavailable-provider"
281+
authority.rename(unavailable)
282+
for command in ("status", "agent-directory"):
283+
failure = subprocess.run([sys.executable, "-m", "loopx.cli", "--registry", str(registry),
284+
"--runtime-root", str(runtime), "--format", "json", command, "--goal-id", "goal-a"],
285+
capture_output=True, text=True, timeout=60)
286+
if failure.stdout.strip():
287+
failed = json.loads(failure.stdout)
288+
if command == "status":
289+
assert not failed["ok"] and not failed.get("todo_index", {}).get("items")
290+
else:
291+
assert all(row["work"] is None for row in failed.get("rows", []))
292+
else:
293+
# Some source-loss paths raise the owning unavailable error.
294+
# That failure is not a successful empty or legacy work view.
295+
assert failure.returncode != 0
296+
assert "LocalCoordinationAuthorityUnavailable:" in failure.stderr
297+
unavailable.rename(authority)
298+
readback("blocked", "released")
299+
readback("deferred", agent="agent-b", target=deferred["todo_id"])
300+
readback("done", agent="agent-c", target=finished["todo_id"])
301+
finally:
302+
restart_effect_runtime()

0 commit comments

Comments
 (0)