
    i~                        U 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m	Z	 d dl
Z
d dl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 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!m"Z"m#Z#m$Z$m%Z%m&Z&m'Z' d dl(m)Z) d dl*m+Z+m,Z,  G d d      Z-d Z. G d de-      Z/ G d de-      Z0 G d de-      Z1d dl2m3Z3 d d
lmZ4 g dZ5 G d d      Z6e$jn                  e$jp                  z  Z9e:e;d<    e:h d      Z<e:e;d<    G d  d!      Z=y)"    N)datetime)AnyDictListOptional)&LITELLM_MAX_STREAMING_DURATION_SECONDSSTREAM_SSE_DONE_STRING)run_async_function)process_response_headers)Logging)get_api_base)update_response_metadata)executor)BaseResponsesAPIConfig)ResponsesAPIRequestUtils)OutputTextDeltaEventResponseAPIUsageResponseCompletedEventResponsesAPIRequestParamsResponsesAPIResponseResponsesAPIStreamEventsResponsesAPIStreamingResponse)	CallTypes)CustomStreamWrapper'async_post_call_success_deployment_hookc                       e Zd ZdZ	 	 	 	 ddej
                  dededede	e
eef      de	e   d	e	e
eef      d
e	e   fdZddZde	e   fdZd Zd Zd Zd ZdefdZdefdZy)!BaseResponsesAPIStreamingIteratorz
    Base class for streaming iterators that process responses from the Responses API.

    This class contains shared logic for both synchronous and asynchronous iterators.
    Nresponsemodelresponses_api_provider_configlogging_objlitellm_metadatacustom_llm_providerrequest_data	call_typec	                 L   || _         || _        || _        d| _        || _        d | _        t        |dt        j                               | _	        d| _
        t        j                         | _        || _        || _        |xs i | _        || _        t#        |xs d| j                  j$                  j'                  di             }	|r|j'                  di       ni }
|
j'                  dd       |	|d| _        t+        | j                   j,                  xs i       | j(                  d	<   y )
NF
start_time litellm_params)r   optional_params
model_infoid)model_idapi_baser#   additional_headers)r   r   r!   finishedr    completed_responsegetattrr   nowr'   _failure_handledtime_stream_created_timer"   r#   r$   r%   r   model_call_detailsget_hidden_paramsr   headers)selfr   r   r    r!   r"   r#   r$   r%   	_api_base_model_infos              u/Users/manta/Documents/Projects/TheRoad-I1/.venv/lib/python3.12/site-packages/litellm/responses/streaming_iterator.py__init__z*BaseResponsesAPIStreamingIterator.__init__-   s    !
&-J*KO!+|X\\^L %+/99;! !1#6 ,8,>B(1 !+2 ,,??CC "
	 7G  r2B 	 $d3!#6

 5MMM!!'R5
01    returnc                     t         yt        j                         | j                  z
  }|t         kD  r@t        j                  dt          d|dd| j
                  xs d| j                  xs d      y)zXRaise litellm.Timeout if the stream has exceeded LITELLM_MAX_STREAMING_DURATION_SECONDS.Nz*Stream exceeded max streaming duration of zs (elapsed z.1fzs)r(   )messager   llm_provider)r   r5   r6   litellmTimeoutr   r#   )r;   elapseds     r>   _check_max_streaming_durationz?BaseResponsesAPIStreamingIterator._check_max_streaming_duration\   s    19))+ 9 99;;//DEkDllwx  AD  xE  EG  Hjj&B!55;  <r@   c                    |syt        j                  |      }|y|t        k(  rd| _        y	 t	        j
                  |      }t        |t              r@| j                  j                  | j                  || j                        }t        |dd      }|r9t        j                  || j                  | j                         }t#        |d|       | j                  r| j                  j%                  d      rt        |dd      }|t&        j(                  t&        j*                  fv rt        |dd      }|r}t        |d	d      }|rnt        |t,              r^| j                  r+| j                  j%                  d
i       j%                  d      nd}	|	r#t        j.                  ||	      }
t#        |d	|
       t        |dd      }|r|t&        j0                  t&        j2                  t&        j4                  fv r|| _        t8        j:                  rV| j                  Jt        |dd      }|r;t        |dd      }|,	 | j                  j=                  |      }|t#        |d|       |t&        j4                  k(  r| jA                          |S | jC                          |S y# t>        $ r Y Cw xY w# t        jD                  $ r Y yt>        $ r}| jG                  |        d}~ww xY w)z.Process a single chunk of data from the streamNT)r   parsed_chunkr!   r   )responses_api_responser"   r#   "encrypted_content_affinity_enabledtypeitemencrypted_contentr+   r,   usageresultcost)$r   _strip_sse_data_from_chunkr	   r0   jsonloads
isinstancedictr    transform_streaming_responser   r!   r2   r   /_update_responses_api_response_id_with_model_idr"   r#   setattrr8   r   OUTPUT_ITEM_ADDEDOUTPUT_ITEM_DONEstr%_wrap_encrypted_content_with_model_idRESPONSE_COMPLETEDRESPONSE_INCOMPLETERESPONSE_FAILEDr1   rE   include_cost_in_streaming_usage_response_cost_calculator	Exception_handle_logging_failed_response"_handle_logging_completed_responseJSONDecodeError_handle_failure)r;   chunkrJ   openai_responses_api_chunkresponse_objectr   
event_typerN   rO   r-   wrapped_content_chunk_typeresponse_obj	usage_objrS   es                   r>   _process_chunkz0BaseResponsesAPIStreamingIterator._process_chunkh   s    $>>uE= ** DMa	::e,L ,-66SS"jj%1$($4$4 T  + #**DjRV"W"7gg/>)-)>)>,0,D,D H
 6
