Skip to content

fix(logging_worker): make flush() survive an event loop change - #42355

Merged
mateo-berri merged 4 commits into
mainfrom
litellm_logging_worker_flush_loop_change
Sep 21, 2026
Merged

mateo-berri merged 4 commits into
mainfrom
litellm_logging_worker_flush_loop_change

Conversation

@devin-ai-integration

@devin-ai-integration devin-ai-integration Bot commented Sep 21, 2026 •

Copy link
Copy Markdown
Contributor

TLDR

Problem this solves:

How it solves it:

  • flush() calls start() first, so it joins the current loop's queue
  • Tasks stranded on the closed loop get carried over and drained
  • Regression tests: enqueue on one loop, flush on another, callback fires

User Flow

Before: a contributor pushes any branch and the required CircleCI unit job comes back red on 12 tests their change never touched

  1. They push a branch and open its workflow at https://app.circleci.com/pipelines/github/BerriAI/litellm
  2. The unit job runs about 6,400 tests and finishes red: 12 failed, all in tests/unit/router_strategy/complexity_router/test_jev_classifier.py
  3. The first failing case reads Failed: Timeout (>90.0s) from pytest-timeout. and the other 11 read RuntimeError: <asyncio.locks.Event ...> is bound to a different event loop
  4. They run uv run pytest tests/unit/router_strategy/complexity_router/test_jev_classifier.py locally and get 43 passed, so nothing points at their change
  5. A re-run of the job fails the same way, and the PR stays unmergeable until someone overrides the check

After: the same push comes back green and the job finishes about 90 seconds sooner

  1. They push a branch and open its workflow at https://app.circleci.com/pipelines/github/BerriAI/litellm
  2. The unit job runs the same tests and finishes green, the 12 jev cases included
  3. They run uv run pytest tests/unit/router_strategy/complexity_router/test_jev_classifier.py locally and get 43 passed
  4. The required check is green on the first run and the PR is mergeable

The same bug reaches SDK users who flush the logging worker themselves

Before: a script that runs litellm.acompletion in one asyncio.run() and then asyncio.run(GLOBAL_LOGGING_WORKER.flush()) never returns from the flush

  1. They run the script; the OpenAI call returns a chatcmpl-... id
  2. The flush call blocks forever (until whatever timeout wraps it), and the success callback for that completion never fires

After: the same script returns from the flush right away

  1. They run the script; the OpenAI call returns a chatcmpl-... id
  2. The flush returns in well under a second and the success callback fires with that completion's id

Relevant issues

Affected release

Linear ticket

Pre-Submission checklist

Please complete all items before asking a LiteLLM maintainer to review your PR

  • I have added meaningful tests
  • The handful of test files covering my change pass locally, e.g. uv run pytest tests/test_litellm/<your_test_file>.py -v. Leave the suites (make test-unit-*, make test-unit) to CI: it finishes in ~15 minutes where a laptop takes an hour or more
  • My PR passes all required CI/CD checks (e.g., lint, schema.d.ts sync check, etc.)
  • My PR's scope is as isolated as possible; it only solves 1 specific problem
  • I have received a Greptile Confidence Score of at least 4/5 before requesting a maintainer review (Greptile reviews automatically once the PR is opened; only comment @greptileai to re-request a review after pushing changes)

Delays in PR merge?

If you're seeing a delay in your PR being merged, ping the LiteLLM Team on Slack (#pr-review).

Screenshots / Proof of Fix

Shared setup: a worktree of this repo after make bootstrap (Python 3.14.3, pytest-asyncio 1.3.0, pytest-xdist 3.8.0, pytest-timeout 2.4.0), every command run from the worktree root. Case 3 needs OPENAI_API_KEY in the environment and uses this script, saved as sdk_flush.py outside the repo:

import asyncio
import time

import litellm
from litellm.integrations.custom_logger import CustomLogger
from litellm.litellm_core_utils.logging_worker import GLOBAL_LOGGING_WORKER


class Recorder(CustomLogger):
    def __init__(self) -> None:
        super().__init__()
        self.ids: list[str] = []

    async def async_log_success_event(self, kwargs, response_obj, start_time, end_time) -> None:
        self.ids.append(response_obj.id)


recorder = Recorder()
litellm.callbacks = [recorder]


async def ask() -> str:
    response = await litellm.acompletion(
        model="gpt-5.4-mini",
        messages=[{"role": "user", "content": "Say hi in two words"}],
    )
    return response.id


async def flush() -> float:
    started = time.monotonic()
    await asyncio.wait_for(GLOBAL_LOGGING_WORKER.flush(), timeout=30)
    return time.monotonic() - started


print("litellm loaded from", litellm.__file__)
response_id = asyncio.run(ask())
print("response id", response_id, "| callback fired before flush:", recorder.ids)
try:
    elapsed = asyncio.run(flush())
    print(f"flush returned in {elapsed:.2f}s | callback fired:", recorder.ids)
except TimeoutError:
    print("flush hung: TimeoutError after 30s | callback fired:", recorder.ids)

Before (3d26a29)

