diff --git a/.sampo/changesets/responses-incomplete-stream-usage.md b/.sampo/changesets/responses-incomplete-stream-usage.md new file mode 100644 index 000000000..4e9f5e71b --- /dev/null +++ b/.sampo/changesets/responses-incomplete-stream-usage.md @@ -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. diff --git a/posthog/ai/openai/openai_converter.py b/posthog/ai/openai/openai_converter.py index 3f3bc869c..0f391fb9f 100644 --- a/posthog/ai/openai/openai_converter.py +++ b/posthog/ai/openai/openai_converter.py @@ -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: @@ -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") @@ -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": @@ -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: diff --git a/posthog/test/ai/openai/test_openai.py b/posthog/test/ai/openai/test_openai.py index 1aee497dd..0d22f393a 100644 --- a/posthog/test/ai/openai/test_openai.py +++ b/posthog/test/ai/openai/test_openai.py @@ -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