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
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
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.
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
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()
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
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
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.
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