From 92e34c81f11619a0ca1be8a19fb97ed2d3f71be1 Mon Sep 17 00:00:00 2001 From: Maxim Date: Sun, 27 Sep 2026 19:50:58 +0200 Subject: [PATCH] fix(agent): close open steps before the step-limit reply finishes the run When the graph hit its recursion limit, the adapter had already started a step for the interrupted node and dropped its run state on the exception, so the recovery path sent RUN_FINISHED with that step still open. @ag-ui/client's verifier rejects that, and Slack showed "Cannot send 'RUN_FINISHED' while steps are still active" instead of the step-limit reply. Track the steps, text messages, and tool calls the stream opens, and close whatever is left innermost-first before the snapshots and the reply. Co-Authored-By: Claude Opus 5.5 --- agent/agui.py | 54 ++++++++++++++++++++ agent/tests/test_agui_recursion.py | 82 +++++++++++++++++++++++++++++- 2 files changed, 135 insertions(+), 1 deletion(-) diff --git a/agent/agui.py b/agent/agui.py index c0805de5..d5ab844b 100644 --- a/agent/agui.py +++ b/agent/agui.py @@ -8,9 +8,11 @@ from ag_ui.core import ( EventType, RunFinishedEvent, + StepFinishedEvent, TextMessageContentEvent, TextMessageEndEvent, TextMessageStartEvent, + ToolCallEndEvent, ) from copilotkit import LangGraphAGUIAgent from langgraph.errors import GraphRecursionError @@ -226,6 +228,49 @@ def build_agui_agent(graph, *, recursion_limit: int | None = None): ) +class OpenEvents: + """The steps, text messages, and tool calls a stream has opened but not closed.""" + + def __init__(self): + self.steps = [] + self.messages = [] + self.tool_calls = [] + + def track(self, event): + kind = getattr(event, "type", None) + if kind == EventType.STEP_STARTED: + self.steps.append(event.step_name) + elif kind == EventType.STEP_FINISHED: + _discard(self.steps, event.step_name) + elif kind == EventType.TEXT_MESSAGE_START: + self.messages.append(event.message_id) + elif kind == EventType.TEXT_MESSAGE_END: + _discard(self.messages, event.message_id) + elif kind == EventType.TOOL_CALL_START: + self.tool_calls.append(event.tool_call_id) + elif kind == EventType.TOOL_CALL_END: + _discard(self.tool_calls, event.tool_call_id) + + def close(self): + """End events for everything still open, innermost first.""" + for tool_call_id in reversed(self.tool_calls): + yield ToolCallEndEvent( + type=EventType.TOOL_CALL_END, tool_call_id=tool_call_id + ) + for message_id in reversed(self.messages): + yield TextMessageEndEvent( + type=EventType.TEXT_MESSAGE_END, message_id=message_id + ) + for step_name in reversed(self.steps): + yield StepFinishedEvent(type=EventType.STEP_FINISHED, step_name=step_name) + self.steps, self.messages, self.tool_calls = [], [], [] + + +def _discard(items, value): + if value in items: + items.remove(value) + + async def iter_agent_events(run, input_data, *, recovery=None): """Yield AG-UI events. A graph step-limit becomes a user message. @@ -239,17 +284,26 @@ async def iter_agent_events(run, input_data, *, recovery=None): request arrives without one and opens the run under that, so echoing the request's own value closes a run nobody opened and leaves the open one hanging. + + Whatever the stream left open is closed first. The adapter starts a step + per node and forgets it when the graph raises, and the client refuses + RUN_FINISHED while a step, text message, or tool call is still active — + the Channel then shows that refusal in Slack where the reply should be. """ started = None + open_events = OpenEvents() try: async for event in run(input_data): if getattr(event, "type", None) == EventType.RUN_STARTED: started = event + open_events.track(event) yield event except GraphRecursionError as error: logger.warning("[agent] the graph hit its step limit: %s", error) thread_id = started.thread_id if started else input_data.thread_id run_id = started.run_id if started else input_data.run_id + for event in open_events.close(): + yield event if recovery is not None: async for event in recovery(thread_id, run_id): yield event diff --git a/agent/tests/test_agui_recursion.py b/agent/tests/test_agui_recursion.py index d735543d..bf89022c 100644 --- a/agent/tests/test_agui_recursion.py +++ b/agent/tests/test_agui_recursion.py @@ -4,7 +4,14 @@ from types import SimpleNamespace import pytest -from ag_ui.core import EventType, RunAgentInput, RunStartedEvent +from ag_ui.core import ( + EventType, + RunAgentInput, + RunStartedEvent, + StepStartedEvent, + TextMessageStartEvent, + ToolCallStartEvent, +) from langchain_core.messages import AIMessage from langgraph.checkpoint.memory import MemorySaver from langgraph.errors import GraphRecursionError @@ -170,3 +177,76 @@ def test_the_recursion_snapshot_keeps_the_messages_the_run_committed(): ) assert [message.content for message in snapshot.messages].count("thinking") >= 1 + + +def still_open_at_run_finished(events): + """What `@ag-ui/client`'s verifier would still see open at RUN_FINISHED. + + The client refuses RUN_FINISHED while any step, text message, or tool call + is active, and the Channel shows that refusal to the person in Slack in + place of the reply. A raw SSE read of the stream looks fine, because + nothing on that side checks the order. + """ + steps, messages, tool_calls = set(), set(), set() + for event in events: + if event.type == EventType.STEP_STARTED: + steps.add(event.step_name) + elif event.type == EventType.STEP_FINISHED: + steps.discard(event.step_name) + elif event.type == EventType.TEXT_MESSAGE_START: + messages.add(event.message_id) + elif event.type == EventType.TEXT_MESSAGE_END: + messages.discard(event.message_id) + elif event.type == EventType.TOOL_CALL_START: + tool_calls.add(event.tool_call_id) + elif event.type == EventType.TOOL_CALL_END: + tool_calls.discard(event.tool_call_id) + elif event.type == EventType.RUN_FINISHED: + return {"steps": steps, "messages": messages, "tool_calls": tool_calls} + raise AssertionError("the run never finished") + + +def test_the_recursion_exit_closes_the_step_the_graph_was_in(): + # The adapter opens a step per node and forgets it when the graph raises, + # so the step the limit interrupted is still open when the reply finishes + # the run. The client rejects that RUN_FINISHED, and Slack shows the + # rejection instead of the reply. + events = recursion_run(looping_agent()) + + assert still_open_at_run_finished(events) == { + "steps": set(), + "messages": set(), + "tool_calls": set(), + } + + +def test_the_recursion_exit_closes_everything_the_stream_left_open(): + async def interrupted(_input): + yield RunStartedEvent(type=EventType.RUN_STARTED, thread_id="t", run_id="r") + yield StepStartedEvent(type=EventType.STEP_STARTED, step_name="model") + yield TextMessageStartEvent( + type=EventType.TEXT_MESSAGE_START, role="assistant", message_id="m1" + ) + yield ToolCallStartEvent( + type=EventType.TOOL_CALL_START, + tool_call_id="c1", + tool_call_name="search", + parent_message_id="m1", + ) + raise GraphRecursionError("Recursion limit of 25 reached") + + events = asyncio.run( + _collect(iter_agent_events(interrupted, SimpleNamespace(thread_id="t", run_id="r"))) + ) + + assert still_open_at_run_finished(events) == { + "steps": set(), + "messages": set(), + "tool_calls": set(), + } + closing = [event.type for event in events[4:7]] + assert closing == [ + EventType.TOOL_CALL_END, + EventType.TEXT_MESSAGE_END, + EventType.STEP_FINISHED, + ], "close innermost first, before the reply opens its own message"