diff --git a/nerve/agent/engine.py b/nerve/agent/engine.py index 019d1035f..9a200e4b2 100644 --- a/nerve/agent/engine.py +++ b/nerve/agent/engine.py @@ -3067,6 +3067,9 @@ async def _open_turn() -> None: await broadcaster.broadcast(session_id, { "type": "auto_turn", "session_id": session_id, }) + # The broadcaster only reaches the web UI; channels such as + # Telegram need their own stream for a turn nobody asked for. + await self._open_autonomous_channel_stream(session_id) def _turn_has_content() -> bool: return st is not None and ( @@ -3091,6 +3094,9 @@ async def _close_turn() -> None: # Empty turn (init arrived but content never did) — drop it; # the finally backstop ships a synthetic done if framing opened. st = None + # After finalize's "done" reached the channel. For an empty + # turn this removes the channel's placeholder instead. + await self._close_autonomous_channel_stream(session_id) try: while True: @@ -3237,9 +3243,38 @@ async def _close_turn() -> None: broadcaster.stop_buffering(session_id) with contextlib.suppress(Exception): await self._broadcast_session_running(session_id, False) + # A turn cut short (cancelled, stream failure) never reached + # _close_turn; flush and release its channel stream. + await self._close_autonomous_channel_stream(session_id) return turns + async def _open_autonomous_channel_stream(self, session_id: str) -> None: + """Ask the channel router to stream an autonomous turn (best effort).""" + router = getattr(self, "_router", None) + if router is None: + return + try: + await router.open_autonomous_stream(session_id) + except Exception as e: + logger.warning( + "Could not open channel stream for autonomous turn in %s: %s", + session_id, e, + ) + + async def _close_autonomous_channel_stream(self, session_id: str) -> None: + """Release a session's autonomous-turn channel stream (idempotent).""" + router = getattr(self, "_router", None) + if router is None: + return + try: + await router.close_autonomous_stream(session_id) + except Exception as e: + logger.warning( + "Could not close channel stream for autonomous turn in %s: %s", + session_id, e, + ) + def _start_idle_watcher( self, session_id: str, client: Any, source: str, ) -> None: diff --git a/nerve/channels/router.py b/nerve/channels/router.py index 5469fba33..b5e1f7914 100644 --- a/nerve/channels/router.py +++ b/nerve/channels/router.py @@ -58,6 +58,11 @@ def __init__(self, engine: AgentEngine): self._pending_batches: dict[ str, list[tuple[InboundMessage, asyncio.Future[str]]] ] = {} + # Where each session was last messaged from: session_id -> + # (channel_name, target). Autonomous turns stream back there. + self._last_inbound: dict[str, tuple[str, str]] = {} + # Open autonomous-turn streams: session_id -> (listener_id, adapter) + self._autonomous: dict[str, tuple[str, StreamAdapter]] = {} # ------------------------------------------------------------------ # # Channel registry # @@ -114,6 +119,8 @@ async def handle_message(self, msg: InboundMessage) -> str: msg.channel_key, source=msg.channel_name, actor=None, ) + self._last_inbound[session_id] = (msg.channel_name, msg.sender_id) + # Store message context for reaction support msg_id = msg.metadata.get("message_id") if msg.metadata else None if msg_id is not None: @@ -491,6 +498,57 @@ async def deliver( session_id=session_id or "", )) + # ------------------------------------------------------------------ # + # Autonomous turns: engine → channel # + # ------------------------------------------------------------------ # + # + # When a background task settles, the CLI runs an autonomous turn + # that no inbound message started. The engine streams it to the + # broadcaster, which reaches the web UI, but a channel's adapter only + # lives for the run of the message it answers. Without the methods + # below, a Telegram session never saw anything said in an autonomous + # turn. + + async def open_autonomous_stream(self, session_id: str) -> None: + """Stream an autonomous turn to the channel that last messaged the + session. + + A no-op when the session was never messaged through a router + channel (web UI, cron, workflow runs), or when a user run is + already streaming to the same target (the turn then reaches it + through that run's adapter). + """ + if session_id in self._autonomous: + return + last = self._last_inbound.get(session_id) + if last is None: + return + channel_name, target = last + channel = self._channels.get(channel_name) + if channel is None or (channel_name, target) in self._adapters: + return + + adapter = StreamAdapter(channel, target, session_id, send_empty=False) + listener_id = f"{channel_name}:{target}:autonomous" + self._autonomous[session_id] = (listener_id, adapter) + await adapter.initialize() + await broadcaster.register(session_id, listener_id, adapter.on_event) + + async def close_autonomous_stream(self, session_id: str) -> None: + """End a session's autonomous-turn stream (idempotent). + + If the turn ended without a ``done`` event (cancelled, failed), + the adapter is finished here: whatever text arrived is sent, and + an empty placeholder is deleted. + """ + entry = self._autonomous.pop(session_id, None) + if entry is None: + return + listener_id, adapter = entry + await broadcaster.unregister(session_id, listener_id) + if not adapter.finished: + await adapter.on_event(session_id, {"type": "done"}) + # ------------------------------------------------------------------ # # Streaming adapter lifecycle # # ------------------------------------------------------------------ # diff --git a/nerve/channels/stream_adapter.py b/nerve/channels/stream_adapter.py index 9a41661aa..d890dc5b5 100644 --- a/nerve/channels/stream_adapter.py +++ b/nerve/channels/stream_adapter.py @@ -29,6 +29,8 @@ class StreamAdapter: Created per inbound message by the ChannelRouter, registered as a broadcaster listener, and torn down after the agent run completes. + The router also creates one for each autonomous turn (see + ``ChannelRouter.open_autonomous_stream``). """ def __init__( @@ -36,10 +38,16 @@ def __init__( channel: BaseChannel, target: str, session_id: str, + *, + send_empty: bool = True, ): self.channel = channel self.target = target self.session_id = session_id + # A user run always answers, even if only with "(no response)". An + # autonomous turn that produced nothing should leave no trace. + self._send_empty = send_empty + self.finished = False # Streaming state self._buffer: str = "" @@ -143,6 +151,16 @@ async def _handle_token(self, content: str) -> None: pass # Edit failures are non-fatal async def _handle_done(self) -> None: + if self.finished: + return # a backstop "done" after the real one must not resend + self.finished = True + if not self._send_empty and not self._normalize_text(self._buffer): + if self._placeholder_id: + try: + await self.channel.delete_message(self.target, self._placeholder_id) + except Exception: + pass + return if self._supports_streaming and self._supports_edit and self._placeholder_id: # Send final text as a new message (triggers notification), # then delete the streaming placeholder. diff --git a/tests/test_autonomous_turns.py b/tests/test_autonomous_turns.py index 9033f19b0..baed32342 100644 --- a/tests/test_autonomous_turns.py +++ b/tests/test_autonomous_turns.py @@ -742,3 +742,194 @@ async def test_idle_watcher_resumes_parked_session_on_background_completion(): ) # Task settled → no longer live → the idle sweep may now reap the client. assert engine._has_live_background_tasks("s1") is False + + +# --------------------------------------------------------------------------- +# Autonomous turns → channel router (Telegram and other non-web channels) +# --------------------------------------------------------------------------- + + +def _record_router(engine: AgentEngine, calls: list) -> None: + async def _open(session_id): + calls.append(("open", session_id)) + + async def _close(session_id): + calls.append(("close", session_id)) + + engine._router = SimpleNamespace( + open_autonomous_stream=_open, close_autonomous_stream=_close, + ) + + +def _patch_finalize_into(engine: AgentEngine, calls: list) -> None: + async def _record(session_id, st, channel, bump_updated_at=True): + calls.append(("finalize", st.full_response_text)) + + engine._finalize_turn = _record # type: ignore[method-assign] + + +@pytest.mark.asyncio +async def test_drain_opens_channel_stream_and_closes_it_after_finalize(): + """The router gets the turn's events between open and close; close comes + after finalize, whose ``done`` sends the channel its final message.""" + engine = _make_engine() + calls: list = [] + _record_router(engine, calls) + _patch_finalize_into(engine, calls) + stream = _FakeStream([ + _sys_msg("init", cwd="/tmp"), + _assistant_text("Background job finished."), + _result_msg(), + ]) + + with patch("nerve.agent.engine.broadcaster") as bc: + bc.broadcast = AsyncMock() + bc.broadcast_token = AsyncMock() + bc.mark_turn_open = lambda sid: None + turns = await engine._drain_pending_messages( + "s1", _fake_client(stream), "web", None, + ) + + assert turns == 1 + assert calls[:3] == [ + ("open", "s1"), + ("finalize", "Background job finished."), + ("close", "s1"), + ] + # Any later close (the drain's finally) is a harmless repeat. + assert set(calls[3:]) <= {("close", "s1")} + + +@pytest.mark.asyncio +async def test_drain_closes_channel_stream_of_an_empty_turn(): + """An empty turn is dropped without finalize; its channel stream must + still close so the channel's placeholder goes away.""" + engine = _make_engine() + calls: list = [] + _record_router(engine, calls) + _patch_finalize_into(engine, calls) + stream = _FakeStream([_sys_msg("init", cwd="/tmp")]) + + with patch("nerve.agent.engine.broadcaster") as bc: + bc.broadcast = AsyncMock() + bc.mark_turn_open = lambda sid: None + turns = await engine._drain_pending_messages( + "s1", _fake_client(stream), "web", None, first_content_timeout=0.05, + ) + + assert turns == 0 + assert calls[0] == ("open", "s1") + assert ("close", "s1") in calls + assert not any(c[0] == "finalize" for c in calls) + + +@pytest.mark.asyncio +async def test_drain_closes_channel_stream_when_the_turn_times_out(): + engine = _make_engine() + calls: list = [] + _record_router(engine, calls) + _patch_finalize_into(engine, calls) + stream = _FakeStream([ + _sys_msg("init", cwd="/tmp"), + _assistant_text("partial"), + ]) + + with patch("nerve.agent.engine.broadcaster") as bc: + bc.broadcast = AsyncMock() + bc.broadcast_token = AsyncMock() + bc.mark_turn_open = lambda sid: None + engine.config.agent.cli_idle_timeout_seconds = 0.05 + with pytest.raises(asyncio.TimeoutError): + await engine._drain_pending_messages("s1", _fake_client(stream), "web", None) + + assert calls[0] == ("open", "s1") + assert calls[-1] == ("close", "s1") + + +@pytest.mark.asyncio +async def test_drain_without_a_router_is_unchanged(): + """Engines that never built a router (tests, headless) skip the hooks.""" + engine = _make_engine() + finalized = _patch_finalize(engine) + stream = _FakeStream([ + _sys_msg("init", cwd="/tmp"), + _assistant_text("hello"), + _result_msg(), + ]) + + with patch("nerve.agent.engine.broadcaster") as bc: + bc.broadcast = AsyncMock() + bc.broadcast_token = AsyncMock() + bc.mark_turn_open = lambda sid: None + turns = await engine._drain_pending_messages( + "s1", _fake_client(stream), "web", None, + ) + + assert turns == 1 + assert finalized[0].full_response_text == "hello" + + +@pytest.mark.asyncio +async def test_drain_gives_each_autonomous_turn_its_own_channel_stream(): + """A stream is closed after every turn: a finished adapter ignores + further events, so a second turn reusing it would never be sent.""" + engine = _make_engine() + calls: list = [] + _record_router(engine, calls) + _patch_finalize_into(engine, calls) + stream = _FakeStream([ + _sys_msg("init", cwd="/tmp"), + _assistant_text("first"), + _result_msg(), + _sys_msg("init", cwd="/tmp"), + _assistant_text("second"), + _result_msg(), + ]) + + with patch("nerve.agent.engine.broadcaster") as bc: + bc.broadcast = AsyncMock() + bc.broadcast_token = AsyncMock() + bc.mark_turn_open = lambda sid: None + turns = await engine._drain_pending_messages( + "s1", _fake_client(stream), "web", None, + ) + + assert turns == 2 + assert calls[:6] == [ + ("open", "s1"), ("finalize", "first"), ("close", "s1"), + ("open", "s1"), ("finalize", "second"), ("close", "s1"), + ] + + +@pytest.mark.asyncio +async def test_drain_closes_channel_stream_when_cancelled_mid_turn(): + """/stop cancels the drain without _close_turn; the finally must still + flush the channel stream.""" + engine = _make_engine() + engine.config.agent.cli_idle_timeout_seconds = 30 + calls: list = [] + _record_router(engine, calls) + _patch_finalize_into(engine, calls) + stream = _FakeStream([ + _sys_msg("init", cwd="/tmp"), + _assistant_text("partial"), + ]) + + with patch("nerve.agent.engine.broadcaster") as bc: + bc.broadcast = AsyncMock() + bc.broadcast_token = AsyncMock() + bc.mark_turn_open = lambda sid: None + task = asyncio.create_task( + engine._drain_pending_messages("s1", _fake_client(stream), "web", None), + ) + for _ in range(200): + if ("open", "s1") in calls: + break + await asyncio.sleep(0.01) + await asyncio.sleep(0.05) # parked waiting for the rest of the turn + task.cancel() + with pytest.raises(asyncio.CancelledError): + await task + + assert calls[0] == ("open", "s1") + assert calls[-1] == ("close", "s1") diff --git a/tests/test_router_autonomous.py b/tests/test_router_autonomous.py new file mode 100644 index 000000000..136d2854f --- /dev/null +++ b/tests/test_router_autonomous.py @@ -0,0 +1,228 @@ +"""Autonomous turns reach the channel that last messaged the session. + +When a background task settles, the CLI runs an autonomous turn that no +inbound message started. The engine broadcasts it, but a channel's stream +adapter only lives for the run of the message it answers, so a Telegram +session used to see nothing of it. ``ChannelRouter.open_autonomous_stream`` +and ``close_autonomous_stream`` give such a turn its own adapter. +""" + +from types import SimpleNamespace +from unittest.mock import AsyncMock, patch + +import pytest + +from nerve.agent.streaming import StreamBroadcaster +from nerve.channels.base import ( + BaseChannel, + ChannelCapability, + ChannelConstraints, + InboundMessage, + OutboundMessage, +) +from nerve.channels.router import ChannelRouter + + +class _FakeChannel(BaseChannel): + """A streaming, editable channel that records what it was asked to do.""" + + def __init__(self, name: str = "chat"): + self._name = name + self.sent: list[str] = [] + self.placeholders: list[str] = [] + self.deleted: list[str] = [] + + @property + def name(self) -> str: + return self._name + + @property + def capabilities(self) -> ChannelCapability: + return ChannelCapability.SEND_TEXT | ChannelCapability.STREAMING + + @property + def constraints(self) -> ChannelConstraints: + return ChannelConstraints( + max_message_length=4096, min_edit_interval=0.0, + supports_message_edit=True, + ) + + async def start(self) -> None: + pass + + async def stop(self) -> None: + pass + + async def send(self, message: OutboundMessage) -> None: + self.sent.append(message.text) + + async def send_placeholder(self, target: str, session_id: str) -> str | None: + placeholder_id = f"ph-{len(self.placeholders) + 1}" + self.placeholders.append(placeholder_id) + return placeholder_id + + async def edit_message(self, target: str, message_id: str, text: str) -> None: + pass + + async def delete_message(self, target: str, message_id: str) -> None: + self.deleted.append(message_id) + + +def _router(channel: _FakeChannel) -> ChannelRouter: + """A router whose engine answers every message instantly.""" + engine = SimpleNamespace( + sessions=SimpleNamespace( + get_active_session=AsyncMock(return_value="s1"), + set_active_session=AsyncMock(), + ), + run=AsyncMock(return_value="ok"), + register_task=lambda session_id, task: None, + ) + router = ChannelRouter(engine) + router.register(channel) + return router + + +async def _message(router: ChannelRouter, channel: _FakeChannel) -> None: + """One inbound user message, answered and torn down like in production.""" + with patch.object(ChannelRouter, "BATCH_DEBOUNCE", 0): + await router.handle_message(InboundMessage( + channel_name=channel.name, channel_key=f"{channel.name}:42", + sender_id="42", text="hi", + )) + + +@pytest.fixture +def bc(): + """A fresh broadcaster wired into the router module.""" + fresh = StreamBroadcaster() + with patch("nerve.channels.router.broadcaster", fresh): + yield fresh + + +@pytest.mark.asyncio +async def test_autonomous_turn_is_sent_to_the_last_inbound_channel(bc): + channel = _FakeChannel() + router = _router(channel) + await _message(router, channel) + channel.sent.clear() # the user run's own answer is not under test + + await router.open_autonomous_stream("s1") + await bc.broadcast_token("s1", "Background job finished.") + await bc.broadcast_done("s1") + await router.close_autonomous_stream("s1") + + assert channel.sent == ["Background job finished."] + # The streaming placeholder is replaced by the final message. + assert channel.deleted == [channel.placeholders[-1]] + assert "s1" not in bc._listeners + + +@pytest.mark.asyncio +async def test_session_never_messaged_through_a_channel_is_left_alone(bc): + channel = _FakeChannel() + router = _router(channel) # no inbound message: web UI, cron, workflow + + await router.open_autonomous_stream("s1") + await bc.broadcast_token("s1", "nobody is listening on the channel") + await bc.broadcast_done("s1") + await router.close_autonomous_stream("s1") + + assert channel.sent == [] + assert channel.placeholders == [] + + +@pytest.mark.asyncio +async def test_empty_turn_leaves_no_message(bc): + """No "(no response)" for a turn that produced nothing.""" + channel = _FakeChannel() + router = _router(channel) + await _message(router, channel) + channel.sent.clear() + + await router.open_autonomous_stream("s1") + placeholder = channel.placeholders[-1] + await router.close_autonomous_stream("s1") # no tokens, no done + + assert channel.sent == [] + assert channel.deleted[-1] == placeholder + + +@pytest.mark.asyncio +async def test_turn_cut_short_still_delivers_what_arrived(bc): + channel = _FakeChannel() + router = _router(channel) + await _message(router, channel) + channel.sent.clear() + + await router.open_autonomous_stream("s1") + await bc.broadcast_token("s1", "Half a thought") + await router.close_autonomous_stream("s1") # cancelled before done + + assert channel.sent == ["Half a thought"] + + +@pytest.mark.asyncio +async def test_backstop_done_after_the_real_one_does_not_resend(bc): + channel = _FakeChannel() + router = _router(channel) + await _message(router, channel) + channel.sent.clear() + + await router.open_autonomous_stream("s1") + await bc.broadcast_token("s1", "Once") + await bc.broadcast_done("s1") + await bc.broadcast_done("s1") + await router.close_autonomous_stream("s1") + + assert channel.sent == ["Once"] + + +@pytest.mark.asyncio +async def test_no_second_stream_while_a_user_run_streams_to_the_target(bc): + """A turn drained inside run() reaches the channel through the run's own + adapter; a second adapter would send everything twice.""" + channel = _FakeChannel() + router = _router(channel) + await _message(router, channel) + user_run_adapter = await router._setup_streaming(channel, "42", "s1") + placeholders_before = list(channel.placeholders) + + await router.open_autonomous_stream("s1") + + assert channel.placeholders == placeholders_before + assert "s1" not in router._autonomous + assert router._adapters[(channel.name, "42")] is user_run_adapter + + +@pytest.mark.asyncio +async def test_open_and_close_are_idempotent(bc): + channel = _FakeChannel() + router = _router(channel) + await _message(router, channel) + channel.sent.clear() + + await router.open_autonomous_stream("s1") + await router.open_autonomous_stream("s1") + await bc.broadcast_token("s1", "One message") + await bc.broadcast_done("s1") + await router.close_autonomous_stream("s1") + await router.close_autonomous_stream("s1") + + assert channel.sent == ["One message"] + assert len(bc._listeners.get("s1", [])) == 0 + + +@pytest.mark.asyncio +async def test_closing_does_not_unregister_a_user_run_listener(bc): + """The autonomous listener has its own id, so tearing it down leaves a + concurrently registered user-run adapter in place.""" + channel = _FakeChannel() + router = _router(channel) + await _message(router, channel) + + await router.open_autonomous_stream("s1") + await bc.register("s1", f"{channel.name}:42", AsyncMock()) + await router.close_autonomous_stream("s1") + + assert [cid for cid, _ in bc._listeners["s1"]] == [f"{channel.name}:42"]