Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
132 changes: 88 additions & 44 deletions langfuse/_client/resource_manager.py
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand Down Expand Up @@ -485,46 +486,54 @@ 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:
# 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(
Expand All @@ -546,12 +555,19 @@ def add_trace_task(
event: dict,
) -> None:
try:
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)
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)

except Full:
langfuse_logger.warning(
Expand Down Expand Up @@ -635,13 +651,41 @@ def flush(self) -> None:
langfuse_logger.debug("Successfully flushed media upload queue")

def shutdown(self) -> None:
self._shutdown = True

# Unregister the atexit handler first
atexit.unregister(self.shutdown)

self.flush()
self._stop_and_join_consumer_threads()
with self._lock:
if getattr(self, "_shutdown", False):
return

self._shutdown = True
Comment on lines +655 to +658

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.

P1 Concurrent shutdown returns early

The first caller sets _shutdown before it performs the actual flush and thread joins, so a concurrent caller sees the flag and returns while teardown is still running. This breaks the public shutdown contract that pending data has been flushed and background threads have terminated when the call returns, and can let the second caller release dependent resources too early.

Knowledge Base Used:

Prompt To Fix With AI
This is a comment left during a code review.
Path: langfuse/_client/resource_manager.py
Line: 652-655

Comment:
**Concurrent shutdown returns early**

The first caller sets `_shutdown` before it performs the actual flush and thread joins, so a concurrent caller sees the flag and returns while teardown is still running. This breaks the public shutdown contract that pending data has been flushed and background threads have terminated when the call returns, and can let the second caller release dependent resources too early.

**Knowledge Base Used:**
- [SDK client lifeycle and configuration](https://app.greptile.com/personal-org-4986/-/custom-context/knowledge-base/langfuse/langfuse-python/-/docs/sdk-client-lifecycle.md)
- [Client initialization and resource management](https://app.greptile.com/personal-org-4986/-/custom-context/knowledge-base/langfuse/langfuse-python/-/docs/client-initialization-and-resources.md)

---

For each issue above, determine whether it is valid and should be fixed. If so, fix it directly.


# 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]
Comment on lines +663 to +667

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.

P1 Old processor remains active

After this eviction, creating another client with the same key registers a new LangfuseSpanProcessor on the shared OpenTelemetry provider, but shutdown never removes or shuts down the old processor. Both processors accept spans for that key, so spans created after re-instantiation can be exported twice and the old exporter remains alive.

Knowledge Base Used:

Prompt To Fix With AI
This is a comment left during a code review.
Path: langfuse/_client/resource_manager.py
Line: 660-664

Comment:
**Old processor remains active**

After this eviction, creating another client with the same key registers a new `LangfuseSpanProcessor` on the shared OpenTelemetry provider, but shutdown never removes or shuts down the old processor. Both processors accept spans for that key, so spans created after re-instantiation can be exported twice and the old exporter remains alive.

**Knowledge Base Used:**
- [SDK client lifeycle and configuration](https://app.greptile.com/personal-org-4986/-/custom-context/knowledge-base/langfuse/langfuse-python/-/docs/sdk-client-lifecycle.md)
- [Client initialization and resource management](https://app.greptile.com/personal-org-4986/-/custom-context/knowledge-base/langfuse/langfuse-python/-/docs/client-initialization-and-resources.md)

---

For each issue above, determine whether it is valid and should be fixed. If so, fix it directly.


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(
Expand Down
4 changes: 2 additions & 2 deletions tests/unit/test_prompt.py
Original file line number Diff line number Diff line change
Expand Up @@ -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


Expand Down
94 changes: 94 additions & 0 deletions tests/unit/test_resource_manager.py
Original file line number Diff line number Diff line change
Expand Up @@ -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()