From 0c95eb8360b5acdc4f6fc82a556b25fa4db0e862 Mon Sep 17 00:00:00 2001 From: Himanshu Joshi Date: Wed, 23 Sep 2026 01:08:18 +0530 Subject: [PATCH 1/2] fix(client): evict singleton on shutdown and prevent deadlocks on re-instantiation Signed-off-by: Himanshu Joshi --- langfuse/_client/resource_manager.py | 33 ++++++++-- tests/unit/test_prompt.py | 4 +- tests/unit/test_resource_manager.py | 94 ++++++++++++++++++++++++++++ 3 files changed, 124 insertions(+), 7 deletions(-) diff --git a/langfuse/_client/resource_manager.py b/langfuse/_client/resource_manager.py index a61101fe3..e5c721974 100644 --- a/langfuse/_client/resource_manager.py +++ b/langfuse/_client/resource_manager.py @@ -485,13 +485,20 @@ def _at_fork_reinit(self) -> None: @classmethod def reset(cls) -> None: with cls._lock: - for key in cls._instances: - cls._instances[key].shutdown() + for key in list(cls._instances.keys()): + if key in cls._instances: + cls._instances[key].shutdown() cls._instances.clear() def add_score_task(self, event: dict, *, force_sample: bool = False) -> None: try: + if getattr(self, "_shutdown", False): + langfuse_logger.warning( + "Langfuse client is already shut down. Dropping score event." + ) + return + # Sample scores with the same sampler that is used for tracing tracer_provider = cast(TracerProvider, otel_trace_api.get_tracer_provider()) should_sample = ( @@ -546,6 +553,12 @@ def add_trace_task( event: dict, ) -> None: try: + if getattr(self, "_shutdown", False): + langfuse_logger.warning( + "Langfuse client is already shut down. Dropping trace event." + ) + return + langfuse_logger.debug( "Trace: Enqueuing event type=%s for trace_id=%s", event["type"], @@ -635,10 +648,20 @@ def flush(self) -> None: langfuse_logger.debug("Successfully flushed media upload queue") def shutdown(self) -> None: - self._shutdown = True + with self._lock: + if getattr(self, "_shutdown", False): + return + + self._shutdown = True + + # Unregister the atexit handler first + atexit.unregister(self.shutdown) - # Unregister the atexit handler first - atexit.unregister(self.shutdown) + # Evict from singleton registry so subsequent client initializations + # construct a fresh, active manager instead of reusing a shut down one + if hasattr(self, "public_key") and self.public_key in self._instances: + if self._instances[self.public_key] is self: + del self._instances[self.public_key] self.flush() self._stop_and_join_consumer_threads() diff --git a/tests/unit/test_prompt.py b/tests/unit/test_prompt.py index eadfb8221..e0d14b765 100644 --- a/tests/unit/test_prompt.py +++ b/tests/unit/test_prompt.py @@ -130,12 +130,12 @@ def langfuse(): from langfuse._client.resource_manager import LangfuseResourceManager langfuse_instance = Langfuse() - langfuse_instance.api = Mock() - if langfuse_instance._resources is None: langfuse_instance._resources = Mock(spec=LangfuseResourceManager) langfuse_instance._resources.prompt_cache = PromptCache() + langfuse_instance.api = Mock() + return langfuse_instance diff --git a/tests/unit/test_resource_manager.py b/tests/unit/test_resource_manager.py index f66a1e052..98782c0ed 100644 --- a/tests/unit/test_resource_manager.py +++ b/tests/unit/test_resource_manager.py @@ -442,3 +442,97 @@ def signal_shutdown(self, *, count): ("join", 0), ("join", 1), ] + + +def test_shutdown_evicts_singleton_from_instances(): + """Test that shutdown() removes the manager from _instances to prevent stale singleton reuse.""" + public_key = "pk-test-shutdown-evict" + client = Langfuse( + public_key=public_key, + secret_key="sk-test-secret", + base_url="http://localhost:3000", + span_exporter=NoOpSpanExporter(), + ) + + with LangfuseResourceManager._lock: + assert public_key in LangfuseResourceManager._instances + + rm_first = client._resources + assert rm_first is not None + assert not rm_first._shutdown + + client.shutdown() + + assert rm_first._shutdown + with LangfuseResourceManager._lock: + assert public_key not in LangfuseResourceManager._instances + + # Re-instantiating client with same public_key should create a fresh, active manager + client_new = Langfuse( + public_key=public_key, + secret_key="sk-test-secret", + base_url="http://localhost:3000", + span_exporter=NoOpSpanExporter(), + ) + rm_second = client_new._resources + assert rm_second is not None + assert rm_second is not rm_first + assert not rm_second._shutdown + + client_new.shutdown() + + +def test_reinstantiation_after_shutdown_processes_scores_without_deadlock(): + """Test that re-instantiating a client after shutdown spawns active consumer threads and does not deadlock.""" + public_key = "pk-test-reinstantiation-deadlock" + client1 = Langfuse( + public_key=public_key, + secret_key="sk-test-secret", + base_url="http://localhost:3000", + span_exporter=NoOpSpanExporter(), + ) + client1.shutdown() + + # Create second client with same key (simulates fixture or worker lifecycle) + client2 = Langfuse( + public_key=public_key, + secret_key="sk-test-secret", + base_url="http://localhost:3000", + span_exporter=NoOpSpanExporter(), + ) + + rm2 = client2._resources + assert rm2 is not None + assert not rm2._shutdown + assert len(rm2._ingestion_consumers) > 0 + + # Ensure add_score_task / flush / shutdown completes immediately without blocking indefinitely + fake_score = { + "type": "score-create", + "body": SimpleNamespace(trace_id="0" * 32, name="test_metric", value=1.0), + } + rm2.add_score_task(fake_score, force_sample=True) + + # Calling flush and shutdown must complete without hanging + client2.flush() + client2.shutdown() + + +def test_shutdown_idempotency(): + """Test that multiple calls to shutdown() are safe and idempotent.""" + public_key = "pk-test-shutdown-idempotent" + client = Langfuse( + public_key=public_key, + secret_key="sk-test-secret", + base_url="http://localhost:3000", + span_exporter=NoOpSpanExporter(), + ) + rm = client._resources + assert rm is not None + + client.shutdown() + assert rm._shutdown + + # Calling shutdown again must not raise any exceptions + rm.shutdown() + client.shutdown() From 22e16afe22b50b261e2e28e090659f6d191d6a34 Mon Sep 17 00:00:00 2001 From: Himanshu Joshi Date: Thu, 24 Sep 2026 00:16:37 +0530 Subject: [PATCH 2/2] fix(client): thread-safe task enqueue and span processor shutdown Signed-off-by: Himanshu Joshi --- langfuse/_client/resource_manager.py | 117 ++++++++++++++++----------- 1 file changed, 69 insertions(+), 48 deletions(-) diff --git a/langfuse/_client/resource_manager.py b/langfuse/_client/resource_manager.py index e5c721974..0641617ac 100644 --- a/langfuse/_client/resource_manager.py +++ b/langfuse/_client/resource_manager.py @@ -262,6 +262,7 @@ def _initialize_instance( media_manager=self._media_manager, mask_otel_spans=mask_otel_spans, ) + self._span_processor = langfuse_processor tracer_provider.add_span_processor(langfuse_processor) self._otel_tracer = tracer_provider.get_tracer( @@ -493,45 +494,46 @@ def reset(cls) -> None: def add_score_task(self, event: dict, *, force_sample: bool = False) -> None: try: - if getattr(self, "_shutdown", False): - langfuse_logger.warning( - "Langfuse client is already shut down. Dropping score event." - ) - return - - # Sample scores with the same sampler that is used for tracing - tracer_provider = cast(TracerProvider, otel_trace_api.get_tracer_provider()) - should_sample = ( - force_sample - or isinstance( - tracer_provider, otel_trace_api.ProxyTracerProvider - ) # default to in-sample if otel sampler is not available - or ( - ( - tracer_provider.sampler.should_sample( - parent_context=None, - trace_id=int(event["body"].trace_id, 16), - name="score", - ).decision - == Decision.RECORD_AND_SAMPLE - if hasattr(event["body"], "trace_id") + with self._lock: + if getattr(self, "_shutdown", False): + langfuse_logger.warning( + "Langfuse client is already shut down. Dropping score event." + ) + return + + # Sample scores with the same sampler that is used for tracing + tracer_provider = cast(TracerProvider, otel_trace_api.get_tracer_provider()) + should_sample = ( + force_sample + or isinstance( + tracer_provider, otel_trace_api.ProxyTracerProvider + ) # default to in-sample if otel sampler is not available + or ( + ( + tracer_provider.sampler.should_sample( + parent_context=None, + trace_id=int(event["body"].trace_id, 16), + name="score", + ).decision + == Decision.RECORD_AND_SAMPLE + if hasattr(event["body"], "trace_id") + else True + ) + if event["body"].trace_id + is not None # do not sample out session / dataset run scores else True ) - if event["body"].trace_id - is not None # do not sample out session / dataset run scores - else True ) - ) - if should_sample: - langfuse_logger.debug( - "Score: Enqueuing event type=%s for trace_id=%s name=%s value=%s", - event["type"], - event["body"].trace_id, - event["body"].name, - event["body"].value, - ) - self._score_ingestion_queue.put(event, block=False) + if should_sample: + langfuse_logger.debug( + "Score: Enqueuing event type=%s for trace_id=%s name=%s value=%s", + event["type"], + event["body"].trace_id, + event["body"].name, + event["body"].value, + ) + self._score_ingestion_queue.put(event, block=False) except Full: langfuse_logger.warning( @@ -553,18 +555,19 @@ def add_trace_task( event: dict, ) -> None: try: - if getattr(self, "_shutdown", False): - langfuse_logger.warning( - "Langfuse client is already shut down. Dropping trace event." - ) - return + with self._lock: + if getattr(self, "_shutdown", False): + langfuse_logger.warning( + "Langfuse client is already shut down. Dropping trace event." + ) + return - langfuse_logger.debug( - "Trace: Enqueuing event type=%s for trace_id=%s", - event["type"], - event["body"].id, - ) - self._score_ingestion_queue.put(event, block=False) + langfuse_logger.debug( + "Trace: Enqueuing event type=%s for trace_id=%s", + event["type"], + event["body"].id, + ) + self._score_ingestion_queue.put(event, block=False) except Full: langfuse_logger.warning( @@ -663,8 +666,26 @@ def shutdown(self) -> None: if self._instances[self.public_key] is self: del self._instances[self.public_key] - self.flush() - self._stop_and_join_consumer_threads() + self.flush() + self._stop_and_join_consumer_threads() + + # Shut down and detach span processor + if hasattr(self, "_span_processor") and self._span_processor is not None: + try: + self._span_processor.shutdown() + except Exception as e: + langfuse_logger.debug("Error shutting down span processor: %s", e) + + if self.tracer_provider is not None and hasattr( + self.tracer_provider, "_active_span_processor" + ): + active_proc = self.tracer_provider._active_span_processor + if hasattr(active_proc, "_span_processors"): + active_proc._span_processors = tuple( + p + for p in active_proc._span_processors + if p is not self._span_processor + ) def _init_tracer_provider(