diff --git a/litellm/integrations/rubrik.py b/litellm/integrations/rubrik.py index 665d1209a30..df5b020d61f 100644 --- a/litellm/integrations/rubrik.py +++ b/litellm/integrations/rubrik.py @@ -37,6 +37,8 @@ if TYPE_CHECKING: _ENDPOINT_ANTHROPIC_MESSAGES = "/v1/messages" _WEBHOOK_PATH_TOOL_BLOCKING = "/v1/after_completion/openai/v1" _WEBHOOK_PATH_LOGGING_BATCH = "/v1/litellm/batch" +_DEFAULT_MAX_QUEUE_SIZE = 10_000 +_DROP_WARNING_INTERVAL_SECONDS = 60.0 class RubrikLogger(CustomGuardrail, CustomBatchLogger): @@ -95,6 +97,30 @@ class RubrikLogger(CustomGuardrail, CustomBatchLogger): f"Invalid RUBRIK_BATCH_SIZE: {_batch_size!r}, using default" ) + # Cap the in-memory retry queue so a Rubrik webhook outage cannot let + # authenticated traffic accumulate prompt/response payloads until the + # proxy runs out of memory. Once the cap is reached, oldest events are + # dropped to make room for fresh ones (drop-oldest backpressure). + self.max_queue_size = _DEFAULT_MAX_QUEUE_SIZE + _max_queue_size = os.getenv("RUBRIK_MAX_QUEUE_SIZE") + if _max_queue_size: + try: + parsed_max = int(_max_queue_size) + if parsed_max > 0: + self.max_queue_size = parsed_max + else: + verbose_logger.warning( + f"RUBRIK_MAX_QUEUE_SIZE={_max_queue_size!r} must be > 0; " + f"using default {_DEFAULT_MAX_QUEUE_SIZE}" + ) + except ValueError: + verbose_logger.warning( + f"Invalid RUBRIK_MAX_QUEUE_SIZE: {_max_queue_size!r}, " + f"using default {_DEFAULT_MAX_QUEUE_SIZE}" + ) + self._dropped_since_warning = 0 + self._last_drop_warning_time = 0.0 + _webhook_url = api_base or os.getenv("RUBRIK_WEBHOOK_URL") if _webhook_url is None: @@ -391,6 +417,7 @@ class RubrikLogger(CustomGuardrail, CustomBatchLogger): return self.log_queue.append(payload) + self._enforce_max_queue_size() if len(self.log_queue) >= self.batch_size: await self.flush_queue() @@ -401,6 +428,24 @@ class RubrikLogger(CustomGuardrail, CustomBatchLogger): exc_info=True, ) + def _enforce_max_queue_size(self) -> None: + overflow = len(self.log_queue) - self.max_queue_size + if overflow <= 0: + return + del self.log_queue[:overflow] + self._dropped_since_warning += overflow + now = time.time() + if now - self._last_drop_warning_time >= _DROP_WARNING_INTERVAL_SECONDS: + verbose_logger.warning( + "Rubrik: log queue exceeded max_queue_size=%s; dropped %s " + "oldest events since the last warning. The Rubrik webhook may " + "be unhealthy or undersized for current traffic.", + self.max_queue_size, + self._dropped_since_warning, + ) + self._dropped_since_warning = 0 + self._last_drop_warning_time = now + async def async_log_success_event(self, kwargs, response_obj, start_time, end_time): await self._enqueue_log_event(kwargs, "success") diff --git a/tests/test_litellm/integrations/test_rubrik.py b/tests/test_litellm/integrations/test_rubrik.py index 56799641d90..922d2fe8a15 100644 --- a/tests/test_litellm/integrations/test_rubrik.py +++ b/tests/test_litellm/integrations/test_rubrik.py @@ -349,6 +349,27 @@ class TestBatchLogging: await handler.flush_queue() assert handler.log_queue == [{"msg": "a"}, {"msg": "b"}] + async def test_enqueue_drops_oldest_when_queue_exceeds_max_size(self, handler): + """A sustained Rubrik webhook outage must not let the in-memory retry + queue grow without bound. Once max_queue_size is exceeded, the oldest + events are dropped to make room for new ones.""" + handler.max_queue_size = 3 + handler.batch_size = 10**6 # disable size-triggered flush + handler.flush_queue = AsyncMock() + for i in range(5): + await handler._enqueue_log_event( + kwargs={ + "standard_logging_object": { + "messages": [{"role": "user", "content": f"hi-{i}"}], + "response": "hello", + }, + }, + event_type="success", + ) + assert len(handler.log_queue) == 3 + retained = [item["messages"][0]["content"] for item in handler.log_queue] + assert retained == ["hi-2", "hi-3", "hi-4"] + async def test_log_batch_failure_preserves_events_added_during_send(self, handler): """Failure must preserve both the snapshot AND events appended mid-flush.""" handler.log_queue = [{"msg": "a"}, {"msg": "b"}]