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
15 changes: 12 additions & 3 deletions docs/reference/protocols/agent-management-projection-v0.md
Original file line number Diff line number Diff line change
Expand Up @@ -112,9 +112,18 @@ 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;
- `session_binding_candidates` and `session_binding_count`: present together
when the Agent has at least one entry in existing
`run_history.goals[].coordination.thread_agent_bindings`. The list holds up
to three distinct `{thread_id, host_surface}` candidates in first-seen order
and `session_binding_count` is the number of distinct bindings observed, so a
shorter list than the count means candidates were capped, not that the
remainder is invalid. Distinctness is decided on the full binding identity, not
on the rendered text: the owner accepts wider identifiers than this projection
displays, so two counted bindings can render alike. Neither field selects an execution route: absence is
not evidence of a stopped host, a binding does not prove executable capacity,
and several candidates for one Agent are resolved by the owning binding
resolver, not by this display projection;
- `stale_claim_hint`;
- `blocked_on`;
- `recent_events`;
Expand Down
64 changes: 43 additions & 21 deletions loopx/control_plane/agents/management_projection.py
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@
normalize_todo_id,
)
from ..todos.summary_item import TODO_SUMMARY_SOURCE_KEYS
from ...thread_agent_binding import collect_accepted_bindings
from .material_frontier import AGENT_MATERIAL_FRONTIER_SCHEMA_VERSION
from .material_handoff import (
build_material_handoff_note_v1,
Expand All @@ -25,6 +26,7 @@
MAX_AGENT_ROWS = 24
MAX_AGENT_TODOS = 8
MAX_REFS = 1
MAX_SESSION_BINDING_CANDIDATES = 3
MAX_WORKSPACE_SCOPES = 4
STALE_CLAIM_THRESHOLD_HOURS = 36
EXECUTING_ACTIVITY_THRESHOLD_HOURS = 8
Expand Down Expand Up @@ -559,6 +561,35 @@ def _safe_next_action(todo: dict[str, Any] | None) -> str:
return "Inspect projected todo."


def _collect_session_binding_candidates(
status_payload: dict[str, Any],
) -> dict[str, list[dict[str, str]]]:
"""Group one Agent's accepted session bindings by the owner's agent identity.

Session bindings come from run_history.coordination.thread_agent_bindings
and are the only source for addressable/bound lifecycle states. They are
read through the binding owner's `collect_accepted_bindings`, so this row
counts the same bindings the peer directory publishes as routes instead of
normalising a second time here; that also keys the candidates by the same
agent identity this projection uses for its rows. The display budget is not
an identity key: the owner accepts thread identifiers to 128 characters and
host surfaces to 64, both wider than what this projection shows, so two
accepted bindings that share a visible prefix are still two routes.
Identity stays whole, and `_compact` only renders it.
"""

candidates: dict[str, list[dict[str, str]]] = {}
run_history = _as_dict(status_payload.get("run_history"))
for binding in collect_accepted_bindings(_as_list(run_history.get("goals"))):
candidates.setdefault(binding["agent_id"], []).append(
{
"thread_id": _compact(binding["thread_id"], limit=120),
"host_surface": _compact(binding["host_surface"], limit=60),
}
)
return candidates


def build_agent_management_projection(
status_payload: dict[str, Any],
*,
Expand All @@ -583,24 +614,11 @@ def build_agent_management_projection(
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",
}
# This is the only source for addressable/bound lifecycle states. One Agent
# can hold several historical bindings, and an omitted one is not disproved:
# 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)

seen_todos: set[tuple[str, str, str, str]] = set()
for todo in _iter_status_todos(status_payload):
Expand Down Expand Up @@ -652,7 +670,7 @@ def build_agent_management_projection(
agent_state = _agent_state(
all_todos,
current=current,
has_session_binding=agent_id in session_bindings,
has_session_binding=agent_id in session_binding_candidates,
last_activity_at=_last_activity([current]) if current else None,
)
agent_row: dict[str, Any] = {
Expand All @@ -666,8 +684,12 @@ def build_agent_management_projection(
"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]
binding_candidates = session_binding_candidates.get(agent_id)
if binding_candidates:
agent_row["session_binding_candidates"] = binding_candidates[
:MAX_SESSION_BINDING_CANDIDATES
]
agent_row["session_binding_count"] = len(binding_candidates)
material_frontier_key = _agent_material_frontier_key(
raw_row=raw_row,
current=current,
Expand Down
2 changes: 1 addition & 1 deletion tests/control_plane/test_agent_lifecycle_state.py
Original file line number Diff line number Diff line change
Expand Up @@ -51,7 +51,7 @@ def test_projected_state(monkeypatch, kwargs, expected):
row = packet["agents"][0]
assert row["state"] == expected
assert "lifecycle_state" not in row
assert ("session_binding" in row) == kwargs.get("binding", False)
assert ("session_binding_candidates" in row) == kwargs.get("binding", False)
assert packet["truth_contract"]["projection_is_writable"] is False


Expand Down
221 changes: 221 additions & 0 deletions tests/control_plane/test_agent_session_binding_candidates.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,221 @@
"""Pin the session-binding candidate summary on an Agent management row.

A peer can keep several historical thread bindings. The row must carry them as
a bounded summary plus the observed count instead of letting the binding walked
last replace the earlier ones: one surviving row invites a reader to treat it as
the only route to that peer, which is how a live peer review request got
substituted for a temporary host sub-agent.
"""

from __future__ import annotations

from datetime import datetime, timezone
from typing import Any

import pytest

from loopx.control_plane.agents import management_projection as projection

NOW = datetime(2026, 9, 25, 12, tzinfo=timezone.utc)


def _binding(
thread_id: str,
*,
agent_id: str = "peer",
host_surface: str = "codex-app",
) -> dict[str, str]:
return {
"agent_id": agent_id,
"thread_id": thread_id,
"host_surface": host_surface,
}


def _row(
monkeypatch: pytest.MonkeyPatch,
goals: list[list[dict[str, str]]],
) -> dict[str, Any]:
monkeypatch.setattr(projection, "now_utc", lambda: NOW)
payload: dict[str, Any] = {
"goal_filter": "test-goal",
"run_history": {
"goals": [
{
"id": f"test-goal-{index}",
"coordination": {
"registered_agents": ["peer", "reviewer"],
"thread_agent_bindings": bindings,
},
}
for index, bindings in enumerate(goals)
]
},
"todo_index": {"items": []},
}
packet = projection.build_agent_management_projection(payload)
rows = {row["agent_id"]: row for row in packet["agents"]}
assert "peer" in rows
return rows["peer"]


def test_a_later_binding_does_not_replace_an_earlier_one(
monkeypatch: pytest.MonkeyPatch,
) -> None:
row = _row(
monkeypatch,
[[_binding("thread-old"), _binding("thread-new", host_surface="codex-cli")]],
)

assert row["session_binding_count"] == 2
assert row["session_binding_candidates"] == [
{"thread_id": "thread-old", "host_surface": "codex-app"},
{"thread_id": "thread-new", "host_surface": "codex-cli"},
]
assert "session_binding" not in row


def test_a_republished_binding_counts_once(
monkeypatch: pytest.MonkeyPatch,
) -> None:
row = _row(
monkeypatch,
[
[_binding("thread-shared")],
[_binding("thread-shared"), _binding("thread-extra")],
],
)

assert row["session_binding_count"] == 2
assert [item["thread_id"] for item in row["session_binding_candidates"]] == [
"thread-shared",
"thread-extra",
]


def test_the_candidate_list_is_bounded_without_losing_the_count(
monkeypatch: pytest.MonkeyPatch,
) -> None:
bindings = [_binding(f"thread-{index}") for index in range(5)]

row = _row(monkeypatch, [bindings])

assert row["session_binding_count"] == 5
assert len(row["session_binding_candidates"]) == (
projection.MAX_SESSION_BINDING_CANDIDATES
)
assert [item["thread_id"] for item in row["session_binding_candidates"]] == [
"thread-0",
"thread-1",
"thread-2",
]


def test_candidates_do_not_leak_across_agents(
monkeypatch: pytest.MonkeyPatch,
) -> None:
row = _row(
monkeypatch,
[[_binding("thread-peer"), _binding("thread-reviewer", agent_id="reviewer")]],
)

assert row["session_binding_count"] == 1
assert row["session_binding_candidates"] == [
{"thread_id": "thread-peer", "host_surface": "codex-app"}
]


def test_several_candidates_keep_the_unbound_lifecycle_state(
monkeypatch: pytest.MonkeyPatch,
) -> None:
single = _row(monkeypatch, [[_binding("thread-one")]])
several = _row(monkeypatch, [[_binding("thread-one"), _binding("thread-two")]])
none = _row(monkeypatch, [[]])

assert single["state"] == "addressable"
assert several["state"] == single["state"]
assert none["state"] == "registered"
assert "session_binding_candidates" not in none


def test_two_long_thread_ids_sharing_a_visible_prefix_are_two_candidates(
monkeypatch: pytest.MonkeyPatch,
) -> None:
"""The display budget is 120 characters; the owner accepts 128.

Two identifiers that differ only past the visible edge collapse into the
same rendered text, so counting rendered text would report one route where
the registry holds two.
"""

prefix = "t" * 120
row = _row(
monkeypatch,
[
[
_binding(prefix + "AAA"),
_binding(prefix + "BBB"),
]
],
)

assert row["session_binding_count"] == 2
assert len(row["session_binding_candidates"]) == 2
assert (
row["session_binding_candidates"][0]["thread_id"]
== row["session_binding_candidates"][1]["thread_id"]
)
# Counting uses the full identity; rendering stays inside the display budget.
assert len(row["session_binding_candidates"][0]["thread_id"]) == 120


def test_two_long_host_surfaces_sharing_a_visible_prefix_are_two_candidates(
monkeypatch: pytest.MonkeyPatch,
) -> None:
row = _row(
monkeypatch,
[
[
_binding("thread-one", host_surface="s" * 60 + "X"),
_binding("thread-one", host_surface="s" * 60 + "Y"),
]
],
)

assert row["session_binding_count"] == 2
assert len(row["session_binding_candidates"]) == 2


def test_an_identical_binding_republished_under_a_long_id_still_counts_once(
monkeypatch: pytest.MonkeyPatch,
) -> None:
prefix = "t" * 120
row = _row(
monkeypatch,
[
[_binding(prefix + "AAA")],
[_binding(prefix + "AAA"), _binding(prefix + "BBB")],
],
)

assert row["session_binding_count"] == 2


def test_padding_a_thread_id_does_not_create_a_second_route(
monkeypatch: pytest.MonkeyPatch,
) -> None:
"""The owner strips a thread identifier before matching it.

Keying the count on the untouched value would therefore split one binding
into two candidates and offer a route that no resolver would confirm.
"""

row = _row(
monkeypatch,
[[_binding("thread-one")], [_binding(" thread-one ")]],
)

assert row["session_binding_count"] == 1
assert row["session_binding_candidates"] == [
{"thread_id": "thread-one", "host_surface": "codex-app"}
]
Loading