From e174598e20b6509f7a2c54f6764603179e9bcbb9 Mon Sep 17 00:00:00 2001 From: Bortlesboat <169967362+Bortlesboat@users.noreply.github.com> Date: Tue, 8 Sep 2026 09:56:34 -0400 Subject: [PATCH 1/2] fix(experiments): await future and custom awaitable results --- langfuse/experiment.py | 6 ++-- tests/unit/test_experiment.py | 65 +++++++++++++++++++++++++++++++++++ 2 files changed, 68 insertions(+), 3 deletions(-) diff --git a/langfuse/experiment.py b/langfuse/experiment.py index aa5481829..4f6be08c5 100644 --- a/langfuse/experiment.py +++ b/langfuse/experiment.py @@ -5,7 +5,7 @@ and result formatting. """ -import asyncio +import inspect from datetime import datetime from typing import ( TYPE_CHECKING, @@ -1007,7 +1007,7 @@ async def _run_evaluator( result = evaluator(**kwargs) # Handle async evaluators - if asyncio.iscoroutine(result): + if inspect.isawaitable(result): result = await result return _normalize_evaluator_result(result) @@ -1026,7 +1026,7 @@ async def _run_task(task: TaskFunction, item: ExperimentItem) -> Any: result = task(item=item) # Handle async tasks - if asyncio.iscoroutine(result): + if inspect.isawaitable(result): result = await result return result diff --git a/tests/unit/test_experiment.py b/tests/unit/test_experiment.py index 3eb55ad40..5f90a1e87 100644 --- a/tests/unit/test_experiment.py +++ b/tests/unit/test_experiment.py @@ -1,5 +1,6 @@ """Tests for ``langfuse.experiment`` — ``RunnerContext`` and ``RegressionError``.""" +import asyncio import inspect import typing from datetime import datetime @@ -15,6 +16,70 @@ from langfuse.batch_evaluation import CompositeEvaluatorFunction +@pytest.fixture(params=["value", "coroutine", "future", "task", "awaitable"]) +def wrap_result(request): + """Return equivalent results through each supported awaitable shape.""" + + def wrap(value): + async def resolve(): + return value + + class CustomAwaitable: + def __await__(self): + return resolve().__await__() + + if request.param == "value": + return value + if request.param == "coroutine": + return resolve() + if request.param == "task": + return asyncio.create_task(resolve()) + if request.param == "awaitable": + return CustomAwaitable() + loop = asyncio.get_running_loop() + future = loop.create_future() + loop.call_soon(future.set_result, value) + return future + + return wrap + + +class TestExperimentAwaitableResults: + def test_task_output_is_resolved(self, langfuse_memory_client, wrap_result): + def task(*, item): + return wrap_result("answer") + + result = langfuse_memory_client.run_experiment( + name="awaitable-task", data=[{"input": "question"}], task=task + ) + + assert result.item_results[0].output == "answer" + + def test_item_and_run_evaluations_are_resolved( + self, langfuse_memory_client, wrap_result, monkeypatch + ): + monkeypatch.setattr(langfuse_memory_client, "create_score", MagicMock()) + + def task(*, item): + return "answer" + + def evaluator(**kwargs): + return wrap_result({"name": "quality", "value": 1.0}) + + result = langfuse_memory_client.run_experiment( + name="awaitable-evaluators", + data=[{"input": "question"}], + task=task, + evaluators=[evaluator], + run_evaluators=[evaluator], + ) + + assert [(e.name, e.value) for e in result.item_results[0].evaluations] == [ + ("quality", 1.0) + ] + assert [(e.name, e.value) for e in result.run_evaluations] == [("quality", 1.0)] + + def _noop_task(*, item, **kwargs): # pragma: no cover - never invoked via mock return None From 46db9afeafcfb2cf5485bd7e6dba7dfbaefe0057 Mon Sep 17 00:00:00 2001 From: Bortlesboat <169967362+Bortlesboat@users.noreply.github.com> Date: Tue, 8 Sep 2026 13:16:27 -0400 Subject: [PATCH 2/2] fix(experiments): await composite evaluator results Address review feedback on #1869 and preserve composite scores from Future, Task and custom awaitable callbacks. Note: pre-existing failures in prompt, prompt atexit, and Windows path tests are not addressed by this PR. --- langfuse/_client/client.py | 3 ++- tests/unit/test_experiment.py | 32 ++++++++++++++++++++++++++++++++ 2 files changed, 34 insertions(+), 1 deletion(-) diff --git a/langfuse/_client/client.py b/langfuse/_client/client.py index 42d861fd4..e19445e0a 100644 --- a/langfuse/_client/client.py +++ b/langfuse/_client/client.py @@ -4,6 +4,7 @@ """ import asyncio +import inspect import logging import os import re @@ -3183,7 +3184,7 @@ async def _process_experiment_item( evaluations=evaluations, ) - if asyncio.iscoroutine(result): + if inspect.isawaitable(result): result = await result composite_evals = _normalize_evaluator_result( diff --git a/tests/unit/test_experiment.py b/tests/unit/test_experiment.py index 5f90a1e87..55d23b8f1 100644 --- a/tests/unit/test_experiment.py +++ b/tests/unit/test_experiment.py @@ -79,6 +79,38 @@ def evaluator(**kwargs): ] assert [(e.name, e.value) for e in result.run_evaluations] == [("quality", 1.0)] + def test_composite_evaluation_is_resolved_and_scored( + self, langfuse_memory_client, wrap_result, monkeypatch + ): + create_score = MagicMock() + monkeypatch.setattr(langfuse_memory_client, "create_score", create_score) + + def task(*, item): + return "answer" + + def evaluator(**kwargs): + return {"name": "quality", "value": 1.0} + + def composite_evaluator(*, evaluations, **kwargs): + return wrap_result({"name": "aggregate", "value": evaluations[0].value}) + + result = langfuse_memory_client.run_experiment( + name="awaitable-composite", + data=[{"input": "question"}], + task=task, + evaluators=[evaluator], + composite_evaluator=composite_evaluator, + ) + + assert [(e.name, e.value) for e in result.item_results[0].evaluations] == [ + ("quality", 1.0), + ("aggregate", 1.0), + ] + assert [ + (call.kwargs["name"], call.kwargs["value"]) + for call in create_score.call_args_list + ] == [("quality", 1.0), ("aggregate", 1.0)] + def _noop_task(*, item, **kwargs): # pragma: no cover - never invoked via mock return None