kiln_ai.adapters.eval.eval_runner

  1import json
  2import logging
  3from dataclasses import dataclass
  4from typing import AsyncGenerator, Dict, List, Literal, Set
  5
  6import litellm
  7
  8from kiln_ai.adapters.adapter_registry import load_skills_for_task
  9from kiln_ai.adapters.errors import KilnRunError
 10from kiln_ai.adapters.eval.base_eval import BaseEval
 11from kiln_ai.adapters.eval.registry import eval_adapter_from_type
 12from kiln_ai.adapters.model_adapters.base_adapter import SkillsDict
 13from kiln_ai.datamodel.basemodel import ID_TYPE
 14from kiln_ai.datamodel.dataset_filters import DatasetFilterId, dataset_filter_from_id
 15from kiln_ai.datamodel.eval import EvalConfig, EvalDataType, EvalRun, EvalScores
 16from kiln_ai.datamodel.task import TaskRunConfig
 17from kiln_ai.datamodel.task_run import TaskRun, Usage
 18from kiln_ai.utils.async_job_runner import AsyncJobRunner, Progress, RetryableError
 19from kiln_ai.utils.git_sync_protocols import SaveContext, default_save_context
 20
 21logger = logging.getLogger(__name__)
 22
 23
 24@dataclass
 25class EvalJob:
 26    item: TaskRun
 27    type: Literal["task_run_eval", "eval_config_eval"]
 28    # If type == "task_run_eval", both of these should be set. If type == "eval_config_eval", only eval_config should be set.
 29    eval_config: EvalConfig
 30    task_run_config: TaskRunConfig | None = None
 31
 32
 33class EvalRunner:
 34    """
 35    Runs an eval. Async execution is supported to make it faster when using remote/fast model providers.
 36
 37    Can run an eval in 2 modes:
 38    1) eval_config_eval: evaluate an eval config using existing dataset items.
 39    2) task_run_eval: evaluate a range of task run configs, generating new run output using existing dataset item input.
 40    """
 41
 42    def __init__(
 43        self,
 44        eval_configs: List[EvalConfig],
 45        run_configs: List[TaskRunConfig] | None,
 46        eval_run_type: Literal["eval_config_eval", "task_run_eval"],
 47        save_context: SaveContext | None = None,
 48    ):
 49        if len(eval_configs) == 0:
 50            raise ValueError("Eval runner requires at least one eval config")
 51        target_eval = eval_configs[0].parent_eval()
 52        if target_eval is None:
 53            raise ValueError("Eval config requires a parent eval")
 54        for eval_config in eval_configs:
 55            parent_eval = eval_config.parent_eval()
 56            if parent_eval is None:
 57                raise ValueError("Eval config requires a parent eval")
 58            if parent_eval.id != target_eval.id:
 59                raise ValueError("All eval configs must have the same parent eval")
 60
 61        target_task = target_eval.parent_task()
 62        if target_task is None:
 63            raise ValueError("Eval config requires a (grand)parent task")
 64
 65        # Check that run_configs is compatible
 66        if eval_run_type == "task_run_eval":
 67            if run_configs is None or len(run_configs) == 0:
 68                raise ValueError("Task run eval requires run configs")
 69            for run_config in run_configs:
 70                parent_task = run_config.parent_task()
 71                if parent_task is None:
 72                    raise ValueError("All run configs must have a parent task")
 73                if parent_task.id != target_task.id:
 74                    raise ValueError(
 75                        "Run config is not for the same task as the eval configs"
 76                    )
 77        else:
 78            if run_configs is not None:
 79                raise ValueError("Mode 'eval_config_eval' does not support run configs")
 80
 81        self.eval_run_type = eval_run_type
 82        self.eval_configs = eval_configs
 83        self.run_configs = run_configs
 84        self.task = target_task
 85        self.eval = target_eval
 86        self._skills: SkillsDict = self._preload_skills()
 87        self._save_context: SaveContext = save_context or default_save_context
 88
 89    def collect_tasks(self) -> List[EvalJob]:
 90        if self.eval_run_type == "eval_config_eval":
 91            if self.eval.eval_configs_filter_id is not None:
 92                return self.collect_tasks_for_eval_config_eval(
 93                    self.eval.eval_configs_filter_id
 94                )
 95            else:
 96                raise ValueError(
 97                    "Eval configs filter ID is required for eval runs of type 'eval_config_eval'"
 98                )
 99
