Skip to content
Draft
Show file tree
Hide file tree
Changes from 3 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
2 changes: 2 additions & 0 deletions python/packages/foundry/agent_framework_foundry/_agent.py
Original file line number Diff line number Diff line change
Expand Up @@ -439,6 +439,7 @@ def _parse_chunk_from_openai(
options: dict[str, Any],
function_call_ids: dict[int, tuple[str, str]],
seen_reasoning_delta_item_ids: set[str] | None = None,
seen_function_call_output_ids: set[str] | None = None,
) -> ChatResponseUpdate:
"""Parse streaming events while preserving hosted-agent session state."""
update = try_parse_oauth_consent_event(event, self.model)
Expand All @@ -448,6 +449,7 @@ def _parse_chunk_from_openai(
options,
function_call_ids,
seen_reasoning_delta_item_ids,
seen_function_call_output_ids,
)
if agent_session_id := _extract_foundry_hosted_agent_session_id(getattr(event, "response", None)):
if update.additional_properties is None:
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -292,12 +292,19 @@ def _parse_chunk_from_openai(
options: dict[str, Any],
function_call_ids: dict[int, tuple[str, str]],
seen_reasoning_delta_item_ids: set[str] | None = None,
seen_function_call_output_ids: set[str] | None = None,
) -> ChatResponseUpdate:
"""Parse streaming event, intercepting oauth_consent_request items."""
update = try_parse_oauth_consent_event(event, self.model)
if update is not None:
return update
return super()._parse_chunk_from_openai(event, options, function_call_ids, seen_reasoning_delta_item_ids)
return super()._parse_chunk_from_openai(
event,
options,
function_call_ids,
seen_reasoning_delta_item_ids,
seen_function_call_output_ids,
)

async def configure_azure_monitor(
self,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1846,7 +1846,7 @@ def test_parse_chunk_delegates_non_oauth_events_to_super() -> None:
return_value=MagicMock(),
) as mock_super:
client._parse_chunk_from_openai(mock_event, {}, {})
mock_super.assert_called_once_with(mock_event, {}, {}, None)
mock_super.assert_called_once_with(mock_event, {}, {}, None, None)


def test_parse_chunk_surfaces_oauth_consent_requested_event() -> None:
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1589,7 +1589,7 @@ def test_parse_chunk_delegates_non_oauth_events_to_super() -> None:
return_value=MagicMock(),
) as mock_super:
client._parse_chunk_from_openai(mock_event, {}, {})
mock_super.assert_called_once_with(mock_event, {}, {}, None)
mock_super.assert_called_once_with(mock_event, {}, {}, None, None)


def test_parse_chunk_surfaces_oauth_consent_requested_event() -> None:
Expand Down
89 changes: 89 additions & 0 deletions python/packages/openai/agent_framework_openai/_chat_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -76,6 +76,7 @@
from openai.types.responses import (
FunctionShellToolParam,
ResponseCustomToolCall,
ResponseFunctionToolCallOutputItem,
ResponseToolSearchCall,
response_create_params,
)
Expand Down Expand Up @@ -344,6 +345,42 @@ async def _open_event_stream(raw_response: Any) -> AsyncGenerator[Any]:
yield raw_response


def _function_call_output_has_result(item: Any) -> bool:
"""True when a ``function_call_output`` item carries a result that can be paired.

``.added`` can precede a populated ``output``, so an in-progress skeleton must not synthesize
an empty result. A blank ``call_id`` is rejected as well: these items are synthesized by the
hosting layer, and a result that cannot be paired to its call is dropped by transports and,
worse, re-sent as an unpairable ``function_call_output`` input item on the next turn.
"""
if getattr(item, "output", None) is None:
return False
if not getattr(item, "call_id", None):
logger.debug("Skipping function_call_output with no call_id: item_id=%s", getattr(item, "id", None))
return False
return True


def _claim_function_call_output(seen_item_ids: set[str] | None, item: Any) -> bool:
"""Claim a ``function_call_output`` item for emission; return ``False`` if already claimed.

The Responses stream can surface the same output item on both ``response.output_item.added``
and ``response.output_item.done``. Whichever event first carries a populated ``output`` emits
the result, and this test-and-set keeps the other from producing a duplicate one. Keyed on the
item id rather than ``call_id``, which is not guaranteed to be unique forever. When no set is
supplied the item is always claimable, so a single event parsed on its own still yields output.
"""
if seen_item_ids is None:
return True
item_id = getattr(item, "id", None)
if not isinstance(item_id, str) or not item_id:
return True
if item_id in seen_item_ids:
return False
seen_item_ids.add(item_id)
return True


