Skip to content
Closed
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
1 change: 0 additions & 1 deletion .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -70,7 +70,6 @@ jobs:
fail-fast: false
matrix:
python-version:
- "3.10"
- "3.11"
- "3.12"
- "3.13"
Expand Down
2 changes: 1 addition & 1 deletion AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -109,7 +109,7 @@ Minimum verification matrix:
The main CI workflow currently runs:
- linting on Python 3.13
- mypy on Python 3.13
- `tests/unit` on a Python 3.10-3.14 matrix
- `tests/unit` on a Python 3.11-3.14 matrix
- `tests/e2e` in 2 mechanical shards plus a serial subset inside each shard
- `tests/live_provider` as one always-on suite
- PR title validation for Conventional Commits
Expand Down
29 changes: 6 additions & 23 deletions langfuse/_client/observe.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,6 @@
import contextvars
import inspect
import os
import sys
from functools import wraps
from typing import (
Any,
Expand Down Expand Up @@ -49,8 +48,6 @@
P = ParamSpec("P")
R = TypeVar("R")

_ASYNCIO_CREATE_TASK_SUPPORTS_CONTEXT = sys.version_info >= (3, 11)


class LangfuseDecorator:
"""Implementation of the @observe decorator for seamless Langfuse tracing integration.
Expand Down Expand Up @@ -217,7 +214,7 @@ def decorator(func: F) -> F:
capture_output=should_capture_output,
transform_to_string=transform_to_string,
)
if asyncio.iscoroutinefunction(func)
if inspect.iscoroutinefunction(func)
else self._sync_observe(
func,
name=name,
Expand Down Expand Up @@ -747,15 +744,7 @@ async def aclose(self) -> None:
self._finalize()

async def _close_generator(self) -> None:
if _ASYNCIO_CREATE_TASK_SUPPORTS_CONTEXT:
close_task = asyncio.create_task(
self.generator.aclose(),
context=self.context,
) # type: ignore
else:
close_task = self.context.run(asyncio.create_task, self.generator.aclose())

await close_task
await asyncio.create_task(self.generator.aclose(), context=self.context)

async def close(self) -> None:
await self.aclose()
Expand All @@ -769,16 +758,10 @@ def __del__(self) -> None:
async def __anext__(self) -> Any:
try:
# Run the generator's __anext__ in the preserved context
if _ASYNCIO_CREATE_TASK_SUPPORTS_CONTEXT:
item = await asyncio.create_task(
self.generator.__anext__(), # type: ignore
context=self.context,
) # type: ignore
else:
item = await self.context.run(
asyncio.create_task,
self.generator.__anext__(), # type: ignore
)
item = await asyncio.create_task(
self.generator.__anext__(), # type: ignore
context=self.context,
)

if self.capture_output:
self.items.append(item)
Expand Down
7 changes: 1 addition & 6 deletions langfuse/types.py
Original file line number Diff line number Diff line change
Expand Up @@ -24,19 +24,14 @@ def my_evaluator(*, output: str, **kwargs) -> Evaluation:
Dict,
Literal,
Mapping,
NotRequired,
Optional,
Protocol,
Sequence,
TypedDict,
Union,
)

try:
from typing import NotRequired # type: ignore
except ImportError:
from typing_extensions import NotRequired


from langfuse.api import MediaContentType

# Span attribute values accepted by the OpenTelemetry trace API. OpenTelemetry 1.45
Expand Down
4 changes: 2 additions & 2 deletions pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,7 @@ readme = "README.md"
authors = [{ name = "langfuse", email = "developers@langfuse.com" }]
license = "MIT"
license-files = ["LICENSE"]
requires-python = ">=3.10,<4.0"
requires-python = ">=3.11,<4.0"
keywords = [
"langfuse",
"llm",
Expand Down Expand Up @@ -139,7 +139,7 @@ module = "langfuse.api.*"
disable_error_code = ["redundant-cast", "no-untyped-def"]

[tool.ruff]
target-version = "py310"
target-version = "py311"

[tool.ruff.lint]
extend-select = [
Expand Down
4 changes: 0 additions & 4 deletions tests/e2e/test_decorators.py
Original file line number Diff line number Diff line change
@@ -1,6 +1,5 @@
import asyncio
import os
import sys
from collections import defaultdict
from concurrent.futures import ThreadPoolExecutor
from time import sleep
Expand Down Expand Up @@ -1880,7 +1879,6 @@ def root_function():


@pytest.mark.asyncio
@pytest.mark.skipif(sys.version_info < (3, 11), reason="requires python3.11 or higher")
async def test_async_generator_context_preservation():
"""Test that async generators preserve context when consumed later (e.g., by streaming responses)"""
langfuse = get_client()
Expand Down Expand Up @@ -1946,7 +1944,6 @@ async def root_function():


@pytest.mark.asyncio
@pytest.mark.skipif(sys.version_info < (3, 11), reason="requires python3.11 or higher")
async def test_async_generator_context_preservation_with_trace_hierarchy():
"""Test that async generators maintain proper parent-child span relationships"""
langfuse = get_client()
Expand Down Expand Up @@ -2009,7 +2006,6 @@ async def parent_function():


@pytest.mark.asyncio
@pytest.mark.skipif(sys.version_info < (3, 11), reason="requires python3.11 or higher")
async def test_async_generator_exception_handling_with_context():
"""Test that exceptions in async generators are properly handled while preserving context"""
langfuse = get_client()
Expand Down
2 changes: 1 addition & 1 deletion tests/unit/test_media.py
Original file line number Diff line number Diff line change
Expand Up @@ -263,6 +263,7 @@ def test_resolve_media_references_uses_configured_httpx_client():
"https://example.com/test.jpg", timeout=fetch_timeout_seconds
)


def test_init_with_urlsafe_base64_data_uri():
original_bytes = b"\xfb\xff"
urlsafe_base64 = base64.urlsafe_b64encode(original_bytes).decode()
Expand All @@ -274,4 +275,3 @@ def test_init_with_urlsafe_base64_data_uri():
assert media._source == "base64_data_uri"
assert media._content_type == "application/octet-stream"
assert media._content_bytes == original_bytes

38 changes: 0 additions & 38 deletions tests/unit/test_observe.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,13 +3,11 @@
import gc
import inspect
import json
import sys
from typing import Any, AsyncGenerator, Generator, cast

import pytest

from langfuse import observe
from langfuse._client import observe as observe_module
from langfuse._client.attributes import LangfuseOtelSpanAttributes
from langfuse._client.observe import (
_ContextPreservedAsyncGeneratorWrapper,
Expand Down Expand Up @@ -96,7 +94,6 @@ def body() -> Generator[str, None, None]:


@pytest.mark.asyncio
@pytest.mark.skipif(sys.version_info < (3, 11), reason="requires python3.11 or higher")
async def test_streaming_response_preserves_context_without_output_capture(
langfuse_memory_client: Any, memory_exporter: Any
) -> None:
Expand Down Expand Up @@ -479,41 +476,6 @@ async def generator() -> AsyncGenerator[str, None]:
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,
) -> None:
marker = contextvars.ContextVar("marker", default="ambient")
seen: list[str] = []
monkeypatch.setattr(observe_module, "_ASYNCIO_CREATE_TASK_SUPPORTS_CONTEXT", False)

async def generator() -> AsyncGenerator[str, None]:
try:
yield marker.get()
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__() == "preserved"
marker.set("ambient-now")

await wrapper.aclose()

assert seen == ["preserved"]
assert span.ended == 1


@pytest.mark.asyncio
async def test_async_generator_wrapper_del_ends_span_when_abandoned() -> None:
async def generator() -> AsyncGenerator[str, None]:
Expand Down
Loading
Loading