
    ^ju                        d Z ddlZddlZddlmZ ddlmZ ddlmZ ddl	m
Z
  e
       rJddlmZ dd	lmZmZ dd
lmZmZmZmZmZmZmZmZmZmZmZmZmZ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)m*Z* ddlm+Z+m,Z,m-Z-m.Z.m/Z/m0Z0m1Z1m2Z2m3Z3 erddl4m5Z5m6Z6m7Z7m8Z8  ejr                  e:      Z; G d de&d      Z<h dZ= G d d      Z> G d de,      Z?de@de@de*fdZAy)z~
Handler for the /v1/responses endpoint (OpenAI Responses API).

Supports streaming (SSE) and non-streaming (JSON) responses.
    N)AsyncGenerator)TYPE_CHECKING   )logging)is_serve_available)HTTPException)JSONResponseStreamingResponse)ResponseResponseCompletedEventResponseContentPartAddedEventResponseContentPartDoneEventResponseCreatedEventResponseErrorResponseErrorEventResponseFailedEvent&ResponseFunctionCallArgumentsDoneEventResponseFunctionToolCallResponseInProgressEventResponseOutputItemAddedEventResponseOutputItemDoneEventResponseOutputMessageResponseOutputTextResponseReasoningItemResponseReasoningTextDeltaEventResponseReasoningTextDoneEventResponseTextDeltaEventResponseTextDoneEvent)ResponseCreateParamsStreaming)InputTokensDetailsOutputTokensDetailsResponseUsage   )	BaseGenerateManagerBaseHandlerModalityReasoningText_StreamErrorget_reasoning_configget_tool_call_configparse_reasoningparse_tool_calls)GenerationConfigPreTrainedModelPreTrainedTokenizerFastProcessorMixinc                   "    e Zd ZU eed<   eed<   y))TransformersResponseCreateParamsStreaminggeneration_configseedN)__name__
__module____qualname__str__annotations__int     l/var/www/ramen.bs-engineer-server.com/venv/lib/python3.12/site-packages/transformers/cli/serving/response.pyr2   r2   N   s    
Ir<   r2   F)total>   textuserstorepromptinclude
background
truncationtool_choiceservice_tiertop_logprobsmax_tool_callsprevious_response_idc            	          e Zd ZdZdedefdZdefdZdeddfd	Zde	e   fd
Z
de	e   fdZdede	e   fdZde	e   fdZde	e   fdZdede	e   fdZde	e   fdZdededede	e   fdZdede	e   fdZde	e   fdZy)_ResponseStreamBuilderz=Builds SSE events for one streaming Responses API generation.
request_idresponse_defaultsc                    || _         d| | _        d| | _        d| | _        d| _        d| _        d| _        d| _        d| _        d| _	        d | _
        d | _        g | _        y )Nresp_msg_rs_r    F)_response_defaultsresp_idmsg_idreasoning_idseqoutput_index	full_textfull_reasoningreasoning_openmessage_openreasoning_itemmessage_item
tool_calls)selfrM   rN   s      r=   __init__z_ResponseStreamBuilder.__init__f   sz    "3zl+ZL)!*. #!<@:>:<r<   returnc                 Z    t        j                  |      }| xj                  dz  c_        |S )Nr#   )r%   chunk_to_sserX   )ra   eventsses      r=   _emitz_ResponseStreamBuilder._emitu   s$    &&u-A
r<   statusr   c                 8    t        di | j                  d|i|S )Nri   r;   )r   rT   )ra   ri   extras      r=   	_responsez _ResponseStreamBuilder._responsez   s     J$11J&JEJJr<   c                     | j                  t        d| j                  | j                  dg                   | j                  t	        d| j                  | j                  dg                   gS )Nzresponse.createdqueued)outputtypesequence_numberresponsezresponse.in_progressin_progress)rh   r   rX   rl   r   ra   s    r=   start_responsez%_ResponseStreamBuilder.start_response}   sj    JJ$+$(HH!^^HR^@ JJ'/$(HH!^^M"^E
 	
r<   c                     d| _         | j                  t        d| j                  | j                  t        | j                  dg g d                  gS )NTresponse.output_item.added	reasoningrt   idrq   summarycontentri   rq   rr   rY   item)r\   rh   r   rX   rY   r   rW   ru   s    r=   start_reasoningz&_ResponseStreamBuilder.start_reasoning   sV    "JJ,5$(HH!%!2!2.,,;TV_l		
 	
