mirror of
https://github.com/BerriAI/litellm.git
synced 2026-09-26 01:12:21 +00:00
fix(langfuse): claim the cache entry before releasing its slot and channel hold on eviction
Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
parent
21ee0008a7
commit
548c486372
2 changed files with 47 additions and 12 deletions
|
|
@ -30,22 +30,23 @@ class LangfuseInMemoryCache(InMemoryCache):
|
|||
def _remove_key(self, key: str) -> None:
|
||||
from litellm.integrations.langfuse.langfuse import LangFuseLogger
|
||||
|
||||
if isinstance(self.cache_dict[key], LangFuseLogger):
|
||||
evicted: Final = self.cache_dict.pop(key, None)
|
||||
self.ttl_dict.pop(key, None)
|
||||
if evicted is None:
|
||||
return
|
||||
|
||||
if isinstance(evicted, LangFuseLogger):
|
||||
litellm.initialized_langfuse_clients -= 1
|
||||
|
||||
# Loggers with a periodic flush task (e.g. NewRelicMetricsLogger) expose
|
||||
# stop() so eviction actually ends the task instead of leaking it.
|
||||
_evicted_stop: Final = getattr(self.cache_dict[key], "stop", None)
|
||||
if callable(_evicted_stop):
|
||||
try:
|
||||
_evicted_stop()
|
||||
except Exception: # noqa: BLE001 # a failing stop() must not block eviction
|
||||
verbose_logger.debug("DynamicLoggingCache: stop() raised during eviction", exc_info=True)
|
||||
|
||||
#########################################################
|
||||
# Call parent class to remove key from cache
|
||||
#########################################################
|
||||
return super()._remove_key(key)
|
||||
_evicted_stop: Final = getattr(evicted, "stop", None)
|
||||
if not callable(_evicted_stop):
|
||||
return
|
||||
try:
|
||||
_evicted_stop()
|
||||
except Exception: # noqa: BLE001 # a failing stop() must not block eviction
|
||||
verbose_logger.debug("DynamicLoggingCache: stop() raised during eviction", exc_info=True)
|
||||
|
||||
|
||||
class DynamicLoggingCache:
|
||||
|
|
|
|||
|
|
@ -82,3 +82,37 @@ class TestLangfuseInMemoryCache:
|
|||
|
||||
release_langfuse_tracing(sibling, grace_seconds=0.0)
|
||||
assert acquire() is not logger.tracing, "eviction did not release the evicted logger's hold"
|
||||
|
||||
@patch("litellm.initialized_langfuse_clients", 3)
|
||||
def test_second_evictor_of_the_same_entry_releases_nothing(self):
|
||||
"""Two callers can expire the same entry at once (a request thread and the reaper). Only the one that
|
||||
claims the entry may give its slot and channel hold back, or a sibling logger loses its channel."""
|
||||
from litellm.integrations.langfuse.langfuse import LangFuseLogger
|
||||
from litellm.integrations.langfuse.langfuse_sdk import acquire_langfuse_tracing, release_langfuse_tracing
|
||||
|
||||
def acquire():
|
||||
return acquire_langfuse_tracing(
|
||||
public_key="pk-double-eviction-test",
|
||||
secret_key="sk",
|
||||
base_url="http://127.0.0.1:1",
|
||||
environment=None,
|
||||
release=None,
|
||||
flush_interval=1.0,
|
||||
mock_mode=True,
|
||||
)
|
||||
|
||||
logger = LangFuseLogger.__new__(LangFuseLogger)
|
||||
logger.api_client = MagicMock()
|
||||
logger.tracing = acquire()
|
||||
sibling = acquire()
|
||||
self.cache.cache_dict["test_key"] = logger
|
||||
self.cache.ttl_dict["test_key"] = time.time() + 100
|
||||
|
||||
self.cache._remove_key("test_key")
|
||||
self.cache._remove_key("test_key")
|
||||
|
||||
assert litellm.initialized_langfuse_clients == 2
|
||||
assert "test_key" not in self.cache.cache_dict and "test_key" not in self.cache.ttl_dict
|
||||
assert acquire() is sibling, "the second evictor took the sibling logger's hold on the channel"
|
||||
release_langfuse_tracing(sibling)
|
||||
release_langfuse_tracing(sibling, grace_seconds=0.0)
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue