codekingpro/portable-devtools
114k
1"""V2 Evaluation Interface."""2 3from __future__ import annotations4 5import asyncio6import concurrent.futures as cf7import contextvars8import io9import logging10import pathlib11import uuid12from collections.abc import AsyncIterable, AsyncIterator, Awaitable, Iterable, Sequence13from typing import (14 TYPE_CHECKING,15 Any,16 Callable,17 Literal,18 Optional,19 TypeVar,20 Union,21 cast,22)23 24import langsmith25from langsmith import run_helpers as rh26from langsmith import run_trees, schemas27from langsmith import run_trees as rt28from langsmith import utils as ls_utils29from langsmith._internal import _aiter as aitertools30from langsmith._internal._beta_decorator import _warn_once31from langsmith.evaluation._runner import (32 AEVALUATOR_T,33 DATA_T,34 EVALUATOR_T,35 ExperimentResultRow,36 _evaluators_include_attachments,37 _ExperimentManagerMixin,38 _extract_feedback_keys,39 _ForwardResults,40 _get_target_args,41 _is_langchain_runnable,42 _load_examples_map,43 _load_experiment,44 _load_tqdm,45 _load_traces,46 _resolve_data,47 _resolve_evaluators,48 _resolve_experiment,49 _target_include_attachments,50 _to_pandas,51 _wrap_summary_evaluators,52)53from langsmith.evaluation.evaluator import (54 SUMMARY_EVALUATOR_T,55 EvaluationResult,56 EvaluationResults,57 RunEvaluator,58)59 60if TYPE_CHECKING:61 import pandas as pd62 from langchain_core.runnables import Runnable63 64 DataFrame = pd.DataFrame65else:66 DataFrame = Any67 68logger = logging.getLogger(__name__)69 70ATARGET_T = Union[71 Callable[[dict], Awaitable[dict]], Callable[[dict, dict], Awaitable[dict]]72]73 74 75async def aevaluate(76 target: Union[77 ATARGET_T, AsyncIterable[dict], Runnable, str, uuid.UUID, schemas.TracerSession78 ],79 /,80 data: Union[81 DATA_T, AsyncIterable[schemas.Example], Iterable[schemas.Example], None82 ] = None,83 evaluators: Optional[Sequence[Union[EVALUATOR_T, AEVALUATOR_T]]] = None,84 summary_evaluators: Optional[Sequence[SUMMARY_EVALUATOR_T]] = None,85 metadata: Optional[dict] = None,86 experiment_prefix: Optional[str] = None,87 description: Optional[str] = None,88 max_concurrency: Optional[int] = 0,89 num_repetitions: int = 1,90 client: Optional[langsmith.Client] = None,91 blocking: bool = True,92 experiment: Optional[Union[schemas.TracerSession, str, uuid.UUID]] = None,93 upload_results: bool = True,94 error_handling: Literal["log", "ignore"] = "log",95 **kwargs: Any,96) -> AsyncExperimentResults:97 r"""Evaluate an async target system on a given dataset.98 99 Args:100 target (AsyncCallable[[dict], dict] | AsyncIterable[dict] | Runnable | EXPERIMENT_T | Tuple[EXPERIMENT_T, EXPERIMENT_T]):101 The target system or experiment(s) to evaluate.102 103 Can be an async function that takes a `dict` and returns a `dict`, a104 langchain `Runnable`, an existing experiment ID, or a two-tuple of experiment IDs.105 data (Union[DATA_T, AsyncIterable[schemas.Example]]): The dataset to evaluate on.106 107 Can be a dataset name, a list of examples, an async generator of examples, or an async iterable of examples.108 evaluators (Optional[Sequence[EVALUATOR_T]]): A list of evaluators to run109 on each example.110 summary_evaluators (Optional[Sequence[SUMMARY_EVALUATOR_T]]): A list of summary111 evaluators to run on the entire dataset.112 metadata (Optional[dict]): Metadata to attach to the experiment.113 experiment_prefix (Optional[str]): A prefix to provide for your experiment name.114 description (Optional[str]): A description of the experiment.115 max_concurrency (int | None): The maximum number of concurrent116 evaluations to run.117 118 If `None` then no limit is set. If `0` then no concurrency.119 num_repetitions (int): The number of times to run the evaluation.120 Each item in the dataset will be run and evaluated this many times.121 client (Optional[langsmith.Client]): The LangSmith client to use.122 blocking (bool): Whether to block until the evaluation is complete.123 experiment (Optional[schemas.TracerSession]): An existing experiment to124 extend.125 126 If provided, `experiment_prefix` is ignored. For advanced usage only.127 error_handling (str, default="log"): How to handle individual run errors.128 129 `'log'` will trace the runs with the error message as part of the130 experiment, `'ignore'` will not count the run as part of the experiment at131 all.132 133 Returns:134 An async iterator over the experiment results.135 136 Environment:137 - `LANGSMITH_TEST_CACHE`: If set, API calls will be cached to disk to save time and138 cost during testing.139 140 Recommended to commit the cache files to your repository for faster CI/CD runs.141 142 Requires the `'langsmith[vcr]'` package to be installed.143 144 Examples:145 >>> from typing import Sequence146 >>> from langsmith import Client, aevaluate147 >>> from langsmith.schemas import Example, Run148 >>> client = Client()149 >>> dataset = client.clone_public_dataset(150 ... "https://smith.langchain.com/public/419dcab2-1d66-4b94-8901-0357ead390df/d"151 ... )152 >>> dataset_name = "Evaluate Examples"153 154 Basic usage:155 156 >>> def accuracy(run: Run, example: Example):157 ... # Row-level evaluator for accuracy.158 ... pred = run.outputs["output"]159 ... expected = example.outputs["answer"]160 ... return {"score": expected.lower() == pred.lower()}161 162 >>> def precision(runs: Sequence[Run], examples: Sequence[Example]):163 ... # Experiment-level evaluator for precision.164 ... # TP / (TP + FP)165 ... predictions = [run.outputs["output"].lower() for run in runs]166 ... expected = [example.outputs["answer"].lower() for example in examples]167 ... # yes and no are the only possible answers168 ... tp = sum([p == e for p, e in zip(predictions, expected) if p == "yes"])169 ... fp = sum([p == "yes" and e == "no" for p, e in zip(predictions, expected)])170 ... return {"score": tp / (tp + fp)}171 172 >>> import asyncio173 >>> async def apredict(inputs: dict) -> dict:174 ... # This can be any async function or just an API call to your app.175 ... await asyncio.sleep(0.1)176 ... return {"output": "Yes"}177 >>> results = asyncio.run(178 ... aevaluate(179 ... apredict,180 ... data=dataset_name,181 ... evaluators=[accuracy],182 ... summary_evaluators=[precision],183 ... experiment_prefix="My Experiment",184 ... description="Evaluate the accuracy of the model asynchronously.",185 ... metadata={186 ... "my-prompt-version": "abcd-1234",187 ... },188 ... )189 ... ) # doctest: +ELLIPSIS190 View the evaluation results for experiment:...191 192 Evaluating over only a subset of the examples using an async generator:193 194 >>> async def example_generator():195 ... examples = client.list_examples(dataset_name=dataset_name, limit=5)196 ... for example in examples:197 ... yield example198 >>> results = asyncio.run(199 ... aevaluate(200 ... apredict,201 ... data=example_generator(),202 ... evaluators=[accuracy],203 ... summary_evaluators=[precision],204 ... experiment_prefix="My Subset Experiment",205 ... description="Evaluate a subset of examples asynchronously.",206 ... )207 ... ) # doctest: +ELLIPSIS208 View the evaluation results for experiment:...209 210 Streaming each prediction to more easily + eagerly debug.211 212 >>> results = asyncio.run(213 ... aevaluate(214 ... apredict,215 ... data=dataset_name,216 ... evaluators=[accuracy],217 ... summary_evaluators=[precision],218 ... experiment_prefix="My Streaming Experiment",219 ... description="Streaming predictions for debugging.",220 ... blocking=False,221 ... )222 ... ) # doctest: +ELLIPSIS223 View the evaluation results for experiment:...224 225 >>> async def aenumerate(iterable):226 ... async for elem in iterable:227 ... print(elem)228 >>> asyncio.run(aenumerate(results))229 230 Running without concurrency:231 232 >>> results = asyncio.run(233 ... aevaluate(234 ... apredict,235 ... data=dataset_name,236 ... evaluators=[accuracy],237 ... summary_evaluators=[precision],238 ... experiment_prefix="My Experiment Without Concurrency",239 ... description="This was run without concurrency.",240 ... max_concurrency=0,241 ... )242 ... ) # doctest: +ELLIPSIS243 View the evaluation results for experiment:...244 245 Using Async evaluators:246 247 >>> async def helpfulness(run: Run, example: Example):248 ... # Row-level evaluator for helpfulness.249 ... await asyncio.sleep(5) # Replace with your LLM API call250 ... return {"score": run.outputs["output"] == "Yes"}251 252 >>> results = asyncio.run(253 ... aevaluate(254 ... apredict,255 ... data=dataset_name,256 ... evaluators=[helpfulness],257 ... summary_evaluators=[precision],258 ... experiment_prefix="My Helpful Experiment",259 ... description="Applying async evaluators example.",260 ... )261 ... ) # doctest: +ELLIPSIS262 View the evaluation results for experiment:...263 264 265 !!! warning "Behavior changed in `langsmith` 0.2.0"266 267 'max_concurrency' default updated from None (no limit on concurrency)268 to 0 (no concurrency at all).269 """ # noqa: E501270 if isinstance(target, (str, uuid.UUID, schemas.TracerSession)):271 invalid_args = {272 "num_repetitions": num_repetitions > 1,273 "experiment": bool(experiment),274 "upload_results": not upload_results,275 "experiment_prefix": bool(experiment_prefix),276 "data": bool(data),277 }278 if any(invalid_args.values()):279 msg = (280 f"Received invalid arguments. "281 f"{tuple(k for k, v in invalid_args.items() if v)} should not be "282 f"specified when target is an existing experiment."283 )284 raise ValueError(msg)285 target_id = target if isinstance(target, (str, uuid.UUID)) else target.id286 logger.debug(f"Running evaluation over existing experiment {target_id}...")287 return await aevaluate_existing(288 target,289 evaluators=evaluators,290 summary_evaluators=summary_evaluators,291 metadata=metadata,292 max_concurrency=max_concurrency,293 client=client,294 blocking=blocking,295 **kwargs,296 )297 elif isinstance(target, (list, tuple)):298 msg = (299 "Running a comparison of two existing experiments asynchronously is not "300 "currently supported. Please use the `evaluate()` method instead and make "301 "sure that your evaluators are defined as synchronous functions."302 )303 raise ValueError(msg)304 elif kwargs:305 msg = (306 f"Received unsupported arguments {kwargs}. These arguments are not "307 f"supported when creating a new experiment."308 )309 raise ValueError(msg)310 elif not data:311 msg = "Must specify 'data' when running evaluations over a target function."312 raise ValueError(msg)313 elif experiment and experiment_prefix:314 msg = (315 "Expected at most one of 'experiment' or 'experiment_prefix',"316 " but both were provided. "317 f"Got: experiment={experiment}, experiment_prefix={experiment_prefix}"318 )319 raise ValueError(msg)320 else:321 if not upload_results:322 _warn_once("'upload_results' parameter is in beta.")323 logger.debug(f"Running evaluation over target system {target}...")324 return await _aevaluate(325 target,326 data=data,327 evaluators=evaluators,328 summary_evaluators=summary_evaluators,329 metadata=metadata,330 experiment_prefix=experiment_prefix,331 description=description,332 max_concurrency=max_concurrency,333 num_repetitions=num_repetitions,334 client=client,335 blocking=blocking,336 experiment=experiment,337 upload_results=upload_results,338 error_handling=error_handling,339 )340 341 342async def aevaluate_existing(343 experiment: Union[str, uuid.UUID, schemas.TracerSession],344 /,345 evaluators: Optional[Sequence[Union[EVALUATOR_T, AEVALUATOR_T]]] = None,346 summary_evaluators: Optional[Sequence[SUMMARY_EVALUATOR_T]] = None,347 metadata: Optional[dict] = None,348 max_concurrency: Optional[int] = 0,349 client: Optional[langsmith.Client] = None,350 load_nested: bool = False,351 blocking: bool = True,352) -> AsyncExperimentResults:353 r"""Evaluate existing experiment runs asynchronously.354 355 Args:356 experiment (Union[str, uuid.UUID]): The identifier of the experiment to evaluate.357 evaluators (Optional[Sequence[EVALUATOR_T]]): Optional sequence of evaluators to use for individual run evaluation.358 summary_evaluators (Optional[Sequence[SUMMARY_EVALUATOR_T]]): Optional sequence of evaluators359 to apply over the entire dataset.360 metadata (Optional[dict]): Optional metadata to include in the evaluation results.361 max_concurrency (int | None): The maximum number of concurrent362 evaluations to run.363 364 If `None` then no limit is set. If `0` then no concurrency.365 client (Optional[langsmith.Client]): Optional Langsmith client to use for evaluation.366 load_nested: Whether to load all child runs for the experiment.367 368 Default is to only load the top-level root runs.369 blocking (bool): Whether to block until evaluation is complete.370 371 Returns:372 An async iterator over the experiment results.373 374 Examples:375 Define your evaluators376 377 >>> from typing import Sequence378 >>> from langsmith.schemas import Example, Run379 >>> def accuracy(run: Run, example: Example):380 ... # Row-level evaluator for accuracy.381 ... pred = run.outputs["output"]382 ... expected = example.outputs["answer"]383 ... return {"score": expected.lower() == pred.lower()}384 >>> def precision(runs: Sequence[Run], examples: Sequence[Example]):385 ... # Experiment-level evaluator for precision.386 ... # TP / (TP + FP)387 ... predictions = [run.outputs["output"].lower() for run in runs]388 ... expected = [example.outputs["answer"].lower() for example in examples]389 ... # yes and no are the only possible answers390 ... tp = sum([p == e for p, e in zip(predictions, expected) if p == "yes"])391 ... fp = sum([p == "yes" and e == "no" for p, e in zip(predictions, expected)])392 ... return {"score": tp / (tp + fp)}393 394 Load the experiment and run the evaluation.395 396 >>> import asyncio397 >>> import uuid398 >>> from langsmith import Client, aevaluate, aevaluate_existing399 >>> client = Client()400 >>> dataset_name = "__doctest_aevaluate_existing_" + uuid.uuid4().hex[:8]401 >>> dataset = client.create_dataset(dataset_name)402 >>> example = client.create_example(403 ... inputs={"question": "What is 2+2?"},404 ... outputs={"answer": "4"},405 ... dataset_id=dataset.id,406 ... )407 >>> async def apredict(inputs: dict) -> dict:408 ... await asyncio.sleep(0.001)409 ... return {"output": "4"}410 >>> results = asyncio.run(411 ... aevaluate(412 ... apredict, data=dataset_name, experiment_prefix="doctest_experiment"413 ... )414 ... ) # doctest: +ELLIPSIS415 View the evaluation results for experiment:...416 >>> experiment_id = results.experiment_name417 >>> # Consume all results to ensure evaluation is complete418 >>> async def consume_results():419 ... result_list = [r async for r in results]420 ... return len(result_list) > 0421 >>> asyncio.run(consume_results())422 True423 >>> import time424 >>> time.sleep(3)425 >>> results = asyncio.run(426 ... aevaluate_existing(427 ... experiment_id,428 ... evaluators=[accuracy],429 ... summary_evaluators=[precision],430 ... )431 ... ) # doctest: +ELLIPSIS432 View the evaluation results for experiment:...433 >>> client.delete_dataset(dataset_id=dataset.id)434 435 436 """ # noqa: E501437 client = client or run_trees.get_cached_client()438 project = (439 experiment440 if isinstance(experiment, schemas.TracerSession)441 else (442 await aitertools.aio_to_thread(443 contextvars.copy_context(), _load_experiment, experiment, client444 )445 )446 )447 runs = await aitertools.aio_to_thread(448 contextvars.copy_context(),449 _load_traces,450 experiment,451 client,452 load_nested=load_nested,453 )454 data_map = await aitertools.aio_to_thread(455 contextvars.copy_context(), _load_examples_map, client, project456 )457 data = [data_map[run.reference_example_id] for run in runs]458 return await _aevaluate(459 runs,460 data=data,461 evaluators=evaluators,462 summary_evaluators=summary_evaluators,463 metadata=metadata,464 max_concurrency=max_concurrency,465 client=client,466 blocking=blocking,467 experiment=project,468 )469 470 471async def _aevaluate(472 target: Union[ATARGET_T, AsyncIterable[dict], Iterable[schemas.Run], Runnable],473 /,474 data: Union[DATA_T, AsyncIterable[schemas.Example]],475 evaluators: Optional[Sequence[Union[EVALUATOR_T, AEVALUATOR_T]]] = None,476 summary_evaluators: Optional[Sequence[SUMMARY_EVALUATOR_T]] = None,477 metadata: Optional[dict] = None,478 experiment_prefix: Optional[str] = None,479 description: Optional[str] = None,480 max_concurrency: Optional[int] = None,481 num_repetitions: int = 1,482 client: Optional[langsmith.Client] = None,483 blocking: bool = True,484 experiment: Optional[Union[schemas.TracerSession, str, uuid.UUID]] = None,485 upload_results: bool = True,486 error_handling: Literal["log", "ignore"] = "log",487) -> AsyncExperimentResults:488 is_async_target = (489 asyncio.iscoroutinefunction(target)490 or (hasattr(target, "__aiter__") and asyncio.iscoroutine(target.__aiter__()))491 or _is_langchain_runnable(target)492 )493 client = client or rt.get_cached_client()494 runs = None if is_async_target else cast(Iterable[schemas.Run], target)495 experiment_, runs = await aitertools.aio_to_thread(496 contextvars.copy_context(),497 _resolve_experiment,498 experiment,499 runs,500 client,501 )502 num_include_attachments = int(503 _target_include_attachments(target)504 ) + _evaluators_include_attachments(evaluators)505 manager = await _AsyncExperimentManager(506 data,507 client=client,508 metadata=metadata,509 experiment=experiment_ or experiment_prefix,510 description=description,511 num_repetitions=num_repetitions,512 runs=runs,513 include_attachments=num_include_attachments > 0,514 reuse_attachments=num_repetitions * num_include_attachments > 1,515 upload_results=upload_results,516 error_handling=error_handling,517 ).astart()518 cache_dir = ls_utils.get_cache_dir(None)519 if cache_dir is not None:520 dsid = await manager.get_dataset_id()521 cache_path = pathlib.Path(cache_dir) / f"{dsid}.yaml"522 else:523 cache_path = None524 with ls_utils.with_optional_cache(cache_path, ignore_hosts=[client.api_url]):525 if is_async_target:526 if evaluators:527 # Run predictions and evaluations in a single pipeline528 manager = await manager.awith_predictions_and_evaluators(529 cast(ATARGET_T, target), evaluators, max_concurrency=max_concurrency530 )531 else:532 manager = await manager.awith_predictions(533 cast(ATARGET_T, target), max_concurrency=max_concurrency534 )535 if summary_evaluators:536 manager = await manager.awith_summary_evaluators(summary_evaluators)537 else:538 if evaluators:539 manager = await manager.awith_evaluators(540 evaluators, max_concurrency=max_concurrency541 )542 if summary_evaluators:543 manager = await manager.awith_summary_evaluators(summary_evaluators)544 results = AsyncExperimentResults(manager)545 if blocking:546 await results.wait()547 return results548 549 550class _AsyncExperimentManager(_ExperimentManagerMixin):551 """Manage the execution of experiments asynchronously.552 553 Supports lazily running predictions and evaluations in parallel to facilitate554 result streaming and early debugging.555 556 Args:557 data (DATA_T): The data used for the experiment. Can be a dataset name or ID OR558 a generator of examples.559 runs (Optional[Iterable[schemas.Run]]): The runs associated with the experiment560 predictions.561 experiment (Optional[schemas.TracerSession]): The tracer session562 associated with the experiment.563 experiment_prefix (Optional[str]): The prefix for the experiment name.564 description (Optional[str]): The description for the experiment.565 metadata (Optional[dict]): Additional metadata for the experiment.566 client (Optional[langsmith.Client]): The Langsmith client used for567 the experiment.568 evaluation_results (Optional[Iterable[EvaluationResults]]): The evaluation569 sresults for the experiment.570 summary_results (Optional[Iterable[EvaluationResults]]): The aggregate results571 for the experiment.572 num_repetitions (Optional[int], default=1): The number of repetitions for573 the experiment.574 include_attachments (Optional[bool], default=False): Whether to include575 attachments. This is used for when we pull the examples for the experiment.576 reuse_attachments (Optional[bool], default=False): Whether to reuse attachments577 from examples. This is True if we need to reuse attachments across multiple578 target/evaluator functions.579 upload_results (Optional[bool], default=True): Whether to upload results580 to Langsmith.581 attachment_raw_data_dict (Optional[dict]): A dictionary to store raw data582 for attachments. Only used if we reuse attachments across multiple583 target/evaluator functions.584 error_handling (str, default="log"): How to handle individual run errors.585 586 `'log'` will trace the runs with the error message as part of the587 experiment, `'ignore'` will not count the run as part of the experiment at588 all.589 """590 591 def __init__(592 self,593 data: Union[DATA_T, AsyncIterable[schemas.Example]],594 /,595 experiment: Optional[Union[schemas.TracerSession, str]] = None,596 metadata: Optional[dict] = None,597 runs: Optional[Union[Iterable[schemas.Run], AsyncIterable[schemas.Run]]] = None,598 client: Optional[langsmith.Client] = None,599 evaluation_results: Optional[AsyncIterable[EvaluationResults]] = None,600 summary_results: Optional[AsyncIterable[EvaluationResults]] = None,601 description: Optional[str] = None,602 num_repetitions: int = 1,603 include_attachments: bool = False,604 reuse_attachments: bool = False,605 upload_results: bool = True,606 attachment_raw_data_dict: Optional[dict] = None,607 error_handling: Literal["log", "ignore"] = "log",608 ):609 super().__init__(610 experiment=experiment,611 metadata=metadata,612 client=client,613 description=description,614 )615 self._data = data616 self._examples: Optional[AsyncIterable[schemas.Example]] = None617 self._runs = (618 aitertools.ensure_async_iterator(runs) if runs is not None else None619 )620 self._evaluation_results = evaluation_results621 self._summary_results = summary_results622 self._num_repetitions = num_repetitions623 self._include_attachments = include_attachments624 self._reuse_attachments = reuse_attachments625 self._upload_results = upload_results626 self._attachment_raw_data_dict = attachment_raw_data_dict627 self._error_handling = error_handling628 629 def _reset_example_attachments(self, example: schemas.Example) -> schemas.Example:630 """Reset attachment readers for an example.631 632 This is only in the case that an attachment is going to be used by more633 than 1 callable (target + evaluators). In that case we keep a single copy634 of the attachment data in self._attachment_raw_data_dict, and create635 readers from that data. This makes it so that we don't have to keep636 copies of the same data in memory, instead we can just create readers637 from the same data.638 """639 if not hasattr(example, "attachments") or not example.attachments:640 return example641 642 new_attachments: dict[str, schemas.AttachmentInfo] = {}643 for name, attachment in example.attachments.items():644 if (645 self._attachment_raw_data_dict is not None646 and str(example.id) + name in self._attachment_raw_data_dict647 ):648 new_attachments[name] = {649 "presigned_url": attachment["presigned_url"],650 "reader": io.BytesIO(651 self._attachment_raw_data_dict[str(example.id) + name]652 ),653 "mime_type": attachment["mime_type"],654 }655 else:656 new_attachments[name] = attachment657 658 # Create a new Example instance with the updated attachments659 return schemas.Example(660 id=example.id,661 created_at=example.created_at,662 dataset_id=example.dataset_id,663 inputs=example.inputs,664 outputs=example.outputs,665 metadata=example.metadata,666 modified_at=example.modified_at,667 source_run_id=example.source_run_id,668 attachments=new_attachments,669 _host_url=example._host_url,670 _tenant_id=example._tenant_id,671 )672 673 async def aget_examples(self) -> AsyncIterator[schemas.Example]:674 if self._examples is None:675 self._examples = _aresolve_data(676 self._data,677 client=self.client,678 include_attachments=self._include_attachments,679 )680 if self._reuse_attachments and self._attachment_raw_data_dict is None:681 examples_copy, self._examples = aitertools.atee(self._examples)682 self._attachment_raw_data_dict = {683 str(e.id) + name: value["reader"].read()684 async for e in examples_copy685 for name, value in (e.attachments or {}).items()686 }687 if self._num_repetitions > 1:688 examples_list = [example async for example in self._examples]689 self._examples = async_chain_from_iterable(690 [691 async_iter_from_list(692 [693 self._reset_example_attachments(example)694 for example in examples_list695 ]696 )697 for _ in range(self._num_repetitions)698 ]699 )700 701 self._examples, examples_iter = aitertools.atee(702 aitertools.ensure_async_iterator(self._examples), 2, lock=asyncio.Lock()703 )704 return examples_iter705 706 async def get_dataset_id(self) -> str:707 if self._experiment is None or not getattr(708 self._experiment, "reference_dataset_id", None709 ):710 example = await aitertools.py_anext(await self.aget_examples())711 if example is None:712 raise ValueError("No examples found in the dataset.")713 return str(example.dataset_id)714 return str(self._experiment.reference_dataset_id)715 716 async def aget_runs(self) -> AsyncIterator[schemas.Run]:717 if self._runs is None:718 raise ValueError("Runs not loaded yet.")719 self._runs, runs = aitertools.atee(720 aitertools.ensure_async_iterator(self._runs), 2, lock=asyncio.Lock()721 )722 async for run in runs:723 yield run724 725 async def aget_evaluation_results(self) -> AsyncIterator[EvaluationResults]:726 if self._evaluation_results is None:727 async for _ in await self.aget_examples():728 yield {"results": []}729 else:730 self._evaluation_results, evaluation_results = aitertools.atee(731 aitertools.ensure_async_iterator(self._evaluation_results),732 2,733 lock=asyncio.Lock(),734 )735 async for result in evaluation_results:736 yield result737 738 async def astart(self) -> _AsyncExperimentManager:739 try:740 first_example = await aitertools.py_anext(await self.aget_examples())741 except StopAsyncIteration:742 raise ValueError(743 "No examples found in the dataset. "744 "Please ensure the data provided to aevaluate is not empty."745 )746 if not first_example:747 raise ValueError(748 "No examples found in the dataset."749 "Please ensure the data provided to aevaluate is not empty."750 )751 project = self._get_project(first_example) if self._upload_results else None752 self._print_experiment_start(project, first_example)753 self._metadata["num_repetitions"] = self._num_repetitions754 return self._copy(755 await self.aget_examples(),756 experiment=project,757 )758 759 def _get_example_with_readers(self, example: schemas.Example) -> schemas.Example:760 new_attachments: dict[str, schemas.AttachmentInfo] = {}761 for name, attachment in (example.attachments or {}).items():762 if (763 self._attachment_raw_data_dict is not None764 and str(example.id) + name in self._attachment_raw_data_dict765 ):766 reader = io.BytesIO(767 self._attachment_raw_data_dict[str(example.id) + name]768 )769 new_attachments[name] = {770 "presigned_url": attachment["presigned_url"],771 "reader": reader,772 "mime_type": attachment["mime_type"],773 }774 else:775 new_attachments[name] = attachment776 777 return schemas.Example(778 id=example.id,779 created_at=example.created_at,780 dataset_id=example.dataset_id,781 inputs=example.inputs,782 outputs=example.outputs,783 metadata=example.metadata,784 modified_at=example.modified_at,785 source_run_id=example.source_run_id,786 attachments=new_attachments,787 _host_url=example._host_url,788 _tenant_id=example._tenant_id,789 )790 791 async def awith_predictions_and_evaluators(792 self,793 target: ATARGET_T,794 evaluators: Sequence[Union[EVALUATOR_T, AEVALUATOR_T]],795 /,796 max_concurrency: Optional[int] = None,797 ) -> _AsyncExperimentManager:798 """Run predictions and evaluations in a single pipeline.799 800 This allows evaluators to process results as soon as they're available from801 the target function, rather than waiting for all predictions to complete first.802 """803 evaluators = _resolve_evaluators(evaluators)804 805 if not hasattr(self, "_evaluation_feedback_executor"):806 self._evaluation_feedback_executor = cf.ThreadPoolExecutor(max_workers=4)807 808 traceable_target = _ensure_async_traceable(target)809 810 async def process_example(example: schemas.Example):811 # Yield the coroutine to be awaited later812 pred = await _aforward(813 traceable_target,814 self._get_example_with_readers(example),815 self.experiment_name,816 self._metadata,817 self.client,818 _target_include_attachments(target),819 self._error_handling,820 )821 example, run = pred["example"], pred["run"]822 result = await self._arun_evaluators(823 evaluators,824 {825 "run": run,826 "example": example,827 "evaluation_results": {"results": []},828 },829 feedback_executor=self._evaluation_feedback_executor,830 )831 return result832 833 async def process_examples():834 """Create a single task per example.835 836 That task is to run the target function and all the evaluators837 sequentially.838 """839 async for example in await self.aget_examples():840 yield process_example(example)841 842 await self._aend()843 844 # Run the per-example tasks with max-concurrency845 # This guarantees that max_concurrency is the upper limit846 # for the number of target/evaluators that can be run in parallel847 experiment_results = aitertools.aiter_with_concurrency(848 max_concurrency,849 process_examples(),850 _eager_consumption_timeout=0.001,851 )852 853 r1, r2, r3 = aitertools.atee(experiment_results, 3, lock=asyncio.Lock())854 855 return self._copy(856 (result["example"] async for result in r1),857 runs=(result["run"] async for result in r2),858 evaluation_results=(result["evaluation_results"] async for result in r3),859 )860 861 async def awith_predictions(862 self,863 target: ATARGET_T,864 /,865 max_concurrency: Optional[int] = None,866 ) -> _AsyncExperimentManager:867 _experiment_results = self._apredict(868 target,869 max_concurrency=max_concurrency,870 include_attachments=_target_include_attachments(target),871 )872 r1, r2 = aitertools.atee(_experiment_results, 2, lock=asyncio.Lock())873 return self._copy(874 (pred["example"] async for pred in r1),875 runs=(pred["run"] async for pred in r2),876 )877 878 async def awith_evaluators(879 self,880 evaluators: Sequence[Union[EVALUATOR_T, AEVALUATOR_T]],881 *,882 max_concurrency: Optional[int] = None,883 ) -> _AsyncExperimentManager:884 evaluators = _resolve_evaluators(evaluators)885 experiment_results = self._ascore(evaluators, max_concurrency=max_concurrency)886 r1, r2, r3 = aitertools.atee(experiment_results, 3, lock=asyncio.Lock())887 return self._copy(888 (result["example"] async for result in r1),889 runs=(result["run"] async for result in r2),890 evaluation_results=(result["evaluation_results"] async for result in r3),891 )892 893 async def awith_summary_evaluators(894 self,895 summary_evaluators: Sequence[SUMMARY_EVALUATOR_T],896 ) -> _AsyncExperimentManager:897 wrapped_evaluators = _wrap_summary_evaluators(summary_evaluators)898 aggregate_feedback_gen = self._aapply_summary_evaluators(wrapped_evaluators)899 return self._copy(900 await self.aget_examples(),901 runs=self.aget_runs(),902 summary_results=aggregate_feedback_gen,903 )904 905 async def aget_results(self) -> AsyncIterator[ExperimentResultRow]:906 async for run, example, evaluation_results in aitertools.async_zip(907 self.aget_runs(), await self.aget_examples(), self.aget_evaluation_results()908 ):909 yield ExperimentResultRow(910 run=run,911 example=example,912 evaluation_results=evaluation_results,913 )914 915 async def aget_summary_scores(self) -> dict[str, list[dict]]:916 if self._summary_results is None:917 return {"results": []}918 return {919 "results": [920 res # type: ignore[misc]921 async for results in self._summary_results922 for res in results["results"]923 ]924 }925 926 ## Private methods927 928 async def _apredict(929 self,930 target: ATARGET_T,931 /,932 max_concurrency: Optional[int] = None,933 include_attachments: bool = False,934 ) -> AsyncIterator[_ForwardResults]:935 fn = _ensure_async_traceable(target)936 937 async def predict_all():938 async for example in await self.aget_examples():939 # Yield the coroutine to be awaited later940 yield _aforward(941 fn,942 self._get_example_with_readers(example),943 self.experiment_name,944 self._metadata,945 self.client,946 include_attachments,947 self._error_handling,948 )949 950 async for result in aitertools.aiter_with_concurrency(951 max_concurrency, predict_all(), _eager_consumption_timeout=0.001952 ):953 yield result954 955 await self._aend()956 957 async def _ascore(958 self,959 evaluators: Sequence[RunEvaluator],960 max_concurrency: Optional[int] = None,961 ) -> AsyncIterator[ExperimentResultRow]:962 with cf.ThreadPoolExecutor(max_workers=4) as feedback_executor:963 964 async def score_all():965 async for current_results in self.aget_results():966 # Yield the coroutine to be awaited later in aiter_with_concurrency967 yield self._arun_evaluators(968 evaluators, current_results, feedback_executor=feedback_executor969 )970 971 async for result in aitertools.aiter_with_concurrency(972 max_concurrency, score_all(), _eager_consumption_timeout=0.001973 ):974 yield result975 976 async def _arun_evaluators(977 self,978 evaluators: Sequence[RunEvaluator],979 current_results: ExperimentResultRow,980 feedback_executor: cf.ThreadPoolExecutor,981 ) -> ExperimentResultRow:982 current_context = rh.get_tracing_context()983 metadata = {984 **(current_context["metadata"] or {}),985 **{"experiment": self.experiment_name},986 }987 with rh.tracing_context(988 **{989 **current_context,990 "project_name": "evaluators",991 "metadata": metadata,992 "enabled": "local" if not self._upload_results else True,993 "client": self.client,994 }995 ):996 run = current_results["run"]997 example = current_results["example"]998 eval_results = current_results["evaluation_results"]999 1000 async def _run_single_evaluator(evaluator: RunEvaluator):1001 evaluator_run_id = uuid.uuid4()1002 try:1003 evaluator_response = await evaluator.aevaluate_run( # type: ignore[call-arg]1004 run=run,1005 example=self._get_example_with_readers(example),1006 evaluator_run_id=evaluator_run_id,1007 )1008 selected_results = self.client._select_eval_results(1009 evaluator_response1010 )1011 1012 if self._upload_results:1013 self.client._log_evaluation_feedback(1014 evaluator_response, run=run, _executor=feedback_executor1015 )1016 return selected_results1017 except Exception as e:1018 try:1019 feedback_keys = _extract_feedback_keys(evaluator)1020 1021 error_response = EvaluationResults(1022 results=[1023 EvaluationResult(1024 key=key,1025 source_run_id=evaluator_run_id,1026 comment=repr(e),1027 extra={"error": True},1028 )1029 for key in feedback_keys1030 ]1031 )1032 selected_results = self.client._select_eval_results(1033 error_response1034 )1035 if self._upload_results:1036 self.client._log_evaluation_feedback(1037 error_response, run=run, _executor=feedback_executor1038 )1039 return selected_results1040 except Exception as e2:1041 logger.debug(f"Error parsing feedback keys: {e2}")1042 pass1043 logger.error(1044 f"Error running evaluator {repr(evaluator)} on"1045 f" run {run.id}: {repr(e)}",1046 exc_info=True,1047 )1048 1049 all_results = []1050 for evaluator in evaluators:1051 all_results.append(await _run_single_evaluator(evaluator))1052 1053 for result in all_results:1054 if result is not None:1055 eval_results["results"].extend(result)1056 return ExperimentResultRow(1057 run=run,1058 example=example,1059 evaluation_results=eval_results,1060 )1061 1062 async def _aapply_summary_evaluators(1063 self, summary_evaluators: Sequence[SUMMARY_EVALUATOR_T]1064 ) -> AsyncIterator[EvaluationResults]:1065 runs, examples = [], []1066 async_examples = aitertools.ensure_async_iterator(await self.aget_examples())1067 async for run, example in aitertools.async_zip(1068 self.aget_runs(), async_examples1069 ):1070 runs.append(run)1071 examples.append(example)1072 aggregate_feedback = []1073 project_id = self._get_experiment().id if self._upload_results else None1074 current_context = rh.get_tracing_context()1075 metadata = {1076 **(current_context["metadata"] or {}),1077 **{1078 "experiment": self.experiment_name,1079 "experiment_id": project_id,1080 },1081 }1082 with rh.tracing_context(1083 **{1084 **current_context,1085 "project_name": "evaluators",1086 "metadata": metadata,1087 "enabled": "local" if not self._upload_results else True,1088 "client": self.client,1089 }1090 ):1091 for evaluator in summary_evaluators:1092 try:1093 summary_eval_result = evaluator(runs, examples)1094 flattened_results = self.client._select_eval_results(1095 summary_eval_result,1096 fn_name=evaluator.__name__,1097 )1098 aggregate_feedback.extend(flattened_results)1099 if self._upload_results:1100 for result in flattened_results:1101 feedback = result.model_dump(exclude={"target_run_id"})1102 evaluator_info = feedback.pop("evaluator_info", None)1103 await aitertools.aio_to_thread(1104 contextvars.copy_context(),1105 self.client.create_feedback,1106 **feedback,1107 run_id=None,1108 project_id=project_id,1109 source_info=evaluator_info,1110 )1111 except Exception as e:1112 logger.error(1113 f"Error running summary evaluator {repr(evaluator)}: {e}",1114 exc_info=True,1115 )1116 yield {"results": aggregate_feedback}1117 1118 async def _get_dataset_version(self) -> Optional[str]:1119 modified_at = []1120 async for example in await self.aget_examples():1121 if example.modified_at:1122 # Should always be defined in practice when fetched,1123 # but the typing permits None1124 modified_at.append(example.modified_at)1125 1126 max_modified_at = max(modified_at) if modified_at else None1127 return max_modified_at.isoformat() if max_modified_at else None1128 1129 async def _get_dataset_splits(self) -> Optional[list[str]]:1130 splits = set()1131 async for example in await self.aget_examples():1132 if (1133 example.metadata1134 and example.metadata.get("dataset_split")1135 and isinstance(example.metadata["dataset_split"], list)1136 ):1137 for split in example.metadata["dataset_split"]:1138 if isinstance(split, str):1139 splits.add(split)1140 else:1141 splits.add("base")1142 1143 return list(splits)1144 1145 async def _aend(self) -> None:1146 if not self._upload_results:1147 return1148 experiment = self._experiment1149 if experiment is None:1150 raise ValueError("Experiment not started yet.")1151 1152 project_metadata = self._get_experiment_metadata()1153 project_metadata["dataset_version"] = await self._get_dataset_version()1154 project_metadata["dataset_splits"] = await self._get_dataset_splits()1155 self.client.update_project(1156 experiment.id,1157 metadata={1158 **experiment.metadata,1159 **project_metadata,1160 },1161 )1162 1163 def _copy(self, *args: Any, **kwargs: Any) -> _AsyncExperimentManager:1164 default_args = (self._data,)1165 default_kwargs = {1166 "experiment": self._experiment,1167 "metadata": self._metadata,1168 "runs": self._runs,1169 "client": self.client,1170 "evaluation_results": self._evaluation_results,1171 "summary_results": self._summary_results,1172 "include_attachments": self._include_attachments,1173 "reuse_attachments": self._reuse_attachments,1174 "upload_results": self._upload_results,1175 "attachment_raw_data_dict": self._attachment_raw_data_dict,1176 "error_handling": self._error_handling,1177 }1178 full_args = list(args) + list(default_args[len(args) :])1179 full_kwargs = {**default_kwargs, **kwargs}1180 return self.__class__(*full_args, **full_kwargs)1181 1182 1183class AsyncExperimentResults:1184 def __init__(1185 self,1186 experiment_manager: _AsyncExperimentManager,1187 ):1188 self._manager = experiment_manager1189 self._results: list[ExperimentResultRow] = []1190 self._condition = asyncio.Condition()1191 self._task = asyncio.create_task(self._process_data(self._manager))1192 self._processed_count = 01193 self._comparison_url: Optional[str] = None1194 1195 @property1196 def experiment_name(self) -> str:1197 return self._manager.experiment_name1198 1199 @property1200 def experiment_id(self) -> uuid.UUID: