Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
30 changes: 25 additions & 5 deletions litellm/integrations/datadog/datadog.py
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@
import asyncio
import datetime
import os
import time
import traceback
from datetime import datetime as datetimeObj
from typing import Any, Dict, List, Optional, Union
Expand Down Expand Up @@ -301,7 +302,7 @@ async def async_post_call_failure_hook(
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()}"
Expand All @@ -324,9 +325,12 @@ async def async_send_batch(self):
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,
)

Expand All @@ -335,9 +339,10 @@ async def async_send_batch(self):
"[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:
Comment thread
emerzon marked this conversation as resolved.
verbose_logger.exception(DD_ERRORS.DATADOG_413_ERROR.value)
self.log_queue = batch_to_send + self.log_queue
return
Comment thread
emerzon marked this conversation as resolved.

response.raise_for_status()
Expand All @@ -348,19 +353,34 @@ async def async_send_batch(self):

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(
"Datadog: Response from datadog API status_code: %s, text: %s",
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
Expand Down Expand Up @@ -429,7 +449,7 @@ async def _log_async_event(self, kwargs, response_obj, start_time, end_time):
)

if len(self.log_queue) >= self.batch_size:
await self.async_send_batch()
await self.flush_queue()

def _create_datadog_logging_payload_helper(
self,
Expand Down
Original file line number Diff line number Diff line change
@@ -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()
Loading