From f049cb58c3cd5ed31b9a4f1bbb757e73ddd8b716 Mon Sep 17 00:00:00 2001 From: Manohar Paturi <186662190+ManoharPaturi@users.noreply.github.com> Date: Thu, 24 Sep 2026 21:01:15 +0530 Subject: [PATCH 1/6] Python: filter non-assistant messages from workflow agent responses When running workflows converted to an agent via workflow.as_agent(), WorkflowAgent._convert_workflow_events_to_agent_response and _convert_workflow_event_to_agent_response_updates forwarded messages without checking their role. When underlying executors or orchestrations such as GroupChat yield conversation history containing user inputs, system instructions, or tool outputs, these non-assistant messages were re-emitted as caller-facing agent responses and streaming chunks. This change filters out non-assistant messages across AgentResponse, list[Message], Message, and AgentResponseUpdate workflow outputs so that WorkflowAgent exclusively returns assistant output. If an event contains only non-assistant messages, it is safely excluded from agent output without altering internal workflow execution. Signed-off-by: Manohar Paturi <186662190+ManoharPaturi@users.noreply.github.com> --- .../core/agent_framework/_workflows/_agent.py | 50 +++-- .../tests/workflow/test_workflow_agent.py | 194 +++++++++++++++++- 2 files changed, 227 insertions(+), 17 deletions(-) diff --git a/python/packages/core/agent_framework/_workflows/_agent.py b/python/packages/core/agent_framework/_workflows/_agent.py index f4aa653310d..82293770a39 100644 --- a/python/packages/core/agent_framework/_workflows/_agent.py +++ b/python/packages/core/agent_framework/_workflows/_agent.py @@ -587,23 +587,32 @@ def _convert_workflow_events_to_agent_response( ) if isinstance(data, AgentResponse): - messages.extend(data.messages) - raw_representations.append(data.raw_representation) - merged_usage = add_usage_details(merged_usage, data.usage_details) - latest_created_at = ( - data.created_at - if not latest_created_at - else max(latest_created_at, data.created_at) - if data.created_at - else latest_created_at - ) + # Filter to only assistant messages — system, tool, and user messages + # are intentionally excluded. System prompts and tool results are + # internal workflow artifacts; user messages would be re-emitted + # (e.g., from GroupChat orchestrators that include full conversation history). + assistant_messages = [msg for msg in data.messages if msg.role == "assistant"] + if assistant_messages: + messages.extend(assistant_messages) + raw_representations.append(data.raw_representation) + merged_usage = add_usage_details(merged_usage, data.usage_details) + latest_created_at = ( + data.created_at + if not latest_created_at + else max(latest_created_at, data.created_at) + if data.created_at + else latest_created_at + ) elif isinstance(data, Message): - messages.append(data) - raw_representations.append(data.raw_representation) + if data.role == "assistant": + messages.append(data) + raw_representations.append(data.raw_representation) elif is_instance_of(data, list[Message]): chat_messages = cast(list[Message], data) - messages.extend(chat_messages) - raw_representations.append(data) + assistant_messages = [msg for msg in chat_messages if msg.role == "assistant"] + if assistant_messages: + messages.extend(assistant_messages) + raw_representations.append(data) else: contents = self._extract_contents(data) if not contents: @@ -654,6 +663,9 @@ def _convert_workflow_event_to_agent_response_updates( executor_id = event.executor_id if isinstance(data, AgentResponseUpdate): + # Filter out non-assistant updates (e.g. user input echoed back) + if data.role is not None and data.role != "assistant": + return [] # Construct a fresh AgentResponseUpdate so we don't mutate a payload # that AgentExecutor still holds a reference to in its `updates` list. return [ @@ -676,9 +688,11 @@ def _convert_workflow_event_to_agent_response_updates( ) ] if isinstance(data, AgentResponse): - # Convert each message in AgentResponse to an AgentResponseUpdate + # Convert each assistant message in AgentResponse to an AgentResponseUpdate updates: list[AgentResponseUpdate] = [] for msg in data.messages: + if msg.role != "assistant": + continue updates.append( AgentResponseUpdate( contents=list(msg.contents), @@ -698,6 +712,8 @@ def _convert_workflow_event_to_agent_response_updates( updates[-1].additional_properties = dict(data.additional_properties) return updates if isinstance(data, Message): + if data.role != "assistant": + return [] return [ AgentResponseUpdate( contents=list(data.contents), @@ -710,10 +726,12 @@ def _convert_workflow_event_to_agent_response_updates( ) ] if is_instance_of(data, list[Message]): - # Convert each Message to an AgentResponseUpdate + # Convert each assistant Message to an AgentResponseUpdate chat_messages = cast(list[Message], data) updates = [] for msg in chat_messages: + if msg.role != "assistant": + continue updates.append( AgentResponseUpdate( contents=list(msg.contents), diff --git a/python/packages/core/tests/workflow/test_workflow_agent.py b/python/packages/core/tests/workflow/test_workflow_agent.py index 0ef70ca2f6e..a09d8f52ecc 100644 --- a/python/packages/core/tests/workflow/test_workflow_agent.py +++ b/python/packages/core/tests/workflow/test_workflow_agent.py @@ -1105,7 +1105,7 @@ async def test_workflow_as_agent_yield_output_with_list_of_chat_messages(self) - async def list_yielding_executor(messages: list[Message], ctx: WorkflowContext[Never, list[Message]]) -> None: # type: ignore[valid-type] # Yield a list of Messages (as SequentialBuilder does) msg_list = [ - Message(role="user", contents=["first message"]), + Message(role="assistant", contents=["first message"]), Message(role="assistant", contents=["second message"]), Message( role="assistant", @@ -2597,3 +2597,195 @@ async def start(messages: list[Message], ctx: WorkflowContext[AgentExecutorReque pending = await workflow._runner_context.get_pending_request_info_events() # The agent's approval id is used as the workflow's pending request id. assert list(pending.keys()) == [approval_id] + + async def test_workflow_as_agent_filters_non_assistant_messages_from_agent_response(self) -> None: + """Verify WorkflowAgent filters user, system, and tool messages from AgentResponse.""" + + @executor + async def mixed_agent_response_executor( + messages: list[Message], ctx: WorkflowContext[Never, AgentResponse] + ) -> None: + response = AgentResponse( + messages=[ + Message(role="system", contents=["System instructions"]), + Message(role="user", contents=["User input question"]), + Message(role="assistant", contents=["Assistant answer"], author_name="Teacher"), + Message(role="tool", contents=["Tool execution result"]), + ] + ) + await ctx.yield_output(response) + + workflow = WorkflowBuilder(start_executor=mixed_agent_response_executor).build() + agent = workflow.as_agent("mixed-response-agent") + + # Test streaming path + updates: list[AgentResponseUpdate] = [] + async for chunk in agent.run("hello", stream=True): + updates.append(chunk) + + assert len(updates) == 1 + assert updates[0].role == "assistant" + assert updates[0].text == "Assistant answer" + assert updates[0].author_name == "Teacher" + + # Test non-streaming path + result = await agent.run("hello") + assert len(result.messages) == 1 + assert result.messages[0].role == "assistant" + assert result.messages[0].text == "Assistant answer" + assert result.messages[0].author_name == "Teacher" + + async def test_workflow_as_agent_filters_non_assistant_messages_from_list_of_messages(self) -> None: + """Verify WorkflowAgent filters user, system, and tool messages from list[Message].""" + + @executor + async def mixed_list_executor(messages: list[Message], ctx: WorkflowContext[Never, list[Message]]) -> None: + await ctx.yield_output([ + Message(role="user", contents=["what is 2+2?"]), + Message(role="assistant", contents=["4"], author_name="Maths"), + Message(role="system", contents=["system prompt"]), + Message(role="assistant", contents=["four"], author_name="English"), + ]) + + workflow = WorkflowBuilder(start_executor=mixed_list_executor).build() + agent = workflow.as_agent("mixed-list-agent") + + # Test streaming path + updates: list[AgentResponseUpdate] = [] + async for chunk in agent.run("calc", stream=True): + updates.append(chunk) + + assert len(updates) == 2 + for update in updates: + assert update.role == "assistant" + assert updates[0].author_name == "Maths" + assert updates[0].text == "4" + assert updates[1].author_name == "English" + assert updates[1].text == "four" + + # Test non-streaming path + result = await agent.run("calc") + assert len(result.messages) == 2 + for message in result.messages: + assert message.role == "assistant" + assert result.messages[0].author_name == "Maths" + assert result.messages[0].text == "4" + assert result.messages[1].author_name == "English" + assert result.messages[1].text == "four" + + async def test_workflow_as_agent_filters_single_non_assistant_message(self) -> None: + """Verify WorkflowAgent filters a single Message when role is not assistant.""" + + @executor + async def user_message_executor(messages: list[Message], ctx: WorkflowContext[Never, Message]) -> None: + await ctx.yield_output(Message(role="user", contents=["echoed user message"])) + + workflow = WorkflowBuilder(start_executor=user_message_executor).build() + agent = workflow.as_agent("user-msg-agent") + + # Streaming should yield no updates + updates: list[AgentResponseUpdate] = [] + async for chunk in agent.run("test", stream=True): + updates.append(chunk) + assert len(updates) == 0 + + # Non-streaming should produce empty messages list + result = await agent.run("test") + assert len(result.messages) == 0 + + async def test_workflow_as_agent_filters_user_agent_response_update(self) -> None: + """Verify WorkflowAgent drops AgentResponseUpdate when role is user.""" + + @executor + async def update_yielding_executor( + messages: list[Message], ctx: WorkflowContext[Never, AgentResponseUpdate] + ) -> None: + await ctx.yield_output(AgentResponseUpdate(contents=[Content.from_text(text="echo")], role="user")) + await ctx.yield_output(AgentResponseUpdate(contents=[Content.from_text(text="answer")], role="assistant")) + + workflow = WorkflowBuilder(start_executor=update_yielding_executor).build() + agent = workflow.as_agent("update-agent") + + updates: list[AgentResponseUpdate] = [] + async for chunk in agent.run("test", stream=True): + updates.append(chunk) + + assert len(updates) == 1 + assert updates[0].role == "assistant" + assert updates[0].text == "answer" + + async def test_workflow_as_agent_empty_after_filtering(self) -> None: + """Verify WorkflowAgent handles all non-assistant messages without crashing.""" + + @executor + async def non_assistant_only_executor( + messages: list[Message], ctx: WorkflowContext[Never, AgentResponse] + ) -> None: + response = AgentResponse( + messages=[ + Message(role="user", contents=["user msg"]), + Message(role="system", contents=["system msg"]), + Message(role="tool", contents=["tool msg"]), + ] + ) + await ctx.yield_output(response) + + workflow = WorkflowBuilder(start_executor=non_assistant_only_executor).build() + agent = workflow.as_agent("all-filtered-agent") + + result = await agent.run("test") + assert len(result.messages) == 0 + assert len(result.raw_representation) == 0 + + async def test_workflow_as_agent_multi_turn_user_input_not_compounded(self) -> None: + """Verify user messages in conversation history do not compound into responses across turns.""" + + class HistoryYieldingExecutor(Executor): + @handler + async def handle_messages( + self, + messages: list[Message], + ctx: WorkflowContext[Never, AgentResponse], + ) -> None: + user_text = messages[-1].text or "" + # Simulates orchestrators that include full conversation history in output + full_history = [ + Message(role="user", contents=[user_text]), + Message(role="assistant", contents=[f"Answer: {user_text}"], author_name="Agent"), + ] + await ctx.yield_output(AgentResponse(messages=full_history)) + + workflow = WorkflowBuilder(start_executor=HistoryYieldingExecutor(id="history-exec")).build() + agent = workflow.as_agent("history-agent") + session = AgentSession() + + # Turn 1 non-streaming + resp1 = await agent.run("first_query", session=session) + assert len(resp1.messages) == 1 + assert resp1.messages[0].role == "assistant" + assert resp1.text == "Answer: first_query" + assert "first_query" not in (resp1.text.replace("Answer: first_query", "")) + + # Turn 2 non-streaming: first_query must not bleed into turn 2 + resp2 = await agent.run("second_query", session=session) + assert len(resp2.messages) == 1 + assert resp2.messages[0].role == "assistant" + assert resp2.text == "Answer: second_query" + + # Streaming check + streaming_agent = workflow.as_agent("streaming-history-agent") + streaming_session = AgentSession() + + chunks1: list[AgentResponseUpdate] = [] + async for chunk in streaming_agent.run("stream_q1", stream=True, session=streaming_session): + chunks1.append(chunk) + assert len(chunks1) == 1 + assert chunks1[0].role == "assistant" + assert chunks1[0].text == "Answer: stream_q1" + + chunks2: list[AgentResponseUpdate] = [] + async for chunk in streaming_agent.run("stream_q2", stream=True, session=streaming_session): + chunks2.append(chunk) + assert len(chunks2) == 1 + assert chunks2[0].role == "assistant" + assert chunks2[0].text == "Answer: stream_q2" From a6eff9e8e7463e3cd17452975fc9e99352526a3e Mon Sep 17 00:00:00 2001 From: Manohar Paturi <186662190+ManoharPaturi@users.noreply.github.com> Date: Fri, 25 Sep 2026 18:36:51 +0530 Subject: [PATCH 2/6] Fix test typing checks in workflow agent tests Signed-off-by: Manohar Paturi <186662190+ManoharPaturi@users.noreply.github.com> --- .../tests/workflow/test_workflow_agent.py | 23 +++++++++++++------ 1 file changed, 16 insertions(+), 7 deletions(-) diff --git a/python/packages/core/tests/workflow/test_workflow_agent.py b/python/packages/core/tests/workflow/test_workflow_agent.py index a09d8f52ecc..d04f01499af 100644 --- a/python/packages/core/tests/workflow/test_workflow_agent.py +++ b/python/packages/core/tests/workflow/test_workflow_agent.py @@ -2603,7 +2603,8 @@ async def test_workflow_as_agent_filters_non_assistant_messages_from_agent_respo @executor async def mixed_agent_response_executor( - messages: list[Message], ctx: WorkflowContext[Never, AgentResponse] + messages: list[Message], + ctx: WorkflowContext[Never, AgentResponse], # type: ignore[valid-type] ) -> None: response = AgentResponse( messages=[ @@ -2639,7 +2640,10 @@ async def test_workflow_as_agent_filters_non_assistant_messages_from_list_of_mes """Verify WorkflowAgent filters user, system, and tool messages from list[Message].""" @executor - async def mixed_list_executor(messages: list[Message], ctx: WorkflowContext[Never, list[Message]]) -> None: + async def mixed_list_executor( + messages: list[Message], + ctx: WorkflowContext[Never, list[Message]], # type: ignore[valid-type] + ) -> None: await ctx.yield_output([ Message(role="user", contents=["what is 2+2?"]), Message(role="assistant", contents=["4"], author_name="Maths"), @@ -2677,7 +2681,10 @@ async def test_workflow_as_agent_filters_single_non_assistant_message(self) -> N """Verify WorkflowAgent filters a single Message when role is not assistant.""" @executor - async def user_message_executor(messages: list[Message], ctx: WorkflowContext[Never, Message]) -> None: + async def user_message_executor( + messages: list[Message], + ctx: WorkflowContext[Never, Message], # type: ignore[valid-type] + ) -> None: await ctx.yield_output(Message(role="user", contents=["echoed user message"])) workflow = WorkflowBuilder(start_executor=user_message_executor).build() @@ -2698,7 +2705,8 @@ async def test_workflow_as_agent_filters_user_agent_response_update(self) -> Non @executor async def update_yielding_executor( - messages: list[Message], ctx: WorkflowContext[Never, AgentResponseUpdate] + messages: list[Message], + ctx: WorkflowContext[Never, AgentResponseUpdate], # type: ignore[valid-type] ) -> None: await ctx.yield_output(AgentResponseUpdate(contents=[Content.from_text(text="echo")], role="user")) await ctx.yield_output(AgentResponseUpdate(contents=[Content.from_text(text="answer")], role="assistant")) @@ -2719,7 +2727,8 @@ async def test_workflow_as_agent_empty_after_filtering(self) -> None: @executor async def non_assistant_only_executor( - messages: list[Message], ctx: WorkflowContext[Never, AgentResponse] + messages: list[Message], + ctx: WorkflowContext[Never, AgentResponse], # type: ignore[valid-type] ) -> None: response = AgentResponse( messages=[ @@ -2735,7 +2744,7 @@ async def non_assistant_only_executor( result = await agent.run("test") assert len(result.messages) == 0 - assert len(result.raw_representation) == 0 + assert not result.raw_representation async def test_workflow_as_agent_multi_turn_user_input_not_compounded(self) -> None: """Verify user messages in conversation history do not compound into responses across turns.""" @@ -2745,7 +2754,7 @@ class HistoryYieldingExecutor(Executor): async def handle_messages( self, messages: list[Message], - ctx: WorkflowContext[Never, AgentResponse], + ctx: WorkflowContext[Never, AgentResponse], # type: ignore[valid-type] ) -> None: user_text = messages[-1].text or "" # Simulates orchestrators that include full conversation history in output From cf2a1c0967eaa4621d0e31990cae8deb36448149 Mon Sep 17 00:00:00 2001 From: Manohar Paturi <186662190+ManoharPaturi@users.noreply.github.com> Date: Mon, 28 Sep 2026 19:00:21 +0530 Subject: [PATCH 3/6] fix(workflows): keep non-assistant messages out of raw_representation the list[Message] branch appended the whole yielded list as the raw_representation even when non-assistant entries were filtered out of the public messages, so user/system/tool content stayed reachable through the non-streaming AgentResponse payload. when the list is filtered, append each assistant message's own raw_representation instead; unfiltered lists keep the original object. Signed-off-by: Manohar Paturi <186662190+ManoharPaturi@users.noreply.github.com> --- .../packages/core/agent_framework/_workflows/_agent.py | 9 ++++++++- .../packages/core/tests/workflow/test_workflow_agent.py | 5 +++++ 2 files changed, 13 insertions(+), 1 deletion(-) diff --git a/python/packages/core/agent_framework/_workflows/_agent.py b/python/packages/core/agent_framework/_workflows/_agent.py index 82293770a39..38fa8ea0a8d 100644 --- a/python/packages/core/agent_framework/_workflows/_agent.py +++ b/python/packages/core/agent_framework/_workflows/_agent.py @@ -612,7 +612,14 @@ def _convert_workflow_events_to_agent_response( assistant_messages = [msg for msg in chat_messages if msg.role == "assistant"] if assistant_messages: messages.extend(assistant_messages) - raw_representations.append(data) + # raw_representation of a filtered list must not leak the + # non-assistant entries the public messages list dropped. + if len(assistant_messages) == len(chat_messages): + raw_representations.append(data) + else: + raw_representations.extend( + msg.raw_representation for msg in assistant_messages + ) else: contents = self._extract_contents(data) if not contents: diff --git a/python/packages/core/tests/workflow/test_workflow_agent.py b/python/packages/core/tests/workflow/test_workflow_agent.py index d04f01499af..0d4299ef543 100644 --- a/python/packages/core/tests/workflow/test_workflow_agent.py +++ b/python/packages/core/tests/workflow/test_workflow_agent.py @@ -2677,6 +2677,11 @@ async def mixed_list_executor( assert result.messages[1].author_name == "English" assert result.messages[1].text == "four" + # raw_representation of the non-streaming result must not leak the + # filtered-out user/system messages through the public payload. + for rep in result.raw_representation or []: + assert rep is None or (isinstance(rep, Message) and rep.role == "assistant") + async def test_workflow_as_agent_filters_single_non_assistant_message(self) -> None: """Verify WorkflowAgent filters a single Message when role is not assistant.""" From f320aed5165f9da51327b6e9fa237efe191bdc09 Mon Sep 17 00:00:00 2001 From: Manohar Paturi <186662190+ManoharPaturi@users.noreply.github.com> Date: Tue, 29 Sep 2026 19:56:54 +0530 Subject: [PATCH 4/6] fix(workflows): drop assistant function calls orphaned by role filtering MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Role-only filtering kept assistant(function_call) while dropping its tool(function_result), producing an invalid call-without-result transcript that providers reject on replay when the response is persisted through a session. Assistant messages whose contents are entirely function-call envelopes are now dropped alongside the non-assistant messages, so the public response carries only user-facing assistant output with no dangling calls. Regression test covers call → tool result → final. Signed-off-by: Manohar Paturi <186662190+ManoharPaturi@users.noreply.github.com> --- .../core/agent_framework/_workflows/_agent.py | 52 ++++++++++++++----- .../tests/workflow/test_workflow_agent.py | 43 ++++++++++++++- 2 files changed, 81 insertions(+), 14 deletions(-) diff --git a/python/packages/core/agent_framework/_workflows/_agent.py b/python/packages/core/agent_framework/_workflows/_agent.py index 38fa8ea0a8d..ea96af26da9 100644 --- a/python/packages/core/agent_framework/_workflows/_agent.py +++ b/python/packages/core/agent_framework/_workflows/_agent.py @@ -3,12 +3,19 @@ from __future__ import annotations import logging -import sys import uuid from collections.abc import AsyncIterable, Awaitable, Callable, Mapping, Sequence from dataclasses import dataclass from datetime import datetime, timezone -from typing import TYPE_CHECKING, Any, ClassVar, Literal, cast, overload +from typing import ( + TYPE_CHECKING, + Any, + ClassVar, + Literal, + TypedDict, # pragma: no cover + cast, + overload, +) from .._agents import BaseAgent from .._sessions import ( @@ -41,17 +48,30 @@ from ._message_utils import normalize_messages_input from ._typing_utils import is_instance_of, is_type_compatible -if sys.version_info >= (3, 11): - from typing import TypedDict # pragma: no cover -else: - from typing_extensions import TypedDict # pragma: no cover - if TYPE_CHECKING: from ._workflow import Workflow, WorkflowInvocationKwargs logger = logging.getLogger(__name__) +def _is_orphaned_function_call(message: Message) -> bool: + """Drop a call whose result cannot follow it. + + True when an assistant message is only a function call whose tool result + is not assistant-role, so keeping the call alone would produce an invalid + (unpaired) transcript for providers that validate call/result history. + """ + if message.role != "assistant": + return False + contents = list(getattr(message, "contents", []) or []) + if not contents: + return False + # Orphaned only when every content is a function-call envelope + return all( + getattr(content, "type", None) in ("function_call", "function_approval_response") for content in contents + ) + + class WorkflowAgent(BaseAgent): """An `Agent` subclass that wraps a workflow and exposes it as an agent.""" @@ -591,7 +611,13 @@ def _convert_workflow_events_to_agent_response( # are intentionally excluded. System prompts and tool results are # internal workflow artifacts; user messages would be re-emitted # (e.g., from GroupChat orchestrators that include full conversation history). - assistant_messages = [msg for msg in data.messages if msg.role == "assistant"] + # Assistant messages that are bare function calls are also dropped + # when their tool result is not assistant-role: a call without its + # result is an invalid transcript for providers that validate + # call/result pairing on replay. + assistant_messages = [ + msg for msg in data.messages if msg.role == "assistant" and not _is_orphaned_function_call(msg) + ] if assistant_messages: messages.extend(assistant_messages) raw_representations.append(data.raw_representation) @@ -609,7 +635,11 @@ def _convert_workflow_events_to_agent_response( raw_representations.append(data.raw_representation) elif is_instance_of(data, list[Message]): chat_messages = cast(list[Message], data) - assistant_messages = [msg for msg in chat_messages if msg.role == "assistant"] + # Keep tool results that pair with surviving assistant calls so the + # transcript stays valid; drop user/system and orphaned calls. + assistant_messages = [ + msg for msg in chat_messages if msg.role == "assistant" and not _is_orphaned_function_call(msg) + ] if assistant_messages: messages.extend(assistant_messages) # raw_representation of a filtered list must not leak the @@ -617,9 +647,7 @@ def _convert_workflow_events_to_agent_response( if len(assistant_messages) == len(chat_messages): raw_representations.append(data) else: - raw_representations.extend( - msg.raw_representation for msg in assistant_messages - ) + raw_representations.extend(msg.raw_representation for msg in assistant_messages) else: contents = self._extract_contents(data) if not contents: diff --git a/python/packages/core/tests/workflow/test_workflow_agent.py b/python/packages/core/tests/workflow/test_workflow_agent.py index 0d4299ef543..a04c572a21f 100644 --- a/python/packages/core/tests/workflow/test_workflow_agent.py +++ b/python/packages/core/tests/workflow/test_workflow_agent.py @@ -3,10 +3,9 @@ import uuid from collections.abc import Awaitable, Sequence from dataclasses import dataclass -from typing import Any, Literal, cast, overload +from typing import Any, Literal, Never, assert_type, cast, overload import pytest -from typing_extensions import Never, assert_type from agent_framework import ( AgentExecutorRequest, @@ -2682,6 +2681,46 @@ async def mixed_list_executor( for rep in result.raw_representation or []: assert rep is None or (isinstance(rep, Message) and rep.role == "assistant") + async def test_workflow_as_agent_drops_orphaned_function_calls(self) -> None: + """assistant(function_call) whose tool result is tool-role must not survive filtering. + + Keeping the call without its result would produce an invalid transcript for + providers that validate call/result pairing on replay (eavanvalkenburg's review). + """ + + @executor + async def tool_transcript_executor( + messages: list[Message], + ctx: WorkflowContext[Never, list[Message]], # type: ignore[valid-type] + ) -> None: + await ctx.yield_output([ + Message( + role="assistant", + contents=[ + Content.from_function_call(call_id="call-1", name="get_weather", arguments={"city": "Paris"}), + ], + ), + Message( + role="tool", + contents=[ + Content.from_function_result(call_id="call-1", result="18C"), + ], + ), + Message(role="assistant", contents=[Content.from_text("It is 18C in Paris.")]), + ]) + + workflow = WorkflowBuilder(start_executor=tool_transcript_executor).build() + agent = workflow.as_agent("tool-transcript-agent") + + result = await agent.run("weather") + + # The orphaned call is dropped; the final user-facing answer survives. + assert all(msg.role == "assistant" for msg in result.messages) + assert not any( + getattr(content, "type", None) == "function_call" for msg in result.messages for content in msg.contents + ) + assert any("18C" in (getattr(content, "text", "") or "") for msg in result.messages for content in msg.contents) + async def test_workflow_as_agent_filters_single_non_assistant_message(self) -> None: """Verify WorkflowAgent filters a single Message when role is not assistant.""" From bd8f1e82df0fb21c8b1969035326ca1d72ed1e6c Mon Sep 17 00:00:00 2001 From: Manohar Paturi Date: Fri, 2 Oct 2026 07:37:00 +0530 Subject: [PATCH 5/6] fix(workflows): unify caller-facing message policy and track stream roles One helper now decides what a workflow-as-agent caller sees, for both response modes: _caller_facing_messages drops non-assistant messages and assistant messages carrying function-call envelopes whole, reasoning included, since their tool-role results are excluded by the role policy and a call without its result is an invalid replay transcript. Streaming gains a _StreamingRoleGate that carries the last declared role across updates, so role-less continuation chunks of a user stream can no longer leak through the assistant-only filter. The foundry-hosting checkpoint resume path uses the same gate. Also imports Never/assert_type from typing_extensions; they do not exist in Python 3.10's typing module and broke the 3.10 CI matrix. --- .../core/agent_framework/_workflows/_agent.py | 105 ++++++++++------- .../tests/workflow/test_workflow_agent.py | 110 +++++++++++++++++- .../_responses.py | 4 +- 3 files changed, 173 insertions(+), 46 deletions(-) diff --git a/python/packages/core/agent_framework/_workflows/_agent.py b/python/packages/core/agent_framework/_workflows/_agent.py index ea96af26da9..20ba723746b 100644 --- a/python/packages/core/agent_framework/_workflows/_agent.py +++ b/python/packages/core/agent_framework/_workflows/_agent.py @@ -4,7 +4,7 @@ import logging import uuid -from collections.abc import AsyncIterable, Awaitable, Callable, Mapping, Sequence +from collections.abc import AsyncIterable, Awaitable, Callable, Iterable, Mapping, Sequence from dataclasses import dataclass from datetime import datetime, timezone from typing import ( @@ -54,22 +54,48 @@ logger = logging.getLogger(__name__) -def _is_orphaned_function_call(message: Message) -> bool: - """Drop a call whose result cannot follow it. +_INTERNAL_CALL_CONTENT_TYPES = frozenset({"function_call", "function_approval_response"}) - True when an assistant message is only a function call whose tool result - is not assistant-role, so keeping the call alone would produce an invalid - (unpaired) transcript for providers that validate call/result history. + +def _contains_internal_call_content(contents: Sequence[Content] | None) -> bool: + """True when any content is a function-call envelope (a call or an approval response).""" + return any(getattr(content, "type", None) in _INTERNAL_CALL_CONTENT_TYPES for content in contents or ()) + + +def _caller_facing_messages(messages: Iterable[Message]) -> list[Message]: + """Return the messages a workflow-as-agent caller should see. + + Non-assistant messages are dropped: system prompts and tool results are internal + workflow artifacts, and user messages would be re-emitted (e.g. from GroupChat + orchestrators that yield full conversation history). Assistant messages carrying + function-call envelopes are dropped whole, reasoning included: their tool-role + results are excluded by the role policy, and a call without its result is an + invalid transcript for providers that validate call/result pairing on replay. + """ + return [ + message + for message in messages + if message.role == "assistant" and not _contains_internal_call_content(message.contents) + ] + + +class _StreamingRoleGate: + """Track the last declared role across a stream of response updates. + + Non-assistant streams may declare ``role`` only on their first delta and use + ``role=None`` continuation chunks. Updates without a declared role inherit the + last declared one, so user content streamed as role-less chunks cannot leak + through the assistant-only filter. """ - if message.role != "assistant": - return False - contents = list(getattr(message, "contents", []) or []) - if not contents: - return False - # Orphaned only when every content is a function-call envelope - return all( - getattr(content, "type", None) in ("function_call", "function_approval_response") for content in contents - ) + + def __init__(self) -> None: + self._last_declared_role: str | None = None + + def should_forward(self, role: str | None) -> bool: + if role is not None: + self._last_declared_role = role + effective_role = role if role is not None else self._last_declared_role + return effective_role is None or effective_role == "assistant" class WorkflowAgent(BaseAgent): @@ -434,6 +460,7 @@ async def _run_stream_impl( session_messages: list[Message] = session_context.get_messages(include_input=True) all_updates: list[AgentResponseUpdate] = [] + role_gate = _StreamingRoleGate() async for event in self._run_core( session_messages, checkpoint_id, @@ -443,7 +470,7 @@ async def _run_stream_impl( function_invocation_kwargs=function_invocation_kwargs, client_kwargs=client_kwargs, ): - updates = self._convert_workflow_event_to_agent_response_updates(response_id, event) + updates = self._convert_workflow_event_to_agent_response_updates(response_id, event, role_gate) for update in updates: all_updates.append(update) yield update @@ -607,17 +634,9 @@ def _convert_workflow_events_to_agent_response( ) if isinstance(data, AgentResponse): - # Filter to only assistant messages — system, tool, and user messages - # are intentionally excluded. System prompts and tool results are - # internal workflow artifacts; user messages would be re-emitted - # (e.g., from GroupChat orchestrators that include full conversation history). - # Assistant messages that are bare function calls are also dropped - # when their tool result is not assistant-role: a call without its - # result is an invalid transcript for providers that validate - # call/result pairing on replay. - assistant_messages = [ - msg for msg in data.messages if msg.role == "assistant" and not _is_orphaned_function_call(msg) - ] + # Only caller-facing assistant messages survive; see + # _caller_facing_messages for the full policy. + assistant_messages = _caller_facing_messages(data.messages) if assistant_messages: messages.extend(assistant_messages) raw_representations.append(data.raw_representation) @@ -630,16 +649,14 @@ def _convert_workflow_events_to_agent_response( else latest_created_at ) elif isinstance(data, Message): - if data.role == "assistant": + if data.role == "assistant" and not _contains_internal_call_content(data.contents): messages.append(data) raw_representations.append(data.raw_representation) elif is_instance_of(data, list[Message]): chat_messages = cast(list[Message], data) - # Keep tool results that pair with surviving assistant calls so the - # transcript stays valid; drop user/system and orphaned calls. - assistant_messages = [ - msg for msg in chat_messages if msg.role == "assistant" and not _is_orphaned_function_call(msg) - ] + # Only caller-facing assistant messages survive; see + # _caller_facing_messages for the full policy. + assistant_messages = _caller_facing_messages(chat_messages) if assistant_messages: messages.extend(assistant_messages) # raw_representation of a filtered list must not leak the @@ -676,6 +693,7 @@ def _convert_workflow_event_to_agent_response_updates( self, response_id: str, event: WorkflowEvent[Any], + role_gate: _StreamingRoleGate, ) -> list[AgentResponseUpdate]: """Convert a workflow event to a list of AgentResponseUpdate objects. @@ -698,8 +716,11 @@ def _convert_workflow_event_to_agent_response_updates( executor_id = event.executor_id if isinstance(data, AgentResponseUpdate): - # Filter out non-assistant updates (e.g. user input echoed back) - if data.role is not None and data.role != "assistant": + # Filter out non-assistant updates (e.g. user input echoed back). + # Role-less continuation chunks inherit the last declared role from + # the gate, and updates carrying function-call envelopes are dropped + # so their orphaned calls cannot be persisted by from_updates. + if not role_gate.should_forward(data.role) or _contains_internal_call_content(data.contents): return [] # Construct a fresh AgentResponseUpdate so we don't mutate a payload # that AgentExecutor still holds a reference to in its `updates` list. @@ -723,11 +744,9 @@ def _convert_workflow_event_to_agent_response_updates( ) ] if isinstance(data, AgentResponse): - # Convert each assistant message in AgentResponse to an AgentResponseUpdate + # Convert each caller-facing assistant message to an AgentResponseUpdate updates: list[AgentResponseUpdate] = [] - for msg in data.messages: - if msg.role != "assistant": - continue + for msg in _caller_facing_messages(data.messages): updates.append( AgentResponseUpdate( contents=list(msg.contents), @@ -747,7 +766,7 @@ def _convert_workflow_event_to_agent_response_updates( updates[-1].additional_properties = dict(data.additional_properties) return updates if isinstance(data, Message): - if data.role != "assistant": + if data.role != "assistant" or _contains_internal_call_content(data.contents): return [] return [ AgentResponseUpdate( @@ -761,12 +780,10 @@ def _convert_workflow_event_to_agent_response_updates( ) ] if is_instance_of(data, list[Message]): - # Convert each assistant Message to an AgentResponseUpdate + # Convert each caller-facing assistant Message to an AgentResponseUpdate chat_messages = cast(list[Message], data) updates = [] - for msg in chat_messages: - if msg.role != "assistant": - continue + for msg in _caller_facing_messages(chat_messages): updates.append( AgentResponseUpdate( contents=list(msg.contents), diff --git a/python/packages/core/tests/workflow/test_workflow_agent.py b/python/packages/core/tests/workflow/test_workflow_agent.py index a04c572a21f..da31cbba7f8 100644 --- a/python/packages/core/tests/workflow/test_workflow_agent.py +++ b/python/packages/core/tests/workflow/test_workflow_agent.py @@ -3,9 +3,10 @@ import uuid from collections.abc import Awaitable, Sequence from dataclasses import dataclass -from typing import Any, Literal, Never, assert_type, cast, overload +from typing import Any, Literal, cast, overload import pytest +from typing_extensions import Never, assert_type from agent_framework import ( AgentExecutorRequest, @@ -2721,6 +2722,113 @@ async def tool_transcript_executor( ) assert any("18C" in (getattr(content, "text", "") or "") for msg in result.messages for content in msg.contents) + async def test_workflow_as_agent_drops_reasoning_call_group_atomically(self) -> None: + """An assistant message mixing reasoning with a function call is dropped whole. + + The call's tool-role result is excluded by the role policy, so keeping the + reasoning-bearing message would leave an orphaned call in the transcript + (moonbox3's review: remove the reasoning/call/result group atomically). + """ + + @executor + async def reasoning_call_executor( + messages: list[Message], + ctx: WorkflowContext[Never, list[Message]], # type: ignore[valid-type] + ) -> None: + await ctx.yield_output([ + Message( + role="assistant", + contents=[ + Content.from_text(text="The user wants the weather, I should call the tool."), + Content.from_function_call(call_id="call-2", name="get_weather", arguments={"city": "Rome"}), + ], + ), + Message( + role="tool", + contents=[Content.from_function_result(call_id="call-2", result="21C")], + ), + Message(role="assistant", contents=[Content.from_text("It is 21C in Rome.")]), + ]) + + workflow = WorkflowBuilder(start_executor=reasoning_call_executor).build() + agent = workflow.as_agent("reasoning-call-agent") + + # Streaming path + updates: list[AgentResponseUpdate] = [] + async for chunk in agent.run("weather", stream=True): + updates.append(chunk) + assert len(updates) == 1 + assert updates[0].text == "It is 21C in Rome." + + # Non-streaming path: neither the call nor its reasoning survives, only the answer. + result = await agent.run("weather") + assert len(result.messages) == 1 + assert result.messages[0].text == "It is 21C in Rome." + + async def test_workflow_as_agent_stream_drops_function_call_updates(self) -> None: + """Streamed updates carrying function-call envelopes are not forwarded. + + Forwarding one would let AgentResponse.from_updates persist an orphaned + call that later session turns replay to providers as invalid history. + """ + + @executor + async def call_update_executor( + messages: list[Message], + ctx: WorkflowContext[Never, AgentResponseUpdate], # type: ignore[valid-type] + ) -> None: + await ctx.yield_output( + AgentResponseUpdate( + contents=[ + Content.from_function_call(call_id="call-3", name="get_weather", arguments={"city": "Oslo"}), + ], + role="assistant", + ) + ) + await ctx.yield_output( + AgentResponseUpdate(contents=[Content.from_text("It is 3C in Oslo.")], role="assistant") + ) + + workflow = WorkflowBuilder(start_executor=call_update_executor).build() + agent = workflow.as_agent("call-update-agent") + + updates: list[AgentResponseUpdate] = [] + async for chunk in agent.run("weather", stream=True): + updates.append(chunk) + + assert len(updates) == 1 + assert updates[0].text == "It is 3C in Oslo." + + async def test_workflow_as_agent_stream_roleless_user_continuation_not_leaked(self) -> None: + """Role-less continuation chunks of a user stream inherit the user role. + + A stream may declare role only on its first delta; without tracking the + declared role across updates, every continuation chunk would be forwarded + as caller-visible output (moonbox3's review). + """ + + @executor + async def user_stream_executor( + messages: list[Message], + ctx: WorkflowContext[Never, AgentResponseUpdate], # type: ignore[valid-type] + ) -> None: + await ctx.yield_output( + AgentResponseUpdate(contents=[Content.from_text(text="user said: hello")], role="user") + ) + await ctx.yield_output(AgentResponseUpdate(contents=[Content.from_text(text=" and this too")], role=None)) + await ctx.yield_output(AgentResponseUpdate(contents=[Content.from_text("the answer")], role="assistant")) + + workflow = WorkflowBuilder(start_executor=user_stream_executor).build() + agent = workflow.as_agent("user-stream-agent") + + updates: list[AgentResponseUpdate] = [] + async for chunk in agent.run("test", stream=True): + updates.append(chunk) + + assert len(updates) == 1 + assert updates[0].role == "assistant" + assert updates[0].text == "the answer" + async def test_workflow_as_agent_filters_single_non_assistant_message(self) -> None: """Verify WorkflowAgent filters a single Message when role is not assistant.""" diff --git a/python/packages/foundry_hosting/agent_framework_foundry_hosting/_responses.py b/python/packages/foundry_hosting/agent_framework_foundry_hosting/_responses.py index 325bf264cd1..22e462edfdf 100644 --- a/python/packages/foundry_hosting/agent_framework_foundry_hosting/_responses.py +++ b/python/packages/foundry_hosting/agent_framework_foundry_hosting/_responses.py @@ -42,6 +42,7 @@ add_usage_details, ) from agent_framework._telemetry import mark_feature_used +from agent_framework._workflows._agent import _StreamingRoleGate from agent_framework.exceptions import AgentFrameworkException from azure.ai.agentserver.core import get_request_context from azure.ai.agentserver.responses import ( @@ -1278,13 +1279,14 @@ async def _resume_workflow_from_checkpoint( TODO(@taochen): #7677 """ + role_gate = _StreamingRoleGate() async for event in agent.workflow.run( stream=True, checkpoint_id=checkpoint_id, checkpoint_storage=checkpoint_storage, ): for update in agent._convert_workflow_event_to_agent_response_updates( # pyright: ignore[reportPrivateUsage] - response_id, event + response_id, event, role_gate ): yield update From c660761cd7b2d097a6a755ed300d5b66a601604b Mon Sep 17 00:00:00 2001 From: Manohar Paturi Date: Fri, 2 Oct 2026 15:46:27 +0530 Subject: [PATCH 6/6] fix(foundry-hosting): silence pyright private-usage on gate import Same reportPrivateUsage suppression the existing call site already carries. --- .../agent_framework_foundry_hosting/_responses.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/python/packages/foundry_hosting/agent_framework_foundry_hosting/_responses.py b/python/packages/foundry_hosting/agent_framework_foundry_hosting/_responses.py index 22e462edfdf..3f1102c2472 100644 --- a/python/packages/foundry_hosting/agent_framework_foundry_hosting/_responses.py +++ b/python/packages/foundry_hosting/agent_framework_foundry_hosting/_responses.py @@ -42,7 +42,7 @@ add_usage_details, ) from agent_framework._telemetry import mark_feature_used -from agent_framework._workflows._agent import _StreamingRoleGate +from agent_framework._workflows._agent import _StreamingRoleGate # pyright: ignore[reportPrivateUsage] from agent_framework.exceptions import AgentFrameworkException from azure.ai.agentserver.core import get_request_context from azure.ai.agentserver.responses import (