r<   r?   c           
          | xj                   |z  c_         | j                  t        d| j                  | j                  | j
                  d|            gS )Nzresponse.reasoning_text.deltar   )rq   item_idrr   rY   content_indexdelta)r[   rh   r   rW   rX   rY   ra   r?   s     r=   reasoning_deltaz&_ResponseStreamBuilder.reasoning_delta   sS    t#JJ/8 --$(HH!%!2!2"#	
 	
r<   c           
         t        | j                  dg d| j                  dgd      | _        | j	                  t        d| j                  | j                  | j                  d| j                              | j	                  t        d	| j                  | j                  | j                  
            g}d| _	        | xj                  dz  c_        |S )Nry   reasoning_textrq   r?   	completedrz   zresponse.reasoning_text.doner   )rq   r   rr   rY   r   r?   response.output_item.doner~   Fr#   )
r   rW   r[   r^   rh   r   rX   rY   r   r\   )ra   partss     r=   finish_reasoningz'_ResponseStreamBuilder.finish_reasoning   s    3  .8K8KLM
 JJ.7 --$(HH!%!2!2"#,,	 JJ+4$(HH!%!2!2,,	
( $Qr<   c                 8   d| _         | j                  t        d| j                  | j                  t        | j                  dddg                   | j                  t        d| j                  | j                  | j                  d	t        d
dg                   gS )NTrx   messagert   	assistant)r{   rq   ri   roler}   r~   zresponse.content_part.addedr   output_textrS   rq   r?   annotationsrq   r   rr   rY   r   part)	r]   rh   r   rX   rY   r   rV   r   r   ru   s    r=   start_messagez$_ResponseStreamBuilder.start_message   s     JJ,5$(HH!%!2!2.;;Y}S^hj		 JJ-6 KK$(HH!%!2!2"#+RUWX	
 	
r<   c                     | xj                   |z  c_         | j                  t        d| j                  | j                  | j
                  d|g             gS )Nzresponse.output_text.deltar   )rq   r   rr   rY   r   r   logprobs)rZ   rh   r   rV   rX   rY   r   s     r=   
text_deltaz!_ResponseStreamBuilder.text_delta   sQ    $JJ&5 KK$(HH!%!2!2"#

 	
r<   c                 
   t        d| j                  g       }t        | j                  ddd|gg       | _        | j                  t        d| j                  | j                  | j                  d| j                  g 	            | j                  t        d
| j                  | j                  | j                  d|            | j                  t        d| j                  | j                  | j                              g}d| _        |S )Nr   r   r   r   r   r{   rq   ri   r   r}   r   zresponse.output_text.doner   )rq   r   rr   rY   r   r?   r   zresponse.content_part.doner   r   r~   F)r   rZ   r   rV   r_   rh   r   rX   rY   r   r   r]   )ra   output_text_partr   s      r=   finish_messagez%_ResponseStreamBuilder.finish_message   s    -=t~~cef1{{%&
 JJ%4 KK$(HH!%!2!2"#
 JJ,5 KK$(HH!%!2!2"#)	 JJ+4$(HH!%!2!2**	-
> "r<   tc_idname	argumentsc                    | xj                   dz  c_         t        ||d||d      }| j                  j                  |       | j	                  t        d| j                  | j                   |            | j	                  t        d| j                  || j                   ||            | j	                  t        d	| j                  | j                   |            gS )
Nr#   function_callr   r{   call_idrq   r   r   ri   rx   r~   z%response.function_call_arguments.done)rq   rr   r   rY   r   r   r   )	rY   r   r`   appendrh   r   rX   r   r   )ra   r   r   r   r   s        r=   	tool_callz _ResponseStreamBuilder.tool_call"  s    Q' 
 	t$JJ,5$(HH!%!2!2	 JJ6@$(HH!!%!2!2'	 JJ+4$(HH!%!2!2	'
 	
r<   msgc                     | j                  t        d| j                  |            | j                  t        d| j                  | j	                  dg t        d|                        gS )	Nerror)rq   rr   r   zresponse.failedfailedserver_error)coder   )ro   r   rp   )rh   r   rX   r   rl   r   )ra   r   s     r=   r   z_ResponseStreamBuilder.errorJ  se    JJ)wZ]^_JJ#*$(HH!^^ =n^a3b , 
 	
