
    i=                        d dl Z d dlZd dlZd dlZd dlZd dlmZ d dlmZmZm	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 d dlmZ d dlmZ d d	lmZ d d
lmZ d dlmZ d dl d dl m!Z! erd dl"m#Z# neZ# G d dee      Z$y)    N)uuid)datetime	timedeltatimezone)TYPE_CHECKINGAnyDictListOptionalTuple)quote)verbose_logger)LITELLM_ASYNCIO_QUEUE_MAXSIZE)AdditionalLoggingUtils)GCSBucketBase)CommonProxyErrors)IntegrationHealthCheckStatus)*)StandardLoggingPayload)
VertexBasec            	           e Zd Zd'dee   ddf fdZd Zd Zdee	   fdZ
ded	edefd
Zdeeef   defdZdedefdZdee	   deeee	   f   fdZdee	   defdZdee	   dedeeef   fdZdee	   ddfdZde	ddfdZd ZdedededefdZdedee   dee   dee   fdZdededefd Zdedefd!Zd"edefd#Zd$ Z d% Z!de"fd&Z# xZ$S )(GCSBucketLoggerNbucket_namereturnc                    ddl m} t        |   |       t	        t        j                  dt                    | _        t	        t        j                  dt                    | _
        t        j                  dt        t              j                               j                         dk(  | _        t        j                          | _        t        |   | j"                  | j                  | j                         t        j$                  t&        	      | _        t        j*                  | j-                                t/        j                  |        |d
ur&t1        dt2        j4                  j6                         y )Nr   premium_user)r   GCS_BATCH_SIZEGCS_FLUSH_INTERVALGCS_USE_BATCHED_LOGGINGtrue)
flush_lock
batch_sizeflush_interval)maxsizeTCGCS Bucket logging is a premium feature. Please upgrade to use it. )litellm.proxy.proxy_serverr   super__init__intosgetenvGCS_DEFAULT_BATCH_SIZEr#   "GCS_DEFAULT_FLUSH_INTERVAL_SECONDSr$   strGCS_DEFAULT_USE_BATCHED_LOGGINGloweruse_batched_loggingasyncioLockr"   Queuer   	log_queuecreate_taskperiodic_flushr   
ValueErrorr   not_premium_uservalue)selfr   r   	__class__s      {/Users/manta/Documents/Projects/TheRoad-I1/.venv/lib/python3.12/site-packages/litellm/integrations/gcs_bucket/gcs_bucket.pyr)   zGCSBucketLogger.__init__   s/   ;[1bii(8:PQR!II*,NO
 II)3/N+O+U+U+Weg 	  ",,... 	 	

 :A1:
 	D//12''-t#UVgVxVxV~V~U  A  $    c                   K   ddl m} |dur&t        dt        j                  j
                         	 t        j                  d||       |j                  dd       }|t        d      | j                  j                         r| j                          d {    | j                  j                  t        |||             d {    y 7 47 # t        $ r+}t        j                  d	t!        |              Y d }~y d }~ww xY ww)
Nr   r   Tr&   zHGCS Logger: async_log_success_event logging kwargs: %s, response_obj: %sstandard_logging_object+standard_logging_object not found in kwargspayloadkwargsresponse_objGCS Bucket logging error: )r'   r   r9   r   r:   r;   r   debuggetr6   fullflush_queueputGCSLogQueueItem	Exception	exceptionr/   )r<   rE   rF   
start_timeend_timer   logging_payloades           r>   async_log_success_eventz'GCSBucketLogger.async_log_success_event<   s
    ;t#UVgVxVxV~V~U  A 	L  Z
 AG

)4AO & !NOO~~""$&&(((..$$+F   )  	L$$'A#a&%JKK	LsS   1DA#C C.C CC DC C 	D!C?:D?DDc                   K   	 t        j                  d||       |j                  dd       }|t        d      | j                  j                         r| j                          d {    | j                  j                  t        |||             d {    y 7 47 # t        $ r+}t        j                  dt        |              Y d }~y d }~ww xY ww)NzHGCS Logger: async_log_failure_event logging kwargs: %s, response_obj: %srA   rB   rC   rG   )r   rH   rI   r9   r6   rJ   rK   rL   rM   rN   rO   r/   )r<   rE   rF   rP   rQ   rR   rS   s          r>   async_log_failure_eventz'GCSBucketLogger.async_log_failure_eventZ   s     	L  Z AG

)4AO & !NOO~~""$&&(((..$$+F   )  	L$$'A#a&%JKK	LsS   CA#B  'B(.B  BB  CB  B   	C)!C
CCCc                     g }t        |      | j                  k  rC	 |j                  | j                  j	                                t        |      | j                  k  rC|S # t
        j                  $ r Y |S w xY w)a  
        Drain items from the queue (non-blocking), respecting batch_size limit.

        This prevents unbounded queue growth when processing is slower than log accumulation.

        Returns:
            List of items to process, up to batch_size items
        )lenr#   appendr6   
