Coverage for .venv/lib/python3.13/site-packages/litellm/proxy/pass_through_endpoints/llm_provider_handlers/anthropic_passthrough_logging_handler.py: 11%
523 statements
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-10 12:01 +0000
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-10 12:01 +0000
1import asyncio
2import json
3from collections.abc import Mapping, Sequence
4from datetime import datetime
5from typing import TYPE_CHECKING, Any, Final, cast
7import httpx
9import litellm
10from litellm._logging import verbose_proxy_logger
11from litellm.constants import ANTHROPIC_BATCHES_ROUTE
12from litellm.litellm_core_utils.core_helpers import map_finish_reason
13from litellm.litellm_core_utils.litellm_logging import Logging as LiteLLMLoggingObj
14from litellm.litellm_core_utils.litellm_logging import use_custom_pricing_for_model
15from litellm.litellm_core_utils.prompt_templates.common_utils import (
16 get_content_from_model_response,
17)
18from litellm.llms.anthropic import get_anthropic_config
19from litellm.llms.anthropic.chat.handler import (
20 ModelResponseIterator as AnthropicModelResponseIterator,
21)
22from litellm.llms.anthropic.chat.transformation import AnthropicConfig
23from litellm.proxy._types import PassThroughEndpointLoggingTypedDict
24from litellm.proxy.auth.auth_utils import get_end_user_id_from_request_body
25from litellm.proxy.pass_through_endpoints.llm_provider_handlers.batch_attribution import (
26 is_collection_route,
27 log_batch_registration_result,
28 optional_str,
29 request_tags_from_metadata,
30)
31from litellm.types.passthrough_endpoints.pass_through_endpoints import (
32 PassthroughStandardLoggingPayload,
33)
34from litellm.types.utils import (
35 Choices,
36 LiteLLMBatch,
37 Message,
38 ModelResponse,
39 TextCompletionResponse,
40 Usage,
41)
43if TYPE_CHECKING: 43 ↛ 44line 43 didn't jump to line 44 because the condition on line 43 was never true
44 from litellm.types.passthrough_endpoints.pass_through_endpoints import EndpointType
46 from ..success_handler import PassThroughEndpointLogging
47else:
48 PassThroughEndpointLogging = Any
49 EndpointType = Any
52class AnthropicPassthroughLoggingHandler:
53 @staticmethod
54 def anthropic_passthrough_handler(
55 httpx_response: httpx.Response,
56 response_body: dict,
57 logging_obj: LiteLLMLoggingObj,
58 url_route: str,
59 result: str,
60 start_time: datetime,
61 end_time: datetime,
62 cache_hit: bool,
63 request_body: dict | None = None,
64 **kwargs,
65 ) -> PassThroughEndpointLoggingTypedDict:
66 """
67 Transforms Anthropic response to OpenAI response, generates a standard logging object so downstream logging can be handled
68 """
69 # Check if this is a batch creation request
70 if "/v1/messages/batches" in url_route and httpx_response.status_code == 200:
71 # Get request body from parameter or kwargs
72 request_body = request_body or kwargs.get("request_body", {})
73 return AnthropicPassthroughLoggingHandler.batch_creation_handler(
74 httpx_response=httpx_response,
75 logging_obj=logging_obj,
76 url_route=url_route,
77 result=result,
78 start_time=start_time,
79 end_time=end_time,
80 cache_hit=cache_hit,
81 request_body=request_body,
82 **kwargs,
83 )
85 model: Final = response_body.get("model", "")
86 speed: Final = AnthropicPassthroughLoggingHandler._cost_relevant_speed(
87 request_body or kwargs.get("request_body")
88 )
89 anthropic_config: Final = get_anthropic_config(url_route)
90 litellm_model_response: Final[ModelResponse] = anthropic_config().transform_response(
91 raw_response=httpx_response,
92 model_response=litellm.ModelResponse(),
93 model=model,
94 messages=[],
95 logging_obj=logging_obj,
96 optional_params={"speed": speed} if speed else {},
97 api_key="",
98 request_data={},
99 encoding=litellm.encoding,
100 json_mode=False,
101 litellm_params={},
102 )
104 kwargs = AnthropicPassthroughLoggingHandler._create_anthropic_response_logging_payload(
105 litellm_model_response=litellm_model_response,
106 model=model,
107 kwargs=kwargs,
108 start_time=start_time,
109 end_time=end_time,
110 logging_obj=logging_obj,
111 response_id=optional_str(response_body.get("id")),
112 )
114 return {
115 "result": litellm_model_response,
116 "kwargs": kwargs,
117 }
119 @staticmethod
120 def _cost_relevant_speed(request_body: Mapping[str, object] | None) -> str | None:
121 """
122 Anthropic's ``speed=fast`` multiplies token cost. The response usage carries the
123 served ``speed`` when the request asked for one, and ``calculate_usage`` prefers
124 that served value; this request-side value is the fallback when the response
125 omits it, so it still has to reach the usage-building paths.
126 """
127 speed: Final = (request_body or {}).get("speed")
128 return speed if isinstance(speed, str) else None
130 @staticmethod
131 def _get_user_from_metadata(
132 passthrough_logging_payload: PassthroughStandardLoggingPayload,
133 ) -> str | None:
134 request_body: Final = passthrough_logging_payload.get("request_body")
135 if request_body:
136 return get_end_user_id_from_request_body(request_body)
137 return None
139 @staticmethod
140 def _resolve_costing_model(model: str, logging_obj: LiteLLMLoggingObj) -> str:
141 if model and model != "unknown":
142 return model
143 litellm_params: Final = (getattr(logging_obj, "model_call_details", {}) or {}).get("litellm_params", {}) or {}
144 deployment_model: Final = litellm_params.get("model")
145 if deployment_model and deployment_model != "unknown":
146 return deployment_model
147 model_group: Final = (litellm_params.get("metadata", {}) or {}).get("model_group")
148 if model_group:
149 return model_group.removeprefix("passthrough/")
150 return model
152 @staticmethod
153 def _resolve_logged_model(
154 litellm_logging_obj: LiteLLMLoggingObj,
155 request_body: Mapping[str, object],
156 all_chunks: Sequence[str | bytes],
157 ) -> str:
158 request_model: Final = request_body.get("model")
159 logged_model: Final = (
160 request_model
161 if isinstance(request_model, str) and request_model
162 else str(litellm_logging_obj.model_call_details.get("model") or "")
163 )
164 if logged_model and logged_model != "unknown":
165 return logged_model
166 return AnthropicPassthroughLoggingHandler._extract_model_from_anthropic_chunks(all_chunks) or logged_model
168 @staticmethod
169 def _usage_only_response_or_none(
170 all_chunks: Sequence[str | bytes], model: str, speed: str | None
171 ) -> ModelResponse | None:
172 try:
173 return AnthropicPassthroughLoggingHandler._build_usage_only_response_from_chunks(
174 all_chunks=all_chunks, model=model, speed=speed
175 )
176 except Exception as e: # noqa: BLE001 # the usage-only fallback must never raise out of failure logging
177 verbose_proxy_logger.warning("Anthropic passthrough: usage-only fallback failed (model=%s): %s", model, e)
178 return None
180 @staticmethod
181 def _assemble_streaming_response(
182 all_chunks: Sequence[str | bytes],
183 litellm_logging_obj: LiteLLMLoggingObj,
184 model: str,
185 speed: str | None,
186 ) -> ModelResponse | TextCompletionResponse | None:
187 try:
188 assembled: Final = AnthropicPassthroughLoggingHandler._build_complete_streaming_response(
189 all_chunks=all_chunks,
190 litellm_logging_obj=litellm_logging_obj,
191 model=model,
192 speed=speed,
193 )
194 except Exception as e: # noqa: BLE001 # any assembly error falls back to usage-only cost
195 verbose_proxy_logger.warning(
196 "Anthropic passthrough: stream assembly raised (model=%s): %s; falling "
197 "back to usage-only cost from raw SSE events.",
198 model,
199 e,
200 )
201 return AnthropicPassthroughLoggingHandler._usage_only_response_or_none(all_chunks, model, speed)
202 if assembled is not None:
203 return assembled
204 return AnthropicPassthroughLoggingHandler._usage_only_response_or_none(all_chunks, model, speed)
206 @staticmethod
207 def _build_streaming_response_for_logging(
208 litellm_logging_obj: LiteLLMLoggingObj,
209 request_body: Mapping[str, object],
210 all_chunks: Sequence[str | bytes],
211 model: str,
212 ) -> ModelResponse | TextCompletionResponse | None:
213 response: Final = AnthropicPassthroughLoggingHandler._assemble_streaming_response(
214 all_chunks=all_chunks,
215 litellm_logging_obj=litellm_logging_obj,
216 model=model,
217 speed=AnthropicPassthroughLoggingHandler._cost_relevant_speed(request_body),
218 )
219 if not isinstance(response, ModelResponse):
220 return response
221 recovered_usage: Final = AnthropicPassthroughLoggingHandler._recover_interrupted_stream_output_tokens(
222 response=response, all_chunks=all_chunks, model=model
223 )
224 if recovered_usage is None:
225 return response
226 AnthropicPassthroughLoggingHandler._clear_placeholder_cost(response=response, usage=recovered_usage)
227 return response
229 @staticmethod
230 def record_partial_usage_for_failure(
231 litellm_logging_obj: LiteLLMLoggingObj,
232 request_body: Mapping[str, object],
233 all_chunks: Sequence[str | bytes],
234 ) -> None:
235 if not all_chunks:
236 return
237 model: Final = AnthropicPassthroughLoggingHandler._resolve_logged_model(
238 litellm_logging_obj, request_body, all_chunks
239 )
240 partial_response: Final = AnthropicPassthroughLoggingHandler._build_streaming_response_for_logging(
241 litellm_logging_obj=litellm_logging_obj, request_body=request_body, all_chunks=all_chunks, model=model
242 )
243 usage: Final = cast(Usage | None, getattr(partial_response, "usage", None))
244 if partial_response is None or usage is None:
245 return
246 litellm_logging_obj.record_partial_usage_for_failure(
247 usage=usage,
248 response_cost=AnthropicPassthroughLoggingHandler._cost_partial_stream_or_zero(
249 partial_response=partial_response, model=model, logging_obj=litellm_logging_obj
250 ),
251 )
253 @staticmethod
254 def _cost_partial_stream_or_zero(
255 partial_response: ModelResponse | TextCompletionResponse, model: str, logging_obj: LiteLLMLoggingObj
256 ) -> float:
257 try:
258 return AnthropicPassthroughLoggingHandler._compute_response_cost(
259 litellm_model_response=partial_response,
260 model=AnthropicPassthroughLoggingHandler._resolve_costing_model(model, logging_obj),
261 logging_obj=logging_obj,
262 )
263 except Exception as e: # noqa: BLE001 # an uncostable partial stream still bills its tokens, at zero cost
264 verbose_proxy_logger.warning(
265 "Anthropic passthrough: could not cost the partial usage of an interrupted stream (model=%s): %s",
266 model,
267 e,
268 )
269 return 0.0
271 @staticmethod
272 def _compute_response_cost(
273 litellm_model_response: ModelResponse | TextCompletionResponse,
274 model: str,
275 logging_obj: LiteLLMLoggingObj,
276 ) -> float:
277 if logging_obj.model_call_details.get("cache_hit") is True:
278 return 0.0
279 custom_llm_provider: Final = logging_obj.model_call_details.get("custom_llm_provider")
280 model_for_cost: Final = (
281 f"{custom_llm_provider}/{model}"
282 if custom_llm_provider and not model.startswith(f"{custom_llm_provider}/")
283 else model
284 )
285 return litellm.completion_cost(
286 completion_response=litellm_model_response,
287 model=model_for_cost,
288 custom_llm_provider=custom_llm_provider,
289 custom_pricing=use_custom_pricing_for_model(
290 litellm_params=(logging_obj.litellm_params if hasattr(logging_obj, "litellm_params") else None)
291 ),
292 router_model_id=logging_obj.get_router_model_id(),
293 )
295 @staticmethod
296 def _extract_message_start_field(
297 all_chunks: Sequence[str | bytes],
298 field: str,
299 ) -> str | None:
300 for raw in all_chunks:
301 text = raw.decode("utf-8") if isinstance(raw, bytes) else raw
302 for line in text.splitlines():
303 if not line.startswith("data:"):
304 continue
305 try:
306 data = json.loads(line[len("data:") :].strip())
307 except (json.JSONDecodeError, ValueError):
308 continue
309 if not isinstance(data, dict):
310 continue
311 if data.get("type") == "message_start":
312 value = (data.get("message") or {}).get(field)
313 if isinstance(value, str) and value:
314 return value
315 return None
317 @staticmethod
318 def _extract_model_from_anthropic_chunks(
319 all_chunks: Sequence[str | bytes],
320 ) -> str | None:
321 return AnthropicPassthroughLoggingHandler._extract_message_start_field(all_chunks, "model")
323 @staticmethod
324 def _extract_response_id_from_anthropic_chunks(
325 all_chunks: Sequence[str | bytes],
326 ) -> str | None:
327 return AnthropicPassthroughLoggingHandler._extract_message_start_field(all_chunks, "id")
329 @staticmethod
330 def _stream_was_interrupted(
331 all_chunks: Sequence[str | bytes],
332 ) -> bool:
333 """
334 Anthropic ends a stream with ``content_block_stop`` -> ``message_delta``
335 -> ``message_stop``; a client disconnect leaves the last event mid
336 ``content_block_delta``. Scan from the tail and decide on the first
337 terminal-region event, so the common completed case is O(1) rather than
338 re-deserializing every line of the stream.
339 """
340 for raw in reversed(all_chunks):
341 text = raw.decode("utf-8") if isinstance(raw, bytes) else raw
342 for line in reversed(text.splitlines()):
343 if not line.startswith("data:"):
344 continue
345 try:
346 data = json.loads(line[len("data:") :].strip())
347 except (json.JSONDecodeError, ValueError):
348 continue
349 if not isinstance(data, dict):
350 continue
351 etype = data.get("type")
352 if etype == "message_delta":
353 return False
354 if etype in (
355 "content_block_delta",
356 "content_block_stop",
357 "message_start",
358 ):
359 return True
360 return True
362 @staticmethod
363 def _recover_interrupted_stream_output_tokens(
364 response: ModelResponse | TextCompletionResponse,
365 all_chunks: Sequence[str | bytes],
366 model: str,
367 ) -> Usage | None:
368 """
369 An Anthropic stream interrupted before its terminal ``message_delta``
370 (client disconnect) carries only the ``message_start`` ``output_tokens``
371 placeholder (typically 1-3), so completion tokens and spend are
372 undercounted ~20x. Re-tokenize the buffered output text to recover a
373 realistic ``output_tokens`` for usage/cost. Completed streams are
374 untouched because their terminal ``message_delta`` short-circuits here.
375 """
376 if not isinstance(response, ModelResponse):
377 return None
378 if not AnthropicPassthroughLoggingHandler._stream_was_interrupted(all_chunks):
379 return None
380 usage: Final = getattr(response, "usage", None)
381 if not isinstance(usage, Usage):
382 return None
383 output_text: Final = get_content_from_model_response(response)
384 if not output_text:
385 return None
386 try:
387 recovered_output_tokens = litellm.token_counter(model=model, text=output_text, count_response_tokens=True)
388 except Exception:
389 verbose_proxy_logger.warning(
390 "Could not re-tokenize interrupted stream output; keeping placeholder completion token count."
391 )
392 return None
393 if recovered_output_tokens <= (usage.completion_tokens or 0):
394 return None
395 usage.completion_tokens = recovered_output_tokens
396 usage.total_tokens = (usage.prompt_tokens or 0) + recovered_output_tokens
397 # Anthropic costing reads completion_tokens_details.text_tokens, so the
398 # stale message_start placeholder there must be corrected too or spend
399 # stays undercounted even after completion_tokens is fixed.
400 details: Final = getattr(usage, "completion_tokens_details", None)
401 if details is not None and getattr(details, "text_tokens", None) is not None:
402 details.text_tokens = recovered_output_tokens
403 return usage
405 @staticmethod
406 def _clear_placeholder_cost(response: ModelResponse, usage: Usage) -> None:
407 usage.cost = None
408 response._hidden_params.pop("response_cost", None) # pyright: ignore[reportPrivateUsage] # no public accessor
410 @staticmethod
411 def _create_anthropic_response_logging_payload(
412 litellm_model_response: ModelResponse | TextCompletionResponse,
413 model: str,
414 kwargs: dict,
415 start_time: datetime,
416 end_time: datetime,
417 logging_obj: LiteLLMLoggingObj,
418 response_id: str | None = None,
419 ):
420 """
421 Create the standard logging object for Anthropic passthrough
423 handles streaming and non-streaming responses
424 """
425 # Only record complete_streaming_response for actual streaming responses.
426 # perform_redaction scrubs this field only when stream is True, so setting
427 # it on a non-streaming response would bypass message redaction.
428 if logging_obj.model_call_details.get("stream") is True:
429 logging_obj.model_call_details["complete_streaming_response"] = litellm_model_response
430 try:
431 model = AnthropicPassthroughLoggingHandler._resolve_costing_model(model, logging_obj)
432 response_cost: Final = AnthropicPassthroughLoggingHandler._compute_response_cost(
433 litellm_model_response=litellm_model_response, model=model, logging_obj=logging_obj
434 )
436 kwargs["response_cost"] = response_cost
437 kwargs["model"] = model
438 # the pass-through success path reads spend from
439 # model_call_details["response_cost"], not from kwargs
440 logging_obj.model_call_details["response_cost"] = response_cost
441 passthrough_logging_payload: Final[PassthroughStandardLoggingPayload | None] = kwargs.get(
442 "passthrough_logging_payload"
443 )
444 if passthrough_logging_payload:
445 user: Final = AnthropicPassthroughLoggingHandler._get_user_from_metadata(
446 passthrough_logging_payload=passthrough_logging_payload,
447 )
448 if user:
449 kwargs.setdefault("litellm_params", {})
450 kwargs["litellm_params"].update({"proxy_server_request": {"body": {"user": user}}})
452 # pretty print standard logging object
453 verbose_proxy_logger.debug(
454 "kwargs= %s",
455 json.dumps(kwargs, indent=4, default=str),
456 )
458 litellm_model_response.id = response_id or logging_obj.litellm_call_id
459 litellm_model_response.model = model
460 logging_obj.model_call_details["model"] = model
461 if not logging_obj.model_call_details.get("custom_llm_provider"):
462 logging_obj.model_call_details["custom_llm_provider"] = litellm.LlmProviders.ANTHROPIC.value
463 return kwargs
464 except Exception as e:
465 verbose_proxy_logger.exception("Error creating Anthropic response logging payload: %s", e)
466 return kwargs
468 @staticmethod
469 def _handle_logging_anthropic_collected_chunks(
470 litellm_logging_obj: LiteLLMLoggingObj,
471 passthrough_success_handler_obj: PassThroughEndpointLogging,
472 url_route: str,
473 request_body: dict,
474 endpoint_type: EndpointType,
475 start_time: datetime,
476 all_chunks: list[str],
477 end_time: datetime,
478 ) -> PassThroughEndpointLoggingTypedDict:
479 """
480 Takes raw chunks from Anthropic passthrough endpoint and logs them in litellm callbacks
482 - Builds complete response from chunks
483 - Creates standard logging object
484 - Logs in litellm callbacks
485 """
487 model: Final = AnthropicPassthroughLoggingHandler._resolve_logged_model(
488 litellm_logging_obj, request_body, all_chunks
489 )
490 complete_streaming_response: Final = AnthropicPassthroughLoggingHandler._build_streaming_response_for_logging(
491 litellm_logging_obj=litellm_logging_obj, request_body=request_body, all_chunks=all_chunks, model=model
492 )
493 if complete_streaming_response is None:
494 verbose_proxy_logger.error(
495 "Unable to build complete streaming response for Anthropic passthrough endpoint, not logging..."
496 )
497 return {
498 "result": None,
499 "kwargs": {},
500 }
501 kwargs: Final = AnthropicPassthroughLoggingHandler._create_anthropic_response_logging_payload(
502 litellm_model_response=complete_streaming_response,
503 model=model,
504 kwargs={},
505 start_time=start_time,
506 end_time=end_time,
507 logging_obj=litellm_logging_obj,
508 response_id=AnthropicPassthroughLoggingHandler._extract_response_id_from_anthropic_chunks(all_chunks),
509 )
511 return {
512 "result": complete_streaming_response,
513 "kwargs": kwargs,
514 }
516 @staticmethod
517 def _split_sse_chunk_into_events(chunk: str | bytes) -> list[str]:
518 """
519 Split a chunk that may contain multiple SSE events into individual events.
521 SSE format: "event: type\ndata: {...}\n\n"
522 Multiple events in a single chunk are separated by double newlines.
524 Args:
525 chunk: Raw chunk string that may contain multiple SSE events
527 Returns:
528 List of individual SSE event strings (each containing "event: X\ndata: {...}")
529 """
530 # Handle bytes input
531 if isinstance(chunk, bytes):
532 chunk = chunk.decode("utf-8")
534 # Split on double newlines to separate SSE events
535 # Filter out empty strings
536 events: Final = [event.strip() for event in chunk.split("\n\n") if event.strip()]
538 return events
540 @staticmethod
541 def _build_complete_streaming_response(
542 all_chunks: Sequence[str | bytes],
543 litellm_logging_obj: LiteLLMLoggingObj,
544 model: str,
545 speed: str | None = None,
546 ) -> ModelResponse | TextCompletionResponse | None:
547 """
548 Builds complete response from raw Anthropic chunks.
550 Fast path: for the dominant case of a pure-text streaming response
551 (no tool_use / thinking / non-text content blocks), the long run of
552 ``content_block_delta`` text deltas is collapsed into a single
553 equivalent SSE event before conversion. ``chunk_parser`` and
554 ``stream_chunk_builder`` remain the single source of truth for chunk
555 shape, usage math and finish-reason mapping, so the rebuilt response
556 (and therefore the logged/billed payload) is identical -- this is
557 asserted by a parity test. Anything non-trivial falls back to the
558 unchanged legacy reconstruction.
560 Per-event Pydantic ``ModelResponseStream`` construction dominated
561 event-loop CPU under concurrent streaming; collapsing the homogeneous
562 text run removes O(num_output_tokens) of it.
563 """
564 collapsed: Final = AnthropicPassthroughLoggingHandler._collapse_pure_text_chunks(all_chunks)
565 if collapsed is not None:
566 return AnthropicPassthroughLoggingHandler._build_complete_streaming_response_legacy(
567 all_chunks=collapsed,
568 litellm_logging_obj=litellm_logging_obj,
569 model=model,
570 speed=speed,
571 )
572 return AnthropicPassthroughLoggingHandler._build_complete_streaming_response_legacy(
573 all_chunks=all_chunks,
574 litellm_logging_obj=litellm_logging_obj,
575 model=model,
576 speed=speed,
577 )
579 # Anthropic SSE block/delta types that the fast path is NOT allowed to
580 # collapse -- their presence forces the unchanged legacy path so tool
581 # calls, thinking, citations, etc. keep byte-identical reconstruction.
582 _FAST_PATH_DISALLOWED_DELTA_TYPES = frozenset(
583 {
584 "input_json_delta",
585 "thinking_delta",
586 "signature_delta",
587 "citations_delta",
588 }
589 )
591 @staticmethod
592 def _collapse_pure_text_chunks(
593 all_chunks: Sequence[str | bytes],
594 ) -> list[str] | None:
595 """
596 Return a new chunk list with the contiguous run of text-only
597 ``content_block_delta`` events replaced by a single equivalent event,
598 or ``None`` if the stream is not a pure single-text-block response
599 (in which case the caller uses the legacy path unchanged).
601 Only ``message_start`` / ``content_block_start(text)`` /
602 ``content_block_delta(text_delta)`` / ``content_block_stop`` /
603 ``message_delta`` / ``message_stop`` / ``ping`` events are accepted.
604 Any other content-block type or delta type returns ``None``.
605 """
606 normalized: Final[list[str]] = []
607 for raw in all_chunks:
608 line = raw.decode("utf-8") if isinstance(raw, bytes) else raw
609 for ev in line.split("\n\n"):
610 ev = ev.strip()
611 if ev:
612 normalized.append(ev)
614 text_block_indexes: Final[set] = set()
615 out: Final[list[str]] = []
616 pending_text: list[str] = []
617 pending_index: int | None = None
618 saw_any_text_delta = False
620 def flush() -> None:
621 nonlocal pending_text, pending_index
622 if pending_text:
623 merged: Final = {
624 "type": "content_block_delta",
625 "index": pending_index if pending_index is not None else 0,
626 "delta": {"type": "text_delta", "text": "".join(pending_text)},
627 }
628 out.append("data: " + json.dumps(merged))
629 pending_text = []
630 pending_index = None
632 for ev in normalized:
633 idx = ev.find("data:")
634 if idx == -1:
635 # Bare "event: <name>" line. The legacy converter turns this
636 # into an empty ModelResponseStream that contributes nothing
637 # to stream_chunk_builder. Drop the high-frequency interior
638 # markers (content_block_delta / ping); keep every other
639 # bare event line verbatim so chunk ordering and the
640 # load-bearing chunks[0] (event: message_start) are retained.
641 name = ev[len("event:") :].strip() if ev.startswith("event:") else ""
642 if name in ("content_block_delta", "ping"):
643 continue
644 flush()
645 out.append(ev)
646 continue
648 json_str = ev[idx + len("data:") :].strip()
649 try:
650 data = json.loads(json_str)
651 except (json.JSONDecodeError, ValueError):
652 return None
654 etype = data.get("type")
655 if etype == "content_block_start":
656 block = data.get("content_block") or {}
657 if block.get("type") != "text":
658 return None
659 text_block_indexes.add(data.get("index"))
660 flush()
661 out.append(ev)
662 elif etype == "content_block_delta":
663 delta = data.get("delta") or {}
664 dtype = delta.get("type")
665 if dtype in AnthropicPassthroughLoggingHandler._FAST_PATH_DISALLOWED_DELTA_TYPES:
666 return None
667 if dtype != "text_delta":
668 return None
669 cur_index = data.get("index")
670 if cur_index not in text_block_indexes:
671 return None
672 # Defensive: Anthropic sends blocks strictly sequentially
673 # (start/deltas/stop, then next block), so pending_text from
674 # block N must be flushed by content_block_stop before block
675 # N+1's deltas arrive. If we ever see a delta whose index
676 # disagrees with the current pending buffer, the stream is
677 # interleaved -- fall back to legacy rather than risk merging
678 # text from different blocks under a single index.
679 if pending_text and pending_index is not None and cur_index != pending_index:
680 return None
681 saw_any_text_delta = True
682 pending_index = cur_index
683 pending_text.append(delta.get("text") or "")
684 elif etype == "ping":
685 # Interior no-op; legacy maps it to an empty chunk.
686 continue
687 else:
688 # message_start / content_block_stop / message_delta /
689 # message_stop / error: pass through unchanged.
690 flush()
691 out.append(ev)
693 flush()
695 if not saw_any_text_delta:
696 return None
697 return out
699 @staticmethod
700 def _build_complete_streaming_response_legacy(
701 all_chunks: Sequence[str | bytes],
702 litellm_logging_obj: LiteLLMLoggingObj,
703 model: str,
704 speed: str | None = None,
705 ) -> ModelResponse | TextCompletionResponse | None:
706 """
707 Original reconstruction: convert every SSE event to a generic chunk
708 and assemble via stream_chunk_builder. Kept verbatim as the fallback
709 / source of truth for the fast path's parity test.
711 - Splits multi-event chunks into individual SSE events
712 - Converts str chunks to generic chunks
713 - Converts generic chunks to litellm chunks (OpenAI format)
714 - Builds complete response from litellm chunks
715 """
716 verbose_proxy_logger.debug("Building complete streaming response from %d chunks", len(all_chunks))
717 anthropic_model_response_iterator: Final = AnthropicModelResponseIterator(
718 streaming_response=None,
719 sync_stream=False,
720 speed=speed,
721 )
722 all_openai_chunks: Final = []
724 # Process each chunk - a chunk may contain multiple SSE events
725 for _chunk_str in all_chunks:
726 # Split chunk into individual SSE events
727 individual_events = AnthropicPassthroughLoggingHandler._split_sse_chunk_into_events(_chunk_str)
729 # Process each individual event
730 for event_str in individual_events:
731 try:
732 # Skip OpenAI-style [DONE] sentinels some Anthropic-compatible
733 # providers emit. Match the whole SSE line so a valid chunk whose
734 # text payload happens to contain "[DONE]" is not dropped.
735 if any(line.strip() == "data: [DONE]" for line in event_str.split("\n")):
736 continue
737 transformed_openai_chunk = anthropic_model_response_iterator.convert_str_chunk_to_generic_chunk(
738 chunk=event_str
739 )
740 if transformed_openai_chunk is not None:
741 all_openai_chunks.append(transformed_openai_chunk)
743 except (StopIteration, StopAsyncIteration):
744 break
745 except json.JSONDecodeError:
746 # Some upstreams emit non-JSON SSE lines; skip them so the
747 # logging pipeline is not broken by a single bad frame.
748 verbose_proxy_logger.debug(
749 "Skipping non-JSON SSE event: %s",
750 event_str[:200],
751 )
752 continue
754 complete_streaming_response: Final = litellm.stream_chunk_builder(
755 chunks=all_openai_chunks,
756 logging_obj=litellm_logging_obj,
757 )
758 verbose_proxy_logger.debug("Complete streaming response built: %s", complete_streaming_response)
759 return complete_streaming_response
761 @staticmethod
762 def _extract_sse_data(event_str: str) -> dict | None:
763 """Parse the JSON object from the ``data:`` line of an Anthropic SSE event."""
764 for line in event_str.splitlines():
765 stripped = line.strip()
766 if stripped.startswith("data:"):
767 payload = stripped[len("data:") :].strip()
768 if not payload or payload == "[DONE]":
769 return None
770 try:
771 return cast(dict, json.loads(payload))
772 except (ValueError, TypeError):
773 return None
774 return None
776 @staticmethod
777 def _build_usage_only_response_from_chunks(
778 all_chunks: Sequence[str | bytes],
779 model: str,
780 speed: str | None = None,
781 ) -> ModelResponse | None:
782 """
783 Build a usage-bearing ModelResponse from Anthropic SSE token-usage events, for
784 cost tracking when stream_chunk_builder cannot reassemble the stream.
786 Anthropic emits usage in ``message_start`` (uncached input + cache tokens, and an
787 initial output_tokens) and the final ``message_delta`` (cumulative output_tokens)
788 regardless of the content/tool shape, so cost is recoverable even when full
789 content assembly fails. Returns ``None`` if no usage event is found.
790 """
791 input_tokens = 0
792 cache_read = 0
793 cache_creation = 0
794 cache_creation_5m: int | None = None
795 cache_creation_1h: int | None = None
796 output_tokens = 0
797 web_search_requests: int | None = None
798 tool_search_requests: int | None = None
799 inference_geo: str | None = None
800 speed_from_stream: str | None = None
801 stop_reason: str | None = None
802 found_usage = False
803 resolved_model = model
804 for _chunk_str in all_chunks:
805 for event_str in AnthropicPassthroughLoggingHandler._split_sse_chunk_into_events(_chunk_str):
806 data = AnthropicPassthroughLoggingHandler._extract_sse_data(event_str)
807 if not data:
808 continue
809 event_type = data.get("type")
810 if event_type == "message_start":
811 message = data.get("message") or {}
812 if not resolved_model or resolved_model == "unknown":
813 resolved_model = message.get("model") or resolved_model
814 usage = message.get("usage") or {}
815 input_tokens = usage.get("input_tokens") or input_tokens
816 cache_read = usage.get("cache_read_input_tokens") or cache_read
817 cache_creation = usage.get("cache_creation_input_tokens") or cache_creation
818 _cc = usage.get("cache_creation")
819 if isinstance(_cc, dict):
820 cache_creation_5m = _cc.get("ephemeral_5m_input_tokens")
821 cache_creation_1h = _cc.get("ephemeral_1h_input_tokens")
822 if usage.get("inference_geo") is not None:
823 inference_geo = usage.get("inference_geo")
824 if isinstance(usage.get("speed"), str):
825 speed_from_stream = usage.get("speed")
826 if usage.get("output_tokens") is not None:
827 output_tokens = usage.get("output_tokens")
828 found_usage = True
829 elif event_type == "message_delta":
830 _delta_stop = (data.get("delta") or {}).get("stop_reason")
831 if _delta_stop:
832 stop_reason = _delta_stop
833 usage = data.get("usage") or {}
834 if usage.get("output_tokens") is not None:
835 output_tokens = usage.get("output_tokens")
836 _stu = usage.get("server_tool_use")
837 if isinstance(_stu, dict):
838 if _stu.get("web_search_requests") is not None:
839 web_search_requests = _stu.get("web_search_requests")
840 if _stu.get("tool_search_requests") is not None:
841 tool_search_requests = _stu.get("tool_search_requests")
842 if usage.get("cache_read_input_tokens") is not None:
843 cache_read = usage.get("cache_read_input_tokens")
844 if usage.get("inference_geo") is not None:
845 inference_geo = usage.get("inference_geo")
846 if isinstance(usage.get("speed"), str):
847 speed_from_stream = usage.get("speed")
848 found_usage = True
849 if not found_usage:
850 return None
851 # If only the 5m/1h split was provided, derive the cache_creation total from it.
852 if not cache_creation and (cache_creation_5m or cache_creation_1h):
853 cache_creation = (cache_creation_5m or 0) + (cache_creation_1h or 0)
854 # build usage via the same AnthropicConfig.calculate_usage path the success
855 # cases use, so prompt_tokens are cache-inclusive and cache / server_tool_use /
856 # inference_geo tokens are priced instead of left at $0
857 usage_object: Final[dict] = {
858 "input_tokens": input_tokens,
859 "output_tokens": output_tokens,
860 }
861 if cache_read:
862 usage_object["cache_read_input_tokens"] = cache_read
863 if cache_creation:
864 usage_object["cache_creation_input_tokens"] = cache_creation
865 if cache_creation_5m is not None or cache_creation_1h is not None:
866 usage_object["cache_creation"] = {
867 "ephemeral_5m_input_tokens": cache_creation_5m or 0,
868 "ephemeral_1h_input_tokens": cache_creation_1h or 0,
869 }
870 if web_search_requests is not None or tool_search_requests is not None:
871 _server_tool_use: Final[dict] = {}
872 if web_search_requests is not None:
873 _server_tool_use["web_search_requests"] = web_search_requests
874 if tool_search_requests is not None:
875 _server_tool_use["tool_search_requests"] = tool_search_requests
876 usage_object["server_tool_use"] = _server_tool_use
877 if inference_geo is not None:
878 usage_object["inference_geo"] = inference_geo
879 if speed_from_stream is not None:
880 usage_object["speed"] = speed_from_stream
881 usage_obj: Final = AnthropicConfig().calculate_usage(
882 usage_object=usage_object, reasoning_content=None, speed=speed
883 )
884 return ModelResponse(
885 model=resolved_model,
886 choices=[
887 Choices(
888 finish_reason=(map_finish_reason(stop_reason) if stop_reason else "stop"),
889 index=0,
890 message=Message(role="assistant", content=""),
891 )
892 ],
893 usage=usage_obj,
894 )
896 @staticmethod
897 def batch_creation_handler(
898 httpx_response: httpx.Response,
899 logging_obj: LiteLLMLoggingObj,
900 url_route: str,
901 result: str,
902 start_time: datetime,
903 end_time: datetime,
904 cache_hit: bool,
905 request_body: dict | None = None,
906 **kwargs,
907 ) -> PassThroughEndpointLoggingTypedDict:
908 """
909 Handle Anthropic batch creation passthrough logging.
910 Creates a managed object for cost tracking when batch job is successfully created.
911 """
912 import base64
914 from litellm._uuid import uuid
915 from litellm.llms.anthropic.batches.transformation import AnthropicBatchesConfig
916 from litellm.types.utils import Choices, SpecialEnums
918 try:
919 _json_response: Final = httpx_response.json()
921 # Only handle successful batch job creation (POST requests with 201 status)
922 if httpx_response.status_code == 200 and "id" in _json_response:
923 # Transform Anthropic response to LiteLLM batch format
924 anthropic_batches_config: Final = AnthropicBatchesConfig()
925 litellm_batch_response: Final = anthropic_batches_config.transform_retrieve_batch_response(
926 model=None,
927 raw_response=httpx_response,
928 logging_obj=logging_obj,
929 litellm_params={},
930 )
931 # Set status to "validating" for newly created batches so polling mechanism picks them up
932 # The polling mechanism only looks for status="validating" jobs
933 litellm_batch_response.status = "validating"
935 # Extract batch ID from the response
936 batch_id: Final = _json_response.get("id", "")
938 # Get model from request body (batch response doesn't include model)
939 request_body = request_body or {}
940 # Try to extract model from the batch request body, supporting Anthropic's nested structure
941 model_name: str = "unknown"
942 if isinstance(request_body, dict):
943 # Standard: {"model": ...}
944 model_name = request_body.get("model") or "unknown"
945 if model_name == "unknown":
946 # Anthropic batches: look under requests[0].params.model
947 requests_list: Final = request_body.get("requests", [])
948 if isinstance(requests_list, list) and len(requests_list) > 0:
949 first_req: Final = requests_list[0]
950 if isinstance(first_req, dict):
951 params: Final = first_req.get("params", {})
952 if isinstance(params, dict):
953 extracted_model: Final = params.get("model")
954 if extracted_model:
955 model_name = extracted_model
957 # Create unified object ID for tracking
958 # Format: base64(litellm_proxy;model_id:{};llm_batch_id:{})
959 # For Anthropic passthrough, prefix model with "anthropic/" so router can determine provider
960 actual_model_id = AnthropicPassthroughLoggingHandler.get_actual_model_id_from_router(model_name)
962 # If model not in router, use "anthropic/{model_name}" format so router can determine provider
963 if actual_model_id == model_name and not actual_model_id.startswith("anthropic/"):
964 actual_model_id = f"anthropic/{model_name}"
966 unified_id_string: Final = SpecialEnums.LITELLM_MANAGED_BATCH_COMPLETE_STR.value.format(
967 actual_model_id, batch_id
968 )
969 unified_object_id: Final = base64.urlsafe_b64encode(unified_id_string.encode()).decode().rstrip("=")
971 # Store the managed object for cost tracking
972 # This will be picked up by check_batch_cost polling mechanism
973 if is_collection_route(url_route, ANTHROPIC_BATCHES_ROUTE):
974 AnthropicPassthroughLoggingHandler._store_batch_managed_object(
975 unified_object_id=unified_object_id,
976 batch_object=litellm_batch_response,
977 model_object_id=batch_id,
978 logging_obj=logging_obj,
979 **kwargs,
980 )
982 # Create a batch job response for logging
983 litellm_model_response = ModelResponse()
984 litellm_model_response.id = str(uuid.uuid4())
985 litellm_model_response.model = model_name
986 litellm_model_response.object = "batch"
987 litellm_model_response.created = int(start_time.timestamp())
989 # Add batch-specific metadata to indicate this is a pending batch job
990 litellm_model_response.choices = [
991 Choices(
992 finish_reason="stop",
993 index=0,
994 message={
995 "role": "assistant",
996 "content": f"Batch job {batch_id} created and is pending. Status will be updated when the batch completes.",
997 "tool_calls": None,
998 "function_call": None,
999 "provider_specific_fields": {
1000 "batch_job_id": batch_id,
1001 "batch_job_state": "in_progress",
1002 "unified_object_id": unified_object_id,
1003 },
1004 },
1005 )
1006 ]
1008 # Set response cost to 0 initially (will be updated when batch completes)
1009 response_cost: Final = 0.0
1010 kwargs["response_cost"] = response_cost
1011 kwargs["model"] = model_name
1012 kwargs["batch_id"] = batch_id
1013 kwargs["unified_object_id"] = unified_object_id
1014 kwargs["batch_job_state"] = "in_progress"
1016 logging_obj.model = model_name
1017 logging_obj.model_call_details["model"] = logging_obj.model
1018 logging_obj.model_call_details["response_cost"] = response_cost
1019 logging_obj.model_call_details["batch_id"] = batch_id
1021 return {
1022 "result": litellm_model_response,
1023 "kwargs": kwargs,
1024 }
1025 else:
1026 # Handle non-successful responses
1027 litellm_model_response = ModelResponse()
1028 litellm_model_response.id = str(uuid.uuid4())
1029 litellm_model_response.model = "anthropic_batch"
1030 litellm_model_response.object = "batch"
1031 litellm_model_response.created = int(start_time.timestamp())
1033 # Add error-specific metadata
1034 litellm_model_response.choices = [
1035 Choices(
1036 finish_reason="stop",
1037 index=0,
1038 message={
1039 "role": "assistant",
1040 "content": f"Batch job creation failed. Status: {httpx_response.status_code}",
1041 "tool_calls": None,
1042 "function_call": None,
1043 "provider_specific_fields": {
1044 "batch_job_state": "failed",
1045 "status_code": httpx_response.status_code,
1046 },
1047 },
1048 )
1049 ]
1051 kwargs["response_cost"] = 0.0
1052 kwargs["model"] = "anthropic_batch"
1053 kwargs["batch_job_state"] = "failed"
1055 return {
1056 "result": litellm_model_response,
1057 "kwargs": kwargs,
1058 }
1060 except Exception as e:
1061 verbose_proxy_logger.error("Error in batch_creation_handler: %s", e)
1062 # Return basic response on error
1063 litellm_model_response = ModelResponse()
1064 litellm_model_response.id = str(uuid.uuid4())
1065 litellm_model_response.model = "anthropic_batch"
1066 litellm_model_response.object = "batch"
1067 litellm_model_response.created = int(start_time.timestamp())
1069 # Add error-specific metadata
1070 litellm_model_response.choices = [
1071 Choices(
1072 finish_reason="stop",
1073 index=0,
1074 message={
1075 "role": "assistant",
1076 "content": f"Error creating batch job: {e}",
1077 "tool_calls": None,
1078 "function_call": None,
1079 "provider_specific_fields": {
1080 "batch_job_state": "failed",
1081 "error": str(e),
1082 },
1083 },
1084 )
1085 ]
1087 kwargs["response_cost"] = 0.0
1088 kwargs["model"] = "anthropic_batch"
1089 kwargs["batch_job_state"] = "failed"
1091 return {
1092 "result": litellm_model_response,
1093 "kwargs": kwargs,
1094 }
1096 @staticmethod
1097 def _store_batch_managed_object(
1098 unified_object_id: str,
1099 batch_object: LiteLLMBatch,
1100 model_object_id: str,
1101 logging_obj: LiteLLMLoggingObj,
1102 **kwargs,
1103 ) -> None:
1104 """
1105 Register a newly created batch for cost tracking.
1106 This will be picked up by the check_batch_cost polling mechanism.
1108 Only the create reaches here, so the row records the creating key and its tags.
1109 An id-scoped route cannot rebuild the unified object id anyway: the model comes
1110 from the create's request body, which a retrieve does not have.
1111 """
1112 try:
1113 # Get the managed files hook from the logging object
1114 # This is a bit of a hack, but we need access to the proxy logging system
1115 from litellm.proxy.proxy_server import proxy_logging_obj
1117 managed_files_hook: Final = proxy_logging_obj.get_proxy_hook("managed_files")
1118 if managed_files_hook is not None and hasattr(managed_files_hook, "store_unified_object_id"):
1119 # Create a mock user API key dict for the managed object storage
1120 from litellm.proxy._types import LitellmUserRoles, UserAPIKeyAuth
1122 _request_metadata: Final = (kwargs.get("litellm_params", {}) or {}).get("metadata", {}) or {}
1124 user_api_key_dict: Final = UserAPIKeyAuth(
1125 user_id=_request_metadata.get("user_api_key_user_id", "default-user"),
1126 api_key=optional_str(_request_metadata.get("user_api_key")),
1127 team_id=_request_metadata.get("user_api_key_team_id"),
1128 team_alias=None,
1129 user_role=LitellmUserRoles.CUSTOMER, # Use proper enum value
1130 user_email=None,
1131 max_budget=None,
1132 spend=0.0, # Set to 0.0 instead of None
1133 models=[], # Set to empty list instead of None
1134 tpm_limit=None,
1135 rpm_limit=None,
1136 budget_duration=None,
1137 budget_reset_at=None,
1138 max_parallel_requests=None,
1139 allowed_model_region=None,
1140 metadata={}, # Set to empty dict instead of None
1141 key_alias=None,
1142 permissions={}, # Set to empty dict instead of None
1143 model_max_budget={}, # Set to empty dict instead of None
1144 model_spend={}, # Set to empty dict instead of None
1145 )
1147 # Store the unified object for batch cost tracking
1148 task: Final = asyncio.create_task(
1149 managed_files_hook.store_unified_object_id(
1150 unified_object_id=unified_object_id,
1151 file_object=batch_object,
1152 litellm_parent_otel_span=None,
1153 model_object_id=model_object_id,
1154 file_purpose="batch",
1155 user_api_key_dict=user_api_key_dict,
1156 request_tags=request_tags_from_metadata(_request_metadata),
1157 persist_attribution=True,
1158 )
1159 )
1160 task.add_done_callback(
1161 lambda finished: log_batch_registration_result(
1162 finished, "Anthropic", unified_object_id, model_object_id, is_batch_create=True
1163 )
1164 )
1165 else:
1166 verbose_proxy_logger.warning(
1167 "Managed files hook not available, cannot store batch object for cost tracking"
1168 )
1170 except Exception as e:
1171 verbose_proxy_logger.error("Error storing Anthropic batch managed object: %s", e)
1173 @staticmethod
1174 def get_actual_model_id_from_router(model_name: str) -> str:
1175 from litellm.proxy.proxy_server import llm_router
1177 if llm_router is not None:
1178 # Try to find the model in the router by the model name
1179 # Use the existing get_model_ids method from router
1180 model_ids: Final = llm_router.get_model_ids(model_name=model_name)
1181 if model_ids and len(model_ids) > 0:
1182 # Use the first model ID found
1183 actual_model_id = model_ids[0]
1184 verbose_proxy_logger.info("Found model ID in router: %s", actual_model_id)
1185 return actual_model_id
1186 else:
1187 # Fallback to model name
1188 actual_model_id = model_name
1189 verbose_proxy_logger.warning("Model not found in router, using model name: %s", actual_model_id)
1190 return actual_model_id
1191 else:
1192 # Fallback if router is not available
1193 verbose_proxy_logger.warning("Router not available, using model name: %s", model_name)
1194 return model_name