HM ((T-B-B-F-F8. "))CVT!RJ!0BB0AA&   ''A64P07>QSW0X-0Z@QSV5W
 (,'<'< %)$9$9$=$=lB$O$S$S(,%& *. !) $,6N6t6t(987&O %,D2E$W &&@&$O-+,??,@@,<<B 3
 /ID+  ?? ,,8GN6
DH (DK ,gtEI  )4	!) )-(8(8(R(R/; )S )& %)
 (,'7(/	64(H #&>&N&NN<<> 21 ??A11 (1 !)$(!) ## 	 	   #		sH   HJ$ 2+J $J$ J$ 	J!J$  J!!J$ $K9KKKc                      y)z8Base implementation - should be overridden by subclassesN r;   s    r>   rg   zDBaseResponsesAPIStreamingIterator._handle_logging_completed_response   s    r@   c                 V   | j                   rt        | j                   dd      nd}|rt        |dd      nd}d}t        |t              r|j	                  dt        |            }t        j                  d|| j                  xs d| j                  xs d      }| j                  |       y)	an  
        Handle logging for RESPONSE_FAILED events by routing to failure handlers.

        Unlike _handle_logging_completed_response (which calls success handlers),
        this constructs an exception from the response error and routes to
        async_failure_handler / failure_handler so logging integrations correctly
        record the call as failed.
        r   NerrorzResponse failedrC   i  r(   )status_coderC   rD   r   )r1   r2   rW   rX   r8   r^   rE   APIErrorr#   r   ri   )r;   rp   
error_infoerror_message	exceptions        r>   rf   zABaseResponsesAPIStreamingIterator._handle_logging_failed_response   s     && D++Z> 	
 >JW\7D9t
)j$'&NN9c*oFM$$!117R**"	
	 	Y'r@   c                   K   	 d}| j                   	 t        | j                         }|!	 t        t        | j                  dd            }| j                  xs t        | j                  di       }t        t        dd      xs g }d}|D ]2  }t        |d      sd}|j                  |||       d{   }|1|}4 |rt        |d	d       |S # t        $ r d}Y w xY w# t
        $ r d}Y w xY w7 ># t
        $ r |cY S w xY ww)
za
        Allow callbacks to modify streaming chunks before returning (parity with chat).
        Nr%   r7   	callbacksF)async_post_call_streaming_deployment_hookT)r$   response_chunkr%   _post_streaming_hooks_ran)r%   r   
ValueErrorr2   r!   re   r$   rE   hasattrr   r[   )r;   rj   typed_call_typer$   r   	hooks_rancallbackrR   s           r>   $_call_post_streaming_deployment_hookzFBaseResponsesAPIStreamingIterator._call_post_streaming_deployment_hook   s>    #	37O~~)+&/&?O &+&/ 0 0+tD'O  ,,   "61L  d;ArII%8%PQ $I#+#U#U%1',"1 $V $ F
 ) & & :DAL7 " +&*O+ ! +&*O+  	L	s   DC3 C C3  C  AC3 C3 2C13C3 :C3 DCC3 CC3  C.+C3 -C..C3 3D>D DDc                 @   K   | j                  |       d{   S 7 w)zY
        Helper to invoke streaming deployment hooks explicitly (used in tests).
        N)r   )r;   rj   s     r>   %call_post_streaming_hooks_for_testingzGBaseResponsesAPIStreamingIterator.call_post_streaming_hooks_for_testing!  s      >>uEEEEs   end_timec                 j   | j                   yi }t        | j                  t              r|j	                  | j                         	 t        | j                  d      r%|j	                  | j                  j                         d|vr+	 t        | j                  di       j                  di       |d<   	 t        | j                   | j                  | j                  || j                  |       	 d}| j                  	 t        | j                        }|	 t        j"                  }	 t%        t&        || j                   |       y# t        $ r Y w xY w# t        $ r i |d<   Y w xY w# t        $ r Y w xY w# t         $ r d}Y qw xY w# t        $ r d}Y w xY w# t        $ r d}Y w xY w# t        $ r Y yw xY w)z^
        Run post-call deployment hooks and update metadata similar to chat pipeline.
        Nr7   r)   )rR   r!   r   kwargsr'   r   )async_functionr$   r   r%   )r1   rW   r$   rX   updater   r!   r7   re   r2   r8   r   r   r'   r%   r   r   	responsesr
   r   )r;   r   request_payloadr   s       r>   _run_post_success_hooksz9BaseResponsesAPIStreamingIterator._run_post_success_hooks'  s    ""**,d''.""4#4#45	t'')=>&&t'7'7'J'JK ?274;$$&:B5#&+   01	$.. ,,jj&??!	#37O~~)+&/&?O
 "'"+"5"5		F,00)	Q  		  746 017  		 " +&*O+ 	#"O	#
  '"&'  		s   ;E *E 19E$ +F :E3 F #F& 	EEE! E!$	E0/E03F>F  FF FFF#"F#&	F21F2r}   c                    | j                   ryd| _         t        j                         }	 t        | j                  j
                  ||| j                  t        j                                	 t        j                  | j                  j                  ||| j                  t        j                                y# t        $ r Y Vw xY w# t        $ r Y yw xY w)z
        Trigger failure handlers before bubbling the exception.
        Only calls handlers once even if called multiple times.
        NT)r   r}   traceback_exceptionr'   r   )r4   	traceback