100        else:
101            return self.collect_tasks_for_task_run_eval()
102
103    def collect_tasks_for_eval_config_eval(
104        self, eval_configs_filter_id: DatasetFilterId
105    ) -> List[EvalJob]:
106        """
107        Collect all jobs for this run, excluding any that have already been run.
108
109        This variant is used for mode "eval_config_eval", using existing dataset run data (input/output).
110
111        The tasks:
112        - should be in the eval config set filter
113        - should not have already been run for this eval config + dataset item pair
114        """
115        filter = dataset_filter_from_id(eval_configs_filter_id)
116
117        # already_run[eval_config_id][dataset_id]
118        already_run: Dict[ID_TYPE, Set[ID_TYPE]] = {}
119        for eval_config in self.eval_configs:
120            already_run[eval_config.id] = set()
121            for run in eval_config.runs(readonly=True):
122                already_run[eval_config.id].add(run.dataset_id)
123
124        return [
125            EvalJob(
126                item=task_run,
127                eval_config=eval_config,
128                type="eval_config_eval",
129            )
130            for task_run in self.task.runs(readonly=True)
131            if filter(task_run)
132            for eval_config in self.eval_configs
133            if task_run.id not in already_run[eval_config.id]
134        ]
135
136    def collect_tasks_for_task_run_eval(self) -> List[EvalJob]:
137        """
138        Collect all jobs for this run, excluding any that have already been run.
139
140        This variant is used for mode "task_run_eval", generating new run output using existing dataset item input.
141
142        The tasks:
143        - should be in the eval set filter
144        - should not have already been run for this eval config + run config + dataset item
145        """
146        filter = dataset_filter_from_id(self.eval.eval_set_filter_id)
147
148        # already_run[eval_config_id][run_config_id][dataset_id]
149        already_run: Dict[ID_TYPE, Dict[ID_TYPE, Set[ID_TYPE]]] = {}
150        for eval_config in self.eval_configs:
151            already_run[eval_config.id] = {}
152            for run_config in self.run_configs or []:
153                already_run[eval_config.id][run_config.id] = set()
154            for run in eval_config.runs(readonly=True):
155                if (
156                    run.task_run_config_id is not None
157                    and run.task_run_config_id in already_run[eval_config.id]
158                ):
159                    already_run[eval_config.id][run.task_run_config_id].add(
160                        run.dataset_id
161                    )
162
163        return [
164            EvalJob(
165                item=task_run,
166                task_run_config=run_config,
167                type="task_run_eval",
168                eval_config=eval_config,
169            )
170            for task_run in self.task.runs(readonly=True)
171            if filter(task_run)
172            for eval_config in self.eval_configs
173            for run_config in self.run_configs or []
174            if task_run.id not in already_run[eval_config.id][run_config.id]
175        ]
176
177    def _preload_skills(self) -> SkillsDict:
178        """Collect all skill IDs from run configs and bulk-load them once."""
179        if self.run_configs is None:
180            return {}
181        merged: SkillsDict = {}
182        for rc in self.run_configs:
183            skills = load_skills_for_task(self.task, rc.run_config_properties)
184            merged.update(skills)
185        return merged
186
187    async def run(self, concurrency: int = 25) -> AsyncGenerator[Progress, None]:
188        """
189        Runs the configured eval run with parallel workers and yields progress updates.
190        """
191        jobs = self.collect_tasks()
192
193        runner = AsyncJobRunner(
194            concurrency=concurrency,
195            jobs=jobs,
196            run_job_fn=self.run_job,
197            max_retries=2,
198        )
199        async for progress in runner.run():
200            yield progress
201
202    async def run_job(self, job: EvalJob) -> bool:
203        try:
204            # Create the evaluator for this eval config/run config pair
205            evaluator = eval_adapter_from_type(job.eval_config.config_type)(
206                job.eval_config,
207                job.task_run_config.run_config_properties
208                if job.task_run_config
209                else None,
210                skills=self._skills,
211            )
212            if not isinstance(evaluator, BaseEval):
213                raise ValueError("Not able to create evaluator from eval config")
214
215            task_output: str | None = None
216            reference_answer: str | None = None
217            trace: str | None = None
218            scores: EvalScores | None = None
219            intermediate_outputs: Dict[str, str] | None = None
220            task_run_usage: Usage | None = None
221            if job.type == "eval_config_eval":
222                # Eval config eval, we use the saved input from the task run, not invoking the task again
223                scores, intermediate_outputs = await evaluator.run_eval(job.item)
224                task_output = job.item.output.output
225                task_run_usage = job.item.usage
226            else:
227                # Task run eval, we invoke the task again to get a fresh output
228                (
229                    result_task_run,
230                    scores,
231                    intermediate_outputs,
232                ) = await evaluator.run_task_and_eval(job.item)
233                task_output = result_task_run.output.output
234                task_run_usage = result_task_run.usage
235
236                parent_eval = job.eval_config.parent_eval()
237                if (
238                    parent_eval
239                    and parent_eval.evaluation_data_type == EvalDataType.full_trace
240                    and result_task_run.trace
241                ):
242                    trace = json.dumps(result_task_run.trace, indent=2)
243
244                if (
245                    parent_eval
246                    and parent_eval.evaluation_data_type
247                    == EvalDataType.reference_answer
248                ):
249                    reference_answer = job.item.output.output
250
251            # Save the job result
252            async with self._save_context():
253                eval_run = EvalRun(
254                    parent=job.eval_config,
255                    task_run_config_id=job.task_run_config.id
256                    if job.task_run_config
257                    else None,
258                    dataset_id=job.item.id,
259                    eval_config_eval=job.type == "eval_config_eval",
260                    scores=scores,
261                    input=job.item.input,
262                    output=task_output,
263                    reference_answer=reference_answer,
264                    intermediate_outputs=intermediate_outputs,
265                    task_run_trace=trace,
266                    task_run_usage=task_run_usage,
267                )
268                eval_run.save_to_file()
269
270            return True
271        except Exception as e:
272            if _is_retryable_error(e):
273                logger.error(
274                    f"Transient error running eval job for dataset item {job.item.id}: {e}",
275                    exc_info=True,
276                )
277                # KilnRunError's own message is genericized user-facing text; keep
278                # the underlying provider detail for the developer-facing error log.
279                raise RetryableError(str(_unwrap_kiln_run_error(e))) from e
280            logger.error(
281                f"Error running eval job for dataset item {job.item.id}: {e}",
282                exc_info=True,
283            )
284            raise
285
286
287def _unwrap_kiln_run_error(e: BaseException) -> BaseException:
288    """The innermost non-wrapper error.
289
290    The model adapter wraps provider exceptions in KilnRunError (to carry the
291    partial trace), whose own message is genericized user-facing text — so both
292    retry classification and error detail must use the underlying error. The
293    isinstance guard on `original` keeps a (contract-violating) None from
294    escaping as the result."""
295    while isinstance(e, KilnRunError) and isinstance(e.original, BaseException):
296        e = e.original
297    return e
298
299
300def _is_retryable_error(e: BaseException) -> bool:
301    e = _unwrap_kiln_run_error(e)
302
303    if isinstance(
304        e,
305        (
306            litellm.RateLimitError,
307            litellm.APIConnectionError,
308            litellm.InternalServerError,
309            litellm.ServiceUnavailableError,
310            litellm.BadGatewayError,
311            litellm.JSONSchemaValidationError,
312        ),
313    ):
314        return True
315
316    # ValueError thrown by Kiln's adapter when structured output doesn't match schema
317    if isinstance(
318        e, ValueError
319    ) and "This task requires a specific output schema" in str(e):
320        return True
321
322    return False
logger = <Logger kiln_ai.adapters.eval.eval_runner (WARNING)>
@dataclass
class EvalJob:
25@dataclass
26class EvalJob:
27    item: TaskRun
28    type: Literal["task_run_eval", "eval_config_eval"]
29    # If type == "task_run_eval", both of these should be set. If type == "eval_config_eval", only eval_config should be set.
30    eval_config: EvalConfig
31    task_run_config: TaskRunConfig | None = None
EvalJob( item: kiln_ai.datamodel.TaskRun, type: Literal['task_run_eval', 'eval_config_eval'], eval_config: kiln_ai.datamodel.eval.EvalConfig, task_run_config: kiln_ai.datamodel.task.TaskRunConfig | None = None)
type: Literal['task_run_eval', 'eval_config_eval']
task_run_config: kiln_ai.datamodel.task.TaskRunConfig | None = None
class EvalRunner:
 34class EvalRunner:
 35    """
 36    Runs an eval. Async execution is supported to make it faster when using remote/fast model providers.
 37
 38    Can run an eval in 2 modes:
 39    1) eval_config_eval: evaluate an eval config using existing dataset items.
 40    2) task_run_eval: evaluate a range of task run configs, generating new run output using existing dataset item input.
 41    """
 42
 43    def __init__(
 44        self,
 45        eval_configs: List[EvalConfig],
 46        run_configs: List[TaskRunConfig] | None,
 47        eval_run_type: Literal["eval_config_eval", "task_run_eval"],
 48        save_context: SaveContext | None = None,
 49    ):
 50        if len(eval_configs) == 0:
 51            raise ValueError("Eval runner requires at least one eval config")
 52        target_eval = eval_configs[0].parent_eval()
 53        if target_eval is None:
 54            raise ValueError("Eval config requires a parent eval")
 55        for eval_config in eval_configs:
 56            parent_eval = eval_config.parent_eval()
 57            if parent_eval is None:
 58                raise ValueError("Eval config requires a parent eval")
 59            if parent_eval.id != target_eval.id:
 60                raise ValueError("All eval configs must have the same parent eval")
 61
 62        target_task = target_eval.parent_task()
 63        if target_task is None:
 64            raise ValueError("Eval config requires a (grand)parent task")
 65
 66        # Check that run_configs is compatible
 67        if eval_run_type == "task_run_eval":
 68            if run_configs is None or len(run_configs) == 0:
 69                raise ValueError("Task run eval requires run configs")
 70            for run_config in run_configs:
 71                parent_task = run_config.parent_task()
 72                if parent_task is None:
 73                    raise ValueError("All run configs must have a parent task")
 74                if parent_task.id != target_task.id:
 75                    raise ValueError(
 76                        "Run config is not for the same task as the eval configs"
 77                    )
 78        else:
 79            if run_configs is not None:
 80                raise ValueError("Mode 'eval_config_eval' does not support run configs")
 81
 82        self.eval_run_type = eval_run_type
 83        self.eval_configs = eval_configs
 84        self.run_configs = run_configs
 85        self.task = target_task
 86        self.eval = target_eval
 87        self._skills: SkillsDict = self._preload_skills()
 88        self._save_context: SaveContext = save_context or default_save_context
 89
 90    def collect_tasks(self) -> List[EvalJob]:
 91        if self.eval_run_type == "eval_config_eval":
 92            if self.eval.eval_configs_filter_id is not None:
 93                return self.collect_tasks_for_eval_config_eval(
 94                    self.eval.eval_configs_filter_id
 95                )
 96            else:
 97                raise ValueError(
 98                    "Eval configs filter ID is required for eval runs of type 'eval_config_eval'"
 99                )