get_nowaitr3   
QueueEmpty)r<   items_to_processs     r>   _drain_queue_batchz"GCSBucketLogger._drain_queue_batchs   sx     35"#doo5 ''(A(A(CD "#doo5
   %% s   )A   A76A7date_strbatch_idc                     | d| dS )zm
        Generate object name for a batched log file.
        Format: {date}/batch-{batch_id}.ndjson
        z/batch-z.ndjson )r<   r^   r_   s      r>   _generate_batch_object_namez+GCSBucketLogger._generate_batch_object_name   s    
 78*G44r?   rE   c                     |j                  dd      xs i }|j                  dd      xs | j                  xs d}|j                  dd      xs | j                  xs d}| d| S )a  
        Extract a synchronous grouping key from kwargs to group items by GCS config.
        This allows us to batch items with the same bucket/credentials together.

        Returns a string key that uniquely identifies the GCS config combination.
        This key may contain sensitive information (bucket names, paths) - use _sanitize_config_key()
        for logging purposes.
         standard_callback_dynamic_paramsNgcs_bucket_namedefaultgcs_path_service_account|)rI   BUCKET_NAMEpath_service_account_json)r<   rE   rd   r   path_service_accounts        r>   _get_config_keyzGCSBucketLogger._get_config_key   s     JJ94@FB 	)
 -001BDI  	 -001KTR -- 	 a 4566r?   
config_keyc                 v    t        j                  |j                  d            }d|j                         dd  S )z
        Create a sanitized version of the config key for logging.
        Uses a hash to avoid exposing sensitive bucket names or service account paths.

        Returns a short hash prefix for safe logging.
        zutf-8zconfig-N   )hashlibsha256encode	hexdigest)r<   rm   hash_objs      r>   _sanitize_config_keyz$GCSBucketLogger._sanitize_config_key   s;     >>*"3"3G"<=++-bq1233r?   itemsc                 z    i }|D ]3  }| j                  |d         }||vrg ||<   ||   j                  |       5 |S )z
        Group items by their GCS config (bucket + credentials).
        This ensures items with different configs are processed separately.

        Returns a dict mapping config_key -> list of items with that config.
        rE   )rl   rY   )r<   rv   groupeditemrm   s        r>   _group_items_by_configz&GCSBucketLogger._group_items_by_config   sS     57D--d8n=J(&(
#J&&t,	 
 r?   c                     g }|D ]4  }|d   }t        j                  |t        d      }|j                  |       6 dj	                  |      S )z
        Combine multiple log payloads into newline-delimited JSON (NDJSON) format.
        Each line is a valid JSON object representing one log entry.
        rD   F)rf   ensure_ascii
)jsondumpsr/   rY   join)r<   rv   linesry   rR   	json_lines         r>   _combine_payloads_to_ndjsonz+GCSBucketLogger._combine_payloads_to_ndjson   sK    
 D"9oO

?CeTILL#  yyr?   c                   K   |sy|d   d   }	 | j                  |       d{   }| j                  |d   |d          d{   }|d   }| j                  t        j                  t
        j                              }t        t        j                         d	z         d