format_excr
   r!   async_failure_handlerr'   r   r3   re   r   submitfailure_handler)r;   r}   r   s      r>   ri   z1BaseResponsesAPIStreamingIterator._handle_failuree  s        $'224		#//EE#$7??!		OO  00#	  		  		s%   A B5 +A	C 5	C C	CCNNNNrA   N)__name__
__module____qualname____doc__httpxResponser^   r   LiteLLMLoggingObjr   r   r   r?   rH   r   rs   rg   rf   r   r   r   r   re   ri   ru   r@   r>   r   r   &   s     6:-115#'-
..-
 -
 (>	-

 '-
 #4S>2-
 &c]-
 tCH~.-
 C=-
^
px0M'N pd(4'RF< <| r@   r   c                 P   K   t        | dd      }||S  ||       d{   S 7 w)zg
    Module-level helper for tests to ensure hooks can be invoked even if the iterator is wrapped.
    r   N)r2   )iteratorrj   hook_fns      r>   r   r     s1      h FMGs   &$&c                        e Zd ZdZ	 	 	 	 ddej
                  dededede	e
eef      de	e   de	e
eef      d	e	e   f fd
Zd ZdefdZd Z xZS )ResponsesAPIStreamingIteratorzS
    Async iterator for processing streaming responses from the Responses API.
    r   r   r    r!   r"   r#   r$   r%   c	           
      \    t         	|   ||||||||       |j                         | _        y N)superr?   aiter_linesstream_iterator
r;   r   r   r    r!   r"   r#   r$   r%   	__class__s
            r>   r?   z&ResponsesAPIStreamingIterator.__init__  s=     	)		
  (335r@   c                     | S r   ru   rv   s    r>   	__aiter__z'ResponsesAPIStreamingIterator.__aiter__      r@   rA   c                   K   	 | j                          	 	 | j                  j                          d {   }| j                          | j                  |      }| j                  rt        || j                  |       d {   }|S u7 V# t        $ r d| _        t        w xY w7 ## t        $ r  t        j                  $ r}d| _        | j                  |       |d }~wt        $ r}d| _        | j                  |       |d }~ww xY ww)NT)rj   )rH   r   	__anext__StopAsyncIterationr0   rs   r   r   	HTTPErrorri   re   r;   rj   rR   rr   s       r>   r   z'ResponsesAPIStreamingIterator.__anext__  s    #	..0-"&"6"6"@"@"BBE
 224,,U3==,,' $(#L#L$ $M $ F "M'  C) -$(DM,,- " 	 	 DM  #G 	 DM  #G	sm   DB( B B
B A
B( B&B( D	B( 
B B##B( (D	CD	*DD		Dc                 p   | j                   }| j                   St        | j                   d      r=	 t        | j                         j                  | j                   j	                               }t        j                  | j                  j                  || j                  t        j                         d             t        j                  | j                  j                  |d| j                  t        j                                | j!                  t        j                                y# t
        $ r Y w xY w)z7Handle logging for completed responses in async contextN
model_dump)rR   r'   r   	cache_hitrR   r   r'   r   r   )r1   r   rM   model_validater   re   asynciocreate_taskr!   async_success_handlerr'   r   r3   r   r   success_handlerr   r;   logging_responses     r>   rg   z@ResponsesAPIStreamingIterator._handle_logging_completed_response  s      22"".7##\4
#'(?(?#@#O#O++668$  	22'??!	 3 	
 	,,#\\^	
 	$$hlln$=)  s   <D) )	D54D5r   )r   r   r   r   r   r   r^   r   r   r   r   r   r?   r   r   r   rg   __classcell__r   s   @r>   r   r     s     6:-115#'6..6 6 (>	6
 '6 #4S>26 &c]6 tCH~.6 C=6.$!> $L#>r@   r   c                        e Zd ZdZ	 	 	 	 ddej
                  dededede	e
eef      de	e   de	e
eef      d	e	e   f fd
Zd Zd Zd Z xZS )!SyncResponsesAPIStreamingIteratorzY
    Synchronous iterator for processing streaming responses from the Responses API.
    r   r   r    r!   r"   r#   r$   r%   c	           
      \    t         	|   ||||||||       |j                         | _        y r   )r   r?   
