mirror of
https://github.com/BerriAI/litellm.git
synced 2026-09-29 01:42:19 +00:00
fix(logging): dedupe failure callbacks on streaming requests
Logging.has_run_logging returned early for every event type when self.stream was True, so has_logged_async_failure / has_logged_sync_failure were never set on streaming requests. Retried streamed calls re-enter async_failure_handler once per attempt on the same logging object (num_retries, router retries), and each attempt fired the failure callbacks again, duplicating events to sinks like s3_v2 and Datadog. Restrict the streaming early-return to the success event types, which genuinely re-enter once per chunk, and keep the failure markers working the same as non-streaming calls. Resolves LIT-8542 Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
parent
0fa1fe2059
commit
25f257cc51
5 changed files with 102 additions and 5 deletions
|
|
@ -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:
|
||||
|
|
|
|||
|
|
@ -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"}
|
||||
|
|
|
|||
|
|
@ -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}"
|
||||
)
|
||||
|
|
|
|||
|
|
@ -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"]
|
||||
|
|
|
|||
|
|
@ -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."""
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue