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/omit-unreported-token-counts.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
pypi/posthog: patch
---

Omit `$ai_input_tokens` and `$ai_output_tokens` when the provider never reported usage, instead of sending `0`, so an interrupted stream no longer looks like a free call. A zero reported by the provider is still sent, and zero keeps meaning a real report of nothing. Covers the OpenAI, Anthropic, Gemini, LangChain, OpenAI Agents and Claude Agent SDK integrations.
2 changes: 1 addition & 1 deletion posthog/ai/anthropic/_anthropic_stream.py
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,7 @@ class _AnthropicStreamAccumulator:
"""Accumulates sync-neutral capture state from Anthropic stream events."""

def __init__(self) -> None:
self.usage_stats: TokenUsage = TokenUsage(input_tokens=0, output_tokens=0)
self.usage_stats: TokenUsage = TokenUsage()
self.accumulated_content = ""
self.content_blocks: List[StreamingContentBlock] = []
self.tools_in_progress: Dict[str, ToolInProgress] = {}
Expand Down
17 changes: 10 additions & 7 deletions posthog/ai/anthropic/anthropic_converter.py
Original file line number Diff line number Diff line change
Expand Up @@ -252,13 +252,16 @@ def extract_anthropic_usage_from_response(response: Any) -> TokenUsage:
Returns:
TokenUsage with standardized usage
"""
if not hasattr(response, "usage"):
return TokenUsage(input_tokens=0, output_tokens=0)

result = TokenUsage(
input_tokens=getattr(response.usage, "input_tokens", 0),
output_tokens=getattr(response.usage, "output_tokens", 0),
)
if getattr(response, "usage", None) is None:
return TokenUsage()

result = TokenUsage()
input_tokens = getattr(response.usage, "input_tokens", None)
if input_tokens is not None:
result["input_tokens"] = input_tokens
output_tokens = getattr(response.usage, "output_tokens", None)
if output_tokens is not None:
result["output_tokens"] = output_tokens

if hasattr(response.usage, "cache_read_input_tokens"):
cache_read = response.usage.cache_read_input_tokens
Expand Down
48 changes: 32 additions & 16 deletions posthog/ai/claude_agent_sdk/processor.py
Original file line number Diff line number Diff line change
Expand Up @@ -45,10 +45,12 @@ class _GenerationData:
"""Data accumulated for a single LLM generation (one API call)."""

model: Optional[str] = None
input_tokens: int = 0
output_tokens: int = 0
cache_read_input_tokens: int = 0
cache_creation_input_tokens: int = 0
# None when the provider never reported a count: absent means unknown,
# 0 is a report of nothing.
input_tokens: Optional[int] = None
output_tokens: Optional[int] = None
cache_read_input_tokens: Optional[int] = None
cache_creation_input_tokens: Optional[int] = None
raw_usage: Optional[Dict[str, Any]] = None
start_time: float = 0.0
end_time: float = 0.0
Expand Down Expand Up @@ -79,13 +81,11 @@ def process_stream_event(self, event: "StreamEvent") -> None:
message = raw.get("message", {})
self._current.model = message.get("model")
usage = message.get("usage", {})
self._current.input_tokens = usage.get("input_tokens", 0)
self._current.output_tokens = usage.get("output_tokens", 0)
self._current.cache_read_input_tokens = usage.get(
"cache_read_input_tokens", 0
)
self._current.input_tokens = usage.get("input_tokens")
self._current.output_tokens = usage.get("output_tokens")
self._current.cache_read_input_tokens = usage.get("cache_read_input_tokens")
self._current.cache_creation_input_tokens = usage.get(
"cache_creation_input_tokens", 0
"cache_creation_input_tokens"
)
self._current.raw_usage = dict(usage)

Expand Down Expand Up @@ -410,8 +410,16 @@ def _emit_generation(
"$ai_provider": "anthropic",
"$ai_framework": "claude-agent-sdk",
"$ai_model": gen.model,
"$ai_input_tokens": gen.input_tokens,
"$ai_output_tokens": gen.output_tokens,
**(
{"$ai_input_tokens": gen.input_tokens}
if gen.input_tokens is not None
else {}
),
**(
{"$ai_output_tokens": gen.output_tokens}
if gen.output_tokens is not None
else {}
),
"$ai_latency": latency,
**extra_props,
}
Expand Down Expand Up @@ -472,8 +480,16 @@ def _emit_generation_from_result(
"$ai_provider": "anthropic",
"$ai_framework": "claude-agent-sdk",
"$ai_model": model,
"$ai_input_tokens": usage.get("input_tokens", 0),
"$ai_output_tokens": usage.get("output_tokens", 0),
**(
{"$ai_input_tokens": usage["input_tokens"]}
if usage.get("input_tokens") is not None
else {}
),
**(
{"$ai_output_tokens": usage["output_tokens"]}
if usage.get("output_tokens") is not None
else {}
),
"$ai_latency": result.duration_api_ms / 1000.0
if result.duration_api_ms
else 0,
Expand All @@ -494,8 +510,8 @@ def _emit_generation_from_result(
finalize_ai_content(output_choices, self._client),
)

cache_read = usage.get("cache_read_input_tokens", 0)
cache_creation = usage.get("cache_creation_input_tokens", 0)
cache_read = usage.get("cache_read_input_tokens")
cache_creation = usage.get("cache_creation_input_tokens")
if cache_read:
properties["$ai_cache_read_input_tokens"] = cache_read
if cache_creation:
Expand Down
8 changes: 6 additions & 2 deletions posthog/ai/gemini/_shared.py
Original file line number Diff line number Diff line change
Expand Up @@ -195,7 +195,9 @@ def _capture_embedding_outcome(
error: Optional[Exception],
latency: float,
) -> None:
input_tokens = extract_gemini_embedding_token_count(response) if response else 0
input_tokens = (
extract_gemini_embedding_token_count(response) if response else None
)
event_properties = {
"$ai_provider": "gemini",
"$ai_model": model,
Expand All @@ -207,7 +209,9 @@ def _capture_embedding_outcome(
"$ai_http_status": (
getattr(error, "status_code", 0) if error is not None else 200
),
"$ai_input_tokens": input_tokens,
# Omitted when the provider never reported a count: absent means
# unknown, 0 is a report of nothing.
**({"$ai_input_tokens": input_tokens} if input_tokens is not None else {}),
"$ai_latency": latency,
"$ai_trace_id": trace_id,
"$ai_base_url": self._base_url,
Expand Down
2 changes: 1 addition & 1 deletion posthog/ai/gemini/gemini.py
Original file line number Diff line number Diff line change
Expand Up @@ -216,7 +216,7 @@ def _generate_content_streaming(
**kwargs: Any,
):
start_time = time.time()
usage_stats: TokenUsage = TokenUsage(input_tokens=0, output_tokens=0)
usage_stats: TokenUsage = TokenUsage()
accumulated_content = []
stop_reason: Optional[str] = None

Expand Down
2 changes: 1 addition & 1 deletion posthog/ai/gemini/gemini_async.py
Original file line number Diff line number Diff line change
Expand Up @@ -217,7 +217,7 @@ async def _generate_content_streaming(
**kwargs: Any,
):
start_time = time.time()
usage_stats: TokenUsage = TokenUsage(input_tokens=0, output_tokens=0)
usage_stats: TokenUsage = TokenUsage()
accumulated_content = []
stop_reason: Optional[str] = None

Expand Down
10 changes: 6 additions & 4 deletions posthog/ai/gemini/gemini_converter.py
Original file line number Diff line number Diff line change
Expand Up @@ -546,7 +546,7 @@ def extract_gemini_usage_from_response(response: Any) -> TokenUsage:
TokenUsage with standardized usage statistics
"""
if not hasattr(response, "usage_metadata") or not response.usage_metadata:
return TokenUsage(input_tokens=0, output_tokens=0)
return TokenUsage()

