diff --git a/litellm/integrations/generic_api/generic_api_callback.py b/litellm/integrations/generic_api/generic_api_callback.py index dfedc3a3cc9..25648f82a03 100644 --- a/litellm/integrations/generic_api/generic_api_callback.py +++ b/litellm/integrations/generic_api/generic_api_callback.py @@ -7,6 +7,7 @@ Callback to log events to a Generic API Endpoint """ import asyncio +import contextlib import json import os import re @@ -158,7 +159,7 @@ class GenericAPILogger(CustomBatchLogger): "endpoint not set for GenericAPILogger, GENERIC_LOGGER_ENDPOINT not found in environment variables" ) - self.headers: dict[str, str] = self._get_headers(headers) + self.headers: dict[str, str] = self.resolve_headers(headers) self.endpoint: str = endpoint self.event_types: list[API_EVENT_TYPES] | None = event_types self.callback_name: str | None = callback_name @@ -190,10 +191,30 @@ class GenericAPILogger(CustomBatchLogger): ######################################################### self.flush_lock = asyncio.Lock() super().__init__(**kwargs, flush_lock=self.flush_lock) - asyncio.create_task(self.periodic_flush()) + self._flush_task = asyncio.create_task(self._run_periodic_flush()) self.log_queue: list[dict | StandardLoggingPayload] = [] - def _get_headers(self, headers: dict | None = None): + async def _run_periodic_flush(self) -> None: + """Run the periodic flush loop; on cancellation, make a best-effort final + flush of any buffered logs before exiting. Prevents dropping the last batch + when this logger is evicted from the cache and cancelled via shutdown().""" + try: + await self.periodic_flush() + except asyncio.CancelledError: + # suppress(Exception) leaves a nested CancelledError (a BaseException) free to + # propagate, so a second cancellation during the final flush is still honored. + with contextlib.suppress(Exception): + await self.flush_queue() + raise + + def shutdown(self) -> None: + """Cancel the background flush task. Called by LoggingCallbackManager when this + logger is evicted from the cache so its periodic_flush task does not leak.""" + if not self._flush_task.done(): + self._flush_task.cancel() + + @staticmethod + def resolve_headers(headers: dict | None = None) -> dict[str, str]: # mutable-ok: caller owns the returned headers """ Get headers for the Generic API Logger diff --git a/litellm/litellm_core_utils/logging_callback_manager.py b/litellm/litellm_core_utils/logging_callback_manager.py index 9b612993a69..423786e8850 100644 --- a/litellm/litellm_core_utils/logging_callback_manager.py +++ b/litellm/litellm_core_utils/logging_callback_manager.py @@ -204,15 +204,19 @@ class LoggingCallbackManager: if ( isinstance(cached_logger, GenericAPILogger) and cached_logger.endpoint == endpoint - and cached_logger.headers == headers + and cached_logger.headers == GenericAPILogger.resolve_headers(headers) and cached_logger.event_types == event_types - and cached_logger.log_format == log_format + and cached_logger.log_format == (log_format if log_format is not None else "json_array") and cached_logger.max_retries == max_retries and cached_logger.retry_delay == retry_delay and cached_logger.timeout == timeout ): return cached_logger + # Cache miss (config changed or first creation). Build the replacement + # first, so that an invalid config raising in the constructor cannot leave + # the still-cached logger cancelled. Only after a successful build do we + # retire the evicted logger's background flush task. new_logger = GenericAPILogger( endpoint=endpoint, headers=headers, @@ -222,6 +226,8 @@ class LoggingCallbackManager: retry_delay=retry_delay, timeout=timeout, ) + if isinstance(cached_logger, GenericAPILogger): + cached_logger.shutdown() _generic_api_logger_cache[callback] = new_logger return new_logger