From 6aa0fb02b136716987c4007dd9fd83a2abf1cff1 Mon Sep 17 00:00:00 2001 From: Wasim Date: Wed, 2 Sep 2026 10:07:28 +0530 Subject: [PATCH 1/3] fix: prevent AssistantAgent deadlock on tool call cancellation Previously, cancelling an in-flight tool call would cause asyncio.gather to raise without executing put_nowait(None), leaving the consumer loop blocking forever on stream.get(). Now return_exceptions=True ensures put_nowait(None) always executes, allowing the consumer loop to terminate properly. Fixes #7956 --- .../src/autogen_agentchat/agents/_assistant_agent.py | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/python/packages/autogen-agentchat/src/autogen_agentchat/agents/_assistant_agent.py b/python/packages/autogen-agentchat/src/autogen_agentchat/agents/_assistant_agent.py index 8b8316fb0a7c..75626a186857 100644 --- a/python/packages/autogen-agentchat/src/autogen_agentchat/agents/_assistant_agent.py +++ b/python/packages/autogen-agentchat/src/autogen_agentchat/agents/_assistant_agent.py @@ -1208,9 +1208,12 @@ async def _execute_tool_calls( stream=stream_queue, ) for call in function_calls - ] + ], + return_exceptions=True, ) # Signal the end of streaming by putting None in the queue. + # This MUST be in a finally block to ensure the consumer loop + # always terminates, even if a tool call raises (e.g. CancelledError). stream_queue.put_nowait(None) return results From 7b2d7cdfa391cef9fcb47bf0d959e1feb0003ec1 Mon Sep 17 00:00:00 2001 From: Wasim Date: Wed, 2 Sep 2026 10:20:28 +0530 Subject: [PATCH 2/3] fix: handle exceptions in tool call results gracefully When return_exceptions=True and a tool call raises (e.g. CancelledError), the result is an exception object rather than a tuple. This update converts exceptions to error FunctionExecutionResult objects so downstream processing doesn't fail with a TypeError. --- .../agents/_assistant_agent.py | 26 ++++++++++++++++--- 1 file changed, 23 insertions(+), 3 deletions(-) diff --git a/python/packages/autogen-agentchat/src/autogen_agentchat/agents/_assistant_agent.py b/python/packages/autogen-agentchat/src/autogen_agentchat/agents/_assistant_agent.py index 75626a186857..33e8ca05265c 100644 --- a/python/packages/autogen-agentchat/src/autogen_agentchat/agents/_assistant_agent.py +++ b/python/packages/autogen-agentchat/src/autogen_agentchat/agents/_assistant_agent.py @@ -1232,7 +1232,27 @@ async def _execute_tool_calls( # Wait for all tool calls to complete. executed_calls_and_results = await task - exec_results = [result for _, result in executed_calls_and_results] + + # Process results, converting any exceptions to error FunctionExecutionResult + exec_results: List[FunctionExecutionResult] = [] + processed_calls_and_results: List[Tuple[FunctionCall, FunctionExecutionResult]] = [] + for item in executed_calls_and_results: + if isinstance(item, BaseException): + # Tool call raised an exception (e.g. CancelledError) + # Create a placeholder FunctionCall and an error result + placeholder_call = FunctionCall(id="error", arguments="{}", name="error") + error_result = FunctionExecutionResult( + content=f"Tool execution error: {type(item).__name__}: {item}", + call_id=placeholder_call.id, + is_error=True, + name=placeholder_call.name, + ) + exec_results.append(error_result) + processed_calls_and_results.append((placeholder_call, error_result)) + else: + call, result = item + exec_results.append(result) + processed_calls_and_results.append((call, result)) # Yield ToolCallExecutionEvent tool_call_result_msg = ToolCallExecutionEvent( @@ -1247,7 +1267,7 @@ async def _execute_tool_calls( # STEP 4C: Check for handoff handoff_output = cls._check_and_handle_handoff( model_result=current_model_result, - executed_calls_and_results=executed_calls_and_results, + executed_calls_and_results=processed_calls_and_results, inner_messages=inner_messages, handoffs=handoffs, agent_name=agent_name, @@ -1318,7 +1338,7 @@ async def _execute_tool_calls( yield reflection_response else: yield cls._summarize_tool_use( - executed_calls_and_results=executed_calls_and_results, + executed_calls_and_results=processed_calls_and_results, inner_messages=inner_messages, handoffs=handoffs, tool_call_summary_format=tool_call_summary_format, From 527c2b51706dfcc046c2ed059188a75171b26738 Mon Sep 17 00:00:00 2001 From: Wasim Date: Wed, 2 Sep 2026 10:20:43 +0530 Subject: [PATCH 3/3] test: add regression test for stream cleanup after cancellation Verifies that AssistantAgent.on_messages_stream terminates cleanly when a tool call is cancelled mid-flight, preventing the deadlock described in issue #8092. --- .../tests/test_assistant_agent.py | 62 +++++++++++++++++++ 1 file changed, 62 insertions(+) diff --git a/python/packages/autogen-agentchat/tests/test_assistant_agent.py b/python/packages/autogen-agentchat/tests/test_assistant_agent.py index 935f2471045c..6715d1525053 100644 --- a/python/packages/autogen-agentchat/tests/test_assistant_agent.py +++ b/python/packages/autogen-agentchat/tests/test_assistant_agent.py @@ -2816,6 +2816,68 @@ async def test_reset_with_cancellation_token(self) -> None: # Context clear should be called mock_context.clear.assert_called_once() + @pytest.mark.asyncio + async def test_stream_terminates_after_tool_cancellation(self) -> None: + """Regression test: streaming operation must terminate cleanly when + a tool call is cancelled mid-flight. Previously, cancellation could + leave the stream blocking forever on queue.get(). + + See https://github.com/microsoft/autogen/issues/8092""" + + import asyncio + + # A tool that blocks until cancelled + async def blocking_tool() -> str: + # Wait forever - will only return when cancelled + await asyncio.sleep(3600) + + model_client = ReplayChatCompletionClient( + [ + CreateResult( + finish_reason="function_calls", + content=[FunctionCall(id="1", arguments="{}", name="blocking_tool")], + usage=RequestUsage(prompt_tokens=10, completion_tokens=5), + cached=False, + ), + ], + model_info={ + "function_calling": True, + "vision": False, + "json_output": False, + "family": ModelFamily.GPT_4O, + "structured_output": False, + }, + ) + + agent = AssistantAgent( + name="test_agent", + model_client=model_client, + tools=[blocking_tool], + ) + + cancellation_token = CancellationToken() + + # Run the stream in a task so we can cancel it + async def consume_stream() -> None: + async for event in agent.on_messages_stream( + [TextMessage(content="Test", source="user")], cancellation_token + ): + pass # We don't care about events, just that the stream terminates + + task = asyncio.create_task(consume_stream()) + + # Give the stream time to start and reach the tool call + await asyncio.sleep(0.1) + + # Cancel the operation - this should cause the stream to terminate + cancellation_token.cancel() + + # The stream must terminate (not hang forever) + try: + await asyncio.wait_for(task, timeout=2.0) + except (asyncio.TimeoutError, asyncio.CancelledError): + pytest.fail("Stream did not terminate after cancellation - deadlock detected") + class TestAssistantAgentStreamingEdgeCases: """Test suite for streaming edge cases and error scenarios."""