
    jP,                       d dl mZ d dlZd dlZd dlmZmZmZmZm	Z	m
Z
mZ d dlmZ d dl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 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" erd dl#m$Z$  ejJ                  e&      Z'	 	 	 	 	 d	 	 	 	 	 	 	 	 	 	 	 	 	 	 	 ddZ(	 	 	 	 	 	 	 	 ddZ)	 d	 	 	 	 	 	 	 ddZ*ddZ+	 	 	 	 	 	 	 	 	 	 	 	 ddZ,	 d	 	 	 	 	 	 	 	 	 	 	 	 	 d dZ-	 	 	 	 	 	 	 	 	 	 	 	 d!dZ.y)"    )annotationsN)AnyDictIteratorListOptionalSetTYPE_CHECKING)OpikApi)dataset_itemdataset_version_public)streamer)retry_decorator)opik_query_languagerest_stream_parser   )datasetr   execution_policy   )
experiment	constants)experiments_client   )ApiError)	llm_judgec              #     K   t         j                  dd}d}|rt        |      nd}	d|rFt        j                  j                  |      }
|
j                         }|rt        j                  |      |rt        j                  d fd       } |       }t        |      dk(  rd}n^|D ]T  }|j                  }||	||	vr|	j                  |       d}|j                  rM|j                  D cg c]8  }t        j                   |j"                  |j$                  |j&                        : }}d}|j(                  r?t        j*                  |j(                  j,                  |j(                  j.                        }t        j0                  d|j                  |j2                  |j4                  |j6                  |j8                  ||d|j:                  }| |d	z  }	|k\  rd} n|	Ct        |	      dk(  sSd} n |r|	r&t        |	      dkD  rt<        j?                  d
|	       yyyc c}w w)a  
    Stream dataset items from the backend as a generator.

    Args:
        rest_client: The REST API client.
        dataset_name: Name of the dataset to stream items from.
        nb_samples: Maximum number of items to retrieve. If None, all items are streamed.
        batch_size: Maximum number of items to fetch per batch from the backend.
        dataset_item_ids: Optional list of specific item IDs to retrieve.
        filter_string: Optional OQL filter string to filter dataset items.
        dataset_version: Optional dataset version hash to filter items by a specific version.

    Yields:
        DatasetItem objects one at a time.
    NTr   c            	         t        j                  j                  j                         t        j
                        S )N)dataset_namelast_retrieved_idsteam_limitfiltersdataset_version)stream
item_class
nb_samples)r   read_and_parse_streamdatasetsstream_dataset_itemsrest_dataset_item_readDatasetItem)
batch_sizer   r"   r!   r   r%   rest_clients   /Users/manta/Documents/Projects/TheRoad-I1/backend/.venv/lib/python3.12/site-packages/opik/api_objects/dataset/rest_operations.py_fetch_batchz*stream_dataset_items.<locals>._fetch_batchF   sM    %;;"++@@!-&7 *#$3 A  2==%
 
    Fnametypeconfigruns_per_itempass_threshold)idtrace_idspan_idsourcedescription
evaluatorsr   r   z=The following dataset items were not found in the dataset: %s)returnz(List[rest_dataset_item_read.DatasetItem] ) r   DATASET_STREAM_BATCH_SIZEsetr   OpikQueryLanguagefor_dataset_itemsget_filter_expressionsjsondumpsr   opik_rest_retrylenr7   remover<   r   EvaluatorItemr1   r2   r3   r   ExecutionPolicyItemr5   r6   r*   r8   r9   r:   r;   dataLOGGERwarning)r,   r   r%   r+   dataset_item_idsfilter_stringr"   should_retrieve_more_itemsitems_yieldeddataset_items_ids_leftoqlfilter_expressionsr.   dataset_itemsitemitem_idr<   er   reconstructed_itemr!   r   s   ````  `             @@r-   r(   r(      s^    0 88
'+!%M!1t  "G!33EEmT 779jj!34G
$		(	(	 	 
)	 %}").&!DggG '%1"88*11': J "__ - !..VVVV xx
 -    $$$#/#C#C"&"7"7"E"E#'#8#8#G#G$ 
 ".!9!9 	"77{{ ,,%!1	" ))	" %$QM%-:*E-2*%1c:P6QUV6V-2*e "- %T #&<"="AK"	
 #BOs%   C5I>=H=;B>I;II.Ic                    	 | j                   j                  ||      S # t        $ r}|j                  dk(  rY d}~y d}~ww xY w)a.  
    Find a dataset version by version name.

    Args:
        rest_client: The REST API client.
        dataset_id: The dataset ID to search versions in.
        version_name: Version name to search for (e.g., 'v1', 'v2').

    Returns:
        The DatasetVersionPublic if found, None otherwise.
    )r7   version_name  N)r'   retrieve_dataset_versionr   status_code)r,   
