Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
35 changes: 35 additions & 0 deletions nerve/agent/engine.py
Original file line number Diff line number Diff line change
Expand Up @@ -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 (
Expand All @@ -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:
Expand Down Expand Up @@ -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:
Expand Down
58 changes: 58 additions & 0 deletions nerve/channels/router.py
Original file line number Diff line number Diff line change
Expand Up @@ -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 #
Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -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 #
# ------------------------------------------------------------------ #
Expand Down
18 changes: 18 additions & 0 deletions nerve/channels/stream_adapter.py
Original file line number Diff line number Diff line change
Expand Up @@ -29,17 +29,25 @@ 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__(
self,
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 = ""
Expand Down Expand Up @@ -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.
Expand Down
191 changes: 191 additions & 0 deletions tests/test_autonomous_turns.py
Original file line number Diff line number Diff line change
Expand Up @@ -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")
Loading
Loading