
    iC                        d Z ddlmZ ddlZddl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 erdd
lmZ ddlmZ  G d dee      Zy)aG  
Vertex AI-specific RAG Ingestion implementation.

Vertex AI RAG Engine handles embedding and chunking internally when files are uploaded,
so this implementation skips the embedding step and directly uploads files to RAG corpora.

Based on: https://docs.cloud.google.com/vertex-ai/generative-ai/docs/model-reference/rag-api-v1
    )annotationsN)TYPE_CHECKINGAnyDictListOptionalTuple)verbose_logger)get_async_httpx_clienthttpxSpecialProvider)get_vertex_base_url)
VertexBase)BaseRAGIngestion)Router)RAGIngestOptionsc                      e Zd ZdZ	 d
	 	 	 ddZ	 	 	 	 ddZ	 	 	 	 	 	 	 	 	 	 	 	 ddZ	 d
	 	 	 	 	 ddZ	 	 d	 	 	 	 	 	 	 	 	 ddZ	 	 	 	 	 	 	 	 	 	 ddZ		 	 	 	 	 	 dd	Z
y)VertexAIRAGIngestiona  
    Vertex AI RAG Engine ingestion implementation.

    Key differences from base:
    - Embedding is handled by Vertex AI RAG Engine when files are uploaded
    - Files are uploaded using the RAG API (import or upload)
    - Chunking is done by Vertex AI RAG Engine (supports custom chunking config)
    - Supports Google Cloud Storage (GCS) and Google Drive sources
    - Supports custom parsing configurations (layout parser, LLM parser)
    Nc                   t        j                  | ||       t        j                  |        t        | j                        }| j                  |      | _        | j                  |      xs d| _        | j                  |      | _
        y )N)ingest_optionsrouterzus-central1)r   __init__r   dictvector_store_configsafe_get_vertex_ai_project
project_idget_vertex_ai_locationlocationsafe_get_vertex_ai_credentialsvertex_credentials)selfr   r   litellm_paramss       z/Users/manta/Documents/Projects/TheRoad-I1/.venv/lib/python3.12/site-packages/litellm/rag/ingestion/vertex_ai_ingestion.pyr   zVertexAIRAGIngestion.__init__)   st    
 	!!$~fUD! d667 99.I33NCT}"&"E"En"U    c                   K   yw)z
        Vertex AI RAG Engine handles embedding internally - skip this step.

        Returns:
            None (Vertex AI embeds when files are uploaded to RAG corpus)
        N )r    chunkss     r"   embedzVertexAIRAGIngestion.embed9   s      s   c                P  K   | j                   st        d      | j                  j                  d      }|sB| j	                  | j
                  xs d| j                  j                  d             d{   }d}|r!|r|r| j                  ||||       d{   }||fS 7 -7 
w)a
  
        Store content in Vertex AI RAG corpus.

        Vertex AI workflow:
        1. Create RAG corpus (if not provided)
        2. Upload file using RAG API (Vertex AI handles chunking/embedding)

        Args:
            file_content: Raw file bytes
            filename: Name of the file
            content_type: MIME type
            chunks: Ignored - Vertex AI handles chunking
            embeddings: Ignored - Vertex AI handles embedding

        Returns:
            Tuple of (rag_corpus_id, file_id)
        zVvertex_project is required for Vertex AI RAG ingestion. Set it in vector_store config.vector_store_idzlitellm-rag-corpusdescription)display_namer*   N)rag_corpus_idfilenamefile_contentcontent_type)r   
ValueErrorr   get_create_rag_corpusingest_name_upload_file_to_corpus)r    r.   r-   r/   r&   
embeddingsr,   result_file_ids           r"   storezVertexAIRAGIngestion.storeF   s     2 1  00445FG"&"9"9!--E1E 4488G #: # M H#'#>#>+!))	 $? $ N n,,s$   A2B&4B"5$B&B$	B&$B&c                  K   | j                  | j                  | j                  d      \  }}| j                  s|| _        t        | j                        }| d| j                   d| j                   d}d|i}|r||d<   | j
                  j                  d      }|r||d	<   | j
                  j                  d
      }	|	rd	|vri |d	<   dd|	ii|d	   d<   t        j                  d|        t        j                  dt        j                  |d              t        t        j                  ddi      }
