litellm/litellm/_logging.py
Deepanshu Lulla 9ce96c2d34
feat(logging): add opt-in session_id and trace_id correlation to JSON log records via contextvars (#34418)
* feat(logging): add opt-in session_id/trace_id correlation to JSON log records via contextvars

Adds two ContextVar instances (session_id_var, trace_id_var) to litellm/_logging.py and
two setter functions (set_session_id, set_trace_id). Logging.__init__() now calls both
setters after assigning litellm_trace_id so every JSON log record emitted within the
async request context carries trace_id and, when provided, session_id — enabling log
correlation in Loki, CloudWatch Logs Insights, and other structured-log sinks without
any changes to individual log call sites.

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>

* fix(logging): guard session_id/trace_id injection against overwriting caller-supplied extra fields

* fix(logging): always reset session_id_var to empty string when no session_id provided

* feat: gate request correlation IDs in logs behind request_correlation_in_logs flag

* refactor: move correlation ID injection into CorrelationContextFilter

* feat(logging): extend request_correlation_in_logs to plaintext logs and StandardLoggingPayload

Plaintext log lines (json_logs off) now get the same trace_id/session_id
suffix as JSON logs via a new CorrelationPlainFormatter, so the flag has a
visible effect regardless of log format.

StandardLoggingPayload gets a new independent session_id field, populated
from litellm_session_id. trace_id's existing session_id-first fallback is
preserved when request_correlation_in_logs is off; with the flag on, an
explicit litellm_trace_id now takes priority over litellm_session_id so
the two fields carry genuinely independent values.

* fix(logging): restore correlation context after nested calls; sanitize correlation ids

Addresses two review findings on this PR.

CorrelationContextFilter's trace_id/session_id contextvars were set on every
Logging.__init__ but never reset, so a nested LiteLLM call sharing the same
asyncio Task as an outer request (e.g. a guardrail's own LLM-as-judge call,
an MCP sampling call) would leave the outer request's subsequent log lines
stamped with the nested call's ids instead of its own. set_trace_id/
set_session_id now return their contextvars.Token, and Logging stores them
and resets both once its own success/failure handler actually completes,
via a new idempotent _restore_correlation_context() called from all four
terminal handlers.

set_trace_id/set_session_id also now strip control characters and bound
length before storing a caller-controlled trace_id/session_id, since these
values can originate from request input (litellm_session_id, x-litellm-
trace-id) and get interpolated into plain-text log lines - without this, a
caller could embed \r/\n or escape sequences to forge fake log entries.

* fix(logging): restore correlation context after nested calls, not before

The previous commit called _restore_correlation_context() as the first
line of each terminal handler, before that handler's own callback
dispatch loop runs. That's backwards: a nested LiteLLM call triggered
from within a callback (e.g. a guardrail's own LLM-as-judge call) would
then capture the *already-reset* value as its own pre-call baseline,
and its own reset would restore to that instead of the true outer
value - verified live to still leak.

success_handler/async_success_handler/failure_handler/async_failure_handler
are now thin wrappers: the original bodies move to
_success_handler_body/etc, called inside a try/finally that restores
context only once the full body - including any nested calls its own
callback dispatch triggers - has actually finished, mirroring proper
stack-scoped nesting semantics.

* test(logging): cover async_failure_handler's correlation-context restore

Codecov flagged the new async_failure_handler wrapper (try/finally around
_async_failure_handler_body) as uncovered - the method had no direct test
at all before this PR's refactor split it into a wrapper. Adds a test that
awaits it directly and asserts both that async_log_failure_event still
fires and that _restore_correlation_context() puts the pre-call
trace_id/session_id back.

* fix(logging): restore correlation context by value, not by contextvars.Token

veria-ai correctly flagged that contextvars.Token.reset() only works in the
exact Context it was created in, and litellm's async success path (and
streaming failure path) dispatch async_success_handler/async_failure_handler
via asyncio.create_task and the global logging worker - a different Context
than Logging.__init__ ran in. reset_trace_id/reset_session_id silently
swallowed the resulting ValueError, so the restore was a no-op for exactly
those paths. Verified independently: reproduced the raw contextvars
behavior, then confirmed litellm's async success dispatch really does go
through asyncio.create_task + GLOBAL_LOGGING_WORKER (litellm/utils.py).

Logging now captures the pre-call *value* (not a Token) and restores via a
plain set_trace_id()/set_session_id() call, which works regardless of which
Task/Context calls it. reset_trace_id/reset_session_id are removed as
dead/unreliable code. Added a regression test that spawns __init__ and the
restore in different asyncio Tasks - confirmed it fails against the prior
Token-based commit and passes here.

* fix(logging): restore correlation context in the originating task too

Greptile's re-review correctly identified a remaining gap: for a
successful acompletion(), async_success_handler is dispatched via
asyncio.create_task + the global logging worker into a *different* Task
than the one wrapper_async/Logging.__init__ ran in. The prior fix (43c164a)
only restored the handler's own (detached, throwaway) Task - it never
touched the originating request Task, which keeps this call's trace_id/
session_id set for the rest of its own execution (e.g. nested calls made
via the same Task).

wrapper()/wrapper_async() in litellm/utils.py now restore the originating
Task's correlation context in a finally block once the whole call is done,
regardless of what detached logging tasks it spawned along the way. Since
the wrapped body rebinds its own `kwargs` local via function_setup(),
sharing the dict object doesn't work here; a small mutable holder carries
the constructed Logging instance back out to the outer wrapper instead.

_restore_correlation_context() is no longer guarded against repeat calls:
with value-based (not Token-based) restoration, each distinct Task that
calls it needs its own restore to take effect in that Task's own view of
the contextvars, so multiple calls (once per Task involved in an attempt)
are required, not just tolerated.

Added a regression test using mock_response to exercise the real success
dispatch path (asyncio.create_task + GLOBAL_LOGGING_WORKER) without a live
provider call, asserting the *test's own* (originating) task context is
restored after the call - this is exactly the case Greptile flagged and
the prior commit didn't cover.

* fix(logging): restore correlation context when function_setup itself fails

Greptile's 4th finding: if function_setup() constructs Logging() (whose
__init__ already mutates trace_id_var/session_id_var) and then raises
before returning - e.g. update_environment_variables() throws - the
caller's wrapper()/wrapper_async() never receives a logging_obj reference,
so its own restore-on-finally never fires. The correlation ids leak into
every subsequent log line on that thread/task with no way to clear them.

