
    jpK                     |   d dl Z d dlZd dlZd dlmZmZmZmZ d dlZd dl	m
Z
 d dlmZ d dlmZmZmZ d dlmZ d dlmZ d dlmZ d dlmZmZmZ d dlmZmZ d d	lmZ d
dl m!Z!m"Z"m#Z#m$Z$ d
dl%m&Z& ddl'm(Z(m)Z)  ejT                  e+      Z,dZ-dej\                  dej^                  dej^                  fdZ0 G d d      Z1y)    N)ListOptionalIteratorCallable)opik_clienttracelocal_recording)dataset_item)
experiment)execution_policy)rest_operations	test_casetest_result)LLMTaskScoringKeyMappingType)models   )evaluation_tasks_executorexception_analyzerhelpersmetrics_evaluator)EvaluationTask   )base_metricscore_resultevaluation_taskitemdefault_policyreturnc                    | j                   |S | j                   j                  | j                   j                  n|j                  dd      | j                   j                  | j                   j                  dS |j                  dd      dS )aV  
    Get execution policy for a dataset item.

    If the item has its own execution policy, merge it with the default.
    Item-level values override default values.

    Args:
        item: The dataset item.
        default_policy: Default execution policy from suite level.

    Returns:
        Merged execution policy for this item.
    runs_per_itemr   pass_thresholdr!   r"   )r   r!   getr"   )r   r   s     v/Users/manta/Documents/Projects/TheRoad-I1/backend/.venv/lib/python3.12/site-packages/opik/evaluation/engine/engine.pyget_item_execution_policyr&      s    " $
 $$22> !!//##OQ7 $$33? !!00   ##$4a8     c                      e Zd ZdZdej
                  dee   dededdf
dZ	 e
j                  d	g d
      	 d-dej                  deej                      dedee   dedej&                  fd       Z e
j                  dddg      dedej,                  dej                  dej0                  deej4                     f
d       Zdej:                  dededeej@                     deej                      dedee   dej&                  fdZ!de"ej:                     dedeej@                     deej                      dedee   dedee   de#jH                  de%deej&                     fd Z&d!ej&                  d"eejN                     dej0                  dej&                  fd#Z(d$eej&                     d%e)jT                  dej0                  ddfd&Z+de"ej:                     dedeej@                     deej                      d'eej                      dedee   dedee   de#jH                  de%deej&                     fd(Z,	 	 d.de"ej:                     ded)eej                      dee   dee   deej@                     de#jH                  dee   dede%deej&                     fd*Z-d+eej                     d)eej                      dee   deej&                     fd,Z.y)/EvaluationEnginez
    Stateless evaluation executor.

    Only stores configuration (client, workers, verbosity).
    All flow-specific data (metrics, key mappings) is passed as method parameters.
    clientproject_nameworkersverboser   Nc                 <    || _         || _        || _        || _        y )N)_client_project_name_workers_verbose)selfr*   r+   r,   r-   s        r%   __init__zEvaluationEngine.__init__E   s!     )r'   metrics_calculation)regular_metricsscoring_key_mappingevaluator_model)nameignore_arguments
test_case_r6   r7   r8   trial_idc                 L   t        j                  |j                  |||      }|j                  |j                  |j
                        \  }}||_        t        j                  |||      }	t        j                  | j                  ||j                  | j                         |	S )N)r   r6   r7   r8   )dataset_item_contenttask_output)r   score_resultsr<   r*   r@   trace_idr+   )r   build_metrics_evaluatorr
   compute_regular_scoresr>   r?   mapped_scoring_inputsr   
TestResultr   log_test_result_feedback_scoresr/   rB   r0   )
r3   r;   r6   r7   r8   r<   item_evaluatorr@   rE   test_result_s
             r%   "_compute_test_result_for_test_casez3EvaluationEngine._compute_test_result_for_test_caseS   s      +BB((+ 3+	
 0>/T/T!+!@!@".. 0U 0
,, ,A
("-- '

 	77<<'((++		
 r'   task_span_metrics_calculationtask_span_evaluatorrB   	task_spanc                     |j                  |j                  |j                  |      \  }}||_        t	        j
                  | j                  ||| j                         |S )N)r>   r?   rM   rA   )compute_task_span_scoresr>   r?   rE   r   rG   r/   r0   )r3   rB   rM   r;   rL   r@   rE   s          r%   ,_compute_scores_for_test_case_with_task_spanz=EvaluationEngine._compute_scores_for_test_case_with_task_span|   so      88%/%D%D&22# 9  	-, ,A
( 	77<<'++		
 r'   r   taskexperiment_c                 V   t        |d      s6t        |d      r|j                  nd} t        j                  |      |      }|j	                  d      }	t        j                  |	t        d| j                        }
d }|j                  |j                  j                  d	      }t        j                  ||j                  |
| j                  |xs d 
      5  t        j!                  d|	       t#        j$                         }	  ||	      }t#        j$                         |z
  }t        j!                  d|       t3        j4                  |       t7        j8                  |