usage = _extract_usage_from_metadata(response.usage_metadata)

Expand Down Expand Up @@ -715,17 +715,19 @@ def format_gemini_streaming_output(
return [{"role": "assistant", "content": [{"type": "text", "text": ""}]}]


def extract_gemini_embedding_token_count(response) -> int:
def extract_gemini_embedding_token_count(response) -> Optional[int]:
"""
Extract total token count from a Gemini embed_content response.
Token counts are only available per-embedding via Vertex AI's statistics.token_count.
Returns 0 if no token counts are available.
Returns None when no embedding carried a token count.
"""
total = 0
reported = False
if hasattr(response, "embeddings") and response.embeddings:
for embedding in response.embeddings:
if hasattr(embedding, "statistics") and embedding.statistics:
token_count = getattr(embedding.statistics, "token_count", None)
if token_count is not None:
total += int(token_count)
return total
reported = True
return total if reported else None
25 changes: 17 additions & 8 deletions posthog/ai/langchain/callbacks.py
Original file line number Diff line number Diff line change
Expand Up @@ -664,11 +664,16 @@ def _capture_generation(
else:
# Add usage
usage = _parse_usage(output, run.provider, run.model)
event_properties["$ai_input_tokens"] = usage.input_tokens
event_properties["$ai_output_tokens"] = usage.output_tokens
event_properties["$ai_cache_creation_input_tokens"] = (
usage.cache_write_tokens
)
# Omitted when the provider never reported a count: absent means
# unknown, 0 is a report of nothing.
if usage.input_tokens is not None:
event_properties["$ai_input_tokens"] = usage.input_tokens
if usage.output_tokens is not None:
event_properties["$ai_output_tokens"] = usage.output_tokens
if usage.cache_write_tokens is not None:
event_properties["$ai_cache_creation_input_tokens"] = (
usage.cache_write_tokens
)
if (
usage.cache_write_5m_tokens is not None
and usage.cache_write_1h_tokens is not None
Expand All @@ -679,8 +684,12 @@ def _capture_generation(
event_properties["$ai_cache_creation_1h_input_tokens"] = (
usage.cache_write_1h_tokens
)
event_properties["$ai_cache_read_input_tokens"] = usage.cache_read_tokens
event_properties["$ai_reasoning_tokens"] = usage.reasoning_tokens
if usage.cache_read_tokens is not None:
event_properties["$ai_cache_read_input_tokens"] = (
usage.cache_read_tokens
)
if usage.reasoning_tokens is not None:
event_properties["$ai_reasoning_tokens"] = usage.reasoning_tokens

# Generation results
generation_result = output.generations[-1]
Expand Down Expand Up @@ -875,7 +884,7 @@ def _parse_usage_model(
}
normalized_usage = ModelUsage(
**{
dataclass_key: parsed_usage.get(mapped_key) or 0
dataclass_key: parsed_usage.get(mapped_key)
for mapped_key, dataclass_key in field_mapping.items()
},
cache_write_5m_tokens=parsed_usage.get("cache_write_5m"),
Expand Down
6 changes: 4 additions & 2 deletions posthog/ai/openai/_embeddings.py
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,7 @@ def _capture_embedding_event(
) -> None:
"""Build and capture telemetry shared by sync and async embedding wrappers."""
usage = getattr(response, "usage", None)
input_tokens = getattr(usage, "prompt_tokens", 0) if usage else 0
input_tokens = getattr(usage, "prompt_tokens", None) if usage else None

event_properties = {
"$ai_provider": "openai",
Expand All @@ -29,7 +29,9 @@ def _capture_embedding_event(
finalize_ai_content(request_kwargs.get("input"), posthog_client),
),
"$ai_http_status": 200,
"$ai_input_tokens": input_tokens,
# Omitted when the provider never reported a count: absent means
# unknown, 0 is a report of nothing.
**({"$ai_input_tokens": input_tokens} if input_tokens is not None else {}),
"$ai_latency": latency,
"$ai_trace_id": trace_id,
"$ai_base_url": str(base_url),
Expand Down
21 changes: 11 additions & 10 deletions posthog/ai/openai/openai_converter.py
Original file line number Diff line number Diff line change
Expand Up @@ -460,13 +460,13 @@ def extract_openai_usage_from_response(response: Any) -> TokenUsage:
Returns:
TokenUsage with standardized usage statistics
"""
if not hasattr(response, "usage"):
return TokenUsage(input_tokens=0, output_tokens=0)
if not hasattr(response, "usage") or not response.usage:
return TokenUsage()

cached_tokens = 0
input_tokens = 0
output_tokens = 0
reasoning_tokens = 0
cached_tokens = None
input_tokens = None
output_tokens = None
reasoning_tokens = None

# Responses API format
if hasattr(response.usage, "input_tokens"):
Expand Down Expand Up @@ -496,10 +496,11 @@ def extract_openai_usage_from_response(response: Any) -> TokenUsage:
):
reasoning_tokens = response.usage.completion_tokens_details.reasoning_tokens

result = TokenUsage(
input_tokens=input_tokens,
output_tokens=output_tokens,
)
result = TokenUsage()
if input_tokens is not None:
result["input_tokens"] = input_tokens
if output_tokens is not None:
result["output_tokens"] = output_tokens

if cached_tokens is not None and cached_tokens > 0:
result["cache_read_input_tokens"] = cached_tokens
Expand Down
51 changes: 36 additions & 15 deletions posthog/ai/openai_agents/processor.py
Original file line number Diff line number Diff line change
Expand Up @@ -478,10 +478,14 @@ def _handle_generation_span(
"""Handle LLM generation spans - maps to $ai_generation event."""
# Extract token usage
usage = span_data.usage or {}
input_tokens = usage.get("input_tokens") or usage.get("prompt_tokens") or 0
output_tokens = (
usage.get("output_tokens") or usage.get("completion_tokens") or 0
)
# None when the span never reported a count: absent means unknown,
# 0 is a report of nothing.
input_tokens = usage.get("input_tokens")
if input_tokens is None:
input_tokens = usage.get("prompt_tokens")
output_tokens = usage.get("output_tokens")
if output_tokens is None:
output_tokens = usage.get("completion_tokens")

# Extract model config parameters
model_config = span_data.model_config or {}
Expand Down Expand Up @@ -510,9 +514,18 @@ def _handle_generation_span(
_ensure_serializable(span_data.output), self._client
)
),
"$ai_input_tokens": input_tokens,
"$ai_output_tokens": output_tokens,
"$ai_total_tokens": (input_tokens or 0) + (output_tokens or 0),
**({"$ai_input_tokens": input_tokens} if input_tokens is not None else {}),
**(
{"$ai_output_tokens": output_tokens}
if output_tokens is not None
else {}
),
# Sum of the reported sides; omitted when neither side was reported.
**(
{"$ai_total_tokens": (input_tokens or 0) + (output_tokens or 0)}
if input_tokens is not None or output_tokens is not None
else {}
),
}

# Add optional token fields if present
Expand Down Expand Up @@ -661,11 +674,10 @@ def _handle_response_span(
# Try to extract usage from response
usage = getattr(response, "usage", None) if response else None
total_cost_usd = getattr(usage, "cost", None) if usage else None
input_tokens = 0
output_tokens = 0
if usage:
input_tokens = getattr(usage, "input_tokens", 0) or 0
output_tokens = getattr(usage, "output_tokens", 0) or 0
# None when the response never reported a count: absent means unknown,
# 0 is a report of nothing.
input_tokens = getattr(usage, "input_tokens", None) if usage else None
output_tokens = getattr(usage, "output_tokens", None) if usage else None

# Try to extract model from response
model = getattr(response, "model", None) if response else None
Expand All @@ -679,9 +691,18 @@ def _handle_response_span(
"$ai_input": self._with_privacy_mode(
finalize_ai_content(_ensure_serializable(span_data.input), self._client)
),
"$ai_input_tokens": input_tokens,
"$ai_output_tokens": output_tokens,
"$ai_total_tokens": input_tokens + output_tokens,
**({"$ai_input_tokens": input_tokens} if input_tokens is not None else {}),
**(
{"$ai_output_tokens": output_tokens}
if output_tokens is not None
else {}
),
# Sum of the reported sides; omitted when neither side was reported.
**(
{"$ai_total_tokens": (input_tokens or 0) + (output_tokens or 0)}
if input_tokens is not None or output_tokens is not None
else {}
),
}

if total_cost_usd is not None:
Expand Down
Loading
Loading