100
101        else:
102            return self.collect_tasks_for_task_run_eval()
103
104    def collect_tasks_for_eval_config_eval(
105        self, eval_configs_filter_id: DatasetFilterId
106    ) -> List[EvalJob]:
107        """
108        Collect all jobs for this run, excluding any that have already been run.
109
110        This variant is used for mode "eval_config_eval", using existing dataset run data (input/output).
111
112        The tasks:
113        - should be in the eval config set filter
114        - should not have already been run for this eval config + dataset item pair
115        """
116        filter = dataset_filter_from_id(eval_configs_filter_id)
117
118        # already_run[eval_config_id][dataset_id]
119        already_run: Dict[ID_TYPE, Set[ID_TYPE]] = {}
120        for eval_config in self.eval_configs:
121            already_run[eval_config.id] = set()
122            for run in eval_config.runs(readonly=True):
123                already_run[eval_config.id].add(run.dataset_id)
124
125        return [
126            EvalJob(
127                item=task_run,
128                eval_config=eval_config,
129                type="eval_config_eval",
130            )
131            for task_run in self.task.runs(readonly=True)
132            if filter(task_run)
133            for eval_config in self.eval_configs
134            if task_run.id not in already_run[eval_config.id]
135        ]
136
137    def collect_tasks_for_task_run_eval(self) -> List[EvalJob]:
138        """
139        Collect all jobs for this run, excluding any that have already been run.
140
141        This variant is used for mode "task_run_eval", generating new run output using existing dataset item input.
142
143        The tasks:
144        - should be in the eval set filter
145        - should not have already been run for this eval config + run config + dataset item
146        """
147        filter = dataset_filter_from_id(self.eval.eval_set_filter_id)
148
149        # already_run[eval_config_id][run_config_id][dataset_id]
150        already_run: Dict[ID_TYPE, Dict[ID_TYPE, Set[ID_TYPE]]] = {}
151        for eval_config in self.eval_configs:
152            already_run[eval_config.id] = {}
153            for run_config in self.run_configs or []:
154                already_run[eval_config.id][run_config.id] = set()
155            for run in eval_config.runs(readonly=True):
156                if (
157                    run.task_run_config_id is not None
158                    and run.task_run_config_id in already_run[eval_config.id]
159                ):
160                    already_run[eval_config.id][run.task_run_config_id].add(
161                        run.dataset_id
162                    )
163
164        return [
165            EvalJob(
166                item=task_run,
167                task_run_config=run_config,
168                type="task_run_eval",
169                eval_config=eval_config,
170            )
171            for task_run in self.task.runs(readonly=True)
172            if filter(task_run)
173            for eval_config in self.eval_configs
174            for run_config in self.run_configs or []
175            if task_run.id not in already_run[eval_config.id][run_config.id]
176        ]
177
178    def _preload_skills(self) -> SkillsDict:
179        """Collect all skill IDs from run configs and bulk-load them once."""
180        if self.run_configs is None:
181            return {}
182        merged: SkillsDict = {}
183        for rc in self.run_configs:
184            skills = load_skills_for_task(self.task, rc.run_config_properties)
185            merged.update(skills)
186        return merged
187
188    async def run(self, concurrency: int = 25) -> AsyncGenerator[Progress, None]:
189        """
190        Runs the configured eval run with parallel workers and yields progress updates.
191        """
192        jobs = self.collect_tasks()
193
194        runner = AsyncJobRunner(
195            concurrency=concurrency,
196            jobs=jobs,
197            run_job_fn=self.run_job,
198            max_retries=2,
199        )
200        async for progress in runner.run():
201            yield progress
202
203    async def run_job(self, job: EvalJob) -> bool:
204        try:
205            # Create the evaluator for this eval config/run config pair
206            evaluator = eval_adapter_from_type(job.eval_config.config_type)(
207                job.eval_config,
208                job.task_run_config.run_config_properties
209                if job.task_run_config
210                else None,
211                skills=self._skills,
212            )
213            if not isinstance(evaluator, BaseEval):
214                raise ValueError("Not able to create evaluator from eval config")
215
216            task_output: str | None = None
217            reference_answer: str | None = None
218            trace: str | None = None
219            scores: EvalScores | None = None
220            intermediate_outputs: Dict[str, str] | None = None
221            task_run_usage: Usage | None = None
222            if job.type == "eval_config_eval":
223                # Eval config eval, we use the saved input from the task run, not invoking the task again
224                scores, intermediate_outputs = await evaluator.run_eval(job.item)
225                task_output = job.item.output.output
226                task_run_usage = job.item.usage
227            else:
228                # Task run eval, we invoke the task again to get a fresh output
229                (
230                    result_task_run,
231                    scores,
232                    intermediate_outputs,
233                ) = await evaluator.run_task_and_eval(job.item)
234                task_output = result_task_run.output.output
235                task_run_usage = result_task_run.usage
236
237                parent_eval = job.eval_config.parent_eval()
238                if (
239                    parent_eval
240                    and parent_eval.evaluation_data_type == EvalDataType.full_trace
241                    and result_task_run.trace
242                ):
243                    trace = json.dumps(result_task_run.trace, indent=2)
244
245                if (
246                    parent_eval
247                    and parent_eval.evaluation_data_type
248                    == EvalDataType.reference_answer
249                ):
250                    reference_answer = job.item.output.output
251
252            # Save the job result
253            async with self._save_context():
254                eval_run = EvalRun(
255                    parent=job.eval_config,
256                    task_run_config_id=job.task_run_config.id
257                    if job.task_run_config
258                    else None,
259                    dataset_id=job.item.id,
260                    eval_config_eval=job.type == "eval_config_eval",
261                    scores=scores,
262                    input=job.item.input,
263                    output=task_output,
264                    reference_answer=reference_answer,
265                    intermediate_outputs=intermediate_outputs,
266                    task_run_trace=trace,
267                    task_run_usage=task_run_usage,
268                )
269                eval_run.save_to_file()
270
271            return True
272        except Exception as e:
273            if _is_retryable_error(e):
274                logger.error(
275                    f"Transient error running eval job for dataset item {job.item.id}: {e}",
276                    exc_info=True,
277                )
278                # KilnRunError's own message is genericized user-facing text; keep
279                # the underlying provider detail for the developer-facing error log.
280                raise RetryableError(str(_unwrap_kiln_run_error(e))) from e
281            logger.error(
282                f"Error running eval job for dataset item {job.item.id}: {e}",
283                exc_info=True,
284            )
285            raise

