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

1""" 

2Ingest pre-built logging payloads from external producers and replay them 

3through LiteLLM's standard success/failure callback fan-out. 

4 

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. 

12 

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""" 

16 

17import uuid 

18from collections.abc import Mapping 

19from datetime import datetime, timezone 

20from typing import Any, Final 

21 

22from fastapi import APIRouter, Depends, HTTPException 

23 

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) 

34 

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"]) 

39 

40 

41class CallbackLogsReplayer: 

42 """ 

43 Replays finished logging payloads through LiteLLM's callback fan-out. 

44 

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 """ 

50 

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) 

59 

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()) 

74 

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 ) 

90 

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 } 

103 

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 

116 

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 } 

128 

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 ) 

139 

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")) 

143 

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 ) 

158 

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) 

181 

182 

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. 

193 

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 ) 

202 

203 return await CallbackLogsReplayer().replay_batch(body.records)