fix: ACP assistant messageIds are UUIDs, not a counter
The ACP schema (agent-client-protocol 0.9.0, ContentChunk.messageId) says "Both clients and agents MUST use UUID format for message IDs". The ported allocator emitted hermes-assistant-N strings, which a strict client may reject or fail to group. A fresh uuid4 per message keeps the grouping semantics and can never collide with an earlier turn's id, so the counter/prefix state is gone. Tests trimmed to three invariants: chunks share one UUID until the None flush sentinel (empty deltas ignored), thought + text share an id, and the no-allocator shape stays unchanged.
This commit is contained in:
+8
-11
@@ -8,6 +8,7 @@ thread-safely onto the loop.
|
||||
|
||||
import asyncio
|
||||
import logging
|
||||
import uuid
|
||||
from collections import deque
|
||||
from typing import Any, Callable, Deque, Dict
|
||||
|
||||
@@ -135,25 +136,21 @@ class AssistantMessageIdAllocator:
|
||||
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``).
|
||||
One allocator lives per ACP session; a contiguous run of deltas shares
|
||||
``current()`` and ``close()`` marks the message finished so the next delta
|
||||
allocates a fresh id. Ids are UUID4 strings because the ACP schema requires
|
||||
UUID-format message ids, and a fresh UUID can never collide with an earlier
|
||||
turn's id.
|
||||
"""
|
||||
|
||||
def __init__(self, prefix: str = "hermes-assistant") -> None:
|
||||
self._prefix = prefix
|
||||
self._sequence = 0
|
||||
def __init__(self) -> None:
|
||||
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
|
||||
self._active = self._last = str(uuid.uuid4())
|
||||
return self._active
|
||||
|
||||
def last(self) -> str | None:
|
||||
|
||||
@@ -2,6 +2,7 @@
|
||||
|
||||
import asyncio
|
||||
import gc
|
||||
import uuid
|
||||
import warnings
|
||||
from concurrent.futures import Future
|
||||
from unittest.mock import AsyncMock, MagicMock, patch
|
||||
@@ -268,12 +269,9 @@ class TestSendUpdate:
|
||||
|
||||
|
||||
class TestAssistantMessageIds:
|
||||
"""Assistant messageId grouping — ported from prime-agent#1781."""
|
||||
"""Streamed chunks carry a per-message ACP messageId; the None flush sentinel starts a new one."""
|
||||
|
||||
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):
|
||||
def test_deltas_share_one_uuid_until_flush(self, mock_conn, event_loop_fixture):
|
||||
from acp_adapter.events import AssistantMessageIdAllocator
|
||||
|
||||
ids = AssistantMessageIdAllocator()
|
||||
@@ -282,11 +280,14 @@ class TestAssistantMessageIds:
|
||||
with patch("acp_adapter.events._send_update",
|
||||
side_effect=lambda c, s, l, u: sent.append(u)):
|
||||
cb("Hello ")
|
||||
cb("") # empty delta is ignored, not a flush
|
||||
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"
|
||||
assert sent[0].message_id == sent[1].message_id
|
||||
assert sent[2].message_id != sent[0].message_id
|
||||
# ACP requires UUID-format message ids.
|
||||
assert uuid.UUID(sent[0].message_id) and uuid.UUID(sent[2].message_id)
|
||||
|
||||
def test_thought_chunks_carry_id(self, mock_conn, event_loop_fixture):
|
||||
from acp_adapter.events import AssistantMessageIdAllocator
|
||||
@@ -309,29 +310,3 @@ class TestAssistantMessageIds:
|
||||
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
|
||||
|
||||
Reference in New Issue
Block a user