Runs an eval. Async execution is supported to make it faster when using remote/fast model providers.

Can run an eval in 2 modes: 1) eval_config_eval: evaluate an eval config using existing dataset items. 2) task_run_eval: evaluate a range of task run configs, generating new run output using existing dataset item input.

EvalRunner( eval_configs: List[kiln_ai.datamodel.eval.EvalConfig], run_configs: Optional[List[kiln_ai.datamodel.task.TaskRunConfig]], eval_run_type: Literal['eval_config_eval', 'task_run_eval'], save_context: Callable[[], contextlib.AbstractAsyncContextManager[None]] | None = None)
43    def __init__(
44        self,
45        eval_configs: List[EvalConfig],
46        run_configs: List[TaskRunConfig] | None,
47        eval_run_type: Literal["eval_config_eval", "task_run_eval"],
48        save_context: SaveContext | None = None,
49    ):
50        if len(eval_configs) == 0:
51            raise ValueError("Eval runner requires at least one eval config")
52        target_eval = eval_configs[0].parent_eval()
53        if target_eval is None:
54            raise ValueError("Eval config requires a parent eval")
55        for eval_config in eval_configs:
56            parent_eval = eval_config.parent_eval()
57            if parent_eval is None:
58                raise ValueError("Eval config requires a parent eval")
59            if parent_eval.id != target_eval.id:
60                raise ValueError("All eval configs must have the same parent eval")
61
62        target_task = target_eval.parent_task()
63        if target_task is None:
64            raise ValueError("Eval config requires a (grand)parent task")
65
66        # Check that run_configs is compatible
67        if eval_run_type == "task_run_eval":
68            if run_configs is None or len(run_configs) == 0:
69                raise ValueError("Task run eval requires run configs")
70            for run_config in run_configs:
71                parent_task = run_config.parent_task()
72                if parent_task is None:
73                    raise ValueError("All run configs must have a parent task")
74                if parent_task.id != target_task.id:
75                    raise ValueError(
76                        "Run config is not for the same task as the eval configs"
77                    )
78        else:
79            if run_configs is not None:
80                raise ValueError("Mode 'eval_config_eval' does not support run configs")
81
82        self.eval_run_type = eval_run_type
83        self.eval_configs = eval_configs
84        self.run_configs = run_configs
85        self.task = target_task
86        self.eval = target_eval
87        self._skills: SkillsDict = self._preload_skills()
88        self._save_context: SaveContext = save_context or default_save_context
eval_run_type
eval_configs
run_configs
task
eval
def collect_tasks(self) -> List[EvalJob]:
 90    def collect_tasks(self) -> List[EvalJob]:
 91        if self.eval_run_type == "eval_config_eval":
 92            if self.eval.eval_configs_filter_id is not None:
 93                return self.collect_tasks_for_eval_config_eval(
 94                    self.eval.eval_configs_filter_id
 95                )
 96            else:
 97                raise ValueError(
 98                    "Eval configs filter ID is required for eval runs of type 'eval_config_eval'"
 99                )
