feat(reasoning): stream reasoning content as a first-class channel
Reasoning now flows as its own stream — symmetric to the answer's ``delta`` / ``stream_end`` pair — instead of being shipped as one oversized progress message. This lets WebUI render a live "Thinking…" bubble that updates in place, then auto-collapses when the stream closes. Other channels remain plugin no-ops by default. ## Protocol New metadata: ``_reasoning_delta`` (chunk) and ``_reasoning_end`` (close marker). ChannelManager routes both to the dedicated plugin hooks below; the legacy one-shot ``_reasoning`` is kept for back-compat and BaseChannel expands it into a single delta + end pair so plugins only ever implement the streaming primitives. WebSocket emits two new events: - ``reasoning_delta`` (event, chat_id, text, optional stream_id) - ``reasoning_end`` (event, chat_id, optional stream_id) ## BaseChannel surface - ``send_reasoning_delta(chat_id, delta, metadata)`` — no-op default - ``send_reasoning_end(chat_id, metadata)`` — no-op default - ``send_reasoning(msg)`` — back-compat wrapper, base impl forwards to the streaming primitives A channel adds reasoning support by overriding the two streaming primitives. Telegram / Slack / Discord / Feishu / WeChat / Matrix keep the base no-ops until their bubble UIs are adapted; reasoning silently drops at dispatch, never as a stray text message. ## AgentHook Adds ``emit_reasoning_end`` to the hook lifecycle. ``_LoopHook`` tracks whether a reasoning segment is open and closes it on: - the first answer delta arriving (so the UI locks the bubble before the answer renders below), - ``on_stream_end``, - one-shot ``reasoning_content`` / ``thinking_blocks`` after a single non-streaming response. ## WebUI - ``UIMessage.reasoning`` is now a single accumulated string with a companion ``reasoningStreaming`` flag. - ``useNanobotStream`` consumes ``reasoning_delta`` / ``reasoning_end``; legacy ``kind: "reasoning"`` is auto-translated to a delta + end. - New ``ReasoningBubble``: shimmer header + auto-expanded while streaming, collapses to a clickable "Thinking" pill once closed, respects ``prefers-reduced-motion``. - Answer deltas adopt the reasoning placeholder so the bubble and the answer share one assistant row. ## Tests - ``tests/channels/test_channel_manager_reasoning.py`` — manager routes delta + end, drops on channel opt-out, expands one-shot back-compat. - ``tests/channels/test_websocket_channel.py`` — new ``reasoning_delta`` / ``reasoning_end`` frames, empty-chunk safety, no-subscriber safety, back-compat expansion. - ``tests/agent/test_runner_reasoning.py`` — runner closes the segment on streaming answer start and after one-shot reasoning. - WebUI ``useNanobotStream`` + ``message-bubble`` cover the new protocol and the shimmer styling. ## Docs ``docs/configuration.md`` and ``docs/websocket.md`` document the new events and the plugin contract. Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
@@ -24,11 +24,15 @@ class _RecordingHook(AgentHook):
|
||||
def __init__(self) -> None:
|
||||
super().__init__()
|
||||
self.emitted: list[str] = []
|
||||
self.end_calls = 0
|
||||
|
||||
async def emit_reasoning(self, reasoning_content: str | None) -> None:
|
||||
if reasoning_content:
|
||||
self.emitted.append(reasoning_content)
|
||||
|
||||
async def emit_reasoning_end(self) -> None:
|
||||
self.end_calls += 1
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_runner_preserves_reasoning_fields_in_assistant_history():
|
||||
@@ -277,3 +281,41 @@ async def test_runner_does_not_double_emit_when_inline_think_already_streamed():
|
||||
|
||||
assert result.final_content == "The answer."
|
||||
assert hook.emitted == ["working..."]
|
||||
assert hook.end_calls >= 1, "reasoning stream must be closed once the answer starts"
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_runner_closes_reasoning_stream_after_one_shot_response():
|
||||
"""A non-streaming response carrying ``reasoning_content`` must emit
|
||||
both a reasoning delta and an end marker so channels can finalize the
|
||||
in-place bubble."""
|
||||
from nanobot.agent.runner import AgentRunSpec, AgentRunner
|
||||
|
||||
provider = MagicMock()
|
||||
|
||||
async def chat_with_retry(**kwargs):
|
||||
return LLMResponse(
|
||||
content="answer",
|
||||
reasoning_content="hidden thought",
|
||||
tool_calls=[],
|
||||
usage={"prompt_tokens": 5, "completion_tokens": 3},
|
||||
)
|
||||
|
||||
provider.chat_with_retry = chat_with_retry
|
||||
tools = MagicMock()
|
||||
tools.get_definitions.return_value = []
|
||||
|
||||
hook = _RecordingHook()
|
||||
runner = AgentRunner(provider)
|
||||
result = await runner.run(AgentRunSpec(
|
||||
initial_messages=[{"role": "user", "content": "q"}],
|
||||
tools=tools,
|
||||
model="test-model",
|
||||
max_iterations=3,
|
||||
max_tool_result_chars=_MAX_TOOL_RESULT_CHARS,
|
||||
hook=hook,
|
||||
))
|
||||
|
||||
assert result.final_content == "answer"
|
||||
assert hook.emitted == ["hidden thought"]
|
||||
assert hook.end_calls == 1
|
||||
|
||||
@@ -1,14 +1,22 @@
|
||||
"""Tests for ChannelManager routing of model reasoning content.
|
||||
|
||||
Reasoning is delivered as a separate plugin action (``send_reasoning``)
|
||||
rather than a metadata flag on a regular outbound. The manager routes
|
||||
``_reasoning`` messages only to channels that opt in via
|
||||
``channel.show_reasoning``; channels without a low-emphasis UI primitive
|
||||
keep the base no-op and the content silently drops at dispatch.
|
||||
Reasoning is delivered through plugin streaming primitives
|
||||
(``send_reasoning_delta`` / ``send_reasoning_end``) so each channel
|
||||
controls in-place rendering — mirroring the existing answer ``send_delta``
|
||||
/ ``stream_end`` pair. The manager forwards reasoning frames only to
|
||||
channels that opt in via ``channel.show_reasoning``; plugins without a
|
||||
low-emphasis UI primitive keep the base no-op and the content silently
|
||||
drops at dispatch.
|
||||
|
||||
One-shot ``_reasoning`` frames are accepted for back-compat with hooks
|
||||
that haven't migrated yet — ``BaseChannel.send_reasoning`` expands them
|
||||
to a single delta + end pair so plugins only implement the streaming
|
||||
primitives.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
from unittest.mock import AsyncMock
|
||||
|
||||
import pytest
|
||||
@@ -27,7 +35,8 @@ class _MockChannel(BaseChannel):
|
||||
def __init__(self, config, bus):
|
||||
super().__init__(config, bus)
|
||||
self._send_mock = AsyncMock()
|
||||
self._send_reasoning_mock = AsyncMock()
|
||||
self._delta_mock = AsyncMock()
|
||||
self._end_mock = AsyncMock()
|
||||
|
||||
async def start(self): # pragma: no cover - not exercised
|
||||
pass
|
||||
@@ -38,8 +47,11 @@ class _MockChannel(BaseChannel):
|
||||
async def send(self, msg):
|
||||
return await self._send_mock(msg)
|
||||
|
||||
async def send_reasoning(self, msg):
|
||||
return await self._send_reasoning_mock(msg)
|
||||
async def send_reasoning_delta(self, chat_id, delta, metadata=None):
|
||||
return await self._delta_mock(chat_id, delta, metadata)
|
||||
|
||||
async def send_reasoning_end(self, chat_id, metadata=None):
|
||||
return await self._end_mock(chat_id, metadata)
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
@@ -50,17 +62,52 @@ def manager() -> ChannelManager:
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_reasoning_routes_to_send_reasoning_not_send(manager):
|
||||
async def test_reasoning_delta_routes_to_send_reasoning_delta(manager):
|
||||
channel = manager.channels["mock"]
|
||||
msg = OutboundMessage(
|
||||
channel="mock",
|
||||
chat_id="c1",
|
||||
content="step-by-step thinking",
|
||||
content="step-by-step",
|
||||
metadata={"_progress": True, "_reasoning_delta": True, "_stream_id": "r1"},
|
||||
)
|
||||
await manager._send_once(channel, msg)
|
||||
channel._delta_mock.assert_awaited_once()
|
||||
args = channel._delta_mock.await_args.args
|
||||
assert args[0] == "c1"
|
||||
assert args[1] == "step-by-step"
|
||||
channel._send_mock.assert_not_awaited()
|
||||
channel._end_mock.assert_not_awaited()
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_reasoning_end_routes_to_send_reasoning_end(manager):
|
||||
channel = manager.channels["mock"]
|
||||
msg = OutboundMessage(
|
||||
channel="mock",
|
||||
chat_id="c1",
|
||||
content="",
|
||||
metadata={"_progress": True, "_reasoning_end": True, "_stream_id": "r1"},
|
||||
)
|
||||
await manager._send_once(channel, msg)
|
||||
channel._end_mock.assert_awaited_once()
|
||||
channel._delta_mock.assert_not_awaited()
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_legacy_one_shot_reasoning_expands_to_delta_plus_end(manager):
|
||||
"""`_reasoning` (no delta/end pair) falls back through `send_reasoning`
|
||||
which the base class expands to a single delta + end. Hooks that haven't
|
||||
migrated still surface in WebUI as a complete stream segment."""
|
||||
channel = manager.channels["mock"]
|
||||
msg = OutboundMessage(
|
||||
channel="mock",
|
||||
chat_id="c1",
|
||||
content="one-shot reasoning",
|
||||
metadata={"_progress": True, "_reasoning": True},
|
||||
)
|
||||
await manager._send_once(channel, msg)
|
||||
channel._send_reasoning_mock.assert_awaited_once_with(msg)
|
||||
channel._send_mock.assert_not_awaited()
|
||||
channel._delta_mock.assert_awaited_once()
|
||||
channel._end_mock.assert_awaited_once()
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
@@ -71,14 +118,14 @@ async def test_dispatch_drops_reasoning_when_channel_opts_out(manager):
|
||||
channel="mock",
|
||||
chat_id="c1",
|
||||
content="hidden thinking",
|
||||
metadata={"_progress": True, "_reasoning": True},
|
||||
metadata={"_progress": True, "_reasoning_delta": True},
|
||||
)
|
||||
await manager.bus.publish_outbound(msg)
|
||||
|
||||
pumped = await _pump_one(manager)
|
||||
await _pump_one(manager)
|
||||
|
||||
assert pumped is True
|
||||
channel._send_reasoning_mock.assert_not_awaited()
|
||||
channel._delta_mock.assert_not_awaited()
|
||||
channel._end_mock.assert_not_awaited()
|
||||
channel._send_mock.assert_not_awaited()
|
||||
|
||||
|
||||
@@ -86,20 +133,24 @@ async def test_dispatch_drops_reasoning_when_channel_opts_out(manager):
|
||||
async def test_dispatch_delivers_reasoning_when_channel_opts_in(manager):
|
||||
channel = manager.channels["mock"]
|
||||
channel.show_reasoning = True
|
||||
msg = OutboundMessage(
|
||||
for chunk in ("first ", "second"):
|
||||
await manager.bus.publish_outbound(OutboundMessage(
|
||||
channel="mock",
|
||||
chat_id="c1",
|
||||
content=chunk,
|
||||
metadata={"_progress": True, "_reasoning_delta": True, "_stream_id": "r1"},
|
||||
))
|
||||
await manager.bus.publish_outbound(OutboundMessage(
|
||||
channel="mock",
|
||||
chat_id="c1",
|
||||
content="visible thinking",
|
||||
metadata={"_progress": True, "_reasoning": True},
|
||||
)
|
||||
await manager.bus.publish_outbound(msg)
|
||||
content="",
|
||||
metadata={"_progress": True, "_reasoning_end": True, "_stream_id": "r1"},
|
||||
))
|
||||
|
||||
pumped = await _pump_one(manager)
|
||||
await _pump_one(manager)
|
||||
|
||||
assert pumped is True
|
||||
channel._send_reasoning_mock.assert_awaited_once()
|
||||
delivered = channel._send_reasoning_mock.await_args.args[0]
|
||||
assert delivered.content == "visible thinking"
|
||||
assert channel._delta_mock.await_count == 2
|
||||
channel._end_mock.assert_awaited_once()
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
@@ -108,21 +159,19 @@ async def test_dispatch_silently_drops_reasoning_for_unknown_channel(manager):
|
||||
channel="ghost",
|
||||
chat_id="c1",
|
||||
content="nobody home",
|
||||
metadata={"_progress": True, "_reasoning": True},
|
||||
metadata={"_progress": True, "_reasoning_delta": True},
|
||||
)
|
||||
await manager.bus.publish_outbound(msg)
|
||||
|
||||
pumped = await _pump_one(manager)
|
||||
await _pump_one(manager)
|
||||
|
||||
assert pumped is True
|
||||
# Mock channel must not receive anything destined for a different channel.
|
||||
manager.channels["mock"]._send_reasoning_mock.assert_not_awaited()
|
||||
manager.channels["mock"]._delta_mock.assert_not_awaited()
|
||||
manager.channels["mock"]._send_mock.assert_not_awaited()
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_base_channel_send_reasoning_is_noop_safe():
|
||||
"""Plugins that don't override `send_reasoning` must not blow up."""
|
||||
async def test_base_channel_reasoning_primitives_are_noop_safe():
|
||||
"""Plugins that don't override the streaming primitives must not blow up."""
|
||||
|
||||
class _Plain(BaseChannel):
|
||||
name = "plain"
|
||||
@@ -138,7 +187,9 @@ async def test_base_channel_send_reasoning_is_noop_safe():
|
||||
pass
|
||||
|
||||
channel = _Plain({}, MessageBus())
|
||||
# No exception, returns None.
|
||||
assert await channel.send_reasoning_delta("c", "x") is None
|
||||
assert await channel.send_reasoning_end("c") is None
|
||||
# And the one-shot wrapper translates without raising.
|
||||
assert await channel.send_reasoning(
|
||||
OutboundMessage(channel="plain", chat_id="c", content="x", metadata={})
|
||||
) is None
|
||||
@@ -151,26 +202,21 @@ async def test_reasoning_routing_does_not_consult_send_progress(manager):
|
||||
channel = manager.channels["mock"]
|
||||
channel.send_progress = False
|
||||
channel.show_reasoning = True
|
||||
msg = OutboundMessage(
|
||||
await manager.bus.publish_outbound(OutboundMessage(
|
||||
channel="mock",
|
||||
chat_id="c1",
|
||||
content="still surfaces",
|
||||
metadata={"_progress": True, "_reasoning": True},
|
||||
)
|
||||
await manager.bus.publish_outbound(msg)
|
||||
metadata={"_progress": True, "_reasoning_delta": True},
|
||||
))
|
||||
|
||||
pumped = await _pump_one(manager)
|
||||
await _pump_one(manager)
|
||||
|
||||
assert pumped is True
|
||||
channel._send_reasoning_mock.assert_awaited_once()
|
||||
channel._delta_mock.assert_awaited_once()
|
||||
|
||||
|
||||
async def _pump_one(manager: ChannelManager) -> bool:
|
||||
"""Drive the dispatcher for exactly one message, then cancel."""
|
||||
import asyncio
|
||||
|
||||
async def _pump_one(manager: ChannelManager) -> None:
|
||||
"""Drive the dispatcher until the outbound queue drains, then cancel."""
|
||||
task = asyncio.create_task(manager._dispatch_outbound())
|
||||
# Yield control until the queue drains.
|
||||
for _ in range(50):
|
||||
await asyncio.sleep(0.01)
|
||||
if manager.bus.outbound.qsize() == 0:
|
||||
@@ -180,4 +226,3 @@ async def _pump_one(manager: ChannelManager) -> bool:
|
||||
await task
|
||||
except asyncio.CancelledError:
|
||||
pass
|
||||
return True
|
||||
|
||||
@@ -359,30 +359,44 @@ async def test_send_delta_emits_delta_and_stream_end() -> None:
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_send_reasoning_emits_reasoning_kind_frame() -> None:
|
||||
async def test_send_reasoning_delta_emits_streaming_frame() -> None:
|
||||
bus = MagicMock()
|
||||
channel = WebSocketChannel({"enabled": True, "allowFrom": ["*"]}, bus)
|
||||
mock_ws = AsyncMock()
|
||||
channel._attach(mock_ws, "chat-1")
|
||||
|
||||
await channel.send_reasoning(OutboundMessage(
|
||||
channel="websocket",
|
||||
chat_id="chat-1",
|
||||
content="step-by-step thinking",
|
||||
metadata={"_progress": True, "_reasoning": True},
|
||||
))
|
||||
await channel.send_reasoning_delta(
|
||||
"chat-1",
|
||||
"step-by-step thinking",
|
||||
{"_reasoning_delta": True, "_stream_id": "r1"},
|
||||
)
|
||||
|
||||
mock_ws.send.assert_awaited_once()
|
||||
payload = json.loads(mock_ws.send.await_args.args[0])
|
||||
assert payload["event"] == "message"
|
||||
assert payload["event"] == "reasoning_delta"
|
||||
assert payload["chat_id"] == "chat-1"
|
||||
assert payload["text"] == "step-by-step thinking"
|
||||
assert payload["kind"] == "reasoning"
|
||||
assert payload["stream_id"] == "r1"
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_send_reasoning_drops_empty_content() -> None:
|
||||
"""Empty reasoning emits nothing — keeps the frontend bubble clean."""
|
||||
async def test_send_reasoning_end_emits_close_frame() -> None:
|
||||
bus = MagicMock()
|
||||
channel = WebSocketChannel({"enabled": True, "allowFrom": ["*"]}, bus)
|
||||
mock_ws = AsyncMock()
|
||||
channel._attach(mock_ws, "chat-1")
|
||||
|
||||
await channel.send_reasoning_end("chat-1", {"_reasoning_end": True, "_stream_id": "r1"})
|
||||
|
||||
payload = json.loads(mock_ws.send.await_args.args[0])
|
||||
assert payload == {"event": "reasoning_end", "chat_id": "chat-1", "stream_id": "r1"}
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_send_reasoning_one_shot_expands_to_delta_plus_end() -> None:
|
||||
"""``send_reasoning`` is back-compat for hooks that haven't migrated:
|
||||
the base implementation must produce one delta and one end so the
|
||||
WebUI sees the same shape either way."""
|
||||
bus = MagicMock()
|
||||
channel = WebSocketChannel({"enabled": True, "allowFrom": ["*"]}, bus)
|
||||
mock_ws = AsyncMock()
|
||||
@@ -391,10 +405,27 @@ async def test_send_reasoning_drops_empty_content() -> None:
|
||||
await channel.send_reasoning(OutboundMessage(
|
||||
channel="websocket",
|
||||
chat_id="chat-1",
|
||||
content="",
|
||||
content="thinking",
|
||||
metadata={"_reasoning": True},
|
||||
))
|
||||
|
||||
assert mock_ws.send.await_count == 2
|
||||
first = json.loads(mock_ws.send.call_args_list[0][0][0])
|
||||
second = json.loads(mock_ws.send.call_args_list[1][0][0])
|
||||
assert first["event"] == "reasoning_delta"
|
||||
assert first["text"] == "thinking"
|
||||
assert second["event"] == "reasoning_end"
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_send_reasoning_delta_drops_empty_chunks() -> None:
|
||||
bus = MagicMock()
|
||||
channel = WebSocketChannel({"enabled": True, "allowFrom": ["*"]}, bus)
|
||||
mock_ws = AsyncMock()
|
||||
channel._attach(mock_ws, "chat-1")
|
||||
|
||||
await channel.send_reasoning_delta("chat-1", "", {"_reasoning_delta": True})
|
||||
|
||||
mock_ws.send.assert_not_awaited()
|
||||
|
||||
|
||||
@@ -403,12 +434,8 @@ async def test_send_reasoning_without_subscribers_is_noop() -> None:
|
||||
bus = MagicMock()
|
||||
channel = WebSocketChannel({"enabled": True, "allowFrom": ["*"]}, bus)
|
||||
|
||||
await channel.send_reasoning(OutboundMessage(
|
||||
channel="websocket",
|
||||
chat_id="unattached",
|
||||
content="thinking",
|
||||
metadata={"_reasoning": True},
|
||||
))
|
||||
await channel.send_reasoning_delta("unattached", "thinking", None)
|
||||
await channel.send_reasoning_end("unattached", None)
|
||||
# No subscribers, no exception, no send.
|
||||
|
||||
|
||||
|
||||
Reference in New Issue
Block a user