From b2f059c6af2df0254c429af2a8c0cc8978f0a636 Mon Sep 17 00:00:00 2001 From: David Traina <44659830+DavidTraina@users.noreply.github.com> Date: Wed, 29 Jul 2026 23:39:49 -0500 Subject: [PATCH 1/3] fix(observe): close abandoned generators when their span already ended --- langfuse/_client/observe.py | 29 ++++++-- tests/unit/test_observe.py | 129 ++++++++++++++++++++++++++++++++++++ 2 files changed, 151 insertions(+), 7 deletions(-) diff --git a/langfuse/_client/observe.py b/langfuse/_client/observe.py index bd0a3edee..221248149 100644 --- a/langfuse/_client/observe.py +++ b/langfuse/_client/observe.py @@ -616,6 +616,12 @@ def _finalize_with_error(self, error: BaseException) -> None: def close(self) -> None: if self._span_ended: + # The span can end while the generator is still suspended (e.g. + # an error surfaced from __next__ without resuming the generator). + # Still close the generator so its cleanup runs deterministically + # in the preserved context instead of at GC time under an + # arbitrary ambient context. + self.context.run(self.generator.close) return try: @@ -711,22 +717,31 @@ def _finalize_with_error(self, error: BaseException) -> None: async def aclose(self) -> None: if self._span_ended: + # The span can end while the generator is still suspended (e.g. a + # cancellation delivered before the inner __anext__ task resumed + # the generator). Still close the generator so its cleanup runs + # deterministically in the preserved context instead of at GC time + # under an arbitrary ambient context. + await self._close_generator() return try: - try: - await asyncio.create_task( - self.generator.aclose(), - context=self.context, - ) # type: ignore - except TypeError: - await self.context.run(asyncio.create_task, self.generator.aclose()) + await self._close_generator() except (Exception, asyncio.CancelledError) as error: self._finalize_with_error(error) raise else: self._finalize() + async def _close_generator(self) -> None: + try: + await asyncio.create_task( + self.generator.aclose(), + context=self.context, + ) # type: ignore + except TypeError: + await self.context.run(asyncio.create_task, self.generator.aclose()) + async def close(self) -> None: await self.aclose() diff --git a/tests/unit/test_observe.py b/tests/unit/test_observe.py index 5527be9b9..42cc0082d 100644 --- a/tests/unit/test_observe.py +++ b/tests/unit/test_observe.py @@ -1,6 +1,7 @@ import asyncio import contextvars import gc +import inspect import json import sys from typing import Any, AsyncGenerator, Generator, cast @@ -270,6 +271,134 @@ async def generator() -> AsyncGenerator[str, None]: assert span.ended == 1 +def test_sync_generator_wrapper_close_closes_generator_after_span_ended() -> None: + marker = contextvars.ContextVar("marker", default="ambient") + seen: list[str] = [] + + def generator() -> Generator[str, None, None]: + try: + yield "item_0" + yield "item_1" + finally: + seen.append(marker.get()) + + span = SpanRecorder() + context = contextvars.copy_context() + context.run(marker.set, "preserved") + wrapper = _ContextPreservedSyncGeneratorWrapper( + generator(), + context, + cast(Any, span), + False, + None, + ) + + assert next(wrapper) == "item_0" + + # __next__ raising without resuming the generator (here: re-entering the + # preserved context) ends the span while the generator is still suspended. + with pytest.raises(RuntimeError): + context.run(lambda: next(wrapper)) + + assert span.ended == 1 + assert seen == [] + + marker.set("ambient-now") + wrapper.close() + + assert seen == ["preserved"] + assert span.ended == 1 + + +@pytest.mark.asyncio +async def test_async_generator_wrapper_aclose_closes_generator_after_span_ended() -> ( + None +): + marker = contextvars.ContextVar("marker", default="ambient") + seen: list[str] = [] + + async def generator() -> AsyncGenerator[str, None]: + try: + yield "item_0" + yield "item_1" + finally: + seen.append(marker.get()) + + span = SpanRecorder() + context = contextvars.copy_context() + context.run(marker.set, "preserved") + wrapper = _ContextPreservedAsyncGeneratorWrapper( + generator(), + context, + cast(Any, span), + False, + None, + ) + + assert await wrapper.__anext__() == "item_0" + + # Span ends while the generator is still suspended (as happens when a + # cancellation never resumes the generator, see the test below). + wrapper._finalize_with_error(asyncio.CancelledError()) + assert span.ended == 1 + assert seen == [] + + marker.set("ambient-now") + await wrapper.aclose() + + assert seen == ["preserved"] + assert span.ended == 1 + + +@pytest.mark.asyncio +async def test_async_generator_wrapper_closes_generator_cancel_never_resumed() -> None: + marker = contextvars.ContextVar("marker", default="ambient") + seen: list[str] = [] + + async def generator() -> AsyncGenerator[str, None]: + try: + yield "item_0" + yield "item_1" + finally: + seen.append(marker.get()) + + span = SpanRecorder() + context = contextvars.copy_context() + context.run(marker.set, "preserved") + raw = generator() + wrapper = _ContextPreservedAsyncGeneratorWrapper( + raw, + context, + cast(Any, span), + False, + None, + ) + + async def consume() -> None: + async for _ in wrapper: + await asyncio.sleep(0) + + consumer = asyncio.create_task(consume()) + for _ in range(4): + await asyncio.sleep(0) + # The cancel lands between chunks, before the inner __anext__ task's + # first step: the generator is never resumed, but the span is ended. + asyncio.get_running_loop().call_soon(consumer.cancel) + with pytest.raises(asyncio.CancelledError): + await consumer + + assert inspect.getasyncgenstate(raw) == "AGEN_SUSPENDED" + assert span.ended == 1 + assert seen == [] + + marker.set("ambient-now") + await wrapper.aclose() + + assert inspect.getasyncgenstate(raw) == "AGEN_CLOSED" + assert seen == ["preserved"] + assert span.ended == 1 + + @pytest.mark.asyncio async def test_async_generator_wrapper_fallback_preserves_context( monkeypatch: pytest.MonkeyPatch, From 0dfbad85be660cc2a6e7ec685c9a67426cc8fa9e Mon Sep 17 00:00:00 2001 From: David Traina <44659830+DavidTraina@users.noreply.github.com> Date: Thu, 30 Jul 2026 01:00:18 -0500 Subject: [PATCH 2/3] fix(observe): only treat create_task signature errors as fallback trigger --- langfuse/_client/observe.py | 9 +++++++-- tests/unit/test_observe.py | 26 ++++++++++++++++++++++++++ 2 files changed, 33 insertions(+), 2 deletions(-) diff --git a/langfuse/_client/observe.py b/langfuse/_client/observe.py index 221248149..a1528ed7f 100644 --- a/langfuse/_client/observe.py +++ b/langfuse/_client/observe.py @@ -735,12 +735,17 @@ async def aclose(self) -> None: async def _close_generator(self) -> None: try: - await asyncio.create_task( + close_task = asyncio.create_task( self.generator.aclose(), context=self.context, ) # type: ignore except TypeError: - await self.context.run(asyncio.create_task, self.generator.aclose()) + # Python 3.10: create_task has no context parameter. Only the + # create_task call is guarded so a TypeError raised from the + # generator's own cleanup is not mistaken for it. + close_task = self.context.run(asyncio.create_task, self.generator.aclose()) + + await close_task async def close(self) -> None: await self.aclose() diff --git a/tests/unit/test_observe.py b/tests/unit/test_observe.py index 42cc0082d..3f2e9137b 100644 --- a/tests/unit/test_observe.py +++ b/tests/unit/test_observe.py @@ -399,6 +399,32 @@ async def consume() -> None: assert span.ended == 1 +@pytest.mark.asyncio +async def test_async_generator_wrapper_aclose_propagates_cleanup_type_error() -> None: + async def generator() -> AsyncGenerator[str, None]: + try: + yield "item_0" + finally: + raise TypeError("cleanup failed") + + span = SpanRecorder() + wrapper = _ContextPreservedAsyncGeneratorWrapper( + generator(), + contextvars.copy_context(), + cast(Any, span), + False, + None, + ) + + assert await wrapper.__anext__() == "item_0" + + with pytest.raises(TypeError, match="cleanup failed"): + await wrapper.aclose() + + assert span.ended == 1 + assert span.updates[-1] == {"level": "ERROR", "status_message": "cleanup failed"} + + @pytest.mark.asyncio async def test_async_generator_wrapper_fallback_preserves_context( monkeypatch: pytest.MonkeyPatch, From eff94e1afca526897e96cf11d2aae5262da326b5 Mon Sep 17 00:00:00 2001 From: David Traina <44659830+DavidTraina@users.noreply.github.com> Date: Thu, 30 Jul 2026 01:08:20 -0500 Subject: [PATCH 3/3] fix(observe): make cancel-race test version agnostic and trim comments --- langfuse/_client/observe.py | 16 +++------------- tests/unit/test_observe.py | 14 +++++--------- 2 files changed, 8 insertions(+), 22 deletions(-) diff --git a/langfuse/_client/observe.py b/langfuse/_client/observe.py index a1528ed7f..3e3837d8c 100644 --- a/langfuse/_client/observe.py +++ b/langfuse/_client/observe.py @@ -616,11 +616,7 @@ def _finalize_with_error(self, error: BaseException) -> None: def close(self) -> None: if self._span_ended: - # The span can end while the generator is still suspended (e.g. - # an error surfaced from __next__ without resuming the generator). - # Still close the generator so its cleanup runs deterministically - # in the preserved context instead of at GC time under an - # arbitrary ambient context. + # Still close the generator so cleanup runs in the preserved context, not at GC time. self.context.run(self.generator.close) return @@ -717,11 +713,7 @@ def _finalize_with_error(self, error: BaseException) -> None: async def aclose(self) -> None: if self._span_ended: - # The span can end while the generator is still suspended (e.g. a - # cancellation delivered before the inner __anext__ task resumed - # the generator). Still close the generator so its cleanup runs - # deterministically in the preserved context instead of at GC time - # under an arbitrary ambient context. + # Still close the generator so cleanup runs in the preserved context, not at GC time. await self._close_generator() return @@ -740,9 +732,7 @@ async def _close_generator(self) -> None: context=self.context, ) # type: ignore except TypeError: - # Python 3.10: create_task has no context parameter. Only the - # create_task call is guarded so a TypeError raised from the - # generator's own cleanup is not mistaken for it. + # Python 3.10 create_task has no context param; guard only the call itself. close_task = self.context.run(asyncio.create_task, self.generator.aclose()) await close_task diff --git a/tests/unit/test_observe.py b/tests/unit/test_observe.py index 3f2e9137b..4830601fc 100644 --- a/tests/unit/test_observe.py +++ b/tests/unit/test_observe.py @@ -1,7 +1,6 @@ import asyncio import contextvars import gc -import inspect import json import sys from typing import Any, AsyncGenerator, Generator, cast @@ -295,8 +294,7 @@ def generator() -> Generator[str, None, None]: assert next(wrapper) == "item_0" - # __next__ raising without resuming the generator (here: re-entering the - # preserved context) ends the span while the generator is still suspended. + # An error from __next__ that never resumed the generator ends the span. with pytest.raises(RuntimeError): context.run(lambda: next(wrapper)) @@ -337,8 +335,7 @@ async def generator() -> AsyncGenerator[str, None]: assert await wrapper.__anext__() == "item_0" - # Span ends while the generator is still suspended (as happens when a - # cancellation never resumes the generator, see the test below). + # Span ends while the generator is still suspended (e.g. the cancel race below). wrapper._finalize_with_error(asyncio.CancelledError()) assert span.ended == 1 assert seen == [] @@ -381,20 +378,19 @@ async def consume() -> None: consumer = asyncio.create_task(consume()) for _ in range(4): await asyncio.sleep(0) - # The cancel lands between chunks, before the inner __anext__ task's - # first step: the generator is never resumed, but the span is ended. + # Cancel lands before the inner __anext__ task's first step. asyncio.get_running_loop().call_soon(consumer.cancel) with pytest.raises(asyncio.CancelledError): await consumer - assert inspect.getasyncgenstate(raw) == "AGEN_SUSPENDED" + assert raw.ag_frame is not None # still suspended, never resumed assert span.ended == 1 assert seen == [] marker.set("ambient-now") await wrapper.aclose() - assert inspect.getasyncgenstate(raw) == "AGEN_CLOSED" + assert raw.ag_frame is None # closed assert seen == ["preserved"] assert span.ended == 1