From 78c12d60eeb437c4628a00ecd5956cb686a92934 Mon Sep 17 00:00:00 2001 From: 1fanwang <1fannnw@gmail.com> Date: Sat, 19 Sep 2026 12:02:17 -0700 Subject: [PATCH 1/4] fix(observe): preserve generator control methods Signed-off-by: 1fanwang <1fannnw@gmail.com> --- langfuse/_client/observe.py | 43 ++++++++-- tests/unit/test_observe.py | 167 +++++++++++++++++++++++++++++++++++- 2 files changed, 200 insertions(+), 10 deletions(-) diff --git a/langfuse/_client/observe.py b/langfuse/_client/observe.py index 848506a46..e84ee0869 100644 --- a/langfuse/_client/observe.py +++ b/langfuse/_client/observe.py @@ -8,6 +8,7 @@ Any, AsyncGenerator, Callable, + Coroutine, Dict, Generator, Iterable, @@ -561,7 +562,7 @@ def _handle_observe_result( class _ContextPreservedSyncGeneratorWrapper: - """Sync generator wrapper that ensures each iteration runs in preserved context.""" + """Preserve tracing context across synchronous generator operations.""" def __init__( self, @@ -640,9 +641,19 @@ def __del__(self) -> None: pass def __next__(self) -> Any: + return self._advance(method=self.generator.__next__) + + def send(self, value: Any) -> Any: + return self._advance(method=self.generator.send, args=(value,)) + + def throw(self, *args: Any) -> Any: + return self._advance(method=self.generator.throw, args=args) + + def _advance( + self, *, method: Callable[..., Any], args: Tuple[Any, ...] = () + ) -> Any: try: - # Run the generator's __next__ in the preserved context - item = self.context.run(next, self.generator) + item: Any = self.context.run(method, *args) if self.capture_output: self.items.append(item) @@ -658,7 +669,7 @@ def __next__(self) -> Any: class _ContextPreservedAsyncGeneratorWrapper: - """Async generator wrapper that ensures each iteration runs in preserved context.""" + """Preserve tracing context across asynchronous generator operations.""" def __init__( self, @@ -767,17 +778,31 @@ def __del__(self) -> None: self._finalize() async def __anext__(self) -> Any: + return await self._advance(method=self.generator.__anext__) + + async def asend(self, value: Any) -> Any: + return await self._advance(method=self.generator.asend, args=(value,)) + + async def athrow(self, *args: Any) -> Any: + return await self._advance(method=self.generator.athrow, args=args) + + async def _advance( + self, + *, + method: Callable[..., Coroutine[Any, Any, Any]], + args: Tuple[Any, ...] = (), + ) -> Any: try: - # Run the generator's __anext__ in the preserved context + operation: Coroutine[Any, Any, Any] = method(*args) if _ASYNCIO_CREATE_TASK_SUPPORTS_CONTEXT: - item = await asyncio.create_task( - self.generator.__anext__(), # type: ignore + item: Any = await asyncio.create_task( + coro=operation, context=self.context, - ) # type: ignore + ) else: item = await self.context.run( asyncio.create_task, - self.generator.__anext__(), # type: ignore + operation, ) if self.capture_output: diff --git a/tests/unit/test_observe.py b/tests/unit/test_observe.py index f2ff11789..837ef9b0b 100644 --- a/tests/unit/test_observe.py +++ b/tests/unit/test_observe.py @@ -4,17 +4,20 @@ import inspect import json import sys +from contextlib import asynccontextmanager, contextmanager, nullcontext from typing import Any, AsyncGenerator, Generator, cast import pytest +from opentelemetry.trace import StatusCode, get_current_span -from langfuse import observe +from langfuse import Langfuse, observe from langfuse._client import observe as observe_module from langfuse._client.attributes import LangfuseOtelSpanAttributes from langfuse._client.observe import ( _ContextPreservedAsyncGeneratorWrapper, _ContextPreservedSyncGeneratorWrapper, ) +from tests.conftest import InMemorySpanExporter class SpanRecorder: @@ -35,6 +38,168 @@ def _finished_spans_by_name(memory_exporter: Any, name: str) -> list[Any]: return [span for span in memory_exporter.get_finished_spans() if span.name == name] +@pytest.mark.parametrize("suppress", [False, True]) +def test_observed_context_manager_preserves_exception_handling( + langfuse_memory_client: Langfuse, + memory_exporter: InMemorySpanExporter, + suppress: bool, +) -> None: + closed: list[bool] = [] + + @contextmanager + @observe(capture_output=False) + def resource() -> Generator[None, None, None]: + try: + yield + except ValueError: + if not suppress: + raise + finally: + closed.append(True) + + manager = resource() + try: + with ( + nullcontext() + if suppress + else pytest.raises(ValueError, match="application failed") + ): + with manager: + raise ValueError("application failed") + + assert closed == [True] + langfuse_memory_client.flush() + spans = memory_exporter.get_finished_spans() + assert len(spans) == 1 + assert spans[0].status.status_code == ( + StatusCode.UNSET if suppress else StatusCode.ERROR + ) + finally: + manager.gen.close() + + +@pytest.mark.asyncio +@pytest.mark.parametrize("suppress", [False, True]) +async def test_observed_async_context_manager_preserves_exception_handling( + langfuse_memory_client: Langfuse, + memory_exporter: InMemorySpanExporter, + suppress: bool, +) -> None: + closed: list[bool] = [] + + @asynccontextmanager + @observe(capture_output=False) + async def resource() -> AsyncGenerator[None, None]: + try: + yield + except ValueError: + if not suppress: + raise + finally: + closed.append(True) + + manager = resource() + try: + with ( + nullcontext() + if suppress + else pytest.raises(ValueError, match="application failed") + ): + async with manager: + raise ValueError("application failed") + + assert closed == [True] + langfuse_memory_client.flush() + spans = memory_exporter.get_finished_spans() + assert len(spans) == 1 + assert spans[0].status.status_code == ( + StatusCode.UNSET if suppress else StatusCode.ERROR + ) + finally: + await manager.gen.aclose() + + +def test_observed_generator_send_and_throw_preserve_context_and_output( + langfuse_memory_client: Langfuse, + memory_exporter: InMemorySpanExporter, +) -> None: + observation_ids: list[str | None] = [] + + @observe() + def stream() -> Generator[str, str, None]: + observation_ids.append(langfuse_memory_client.get_current_observation_id()) + try: + value = yield "ready" + observation_ids.append(langfuse_memory_client.get_current_observation_id()) + yield value + except ValueError: + observation_ids.append(langfuse_memory_client.get_current_observation_id()) + yield "recovered" + + generator = stream() + try: + assert next(generator) == "ready" + assert generator.send("sent") == "sent" + assert generator.throw(ValueError("recover")) == "recovered" + with pytest.raises(StopIteration): + next(generator) + + langfuse_memory_client.flush() + spans = memory_exporter.get_finished_spans() + assert len(spans) == 1 + assert len(observation_ids) == 3 + assert observation_ids[0] is not None + assert len(set(observation_ids)) == 1 + assert not get_current_span().get_span_context().is_valid + assert ( + (spans[0].attributes[LangfuseOtelSpanAttributes.OBSERVATION_OUTPUT]) + == "readysentrecovered" + ) + finally: + generator.close() + + +@pytest.mark.asyncio +async def test_observed_async_generator_send_and_throw_preserve_context_and_output( + langfuse_memory_client: Langfuse, + memory_exporter: InMemorySpanExporter, +) -> None: + observation_ids: list[str | None] = [] + + @observe() + async def stream() -> AsyncGenerator[str, str]: + observation_ids.append(langfuse_memory_client.get_current_observation_id()) + try: + value = yield "ready" + observation_ids.append(langfuse_memory_client.get_current_observation_id()) + yield value + except ValueError: + observation_ids.append(langfuse_memory_client.get_current_observation_id()) + yield "recovered" + + generator = stream() + try: + assert await generator.__anext__() == "ready" + assert await generator.asend("sent") == "sent" + assert await generator.athrow(ValueError("recover")) == "recovered" + with pytest.raises(StopAsyncIteration): + await generator.__anext__() + + langfuse_memory_client.flush() + spans = memory_exporter.get_finished_spans() + assert len(spans) == 1 + assert len(observation_ids) == 3 + assert observation_ids[0] is not None + assert len(set(observation_ids)) == 1 + assert not get_current_span().get_span_context().is_valid + assert ( + (spans[0].attributes[LangfuseOtelSpanAttributes.OBSERVATION_OUTPUT]) + == "readysentrecovered" + ) + finally: + await generator.aclose() + + @pytest.mark.asyncio async def test_capture_output_false_preserves_type_when_current_span_is_updated( langfuse_memory_client: Any, memory_exporter: Any From 8494efe9650596cc9e23b382931059c658f35bf0 Mon Sep 17 00:00:00 2001 From: 1fanwang <1fannnw@gmail.com> Date: Sat, 19 Sep 2026 12:46:56 -0700 Subject: [PATCH 2/4] test(observe): exercise cleanup with the real SDK client Signed-off-by: 1fanwang <1fannnw@gmail.com> --- tests/unit/test_observe.py | 123 +++++++++++++++++++++++++------------ 1 file changed, 85 insertions(+), 38 deletions(-) diff --git a/tests/unit/test_observe.py b/tests/unit/test_observe.py index 837ef9b0b..2fea13ddf 100644 --- a/tests/unit/test_observe.py +++ b/tests/unit/test_observe.py @@ -3,11 +3,14 @@ import gc import inspect import json +import logging import sys -from contextlib import asynccontextmanager, contextmanager, nullcontext -from typing import Any, AsyncGenerator, Generator, cast +from contextlib import asynccontextmanager, contextmanager +from tempfile import TemporaryFile +from typing import Any, AsyncGenerator, BinaryIO, Generator, cast import pytest +from opentelemetry.sdk.trace import TracerProvider from opentelemetry.trace import StatusCode, get_current_span from langfuse import Langfuse, observe @@ -38,38 +41,70 @@ def _finished_spans_by_name(memory_exporter: Any, name: str) -> list[Any]: return [span for span in memory_exporter.get_finished_spans() if span.name == name] +@pytest.fixture +def native_memory_client( + monkeypatch: pytest.MonkeyPatch, memory_exporter: InMemorySpanExporter +) -> Generator[Langfuse, None, None]: + monkeypatch.setenv(name="LANGFUSE_PUBLIC_KEY", value="test-generator-protocol") + client = Langfuse( + public_key="test-generator-protocol", + secret_key="test-secret-key", + base_url="http://127.0.0.1:9", + tracer_provider=TracerProvider(), + span_exporter=memory_exporter, + tracing_enabled=True, + sample_rate=1.0, + ) + try: + yield client + finally: + client.shutdown() + + @pytest.mark.parametrize("suppress", [False, True]) def test_observed_context_manager_preserves_exception_handling( - langfuse_memory_client: Langfuse, + native_memory_client: Langfuse, memory_exporter: InMemorySpanExporter, suppress: bool, ) -> None: - closed: list[bool] = [] + opened_file: BinaryIO | None = None + caught_error: Exception | None = None @contextmanager @observe(capture_output=False) - def resource() -> Generator[None, None, None]: - try: - yield - except ValueError: - if not suppress: - raise - finally: - closed.append(True) + def resource() -> Generator[BinaryIO, None, None]: + with TemporaryFile() as handle: + try: + yield handle + except ValueError: + if not suppress: + raise manager = resource() try: - with ( - nullcontext() - if suppress - else pytest.raises(ValueError, match="application failed") - ): - with manager: + try: + with manager as handle: + opened_file = handle raise ValueError("application failed") + except Exception as error: + caught_error = error - assert closed == [True] - langfuse_memory_client.flush() + assert opened_file is not None + native_memory_client.flush() spans = memory_exporter.get_finished_spans() + logging.getLogger(__name__).info( + "sync suppress=%s error=%r resource_closed=%s finished_spans=%s", + suppress, + caught_error, + opened_file.closed, + len(spans), + ) + if suppress: + assert caught_error is None + else: + assert isinstance(caught_error, ValueError) + assert str(caught_error) == "application failed" + assert opened_file.closed assert len(spans) == 1 assert spans[0].status.status_code == ( StatusCode.UNSET if suppress else StatusCode.ERROR @@ -81,36 +116,48 @@ def resource() -> Generator[None, None, None]: @pytest.mark.asyncio @pytest.mark.parametrize("suppress", [False, True]) async def test_observed_async_context_manager_preserves_exception_handling( - langfuse_memory_client: Langfuse, + native_memory_client: Langfuse, memory_exporter: InMemorySpanExporter, suppress: bool, ) -> None: - closed: list[bool] = [] + opened_file: BinaryIO | None = None + caught_error: Exception | None = None @asynccontextmanager @observe(capture_output=False) - async def resource() -> AsyncGenerator[None, None]: - try: - yield - except ValueError: - if not suppress: - raise - finally: - closed.append(True) + async def resource() -> AsyncGenerator[BinaryIO, None]: + with TemporaryFile() as handle: + try: + yield handle + except ValueError: + if not suppress: + raise manager = resource() try: - with ( - nullcontext() - if suppress - else pytest.raises(ValueError, match="application failed") - ): - async with manager: + try: + async with manager as handle: + opened_file = handle raise ValueError("application failed") + except Exception as error: + caught_error = error - assert closed == [True] - langfuse_memory_client.flush() + assert opened_file is not None + native_memory_client.flush() spans = memory_exporter.get_finished_spans() + logging.getLogger(__name__).info( + "async suppress=%s error=%r resource_closed=%s finished_spans=%s", + suppress, + caught_error, + opened_file.closed, + len(spans), + ) + if suppress: + assert caught_error is None + else: + assert isinstance(caught_error, ValueError) + assert str(caught_error) == "application failed" + assert opened_file.closed assert len(spans) == 1 assert spans[0].status.status_code == ( StatusCode.UNSET if suppress else StatusCode.ERROR From 968062a1dbe4406f9bdaaa110f89d9ab62bcdc3c Mon Sep 17 00:00:00 2001 From: 1fanwang <1fannnw@gmail.com> Date: Sat, 19 Sep 2026 13:32:45 -0700 Subject: [PATCH 3/4] fix(observe): finalize spans for process-exit exceptions Signed-off-by: 1fanwang <1fannnw@gmail.com> --- langfuse/_client/observe.py | 26 +++++++++++++++++++++----- tests/unit/test_observe.py | 32 ++++++++++++++++++++------------ 2 files changed, 41 insertions(+), 17 deletions(-) diff --git a/langfuse/_client/observe.py b/langfuse/_client/observe.py index e84ee0869..3b4057573 100644 --- a/langfuse/_client/observe.py +++ b/langfuse/_client/observe.py @@ -663,7 +663,7 @@ def _advance( self._finalize() raise # Re-raise StopIteration - except (Exception, asyncio.CancelledError) as e: + except BaseException as e: self._finalize_with_error(e) raise @@ -786,6 +786,15 @@ async def asend(self, value: Any) -> Any: async def athrow(self, *args: Any) -> Any: return await self._advance(method=self.generator.athrow, args=args) + async def _run_operation( + self, operation: Coroutine[Any, Any, Any] + ) -> Tuple[Any, Optional[BaseException]]: + try: + return await operation, None + except (KeyboardInterrupt, SystemExit) as error: + # Tasks otherwise re-raise these before their awaiter can handle them. + return None, error + async def _advance( self, *, @@ -793,18 +802,25 @@ async def _advance( args: Tuple[Any, ...] = (), ) -> Any: try: - operation: Coroutine[Any, Any, Any] = method(*args) + operation: Coroutine[Any, Any, Tuple[Any, Optional[BaseException]]] = ( + self._run_operation(method(*args)) + ) + item: Any + error: Optional[BaseException] if _ASYNCIO_CREATE_TASK_SUPPORTS_CONTEXT: - item: Any = await asyncio.create_task( + item, error = await asyncio.create_task( coro=operation, context=self.context, ) else: - item = await self.context.run( + item, error = await self.context.run( asyncio.create_task, operation, ) + if error is not None: + raise error + if self.capture_output: self.items.append(item) @@ -820,6 +836,6 @@ async def _advance( raise self._finalize_with_error(e) raise - except Exception as e: + except BaseException as e: self._finalize_with_error(e) raise diff --git a/tests/unit/test_observe.py b/tests/unit/test_observe.py index 2fea13ddf..58e965f4f 100644 --- a/tests/unit/test_observe.py +++ b/tests/unit/test_observe.py @@ -62,13 +62,18 @@ def native_memory_client( @pytest.mark.parametrize("suppress", [False, True]) +@pytest.mark.parametrize( + "error_type", [ValueError, BaseException, KeyboardInterrupt, SystemExit] +) def test_observed_context_manager_preserves_exception_handling( native_memory_client: Langfuse, memory_exporter: InMemorySpanExporter, suppress: bool, + error_type: type[BaseException], ) -> None: opened_file: BinaryIO | None = None - caught_error: Exception | None = None + caught_error: BaseException | None = None + application_error = error_type("application failed") @contextmanager @observe(capture_output=False) @@ -76,7 +81,7 @@ def resource() -> Generator[BinaryIO, None, None]: with TemporaryFile() as handle: try: yield handle - except ValueError: + except error_type: if not suppress: raise @@ -85,8 +90,8 @@ def resource() -> Generator[BinaryIO, None, None]: try: with manager as handle: opened_file = handle - raise ValueError("application failed") - except Exception as error: + raise application_error + except BaseException as error: caught_error = error assert opened_file is not None @@ -102,8 +107,7 @@ def resource() -> Generator[BinaryIO, None, None]: if suppress: assert caught_error is None else: - assert isinstance(caught_error, ValueError) - assert str(caught_error) == "application failed" + assert caught_error is application_error assert opened_file.closed assert len(spans) == 1 assert spans[0].status.status_code == ( @@ -115,13 +119,18 @@ def resource() -> Generator[BinaryIO, None, None]: @pytest.mark.asyncio @pytest.mark.parametrize("suppress", [False, True]) +@pytest.mark.parametrize( + "error_type", [ValueError, BaseException, KeyboardInterrupt, SystemExit] +) async def test_observed_async_context_manager_preserves_exception_handling( native_memory_client: Langfuse, memory_exporter: InMemorySpanExporter, suppress: bool, + error_type: type[BaseException], ) -> None: opened_file: BinaryIO | None = None - caught_error: Exception | None = None + caught_error: BaseException | None = None + application_error = error_type("application failed") @asynccontextmanager @observe(capture_output=False) @@ -129,7 +138,7 @@ async def resource() -> AsyncGenerator[BinaryIO, None]: with TemporaryFile() as handle: try: yield handle - except ValueError: + except error_type: if not suppress: raise @@ -138,8 +147,8 @@ async def resource() -> AsyncGenerator[BinaryIO, None]: try: async with manager as handle: opened_file = handle - raise ValueError("application failed") - except Exception as error: + raise application_error + except BaseException as error: caught_error = error assert opened_file is not None @@ -155,8 +164,7 @@ async def resource() -> AsyncGenerator[BinaryIO, None]: if suppress: assert caught_error is None else: - assert isinstance(caught_error, ValueError) - assert str(caught_error) == "application failed" + assert caught_error is application_error assert opened_file.closed assert len(spans) == 1 assert spans[0].status.status_code == ( From 6b06319f633d72bb4e37ee27a0a59765eabc1771 Mon Sep 17 00:00:00 2001 From: 1fanwang <1fannnw@gmail.com> Date: Sat, 19 Sep 2026 13:36:17 -0700 Subject: [PATCH 4/4] fix(observe): defer generator operations until task execution Signed-off-by: 1fanwang <1fannnw@gmail.com> --- langfuse/_client/observe.py | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/langfuse/_client/observe.py b/langfuse/_client/observe.py index 3b4057573..25b752b6f 100644 --- a/langfuse/_client/observe.py +++ b/langfuse/_client/observe.py @@ -787,10 +787,10 @@ async def athrow(self, *args: Any) -> Any: return await self._advance(method=self.generator.athrow, args=args) async def _run_operation( - self, operation: Coroutine[Any, Any, Any] + self, method: Callable[..., Coroutine[Any, Any, Any]], args: Tuple[Any, ...] ) -> Tuple[Any, Optional[BaseException]]: try: - return await operation, None + return await method(*args), None except (KeyboardInterrupt, SystemExit) as error: # Tasks otherwise re-raise these before their awaiter can handle them. return None, error @@ -803,7 +803,7 @@ async def _advance( ) -> Any: try: operation: Coroutine[Any, Any, Tuple[Any, Optional[BaseException]]] = ( - self._run_operation(method(*args)) + self._run_operation(method=method, args=args) ) item: Any error: Optional[BaseException]