Skip to content
Open
Show file tree
Hide file tree
Changes from all 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
5 changes: 5 additions & 0 deletions .sampo/changesets/responses-incomplete-stream-usage.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
pypi/posthog: patch
---

Capture token usage and output for OpenAI Responses streams that end incomplete, such as when `max_output_tokens` is reached.
20 changes: 16 additions & 4 deletions posthog/ai/openai/openai_converter.py
Original file line number Diff line number Diff line change
Expand Up @@ -520,6 +520,15 @@ def extract_openai_usage_from_response(response: Any) -> TokenUsage:
return result


# Stream events that carry the final Responses API response, with its usage and
# output. Exactly one of them ends a stream.
_RESPONSES_TERMINAL_EVENT_TYPES = (
"response.completed",
"response.incomplete",
"response.failed",
)


def extract_openai_usage_from_chunk(
chunk: Any, provider_type: str = "chat"
) -> TokenUsage:
Expand Down Expand Up @@ -577,8 +586,10 @@ def extract_openai_usage_from_chunk(
usage["raw_usage"] = serialized

elif provider_type == "responses":
# For Responses API, usage is only in chunk.response.usage for completed events
if hasattr(chunk, "type") and chunk.type == "response.completed":
# For Responses API, usage is only in chunk.response.usage on the terminal
# event. A run cut short (e.g. by max_output_tokens) ends on
# response.incomplete instead of response.completed, and is still billed.
if getattr(chunk, "type", None) in _RESPONSES_TERMINAL_EVENT_TYPES:
if (
hasattr(chunk, "response")
and hasattr(chunk.response, "usage")
Expand Down Expand Up @@ -634,7 +645,8 @@ def extract_openai_content_from_chunk(
Returns:
For "chat": text content (str), or an audio/refusal delta block (dict),
if present. For "responses": the full `response.output` list on the
`response.completed` event. None otherwise.
terminal (`response.completed`, `response.incomplete` or
`response.failed`) event. None otherwise.
"""

if provider_type == "chat":
Expand Down Expand Up @@ -663,7 +675,7 @@ def extract_openai_content_from_chunk(

elif provider_type == "responses":
# Responses API format
if hasattr(chunk, "type") and chunk.type == "response.completed":
if getattr(chunk, "type", None) in _RESPONSES_TERMINAL_EVENT_TYPES:
if hasattr(chunk, "response") and chunk.response:
res = chunk.response
if res.output:
Expand Down
79 changes: 79 additions & 0 deletions posthog/test/ai/openai/test_openai.py
Original file line number Diff line number Diff line change
Expand Up @@ -1916,6 +1916,85 @@ def test_streaming_responses_api_extracts_model_from_response_object(mock_client
assert props["$ai_model"] == "gpt-4o-mini-stored"


def test_streaming_responses_api_captures_usage_and_output_when_incomplete(
mock_client,
):
"""A stream cut short by max_output_tokens ends on response.incomplete, not
response.completed. Its usage and partial output are still billed and must be
captured."""
from openai.types.responses import ResponseIncompleteEvent
from openai.types.responses.response import IncompleteDetails

incomplete_response = Response(
id="resp_incomplete",
model="gpt-4o-mini",
object="response",
created_at=1741476542,
status="incomplete",
error=None,
incomplete_details=IncompleteDetails(reason="max_output_tokens"),
instructions=None,
max_output_tokens=16,
tools=[],
tool_choice="auto",
output=[
ResponseOutputMessage(
id="msg_123",
type="message",
role="assistant",
status="incomplete",
content=[
ResponseOutputText(
type="output_text",
text="Once upon a time",
annotations=[],
)
],
)
],
parallel_tool_calls=True,
previous_response_id=None,
usage=make_response_usage(
input_tokens=20,
output_tokens=16,
total_tokens=36,
),
user=None,
metadata={},
)
chunks = [
ResponseIncompleteEvent(
type="response.incomplete",
sequence_number=1,
response=incomplete_response,
)
]

with patch("openai.resources.responses.Responses.create") as mock_create:
mock_create.return_value = iter(chunks)

client = OpenAI(api_key="test-key", posthog_client=mock_client)
response_generator = client.responses.create(
model="gpt-4o-mini",
input=[{"role": "user", "content": "Tell me a story"}],
max_output_tokens=16,
stream=True,
posthog_distinct_id="test-id",
)
list(response_generator)

props = mock_client.capture.call_args[1]["properties"]
assert props["$ai_stop_reason"] == "max_output_tokens"
assert props["$ai_input_tokens"] == 20
assert props["$ai_output_tokens"] == 16
assert props["$ai_output_choices"] == [
{
"role": "assistant",
"content": [{"type": "text", "text": "Once upon a time"}],
}
]


def test_non_streaming_extracts_model_from_response(mock_client):
"""Test that non-streaming calls extract model from response when not in kwargs."""
# Create a response with model but we won't pass model in kwargs
Expand Down