iter_linesr   r   s
            r>   r?   z*SyncResponsesAPIStreamingIterator.__init__  s=     	)		
  (224r@   c                     | S r   ru   rv   s    r>   __iter__z*SyncResponsesAPIStreamingIterator.__iter__  r   r@   c                    	 | j                          	 	 t        | j                        }| j                          | j                  |      }| j                  rt        |t        | j                  |      }|S e# t        $ r d| _        t        w xY w# t        $ r  t        j                  $ r}d| _        | j                  |       |d }~wt        $ r}d| _        | j                  |       |d }~ww xY w)NT)r   rj   )rH   nextr   StopIterationr0   rs   r
   r   r   r   ri   re   r   s       r>   __next__z*SyncResponsesAPIStreamingIterator.__next__  s    #	..0( !5!56E
 224,,U3=='''/'+'P'P$F "M'  % ($(DM''($  	 	 DM  #G 	 DM  #G	s@   B A8 AB 7B 8BB C3.CC3C..C3c                 T   | j                   }| j                   St        | j                   d      r=	 t        | j                         j                  | j                   j	                               }t        | j                  j                  || j                  t        j                         d       t        j                  | j                  j                  |d| j                  t        j                                | j                  t        j                                y# t
        $ r Y w xY w)z6Handle logging for completed responses in sync contextNr   )r   rR   r'   r   r   r   r   )r1   r   rM   r   r   re   r
   r!   r   r'   r   r3   r   r   r   r   r   s     r>   rg   zDSyncResponsesAPIStreamingIterator._handle_logging_completed_responseA  s      22"".7##\4
#'(?(?#@#O#O++668$  	++AA#\\^	
 	,,#\\^	
 	$$hlln$='  s   <D 	D'&D'r   )r   r   r   r   r   r   r^   r   r   r   r   r   r?   r   r   rg   r   r   s   @r>   r   r     s     6:-115#'5..5 5 (>	5
 '5 #4S>25 &c]5 tCH~.5 C=5.$L">r@   r   c                        e Zd ZdZdZ	 	 	 	 ddej                  dedede	de
eeef      de
e   d	e
eeef      d
e
e   f fdZd ZdefdZd ZdefdZdedefdZ xZS )!MockResponsesAPIStreamingIteratoru   
    Mock iterator—fake a stream by slicing the full response text into
    5 char deltas, then emit a completed event.

    Models like o1-pro don't support streaming, so we fake it.
       r   r   r    r!   r"   r#   r$   r%   c	           
      p   t         |   ||||||||       | j                  j                  | j                  ||      }	| j                  |	      }
t        dt        |
      | j                        D cg c]:  }t        t        j                  |
||| j                  z    |	j                  dd      < }}t        j                  r3|1t        |	dd       }|"	 |j!                  |	      }|t#        |d|       |t'        t        j(                  |	      gz   | _        d| _        y c c}w # t$        $ r Y <w xY w)	N)r   r   r    r!   r"   r#   r$   r%   )r   raw_responser!   r   )rM   deltaitem_idoutput_indexcontent_indexrP   rQ   rS   )rM   r   )r   r?   r    transform_response_api_responser   _collect_textrangelen
CHUNK_SIZEr   r   OUTPUT_TEXT_DELTAr,   rE   rc   r2   rd   r[   re   r   r`   _events_idx)r;   r   r   r    r!   r"   r#   r$   r%   transformed	full_textideltasrq   rS   r   s                  r>   r?   z*MockResponsesAPIStreamingIterator.__init__p  su    	*G#- 3% 	 		
 ..NNjj%' O  	 &&{3	 1c)ndoo>	
 ? !-??A$78# ? 	 	
 22{7N4;KRV4WI$,7,Q,Q* -R -D '	648 "-@@$!
 
 	A	
* ! s   4?D$!D) )	D54D5c                     | S r   ru   rv   s    r>   r   z+MockResponsesAPIStreamingIterator.__aiter__  r   r@   rA   c                    K   | j                   t        | j                        k\  rt        | j                  | j                      }| xj                   dz  c_         |S wN   )r   r   r   r   r;   evts     r>   r   z+MockResponsesAPIStreamingIterator.__anext__  sE     99DLL))$$ll499%		Q	
s   AAc                     | S r   ru   rv   s    r>   r   z*MockResponsesAPIStreamingIterator.__iter__  r   r@   c                     | j                   t        | j                        k\  rt        | j                  | j                      }| xj                   dz  c_         |S r   )r   r   r   r   r   s     r>   r   z*MockResponsesAPIStreamingIterator.__next__  sA    99DLL))ll499%		Q	
r@   respc                     d}|j                   D ]6  }t        |dd       }|dk(  st        |dg       D ]  }||j                  z  } 8 |S )Nr(   rM   rC   content)outputr2   text)r;   r   outout_item	item_typecs         r>   r   z/MockResponsesAPIStreamingIterator._collect_text  sR    H&$7II% 9b9A166MC : $
 
r@   r   )r   r   r   r   r   r   r   r^   r   r   r   r   r   r?   r   r   r   r   r   r   r   r   r   s   @r>   r   r   f  s     J 6:-115#'A..A A (>	A
 'A #4S>2A &c]A tCH~.A C=AF!> 7 "6 3 r@   r   )verbose_logger)zresponse.createdresponse.completedzresponse.failedzresponse.incompleterx   c                       e Zd ZdZ	 	 ddedededee   dee   f
dZd	e	d
e
fdZded
dfdZded
dfdZded
dfdZddZddZddZddZy)ResponsesWebSocketStreaminga  
    Manages bidirectional WebSocket forwarding for the Responses API
    WebSocket mode (wss://.../v1/responses).

    Unlike the Realtime API, the Responses API WebSocket mode:
    - Uses response.create as the client-to-server event
    - Streams back the same events as the HTTP streaming Responses API
    - Supports previous_response_id for incremental continuation
    - Supports generate: false for warmup
    - One response at a time per connection (sequential, no multiplexing)
    N	websocket
backend_wsr!   user_api_key_dictr$   c                 n    || _         || _        || _        || _        |xs i | _        g | _        g | _        y r   )r   r   r!   r   r$   messagesinput_messages)r;   r   r   r!   r   r$   s         r>   r?   z$ResponsesWebSocketStreaming.__init__  s>     #$&!2"."4"$&46r@   	event_objrA   c                 0    |j                  d      t        v S )NrM   )r8   RESPONSES_WS_LOGGED_EVENT_TYPES)r;   r   s     r>   _should_store_eventz/ResponsesWebSocketStreaming._should_store_event  s    }}V$(GGGr@   eventc                 0   t        |t              r|j                  d      }t        |t              r	 t	        j
                  |      }n|}| j                  |      r| j                  j                  |       y y # t        j                  t        f$ r Y y w xY w)Nutf-8)rW   bytesdecoder^   rU   rV   rh   	TypeErrorr  r   append)r;   r  r   s      r>   _store_eventz(ResponsesWebSocketStreaming._store_event  s    eU#LL)EeS! JJu-	 I##I.MM  + / (()4 s   A9 9BBrC   c                 x   	 t        |t              rt        j                  |      }nt        |t              r|}ny|j                  d      dk7  ry|j                  dg       }t        |t              r| j                  j                  d|d       yt        |t              r|D ]  }t        |t              s|j                  d      dk(  s)|j                  d      dk(  s>|j                  d	g       }t        |t              r| j                  j                  d|d       t        |t              s|D ][  }t        |t              s|j                  d      d
k(  s)|j                  dd      }|s>| j                  j                  d|d       ]  yy# t        j                  t        t        f$ r Y yw xY w)z<Extract user input content from response.create for logging.NrM   response.createinputuser)roler   rC   r  r   
input_textr   r(   )rW   r^   rU   rV   rX   r8   r   r
  listrh   AttributeErrorr	  )r;   rC   msg_objinput_itemsrN   r   r   r   s           r>    _collect_input_from_client_eventz<ResponsesWebSocketStreaming._collect_input_from_client_event  sy   &	'3'**W-GT*!{{6"&77!++gr2K+s+##**F{+ST+t,'D%dD1 xx'94&9IV9S"&((9b"9%gs3 //66)/G D (6%,$.q$$7()f(E+,55+<D'+(,(;(;(B(B5;,M)* &- ( -* $$ni@ 		sH   9F F A F :F F "AF 4F 
F F 4"F F98F9c                 z    | j                  |       | j                  r| j                  j                  |d       y y )Nr(   )r  api_key)r  r!   pre_call)r;   rC   s     r>   _store_inputz(ResponsesWebSocketStreaming._store_input9  s7    --g6%%GR%@ r@   c                 v  K   | j                   sy | j                  r#| j                  | j                   j                  d<   | j                  rmt	        j
                  | j                   j                  | j                               t        j                  | j                   j                  | j                         y y w)Nr   )
r!   r   r7   r   r   r   r   _ws_executorr   r   rv   s    r>   _log_messagesz)ResponsesWebSocketStreaming._log_messages>  s     >B>Q>QD//
;== 0 0 F Ft}} UV 0 0 @ @$--P s   B7B9c                   K   ddl }	 	 	 | j                  j                  d       d{   }t	        |t
              r|j                  d      }n|}| j                  |       | j                  j                  |       d{    ~7 ]# t        $ r& | j                  j                          d{  7  }Y w xY w7 9# |j                  j                  $ r }t        j                  d|       Y d}~n/d}~wt        $ r }t        j                  d|       Y d}~nd}~ww xY w| j!                          d{  7   y# | j!                          d{  7   w xY ww)z4Forward events from backend WebSocket to the client.r   NF)r  r  z*Responses WS backend connection closed: %sz+Error in responses WS backend_to_client: %s)
websocketsr   recvr	  rW   r  r  r  r   	send_text
exceptionsConnectionClosedr   debugre   r}   r  )r;   r  r   response_strrr   s        r>   backend_to_clientz-ResponsesWebSocketStreaming.backend_to_clientG  s%    	'@)-)=)=U)=)K#KL lE2#/#6#6w#?L#/L!!,/nn..|<<< #K  @)-)=)=)?#?#?L@ =$$55 	R  !MqQQ 	W$$%RTUVV	W $$&&&$$$&&&s   EB= B	 BB	 AB= B;B= B	 	&B8/B20B85B= 7B88B= =DC1,D: 1D=DD: DD:  E3D64E:EEEEc                 >  K   	 	 | j                   j                          d{   }| j                  |       | j                  |       | j                  j                  |       d{    h7 J7 # t        $ r }t        j                  d|       Y d}~yd}~ww xY ww)z6Forward response.create events from client to backend.Nz(Responses WS client_to_backend ended: %s)	r   receive_textr  r  r   sendre   r   r$  )r;   rC   rr   s      r>   client_to_backendz-ResponsesWebSocketStreaming.client_to_backenda  s     		P $ ; ; ==!!'*!!'*oo**7333 = 4 	P  !KQOO	PsE   BA1 A-AA1 'A/(A1 /A1 1	B:BBBBc                   K   t        j                  | j                               }	 | j                          d{    |j                         s|j                          	 | d{    	 | j                  j                          d{    y7 S# t        $ r Y \w xY w7 9# t         j                  $ r Y Lw xY w7 1# t        $ r Y yw xY w# |j                         s6|j                          	 | d{  7   n# t         j                  $ r Y nw xY w	 | j                  j                          d{  7   w # t        $ r Y w w xY wxY ww)z,Run both forwarding directions concurrently.N)
r   r   r&  r*  re   donecancelCancelledErrorr   close)r;   forward_tasks     r>   bidirectional_forwardz1ResponsesWebSocketStreaming.bidirectional_forwardn  s0    **4+A+A+CD	((***  $$&##%&&&oo++--- + 		 '--  .   $$&##%&&&-- oo++--- s	  $E	B BB  E	 B! %B&B! +B< B:	B< E	B 	BC BC B! !B74E	6B77E	:B< <	CE	CE	!E-C92C53C98E9DEDED70D31D76E7	E EEEE	)NNr   )r   r   r   r   r   r   r   r   r?   rX   boolr  r  r  r  r  r&  r*  r1  ru   r@   r>   r   r     s    
" ,0'+77 7 '	7
 $C=7 tn7 HT Hd H,# ,$ ,( ( (TAC AD A
Q'4Pr@   r   _RESPONSE_CREATE_PARAMS>   
aresponseslitellm_call_idr   litellm_logging_obj_aresponses_websocket_MANAGED_WS_SKIP_KWARGSc                      e Zd ZdZ	 	 	 	 	 	 d.dededddee   deeeef      d	ee   d
ee   dee   dee   deddfdZ	e
dedee   fd       Zd/dededdfdZdedeeeef      fdZdedeeeef      ddfdZe
deeef   dee   fd       Ze
deeef   deeeef      fd       Ze
dedeeeef      fd       Zdedeeeef      fd Ze
d!eeef   deeef   fd"       Zd#eeef   dee   d$eeeef      d%eeeef      ddf
d&Zd#eeef   d'ee   ddfd(Ze
d#eeef   deddfd)       Zded#eeef   deeeef      fd*Zdeeeef      d%eeeef      d$eeeef      ddfd+Zdeddfd,Zd0d-Zy)1 ManagedResponsesWebSocketHandlera  
    Handles Responses API WebSocket mode for providers that do not expose a
    native ``wss://`` responses endpoint.

    Instead of proxying to a provider WebSocket, this handler:
    - Listens for ``response.create`` events from the client
    - Makes HTTP streaming calls via ``litellm.aresponses(stream=True)``
    - Serialises and forwards every streaming event back over the WebSocket
    - Supports ``previous_response_id`` for multi-turn conversations via
      in-memory session tracking (avoids async DB-write timing issues)
    - Supports sequential requests over a single persistent connection

    This makes every provider that LiteLLM can reach over HTTP available on
    the WebSocket transport without any provider-specific changes.
    Nr   r   r!   r   r   r"   r  r.   timeoutr#   r   rA   c
                    || _         || _        || _        || _        |xs i | _        || _        || _        || _        |	| _        |
j                         D ci c]  \  }}|t        vs|| c}}| _        i | _        y c c}}w r   )r   r   r!   r   r"   r  r.   r;  r#   itemsr8  extra_kwargs_session_history)r;   r   r   r!   r   r"   r  r.   r;  r#   r   kvs                r>   r?   z)ManagedResponsesWebSocketHandler.__init__  s     #
&!20@0FB #6  $\\^-
+TQq8O/OAqD^-
 BD-
s   A>(A>rj   c                    	 t        | d      r| j                  d      S t        | d      r+t        j                  | j	                  d      t
              S t        | t              rt        j                  | t
              S t        j                  t        |             S # t        $ r }t        j                  d|       Y d}~yd}~ww xY w)zHSerialize a streaming chunk to a JSON string for WebSocket transmission.model_dump_jsonT)exclude_noner   )defaultz1ManagedResponsesWS: failed to serialize chunk: %sN)r   rC  rU   dumpsr   r^   rW   rX   re   r   r$  )rj   excs     r>   _serialize_chunkz1ManagedResponsesWebSocketHandler._serialize_chunk  s    	u/0,,$,??ul+zz%"2"2"2"EsSS%&zz%55::c%j)) 	  CS 		s(   B  6B  *B  B   	C	)CC	rC   
error_typec                    K   	 | j                   j                  t        j                  d||dd             d {    y 7 # t        $ r Y y w xY ww)Nrx   )rM   rC   )rM   rx   )r   r!  rU   rF  re   )r;   rC   rI  s      r>   _send_errorz,ManagedResponsesWebSocketHandler._send_error  sS     	..**

$
w/WX  
  		s8   A7A AA  AA 	AAAAprevious_response_idc                     t        j                  |      }|j                  d|      }t        | j                  j                  |g             S )z
        Return accumulated message history for *previous_response_id*.

        The key is the *decoded* response ID (the raw provider response ID before
        LiteLLM base64-encodes it into the ``resp_...`` format).
        response_id)r   !_decode_responses_api_response_idr8   r  r?  )r;   rL  decodedraw_ids       r>   _get_history_messagesz6ManagedResponsesWebSocketHandler._get_history_messages  sH     +LL 
 ],@AD))--fb9::r@   rN  r   c                 "    || j                   |<   y)u   
        Store the complete accumulated message history for *response_id*.

        Replaces any prior value — callers are responsible for passing the full
        history (prior turns + current input + new output).
        N)r?  )r;   rN  r   s      r>   _store_historyz/ManagedResponsesWebSocketHandler._store_history  s     .6k*r@   completed_eventc                     | j                  di       }t        |t              r|j                  d      nd}|syt        j                  |      }|j                  d|      S )z
        Pull the raw (decoded) response ID out of a ``response.completed`` event.
        Returns *None* if the event doesn't contain a usable ID.
        r   r,   NrN  )r8   rW   rX   r   rO  )rU  resp_obj
encoded_idrP  s       r>   _extract_response_idz5ManagedResponsesWebSocketHandler._extract_response_id  s[     #&&z26",Xt"<HLL$ 	 *LLZX{{=*55r@   c                 <   | j                  di       }t        |t              sg S g }|j                  dg       xs g D ]  }t        |t              s|j                  d      }|j                  dd      }|dk(  r|j                  d      xs g }|D cg c]7  }t        |t              r%|j                  d      dv r|j                  d	d
      9 }}d
j                  |      }	|	s|j	                  d|d|	dgd       |dk(  s|j	                  |        |S c c}w )z
        Convert the output items in a ``response.completed`` event into
        Responses API message dicts suitable for the next turn's ``input``.
        r   r   rM   r  	assistantrC   r   )output_textr   r   r(   r\  rM   r   rM   r  r   function_call)r8   rW   rX   joinr
  )
rU  rW  r   rN   r   r  content_partsp
text_partsr   s
             r>   _extract_output_messagesz9ManagedResponsesWebSocketHandler._extract_output_messages  s+    #&&z26(D)I)+LL2.4"4DdD)(I88FK0DI% $ 3 9r +*!!T*quuV}@W/W EE&"%*  
 wwz*OO$-$(1>(M'N o-%- 5. !s   <D	input_valc                     t        | t              rddd| dgdgS t        | t              r!| D cg c]  }t        |t              s| c}S g S c c}w )z
        Normalise the ``input`` field of a ``response.create`` event to a list
        of Responses API message dicts.
        rC   r  r  r]  r^  )rW   r^   r  rX   )re  rN   s     r>   _input_to_messagesz3ManagedResponsesWebSocketHandler._input_to_messages0  se     i% &")5y IJ  i&%.IYT*T42HDYII	 Js   AAraw_messagec                    K   	 t        j                  |      }|j	                  d      dk7  ry|S # t         j                  $ r | j                  dd       d{  7   Y yw xY ww)zOParse raw WS text; return the message dict or None (JSON error / ignored type).z%Invalid JSON in response.create eventinvalid_request_errorNrM   r  )rU   rV   rh   rK  r8   )r;   rh  r  s      r>   _parse_messagez/ManagedResponsesWebSocketHandler._parse_messageF  sn     	jj-G ;;v"33 ## 	""79P   		s1   A$0 A$(A!AA!A$ A!!A$r  c                     | j                  d      }t        |t              r|r|n)| j                         D ci c]  \  }}|dk7  s|| c}}}t        D ci c]  }||v r||   |||    c}S c c}}w c c}w )z
        Extract Responses API params from the event, handling both wire formats:
          Nested: {"type": "response.create", "response": {"input": [...], ...}}
          Flat:   {"type": "response.create", "input": [...], "model": "...", ...}
        r   rM   )r8   rW   rX   r=  r3  )r  nestedr@  rA  response_paramsparams         r>   _build_base_call_kwargsz8ManagedResponsesWebSocketHandler._build_base_call_kwargsT  s     Z( &$'F #*==?B?41aa6k!Q$?B 	 1
0'OE,B,N ?5))0
 	
 C
s   A1A1A7call_kwargscurrent_messagesprior_historyc                     |sy|r)||z   |d<   t        j                  dt        |      |       yt        j                  d|       ||d<   y)zHPrepend in-memory turn history, or fall back to DB-based reconstruction.Nr  zMManagedResponsesWS: prepended %d history messages for previous_response_id=%szuManagedResponsesWS: no in-memory history for previous_response_id=%s; falling back to DB-based session reconstructionrL  )r   r$  r   )r;   rq  rL  rr  rs  s        r>   _apply_historyz/ManagedResponsesWebSocketHandler._apply_historyg  s`     $#03C#CK   _M"$   B$ 3GK./r@   event_modelc                 *   | j                   | j                   |d<   | j                  | j                  |d<   | j                  | j                  |d<   | j                  |s| j                  |d<   | j                  rt        | j                        |d<   yy)zBInject connection-level credentials and metadata into call_kwargs.Nr  r.   r;  r#   r"   )r  r.   r;  r#   r"   rX   )r;   rq  rv  s      r>   _inject_credentialsz4ManagedResponsesWebSocketHandler._inject_credentials  s     <<#%)\\K	"==$&*mmK
#<<#%)\\K	" ##/151I1IK-.  .243H3H.IK*+ !r@   c                    | j                  d      xs i j                  d      xs i }t        |t              syt        |j                  d      xs i       }| j                  d      |d<   | j                  d      |d<   ||d<   dD ]  }|| v s| |   | |   ||<    i |d|i}d| vri | d<   || d   d<   | j                  d	i        || d	   d<   y)
zGUpdate proxy_server_request body so spend logs record the full request.r"   proxy_server_requestNbodyr  storer   )toolstool_choiceinstructionsmetadatar)   )r8   rW   rX   
setdefault)rq  r   rz  r{  r@  s        r>   _update_proxy_requestz6ManagedResponsesWebSocketHandler._update_proxy_request  s    !,0B C IrNN" 
   	 .5(,,V4:;#0W#0WWEAKKN$>%a.Q F  F"6EE[0.0K*+BV&'(>?/4@T$%&<=r@   c                   K   d}t        j                  dd|i| d{   }|2 3 d{   }|t        |dd      xs# t        |t              r|j                  d      nd}| j                  |      }|R|dk(  r|	 t        j                  |      }	 | j                  j                  |       d{    7 7 # t        $ r Y 4w xY w7 # t        $ r$}t        j                  d|       |cY d}~c S d}~ww xY w6 |S w)a>  
        Stream ``litellm.aresponses`` and forward every chunk over the WebSocket.

        Captures the ``response.completed`` event type from the chunk object
        directly (before serialization) to avoid a redundant JSON round-trip on
        every chunk.  Returns the completed event dict, or ``None``.
        Nr   rM   r   z5ManagedResponsesWS: error sending chunk to client: %sru   )rE   r4  r2   rW   rX   r8   rH  rU   rV   re   r   r!  r   r$  )	r;   r   rq  rU  stream_responserj   
chunk_type
serializedsend_excs	            r>   _stream_and_forwardz4ManagedResponsesWebSocketHandler._stream_and_forward  s     59 ' 2 2 N N+ NN* 	'%} 5 %/t%<		&!$  ..u5J!11o6M&*jj&<O'nn..z:::# O	' !  ; '$$KX '&	'# +, s   DB:DC?B<C?AD B>C4C5C9D<C?>	C
D	C

DC	C<C7/C<0D7C<<Dc                     |y| j                  |      }|sy| j                  |      }||z   |z   }| j                  ||       t        j                  dt        |      |       y)zMStore this turn in in-memory history for future previous_response_id lookups.Nz9ManagedResponsesWS: stored %d messages for response_id=%s)rY  rd  rT  r   r$  r   )r;   rU  rs  rr  new_response_idoutput_msgsall_messagess          r>   _save_turn_historyz3ManagedResponsesWebSocketHandler._save_turn_history  so     "33OD33OD$'77+EO\:G	
r@   c                   K   | j                  |       d{   }|y| j                  |      }d|d<   |j                  dd      }|xs | j                  }|j                  dd      }| j	                  |j                  d            }|r| j                  |      ng }| j                  ||||       | j                  ||       | j                  ||       |j                  | j                         	 | j                  ||       d{   }	| j%                  |	||       y7 	7 # t        $ rC}
t        j                  d|
       | j!                  t#        |
             d{  7   Y d}
~
yd}
~
ww xY ww)a  
        Parse one ``response.create`` event, call ``litellm.aresponses(stream=True)``,
        and forward every streaming event to the client.

        Multi-turn support via in-memory session history
        ------------------------------------------------
        When ``previous_response_id`` is present in the event:
        1. Look up the accumulated message history in ``self._session_history``
           (keyed by the decoded provider response ID).
        2. Prepend those messages to the current ``input`` so the model has full
           conversation context.
        3. After the stream completes, extract the new response ID and output
           messages from ``response.completed`` and store them in
           ``self._session_history`` for the next turn.

        This in-memory approach avoids the async DB-write race condition that
        occurs when spend logs haven't been committed by the time the second
        ``response.create`` arrives over the same WebSocket connection.
        NTstreamr   rL  r  z8ManagedResponsesWS: error processing response.create: %s)rk  rp  popr   rg  r8   rR  ru  rx  r  r   r>  r  re   r   r}   rK  r^   r  )r;   rh  r  rq  rv  r   rL  rr  rs  rU  rG  s              r>   _process_response_createz9ManagedResponsesWebSocketHandler._process_response_create  su    ( ++K88?227; $H%0__Wd%C)tzz.9oo"D/
  22;??73KL
 $ &&';< 	 	-/?	
 	  k:"";64,,-	$($<$<UK$PPO 	@PQM 9< Q 	$$JC ""3s8,,,	sX   E4D CE42D% D#D% E4#D% %	E1.3E,!E$"E,'E4,E11E4c                 z  K   	 	 	 | j                   j                          d{   }| j                  |       d{    =7 # t        $ r }t        j                  d|       Y d}~yd}~ww xY w7 3# t        $ r=}t        j                  d|       | j                  d|        d{  7   Y d}~yd}~ww xY ww)z
        Main loop: accept ``response.create`` events sequentially and handle
        each one before waiting for the next message.
        Nz+ManagedResponsesWS: client disconnected: %sz(ManagedResponsesWS: unexpected error: %szInternal server error: )r   r(  re   r   r$  r  r}   rK  )r;   rC   rG  s      r>   runz$ManagedResponsesWebSocketHandler.run'  s     
	D$(NN$?$?$AAG 33G<<< A  "((Es 	 = 	D$$%OQTU""%<SE#BCCC	Ds|   B;A2 A AA A2 A0A2 A 	A-A(#A2 'B;(A--A2 2	B8;-B3(B+)B3.B;3B88B;)NNNNNN)server_errorr   )r   r   r   r   r   r^   r   r   floatr?   staticmethodrH  rK  r   rR  rT  rY  rd  rg  rk  rp  ru  rx  r  r  r  r  r  ru   r@   r>   r:  r:    s   * ,059!%"&#'-1DD D )	D
 $C=D #4S>2D #D 3-D %D &c]D D 
DH      # SW ;# ;$tCQTH~BV ;6# 6d38n9M 6RV 6 6d38n 6# 6 6 "c3h"	d38n	" "H c d4S>.B  * c3h8P  
c3h 
DcN 
 
$G#s(^G 'smG tCH~.	G
 DcN+G 
G6JS>J8@J	J$ U4S> U# U$ U U*""'+CH~"	$sCx.	!"H
!$sCx.1
 DcN+
 tCH~.	

 

2:R# :R$ :R@Dr@   r:  )>r   rU   r5   r   r   typingr   r   r   r   r   rE   litellm.constantsr   r	   #litellm.litellm_core_utils.asyncifyr
   'litellm.litellm_core_utils.core_helpersr   *litellm.litellm_core_utils.litellm_loggingr   r   :litellm.litellm_core_utils.llm_response_utils.get_api_baser   ?litellm.litellm_core_utils.llm_response_utils.response_metadatar   /litellm.litellm_core_utils.thread_pool_executorr   .litellm.llms.base_llm.responses.transformationr   litellm.responses.utilsr   litellm.types.llms.openair   r   r   r   r   r   r   litellm.types.utilsr   litellm.utilsr   r   r   r   r   r   r   litellm._loggingr   r  r  r   __required_keys____optional_keys__r3  	frozenset__annotations__r8  r:  ru   r@   r>   <module>r     s        , ,   C L S S E Q <   * V^ ^B h>$E h>Vg>(I g>Th(I h^ , T# ^ ^L //112  
 &/&  dD dDr@   