r<   c                 4   g }| j                   |j                  | j                          |j                  | j                         |j                  | j                         | j                  t        d| j                  | j                  d||                  gS )Nzresponse.completedr   )ro   usagerp   )	r^   r   r_   extendr`   rh   r   rX   rl   )ra   r   
all_outputs      r=   r   z _ResponseStreamBuilder.completedX  s    
*d112$++,$//*JJ&-$(HH!^^K
RW^X
 	
r<   N)r5   r6   r7   __doc__r8   dictrb   rh   rl   listrv   r   r   r   r   r   r   r   r   r   r;   r<   r=   rL   rL   c   s   G=c =d =c 
K K K
S	 
$
c 

C 
DI 
$s) @
tCy 
2
s 
tCy 
 *S	 *X&
s &
# &
# &
$s) &
P
 
c 

$s) 
r<   rL   c                   X    e Zd ZdZeZeZdede	de
ez  fdZedee   dz  dee   dz  fd       Zededee   fd	       Zed
ee   dee   fd       Z	 	 dde	ddddde	dededddededz  dedz  de
fdZ	 	 dde	ddddde	dededddededz  dedz  defdZddedddef fdZ xZS )ResponseHandlerz+Handler for the ``/v1/responses`` endpoint.bodyrM   rc   c                 |  K   | j                  |       | j                  |      \  }}}| j                  j                  ||      }| j                  j                  ||      }t        j                  d| d|        | j                  j                  ||      }| j                  |      }	| j                  |	|      }
t        d |
D              }i }|rd|d<   |j                  | j                         |j                  |j                  d      xs i        | j                  |j                  d	            } |j                   |
fd
||rdndd
d
|t"        j$                  k(  xr |d|}|s|j'                  |j(                        }| j+                  ||j,                  |      }|r|j/                  ||       |j                  d	      rt1        ||      nd}t3        |||d         }|j                  dd
      }|r| j5                  ||||||||||
      S | j7                  ||||||||||
       d{   S 7 w)az  Validate, load model, dispatch to streaming or non-streaming.

        Args:
            body (`dict`): The raw JSON request body (OpenAI Responses API format).
            request_id (`str`): Unique request identifier (from header or auto-generated).

        Returns:
            `StreamingResponse | JSONResponse`: SSE stream or JSON depending on ``body["stream"]``.
        )	processorz[Request received] Model: z, CB: use_cbc              3      K   | ]O  }t        |j                  d       t              r|j                  d       ng D ]  }|j                  d      dk(    Q yw)r}   rq   videoN)
isinstancegetr   ).0r   cs      r=   	<genexpr>z1ResponseHandler.handle_request.<locals>.<genexpr>  sX      
,6swwy7I4,Pcggi(VX
  EE&MW$
$
s   AA    
num_frameschat_template_kwargstoolsTNpt)add_generation_promptr   return_tensorsreturn_dicttokenizeload_audio_from_video	input_idsstream)gen_managertool_configreasoning_config)_validate_request_resolve_modelmodel_managerget_model_modalitygeneration_stateuse_continuous_batchingloggerwarningget_manager_normalize_input"get_processor_inputs_from_messagesanyupdater   r   _normalize_toolsapply_chat_templater&   
MULTIMODALtodevice_build_generation_configr3   init_cbr*   r)   
_streaming_non_streaming)ra   r   rM   model_idmodelr   modalityr   r   messagesprocessor_inputs	has_videor   r   inputs
gen_configr   r   	streamings                      r=   handle_requestzResponseHandler.handle_requesto  s{     	t$%)%8%8%>"%%%88)8T&&>>uhO3H:VF8LM++777P
 ((.BB8XV 
'
 
	 &(13 .##D$=$=>##DHH-C$D$JK%%dhhw&78...	
"&#)4t"*h.A.A"A"Oi	
 #	
 YYu||,F2249P9PY_2`