The exact CircleCI unit job command over the whole tests/unit tree

  1. Command

    mapfile -t files < <(find tests/unit -name 'test_*.py' | sort)
    LITELLM_LOCAL_MODEL_COST_MAP=True uv run --no-sync pytest "${files[@]}" -p no:rerunfailures -p no:pytest-retry --timeout=90 -n 4 --dist=loadscope --tb=short --junitxml=test-results/unit/junit.xml 2>&1 | tee pytest.log
    grep -c '^FAILED .*test_jev' pytest.log
    
  2. Observed output

    ========== 196 failed, 6223 passed, 423 warnings, 2 errors in 122.12s ==========
    12
    

    The 12 are exactly the cases CircleCI job 2197604 on main fails (4x test_jev_http_errors_do_not_dispatch_successful_usage, 8x test_jev_invalid_usage_never_reaches_spend_callbacks). The other 184 only fail on this laptop (missing vertexai, secret-detection scanner cases) and are not in the CircleCI results

Two modules on one worker, the smallest sequence that reproduces

  1. Command

    LITELLM_LOCAL_MODEL_COST_MAP=True uv run --no-sync pytest tests/unit/llms/chat/test_converse_handler.py tests/unit/router_strategy/complexity_router/test_jev_classifier.py -p no:xdist -p no:rerunfailures -p no:pytest-retry --timeout=90 --tb=short -q
    
  2. Observed output

    FAILED tests/unit/router_strategy/complexity_router/test_jev_classifier.py::test_jev_http_errors_do_not_dispatch_successful_usage[400] - Failed: Timeout (>90.0s) from pytest-timeout.
    FAILED tests/unit/router_strategy/complexity_router/test_jev_classifier.py::test_jev_http_errors_do_not_dispatch_successful_usage[429] - RuntimeError: <asyncio.locks.Event object at 0x11165b620 [unset]> is bound to a different event loop
    ... (9 more RuntimeError cases)
    12 failed, 47 passed, 76 warnings in 90.99s (0:01:30)
    

    The RuntimeError traceback ends in logging_worker.py:490: in flush -> await self._queue.join()

SDK script: one real OpenAI call, then flush from a second asyncio.run()

  1. Command

    PYTHONPATH=$PWD python sdk_flush.py
    
  2. Observed output

    litellm loaded from .../base/litellm/__init__.py
    response id chatcmpl-EQh1ER1gFY4Nnh8tIjP0mNKamxlFm | callback fired before flush: []
    flush hung: TimeoutError after 30s | callback fired: []
    

After (e86ba8b)

The tip is now 9388602: a merge of main (37f1670) that resolves the branch's stale lint script, then a test-only commit that swaps the two new tests' list trackers for AsyncMock. The merge brings in commits touching neither logging_worker.py nor the two test modules above, and the PR's logging_worker.py diff against its base is byte-identical before and after both commits, so the runs below still describe the tip; the last item is CircleCI's own unit job run at 9388602 itself

The exact CircleCI unit job command over the whole tests/unit tree

  1. Same command

  2. Observed output

    ===== 184 failed, 6235 passed, 410 warnings, 2 errors in 105.85s (0:01:45) =====
    0
    

    The failure sets differ by exactly the 12 jev cases (12 only before, 0 only after); the 90-second hang is gone (an earlier run of the same command at 212ab63, the fix commit alone, finished in 38.77s on an idle machine with the same 184 failures)

Two modules on one worker, the smallest sequence that reproduces

  1. Same command

  2. Observed output

    59 passed, 76 warnings in 2.72s
    

SDK script: one real OpenAI call, then flush from a second asyncio.run()

  1. Same command

  2. Observed output

    LoggingWorker: event loop changed; carried 1 pending and revived 0 dequeued logging task(s) onto the new loop
    litellm loaded from .../wt/litellm/__init__.py
    response id chatcmpl-EQhGNqxxsJheho8BGn6PwjThLbaBZ | callback fired before flush: []
    flush returned in 1.60s | callback fired: ['chatcmpl-EQhGNqxxsJheho8BGn6PwjThLbaBZ']
    