def _annotations_to_output_text(annotations: Sequence[Annotation] | None) -> list[dict[str, Any]]:
"""Convert framework `Annotation` objects to Responses API `output_text` annotation dicts.

Expand Down Expand Up @@ -712,6 +749,7 @@ def _inner_get_response(
if stream:
function_call_ids: dict[int, tuple[str, str]] = {}
seen_reasoning_delta_item_ids: set[str] = set()
seen_function_call_output_ids: set[str] = set()
validated_options: dict[str, Any] | None = None
# Captured once request options are validated/prepared so the streaming finalizer can
# still parse the aggregated response into structured output after the stream completes.
Expand Down Expand Up @@ -750,6 +788,7 @@ async def _stream() -> AsyncIterable[ChatResponseUpdate]:
options=validated_options,
function_call_ids=function_call_ids,
seen_reasoning_delta_item_ids=seen_reasoning_delta_item_ids,
seen_function_call_output_ids=seen_function_call_output_ids,
)
if served_model is not None:
update.model = served_model
Expand Down Expand Up @@ -780,6 +819,7 @@ async def _stream() -> AsyncIterable[ChatResponseUpdate]:
options=validated_options,
function_call_ids=function_call_ids,
seen_reasoning_delta_item_ids=seen_reasoning_delta_item_ids,
seen_function_call_output_ids=seen_function_call_output_ids,
)
else:
raw_create_response = await client.responses.with_raw_response.create(
Expand All @@ -794,6 +834,7 @@ async def _stream() -> AsyncIterable[ChatResponseUpdate]:
options=validated_options,
function_call_ids=function_call_ids,
seen_reasoning_delta_item_ids=seen_reasoning_delta_item_ids,
seen_function_call_output_ids=seen_function_call_output_ids,
)
if served_model is not None:
update.model = served_model
Expand Down Expand Up @@ -2628,6 +2669,30 @@ def _parse_hosted_function_call_content(
raw_representation=item,
)

def _parse_function_call_output_content(self, item: ResponseFunctionToolCallOutputItem) -> Content:
"""Create function result content for a Responses ``function_call_output`` item.

A hosted tool that executes server-side -- for example a Foundry Toolbox dispatching
through its generic ``call_tool`` wrapper -- returns its result as a standalone
``function_call_output`` item rather than on the originating call item. Parsing it keeps
the call/result pair intact for transports such as AG-UI, which otherwise sees a tool call
with no result and falls back to treating it as declaration-only (issue #8068).

``output`` is either a string or a list of input-content parts, so it is normalized through
:meth:`_stringify_mcp_output` rather than JSON-encoding provider models.
"""
additional_properties: dict[str, Any] = {"item_type": item.type, "status": item.status}
if item.id:
additional_properties["item_id"] = item.id
if item.name:
additional_properties["name"] = item.name
Comment thread
manjunathshiva marked this conversation as resolved.
Outdated
return Content.from_function_result(
call_id=item.call_id,
result=self._stringify_mcp_output(item.output),
additional_properties=additional_properties,
raw_representation=item,
)

# region Parse methods
def _get_finish_reason_from_openai_response(self, response: Any) -> FinishReason | None:
"""Get the framework finish reason from a terminal Responses API response."""
Expand Down Expand Up @@ -2855,6 +2920,9 @@ def _parse_response_from_openai(
raw_representation=item,
)
)
case "function_call_output": # ResponseFunctionToolCallOutputItem
if _function_call_output_has_result(item):
contents.append(self._parse_function_call_output_content(item))
case "custom_tool_call":
contents.append(
self._parse_hosted_function_call_content(item, name=item.name, arguments=item.input)
Expand Down Expand Up @@ -2966,6 +3034,7 @@ def _parse_chunk_from_openai(
options: dict[str, Any],
function_call_ids: dict[int, tuple[str, str]],
seen_reasoning_delta_item_ids: set[str] | None = None,
seen_function_call_output_ids: set[str] | None = None,
) -> ChatResponseUpdate:
"""Parse an OpenAI Responses API streaming event into a ChatResponseUpdate."""
metadata: dict[str, Any] = {}
Expand Down Expand Up @@ -3330,6 +3399,14 @@ def output_text_properties(output: Any) -> dict[str, Any] | None:
)
case "web_search_call" | "file_search_call":
contents.append(self._parse_search_tool_call_content(event_item))
case "function_call_output": # ResponseFunctionToolCallOutputItem
# Emitted from whichever of `.added` / `.done` first carries a populated
# `output`; the item id is recorded so the other event cannot emit a second
# result for the same item (issue #8068).
if _function_call_output_has_result(event_item) and _claim_function_call_output(
seen_function_call_output_ids, event_item
):
contents.append(self._parse_function_call_output_content(event_item))
case _:
if getattr(event_item, "type", None) != _AZURE_AI_SEARCH_CALL_OUTPUT_TYPE:
logger.debug("Unparsed event of type: %s: %s", event.type, event)
Expand Down Expand Up @@ -3534,6 +3611,18 @@ def _get_ann_value(key: str) -> Any:
arguments=tool_search_call.arguments,
)
)
elif getattr(done_item, "type", None) == "function_call_output":
# Counterpart to the `response.output_item.added` branch: whichever event first
# carries a populated `output` emits the result, and the shared seen-id set
# keeps the other from duplicating it (issue #8068).
if _function_call_output_has_result(done_item) and _claim_function_call_output(
seen_function_call_output_ids, done_item
):
contents.append(
self._parse_function_call_output_content(
cast(ResponseFunctionToolCallOutputItem, done_item)
)
)
elif getattr(done_item, "type", None) == _AZURE_AI_SEARCH_CALL_OUTPUT_TYPE:
pass
case _:
Expand Down
Loading
Loading