100
101        else:
102            return self.collect_tasks_for_task_run_eval()
def collect_tasks_for_eval_config_eval( self, eval_configs_filter_id: Annotated[str, AfterValidator(func=<function <lambda>>)]) -> List[EvalJob]:
104    def collect_tasks_for_eval_config_eval(
105        self, eval_configs_filter_id: DatasetFilterId
106    ) -> List[EvalJob]:
107        """
108        Collect all jobs for this run, excluding any that have already been run.
109
110        This variant is used for mode "eval_config_eval", using existing dataset run data (input/output).
111
112        The tasks:
113        - should be in the eval config set filter
114        - should not have already been run for this eval config + dataset item pair
115        """
116        filter = dataset_filter_from_id(eval_configs_filter_id)
117
118        # already_run[eval_config_id][dataset_id]
119        already_run: Dict[ID_TYPE, Set[ID_TYPE]] = {}
120        for eval_config in self.eval_configs:
121            already_run[eval_config.id] = set()
122            for run in eval_config.runs(readonly=True):
123                already_run[eval_config.id].add(run.dataset_id)
124
125        return [
126            EvalJob(
127                item=task_run,
128                eval_config=eval_config,
129                type="eval_config_eval",
130            )
131            for task_run in self.task.runs(readonly=True)
132            if filter(task_run)
133            for eval_config in self.eval_configs
134            if task_run.id not in already_run[eval_config.id]
135        ]

