From 76c620c9cc5e3cc3e295f06b393b41adfa59dbeb Mon Sep 17 00:00:00 2001 From: JunZ-Leo <100498253+JunZ-Leo@users.noreply.github.com> Date: Fri, 25 Sep 2026 21:59:36 +0800 Subject: [PATCH 1/4] fix(agents): report session binding candidates, not the last one walked `session_bindings` was keyed by Agent and reassigned for every `thread_agent_bindings` entry the goal walk touched, so a peer with two historical bindings published one object under the singular name `session_binding` and dropped the rest without saying so. #5039 records what that cost: a coordinator read the surviving row as the only route and substituted a temporary host sub-agent for an existing peer. The row now carries a first-seen, deduplicated candidate list capped by a named budget plus the distinct count, so a short list reads as a cap and not as a disproved remainder, and no single object remains for a consumer to mistake for the selected route. Membership still drives the lifecycle state, so an Agent's state is unchanged whether it holds one binding or several. The protocol document is updated in the same commit, including the explicit statement that neither field selects an execution route or transfers any authority; `thread_agent_binding.py` remains the resolver owner and this module imports nothing from it. Signed-off-by: JunZ-Leo <100498253+JunZ-Leo@users.noreply.github.com> --- .../agent-management-projection-v0.md | 13 +- .../agents/management_projection.py | 32 ++-- .../test_agent_lifecycle_state.py | 2 +- .../test_agent_session_binding_candidates.py | 138 ++++++++++++++++++ 4 files changed, 171 insertions(+), 14 deletions(-) create mode 100644 tests/control_plane/test_agent_session_binding_candidates.py diff --git a/docs/reference/protocols/agent-management-projection-v0.md b/docs/reference/protocols/agent-management-projection-v0.md index 057acd1a20..5439c40238 100644 --- a/docs/reference/protocols/agent-management-projection-v0.md +++ b/docs/reference/protocols/agent-management-projection-v0.md @@ -112,9 +112,16 @@ 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. 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`; diff --git a/loopx/control_plane/agents/management_projection.py b/loopx/control_plane/agents/management_projection.py index d756766cae..6ebd69e1a1 100644 --- a/loopx/control_plane/agents/management_projection.py +++ b/loopx/control_plane/agents/management_projection.py @@ -25,6 +25,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 @@ -583,8 +584,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]] = {} + # 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: dict[str, list[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): @@ -596,11 +600,15 @@ def build_agent_management_projection( 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", - } + if not agent_id or not thread_id: + continue + candidates = session_binding_candidates.setdefault(agent_id, []) + candidate = { + "thread_id": thread_id, + "host_surface": host_surface or "unknown", + } + if candidate not in candidates: + candidates.append(candidate) seen_todos: set[tuple[str, str, str, str]] = set() for todo in _iter_status_todos(status_payload): @@ -652,7 +660,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] = { @@ -666,8 +674,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, diff --git a/tests/control_plane/test_agent_lifecycle_state.py b/tests/control_plane/test_agent_lifecycle_state.py index 2d5d4cb039..2f82553c32 100644 --- a/tests/control_plane/test_agent_lifecycle_state.py +++ b/tests/control_plane/test_agent_lifecycle_state.py @@ -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 diff --git a/tests/control_plane/test_agent_session_binding_candidates.py b/tests/control_plane/test_agent_session_binding_candidates.py new file mode 100644 index 0000000000..59090cce7c --- /dev/null +++ b/tests/control_plane/test_agent_session_binding_candidates.py @@ -0,0 +1,138 @@ +"""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 From c03250fcea555e6a51c33d7fb607158ddcebccdf Mon Sep 17 00:00:00 2001 From: JunZ-Leo <100498253+JunZ-Leo@users.noreply.github.com> Date: Sat, 26 Sep 2026 21:49:24 +0800 Subject: [PATCH 2/4] fix(agents): count bindings by full identity, render them by budget Deduplication keyed on the already-rendered candidate: the thread identifier was compacted to 120 characters and the host surface to 60 before the pair was compared, while the binding owner accepts 128 and 64. Two valid bindings that differ only past the visible edge therefore rendered to the same text and were counted once, so the field the slice exists to make trustworthy still under-reported the number of routes a peer holds. Identity is now the untouched value; compaction only renders it. `strip` matches the owner's own normalisation, which is what keeps a padded identifier from being counted as a second binding. The rendered values stay inside the display budget, and the protocol document states that distinctness is decided on the full identity, so two counted bindings may legitimately render alike. Adds the long thread identifier and long host surface regressions, a padding case, and asserts the rendered length so a fix cannot pass by publishing raw values instead. Signed-off-by: JunZ-Leo <100498253+JunZ-Leo@users.noreply.github.com> --- .../agent-management-projection-v0.md | 4 +- .../agents/management_projection.py | 32 ++++--- .../test_agent_session_binding_candidates.py | 83 +++++++++++++++++++ 3 files changed, 108 insertions(+), 11 deletions(-) diff --git a/docs/reference/protocols/agent-management-projection-v0.md b/docs/reference/protocols/agent-management-projection-v0.md index 5439c40238..0ba5095cd7 100644 --- a/docs/reference/protocols/agent-management-projection-v0.md +++ b/docs/reference/protocols/agent-management-projection-v0.md @@ -118,7 +118,9 @@ Optional fields: 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. Neither field selects an execution route: absence is + 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; diff --git a/loopx/control_plane/agents/management_projection.py b/loopx/control_plane/agents/management_projection.py index 6ebd69e1a1..46694733a0 100644 --- a/loopx/control_plane/agents/management_projection.py +++ b/loopx/control_plane/agents/management_projection.py @@ -589,6 +589,7 @@ def build_agent_management_projection( # 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: dict[str, list[dict[str, str]]] = {} + seen_binding_identities: set[tuple[str, str, str]] = set() 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): @@ -597,18 +598,29 @@ def build_agent_management_projection( 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) + raw_agent = str(raw_binding.get("agent_id") or "").strip() + raw_thread = str(raw_binding.get("thread_id") or "").strip() + raw_host = str(raw_binding.get("host_surface") or "").strip() + agent_id = _compact(raw_agent, limit=120) + thread_id = _compact(raw_thread, limit=120) + host_surface = _compact(raw_host, limit=60) if not agent_id or not thread_id: continue - candidates = session_binding_candidates.setdefault(agent_id, []) - candidate = { - "thread_id": thread_id, - "host_surface": host_surface or "unknown", - } - if candidate not in candidates: - candidates.append(candidate) + # The display budget is not an identity key. The binding owner accepts + # thread identifiers to 128 characters and host surfaces to 64, both + # wider than what this projection shows, so two valid bindings that + # share a visible prefix are still two routes. Identity is therefore + # the untouched value; compaction only renders it. + identity = (raw_agent, raw_thread, raw_host) + if identity in seen_binding_identities: + continue + seen_binding_identities.add(identity) + session_binding_candidates.setdefault(agent_id, []).append( + { + "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): diff --git a/tests/control_plane/test_agent_session_binding_candidates.py b/tests/control_plane/test_agent_session_binding_candidates.py index 59090cce7c..f045357aa4 100644 --- a/tests/control_plane/test_agent_session_binding_candidates.py +++ b/tests/control_plane/test_agent_session_binding_candidates.py @@ -136,3 +136,86 @@ def test_several_candidates_keep_the_unbound_lifecycle_state( 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"} + ] From fbc5ba81235ec21a4bd3e76f4124e46282665c26 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Sat, 26 Sep 2026 23:32:39 +0800 Subject: [PATCH 3/4] fix(agents): keep the projection ratchet green by extracting the collector The session-binding candidate branch pushed build_agent_management_projection past the maintainability ratchet: 92 statements and 45 decision points against the 90/60 thresholds, so `control-plane-maintainability-ratchet-smoke` failed on this head while base passed. Collecting the bindings now lives in a private _collect_session_binding_candidates helper in the same module; the projection function only consumes the returned summary. Identity, public-safe compaction, full count before the display cap and lifecycle membership are unchanged. Behaviour evidence: the same mixed fixture (duplicate identity, long thread and host values that share a display prefix, whitespace padding, non-dict entries, a second Agent) produces byte-identical agent rows before and after the extraction, with session_binding_count = 5 and three rendered candidates. Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- .../agents/management_projection.py | 79 +++++++++++-------- 1 file changed, 46 insertions(+), 33 deletions(-) diff --git a/loopx/control_plane/agents/management_projection.py b/loopx/control_plane/agents/management_projection.py index 46694733a0..ae972e05e5 100644 --- a/loopx/control_plane/agents/management_projection.py +++ b/loopx/control_plane/agents/management_projection.py @@ -560,6 +560,51 @@ 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 historical session bindings by their full identity. + + Session bindings come from run_history.coordination.thread_agent_bindings + and are the only source for addressable/bound lifecycle states. The display + budget is not an identity key: the binding owner accepts thread identifiers + to 128 characters and host surfaces to 64, both wider than what this + projection shows, so two valid bindings that share a visible prefix are + still two routes. Identity is therefore the untouched value, and compaction + only renders it. + """ + + candidates: dict[str, list[dict[str, str]]] = {} + seen_identities: set[tuple[str, str, str]] = set() + 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 + raw_agent = str(raw_binding.get("agent_id") or "").strip() + raw_thread = str(raw_binding.get("thread_id") or "").strip() + raw_host = str(raw_binding.get("host_surface") or "").strip() + agent_id = _compact(raw_agent, limit=120) + thread_id = _compact(raw_thread, limit=120) + host_surface = _compact(raw_host, limit=60) + if not agent_id or not thread_id: + continue + identity = (raw_agent, raw_thread, raw_host) + if identity in seen_identities: + continue + seen_identities.add(identity) + candidates.setdefault(agent_id, []).append( + { + "thread_id": thread_id, + "host_surface": host_surface or "unknown", + } + ) + return candidates + + def build_agent_management_projection( status_payload: dict[str, Any], *, @@ -588,39 +633,7 @@ def build_agent_management_projection( # 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: dict[str, list[dict[str, str]]] = {} - seen_binding_identities: set[tuple[str, str, str]] = set() - 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 - raw_agent = str(raw_binding.get("agent_id") or "").strip() - raw_thread = str(raw_binding.get("thread_id") or "").strip() - raw_host = str(raw_binding.get("host_surface") or "").strip() - agent_id = _compact(raw_agent, limit=120) - thread_id = _compact(raw_thread, limit=120) - host_surface = _compact(raw_host, limit=60) - if not agent_id or not thread_id: - continue - # The display budget is not an identity key. The binding owner accepts - # thread identifiers to 128 characters and host surfaces to 64, both - # wider than what this projection shows, so two valid bindings that - # share a visible prefix are still two routes. Identity is therefore - # the untouched value; compaction only renders it. - identity = (raw_agent, raw_thread, raw_host) - if identity in seen_binding_identities: - continue - seen_binding_identities.add(identity) - session_binding_candidates.setdefault(agent_id, []).append( - { - "thread_id": thread_id, - "host_surface": host_surface or "unknown", - } - ) + 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): From 725b8fe8750cbf4185a6e867c0e857349b31d3be Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Sat, 26 Sep 2026 23:41:44 +0800 Subject: [PATCH 4/4] fix(agents): read binding candidates through the shared resolver Main now publishes peer routes from `collect_accepted_bindings` in the binding owner module, and that walk normalises the agent lane through `normalize_todo_claimed_by`. The management projection had its own walk that keyed candidates by the raw (strip-only) agent string, so a binding written as `Peer` was grouped under a key no management row ever looks up: the candidate was silently dropped, and a later `peer` row could under-report its routes. The private walk also counted bindings the owner refuses to name (a thread id past 128 characters, a host surface with inner whitespace, a missing host, an unnameable lane), which rendered a route no resolver would confirm. Delegate to the owner's collector instead: one normalisation rule, candidates keyed by the same agent identity the rows use, and `session_binding_count` that matches the `peer_route.candidate_count` the peer directory publishes for the same bindings. Display keeps its own 3-entry cap and 120/60 character budget, so identity still stays whole and `_compact` still only renders it. Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- .../agents/management_projection.py | 51 +++++++------------ 1 file changed, 18 insertions(+), 33 deletions(-) diff --git a/loopx/control_plane/agents/management_projection.py b/loopx/control_plane/agents/management_projection.py index ae972e05e5..b8f9e1d642 100644 --- a/loopx/control_plane/agents/management_projection.py +++ b/loopx/control_plane/agents/management_projection.py @@ -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, @@ -563,45 +564,29 @@ def _safe_next_action(todo: dict[str, Any] | None) -> str: def _collect_session_binding_candidates( status_payload: dict[str, Any], ) -> dict[str, list[dict[str, str]]]: - """Group one Agent's historical session bindings by their full identity. + """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. The display - budget is not an identity key: the binding owner accepts thread identifiers - to 128 characters and host surfaces to 64, both wider than what this - projection shows, so two valid bindings that share a visible prefix are - still two routes. Identity is therefore the untouched value, and compaction - only renders it. + 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]]] = {} - seen_identities: set[tuple[str, str, str]] = set() 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 - raw_agent = str(raw_binding.get("agent_id") or "").strip() - raw_thread = str(raw_binding.get("thread_id") or "").strip() - raw_host = str(raw_binding.get("host_surface") or "").strip() - agent_id = _compact(raw_agent, limit=120) - thread_id = _compact(raw_thread, limit=120) - host_surface = _compact(raw_host, limit=60) - if not agent_id or not thread_id: - continue - identity = (raw_agent, raw_thread, raw_host) - if identity in seen_identities: - continue - seen_identities.add(identity) - candidates.setdefault(agent_id, []).append( - { - "thread_id": thread_id, - "host_surface": host_surface or "unknown", - } - ) + 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