function_setup()'s own except block now restores the context itself in
that case, using whatever logging_obj it managed to construct before
failing (locals().get(), safe against the earlier failure modes where
logging_obj was never assigned at all).

Added a regression test that monkeypatches Logging.update_environment_variables
to raise after construction, confirmed it fails without this fix (the
leaked ids show up directly in the raised exception's own log line) and
passes with it. Broader sweep (test_utils.py, test_router.py,
test_main_module_header.py, streaming handler tests, plus all
logging-specific tests): 722 passed.

* fix(logging): don't assume every litellm_logging_obj is a real Logging instance

CI caught a real regression from the last commit: tests/test_litellm/llms/xai/test_xai_key_fallback.py
injects a minimal FakeLogging stand-in (only implementing
update_from_kwargs) as litellm_logging_obj for a narrow realtime-config
unit test, bypassing the real Logging class entirely. wrapper()/
wrapper_async()'s finally block and function_setup()'s except block both
unconditionally called _restore_correlation_context() on whatever ended up
in the holder, which doesn't exist on that stand-in.

_restore_correlation_context is new plumbing specific to this PR's
feature, not part of any pre-existing stand-in's expected interface, so
callers of it can't assume every object playing the litellm_logging_obj
role implements it. Added _restore_correlation_context_if_supported(),
a small getattr-guarded helper, and used it at all three call sites.

* fix(logging): don't restore context too early on setup failure or streaming

Two more findings from Greptile's 5th review round.

1. function_setup()'s except block restored correlation context *after*
   logging the "Error in function_setup" exception, so that diagnostic log
   line itself was stamped with the doomed call's ids instead of the outer
   ids - misleading, since the failed call never produces anything else to
   attribute those ids to. Restore now happens before the log call.

2. wrapper()/wrapper_async() restored the originating task's context as
   soon as a streaming call returned, before the caller ever starts
   iterating the CustomStreamWrapper it just got back. Any log lines
   emitted while iterating (in the same thread/task) incorrectly showed
   the pre-call ids instead of this call's own ones. The wrapper finally
   block now skips the restore when the return value is a stream wrapper,
   deferring to the terminal handler that already fires once the stream is
   actually assembled/exhausted.

Both verified with tests that fail against the prior commit and pass
against this one. Broader sweep unchanged at 829 passing.

* fix(logging): best-effort correlation cleanup on abandoned streams

Greptile's 7th finding: if a caller returns a streaming response and never
fully consumes it - stops iterating early, drops the reference, cancels
it - the terminal handler that normally restores the originating task's
trace_id/session_id never fires, since it only runs once the stream is
actually assembled/exhausted. The ids leak into every subsequent log line
in that thread/task with no bound.

There's no reliable Python hook for "this was abandoned without being
closed" - CustomStreamWrapper has no close()/__aexit__/context-manager
convention today, and the only automatic option is __del__, whose timing
is inherently unpredictable (delayed by cyclic GC, not guaranteed at
interpreter shutdown, can run on a different thread). This is a best-effort
safety net, not a guarantee, and is documented as such in the docstring.

Testing this via real garbage collection proved unreliable in practice:
per-chunk logging submits work to a thread pool executor whose worker
thread transiently holds its own bound-method reference to the wrapper
until that task completes, so refcount doesn't hit zero on a
deterministic schedule even with polling. Tests call __del__ directly
instead - a plain method, safe to invoke early - which exercises exactly
the restore logic real garbage collection would eventually trigger,
plus a case confirming a broken logging_obj can never make __del__ raise.

* fix(logging): restore consumer's context at every real stream exit point

Two more findings from this round.

Veria AI: even a *fully consumed* stream never restored the actual
consuming thread/task's correlation context. The terminal success dispatch
(dispatch_success_handlers via asyncio.create_task for async, or
success_handler via the shared executor for sync) only restores whatever
detached context it runs in - never the caller's own thread/task that's
running the for/async for loop. Same root cause as the wrapper-level fix
two rounds ago, just missed for the streaming-completion path.

Greptile: explicit aclose() (client disconnect, router fallback aborting
a partial stream) closed the underlying stream without restoring
correlation context either, since request wrappers intentionally skip
restoration for returned streams and no terminal handler runs on this
path.

Added CustomStreamWrapper._restore_consumer_correlation_context(), called
from every point control genuinely returns to the consumer: the final
raise StopIteration/StopAsyncIteration on natural exhaustion (both sync
branches, both async branches), _handle_stream_fallback_error (the shared
choke point for all three failure-raising call sites), and aclose(). __del__
now delegates to the same helper instead of duplicating it.

Verified with tests extending the existing streaming-exhaustion cases to
assert the consuming context is restored after the loop completes (fails
against the prior commit, passes now), plus a dedicated aclose() test.
Broader sweep: 832 passing.

* fix(logging): don't let a delayed __del__ finalizer clobber a newer active call

If an abandoned stream's __del__ fires late (after cyclic GC delay), a
different call may have already taken over the correlation contextvars in
the same Task/thread. Restoring unconditionally would stomp that active
call's trace_id/session_id with the abandoned stream's stale pre-call
snapshot. __del__ now only restores when the contextvars still hold the
ids this call itself set.

* fix(logging): compare sanitized ids in the __del__ ownership guard

set_trace_id()/set_session_id() sanitize (strip control chars, bound length)
before storing, so the contextvar's value can differ from the raw
litellm_trace_id/litellm_session_id. The __del__ ownership guard was
comparing against the raw values, so a caller-supplied id containing control
characters or exceeding 256 chars would never match, permanently skipping
cleanup. Capture what set_trace_id()/set_session_id() actually stored and
compare against that instead.

* fix(logging): restore consumer context on the synthesized finish_reason chunk

Both __next__ and _finalize_completed_stream() have a branch that fires when
the underlying stream ends without ever emitting an explicit finish_reason
chunk: they synthesize one via finish_reason_handler() and return it. A
consumer that stops as soon as it sees finish_reason - a common pattern -
never calls __next__()/__anext__() again, so the existing restore in the
sent_last_chunk-is-True StopIteration branch never runs for them. The
underlying stream is already exhausted at this point regardless of whether
the caller keeps iterating, so restoring here is safe.

* fix(logging): don't restore correlation context before the caller receives the final chunk

The previous fix (5147c69186) restored context immediately before returning
the synthesized finish_reason chunk from __next__/_finalize_completed_stream,
reasoning that completion_stream was already exhausted. But that chunk is
still this call's own data, and the caller's own application-level log
statements processing it run in the same synchronous frame right after the
return - restoring first made those lines carry the wrong (outer) ids,
exactly what wrapper()/wrapper_async() deliberately avoid by not restoring
while a stream is being iterated.

Revert to not restoring there. A caller that keeps iterating still gets a
correct, deterministic restore on its very next __next__()/__anext__() call
(completion_stream is exhausted, so that immediately re-raises
StopIteration/StopAsyncIteration through the already-restoring branch). A
caller that stops right after finish_reason relies on aclose() or the
best-effort __del__ guard, same as any other stream the caller doesn't fully
exhaust.

* refactor(logging): hoist a safely-hoistable function-body import to module top

CorrelationContextFilter.filter()'s `import litellm` was a function-body
import; verified it can move to module top without a circular-import failure
(litellm/__init__.py already imports from litellm._logging before setting
request_correlation_in_logs, but a bare `import litellm` only binds the
already-in-sys.modules module object - the attribute itself isn't read until
filter() actually runs, by which point litellm is fully initialized).

* test(logging): move correlation tests into their conventionally-mapped files

tests/test_litellm/ mirrors litellm/ in a parallel path. Correlation tests
for the Logging class (litellm_logging.py), function_setup/wrapper_async
(utils.py), and CustomStreamWrapper (streaming_handler.py) had all landed in
test_logging.py, which only maps to litellm/_logging.py itself. Moving each
group to its correctly-mapped file: test_litellm_logging.py (Logging class
init/restore), test_utils.py (function_setup, wrapper_async), and
test_streaming_handler.py (CustomStreamWrapper) in the next commit.
test_logging.py keeps only what actually exercises _logging.py's own
contextvars/filters/formatters/sanitization. No behavior change - same
assertions, same coverage, just relocated.

* fix(logging): restore correlation context unconditionally in wrapper()'s sync path

Blocking finding from review: a caller-visible correlation feature was
silently misattributing one request's logs to a different, unrelated one on
the sync/threaded path. wrapper()/wrapper_async() both left trace_id/session_id
"open" across a stream's entire iteration so the caller's own log lines while
consuming it would carry the right ids. That's safe for wrapper_async(): each
async call gets its own asyncio Task with its own copy of the contextvars,
and Tasks are never recycled across requests, so a leftover value can only
ever affect that one already-abandoned Task.

It is not safe for wrapper() (sync): a plain OS thread has no such per-call
isolation, and a thread pool's worker threads *are* recycled across unrelated
requests. If a sync stream was abandoned (client disconnect, early break, an
uncaught exception) without ever being exhausted or closed, nothing restored
its contextvars, and a pool could later hand that same thread to a completely
different call, which would inherit the abandoned request's ids as its own
"pre-call" baseline and then restore back to that poison when it finished -
permanently misattributing every subsequent log line on that thread,
including its own, to the abandoned request. Strengthening the __del__
finalizer already added for this can't fix it: finalizer timing is exactly
what a permanently-reused thread can't rely on.

wrapper() now restores unconditionally in its own finally, before a sync
stream is ever handed back to the caller. The trade-off: a sync stream
consumer's own application-level log statements while iterating no longer
automatically carry this call's ids (litellm's own internal per-chunk
logging is unaffected, since it's dispatched separately). That's an
acceptable cost for eliminating a silent cross-request misattribution bug.
wrapper_async() keeps the existing conditional (skip-if-streaming) behavior,
justified by the Task-isolation argument above; CustomStreamWrapper's
__del__/aclose()/next-iteration restore machinery remains meaningful and
necessary there.

This also simplifies wrapper()/wrapper_async() back toward their original
shape: both previously used a mutable-dict-holder split into a separate
_body function to smuggle logging_obj/result out to an outer finally,
working around function_setup() rebinding its own local `kwargs`. That
restructuring is no longer needed - `logging_obj` (and, for wrapper_async(),
`result`) were already function-level locals in scope for a plain
try/finally; three of wrapper_async()'s retry-return statements now assign
through `result` first so it accurately reflects what's actually returned
even on a retry path.

Regression test: test_abandoned_sync_stream_does_not_contaminate_a_later_call_on_the_same_thread
in test_streaming_handler.py reproduces the exact reported scenario with a
real single-worker ThreadPoolExecutor - confirmed it fails with the prior
(skip-restore-on-stream) wrapper() and passes with this fix.

* refactor(logging): use Mapping instead of bare dict for read-only params

_get_standard_logging_payload_trace_id/_session_id only read litellm_params
(.get() calls, no mutation) - annotate it as Mapping[str, Any] rather than a
bare mutable dict, per the repo's no-mutable-collection-in-annotation rule.

* fix(logging): scope request_correlation_in_logs to the async/proxy path only

Blocking review finding: wrapper() (the sync entry point) used the same
skip-restore-on-stream design as wrapper_async(), but a plain OS thread has
no per-call context isolation the way an asyncio Task does, and a thread
pool's worker threads are recycled across unrelated requests - an abandoned
sync stream could leave its ids stuck on a thread a pool later hands to a
completely different request, misattributing that request's logs. A fix
existed and was tested (restore unconditionally in wrapper()'s own finally),
but it doesn't benefit this feature's primary consumer - the proxy only ever
calls the async entry point - and carries sync-specific complexity this PR
doesn't need.

Scope the feature to async only instead: Logging.__init__() takes a new
supports_correlation_logging parameter (default True), threaded down from a
new function_setup(..., is_async_call: bool = True) parameter. wrapper() is
the one caller that passes is_async_call=False; every other function_setup()
call site (wrapper_async(), the router, and proxy/MCP-internal call sites)
is already async and keeps the default. With
supports_correlation_logging=False, Logging.__init__() never calls
set_trace_id()/set_session_id() at all, so a sync call has nothing to leak
in the first place. wrapper() reverts to its pre-review shape with no
correlation-specific code at all.

StandardLoggingPayload's own trace_id/session_id fields are unaffected
either way - they're a deterministic per-call read of
self.litellm_trace_id/self.litellm_session_id, not ambient contextvar state,
so they were never exposed to the cross-request bug.

Full sync/direct-SDK support (stamping + its own safe-restore mechanism) is
deferred to a follow-up PR; the fix and its regression test already exist in
this branch's history at commit 9f3a20f4b2 and can be resurrected there.

Tests: replaced the two wrapper()-level tests with ones proving the new
invariant (sync calls, streaming and non-streaming, never touch
trace_id_var/session_id_var even when the caller explicitly passes
litellm_trace_id/litellm_session_id), and added a direct unit test for the
supports_correlation_logging=False gate on Logging.__init__ itself. Verified
live: a real proxy (Postgres-backed, real OpenAI calls) shows clean
trace_id/session_id isolation across two concurrent sessions with no
cross-contamination; a standalone script confirms real sync SDK calls
against a real model never touch the correlation contextvars.

* feat(logging): fall back to W3C traceparent/baggage for trace_id/session_id

request_correlation_in_logs previously only resolved trace_id/session_id from
litellm-specific sources: x-litellm-trace-id/x-litellm-session-id headers, a
generic x-<vendor>-session-id header, or Anthropic-style metadata.user_id. If
none were present, trace_id fell back to an auto-generated UUID unrelated to
anything else, and session_id stayed empty - even when the caller already had
real distributed-tracing instrumentation sending the actual industry-standard
headers for this.

Add a fallback to the W3C Trace Context traceparent header (trace-id
component) and W3C Baggage header (session.id entry), so a request already
carrying real OpenTelemetry trace context correlates litellm's own logs with
the same trace in the caller's observability backend (Datadog, Honeycomb,
Tempo, etc.) instead of getting an unrelated generated id. Precedence is
unchanged for existing sources: explicit litellm headers and the Anthropic
metadata path both still win over this new fallback, which only fires when
neither found anything. trace_id and session_id are resolved independently
here (unlike the existing chain_id mechanism, which uses one shared value for
both), since traceparent and baggage are semantically distinct W3C concepts.

New helpers _trace_id_from_traceparent/_session_id_from_baggage in
litellm_pre_call_utils.py parse the header formats directly (no new
dependency - both are simple fixed-width/delimited strings), wired into
LiteLLMProxyRequestSetup.add_litellm_metadata_from_request_headers() only
when the corresponding litellm_trace_id/litellm_session_id key isn't already
set by the existing paths.

Verified live against a real proxy: a bare traceparent header produces a log
trace_id exactly matching its trace-id component; a traceparent alongside an
explicit x-litellm-trace-id header (different value) produces a log showing
the explicit header's value, proving precedence.

* fix(logging): reserve trace_id/session_id in JsonFormatter against message-content spoofing

JsonFormatter merges keys parsed from the message body before applying extra
record attributes, and the extra-attributes loop skips a key that's already
present. A caller-controlled log message that happens to parse as JSON/dict
with a "trace_id"/"session_id" key (e.g. the proxy logging a raw request-header
dict) could therefore make the JSON record carry the attacker-supplied value
instead of the real correlation context set via CorrelationContextFilter.

trace_id/session_id are now applied from the LogRecord's own attributes after
message-content parsing, unconditionally overwriting anything the message body
claimed for those two keys.

* style(logging): fix import order (ruff I001) in _logging.py and litellm_logging.py

- _logging.py: import litellm belongs after the stdlib from-imports, grouped
  with the other litellm.* imports, not before them.
- litellm_logging.py: the refactor to Mapping introduced a second, separate
  `from collections.abc import Mapping` instead of merging it into the
  existing `from collections.abc import Callable` import.

Caught by the strict-rule budget gate (ruff-strict-budget.json caps I001 at
0 new violations); both auto-fixed with `ruff check --fix --select I001`.

* style(logging): freeze mutable-collection constructions flagged by LIT002

Five sites in this PR's diff built a mutable list/dict literal instead of a
frozen value: a plain list of optional strings in CorrelationPlainFormatter,
a `kwargs or {}` fallback, a `metadata or {}` fallback, two `[...]` candidate
orderings, and a `dict(headers)` copy feeding a dict comprehension. Each is
build-once/read-only, so this rewrites them as tuples, MappingProxyType, or a
plain conditional `.get()` instead of seeding then reading a fresh mutable
collection - no behavior change, confirmed by the existing test suite.

Caught by the type-discipline budget gate (LIT002 capped at 0 new
violations).

* fix(logging): reserve trace_id/session_id even when no correlation context is active

Live-proxy verification surfaced a gap in the earlier message-content-spoofing
fix (7f390a57fc): that fix only overwrites trace_id/session_id from the
LogRecord's own attribute, so it does nothing for a log line emitted before
CorrelationContextFilter has stamped anything on this record (e.g. the
"Request Headers" debug line, which fires before Logging.__init__() runs for
the request). On such a record, a caller-supplied header literally named
trace_id/session_id still got promoted into the JSON output via the embedded
JSON/dict-repr parser, since there was no genuine value to protect.

Fixed at the source: trace_id/session_id are now excluded unconditionally from
the message-content-parsing promotion step, not just superseded afterward.
Verified live against a real proxy - the exact adversarial request (headers
literally named trace_id/session_id) no longer leaks into any JSON log record.
Added a regression test for this no-active-context variant specifically,
confirmed it fails against the prior commit and passes now.

Also fixes an unrelated basedpyright regression from an earlier rebase's
conflict resolution: litellm/utils.py's `logging_obj` was incorrectly
re-annotated `Final` at its second assignment in function_setup() (it's first
declared `None` a few lines earlier), which basedpyright correctly rejects.

* fix(proxy): stop logging the raw W3C baggage session_id value

_session_id_from_baggage() extracts the caller-controlled session.id entry
verbatim - it isn't sanitized until set_session_id() runs later in
Logging.__init__(). The debug log line for this extraction interpolated the
raw value directly, so a caller could embed terminal control characters or
ANSI escape sequences that forge/alter plaintext log output for anyone
tailing the proxy's logs.

Verified live: a baggage header with an embedded ANSI escape reached the
terminal as a real, unescaped control sequence before this fix. Drops the
value from the log line entirely (the extraction succeeding is enough signal
on its own) rather than sanitizing-then-logging, matching veria-ai's
suggestion. Added a regression test using caplog that fails against the prior
commit and passes now.

* fix(logging): restore consumer context only after stream-failure exception mapping

_map_anthropic_exception/_map_aleph_alpha_exception synchronously log a debug
diagnostic (the raw status code) as part of exception_type()'s mapping.
_handle_stream_fallback_error restored the consumer's outer correlation
context before calling exception_type(), so that diagnostic log line carried
the outer (or empty) trace_id/session_id instead of the failing stream's own -
flagged by Greptile.

Moved the restore to run after mapping completes, matching the same
restore-after-not-before pattern already applied elsewhere in this file for
success/finish_reason handling. Added a regression test that captures the
correlation context live during a mocked exception_type() call; fails against
the prior commit, passes now.

* fix(logging): restore consumer context only after aclose()'s stream close completes

aclose() restored the consumer's outer correlation context as its first
statement, before awaiting the underlying provider stream's own aclose()/
close(). If that close attempt raises, the except branch's debug diagnostic
ran under the already-restored outer context instead of the closing stream's
own trace_id/session_id - flagged by Greptile, same restore-too-early pattern
as the stream-failure fix in f1cf9589d6.

Moved the restore to the end of aclose(), after the close attempt (and its
diagnostic logging) completes. Added a regression test with a fake stream
whose aclose() raises, capturing the correlation context live during the
diagnostic log call; fails against the prior commit, passes now.

* style(logging): satisfy new strict-lint budgets introduced upstream (Final, ANN401, S110, TRY300, kwargs typing)

Rebasing onto litellm_internal_staging pulled in 116 upstream commits that
introduced/tightened several lint gates this PR's own code now trips:

- LIT010 (every local/module-level variable must be Final): added Final
  annotations across _logging.py, litellm_logging.py, streaming_handler.py,
  litellm_pre_call_utils.py, and utils.py. Where a name is genuinely
  reassigned (logging_obj: starts None, later set to the real object) or
  branch-assigned, either restructured into a single ternary expression
  (ordered_candidates) or suppressed with `# rebind-ok: <reason>` matching
  this repo's documented escape hatch.
- LIT011 (parameter mutation): suppressed the two new `data[key] = value`
  writes in litellm_pre_call_utils.py with `# rebind-ok`, matching the
  unsuppressed precedent already used for every other `data[...]` write in
  the same function - `data` is an intentional out-param there.
- ANN001/ANN003/ANN202 (missing parameter/return type annotations): fully
  typed success_handler/_success_handler_body, their async twins, and
  failure_handler/_failure_handler_body/async variants in litellm_logging.py,
  plus function_setup in utils.py (added Rules to its existing TYPE_CHECKING
  block for the rules_obj: Rules annotation).
- ANN401 (explicit Any disallowed): suppressed with `# noqa: ANN401` on the
  handful of genuinely-heterogeneous result/*args/**kwargs parameters, since
  ordinary suppression is this repo's documented path.
- S110 (try/except/pass): added to the existing BLE001 noqa on the one
  best-effort correlation-cleanup try/except this PR added.
- TRY300 (return inside try): moved two `return result` statements into
  `else:` blocks in the retry-fallback paths this PR's own diff touched.
- reportPrivateUsage (basedpyright): renamed the two new
  StandardLoggingPayloadSetup static methods (get_standard_logging_payload_
  trace_id/session_id) to drop their leading underscore, since they're
  genuinely called from a sibling module-level function in the same file.

No behavior change - confirmed by the full existing test suite (819 passed)
plus all four lint gates (ruff format, ruff-strict, type-discipline,
basedpyright) passing clean.

* fix(lint): stop RUF100 flagging noqa suppressions the strict gate needs

CI's plain "ruff check" job uses the default ruff.toml, a narrower config
than ruff-strict.toml (used only by the strict-rule budget gate). ANN401 and
S110 aren't enabled in the default config, so RUF100 (unused-noqa) flagged
the `# noqa: ANN401`/`# noqa: ...,S110` suppressions this PR added as pointless
under that config, even though they're genuinely needed under ruff-strict.toml.

- ANN401: added to ruff.toml's existing `lint.external` list (same mechanism
  already used for C901/TID251, enforced by the strict gate but not by this
  config) - these Any usages are genuinely dynamic/forwarded, so the
  suppression itself is correct and just needed registering.
- S110: fixed the underlying code instead of registering another external
  code - the try/except/pass in
  CustomStreamWrapper._restore_consumer_correlation_context now logs at
  debug level on failure (matching the existing best-effort-cleanup pattern
  in _record_partial_usage_for_failure elsewhere in this file), which
  satisfies S110's own suggestion directly and needs no suppression at all.

Verified against both ruff.toml and ruff-strict.toml directly, plus all
three other gates (ruff format, type-discipline, basedpyright) and the full
test suite (821 passed).

* fix(lint): scope the ANN401 exemption to file level instead of a repo-wide noqa

Ruff has no per-line-scoped way to register a noqa code across configs (that
requires the default ruff.toml's lint.external list, which is repo-wide in
scope even though the noqa itself is per-line). Since ruff does support
file-level exemptions via per-file-ignores, and ANN401 only needed exempting
in exactly two files, moved the exemption there instead:

- ruff-strict.toml: added [lint.per-file-ignores] disabling ANN401 for
  litellm_logging.py and utils.py specifically, with a comment explaining
  why (heterogeneous response/forwarded-args parameters with no fitting
  concrete type - already verified by trying CostResponseTypes and hitting
  a real basedpyright mismatch).
- ruff.toml: reverted the ANN401 entry from lint.external - no longer
  needed, since there's no `# noqa: ANN401` left anywhere for RUF100 to
  second-guess.
- Removed the now-redundant `# noqa: ANN401` from the 10 affected
  parameters in both files, keeping the existing kwargs-ok reasons and
  adding a short inline comment on the `result`/`*args` lines pointing at
  the ruff-strict.toml exemption for context.

Verified against both configs directly (ANN401 clean under ruff-strict.toml
for these files, RUF100 clean under the default config), all four gates
(ruff format, ruff-strict, type-discipline, basedpyright), and the full
test suite (821 passed).

* fix(logging): redact credential-shaped trace_id/session_id before stamping log records

CorrelationContextFilter stamps trace_id/session_id onto a LogRecord after
SecretRedactionFilter has already run, so a caller-controlled value (e.g. via
x-litellm-trace-id or a W3C baggage header) that happens to look like a real
credential reached JSON and plaintext logs unredacted. Apply the same
credential redaction already used elsewhere in this module at
_sanitize_correlation_id(), the single choke point both set_trace_id() and
set_session_id() route through, so every caller-facing entry point is covered
without depending on filter ordering.

* fix(logging): restore correlation context when a stream's max-duration timeout fires

CustomStreamWrapper.__anext__() called _check_max_streaming_duration() before
entering its try block, so the litellm.Timeout it raises bypassed the except
Exception -> _handle_stream_fallback_error path entirely, leaking the timed-out
stream's own trace_id/session_id into whatever the consumer's task logs next.
Move the check inside the try so it flows through the same restoration path
every other stream failure already uses.

* test(streaming): make dispatch_failure_handlers mock awaitable for the async max-duration test

Moving _check_max_streaming_duration() inside __anext__()'s try block (prior
commit) means a max-duration Timeout now dispatches failure handlers through
the same path every other stream failure already uses, instead of bypassing
it entirely. dispatch_failure_handlers is async on the real Logging class;
the test's plain MagicMock logging_obj made asyncio.create_task() choke on a
non-coroutine return value once that path actually got exercised.

---------

Co-authored-by: Deepanshu <deepanshu.lulla@alpha-sense.com>
Co-authored-by: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-08-10 10:40:13 -07:00

560 lines
20 KiB
Python

import ast
import contextvars
import logging
import os
import sys
from datetime import datetime
from logging import Formatter
from typing import Any, Final
import litellm
from litellm.litellm_core_utils.safe_json_dumps import safe_dumps
from litellm.litellm_core_utils.safe_json_loads import safe_json_loads
from litellm.litellm_core_utils.secret_redaction import redact_string
set_verbose = False
session_id_var: Final[contextvars.ContextVar[str]] = contextvars.ContextVar("session_id", default="")
trace_id_var: Final[contextvars.ContextVar[str]] = contextvars.ContextVar("trace_id", default="")
_MAX_CORRELATION_ID_LENGTH: Final = 256
def _sanitize_correlation_id(value: str) -> str:
"""Strip control characters, bound length, and redact credential-shaped
content before a caller-controlled trace_id/session_id (e.g.
litellm_session_id, x-litellm-trace-id) is stamped into log lines.
Without the first two, a caller could embed \\r/\\n or terminal escape
sequences to forge fake log entries, or submit an oversized value repeated
across every log line for the request. Without the redaction, a caller
could smuggle a real credential (e.g. an sk-... key) through this field:
CorrelationContextFilter stamps trace_id/session_id onto the record after
SecretRedactionFilter has already run, so those two fields never otherwise
pass through credential redaction.
"""
stripped: Final = "".join(ch for ch in value if ch.isprintable())
return _redact_string(stripped[:_MAX_CORRELATION_ID_LENGTH])
def set_session_id(session_id: str) -> "contextvars.Token[str]":
return session_id_var.set(_sanitize_correlation_id(session_id))
def set_trace_id(trace_id: str) -> "contextvars.Token[str]":
return trace_id_var.set(_sanitize_correlation_id(trace_id))
if set_verbose is True:
logging.warning(
"`litellm.set_verbose` is deprecated. Please set `os.environ['LITELLM_LOG'] = 'DEBUG'` for debug logs."
)
_ENABLE_SECRET_REDACTION: Final = os.getenv("LITELLM_DISABLE_REDACT_SECRETS", "").lower() != "true"
def _redact_string(value: str) -> str:
if not _ENABLE_SECRET_REDACTION:
return value
return redact_string(value)
def redact_secrets(value: str) -> str:
"""Public API: redact known secret/credential patterns from an arbitrary string.
Use this for code paths that bypass the logging system — e.g. Slack/Teams
alerting, HTTP error response bodies, or any other string that may contain
secrets and will be sent to an external sink.
Not to be confused with redact_message_input_output_from_logging() in
litellm_core_utils/redact_messages.py, which redacts LLM prompt/response
content for privacy — this function redacts credential patterns (API keys,
PEM blocks, tokens, etc.) by shape.
"""
if not _ENABLE_SECRET_REDACTION:
return value
return _redact_string(value)
class SecretRedactionFilter(logging.Filter):
"""Scrubs known secret/credential patterns from log records."""
_formatter = logging.Formatter()
def filter(self, record: logging.LogRecord) -> bool:
if not _ENABLE_SECRET_REDACTION:
return True
try:
record.msg = _redact_string(record.getMessage())
record.args = None
except Exception:
if isinstance(record.msg, str):
record.msg = _redact_string(record.msg)
# Redact exception tracebacks
if record.exc_info and record.exc_info[1] is not None:
try:
record.exc_text = _redact_string(self._formatter.formatException(record.exc_info))
except Exception:
pass
# Redact extra fields passed via logger.debug("msg", extra={...})
for key, value in list(record.__dict__.items()):
if key not in _STANDARD_RECORD_ATTRS and isinstance(value, str):
setattr(record, key, _redact_string(value))
return True
_secret_filter: Final = SecretRedactionFilter()
class CorrelationContextFilter(logging.Filter):
"""Stamps each log record with the current request's trace_id and session_id from contextvars.
Works in tandem with JsonFormatter: the formatter's record.__dict__ loop picks up these
attributes as first-class JSON fields without any formatter-level code.
"""
def filter(self, record: logging.LogRecord) -> bool:
if not litellm.request_correlation_in_logs:
return True
trace_id: Final = trace_id_var.get()
if trace_id:
record.trace_id = trace_id # rebind-ok: stamping the LogRecord is the Filter interface's contract
session_id: Final = session_id_var.get()
if session_id:
record.session_id = session_id # rebind-ok: stamping the LogRecord is the Filter interface's contract
return True
_correlation_filter: Final = CorrelationContextFilter()
json_logs = bool(os.getenv("JSON_LOGS", False))
# Create a handler for the logger (you may need to adapt this based on your needs)
log_level: Final = os.getenv("LITELLM_LOG", "DEBUG")
numeric_level: Final[str] = getattr(logging, log_level.upper())
handler: Final = logging.StreamHandler()
handler.setLevel(numeric_level)
handler.addFilter(_secret_filter)
handler.addFilter(_correlation_filter)
def _try_parse_json_message(message: str) -> dict[str, Any] | None:
"""
Try to parse a log message as JSON. Returns parsed dict if valid, else None.
Handles messages that are entirely valid JSON (e.g. json.dumps output).
Uses shared safe_json_loads for consistent error handling.
"""
if not message or not isinstance(message, str):
return None
msg_stripped: Final = message.strip()
if not (msg_stripped.startswith("{") or msg_stripped.startswith("[")):
return None
parsed: Final = safe_json_loads(message, default=None)
if parsed is None or not isinstance(parsed, dict):
return None
return parsed
def _try_parse_embedded_python_dict(message: str) -> dict[str, Any] | None:
"""
Try to find and parse a Python dict repr (e.g. str(d) or repr(d)) embedded in
the message. Handles patterns like:
"get_available_deployment for model: X, Selected deployment: {'model_name': '...', ...} for model: X"
Uses ast.literal_eval for safe parsing. Returns the parsed dict or None.
"""
if not message or not isinstance(message, str) or "{" not in message:
return None
i = 0
while i < len(message):
start = message.find("{", i)
if start == -1:
break
depth = 0
for j in range(start, len(message)):
c = message[j]
if c == "{":
depth += 1
elif c == "}":
depth -= 1
if depth == 0:
substr = message[start : j + 1]
try:
result = ast.literal_eval(substr)
if isinstance(result, dict) and len(result) > 0:
return result
except (ValueError, SyntaxError, TypeError):
pass
break
i = start + 1
return None
# Standard LogRecord attribute names - used to identify 'extra' fields.
# Derived at runtime so we automatically include version-specific attrs (e.g. taskName).
def _get_standard_record_attrs() -> frozenset:
"""Standard LogRecord attribute names - excludes extra keys from logger.debug(..., extra={...})."""
return frozenset(logging.LogRecord("", 0, "", 0, "", (), None).__dict__.keys())
_STANDARD_RECORD_ATTRS: Final = _get_standard_record_attrs()
# CorrelationContextFilter is the only legitimate source for these two JSON fields;
# see JsonFormatter.format() for why they're excluded from the generic message-content
# and extra-attribute promotion paths.
_RESERVED_CORRELATION_FIELDS: Final = frozenset(("trace_id", "session_id"))
class JsonFormatter(Formatter):
def __init__(self):
super().__init__()
def formatTime(self, record, datefmt=None):
# Use datetime to format the timestamp in ISO 8601 format
dt: Final = datetime.fromtimestamp(record.created)
return dt.isoformat()
def format(self, record):
message_str: Final = record.getMessage()
json_record: Final[dict[str, Any]] = {
"message": message_str,
"level": record.levelname,
"timestamp": self.formatTime(record),
}
# Parse embedded JSON or Python dict repr in message so sub-fields become first-class properties.
# trace_id/session_id are excluded here unconditionally (not just "if not already
# set") - CorrelationContextFilter is the only legitimate source for these two
# fields, and a message that merely happens to parse as JSON/dict (e.g. a proxy
# log line dumping raw request headers) must never be able to claim them, even on
# a record the filter hasn't stamped yet (no correlation context active for it).
parsed = _try_parse_json_message(message_str)
if parsed is None:
parsed = _try_parse_embedded_python_dict(message_str)
if parsed is not None:
for key, value in parsed.items():
if key not in json_record and key not in _RESERVED_CORRELATION_FIELDS:
json_record[key] = value
# Include extra attributes passed via logger.debug("msg", extra={...})
for key, value in record.__dict__.items():
if key not in _STANDARD_RECORD_ATTRS and key not in json_record:
json_record[key] = value
# trace_id/session_id are reserved: CorrelationContextFilter is the only
# legitimate source for these two fields. Without this, a message string
# that happens to parse as JSON/dict (e.g. a proxy log line dumping raw
# request headers) with a "trace_id"/"session_id" key would have already
# claimed the key at the parsed-message step above, and the extra-attributes
# loop's "key not in json_record" guard would then skip the real value -
# letting a caller-supplied header spoof another request's correlation ids.
for reserved_key in _RESERVED_CORRELATION_FIELDS:
value = getattr(record, reserved_key, None)
if value:
json_record[reserved_key] = value
# Set component/logger only if not already supplied via extra={...}
if "component" not in json_record:
json_record["component"] = record.name
if "logger" not in json_record:
json_record["logger"] = f"{record.filename}:{record.lineno}"
if record.exc_info:
json_record["stacktrace"] = record.exc_text or self.formatException(record.exc_info)
return safe_dumps(json_record)
class CorrelationPlainFormatter(logging.Formatter):
"""Appends trace_id/session_id to plain-text log lines stamped by CorrelationContextFilter.
Mirrors JsonFormatter's handling of these two fields so request_correlation_in_logs
behaves the same whether or not json_logs is enabled.
"""
def format(self, record: logging.LogRecord) -> str:
formatted: Final = super().format(record)
trace_id: Final = getattr(record, "trace_id", None)
session_id: Final = getattr(record, "session_id", None)
if not trace_id and not session_id:
return formatted
parts: Final = tuple(
p
for p in (f"trace_id={trace_id}" if trace_id else None, f"session_id={session_id}" if session_id else None)
if p
)
return f"{formatted} [{' '.join(parts)}]"
# Function to set up exception handlers for JSON logging
def _setup_json_exception_handlers(formatter):
# Create a handler with JSON formatting for exceptions
error_handler: Final = logging.StreamHandler()
error_handler.setFormatter(formatter)
error_handler.addFilter(_secret_filter)
error_handler.addFilter(_correlation_filter)
# Setup excepthook for uncaught exceptions
def json_excepthook(exc_type, exc_value, exc_traceback):
record: Final = logging.LogRecord(
name="LiteLLM",
level=logging.ERROR,
pathname="",
lineno=0,
msg=str(exc_value),
args=(),
exc_info=(exc_type, exc_value, exc_traceback),
)
error_handler.handle(record)
sys.excepthook = json_excepthook
# Configure asyncio exception handler if possible
try:
import asyncio
def async_json_exception_handler(loop, context):
exception: Final = context.get("exception")
if exception:
exc_type: Final = type(exception)
record: Final = logging.LogRecord(
name="LiteLLM",
level=logging.ERROR,
pathname="",
lineno=0,
msg=str(exception),
args=(),
exc_info=(exc_type, exception, exception.__traceback__),
)
error_handler.handle(record)
else:
loop.default_exception_handler(context)
asyncio.get_event_loop().set_exception_handler(async_json_exception_handler)
except Exception:
pass
# Create a formatter and set it for the handler
if json_logs:
handler.setFormatter(JsonFormatter())
_setup_json_exception_handlers(JsonFormatter())
else:
formatter: Final = CorrelationPlainFormatter(
"\033[92m%(asctime)s - %(name)s:%(levelname)s\033[0m: %(filename)s:%(lineno)s - %(message)s",
datefmt="%H:%M:%S",
)
handler.setFormatter(formatter)
verbose_proxy_logger = logging.getLogger("LiteLLM Proxy")
verbose_router_logger = logging.getLogger("LiteLLM Router")
verbose_logger = logging.getLogger("LiteLLM")
# Add the handler to the loggers
verbose_router_logger.addHandler(handler)
verbose_proxy_logger.addHandler(handler)
verbose_logger.addHandler(handler)
def _suppress_loggers():
"""Suppress noisy loggers at INFO level"""
# Suppress httpx request logging at INFO level
httpx_logger: Final = logging.getLogger("httpx")
httpx_logger.setLevel(logging.WARNING)
# Suppress APScheduler logging at INFO level
apscheduler_executors_logger: Final = logging.getLogger("apscheduler.executors.default")
apscheduler_executors_logger.setLevel(logging.WARNING)
apscheduler_scheduler_logger: Final = logging.getLogger("apscheduler.scheduler")
apscheduler_scheduler_logger.setLevel(logging.WARNING)
_REDACTED_THIRD_PARTY_LOGGERS: Final[tuple[str, ...]] = (
"apscheduler.executors.default",
"apscheduler.scheduler",
"asyncio",
"backoff",
"httpx",
"uvicorn.error",
)
def _redact_third_party_loggers() -> None:
"""Extend secret redaction to records litellm does not emit directly.
litellm's own loggers are covered by the filter on their shared handler, but a
litellm value can also reach a log record through a dependency that logs on its
own logger. Those records never pass through a litellm handler.
The filter is attached to each emitting logger rather than to the root logger or
to root's handlers. `Logger.handle` applies the emitting logger's filters before
any handler runs, so redaction happens once, at the earliest point in the
record's life, and covers every downstream handler regardless of who owns it.
The alternatives do not hold: `callHandlers` consults ancestors for handlers but
never for filters, so a filter on the root logger never sees these records at
all, and a filter on a root handler only covers that one handler, leaving
handlers registered earlier or on the emitting logger itself untouched.
Each name is the exact logger a dependency emits on; a parent name would not
cover its children, for the same reason the root logger does not.
"""
for name in _REDACTED_THIRD_PARTY_LOGGERS:
logging.getLogger(name).addFilter(_secret_filter)
# Call the suppression function
_suppress_loggers()
_redact_third_party_loggers()
ALL_LOGGERS: Final = [
logging.getLogger(),
verbose_logger,
verbose_router_logger,
verbose_proxy_logger,
]
def _get_loggers_to_initialize():
"""
Get all loggers that should be initialized with the JSON handler.
Includes third-party integration loggers (like langfuse) if they are
configured as callbacks.
"""
import litellm
loggers: Final = list(ALL_LOGGERS)
# Add langfuse logger if langfuse is being used as a callback
langfuse_callbacks: Final = {"langfuse", "langfuse_otel"}
all_callbacks: Final = set(litellm.success_callback + litellm.failure_callback)
if langfuse_callbacks & all_callbacks:
loggers.append(logging.getLogger("langfuse"))
return loggers
def _initialize_loggers_with_handler(handler: logging.Handler):
"""
Initialize all loggers with a handler
- Adds a handler to each logger
- Prevents bubbling to parent/root (critical to prevent duplicate JSON logs)
"""
handler.addFilter(_secret_filter)
handler.addFilter(_correlation_filter)
for lg in _get_loggers_to_initialize():
lg.handlers.clear() # remove any existing handlers
lg.addHandler(handler) # add JSON formatter handler
lg.propagate = False # prevent bubbling to parent/root
def _get_uvicorn_json_log_config():
"""
Generate a uvicorn log_config dictionary that applies JSON formatting to all loggers.
This ensures that uvicorn's access logs, error logs, and all application logs
are formatted as JSON when json_logs is enabled.
"""
json_formatter_class: Final = "litellm._logging.JsonFormatter"
# Use the module-level log_level variable for consistency
uvicorn_log_level: Final = log_level.upper()
log_config: Final = {
"version": 1,
"disable_existing_loggers": False,
"formatters": {
"json": {
"()": json_formatter_class,
},
"default": {
"()": json_formatter_class,
},
"access": {
"()": json_formatter_class,
},
},
"handlers": {
"default": {
"formatter": "json",
"class": "logging.StreamHandler",
"stream": "ext://sys.stdout",
},
"access": {
"formatter": "access",
"class": "logging.StreamHandler",
"stream": "ext://sys.stdout",
},
},
"loggers": {
"uvicorn": {
"handlers": ["default"],
"level": uvicorn_log_level,
"propagate": False,
},
"uvicorn.error": {
"handlers": ["default"],
"level": uvicorn_log_level,
"propagate": False,
},
"uvicorn.access": {
"handlers": ["access"],
"level": uvicorn_log_level,
"propagate": False,
},
},
}
return log_config
def _turn_on_json():
"""
Turn on JSON logging
- Adds a JSON formatter to all loggers
"""
handler: Final = logging.StreamHandler()
handler.setFormatter(JsonFormatter())
_initialize_loggers_with_handler(handler)
# Set up exception handlers
_setup_json_exception_handlers(JsonFormatter())
def _turn_on_debug():
verbose_logger.setLevel(level=logging.DEBUG) # set package log to debug
verbose_router_logger.setLevel(level=logging.DEBUG) # set router logs to debug
verbose_proxy_logger.setLevel(level=logging.DEBUG) # set proxy logs to debug
def _disable_debugging():
"""Disable the package, router, and proxy verbose loggers."""
verbose_logger.disabled = True
verbose_router_logger.disabled = True
verbose_proxy_logger.disabled = True
def _enable_debugging():
verbose_logger.disabled = False
verbose_router_logger.disabled = False
verbose_proxy_logger.disabled = False
def print_verbose(print_statement):
try:
if set_verbose:
print(redact_secrets(str(print_statement))) # noqa: T201
except Exception:
pass
def _is_debugging_on() -> bool:
"""
Returns True if debugging is on
"""
return verbose_logger.isEnabledFor(logging.DEBUG) or set_verbose is True