|
j                  ||d| dd       d{   }|j                  dvr/d|j                    }t        j"                  |       t%        |      |j                         }t        j                  dt        j                  |d              |j                  d      r#|j                  di       j                  dd       }nE|j                  dd       }t        j                  d!|        | j'                  ||"       d{   }t        j                  d#|        |S 7 7 !w)$a  
        Create a Vertex AI RAG corpus.

        Args:
            display_name: Display name for the corpus
            description: Optional description

        Returns:
            RAG corpus ID (format: projects/{project}/locations/{location}/ragCorpora/{corpus_id})
        	vertex_aicredentialsr   custom_llm_providerz/v1beta1/projects/z/locations/z/ragCorporadisplayNamer*   vector_db_configvectorDbConfigembedding_modelvertexPredictionEndpointendpointragEmbeddingModelConfigzCreating RAG corpus: Request body:    indenttimeout      N@llm_providerparamsBearer application/jsonAuthorizationzContent-TypejsonheadersN      zFailed to create RAG corpus: zCreate corpus response: doneresponsename zPolling operation: )operation_nameaccess_tokenzCreated RAG corpus: )_ensure_access_tokenr   r   r   r   r   r1   r
   debugrR   dumpsr   r   RAGpoststatus_codetexterror	Exception_poll_operation)r    r+   r*   r\   r   base_urlurlrequest_bodyr>   r@   clientrX   	error_msgresponse_datacorpus_namer[   s                   r"   r2   z'VertexAIRAGIngestion._create_rag_corpusy   s      $(#<#<// + $= $
 j (DO 't}}5j (DMM?+O 	 <(
 *5L'  33778JK-=L)* 22667HI|313-.*Z,IIL)*+DE 	4SE:;~djja.P-QRS'-11t$

  #*<.!9 2 % 
 
 z17GI  +I&& &tzz-'J&KL	
 V$'++J;??KK +..vr:N  #6~6F!GH $ 4 4-) !5 ! K
 	3K=ABC
6s%   EI*I%C2I*I(I*(I*c                  K   ddl }t        | j                        }| d| }t        t        j
                  ddi      }t        |      D ]  }	|j                  |dd| i	       d{   }
|
j                  d
k7  r/d|
j                   }t        j                  |       t        |      |
j                         }|j                  d      rMd|v r|d   }t        d|       |j                  di       j                  dd      }|r|c S t        d|       t        j                  d|	dz    d|        |j                  |       d{    	 t        d| d      7 7 w)a  
        Poll a long-running operation until it completes.

        Args:
            operation_name: The operation name (e.g., "operations/123456")
            access_token: Access token for authentication
            max_retries: Maximum number of polling attempts
            retry_delay: Delay between polling attempts in seconds

        Returns:
            The corpus name from the completed operation

        Raises:
            Exception: If operation fails or times out
        r   N	/v1beta1/rH   rI   rJ   rP   rM   )rS   rU   zFailed to poll operation: rW   rd   zOperation failed: rX   rY   rZ   z&No corpus name in operation response: z Operation not done yet, attempt    /zOperation timed out after z	 attempts)asyncior   r   r   r   r`   ranger1   rb   rc   r
   rd   re   rR   r^   sleep)r    r[   r\   max_retriesretry_delayrr   rg   rh   rj   attemptrX   rk   operation_datard   rm   s                  r"   rf   z$VertexAIRAGIngestion._poll_operation   s    , 	&t}}5 
)N#34'-11t$

 [)G#ZZ#w|n%= (  H ##s*8H	$$Y/	**%]]_N!!&)n,*73E#&8$@AA -00R@DDVRP&&#@@PQ    27Q;-qN --,,,C *F 4[MKLLE@ -s%   A(E(*E$+C"E(E&E(&E(c                  K   | j                  | j                  | j                  d      \  }}t        | j                        }| d| d}dd|ii}	| j
                  j                  d      }
|
r|
|	d   d<   | j                  }|rgt        |t              rW|j                  d	      }|j                  d
      }|s|r1d|	vri |	d<   ddi ii|	d   d<   |	d   d   d   d   }|r||d	<   |r||d
<   t        j                  d|        t        j                  dt        j                  |	d              dt        j                  |	      df|||xs dfd}t        t        j                   ddi      }|j#                  ||d| dd       d{   }|j$                  dvr/d|j&                   }t        j(                  |       t+        |      	 |j                         }|j                  d i       j                  d!d"      }|s|j                  d!d"      }t        j                  d#|        |S 7 # t*        $ r"}t        j,                  d$|        Y d}~y%d}~ww xY ww)&aA  
        Upload a file to Vertex AI RAG corpus using multipart upload.

        Args:
            rag_corpus_id: RAG corpus resource name
            filename: Name of the file
            file_content: File content bytes
            content_type: MIME type

        Returns:
            File ID or resource name
        r9   r:   z/upload/v1beta1/z/ragFiles:uploadrag_filer+   file_descriptionr*   
chunk_sizechunk_overlapupload_rag_file_configrag_file_chunking_configfixed_length_chunkingrag_file_transformation_configzUploading file to RAG corpus: z
Metadata: rE   rF   NrN   zapplication/octet-stream)metadatafilerH   g     r@rJ   rM   	multipart)rP   zX-Goog-Upload-Protocol)filesrS   rT   zFailed to upload file: ragFilerY   rZ   zUpload complete. File ID: z!Could not parse upload response: uploaded)r]   r   r   r   r   r   r1   chunking_strategy
isinstancer   r
   r^   rR   r_   r   r   r`   ra   rb   rc   rd   re   warning)r    r,   r-   r.   r/   r\   _rg   rh   r   r*   r   r|   r}   chunking_configr   rj   rX   rk   rl   file_ides                         r"   r4   z+VertexAIRAGIngestion._upload_file_to_corpus"  s    ( 33// + 4 
a 't}}5
*}o=MN $
 ..223EF2=HZ / !22,=t!D*..|<J-11/BM]+8;9;H56 /1H"0MX123ST #++C"D4#,#..E#G 4>OL1 7DOO4=cUCDz$**Xa*H)IJK tzz(35GH: :
 (-11u%

  #*<.!9*5 % 
 
 z11(--AI  +I&&	$MMOM#''	26::62FG'++FB7  #=gY!GHN/
0  	""%Fqc#JK	s>   F
IH/AIAH1 .I1	I:IIIIc                  K   | j                  | j                  | j                  d      \  }}t        | j                        }| d| d}ddd|iii}| j
                  }|rIt        |t              r9|j                  d      }	|j                  d	      }
