diff --git a/litellm/litellm_core_utils/litellm_logging.py b/litellm/litellm_core_utils/litellm_logging.py index a6391a2ae27..833e2be225b 100644 --- a/litellm/litellm_core_utils/litellm_logging.py +++ b/litellm/litellm_core_utils/litellm_logging.py @@ -2265,13 +2265,9 @@ class Logging(LiteLLMLoggingBaseClass): self, event_type: Literal["async_success", "sync_success", "async_failure", "sync_failure"], ) -> None: - if self.stream is not None and self.stream is True: - """ - Ignore check on stream, as there can be multiple chunks - """ + if self.stream is True and event_type in ("async_success", "sync_success"): return self.model_call_details[f"has_logged_{event_type}"] = True - return def should_run_callback(self, callback: litellm.CALLBACK_TYPES, litellm_params: dict, event_hook: str) -> bool: if litellm.global_disable_no_log_param: diff --git a/tests/e2e/coverage_registry/logging.yaml b/tests/e2e/coverage_registry/logging.yaml index 7c83e4d3aea..4db55a5eb9b 100644 --- a/tests/e2e/coverage_registry/logging.yaml +++ b/tests/e2e/coverage_registry/logging.yaml @@ -5,6 +5,7 @@ - {id: logging.datadog.success.exports_metric, module: logging, tier: P0, event: success, assertions: [exports_metric], exercised_on: [chat_completions, messages, responses, embeddings], source: "integrations/datadog/datadog.py", rationale: "Powers dashboards/alerts; cardinality regressions common"} - {id: logging.datadog.stream.exports_metric, module: logging, tier: P0, event: stream, assertions: [exports_metric], exercised_on: [chat_completions, messages, responses], source: "integrations/datadog/datadog.py", rationale: "Streaming aggregates usage after the last chunk; delivery and cost must survive that path"} - {id: logging.datadog.failure.exports_metric, module: logging, tier: P0, event: failure, assertions: [exports_metric], exercised_on: [chat_completions], source: "integrations/datadog/datadog.py", rationale: "Failure metrics for alerting/SLO"} +- {id: logging.datadog.stream_failure.exports_metric, module: logging, tier: P0, event: failure, assertions: [exports_metric], exercised_on: [chat_completions], source: "litellm_core_utils/litellm_logging.py", rationale: "Streamed failures re-enter the failure handler once per retry on the same logging object; dedup must hold on the stream path too (#42988)"} - {id: logging.prometheus.success.exports_metric, module: logging, tier: P0, event: success, assertions: [exports_metric], exercised_on: [chat_completions, messages, embeddings], source: "integrations/prometheus.py", rationale: "Standard OSS metrics; per-key cardinality (existing e2e)"} - {id: logging.prometheus.success.records_queue_time, module: logging, tier: P1, event: success, assertions: [records_queue_time], exercised_on: [chat_completions], source: "integrations/prometheus.py / LIT-2034", fail_before_fix: proven, rationale: "Queue time feeds saturation alerting; the family stayed registered while no observation was ever recorded, so presence alone is not the contract"} - {id: logging.otel.success.exports_metric, module: logging, tier: P0, event: success, assertions: [exports_metric], exercised_on: [chat_completions, messages, responses, embeddings], source: "integrations/otel/logger.py", rationale: "OTEL spans on every call path"} diff --git a/tests/e2e/logging/test_datadog_log_e2e.py b/tests/e2e/logging/test_datadog_log_e2e.py index a4821ed058b..7790b8b3243 100644 --- a/tests/e2e/logging/test_datadog_log_e2e.py +++ b/tests/e2e/logging/test_datadog_log_e2e.py @@ -398,3 +398,50 @@ class TestDataDogFailureDelivery: assert payload.error_str is not None and "AnthropicException" in payload.error_str, ( f"the event must carry the provider error, got error_str={payload.error_str!r}" ) + + @pytest.mark.covers("logging.datadog.stream_failure.exports_metric", exercised_on=["chat_completions"]) + def test_failed_chat_completions_stream_emits_one_error_event( + self, client: LoggingClient, dd_logs: DdLogsReader, resources: ResourceManager + ) -> None: + """A STREAMED /chat/completions call that fails at the provider after + its configured retries must still reach the DataDog logs intake as + exactly one error-grade event: the retry loop invokes the failure + handler once per attempt on the same logging object, so a dedup that + only works for non-streaming calls multiplies every retried stream + failure by its attempt count (issue #42988).""" + _assert_datadog_configured(client) + + model_name = f"dd-err-stream-{unique_marker()}" + model_id = client.create_model( + model_name, + LiteLLMParamsBody(model="anthropic/claude-haiku-4-5", api_key=INVALID_UPSTREAM_API_KEY, num_retries=2), + ) + resources.defer(lambda: client.delete_model(model_id)) + key = client.key_with_alias(f"dd-err-stream-key-{unique_marker()}", models=[model_name]) + resources.defer(lambda: client.delete_key(key)) + + deadline = time.monotonic() + client.proxy.poll_timeout + while True: + outcome = client.chat_raw(key, model_name, "trigger an upstream auth failure", stream=True, max_tokens=16) + assert not outcome.ok, "the call must fail; the deployment's upstream key is invalid" + assert outcome.status_code != -1, ( + "network failure between the test and the proxy while provoking the provider " + "failure; retrying now could double-log the failure payload and falsely trip " + f"the exactly-one assertion - fix the rig connectivity first: {outcome.body[:200]}" + ) + if "AnthropicException" in outcome.body or time.monotonic() >= deadline: + break + time.sleep(client.proxy.poll_interval) + assert "AnthropicException" in outcome.body, ( + "never saw the upstream provider failure before the deadline; the key may still be " + f"propagating - last outcome {outcome.status_code}: {outcome.body[:200]}" + ) + assert outcome.status_code == 401, ( + f"an upstream auth failure must map to 401, got {outcome.status_code}: {outcome.body[:200]}" + ) + + events = dd_logs.poll_events_for_query(f"@model_group:{model_name}") + payload = _assert_exactly_one_failure_event(events, model_group=model_name) + assert payload.error_str is not None and "AnthropicException" in payload.error_str, ( + f"the event must carry the provider error, got error_str={payload.error_str!r}" + ) diff --git a/tests/e2e/models.py b/tests/e2e/models.py index c96f4b0bef1..84dc126acb4 100644 --- a/tests/e2e/models.py +++ b/tests/e2e/models.py @@ -1248,6 +1248,7 @@ class LiteLLMParamsBody(BaseModel): tpm: int | None = None weight: int | None = None order: int | None = None + num_retries: int | None = None ModelMode = Literal["batch", "realtime", "image_generation"] diff --git a/tests/test_litellm/litellm_core_utils/test_litellm_logging.py b/tests/test_litellm/litellm_core_utils/test_litellm_logging.py index 23c01841b1b..c8a3f715837 100644 --- a/tests/test_litellm/litellm_core_utils/test_litellm_logging.py +++ b/tests/test_litellm/litellm_core_utils/test_litellm_logging.py @@ -6,6 +6,7 @@ import logging import os import sys import time +import traceback from collections.abc import Callable, Iterator, Mapping from types import MappingProxyType from typing import Final, Literal @@ -1616,6 +1617,57 @@ def test_logging_prevent_double_logging(logging_obj): assert logging_obj.should_run_logging(event_type="async_failure") == True +def test_has_run_logging_stream_marks_failure_events(logging_obj): + """ + On streaming requests the stream early-return must only skip the success + dedup markers (multiple chunks each log success); failure events must still + be marked so a retried stream failure is logged once. + """ + logging_obj.stream = True + logging_obj.has_run_logging(event_type="async_failure") + logging_obj.has_run_logging(event_type="sync_failure") + assert logging_obj.should_run_logging(event_type="async_failure") == False + assert logging_obj.should_run_logging(event_type="sync_failure") == False + assert logging_obj.should_run_logging(event_type="async_success") == True + assert logging_obj.should_run_logging(event_type="sync_success") == True + + +@pytest.mark.asyncio +async def test_async_failure_handler_stream_runs_failure_callbacks_once(): + """ + A streamed call that retries re-enters async_failure_handler once per + attempt on the same logging object; only the first may run the failure + callbacks (issue #42988). + """ + + class CountingLogger(CustomLogger): + def __init__(self): + self.count = 0 + + async def async_log_failure_event(self, kwargs, response_obj, start_time, end_time): + self.count += 1 + + spy = CountingLogger() + logging_obj = LitellmLogging( + model="anthropic/claude-haiku-4-5", + messages=[{"role": "user", "content": "hi"}], + stream=True, + call_type="acompletion", + start_time=time.time(), + litellm_call_id="stream-failure-dedup", + function_id="stream-failure-dedup", + dynamic_async_failure_callbacks=[spy], + kwargs={"litellm_params": {"metadata": {}}}, + ) + try: + raise ValueError("upstream auth failure") + except ValueError as e: + exception, tb = e, traceback.format_exc() + for _ in range(3): + await logging_obj.async_failure_handler(exception, tb) + assert spy.count == 1 + + @pytest.mark.asyncio async def test_datadog_logger_not_shadowed_by_llm_obs(monkeypatch): """Ensure DataDog logger instantiates even when LLM Obs logger already cached."""