CircleCI's own unit job at the tip (9388602)

  1. Job: https://app.circleci.com/pipelines/github/BerriAI/litellm/90019/workflows/c658d518-d040-4418-b6d9-c876fb3c885f/jobs/2198653 (pipeline 90019, started 2026-09-21 23:22Z)

  2. Observed result, read back from the CircleCI tests API for job 2198653

    6431 tests, 43 test_jev_classifier cases, 0 jev failures, slowest jev case 1.16s
    1 failure: tests/unit/enterprise/enterprise_callbacks/test_secret_detection.py::test_scan_message_stays_linear_on_adversarial_credential_lines[assignment-flood]
      assert (269.98 - 258.53) < 10.0   (a wall-clock bound; the same case fails main's own pipeline 89979)
    

    Same-hour control on branches without this fix: pipeline 90018 (23:01Z, job 2198616) and pipeline 90015 (22:41Z, job 2198465) both still fail the 12 jev cases

Type

🐛 Bug Fix

Caveats (if any)

Low

  • flush() now starts a worker on the calling loop when none runs there
    • nothing in litellm/ or enterprise/ calls flush(), only tests and SDK users
    • the atexit drain (_flush_on_exit) never goes through flush(), so it is untouched
    • chosen over a second inline drain loop like the atexit one (a duplicate of the worker's timeout and concurrency handling) and over raising on a loop change (every caller today wants the wait to finish, not an error)
  • The 12 jev cases still flush without enqueueing first
    • whether they meet a stale queue depends on xdist scheduling
    • the fix makes that harmless rather than deterministic
  • flush() also runs callbacks earlier tests left on the previous loop
    • same carry-over ensure_initialized_and_enqueue already did
    • a test asserting "no callback fired at all" after a flush would be order-dependent; none does today
  • A queue at capacity (50,000 by default) when a loop closes cannot requeue every dequeued task
    • flush() then returns with the leftover parked until the next loop change or atexit
    • that gap predates this change
  • 184 unrelated failures in the local full-suite runs above come from optional deps missing on the laptop; CircleCI has them and shows only the 12
  • Five CircleCI jobs are red at the tip for reasons already red on every branch this hour, none touching logging_worker.py; the GitHub required checks (the ruleset's list, all GitHub Actions jobs) are green
    • unit: only the wall-clock-bound secret detection case above, also red on main pipeline 89979
    • llm_translation_testing (10 Fireworks document-inlining cases, one hosted_vllm unhashable type: 'list', one Bedrock Moonshot JSON stream) and local_testing_part1 (two Bedrock get_model_info cases): identical failure lists on pipelines 90015 and 90018 of other branches
    • integration-cost and integration-providers: red on every pipeline from 90009 (21:54Z) through 90019 across seven branches, green on 90008 and 90010 whose tips predate the newest main commits

Final Attestation

  • The tests check the right things, including the edge cases, and regressions in the respective real-world customer use-cases are not possible after this PR

  • e86ba8b passes /live-pr-risk

flush() awaited join() on whatever queue the worker held, even one bound to
an event loop that has since closed. Its unfinished counter is never
decremented on the new loop, so the first flush() after a loop change hung
until pytest-timeout killed it and every later one raised "is bound to a
different event loop" from the queue's Event. The CircleCI unit job has
been red on every branch since the first tests that flush without
enqueueing landed, and an SDK script that flushes from a second
asyncio.run() hangs the same way.

flush() now goes through start() first, which carries the tasks stranded
on the previous loop onto the current one and guarantees a worker there to
drain them, the same loop-change handling every other entry point already
had.
@devin-ai-integration

devin-ai-integration Bot commented Sep 21, 2026 •

Copy link
Copy Markdown
Contributor Author

I'll fix CI failures and address comments from users with write access. I'll skip comments containing "(aside)".

  • Disable automatic comment, CI, and merge conflict monitoring

@greptile-apps

greptile-apps Bot commented Sep 21, 2026 •

Copy link
Copy Markdown
Contributor

RetriggerConfidence Score: 5/5

The PR appears safe to merge; the current implementation addresses cross-loop flushing and no outstanding findings remain.

Summary

The PR makes LoggingWorker.flush() initialize the worker on the current event loop before joining its queue, allowing work stranded by a closed loop to be transferred and drained.

  • Adds cross-event-loop regression coverage for queued and dequeued logging tasks.
  • Verifies that flushing starts a worker when a current-loop queue has no active consumer.
  • Replaces mutable test trackers with AsyncMock.await_count assertions, addressing the resolved previous review finding.

Reviews (3) · Last reviewed commit: "test(logging_worker): track callback run..."

@codspeed

codspeed Bot commented Sep 21, 2026 •

Copy link
Copy Markdown
Contributor

Merging this PR will not alter performance

✅ 31 untouched benchmarks


Comparing litellm_logging_worker_flush_loop_change (9388602) with main (1baa26d)1

Open in CodSpeed

Footnotes

  1. No successful run was found on main (a543124) during the generation of this report, so 1baa26d was used instead as the comparison base. There might be some changes unrelated to this pull request in this report. ↩

@codecov

codecov Bot commented Sep 21, 2026 •

Copy link
Copy Markdown

Codecov Report

✅ All modified and coverable lines are covered by tests.

📢 Thoughts on this report? Let us know!

@mateo-berri

Copy link
Copy Markdown
Contributor

@greptileai

Comment thread tests/test_litellm/litellm_core_utils/test_logging_worker.py Outdated
@mateo-berri

Copy link
Copy Markdown
Contributor

@greptileai

@mateo-berri

Copy link
Copy Markdown
Contributor

bugbot run

@cursor cursor Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

✅ Bugbot reviewed your changes and found no new issues!

Comment @cursor review or bugbot run to trigger another review on this PR

Reviewed by Cursor Bugbot for commit 9388602. Configure here.

@mateo-berri mateo-berri added run-ci and removed run-ci labels Sep 21, 2026

@mateo-berri mateo-berri left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LGTM

@mateo-berri
mateo-berri merged commit 8c8fb73 into main Sep 21, 2026
134 of 153 checks passed
@mateo-berri
mateo-berri deleted the litellm_logging_worker_flush_loop_change branch September 21, 2026 23:33
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant