fix(veria): cap Rubrik retry queue at 10k events with drop-oldest

A persistent Rubrik webhook outage previously let authenticated traffic
accumulate prompt/response payloads in the in-memory retry queue
without bound. The PR-introduced retry-on-failure behavior in
flush_queue() never trims the queue, so under sustained outage and
high request volume the proxy can run out of memory.

Cap the queue at RUBRIK_MAX_QUEUE_SIZE events (default 10_000) and
drop the oldest events when the cap is exceeded. Emit a throttled
verbose_logger warning so operators can detect a stuck webhook.
This commit is contained in:
mateo-berri 2026-05-21 01:04:49 +00:00
parent e7bc7772aa
commit 198d6cc960
No known key found for this signature in database
2 changed files with 66 additions and 0 deletions

View file

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

View file

@ -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"}]