diff --git a/acp_adapter/events.py b/acp_adapter/events.py index f63c57feb1..7aed3758f8 100644 --- a/acp_adapter/events.py +++ b/acp_adapter/events.py @@ -121,22 +121,83 @@ def make_tool_progress_cb( return _tool_progress -def _make_text_cb(conn: acp.Client, session_id: str, loop: asyncio.AbstractEventLoop, wrap: Callable[[str], Any]) -> Callable: - def _cb(text: str) -> None: +# ------------------------------------------------------------------ +# Assistant message identity +# ------------------------------------------------------------------ + + +class AssistantMessageIdAllocator: + """Allocates stable per-message ids for streamed assistant chunks. + + ACP clients group streamed ``agent_message_chunk`` / ``agent_thought_chunk`` + deltas into one assistant reply by ``messageId`` and use a NEW id to start + the next reply (root-reply replacement semantics). Without ids, a client + that replaces "the current assistant message" on each chunk collapses + separate autonomous turns into one bubble. + + One allocator lives per ACP session so the sequence is monotonic across + turns — two different turns must never reuse an id. A contiguous run of + deltas shares ``current()``; ``close()`` marks the message finished so the + next delta allocates a fresh id. Ported from + PrimeIntellect-ai/prime-agent#1781 (``prime-agent-assistant-N``). + """ + + def __init__(self, prefix: str = "hermes-assistant") -> None: + self._prefix = prefix + self._sequence = 0 + self._active: str | None = None + self._last: str | None = None + + def current(self) -> str: + """Return the active message id, allocating one if none is open.""" + if self._active is None: + self._sequence += 1 + self._active = f"{self._prefix}-{self._sequence}" + self._last = self._active + return self._active + + def last(self) -> str | None: + """Return the most recently allocated id (open or closed).""" + return self._last + + def close(self) -> None: + """End the active message; the next chunk starts a new id.""" + self._active = None + + +def _make_text_cb( + conn: acp.Client, session_id: str, loop: asyncio.AbstractEventLoop, wrap: Callable[[str], Any], + message_ids: AssistantMessageIdAllocator | None = None, +) -> Callable: + # ``None`` is the flush sentinel Hermes core sends between assistant messages + # (before tool execution / at end of stream): it closes the active messageId so + # the next delta opens a new bubble instead of merging into the previous one. + def _cb(text: str | None) -> None: if text: - _send_update(conn, session_id, loop, wrap(text)) + update = wrap(text) + if message_ids is not None: + update.message_id = message_ids.current() + _send_update(conn, session_id, loop, update) + elif text is None and message_ids is not None: + message_ids.close() return _cb -def make_thinking_cb(conn: acp.Client, session_id: str, loop: asyncio.AbstractEventLoop) -> Callable: +def make_thinking_cb( + conn: acp.Client, session_id: str, loop: asyncio.AbstractEventLoop, + message_ids: AssistantMessageIdAllocator | None = None, +) -> Callable: """Create a ``thinking_callback`` for AIAgent.""" - return _make_text_cb(conn, session_id, loop, acp.update_agent_thought_text) + return _make_text_cb(conn, session_id, loop, acp.update_agent_thought_text, message_ids) -def make_message_cb(conn: acp.Client, session_id: str, loop: asyncio.AbstractEventLoop) -> Callable: +def make_message_cb( + conn: acp.Client, session_id: str, loop: asyncio.AbstractEventLoop, + message_ids: AssistantMessageIdAllocator | None = None, +) -> Callable: """Create a callback that streams agent response text to the editor.""" - return _make_text_cb(conn, session_id, loop, acp.update_agent_message_text) + return _make_text_cb(conn, session_id, loop, acp.update_agent_message_text, message_ids) def make_step_cb( diff --git a/acp_adapter/server.py b/acp_adapter/server.py index eb18d79513..f7daaa8d94 100644 --- a/acp_adapter/server.py +++ b/acp_adapter/server.py @@ -28,7 +28,8 @@ from acp_adapter.auth import TERMINAL_SETUP_AUTH_METHOD_ID, build_auth_methods, from acp_adapter.commands import HERMES_VERSION, SlashCommandsMixin, _estimate_tokens from acp_adapter.content import PromptBlock, _content_blocks_to_openai_user_content, _extract_text from acp_adapter.events import ( - _build_plan_update_from_todo_result, make_message_cb, make_step_cb, make_thinking_cb, make_tool_progress_cb, + AssistantMessageIdAllocator, _build_plan_update_from_todo_result, make_message_cb, make_step_cb, + make_thinking_cb, make_tool_progress_cb, ) from acp_adapter.model_catalog import build_model_state, encode_model_choice from acp_adapter.permissions import make_approval_callback @@ -849,9 +850,14 @@ class HermesACPAgent(SlashCommandsMixin, acp.Agent): cbs.tool_progress_cb = make_tool_progress_cb( conn, session_id, loop, tool_call_ids, tool_call_meta, edit_approval_policy_getter=policy_getter ) - cbs.reasoning_cb = make_thinking_cb(conn, session_id, loop) + # Per-session allocator: a new turn must never reuse a previous turn's + # assistant messageId (ACP clients replace the bubble with that id). + if state.message_ids is None: + state.message_ids = AssistantMessageIdAllocator() + state.message_ids.close() # new turn -> next chunk opens a fresh id + cbs.reasoning_cb = make_thinking_cb(conn, session_id, loop, state.message_ids) cbs.step_cb = make_step_cb(conn, session_id, loop, tool_call_ids, tool_call_meta) - message_cb = make_message_cb(conn, session_id, loop) + message_cb = make_message_cb(conn, session_id, loop, state.message_ids) def stream_delta_cb(text: str) -> None: cbs.streamed = cbs.streamed or bool(text) @@ -907,7 +913,16 @@ class HermesACPAgent(SlashCommandsMixin, acp.Agent): suppress = interrupted and final_response.startswith(INTERRUPT_WAITING_FOR_MODEL_PREFIX) # Send the final text unless already streamed — or if a plugin hook transformed it after. if final_response and conn and not suppress and (not streamed_message or result.get("response_transformed")): - await conn.session_update(session_id, acp.update_agent_message_text(final_response)) + update = acp.update_agent_message_text(final_response) + if state.message_ids is not None: + # A plugin-rewritten reply replaces the streamed bubble (same id); an + # unstreamed final response opens its own. + if streamed_message and result.get("response_transformed"): + update.message_id = state.message_ids.last() or state.message_ids.current() + else: + update.message_id = state.message_ids.current() + state.message_ids.close() + await conn.session_update(session_id, update) # Go idle before draining so recursive prompt() calls can acquire the session. with state.runtime_lock: diff --git a/acp_adapter/session.py b/acp_adapter/session.py index a0823e6af7..71be47e8c0 100644 --- a/acp_adapter/session.py +++ b/acp_adapter/session.py @@ -143,6 +143,9 @@ class SessionState: runtime_lock: Any = field(default_factory=threading.Lock) current_prompt_text: str = "" interrupted_prompt_text: str = "" + # Per-session allocator for ACP assistant messageIds (lazily created by + # the server so streamed chunks group into distinct assistant replies). + message_ids: Any = None class SessionManager: diff --git a/tests/acp_adapter/test_events.py b/tests/acp_adapter/test_events.py index a1bcb4bb76..0e417197d0 100644 --- a/tests/acp_adapter/test_events.py +++ b/tests/acp_adapter/test_events.py @@ -265,3 +265,73 @@ class TestSendUpdate: and "_session_update" in str(w.message) ] assert runtime_warnings == [] + + +class TestAssistantMessageIds: + """Assistant messageId grouping — ported from prime-agent#1781.""" + + def _sent_updates(self, mock_rcts): + return [call.args[0] for call in mock_rcts.call_args_list] + + def test_deltas_share_one_id_until_flush(self, mock_conn, event_loop_fixture): + from acp_adapter.events import AssistantMessageIdAllocator + + ids = AssistantMessageIdAllocator() + cb = make_message_cb(mock_conn, "s", event_loop_fixture, ids) + sent = [] + with patch("acp_adapter.events._send_update", + side_effect=lambda c, s, l, u: sent.append(u)): + cb("Hello ") + cb("world") + cb(None) # flush sentinel — closes the message + cb("next turn") + assert sent[0].message_id == sent[1].message_id == "hermes-assistant-1" + assert sent[2].message_id == "hermes-assistant-2" + + def test_thought_chunks_carry_id(self, mock_conn, event_loop_fixture): + from acp_adapter.events import AssistantMessageIdAllocator + + ids = AssistantMessageIdAllocator() + think = make_thinking_cb(mock_conn, "s", event_loop_fixture, ids) + msg = make_message_cb(mock_conn, "s", event_loop_fixture, ids) + sent = [] + with patch("acp_adapter.events._send_update", + side_effect=lambda c, s, l, u: sent.append(u)): + think("pondering") + msg("answer") + # Reasoning and answer of the same reply share one message id. + assert sent[0].message_id == sent[1].message_id + + def test_no_allocator_keeps_legacy_shape(self, mock_conn, event_loop_fixture): + cb = make_message_cb(mock_conn, "s", event_loop_fixture) + sent = [] + with patch("acp_adapter.events._send_update", + side_effect=lambda c, s, l, u: sent.append(u)): + cb("text") + assert sent[0].message_id is None + + def test_ids_monotonic_never_reused(self): + from acp_adapter.events import AssistantMessageIdAllocator + + ids = AssistantMessageIdAllocator() + seen = set() + for _ in range(5): + i = ids.current() + assert i not in seen + seen.add(i) + ids.close() + assert ids.last() == "hermes-assistant-5" + + def test_empty_string_does_not_close_message(self, mock_conn, event_loop_fixture): + """Only the None sentinel ends a message; '' deltas are ignored.""" + from acp_adapter.events import AssistantMessageIdAllocator + + ids = AssistantMessageIdAllocator() + cb = make_message_cb(mock_conn, "s", event_loop_fixture, ids) + sent = [] + with patch("acp_adapter.events._send_update", + side_effect=lambda c, s, l, u: sent.append(u)): + cb("a") + cb("") + cb("b") + assert sent[0].message_id == sent[1].message_id