dataset_idr[   rX   s       r-   find_version_by_namer`      sN     ##<< = 
 	
  ==Cs    	A==Ac                   d}g }d}t        |      |k  r| j                  j                  ||      }t        |j                        dk(  r	 |S |j                  d |t        |      z
   D ]\  }t	        j
                  |j                  |j                  | |j                        }|r|j                          |j                  |       ^ |dz  }t        |      |k  r|S )Nd   r   )pagesizer   )r1   r;   r,   dataset_items_count)rG   r'   find_datasetscontentr   Datasetr1   r;   re   __internal_api__sync_hashes__append)	r,   max_results
sync_items	page_sizer'   rc   page_datasetsdataset_ferndataset_s	            r-   get_datasetsrq      s     I&(HD
h-+
%#,,:: ; 

 }$$%*" O *112Q[3x=5PRL!&&(44'$0$D$D	H 668OOH% S 		- h-+
%0 Or/   c                    	 | j                   j                  |      j                  }|S # t        $ r/}|j                  dk(  rt        j                  d| d      | d }~ww xY w)Nr   r\   zDataset with the name z not found.)r'   get_dataset_by_identifierr7   r   r^   
exceptionsDatasetNotFound)r,   r   r_   rX   s       r-   get_dataset_idrw      s{    	 ))CC% D 

" 	   ==C,,(kB 	s   &* 	A"*AA"c                   d}g }d}t        |      |k  r| j                  j                  |||      }t        |j                        dk(  r	 |S |j                  d |t        |      z
   D ]U  }	|j	                  t        j                  |	j                  |	j                  |	j                  | |||	j                               W |dz  }t        |      |k  r|S )Nrb   r   )rc   rd   r_   r   )r7   r1   r   r,   r   r   tags)rG   experimentsfind_experimentsrg   rj   r   
Experimentr7   r1   r   ry   )
r,   r_   rk   r   r   rm   rz   rc   page_experimentsexperiment_s
             r-   get_dataset_experimentsr      s     I/1KD
k
[
(&22CC! D 
 ''(A-"  ,334TkCDT6TUK%%"~~$))!,!9!9 +%'9$))
 V 		/ k
[
(2 r/   c                   | j                   j                  ||d|       | j                   j                  |      }|xs t        j                  j                         }i }|r?|D 	cg c]0  }	|	j                  d|	j                         j                  d      d2 c}	|d<   |j                  d	d
      |j                  dd
      d|d<   | j                   j                  |j                  |d       |j                  S c c}	w )a  
    Create a dataset of type 'evaluation_suite' and its initial version
    with evaluators and execution_policy persisted to the backend.

    Args:
        rest_client: The REST API client.
        dataset_name: The name of the dataset/suite.
        description: Optional description.
        evaluators: LLMJudge evaluators.
        exec_policy: Execution policy dict.

    Returns:
        The dataset ID.
    evaluation_suite)r1   r;   r2   ry   rs   r   Tby_aliasr0   r<   r5   r   r6   r4   r   r7   requestoverride)r'   create_datasetrt   r   DEFAULT_EXECUTION_POLICYcopyr1   	to_config
model_dumpgetapply_dataset_item_changesr7   )
r,   r   r;   r<   exec_policyry   ro   resolved_policyr   rX   s
             r-   create_evaluation_suite_datasetr     s   , ''{9KRV (  ''AA! B L "U%5%N%N%S%S%UO G  !
  	 #++-22D2A
  !
 ),,_a@)--.>B#G 33??Gd 4  ??!!
s   &5C=c           	        ||D cg c]0  }|j                   d|j                         j                  d      d2 c}|j                  dd      |j                  dd      dd	}| j                  j                  ||d
       yc c}w )a  
    Update suite-level evaluators and execution_policy by creating a new
    dataset version based on the current latest version.

    Args:
        rest_client: The REST API client.
        dataset_id: The dataset ID.
        base_version_id: The current latest version UUID to base the update on.
        evaluators: Suite-level LLMJudge evaluators.
        exec_policy: Execution policy dict.
    r   Tr   r0   r5   r   r6   r4   )base_versionr<   r   Fr   N)r1   r   r   r   r'   r   )r,   r_   base_version_idr<   r   rX   r   s          r-   update_evaluation_suite_datasetr   ;  s    & (  
  	 #++-22D2A
  
 )___a@)oo.>B
G 33w 4 
s   5B)NNNNN)r,   r   r   strr%   Optional[int]r+   r   rN   Optional[List[str]]rO   Optional[str]r"   r   r=   z"Iterator[dataset_item.DatasetItem])r,   r   r_   r   r[   r   r=   z5Optional[dataset_version_public.DatasetVersionPublic])i  T)r,   r   rk   intrl   boolr=   zList[dataset.Dataset])r,   r   r   r   r=   r   )r,   r   r_   r   rk   r   r   zstreamer.Streamerr   z$experiments_client.ExperimentsClientr=   zList[experiment.Experiment])N)r,   r   r   r   r;   r   r<   z"Optional[List[llm_judge.LLMJudge]]r   z*Optional[execution_policy.ExecutionPolicy]ry   r   r=   r   )r,   r   r_   r   r   r   r<   zList[llm_judge.LLMJudge]r   z execution_policy.ExecutionPolicyr=   None)/
__future__r   rD   loggingtypingr   r   r   r   r   r	   r
   opik.rest_apir   opik.rest_api.typesr   r)   r   opik.exceptionsru   opik.message_processingr   opik.rest_client_configuratorr   opik.api_objectsr   r    r   r   r   r   r   rest_api.core.api_errorr    opik.evaluation.suite_evaluatorsr   	getLogger__name__rL   r(   r`   rq   rw   r   r   r   r>   r/   r-   <module>r      s   "   J J J ! % , 9 D 5 5 $ + /:			8	$ !% $,0#'%)w
w
w
 w
 	w

 *w
 !w
 #w
 (w
t  ;	6 GK'*?CD$$$ $  	$
 =$ !$Z !%111 1 3	1
 <1 1 	1h### # )	#
 2# 
#r/   