j                  |j                  ||	|      }t#        j$                         }| j;                  |||||      }||_        t#        j$                         |z
  |_        d d d        |S # t&        $ r>}t)        j*                  |      r#t        j-                  t.        j0                          d }~ww xY w# 1 sw Y   S xY w)Nopik_tracked__name__llm_task)r9   T)
include_id
evaluation)inputr9   
created_byr+   )exclude_none)r   dataset_item_id
trace_datar*   r   zTask started, input: %szTask finished, output: %s)output)rB   r\   r?   r>   r
   )r;   r6   r7   r8   r<   ) hasattrrU   opiktrackget_contentr   	TraceDataEVALUATION_TASK_NAMEr0   r   
model_dumpr   evaluate_llm_task_contextidr/   LOGGERdebugtimeperf_counter	Exceptionr    is_llm_provider_rate_limit_errorerrorlogging_messages;LLM_PROVIDER_RATE_LIMIT_ERROR_DETECTED_IN_EVALUATE_FUNCTIONopik_contextupdate_current_tracer   TestCaserJ   task_execution_timescoring_time)r3   r   rQ   r<   rR   r6   r7   r8   r9   item_contentr]   execution_policy_dict
task_starttask_output_	exceptionrt   r;   scoring_startrI   s                      r%   !_compute_test_result_for_llm_taskz2EvaluationEngine._compute_test_result_for_llm_task   s    t^,$+D*$=4==:D(4::4(.D''4'8__%#++	

 !%  ,$($9$9$D$DRV$D$W!.." GG!<<2:d
 LL2LA**,J#L1 #'"3"3"5
"BLL4lC--\B"++# $(%1!J !--/MBB% /$7 /! C L 0CL,(,(9(9(;m(KL%O
R =  %FFyQLL(dd 
R s1   +HGB>H	H9HHHH(dataset_itemsdescriptiontotal_itemsdefault_execution_policyshow_scores_in_progress_barc                 J   t        j                  t        j                     | j                  | j
                  |||
      5 }|D ]  }t        ||	      }|j                  dd      }t        j                  ||j                  dd            |_
        |j                  |j                  |       t        |      D ]D  }|j                  t        j                   | j"                  |||||||      |j                         F  |j%                         cddd       S # 1 sw Y   yxY w)	z
        Execute tasks with full parallelism and item-based progress.

        Uses StreamingExecutor with group_id to track progress per item
        (not per run), so that runs_per_item > 1 shows correct item count.
        )r,   r-   desctotalshow_score_postfixr!   r   r"   r#   )r   rQ   r<   rR   r6   r7   r8   )group_idN)r   StreamingExecutorr   rF   r1   r2   r&   r$   r
   ExecutionPolicyItemr   set_group_sizerg   rangesubmit	functoolspartialr|   get_results)r3   r}   rQ   rR   r6   r7   r8   r~   r   r   r   executorr   item_policy	item_runsrun_ids                   r%   +_compute_test_results_with_execution_policyz<EvaluationEngine._compute_test_results_with_execution_policy   s   & '889O9OPMMMM:
 %7>VW'OOOQ?	 )5(H(H"+#.??3CQ#G)% ''; $I.FOO!)) BB!%!%%+(3,;0C,;	 "& $  / &: '')I
 
 
s    CDD"evaluation_task_resulttrace_treesc                    |j                   j                  }d }|D ]  }|j                  |k(  s|} n |t        d|       t	        |j
                        dk(  rt        d|       |j
                  d   }t        j                  t        j                  ||j                  |j                  |j                  |j                  |j                  |j                  | j                   d|j"                  |j$                        | j&                        5  | j)                  |||j                   |      }|xj*                  |z  c_        |cd d d        S # 1 sw Y   y xY w)Nz No trace found for test result: r   z}Task trace contains no spans. Task span metrics require at least one span to be present in the execution trace. Test result: rX   )rg   r9   
start_timemetadatarY   r^   tagsr+   rZ   
error_info	thread_id)r]   r*   )rB   rM   r;   rL   )r   rB   rg   
ValueErrorlenspansr   &evaluate_llm_task_result_spans_contextr   rc   r9   r   r   rY   r^   r   r0   r   r   r/   rP   r@   )	r3   r   r   rL   rB   
task_tracetrace_evaluation_spanr@   s	            r%   *_update_test_result_with_task_span_metricsz;EvaluationEngine._update_test_result_with_task_span_metrics  sm    *33<<
!FyyH$#
 "
 23I2JK 
 z A% P  Qg  Ph  i  %**1-;;__%00#,, &&!((__!//'%00$.. <<
  !MM!)1;;$7	 N M #00MA0)1
 
 