z2@D@Q*9e<W[/	5&BUVHHXt,	??''!1 #   ,,''!1 -    s   H3H<5H:6H<r   Nc                     | s| S | D cg c]5  }d|vr-d|j                         D ci c]  \  }}|dk7  s|| c}}dn|7 c}}}S c c}}w c c}}}w )a  Normalize Responses API tool definitions for ``apply_chat_template``.

        The Responses API uses a flat format: ``{"type": "function", "name": ..., "parameters": ...}``
        while ``apply_chat_template`` expects a nested format:
        ``{"type": "function", "function": {"name": ..., "parameters": ...}}``.
        Already-nested tools are passed through unchanged.
        functionrq   )rq   r   )items)r   tkvs       r=   r   z ResponseHandler._normalize_tools  sk     L 
 
 ^hop]pZqwwy-Xtq!AQWKad-XYvww
 	
-X
s   AA
A

A
Ac                 B   | d   }| j                  d      }t        |t              rd|dg}nCt        |t              r&|rd|d   vrd|dg}n#t        j                  |      }nt        dd	      |r,|r|d   d   d
k(  r
||d   d<   |S |j                  dd
|d       |S )uQ  Normalize the Responses API ``input`` field into chat messages.

        The Responses API accepts multiple input formats. This method converts them
        into a structure close to what ``apply_chat_template`` expects (messages with
        ``role``, ``content``, ``tool_calls``, ``tool_call_id``). Further processing
        is done by ``get_processor_inputs_from_messages``.

        NOTE: if this conversion logic grows too complex, consider having separate
        ``get_processor_inputs_from_messages`` implementations for chat completions
        and the Responses API instead of funneling both through the same path.

        Formats handled:
            - **String** → single user message.
            - **Flat content list** (``input_text``, ``input_image``, no ``role``) → user message.
            - **Multi-turn list** — messages and tool call items (``function_call``,
              ``function_call_output``) from a previous response, converted via
              :meth:`_normalize_response_items`.

        If ``instructions`` is present, it is prepended as a system message.
        inputinstructionsr@   r   r}   r   r     z 'input' must be a string or liststatus_codedetailsystemr}   )r   r   r8   r   r   _normalize_response_itemsr   insert)r   inpr   r   s       r=   r   z ResponseHandler._normalize_input  s    , 7mxx/c3!'C89HT"vSV+%+<=*DDSIC8Z[[ HQK/8;)5I&  H#NOr<   r   c                 >   g }d}| D ]  }|j                  d      }|dk(  r,dj                  d |j                  d      xs g D              }Fd|v r;|d   |j                  dd      d}||d   d	k(  r||d
<   d}|j                  |       |dk(  rY|d   |d   |d   dd}|r0|d   d   d	k(  r%|d   j                  dg       j                  |       |j                  d	|gd       |dk(  r|j                  d|d   |d   d       t	        dd|       |S )u  Convert a list of Responses API items into chat messages.

        Input items may be a mix of:
            - Messages (``EasyInputMessageParam`` with ``role``, or ``type: "message"``).
            - ``reasoning`` — buffered and attached as ``reasoning_content`` to the next assistant message.
            - ``function_call`` — merged as ``tool_calls`` onto the preceding assistant message.
            - ``function_call_output`` — converted to ``role: "tool"`` messages.
        Nrq   ry   rS   c              3   &   K   | ]	  }|d      yw)r?   Nr;   )r   r   s     r=   r   z<ResponseHandler._normalize_response_items.<locals>.<genexpr>  s     +Y!AfI+Ys   r}   r   r   r   reasoning_contentr   r   r   r   )r   r   )r{   r   r`   )r   r`   function_call_outputtoolro   )r   tool_call_idr}   r   zUnsupported input item type: r   )r   joinr   