|	s|
r|	xs d
|
xs dd|d   d<   | j                  j                  d      }|r||d   d<   t        j                  d|        t        j                  dt        j                  |d              t        t        j                   ddi      }|j#                  ||d| dd       d{   }|j$                  dvr/d|j&                   }t        j(                  |       t+        |      |j                         }|j                  dd      }t        j                  d |        |S 7 }w)!a  
        Import files from Google Cloud Storage into RAG corpus.

        Args:
            rag_corpus_id: RAG corpus resource name
            gcs_uris: List of GCS URIs (e.g., ["gs://bucket/file.pdf"])

        Returns:
            Operation name for tracking import progress
        r9   r:   ro   z/ragFiles:importimportRagFilesConfig	gcsSourceurisr|   r}   i   rU   )	chunkSizechunkOverlapragFileChunkingConfigmax_embedding_requests_per_minmaxEmbeddingRequestsPerMinzImporting files from GCS: rD   rE   rF   rH   rI   rJ   rM   rN   rO   rQ   NrT   zFailed to import files: rY   rZ   zImport operation started: )r]   r   r   r   r   r   r   r   r1   r   r
   r^   rR   r_   r   r   r`   ra   rb   rc   rd   re   )r    r,   gcs_urisr\   r   rg   rh   ri   r   r|   r}   max_embedding_qpmrj   rX   rk   rl   r[   s                    r"   _import_files_from_gcsz+VertexAIRAGIngestion._import_files_from_gcs  s
      33// + 4 
a 't}}5
)6FG #[682D$E(

 !22,=t!D*..|<J-11/BM]!+!3t$1$8SQ345LM !4488,
  " /0, 	9#?@~djja.P-QRS'-11t$

  #*<.!9 2 % 
 
 z128==/BI  +I&& &**6269.9IJK%
s   EGGA>G)N)r   z'RAGIngestOptions'r   zOptional['Router'])r&   	List[str]returnOptional[List[List[float]]])r.   zOptional[bytes]r-   Optional[str]r/   r   r&   r   r5   r   r   z#Tuple[Optional[str], Optional[str]])r+   strr*   r   r   r   )   g       @)
r[   r   r\   r   ru   intrv   floatr   r   )
r,   r   r-   r   r.   bytesr/   r   r   r   )r,   r   r   r   r   r   )__name__
__module____qualname____doc__r   r'   r7   r2   rf   r4   r   r%   r#   r"   r   r      sS   	 &*V*V #V  
%1-%1-  1- $	1-
 1- 01- 
-1-l &*`` #` 
	`L  EMEM EM 	EM
 EM 
EMNll l 	l
 $l 
l\NN N 
	Nr#   r   )r   
__future__r   rR   typingr   r   r   r   r   r	   litellm._loggingr
   &litellm.llms.custom_httpx.http_handlerr   r   #litellm.llms.vertex_ai.common_utilsr   &litellm.llms.vertex_ai.vertex_llm_baser   $litellm.rag.ingestion.base_ingestionr   litellmr   litellm.types.ragr   r   r%   r#   r"   <module>r      sF    #  B B + D = A2A+Z Ar#   