Collect all jobs for this run, excluding any that have already been run.

This variant is used for mode "eval_config_eval", using existing dataset run data (input/output).

The tasks:

  • should be in the eval config set filter
  • should not have already been run for this eval config + dataset item pair
def collect_tasks_for_task_run_eval(self) -> List[EvalJob]:
137    def collect_tasks_for_task_run_eval(self) -> List[EvalJob]:
138        """
139        Collect all jobs for this run, excluding any that have already been run.
140
141        This variant is used for mode "task_run_eval", generating new run output using existing dataset item input.
142
143        The tasks:
144        - should be in the eval set filter
145        - should not have already been run for this eval config + run config + dataset item
146        """
147        filter = dataset_filter_from_id(self.eval.eval_set_filter_id)
148
149        # already_run[eval_config_id][run_config_id][dataset_id]
150        already_run: Dict[ID_TYPE, Dict[ID_TYPE, Set[ID_TYPE]]] = {}
151        for eval_config in self.eval_configs:
152            already_run[eval_config.id] = {}
153            for run_config in self.run_configs or []:
154                already_run[eval_config.id][run_config.id] = set()
155            for run in eval_config.runs(readonly=True):
156                if (
157                    run.task_run_config_id is not None
158                    and run.task_run_config_id in already_run[eval_config.id]
159                ):
160                    already_run[eval_config.id][run.task_run_config_id].add(
161                        run.dataset_id
162                    )
163
164        return [
165            EvalJob(
166                item=task_run,
167                task_run_config=run_config,
168                type="task_run_eval",
169                eval_config=eval_config,
170            )
171            for task_run in self.task.runs(readonly=True)
172            if filter(task_run)
173            for eval_config in self.eval_configs
174            for run_config in self.run_configs or []
175            if task_run.id not in already_run[eval_config.id][run_config.id]
176        ]