s   6EEtest_results	recordingc           	      R   |j                   }t        |      dk(  rt        j                  d       y|D cg c]%  }t	        j
                  | j                  |||      ' }}t        j                  || j                  | j                  d       t        j                  d       yc c}w )z+Evaluate task spans from a local recording.r   z,No trace trees found in the local recording.N)r   r   rL   zLLM task spans evaluation)evaluation_tasksr,   r-   r   uP   Task evaluation span handling is disabled — the evaluation has been completed.)r   r   rh   warningr   r   r   r   executer1   r2   ri   )r3   r   r   rL   r   rI   span_evaluation_taskss          r%   +_update_test_results_with_task_span_metricsz<EvaluationEngine._update_test_results_with_task_span_metricsR  s      ++{q NNIJ !-O
 !- ??'3'$7	 !- 	 O
 	"))2MMMM,		
 	^	
#O
s   *B$task_span_metricsc                 z   t        j                  | j                  ||||||||	|
|      }|s |       S t        j	                  dt        |             t        j                  ||      }t        j                  | j                        5 } |       }| j                  |||       ddd       |S # 1 sw Y   S xY w)z
        Shared execution logic. Runs the task, scores with regular metrics,
        and optionally scores with task-span metrics.
        )
r}   rQ   rR   r6   r7   r8   r~   r   r   r   u`   Detected %d LLM task span scoring metrics — enabling handling of the LLM task evaluation span.)scoring_metricsr7   )r*   )r   r   rL   N)r   r   r   rh   ri   r   r   MetricsEvaluatorr	   record_traces_locallyr/   r   )r3   r}   rQ   rR   r6   r   r7   r8   r~   r   r   r   computerL   r   r   s                   r%   _execute_evaluationz$EvaluationEngine._execute_evaluationw  s    $ ?H>O>O<<'#+ 3+##%=(C?
 !9n!"	

 0@@- 3

 22$,,G9"9L<<)#$7 =  H  H s   
B00B:r   c                 v    ||ni }t        j                  |      \  }}| j                  ||||||||	|||
      S )a%  
        Run a task on dataset items in parallel, then score results with metrics.

        This is the universal entry point for all evaluation flows that involve
        running a task. The caller is responsible for resolving dataset items
        and building the execution policy.
        )r}   rQ   rR   r6   r   r7   r8   r~   r   r   r   )r   (split_into_regular_and_task_span_metricsr   )r3   r}   rQ   r   r7   r8   rR   r   r   r~   r   resolved_scoring_key_mappingr6   r   s                 r%   run_and_scorezEvaluationEngine.run_and_score  si    * $7#B 	%
 FFW 	+* '''#+/ <+##%=(C ( 
 	
r'   
test_casesc           
         ||ni }t        j                  |      \  }}|D cg c]&  }t        j                  | j                  |||d      ( }}t        j                  || j                  | j                        }	|	S c c}w )z
        Score existing test cases with metrics (no task execution).

        This is the universal entry point for re-scoring existing experiment
        results with new metrics.
        N)r;   r6   r7   r8   )r   r,   r-   )	r   r   r   r   rJ   r   r   r1   r2   )
r3   r   r   r7   r   r6   _r;   r   r   s
             r%   score_test_casesz!EvaluationEngine.score_test_cases  s     $7#B 	% /WW
 )	J
 )
 77% /$@ $ ) 	 	J
 6O5V5V-MMMM6
 #	J
s   +A>)r   )
EvaluationT)/rU   
__module____qualname____doc__r   Opikr   strintr4   r`   ra   r   rs   r   r   
BaseMetricr   r   rF   rJ   r   	SpanModelr   r   r   ScoreResultrP   r
   DatasetItemr   r   
Experimentr|   r   dataset_execution_policyExecutionPolicyboolr   
TraceModelr   r	   _LocalRecordingHandler   r   r   r    r'   r%   r)   r)   =   s   
   
  sm
  	
 
 
  

  TZZ"
 && k445 3	
 "#  
		B TZZ,&(=> ## &&	
 /?? 
l&&	'	6C&&C C 	C
 j334C k445C 3C "#C 
		CN7* 8 897* 7* j334	7*
 k4457* 37* "#7* 7* c]7* #;"J"J7* &*7* 
k$$	%7*v3* + 6 63* &++,3* /??	3*
 
		3*j!
;112!
 #88!
 /??	!

 
!
J5 8 895 5 j334	5
 k4455   6 675 35 "#5 5 c]5 #;"J"J5 &*5 
k$$	%5F (,0(
 8 89(
 (
 k445	(

 &&;<(
 "#(
 j334(
 #;"J"J(
 c](
 (
 &*(
 
k$$	%(
T%++,% k445% &&;<	%
 
k$$	%%r'   r)   )2r   loggingrj   typingr   r   r   r   r`   opik.logging_messagesro   opik.opik_contextrq   opik.api_objectsr   r   r	   opik.api_objects.datasetr
   opik.api_objects.experimentr   r   r   opik.evaluationr   r   r   opik.evaluation.typesr   r   !opik.message_processing.emulationr    r   r   r   r   typesr   metricsr   r   	getLoggerrU   rh   rd   r   r   r&   r)   r   r'   r%   <module>r      s       5 5  0 ( @ @ 1 2 Q C C @ 4 W W ! / 
		8	$( 

"
",<< --DB Br'   