Port from PrimeIntellect-ai/prime-agent#1781: stable assistant messageIds on ACP streamed chunks

ACP clients group streamed agent_message_chunk / agent_thought_chunk
updates into one assistant reply by messageId, and use a new id to
start the next reply (root-reply replacement semantics). Hermes' ACP
adapter sent every chunk without a messageId, so clients that replace
'the current assistant message' per chunk collapsed separate
autonomous turns into one bubble.

- AssistantMessageIdAllocator (per ACP session, monotonic across
  turns): a contiguous run of reasoning + text deltas shares one
  hermes-assistant-N id; the None flush sentinel Hermes core emits
  before tool execution / at end of stream closes it.
- make_message_cb / make_thinking_cb stamp update.message_id when an
  allocator is provided; legacy no-allocator shape unchanged.
- Unstreamed final responses open their own id; plugin-transformed
  responses reuse the streamed message's id (replacement).
- Tests: grouping until flush, thought+text sharing, monotonic ids,
  empty-string vs None sentinel, legacy shape.
This commit is contained in:
Teknium
2026-08-25 21:48:43 -07:00
parent aea84cbb1a
commit 4a1e44dd2f
4 changed files with 160 additions and 11 deletions

View File

@@ -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(

View File

@@ -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:

View File

@@ -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:

View File

@@ -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