Collect all jobs for this run, excluding any that have already been run.

This variant is used for mode "task_run_eval", generating new run output using existing dataset item input.

The tasks:

  • should be in the eval set filter
  • should not have already been run for this eval config + run config + dataset item
async def run( self, concurrency: int = 25) -> AsyncGenerator[kiln_ai.utils.async_job_runner.Progress, NoneType]:
188    async def run(self, concurrency: int = 25) -> AsyncGenerator[Progress, None]:
189        """
190        Runs the configured eval run with parallel workers and yields progress updates.
191        """
192        jobs = self.collect_tasks()
193
194        runner = AsyncJobRunner(
195            concurrency=concurrency,
196            jobs=jobs,
197            run_job_fn=self.run_job,
198            max_retries=2,
199        )
200        async for progress in runner.run():
201            yield progress

Runs the configured eval run with parallel workers and yields progress updates.

async def run_job(self, job: EvalJob) -> bool:
203    async def run_job(self, job: EvalJob) -> bool:
204        try:
205            # Create the evaluator for this eval config/run config pair
206            evaluator = eval_adapter_from_type(job.eval_config.config_type)(
207                job.eval_config,
208                job.task_run_config.run_config_properties
209                if job.task_run_config
210                else None,
211                skills=self._skills,
212            )
213            if not isinstance(evaluator, BaseEval):
214                raise ValueError("Not able to create evaluator from eval config")
215
216            task_output: str | None = None
217            reference_answer: str | None = None
218            trace: str | None = None
219            scores: EvalScores | None = None
220            intermediate_outputs: Dict[str, str] | None = None
221            task_run_usage: Usage | None = None
222            if job.type == "eval_config_eval":
223                # Eval config eval, we use the saved input from the task run, not invoking the task again
224                scores, intermediate_outputs = await evaluator.run_eval(job.item)
225                task_output = job.item.output.output
226                task_run_usage = job.item.usage
227            else:
228                # Task run eval, we invoke the task again to get a fresh output
229                (
230                    result_task_run,
231                    scores,
232                    intermediate_outputs,
233                ) = await evaluator.run_task_and_eval(job.item)
234                task_output = result_task_run.output.output
235                task_run_usage = result_task_run.usage
236
237                parent_eval = job.eval_config.parent_eval()
238                if (
239                    parent_eval
240                    and parent_eval.evaluation_data_type == EvalDataType.full_trace
241                    and result_task_run.trace
242                ):
243                    trace = json.dumps(result_task_run.trace, indent=2)
244
245                if (
246                    parent_eval
247                    and parent_eval.evaluation_data_type
248                    == EvalDataType.reference_answer
249                ):
250                    reference_answer = job.item.output.output
251
252            # Save the job result
253            async with self._save_context():
254                eval_run = EvalRun(
255                    parent=job.eval_config,
256                    task_run_config_id=job.task_run_config.id
257                    if job.task_run_config
258                    else None,
259                    dataset_id=job.item.id,
260                    eval_config_eval=job.type == "eval_config_eval",
261                    scores=scores,
262                    input=job.item.input,
263                    output=task_output,
264                    reference_answer=reference_answer,
265                    intermediate_outputs=intermediate_outputs,
266                    task_run_trace=trace,
267                    task_run_usage=task_run_usage,
268                )
269                eval_run.save_to_file()
270
271            return True
272        except Exception as e:
273            if _is_retryable_error(e):
274                logger.error(
275                    f"Transient error running eval job for dataset item {job.item.id}: {e}",
276                    exc_info=True,
277                )
278                # KilnRunError's own message is genericized user-facing text; keep
279                # the underlying provider detail for the developer-facing error log.
280                raise RetryableError(str(_unwrap_kiln_run_error(e))) from e
281            logger.error(
282                f"Error running eval job for dataset item {job.item.id}: {e}",
283                exc_info=True,
284            )
285            raise