fix(langfuse): propagate interrupts raised during deferred client teardown

Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
yucheng 2026-09-15 00:00:36 +00:00
parent 58c913eca3
commit 45fc3111dc
2 changed files with 33 additions and 12 deletions

View file

@ -342,21 +342,15 @@ def lease_langfuse_client(client: Langfuse) -> Generator[None]:
try:
yield
finally:
_run_teardowns(state, state.release_lease(), propagate_base_exception=False)
_run_teardowns(state, state.release_lease())
def _run_teardowns(
state: _LangfuseLifecycleState,
clients: tuple[Langfuse, ...],
*,
propagate_base_exception: bool = True,
) -> None:
def _run_teardowns(state: _LangfuseLifecycleState, clients: tuple[Langfuse, ...]) -> None:
"""Tear down ``clients``, then whatever eviction queued meanwhile, and hand the flag back.
A failing ordinary teardown is logged and skipped rather than raised: the thread here is usually a
request callback that merely held the last lease, and its request must not fail on eviction's behalf.
Interrupts requeue the unfinished batch and normally propagate, while a callback exception already
in flight takes precedence over an eviction interrupt.
An interrupt requeues the unfinished batch for the next eviction or lease exit and propagates.
"""
batch = clients # rebind-ok: drains each batch queued while the previous one was being torn down
try:
@ -370,9 +364,7 @@ def _run_teardowns(
verbose_logger.exception("Langfuse client teardown failed during cache eviction")
except BaseException:
state.requeue(batch[index:])
if propagate_base_exception:
raise
return
raise
batch = state.next_teardown_batch()
finally:
state.end_teardown()

View file

@ -722,6 +722,35 @@ def test_teardown_failure_does_not_strand_queued_clients(monkeypatch):
assert not state.teardown_in_progress
def test_interrupt_during_deferred_teardown_propagates_and_requeues_the_client(monkeypatch):
"""A Ctrl-C landing in the lease exit's teardown must reach the caller, not be swallowed."""
client = Langfuse(public_key=PUBLIC_KEY, secret_key="sk-original", host="http://127.0.0.1:1")
register_langfuse_client(client)
state = _lifecycle_state(client)
original_teardown = _teardown_langfuse_client
calls = []
def teardown(target):
calls.append(target)
if len(calls) == 1:
raise KeyboardInterrupt
original_teardown(target)
monkeypatch.setattr("litellm.integrations.langfuse.langfuse_sdk._teardown_langfuse_client", teardown)
with pytest.raises(KeyboardInterrupt):
with lease_langfuse_client(client):
shutdown_langfuse_client(client)
assert state.pending_clients == {client}
assert not state.teardown_in_progress
with lease_langfuse_client(client):
pass
assert calls == [client, client]
assert not state.pending_clients
def test_queued_eviction_waits_for_the_last_of_two_overlapping_leases():
exporter = InMemorySpanExporter()
provider = build_isolated_tracer_provider(environment=None, release=None)