fix: drain datadog batches safely (#25663)

* fix: drain datadog batches safely

* fix: preserve datadog batches on 413

* fix: import time in datadog flush queue

* test: cover datadog batching edge cases

* fix: only stamp successful datadog flushes

* test: use sync mock for datadog payload builder
This commit is contained in:
Emerson Gomes 2026-04-13 21:34:58 -05:00 • committed by Sameer Kankute
parent e724e5e07d
commit a302b53980
No known key found for this signature in database
2 changed files with 292 additions and 5 deletions

View file

@ -16,6 +16,7 @@ For batching specific details see CustomBatchLogger class
import asyncio
import datetime
import os
import time
import traceback
from datetime import datetime as datetimeObj
from typing import Any, Dict, List, Optional, Union
@ -301,7 +302,7 @@ class DataDogLogger(
self.log_queue.append(dd_payload)
if len(self.log_queue) >= self.batch_size:
await self.async_send_batch()
await self.flush_queue()
except Exception as e:
verbose_logger.exception(
f"Datadog: async_post_call_failure_hook - {str(e)}\n{traceback.format_exc()}"
@ -324,9 +325,12 @@ class DataDogLogger(
verbose_logger.exception("Datadog: log_queue does not exist")
return
batch_to_send = self.log_queue[:]
self.log_queue = []
verbose_logger.debug(
"Datadog - about to flush %s events on %s",
len(self.log_queue),
len(batch_to_send),
self.intake_url,
)
@ -335,9 +339,10 @@ class DataDogLogger(
"[DATADOG MOCK] Mock mode enabled - API calls will be intercepted"
)
response = await self.async_send_compressed_data(self.log_queue)
response = await self.async_send_compressed_data(batch_to_send)
if response.status_code == 413:
verbose_logger.exception(DD_ERRORS.DATADOG_413_ERROR.value)
self.log_queue = batch_to_send + self.log_queue
return
response.raise_for_status()
@ -348,7 +353,7 @@ class DataDogLogger(
if self.is_mock_mode:
verbose_logger.debug(
f"[DATADOG MOCK] Batch of {len(self.log_queue)} events successfully mocked"
f"[DATADOG MOCK] Batch of {len(batch_to_send)} events successfully mocked"
)
else:
verbose_logger.debug(
@ -356,11 +361,26 @@ class DataDogLogger(
response.status_code,
response.text,
)
except Exception as e:
self.log_queue = batch_to_send + self.log_queue
verbose_logger.exception(
f"Datadog Error sending batch API - {str(e)}\n{traceback.format_exc()}"
)
async def flush_queue(self):
if self.flush_lock is None:
return
async with self.flush_lock:
if self.log_queue:
verbose_logger.debug(
"Datadog: Flushing batch of %s events", len(self.log_queue)
)
await self.async_send_batch()
if not self.log_queue:
self.last_flush_time = time.time()
def log_success_event(self, kwargs, response_obj, start_time, end_time):
"""
Sync Log success events to Datadog
@ -429,7 +449,7 @@ class DataDogLogger(
)
if len(self.log_queue) >= self.batch_size:
await self.async_send_batch()
await self.flush_queue()
def _create_datadog_logging_payload_helper(
self,

View file

@ -0,0 +1,267 @@
from unittest.mock import AsyncMock, Mock, patch
import pytest
from httpx import Request, Response
from litellm.integrations.datadog.datadog import DataDogLogger
from litellm.types.integrations.datadog import DatadogPayload
@pytest.fixture
def datadog_env(monkeypatch):
monkeypatch.setenv("DD_API_KEY", "test_api_key")
monkeypatch.setenv("DD_SITE", "test.datadoghq.com")
@pytest.mark.asyncio
async def test_async_send_batch_keeps_events_appended_during_send(datadog_env):
with patch("asyncio.create_task"):
logger = DataDogLogger()
logger.log_queue = [
DatadogPayload(
ddsource="litellm",
ddtags="env:test",
hostname="host",
message=f'{{"event": {i}}}',
service="svc",
status="info",
)
for i in range(2)
]
async def _mock_send(data):
logger.log_queue.append(
DatadogPayload(
ddsource="litellm",
ddtags="env:test",
hostname="host",
message='{"event": 2}',
service="svc",
status="info",
)
)
return Response(
202, request=Request("POST", "https://example.com"), text="Accepted"
)
logger.async_send_compressed_data = AsyncMock(side_effect=_mock_send)
await logger.async_send_batch()
assert logger.async_send_compressed_data.await_count == 1
sent_batch = logger.async_send_compressed_data.await_args.args[0]
assert len(sent_batch) == 2
assert len(logger.log_queue) == 1
assert logger.log_queue[0]["message"] == '{"event": 2}'
@pytest.mark.asyncio
async def test_failure_hook_threshold_flush_uses_flush_queue(datadog_env):
with patch("asyncio.create_task"):
logger = DataDogLogger()
logger.batch_size = 1
logger.flush_queue = AsyncMock()
await logger.async_post_call_failure_hook(
request_data={},
original_exception=Exception("boom"),
user_api_key_dict=type("UserKey", (), {})(),
traceback_str="trace",
)
logger.flush_queue.assert_awaited_once()
@pytest.mark.asyncio
async def test_async_send_batch_requeues_events_on_413(datadog_env):
with patch("asyncio.create_task"):
logger = DataDogLogger()
logger.log_queue = [
DatadogPayload(
ddsource="litellm",
ddtags="env:test",
hostname="host",
message=f'{{"event": {i}}}',
service="svc",
status="info",
)
for i in range(2)
]
logger.async_send_compressed_data = AsyncMock(
return_value=Response(
413,
request=Request("POST", "https://example.com"),
text="Payload Too Large",
)
)
await logger.async_send_batch()
assert logger.async_send_compressed_data.await_count == 1
assert len(logger.log_queue) == 2
assert [event["message"] for event in logger.log_queue] == [
'{"event": 0}',
'{"event": 1}',
]
@pytest.mark.asyncio
async def test_async_send_batch_handles_empty_queue(datadog_env):
with patch("asyncio.create_task"):
logger = DataDogLogger()
logger.log_queue = []
logger.async_send_compressed_data = AsyncMock()
await logger.async_send_batch()
logger.async_send_compressed_data.assert_not_awaited()
@pytest.mark.asyncio
async def test_async_send_batch_requeues_events_on_exception(datadog_env):
with patch("asyncio.create_task"):
logger = DataDogLogger()
logger.log_queue = [
DatadogPayload(
ddsource="litellm",
ddtags="env:test",
hostname="host",
message=f'{{"event": {i}}}',
service="svc",
status="info",
)
for i in range(2)
]
logger.async_send_compressed_data = AsyncMock(side_effect=RuntimeError("boom"))
await logger.async_send_batch()
assert [event["message"] for event in logger.log_queue] == [
'{"event": 0}',
'{"event": 1}',
]
@pytest.mark.asyncio
async def test_log_async_event_threshold_flush_uses_flush_queue(datadog_env):
with patch("asyncio.create_task"):
logger = DataDogLogger()
logger.batch_size = 1
logger.flush_queue = AsyncMock()
logger.create_datadog_logging_payload = Mock(
return_value=DatadogPayload(
ddsource="litellm",
ddtags="env:test",
hostname="host",
message='{"event": 0}',
service="svc",
status="info",
)
)
await logger._log_async_event(
kwargs={},
response_obj={},
start_time=None,
end_time=None,
)
logger.flush_queue.assert_awaited_once()
@pytest.mark.asyncio
async def test_flush_queue_updates_last_flush_time(datadog_env):
with patch("asyncio.create_task"):
logger = DataDogLogger()
logger.log_queue = [
DatadogPayload(
ddsource="litellm",
ddtags="env:test",
hostname="host",
message='{"event": 0}',
service="svc",
status="info",
)
]
logger.last_flush_time = 0
async def _successful_send():
logger.log_queue = []
logger.async_send_batch = AsyncMock(side_effect=_successful_send)
await logger.flush_queue()
logger.async_send_batch.assert_awaited_once()
assert logger.last_flush_time > 0
@pytest.mark.asyncio
async def test_flush_queue_does_not_update_last_flush_time_when_send_requeues(
datadog_env,
):
with patch("asyncio.create_task"):
logger = DataDogLogger()
logger.log_queue = [
DatadogPayload(
ddsource="litellm",
ddtags="env:test",
hostname="host",
message='{"event": 0}',
service="svc",
status="info",
)
]
logger.last_flush_time = 123.0
async def _requeue_batch():
logger.log_queue = [
DatadogPayload(
ddsource="litellm",
ddtags="env:test",
hostname="host",
message='{"event": 0}',
service="svc",
status="info",
)
]
logger.async_send_batch = AsyncMock(side_effect=_requeue_batch)
await logger.flush_queue()
logger.async_send_batch.assert_awaited_once()
assert logger.last_flush_time == 123.0
@pytest.mark.asyncio
async def test_flush_queue_returns_without_lock(datadog_env):
with patch("asyncio.create_task"):
logger = DataDogLogger()
logger.flush_lock = None
logger.log_queue = [
DatadogPayload(
ddsource="litellm",
ddtags="env:test",
hostname="host",
message='{"event": 0}',
service="svc",
status="info",
)
]
logger.async_send_batch = AsyncMock()
await logger.flush_queue()
logger.async_send_batch.assert_not_awaited()