From 685d7f79b10da72a47c50d28d4be2e3b17cce38f Mon Sep 17 00:00:00 2001 From: Hassieb Pakzad <68423100+hassiebp@users.noreply.github.com> Date: Wed, 5 Aug 2026 18:02:08 +0200 Subject: [PATCH] fix(resource-manager): prevent shutdown queue deadlock --- langfuse/_client/resource_manager.py | 66 +++++++++++---- langfuse/_task_manager/media_manager.py | 37 ++++++-- tests/unit/test_resource_manager.py | 108 +++++++++++++++++++++++- 3 files changed, 186 insertions(+), 25 deletions(-) diff --git a/langfuse/_client/resource_manager.py b/langfuse/_client/resource_manager.py index 67c44920a..9837f095c 100644 --- a/langfuse/_client/resource_manager.py +++ b/langfuse/_client/resource_manager.py @@ -134,10 +134,14 @@ def __new__( id_generator: Optional[IdGenerator] = None, span_exporter: Optional[SpanExporter] = None, ) -> "LangfuseResourceManager": - if public_key in cls._instances: - return cls._instances[public_key] - with cls._lock: + cached_instance = cls._instances.get(public_key) + if cached_instance is not None and not cached_instance._shutdown: + return cached_instance + + if cached_instance is not None: + cls._instances.pop(public_key, None) + if public_key not in cls._instances: instance = super(LangfuseResourceManager, cls).__new__(cls) @@ -226,6 +230,7 @@ def _initialize_instance( self._custom_httpx_client = httpx_client self._init_api_clients() + self._span_processor: Optional[LangfuseSpanProcessor] = None # Media self._media_upload_enabled = os.environ.get( @@ -263,6 +268,7 @@ def _initialize_instance( mask_otel_spans=mask_otel_spans, ) tracer_provider.add_span_processor(langfuse_processor) + self._span_processor = langfuse_processor self._otel_tracer = tracer_provider.get_tracer( LANGFUSE_TRACER_NAME, @@ -476,11 +482,12 @@ def _at_fork_reinit(self) -> None: @classmethod def reset(cls) -> None: with cls._lock: - for key in cls._instances: - cls._instances[key].shutdown() - + instances = list(cls._instances.values()) cls._instances.clear() + for instance in instances: + instance.shutdown() + def add_score_task(self, event: dict, *, force_sample: bool = False) -> None: try: # Sample scores with the same sampler that is used for tracing @@ -508,10 +515,17 @@ def add_score_task(self, event: dict, *, force_sample: bool = False) -> None: ) if should_sample: - langfuse_logger.debug( - f"Score: Enqueuing event type={event['type']} for trace_id={event['body'].trace_id} name={event['body'].name} value={event['body'].value}" - ) - self._score_ingestion_queue.put(event, block=False) + with self._lock: + if self._shutdown: + langfuse_logger.warning( + "Score: Dropping event because the Langfuse client has already been shut down." + ) + return + + langfuse_logger.debug( + f"Score: Enqueuing event type={event['type']} for trace_id={event['body'].trace_id} name={event['body'].name} value={event['body'].value}" + ) + self._score_ingestion_queue.put(event, block=False) except Full: langfuse_logger.warning( @@ -531,10 +545,17 @@ def add_trace_task( event: dict, ) -> None: try: - langfuse_logger.debug( - f"Trace: Enqueuing event type={event['type']} for trace_id={event['body'].id}" - ) - self._score_ingestion_queue.put(event, block=False) + with self._lock: + if self._shutdown: + langfuse_logger.warning( + "Trace: Dropping event because the Langfuse client has already been shut down." + ) + return + + langfuse_logger.debug( + f"Trace: Enqueuing event type={event['type']} for trace_id={event['body'].id}" + ) + self._score_ingestion_queue.put(event, block=False) except Full: langfuse_logger.warning( @@ -612,13 +633,24 @@ def flush(self) -> None: langfuse_logger.debug("Successfully flushed media upload queue") def shutdown(self) -> None: - self._shutdown = True + with self._lock: + if self._shutdown: + return + + self._shutdown = True + if self._instances.get(self.public_key) is self: + self._instances.pop(self.public_key) + self._media_manager.begin_shutdown() # Unregister the atexit handler first atexit.unregister(self.shutdown) - self.flush() - self._stop_and_join_consumer_threads() + try: + self.flush() + finally: + self._stop_and_join_consumer_threads() + if self._span_processor is not None: + self._span_processor.shutdown() def _init_tracer_provider( diff --git a/langfuse/_task_manager/media_manager.py b/langfuse/_task_manager/media_manager.py index 14dceec19..ad2bc11c7 100644 --- a/langfuse/_task_manager/media_manager.py +++ b/langfuse/_task_manager/media_manager.py @@ -1,4 +1,5 @@ import os +import threading import time from queue import Empty, Full, Queue from typing import Any, Callable, Optional, TypeVar, cast @@ -42,6 +43,8 @@ def __init__( self._httpx_client = httpx_client self._queue = media_upload_queue self._max_retries = max_retries + self._state_lock = threading.Lock() + self._shutdown = False self._enabled = os.environ.get( LANGFUSE_MEDIA_UPLOAD_ENABLED, "True" ).lower() not in ("false", "0") @@ -53,9 +56,15 @@ def reinitialize( httpx_client: httpx.Client, media_upload_queue: Queue, ) -> None: - self._api_client = api_client - self._httpx_client = httpx_client - self._queue = media_upload_queue + with self._state_lock: + self._api_client = api_client + self._httpx_client = httpx_client + self._queue = media_upload_queue + self._shutdown = False + + def begin_shutdown(self) -> None: + with self._state_lock: + self._shutdown = True def process_next_media_upload(self) -> None: try: @@ -99,6 +108,13 @@ def _find_and_process_media( if not self._enabled: return data + with self._state_lock: + if self._shutdown: + logger.warning( + "Media: Skipping upload because the Langfuse client has already been shut down." + ) + return data + seen = set() max_levels = 10 @@ -279,10 +295,17 @@ def _process_media( field=field, ) - self._queue.put( - item=upload_media_job, - block=False, - ) + with self._state_lock: + if self._shutdown: + logger.warning( + f"Media: Skipping upload for media_id={media._media_id} because the Langfuse client has already been shut down." + ) + return + + self._queue.put( + item=upload_media_job, + block=False, + ) logger.debug( f"Queue: Enqueued media ID {media._media_id} for upload processing | trace_id={trace_id} | field={field}" ) diff --git a/tests/unit/test_resource_manager.py b/tests/unit/test_resource_manager.py index f66a1e052..ccf3a28bb 100644 --- a/tests/unit/test_resource_manager.py +++ b/tests/unit/test_resource_manager.py @@ -20,11 +20,14 @@ class NoOpSpanExporter(SpanExporter): """Minimal exporter used to verify configuration propagation.""" + def __init__(self) -> None: + self.shutdown_count = 0 + def export(self, spans: Sequence[ReadableSpan]) -> SpanExportResult: return SpanExportResult.SUCCESS def shutdown(self) -> None: - pass + self.shutdown_count += 1 def test_get_client_preserves_all_settings(monkeypatch): @@ -172,6 +175,109 @@ def test_media_upload_consumer_signal_shutdown_wakes_blocked_thread(): assert not consumer.is_alive() +def test_shutdown_evicts_manager_and_rejects_stale_client_tasks(monkeypatch): + monkeypatch.setenv("LANGFUSE_MEDIA_UPLOAD_ENABLED", "false") + + with LangfuseResourceManager._lock: + LangfuseResourceManager._instances.clear() + + old_exporter = NoOpSpanExporter() + settings = { + "public_key": "pk-shutdown-reinit", + "secret_key": "sk-shutdown-reinit", + "span_exporter": old_exporter, + } + first_client = Langfuse(**settings) + stale_client = Langfuse(**settings) + old_manager = first_client._resources + + assert old_manager is not None + assert stale_client._resources is old_manager + + first_client.shutdown() + + assert old_manager._shutdown + assert settings["public_key"] not in LangfuseResourceManager._instances + assert not old_manager._ingestion_consumers[0].is_alive() + assert old_exporter.shutdown_count == 1 + + stale_client.create_score(name="quality", value=1.0) + stale_client._create_trace_tags_via_ingestion( + trace_id="0" * 32, + tags=["after-shutdown"], + ) + + assert old_manager._score_ingestion_queue.unfinished_tasks == 0 + stale_client.shutdown() + + fresh_client = Langfuse( + public_key=settings["public_key"], + secret_key=settings["secret_key"], + span_exporter=NoOpSpanExporter(), + ) + fresh_manager = fresh_client._resources + + assert fresh_manager is not None + assert fresh_manager is not old_manager + assert fresh_manager._ingestion_consumers[0].is_alive() + + fresh_client.shutdown() + + +def test_shutdown_rejects_stale_media_tasks(monkeypatch): + monkeypatch.setenv("LANGFUSE_MEDIA_UPLOAD_ENABLED", "true") + + with LangfuseResourceManager._lock: + LangfuseResourceManager._instances.clear() + + client = Langfuse( + public_key="pk-media-shutdown", + secret_key="sk-media-shutdown", + span_exporter=NoOpSpanExporter(), + ) + manager = client._resources + assert manager is not None + + client.shutdown() + + data_uri = "data:text/plain;base64,SGVsbG8=" + processed = manager._media_manager._find_and_process_media( + data=data_uri, + trace_id="0" * 32, + observation_id="0" * 16, + field="input", + ) + + assert processed == data_uri + assert manager._media_upload_queue.unfinished_tasks == 0 + + +def test_reset_handles_shutdown_eviction(monkeypatch): + monkeypatch.setenv("LANGFUSE_MEDIA_UPLOAD_ENABLED", "false") + + with LangfuseResourceManager._lock: + LangfuseResourceManager._instances.clear() + + first_client = Langfuse( + public_key="pk-reset-first", + secret_key="sk-reset-first", + span_exporter=NoOpSpanExporter(), + ) + second_client = Langfuse( + public_key="pk-reset-second", + secret_key="sk-reset-second", + span_exporter=NoOpSpanExporter(), + ) + + LangfuseResourceManager.reset() + + assert LangfuseResourceManager._instances == {} + assert first_client._resources is not None + assert first_client._resources._shutdown + assert second_client._resources is not None + assert second_client._resources._shutdown + + def test_at_fork_reinit_creates_new_queues_and_consumers(monkeypatch): """_at_fork_reinit() must replace queues and start fresh consumer threads.""" monkeypatch.setenv("LANGFUSE_MEDIA_UPLOAD_ENABLED", "false")