diff --git a/litellm/integrations/datadog/datadog_llm_obs.py b/litellm/integrations/datadog/datadog_llm_obs.py index 8f13b5602b0..2cda5e3f088 100644 --- a/litellm/integrations/datadog/datadog_llm_obs.py +++ b/litellm/integrations/datadog/datadog_llm_obs.py @@ -148,7 +148,7 @@ class DataDogLLMObsLogger(CustomBatchLogger): self.log_queue.append(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"DataDogLLMObs: Error logging success event - {str(e)}") @@ -160,7 +160,7 @@ class DataDogLLMObsLogger(CustomBatchLogger): self.log_queue.append(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"DataDogLLMObs: Error logging failure event - {str(e)}") @@ -186,6 +186,14 @@ class DataDogLLMObsLogger(CustomBatchLogger): undelivered = await self._send_with_413_split(batch_to_send) if undelivered: self.log_queue = undelivered + self.log_queue + overflow = len(self.log_queue) - self.max_queue_size + if overflow > 0: + del self.log_queue[:overflow] + verbose_logger.warning( + "DataDogLLMObs: log queue exceeded max_queue_size=%s; dropped %s oldest spans", + self.max_queue_size, + overflow, + ) if self.is_mock_mode: verbose_logger.debug(f"[DATADOG MOCK] Batch of {len(batch_to_send)} events successfully mocked") @@ -253,7 +261,7 @@ class DataDogLLMObsLogger(CustomBatchLogger): pending.append(chunk[:mid]) continue - if response.status_code != 202: + if response.status_code >= 300: verbose_logger.error( "DataDogLLMObs: unexpected response status_code=%s, text=%s", response.status_code, diff --git a/tests/test_litellm/integrations/datadog/test_datadog_llm_observability.py b/tests/test_litellm/integrations/datadog/test_datadog_llm_observability.py index 2fd25efa447..c53031955c8 100644 --- a/tests/test_litellm/integrations/datadog/test_datadog_llm_observability.py +++ b/tests/test_litellm/integrations/datadog/test_datadog_llm_observability.py @@ -1427,3 +1427,76 @@ class TestDataDogLLMObsBatchSplit: assert delivered == ["0", "1"] assert [s["span_id"] for s in logger.log_queue] == ["2", "3"] + + @pytest.mark.asyncio + async def test_requeue_drops_oldest_beyond_max_queue_size(self): + """Re-queued spans must respect max_queue_size, dropping oldest first.""" + logger = self._make_logger() + logger.max_queue_size = 2 + logger.log_queue = [self._make_span(span_id=str(i)) for i in range(3)] + + async def _mock_post_spans(spans): + raise Exception("Connection refused") + + with patch.object(logger, "_post_spans", side_effect=_mock_post_spans): + await logger.async_send_batch() + + assert [s["span_id"] for s in logger.log_queue] == ["1", "2"] + + @pytest.mark.asyncio + async def test_returned_non_2xx_response_requeues(self): + """A client that returns a non-2xx response instead of raising must + also re-queue the undelivered spans.""" + logger = self._make_logger() + logger.log_queue = [self._make_span(span_id=str(i)) for i in range(2)] + + async def _mock_post_spans(spans): + resp = MagicMock() + resp.status_code = 500 + resp.text = "Internal Server Error" + return resp + + with patch.object(logger, "_post_spans", side_effect=_mock_post_spans): + await logger.async_send_batch() + + assert [s["span_id"] for s in logger.log_queue] == ["0", "1"] + + @pytest.mark.asyncio + async def test_returned_2xx_non_202_counts_as_delivered(self): + """A 2xx response other than 202 is a delivery, never a re-queue.""" + logger = self._make_logger() + logger.log_queue = [self._make_span(span_id=str(i)) for i in range(2)] + + async def _mock_post_spans(spans): + resp = MagicMock() + resp.status_code = 200 + resp.text = "OK" + return resp + + with patch.object(logger, "_post_spans", side_effect=_mock_post_spans): + await logger.async_send_batch() + + assert logger.log_queue == [] + + @pytest.mark.asyncio + async def test_batch_size_trigger_flushes_under_lock(self): + """The batch-size trigger goes through flush_queue so it takes the + flush lock and updates last_flush_time on a full drain.""" + logger = self._make_logger() + logger.batch_size = 2 + logger.last_flush_time = 0.0 + + async def _mock_post_spans(spans): + resp = MagicMock() + resp.status_code = 202 + resp.text = "OK" + return resp + + with patch.object(logger, "_post_spans", side_effect=_mock_post_spans), patch.object( + logger, "create_llm_obs_payload", side_effect=lambda *a, **k: self._make_span() + ): + await logger.async_log_success_event({}, None, None, None) + await logger.async_log_success_event({}, None, None, None) + + assert logger.log_queue == [] + assert logger.last_flush_time > 0.0