t        j                         j                  dd  }| j                  ||      }	| j                  |      }
| j                  |||	|
       d{    t        |      }d}||fS 7 7 7 # t         $ r<}d}t        |      }t#        j$                  dt'        |              ||fcY d}~S d}~ww xY ww)z
        Send a batch of items that share the same GCS config.

        Returns:
            (success_count, error_count)
        )r   r   r   rE   Nvertex_instancerk   r   service_account_jsonr   i  -ro   headersr   object_namerR   z6GCS Bucket error logging batch payload to GCS bucket: )get_gcs_logging_configconstruct_request_headers_get_object_date_from_datetimer   nowr   utcr*   timer   uuid4hexrb   r   _log_json_data_on_gcsrX   rN   r   rO   r/   )r<   rv   rm   first_kwargsgcs_logging_configr   r   current_dater_   r   combined_payloadsuccess_counterror_countrS   s                 r>   _send_grouped_batchz#GCSBucketLogger._send_grouped_batch   s     Qx)#	09=9T9T: 4 !:: 23D E%78N%O ;  G -];K>>X\\*L diikD012!DJJL4D4DRa4H3IJH::<RK#??F,,'' 0	 -     JMK!;//54  	0Me*K$$HQQ ";//	0si   E#D D D DB:D ?D D E#D D D 	E $1EE E#E  E#c                 P   K   |D ]  }| j                  |       d{     y7 w)z
        Send each log individually as separate GCS objects (legacy behavior).
        This is used when GCS_USE_BATCHED_LOGGING is disabled.
        N)_send_single_log_item)r<   rv   ry   s      r>   _send_individual_logsz%GCSBucketLogger._send_individual_logs   s)     
 D,,T222 2s   &$&ry   c                   K   	 | j                  |d          d{   }| j                  |d   |d          d{   }|d   }| j                  |d   |d   |d   	      }| j                  ||||d   
       d{    y7 h7 I7 	# t        $ r+}t        j                  dt        |              Y d}~yd}~ww xY ww)zH
        Send a single log item to GCS as an individual object.
        rE   Nr   rk   r   r   rD   rF   )rE   rR   rF   r   z;GCS Bucket error logging individual payload to GCS bucket: )r   r   _get_object_namer   rN   r   rO   r/   )r<   ry   r   r   r   r   rS   s          r>   r   z%GCSBucketLogger._send_single_log_item  s     	9=9T9TX: 4 !:: 23D E%78N%O ;  G -];K//H~ $Y!.1 0 K ,,'' $Y	 -   !4  	$$McRSfXV 	sa   CB
 B B
 BAB
 >B?B
 CB
 B
 B
 
	B>!B94C9B>>Cc                   K   | j                         }|sy| j                  rD| j                  |      }|j                         D ]  \  }}| j	                  ||       d{    ! y| j                  |       d{    y7 !7 w)aQ  
        Process queued logs - sends logs to GCS Bucket.

        If `GCS_USE_BATCHED_LOGGING` is enabled (default), batches multiple log payloads
        into single GCS object uploads (NDJSON format), dramatically reducing API calls.

        If disabled, sends each log individually as separate GCS objects (legacy behavior).
        N)r]   r2   rz   rv   r   r   )r<   r\   grouped_itemsrm   group_itemss        r>   async_send_batchz GCSBucketLogger.async_send_batch'  s       224## 778HIM+8+>+>+@'
K..{JGGG ,A ,,-=>>> H>s$   ABB B:B;BBrR   rF   c                 d   | j                  t        j                  t        j                              }|j                  dd      | j                  |      }n#| j                  ||j                  dd            }|j                  dd      xs i }|j                  dd      xs i }d	|v r|d	   }|S )
zD
        Get the object name to use for the current payload
        	error_strN)request_date_strid r   response_idlitellm_paramsmetadata