setdefaultr   )r   r   pending_reasoningr   	item_typer   tcs          r=   r  z)ResponseHandler._normalize_response_items  sf    (, "	kD(IK'$&GG+Ytxx	?R?XVX+Y$Y!~#F|B8OP$0T&\[5P/@C+,(,%$o-y/)-fDDU V V 4 CRL++L"=DDRHOO[$MN44 &(,Y#'> $>[\e[h<ijjE"	kH r<   r   r.   r   z(ProcessorMixin | PreTrainedTokenizerFastr   r   r   r-   r   r   r   c           	      \  	 |j                  |||	|
      \  |d   }t        |t              rt        |      n|j                  d   d t        j
                         |dg |j                  dd      dd	d
t        t        df   f	fd}t         |       d      S )zDGenerate a streaming Responses API reply (SSE) using DirectStreamer.)rM   r   r   r   r  rP   rs   parallel_tool_callsFauto)r{   
created_atr   objectr   r  rF   rc   Nc            	       K   t              } 	 dj                  | j                                d}|s;
j                          d {   g}	 	 |j	                  
j                                 | j                  r"dj                  | j%                                | j&                  s"dj                  | j)                                dj                  | j-                                r_t/        	j0                  d         }|rCt3        |      D ]5  \  }}dj                  | j5                   d| |d   |d	                7 dj                  | j7                  t9        j:                                     y 7 ?# t        j                  $ r Y nw xY wg }|D ]N  }|d} nGt        |t              rbt        j                  d|j                          |j                  | j                  |j                               dj                  |        y t        |t              rL| j                  s|j                  | j!                                |j                  | j#                  |             | j                  r|j                  | j%                                | j&                  s|j                  | j)                                |j                  | j+                  |             Q |rdj                  |       |sސ# t<        t        j>                  f$ r jA                           w xY ww)
N)rM   rN   rS   FTz"Exception in response generation: schema_tool_call_r   r   )!rL   r  rv   r   r   
get_nowaitasyncio
QueueEmptyr   r(   r   r   r   r   r'   r\   r   r   r   r]   r   r   r   r,   generated_token_ids	enumerater   r   compute_usagetotal_tokensGeneratorExitCancelledErrorcancel)builderdonebatchr   r?   parsedir  	input_lenr   queuerM   rN   streamerr   s           r=   event_streamz0ResponseHandler._streaming.<locals>.event_stream`  s    ,
VghG<ggg44677 #(99;./E"!LL)9)9);< #< ))'''":":"<==++'''"7"7"9::ggg44677 -i9U9UWbckWlmF%.v%6 EAr"$'' ' 1 1ZLA32OQSTZQ[]_`k]l m# 
 ggg//iI^I^0_`aaa / #--  (*E % C<#'D!%dL9"LL+MdhhZ)XY!LLtxx)@A"$''%.0"%dM:#*#9#9 %W-D-D-F G!LL)@)@)FG&55 %W-E-E-G H#*#7#7 %W-B-B-D E!LL););D)AB%C(  ggen,= d "7#9#9:  !	sa   M9L F
L !F 3DL 	M
L F# L "F##BL &M'C-L +MMztext/event-stream)
media_type)
generate_streamingr   r   lenshapetimer   r   r8   r
   )ra   rM   r   r   r   r   r   r   r   r   r   r   r*  r'  r(  rN   r)  s    ` `     `   @@@@r=   r   zResponseHandler._streaming:  s     &88!#- 9 
