Stream native response tool output
Forward incremental response function arguments through the canonical Agent Zero tool envelope so final Codex answers render progressively while preserving buffered execution for non-response tool calls. Cover partial DirtyJson rendering, completed envelope assembly, duplicate completion suppression, and the transport contract.
Alessandro committed
Aug 27, 2026 at 07:27 UTC
117ae265f9bb78e9d230c6832ec4aee3395e1236
3 files changed
+91
-6
helpers/litellm_transport.py
+26
-6
@@ -1166,6 +1166,7 @@ class ResponsesEventParser:
1166
self.function_calls: dict[str, dict[str, Any]] = {}
1167
self.output_index_keys: dict[str, str] = {}
1168
self.emitted_function_calls: set[str] = set()
1169
+ self.streamed_response_calls: dict[str, str] = {}
1170
self.seen_response_delta = False
1171
self.seen_reasoning_delta = False
1172
self.completed_response: Any = None
@@ -1189,7 +1190,7 @@ class ResponsesEventParser:
1190
elif event_type == "response.output_item.added":
1191
self._remember_function_call(_get_value(event, "item"), event)
1192
elif event_type == "response.function_call_arguments.delta":
1192
- self._append_function_call_arguments(event)
1193
+ response_delta = self._append_function_call_arguments(event)
1194
elif event_type == "response.function_call_arguments.done":
1195
response_delta = self._complete_function_call(event)
1196
elif event_type == "response.output_item.done":
@@ -1224,14 +1225,22 @@ class ResponsesEventParser:
1225
self.output_index_keys[str(output_index)] = key
1226
return key
1227
1227
- def _append_function_call_arguments(self, event: Any) -> None:
1228
+ def _append_function_call_arguments(self, event: Any) -> str:
1229
key = self._event_key(event)
1230
if not key:
1230
- return
1231
+ return ""
1232
current = self.function_calls.setdefault(key, {"type": "function_call"})
1232
- current["arguments"] = str(current.get("arguments") or "") + str(
1233
- _get_value(event, "delta") or ""
1234
- )
1233
+ delta = str(_get_value(event, "delta") or "")
1234
+ current["arguments"] = str(current.get("arguments") or "") + delta
1235
+ if current.get("name") != "response":
1236
+ return ""
1237
+ if key not in self.streamed_response_calls:
1238
+ self.streamed_response_calls[key] = str(current["arguments"])
1239
+ return '{"tool_name":"response","tool_args":' + str(
1240
+ current["arguments"]
1241
+ )
1242
+ self.streamed_response_calls[key] += delta
1243
+ return delta
1244
1245
def _complete_function_call(self, event: Any) -> str:
1246
key = self._event_key(event)
@@ -1242,14 +1251,25 @@ class ResponsesEventParser:
1251
current["arguments"] = _get_value(event, "arguments")
1252
if _get_value(event, "name"):
1253
current["name"] = _get_value(event, "name")
1254
+ if key in self.streamed_response_calls:
1255
+ return self._finish_response_call(key, current)
1256
return self._emit_function_call(key, current)
1257
1258
def _complete_output_item(self, item: Any, event: Any) -> str:
1259
key = self._remember_function_call(item, event)
1260
if not key:
1261
return ""
1262
+ if key in self.streamed_response_calls:
1263
+ return self._finish_response_call(key, self.function_calls[key])
1264
return self._emit_function_call(key, self.function_calls[key])
1265
1266
+ def _finish_response_call(self, key: str, item: Any) -> str:
1267
+ streamed = self.streamed_response_calls.pop(key)
1268
+ arguments = str(_get_value(item, "arguments") or "")
1269
+ self.emitted_function_calls.add(key)
1270
+ tail = arguments[len(streamed) :] if arguments.startswith(streamed) else ""
1271
+ return tail + "}"
1272
+
1273
def _complete_response(self, event: Any) -> tuple[str, str]:
1274
self.completed_response = _get_value(event, "response")
1275
if self.seen_response_delta or self.emitted_function_calls:
helpers/litellm_transport.py.dox.md
+1
@@ -35,6 +35,7 @@
35
- Preserve Chat Completions tool calls from both non-streaming responses and streaming deltas as canonical `LLMResult` function-call items.
36
- Preserve provider usage and LiteLLM response cost for both transports only when the response or stream actually supplies them; do not synthesize unavailable provider accounting.
37
- Preserve Responses function calls collected from stream events when a terminal completed event omits them.
38
+- Stream native `response` function arguments through a canonical response-tool envelope while continuing to buffer other function calls until completion.
39
- Serialize synthesized Responses function-call JSON with literal Unicode so streamed raw-response logs preserve tool arguments.
40
- Preserve provider-state metadata when Responses API calls succeed, and fall back to local replay when provider state is unsupported.
41
- Keep prompt-cache markers only for providers that accept them.
tests/test_stream_tool_early_stop.py
+64
@@ -13,6 +13,7 @@ if str(PROJECT_ROOT) not in sys.path:
13
import models
14
from helpers import extract_tools
15
from helpers import litellm_transport
16
+from helpers.dirty_json import DirtyJson
17
18
19
@pytest.fixture(autouse=True)
@@ -1704,6 +1705,69 @@ def test_responses_stream_parser_accumulates_function_call_arguments():
1705
) == {"reasoning_delta": "", "response_delta": ""}
1706
1707
1708
+def test_responses_stream_parser_streams_response_function_arguments():
1709
+ parser = litellm_transport.ResponsesEventParser()
1710
+
1711
+ parser.parse(
1712
+ {
1713
+ "type": "response.output_item.added",
1714
+ "output_index": 0,
1715
+ "item": {
1716
+ "type": "function_call",
1717
+ "id": "fc_1",
1718
+ "name": "response",
1719
+ "arguments": "",
1720
+ },
1721
+ }
1722
+ )
1723
+ chunks = [
1724
+ parser.parse(
1725
+ {
1726
+ "type": "response.function_call_arguments.delta",
1727
+ "item_id": "fc_1",
1728
+ "delta": '{"text":"Hello',
1729
+ }
1730
+ )["response_delta"],
1731
+ parser.parse(
1732
+ {
1733
+ "type": "response.function_call_arguments.delta",
1734
+ "item_id": "fc_1",
1735
+ "delta": ' world"}',
1736
+ }
1737
+ )["response_delta"],
1738
+ parser.parse(
1739
+ {
1740
+ "type": "response.function_call_arguments.done",
1741
+ "item_id": "fc_1",
1742
+ "name": "response",
1743
+ "arguments": '{"text":"Hello world"}',
1744
+ }
1745
+ )["response_delta"],
1746
+ ]
1747
+
1748
+ assert chunks[0] == '{"tool_name":"response","tool_args":{"text":"Hello'
1749
+ assert DirtyJson.parse_string(chunks[0]) == {
1750
+ "tool_name": "response",
1751
+ "tool_args": {"text": "Hello"},
1752
+ }
1753
+ assert extract_tools.json_parse_dirty("".join(chunks)) == {
1754
+ "tool_name": "response",
1755
+ "tool_args": {"text": "Hello world"},
1756
+ }
1757
+ assert chunks[-1] == "}"
1758
+ assert parser.parse(
1759
+ {
1760
+ "type": "response.output_item.done",
1761
+ "item": {
1762
+ "type": "function_call",
1763
+ "id": "fc_1",
1764
+ "name": "response",
1765
+ "arguments": '{"text":"Hello world"}',
1766
+ },
1767
+ }
1768
+ ) == {"reasoning_delta": "", "response_delta": ""}
1769
+
1770
+
1771
def test_responses_stream_parser_uses_completed_response_when_no_deltas_arrive():
1772
parser = litellm_transport.ResponsesEventParser()
1773