gcs_log_id)r   r   r   r   r   rI   _generate_failure_object_name_generate_success_object_name)r<   rE   rR   rF   r   r   _litellm_params	_metadatas           r>   r   z GCSBucketLogger._get_object_name=  s     ::8<<;UV{D1=<<!- = K <<!-(,,T26 = K !**%5t<B#''
D9?R	9$#L1Kr?   
request_idstart_time_utcend_time_utcc           
        K   |t        d      ||t        d      z   |t        d      z
  g}d}|D ]i  }	 | j                  |      }| j                  ||      }t	        |d      }| j                  |       d{   }	|	t        j                  |	      }
|
c S k y7 "# t        $ r.}t        j                  d	| d
t        |              Y d}~d}~ww xY ww)z
        Get the request and response payload for a given `request_id`
        Tries current day, next day, and previous day until it finds the payload
        Nz@start_time_utc is required for getting a payload from GCS Bucket   )days)datetime_objr   r   )safez!Failed to fetch payload for date z: )r9   r   r   r   r   download_gcs_objectr~   loadsrN   r   rH   r/   )r<   r   r   r   dates_to_tryr^   dater   encoded_object_nameresponseloaded_responserS   s               r>   get_request_response_payloadz,GCSBucketLogger.get_request_response_payloadV  s     !R 
 YA..YA..

  D>>D>Q"@@%- * A  ',Kb&A#!%!9!9:M!NN'&*jj&:O** ( !&  O
  $$7zCF8L 	sA   4CAB"=B >B"C B""	C+$CCCCr   r   c                     | d| S )N/ra   )r<   r   r   s      r>   r   z-GCSBucketLogger._generate_success_object_name  s    
 ##1[M22r?   c                 H    | dt        j                         j                   S )Nz	/failure-)r   r   r   )r<   r   s     r>   r   z-GCSBucketLogger._generate_failure_object_name  s#     ##9TZZ\-=-=,>??r?   r   c                 $    |j                  d      S )Nz%Y-%m-%d)strftime)r<   r   s     r>   r   z.GCSBucketLogger._get_object_date_from_datetime  s    $$Z00r?   c                 r   K   | j                          d{    t        j                         | _        y7 w)zB
        Override flush_queue to work with asyncio.Queue.
        N)r   r   last_flush_timer<   s    r>   rK   zGCSBucketLogger.flush_queue  s-      ##%%%#yy{ 	&s   757c                    K   	 t        j                  | j                         d{    t        j                  d| j                   d       | j                          d{    c7 @7 w)zE
        Override periodic_flush to work with asyncio.Queue.
        Nz GCS Bucket periodic flush after z seconds)r3   sleepr$   r   rH   rK   r   s    r>   r8   zGCSBucketLogger.periodic_flush  sf      -- 3 3444  243F3F2GxP ""$$$ 4 %s!   $A+A':A+!A)"A+)A+c                     K   t        d      w)Nz(GCS Bucket does not support health check)NotImplementedErrorr   s    r>   async_health_checkz"GCSBucketLogger.async_health_check  s     !"LMMs   )N)%__name__
__module____qualname__r   r/   r)   rT   rV   r
   rM   r]   rb   r	   r   rl   ru   rz   r   r   r*   r   r   r   r   r   r   r   dictr   r   r   r   rK   r8   r   r   __classcell__)r=   s   @r>   r   r      s   HSM T BL<L2 D$9  "5C 53 53 57d38n 7 744s 4s 4/*	c4((	)"
 o1F 
 3 
 10/*108;10	sCx10f3o1F 34 3 D @?,-CSV	2(( !*( x(	(
 
$(T33 3 
	3@@ 
@18 1 1+	%N*F Nr?   r   )%r3   rp   r~   r+   r   litellm._uuidr   r   r   r   typingr   r   r	   r
   r   r   urllib.parser   litellm._loggingr   litellm.constantsr   -litellm.integrations.additional_logging_utilsr   /litellm.integrations.gcs_bucket.gcs_bucket_baser   litellm.proxy._typesr   ,litellm.types.integrations.base_health_checkr   %litellm.types.integrations.gcs_bucketlitellm.types.utilsr   &litellm.llms.vertex_ai.vertex_llm_baser   r   ra   r?   r>   <module>r      s\       	   2 2 B B  + ; P I 2 U 3 6AJINm%; INr?   