Coverage for .venv/lib/python3.13/site-packages/litellm/proxy/logging_endpoints/callback_logs_endpoints.py: 87%
61 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
1"""
2Ingest pre-built logging payloads from external producers and replay them
3through LiteLLM's standard success/failure callback fan-out.
5This exists for hosts that own a request outside the Python process — e.g. the
6`litellm-rust` gateway proxying realtime websockets. Those hosts can't use the
7in-process logging object, so they POST a finished `StandardLoggingPayload` here
8and Python replays it through the exact same path a normal completion uses:
9`Logging.async_success_handler` / `async_failure_handler`. Every registered
10callback (spend logs, Langfuse, Datadog, ...) fires unchanged — there is no
11spend-logs-specific or callback-specific code here, only the replay.
13The endpoint is generic: realtime is the first producer, but the contract is the
14self-describing `StandardLoggingPayload`, so completions/responses can use it too.
15"""
17import uuid
18from collections.abc import Mapping
19from datetime import datetime, timezone
20from typing import Any, Final
22from fastapi import APIRouter, Depends, HTTPException
24from litellm._logging import verbose_proxy_logger
25from litellm.litellm_core_utils.litellm_logging import Logging as LiteLLMLogging
26from litellm.proxy._types import LitellmUserRoles, UserAPIKeyAuth
27from litellm.proxy.auth.user_api_key_auth import user_api_key_auth
28from litellm.types.proxy.callback_logs_endpoints import (
29 CallbackLogFailure,
30 CallbackLogRecord,
31 CallbackLogsRequest,
32 CallbackLogsResponse,
33)
35# Routes the Python proxy exposes for the Rust data-plane gateway to call into
36# (logging today; auth/budgets later). Namespaced under /v1/rust_control_plane so
37# they're clearly distinct from the proxy's own control-plane/management routes.
38rust_control_plane_router: Final = APIRouter(prefix="/v1/rust_control_plane", tags=["rust control plane"])
41class CallbackLogsReplayer:
42 """
43 Replays finished logging payloads through LiteLLM's callback fan-out.
45 Each helper is small and pure so the replay path is easy to read and test:
46 rebuild a `Logging` object from the payload, seed `model_call_details` with
47 exactly what the callbacks read, then dispatch to the success/failure
48 handler. No spend/callback logic lives here — only the replay.
49 """
51 @staticmethod
52 def _epoch_to_datetime(value: object) -> datetime:
53 """`StandardLoggingPayload` stores startTime/endTime as float epoch seconds."""
54 if isinstance(value, (int, float)): 54 ↛ 55line 54 didn't jump to line 55 because the condition on line 54 was never true
55 return datetime.fromtimestamp(float(value), tz=timezone.utc)
56 if isinstance(value, datetime): 56 ↛ 57line 56 didn't jump to line 57 because the condition on line 56 was never true
57 return value
58 return datetime.now(tz=timezone.utc)
60 @staticmethod
61 def _build_logging_obj(payload: dict[str, Any]) -> LiteLLMLogging:
62 """
63 Reconstruct a `Logging` object from a finished payload and seed
64 `model_call_details` with exactly what the success/failure callbacks
65 read: the prebuilt `standard_logging_object`, the resolved
66 `response_cost`, and the `litellm_params.metadata` keys used for cost
67 attribution. Setting `standard_logging_object` up front makes the handler
68 skip rebuilding it.
69 """
70 model: Final = payload.get("model") or ""
71 call_type: Final = payload.get("call_type") or "acompletion"
72 start_time: Final = CallbackLogsReplayer._epoch_to_datetime(payload.get("startTime"))
73 call_id: Final = payload.get("litellm_call_id") or payload.get("id") or str(uuid.uuid4())
75 logging_obj: Final = LiteLLMLogging(
76 model=model,
77 messages=payload.get("messages") or [],
78 # A replayed payload is always a *terminal*, fully-aggregated event —
79 # the producer (e.g. the rust gateway) already collected the whole
80 # session before POSTing. Never mark it streaming: a streaming
81 # Logging object makes async_success_handler wait for a
82 # complete_streaming_response that will never arrive, so the spend
83 # log is never written.
84 stream=False,
85 call_type=call_type,
86 start_time=start_time,
87 litellm_call_id=call_id,
88 function_id="",
89 )
91 metadata: Final[dict[str, Any]] = payload.get("metadata") or {}
92 user_api_key_hash: Final = metadata.get("user_api_key_hash")
93 litellm_metadata: Final[dict[str, Any]] = {
94 "user_api_key": user_api_key_hash,
95 "user_api_key_hash": user_api_key_hash,
96 "user_api_key_alias": metadata.get("user_api_key_alias"),
97 "user_api_key_user_id": metadata.get("user_api_key_user_id"),
98 "user_api_key_team_id": metadata.get("user_api_key_team_id"),
99 "user_api_key_org_id": metadata.get("user_api_key_org_id"),
100 "user_api_key_end_user_id": metadata.get("user_api_key_end_user_id"),
101 "spend_logs_metadata": metadata.get("spend_logs_metadata"),
102 }
104 logging_obj.model_call_details.update(
105 {
106 "model": model,
107 "call_type": call_type,
108 "custom_llm_provider": payload.get("custom_llm_provider"),
109 "response_cost": payload.get("response_cost") or 0.0,
110 "standard_logging_object": payload,
111 "litellm_params": {"metadata": litellm_metadata},
112 "cache_hit": payload.get("cache_hit") or False,
113 }
114 )
115 return logging_obj
117 @staticmethod
118 def _response_obj_from_payload(payload: Mapping[str, object]) -> dict[str, object]:
119 """Minimal response object so usage-derived spend-log fields resolve."""
120 return {
121 "id": payload.get("id"),
122 "usage": {
123 "prompt_tokens": payload.get("prompt_tokens", 0),
124 "completion_tokens": payload.get("completion_tokens", 0),
125 "total_tokens": payload.get("total_tokens", 0),
126 },
127 }
129 async def replay(self, record: CallbackLogRecord) -> None:
130 """Replay one record through the matching success/failure handler."""
131 payload: Final = record.standard_logging_payload
132 verbose_proxy_logger.debug(
133 "CallbackLogsReplayer: replaying %s record id=%s model=%s call_type=%s",
134 record.status,
135 payload.get("id"),
136 payload.get("model"),
137 payload.get("call_type"),
138 )
140 logging_obj: Final = self._build_logging_obj(payload)
141 start_time: Final = self._epoch_to_datetime(payload.get("startTime"))
142 end_time: Final = self._epoch_to_datetime(payload.get("endTime"))
144 if record.status == "success":
145 await logging_obj.async_success_handler(
146 result=self._response_obj_from_payload(payload),
147 start_time=start_time,
148 end_time=end_time,
149 )
150 else:
151 error_str: Final = record.error or payload.get("error_str") or "replayed failure"
152 await logging_obj.async_failure_handler(
153 Exception(error_str),
154 traceback_exception="",
155 start_time=start_time,
156 end_time=end_time,
157 )
159 async def replay_batch(self, records: list[CallbackLogRecord]) -> CallbackLogsResponse:
160 """Replay a batch; a single bad record never sinks the rest. Each failure
161 is reported back with its batch index so the caller can retry/triage it."""
162 processed = 0
163 failures: Final[list[CallbackLogFailure]] = []
164 for index, record in enumerate(records):
165 try:
166 await self.replay(record)
167 processed += 1
168 except Exception as e:
169 failures.append(CallbackLogFailure(index=index, error=str(e)))
170 verbose_proxy_logger.exception(
171 "CallbackLogsReplayer: failed to replay record %s: %s",
172 index,
173 str(e),
174 )
175 verbose_proxy_logger.debug(
176 "CallbackLogsReplayer: batch done processed=%s failed=%s",
177 processed,
178 len(failures),
179 )
180 return CallbackLogsResponse(processed=processed, failed=len(failures), failures=failures)
183@rust_control_plane_router.post(
184 "/logs",
185 dependencies=[Depends(user_api_key_auth)],
186)
187async def ingest_callback_logs(
188 body: CallbackLogsRequest,
189 user_api_key_dict: UserAPIKeyAuth = Depends(user_api_key_auth),
190) -> CallbackLogsResponse:
191 """
192 Replay a batch of finished logging payloads through the callback fan-out.
194 Admin-only: the payloads write spend logs and trigger every callback, so this
195 is a trusted internal route, not a public surface.
196 """
197 if user_api_key_dict.user_role != LitellmUserRoles.PROXY_ADMIN: 197 ↛ 198line 197 didn't jump to line 198 because the condition on line 197 was never true
198 raise HTTPException(
199 status_code=403,
200 detail="/v1/rust_control_plane/logs is admin-only (proxy admin key required).",
201 )
203 return await CallbackLogsReplayer().replay_batch(body.records)