Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
_arunner.py1396 linesDownload Raw Back to evaluation
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:

Showing the first 1,200 of 1396 lines. Download the file for the rest.

codekingpro/portable-devtools · Team Ai