x ;'	&0D&AC	NyWYGZ	 *&))+ #'88,A5#I!	
>	N39$= >	 >	@ !<OPPr<   c                   K   |j                  |||||       d{   \  }}}g }|
9t        ||||
      \  }}|&|j                  t        d| dg d|dgd             |j                  t	        d	| d
ddt        d|g       gg              |	Rt        |||	d         }|r@t        |      D ]2  \  }}| d| }|j                  t        ||d|d   |d   d             4 t        |t        |            }t        d| t        j                         d||d|g |j                  dd      d
      }t        |j                  d            S 7 4w)z;Generate a non-streaming Responses API reply (single JSON).)rM   NrR   ry   r   r   r   rz   rQ   r   r   r   r   r   r  r  r   r   r   r   rP   rs   r  Fr  )
r{   r  ri   r   ro   r  r   r   r  rF   T)exclude_none)generate_non_streamingr+   r   r   r   r   r,   r  r   r  r-  r   r/  r   r	   
model_dump)ra   rM   r   r   r   r   r   r   r   r   r   rZ   r'  generated_idsoutput_itemsr  r%  r&  r  r   r   rs   s                         r=   r   zResponseHandler._non_streaming  s     5@4V4V9fjZ 5W 5
 /
+	9m '+:9mU^`p+q(I( ,##) -( "*:DU!V W* 	!*&" +Y\^_`		
 "%iH@UVF&v. EAr)l+aS9E ''0$$)!0!#F&(o#.	 i]);<zl#yy{ $)> F
 H//T/BCCw/
s   EED4Emodel_generation_configr   c                 t    t         |   |||      }|j                  d      t        |d         |_        |S )zXApply Responses API params (``max_output_tokens``) on top of the base generation config.r   max_output_tokens)superr   r   r:   max_new_tokens)ra   r   r6  r   r3   	__class__s        r=   r   z(ResponseHandler._build_generation_config  sF    !G<TCZci<j88'(4/248K3L/M,  r<   )NN)F)r5   r6   r7   r   r2   _valid_params_classUNUSED_RESPONSE_FIELDS_unused_fieldsr   r8   r
   r	   r   staticmethodr   r   r   r  r$   r   r   boolr   __classcell__)r;  s   @r=   r   r   i  s   5C+NU U3 UCTWcCc Ur 
T
T 1 
d4j46G 
 
 *t *T
 * *X 0d 0T
 0 0| $((,fQfQ !fQ >	fQ
 fQ fQ fQ 'fQ )fQ D[fQ +fQ 
fQh $((,IDID !ID >	ID
 ID ID ID 'ID )ID D[ID +ID 
IDZ!T !L^ !hl ! !r<   r   input_tokensoutput_tokensrc   c                 ~    t        | || |z   t        dddidt        j                  v rddini t        d            S )a  Build a ``ResponseUsage`` object for a Responses API reply.

    Args:
        input_tokens (`int`): Number of prompt tokens.
        output_tokens (`int`): Number of generated tokens.

    Returns:
        `ResponseUsage`: Usage statistics with zero-filled detail fields.
    cached_tokensr   cache_write_tokens)reasoning_tokens)rB  rC  r  input_tokens_detailsoutput_tokens_detailsr;   )r"   r    model_fieldsr!   )rB  rC  s     r=   r  r    sX     !#!M1/ 

+?CUCbCb+b#Q'hj
 21E	 	r<   )Br   r  r/  collections.abcr   typingr   utilsr   utils.import_utilsr   fastapir   fastapi.responsesr	   r
   openai.types.responsesr   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   -openai.types.responses.response_create_paramsr   %openai.types.responses.response_usager    r!   r"   r$   r%   r&   r'   r(   r)   r*   r+   r,   transformersr-   r.   r/   r0   
get_loggerr5   r   r2   r=  rL   r   r:   r  r;   r<   r=   <module>rV     s      *    4 %A     , \ll
 
 
 gg 
		H	%0MUZ 
  C
 C
LO!k O!d C M r<   