+
    QV-jÏ  ã                   ór  € R t ^ RIt^ RIt^ RIt^ RIt^ RIt^ RIHtHt ^ RI	H
t
 ^ RIHt ^ RIHt ^ RIHt ^ RIHt ]'       d3   ^ RIt^ RIt^ RIt^ RIHtHtHtHtHt ^ R	IHt ^ R
IHt ^ RI H!t! ^RI"H#t# ]PH                  ! ]%4      t&Rt' ! R R]PP                  4      t) ! R R4      t* ! R R]+4      t, ! R R]-4      t. ! R R]/4      t0R\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!RR"R&R'RR(/R)RR"R*R+//////t1R, R- lt2R. R/ lt3R0 R1 lt4R2R3.R4R5RRR"R&R6RR(/R7RR(//R8R9//t5R:R2. R^OR4R;//t6R_R< R= llt7R> R? lt8R@ RA lt9RB RC lt: ! RD RE4      t;RF RG lt< ! RH RI4      t= ! RJ RK4      t>RL RM lt?RN RO lt@ ! RP RQ4      tA ! RR RS]4      tB ! RT RU]B4      tC ! RV RW]B4      tD ! RX RY4      tE ! RZ R[4      tFR# )`z?
Shared types, constants, and utilities for the serving layer.
N)ÚABCÚabstractmethod)ÚCallable©ÚFuture)ÚQueue)ÚTYPE_CHECKING)Úlogging)ÚContinuousBatchingConfigÚGenerationConfigÚPreTrainedModelÚPreTrainedTokenizerFastÚProcessorMixin)ÚContinuousBatchingManager)ÚGenerationOutput)Ú	Scheduler)ÚModelManagerzx-request-idc                   ó*   € ] tR t^9tRtRtRtRtRtRt	R# )ÚModalityÚLLMÚVLMÚ
MULTIMODALÚSTTÚTTS© N)
Ú__name__Ú
__module__Ú__qualname__Ú__firstlineno__r   r   r   r   r   Ú__static_attributes__r   ó    Úo/Volumes/fast/ai/experiments/ui-tars-smoke/.venv/lib/python3.14/site-packages/transformers/cli/serving/utils.pyr   r   9   s   † Ø
€CØ
€CØ€JØ
€CØ
„Cr    r   c                   ó6   a € ] tR t^At o RtV 3R lR ltRtV tR# )Ú_StreamErrorz5Sentinel to signal an error from the generate thread.c                ó    <€ V ^8„  d   QhRS[ /# )é   Úmsg©Ústr)ÚformatÚ__classdict__s   "€r!   Ú__annotate__Ú_StreamError.__annotate__D   s   ø€ ÷ ñ ™Cñ r    c                ó   € Wn         R # ©N©r&   )Úselfr&   s   &&r!   Ú__init__Ú_StreamError.__init__D   s   € ØŽr    r/   N)r   r   r   r   Ú__doc__r1   r   Ú__classdictcell__©r*   s   @r!   r#   r#   A   s   ø‡ € Ù?÷ö r    r#   c                   ó   € ] tR t^HtRtRtR# )Ú_GenerationCancelledzERaised inside ``DirectStreamer.put()`` to abort ``model.generate()``.r   N©r   r   r   r   r3   r   r   r    r!   r7   r7   H   s   † ÝOr    r7   c                   ó   € ] tR t^LtRtRtR# )ÚReasoningTextzÃTagged str subclass: text chunk belonging to a thinking/reasoning block.

Streamers wrap reasoning text with this so handlers can route it to
``reasoning_content`` deltas instead of ``content``.
r   Nr8   r   r    r!   r:   r:   L   ó   † õr    r:   c                   ó   € ] tR t^TtRtRtR# )ÚCBWorkerDeadErrorzãRaised when a request is submitted to a CB worker that has died.

Surfaced as 503 by the FastAPI exception handler. Carries the original error message
that killed the worker so the client knows why the server is in this state.
r   Nr8   r   r    r!   r=   r=   T   r;   r    r=   Ústcz<tool_call>Úetcz</tool_call>Úschemazx-regex-iteratorz<tool_call>(.*?)</tool_call>ÚtypeÚarrayÚitemsÚobjectzx-parserÚjsonz9<function=(?P<name>[^>\n]+)>(?P<arguments>.*?)</function>Ú
propertiesÚnameÚstringÚ	argumentszx-regex-key-valuez<<parameter=(?P<key>[^>\n]+)>\s*(?P<value>.*?)\s*</parameter>c                ó6   € V ^8„  d   QhRRR\         R,          /# ©r%   Úmodelr   ÚreturnN©Údict)r)   s   "r!   r+   r+   Œ   s#   € ÷ Bñ BÐ+<ð BÄÈÅñ Br    c                óÎ  a
€ \        V RV 4      p\        VRR4      p\        VRR4      p\        VRR4      pV'       d"   V'       d   V'       d   VR,          R,          pM^VP                  P                  o
\        V
3R l\        P                  4        4       R4      pVf   R# VR	,          VR
,          VR,          rdpVP                  V4      pVP                  V4      p	RVRVRV	/# )aR  Return tool call config for the model, or ``None`` if tool calls are not supported.

Returns a dict with:
    - ``schema`` (`dict`): Schema to pass to ``tokenizer.parse_response(block, schema)``.
    - ``stc_id`` (`int`): Token ID of the start-of-tool-call delimiter.
    - ``etc_id`` (`int`): Token ID of the end-of-tool-call delimiter.
Ú	tokenizerÚ	stc_tokenNÚ	etc_tokenÚresponse_schemarF   Ú
tool_callsc              3   ó>   <"  € T F  w  rSV9   g   K  Vx € K  	  R # 5ir.   r   )Ú.0ÚtypesÚvÚ
model_types   &  €r!   Ú	<genexpr>Ú'get_tool_call_config.<locals>.<genexpr>    s   øé € Ð_Ñ+G™x˜uÈ:ÐY^ÑK^ŸšÓ+Gùó   ƒ“
r>   r?   r@   Ústc_idÚetc_id)ÚgetattrÚconfigrZ   ÚnextÚ_TOOL_CALL_FALLBACKSrC   Úconvert_tokens_to_ids)Ú	processorrL   rQ   r>   r?   rT   r@   Úfallbackr^   r_   rZ   s   &&        @r!   Úget_tool_call_configrg   Œ   sØ   ø€ ô ˜	 ;°	Ó:€IÜ
�)˜[¨$Ó
/€CÜ
�)˜[¨$Ó
/€CÜ˜iÐ):¸DÓA€O÷ �s—Ø  Õ.¨|Õ<‰ð —\‘\×,Ñ,ˆ
ÜÔ_Ô+?×+EÑ+EÔ+GÓ_ÐaeÓfˆØÒÙØ# E�?¨H°U­O¸XÀhÕ=O�&ˆà×,Ñ,¨SÓ1€FØ×,Ñ,¨SÓ1€FØ�f˜h¨°¸&ÐAÐAr    c                ó0   € V ^8„  d   QhR\         R\         /# )r%   Ú	tool_callrM   rN   )r)   s   "r!   r+   r+   ª   s   € ÷ ñ ¤Dð ¬Tñ r    c                ó¾   € V P                  RV 4      pVP                  R/ 4      pRVR,          R\        V\        4      '       g   \        P                  ! V4      /# T/# )aƒ  Normalize a parsed tool call to ``{"name": str, "arguments": str}``.

Different models return different structures from ``parse_response``:
- Gemma: ``{"function": {"name": ..., "arguments": {...}}}`` (nested, arguments as dict)
- Qwen:  ``{"name": ..., "arguments": {...}}`` (flat, arguments as dict)

The OpenAI API expects ``arguments`` as a JSON **string**, so we ``json.dumps`` it.
ÚfunctionrI   rG   )ÚgetÚ
isinstancer(   rE   Údumps)ri   rk   rI   s   &  r!   Ú_normalize_tool_callro   ª   s_   € ð �}‰}˜Z¨Ó3€HØ—‘˜[¨"Ó-€Ià�˜Õ Ø´*¸YÌ×2LÒ2L”T—Z’Z 	Ó*ðð àR[ðð r    c                óT   € V ^8„  d   QhR\         R\        \         ,          R,          /# )r%   r@   rM   N)rO   Úlist)r)   s   "r!   r+   r+   »   s#   € ÷ .ñ .´tð .ÄÄTÅ
ÈTÕ@Qñ .r    c                óÀ   € V P                  W4      pV'       g   R# \        V\        4      '       g   V.pV Uu. uF  p\        V4      NK  	  ppV'       d   V# R# u upi )a  Parse tool calls from generated token IDs using ``tokenizer.parse_response``.

Args:
    processor: The processor or tokenizer.
    generated_ids: Token IDs from generation. Passed directly to ``parse_response``
        which decodes them internally, preserving special tokens that
        ``skip_special_tokens=True`` would strip (e.g. Gemma's ``<|tool_call>``).
    schema: The tool call schema (from ``response_schema`` or ``_TOOL_CALL_FALLBACKS``).

Returns a list of ``{"name": str, "arguments": str}`` dicts, or ``None`` if none found.
N)Úparse_responserm   rq   ro   )re   Úgenerated_idsr@   Úparsedri   rU   s   &&&   r!   Úparse_tool_callsrv   »   s[   € ð ×%Ñ% mÓ<€FßÙÜ�fœd×#Ò#Ø�ˆÙCIÓJÁ6°iÔ& yÖ1Á6€JÐJß#ˆ:Ð-¨Ð-ùò Ks   ¹AÚstartz<think>Úendz</think>ÚthinkingÚcontentzx-regexzK(?:<think>)?(?P<thinking>.*?)</think>(?P<content>.*?)(?:<\|[^|<>\s]+\|>)?\ZÚgemma4z
<channel|>c                ó6   € V ^8„  d   QhRRR\         R,          /# rK   rN   )r)   s   "r!   r+   r+   ì   s"   € ÷ ñ Ð+<ð ÔQUÐX\ÕQ\ñ r    c                óŠ  a	a
€ \        V RV 4      o
VP                  P                  P                  4       o	\	        V	3R l\
        P                  4        4       \        4      pVR,           Uu. uF  pS
P                  V4      NK  	  ppS
P                  VR,          4      p\        ;QJ d    V
3R lV 4       F  '       g   K   RM	  RM! V
3R lV 4       4      '       g   VRS
P                  39   d   R# \        S
R	R4      pV'       d   R
VR,          9   g   \        R,          pRVRVRV/pVe   \        W%4      VR&   V# u upi )a¥  Return reasoning config for the model, or ``None`` if not supported.

The config drives both streaming detection (token IDs) and post-hoc parsing
(response schema). Returns a dict with:
    - ``start_ids`` (`list[int]`): Token ID sequence that opens a thinking block.
    - ``end_id`` (`int`): Token ID that closes the block.
    - ``schema`` (`dict`): Response schema with ``thinking`` / ``content``
      properties for :func:`parse_reasoning`.
    - ``start_in_thinking`` (`bool`, only when ``input_ids`` is given): Whether
      the rendered prompt already opened an unclosed thinking block (prefilled
      by the template), so the model's output begins inside the block.
rQ   c              3   ó>   <"  € T F  w  rVS8X  g   K  Vx € K  	  R # 5ir.   r   )rW   ÚkrY   rZ   s   &  €r!   r[   Ú'get_reasoning_config.<locals>.<genexpr>ü   s   øé € ÐCÑ/‰tˆq°1¸
±?�ŠÓ/ùr]   rw   rx   c              3   óD   <"  € T F  qR SP                   39   x € K  	  R # 5ir.   )Úunk_token_id)rW   ÚtidrQ   s   & €r!   r[   r€     s   øé € Ð
F¹I°S�4˜×/Ñ/Ð0Ö0»Iùs   ƒ TFNrT   ry   rF   r@   Ú	start_idsÚend_idÚstart_in_thinking)r`   ra   rZ   Úlowerrb   Ú_THINKING_TOKENSrC   Ú_DEFAULT_THINKING_TOKENSrd   Úanyr‚   Ú_starts_in_thinking)re   rL   Ú	input_idsÚthinking_tokensÚtr„   r…   r@   ra   rZ   rQ   s   &&&      @@r!   Úget_reasoning_configr�   ì   s  ù€ ô ˜	 ;°	Ó:€IØ—‘×(Ñ(×.Ñ.Ó0€JÜÜCÔ'×-Ñ-Ô/ÓCÜ ó€Oð >MÈWÖ=UÓVÑ=U¸�×0Ñ0°Ö3Ñ=U€IÐVØ×,Ñ,¨_¸UÕ-CÓD€Fß
ƒsÔ
F¹IÓ
F‡s‡s‚sÔ
F¹IÓ
F×FÒFÈ&ÐUYÐ[d×[qÑ[qÐTrÔJrÙô �YÐ 1°4Ó8€Fß�z V¨LÕ%9Ô9Ü)¨(Õ3ˆØ ¨H°f¸hÈÐO€FØÒÜ&9¸)Ó&OˆÐ"Ñ#Ø€Mùò Ws   Á+E c          	      ól   € V ^8„  d   QhR\         R\        R\        \         \         R,          3,          /# )r%   rz   Úreasoning_configrM   N)r(   rO   Útuple)r)   s   "r!   r+   r+     s4   € ÷ ñ ´sð Ìdð ÔW\Ô]`ÔbeÐhlÕblÐ]lÕWmñ r    c                óÜ   € V P                  WR,          4      pV'       d/   VP                  RR4      pV'       d   VP                  RR4      V3# VP                  R4      '       d   RV3# VR3# )u‚  Split generated output into ``(content, reasoning_content)`` via ``parse_response``.

If the schema's regex matches (closing marker present), use it. For prompts
that prefill the opener (QwQ-32B, DeepSeek-R1) the entire output is reasoning
until ``</think>`` arrives â€” when that's truncated, fall back to treating
all decoded text as reasoning. Returns ``(content, None)`` otherwise.
r@   ry   Ú rz   r†   N)rs   rl   )re   rt   rz   r‘   ru   Ú	reasonings   &&&&  r!   Úparse_reasoningr–     sm   € ð ×%Ñ% mÀhÕ5OÓP€FßØ—J‘J˜z¨2Ó.ˆ	ßØ—:‘:˜i¨Ó,¨iÐ7Ð7ð ×ÑÐ/×0Ò0Ø�7ˆ{ÐØ�Dˆ=Ðr    c                óF   € V ^8„  d   QhR\         \        ,          R\        /# )r%   r„   rM   )rq   ÚintÚbool)r)   s   "r!   r+   r+   "  s   € ÷ ñ ¬d´3­ið ¼Dñ r    c                ól  € \        V R4      '       d   V P                  4       p V '       d9   \        V ^ ,          \        4      '       d   \	        V 4      ^8w  d   R# V ^ ,          p \	        V4      pR F@  p\	        V 4      W#,           8¼  g   K  \	        V 4      V,
          pWV,
          V V8X  g   K?   R# 	  R# )uY  True if the rendered prompt ends with an unclosed thinking block.

Some reasoning-model chat templates prefill the thinking opener as the final
prompt tokens (e.g. DeepSeek-R1, QwQ-32B emit ``<think>\n`` at the end when
``add_generation_prompt=True``). In those cases the model resumes *inside*
the block, so its output contains only ``...reasoning</think>answer`` with
no opening tag â€” the streamer must start with ``_inside_thinking=True``.

The prefill always lands at the tail of the prompt (optionally followed by a
single whitespace token like ``\n``), so we only inspect the last few tokens.
ÚtolistFT)é    é   )Úhasattrr›   rm   rq   Úlen)rŒ   r„   ÚnÚtrailingrx   s   &&   r!   r‹   r‹   "  s’   € ô ˆy˜(×#Ò#Ø×$Ñ$Ó&ˆ	ß”Z 	¨!¥¬d×3Ò3Üˆy‹>˜QÔÙØ˜a•Lˆ	ÜˆI‹€AãˆÜˆy‹>˜Q�\Ö)Ü�i“. 8Õ+ˆCØ˜q� 3Ð'¨9Ö4Úñ	 ñ
 r    c                ó0   € V ^8„  d   QhR\         R\        /# )r%   Útoken_idrM   )r˜   r™   )r)   s   "r!   r+   r+   >  s   € ÷ ñ ´ð ¼ñ r    c                ó–  € V P                   f   R# V P                  '       d   WP                  8X  d
   RV n        R# R# V P                   \        V P                  4      ,          pW8w  d
   . V n        R# V P                  P                  V4       \        V P                  4      \        V P                   4      8X  d   RV n        . V n        R# )uE  Mutate ``streamer``'s thinking state; return ``True`` if ``token_id`` is a start or end token.

Shared between :class:`DirectStreamer` and :class:`CBStreamer` â€” both track the
same four attributes (``_thinking_start_ids``, ``_thinking_end_id``,
``_inside_thinking``, ``_thinking_prefix``) and need identical edge handling.
FT)Ú_thinking_start_idsÚ_inside_thinkingÚ_thinking_end_idrŸ   Ú_thinking_prefixÚappend)Ústreamerr£   Úexpecteds   && r!   Ú_advance_thinking_stater¬   >  sª   € ð ×#Ñ#Ò+ÙØ× × Ð Ø×0Ñ0Ô0Ø(-ˆHÔ%ÙÙØ×+Ñ+¬C°×0IÑ0IÓ,JÕK€HØÔØ$&ˆÔ!ÙØ×Ñ×$Ñ$ XÔ.Ü
ˆ8×$Ñ$Ó%¬¨X×-IÑ-IÓ)JÔJØ$(ˆÔ!Ø$&ˆÔ!Ùr    c                   ó~   a € ] tR tRt o RtV 3R lR ltV 3R lR ltV 3R lR ltV 3R	 lR
 ltV 3R lR lt	Rt
V tR# )ÚDownloadAggregatoriW  zýAggregates byte-progress across multiple concurrent download tqdm bars.

huggingface_hub opens one tqdm bar per file shard. This class tracks them all and emits
a single aggregate ``{"stage": "download", "progress": {...}}`` event whenever any updates.
c                ó&   <€ V ^8„  d   QhRS[ RS[/# )r%   ÚenqueueÚmodel_id)r   r(   )r)   r*   s   "€r!   r+   ÚDownloadAggregator.__annotate__^  s   ø€ ÷ 5ñ 5¡ð 5±Cñ 5r    c                ó:   € Wn         W n        / V n        R V n        R # r.   )r°   rL   ÚbarsÚlast_emitted_current)r0   r°   r±   s   &&&r!   r1   ÚDownloadAggregator.__init__^  s   € ØŒØŒ
Ø79ˆŒ	Ø04ˆÖ!r    c                ó8   <€ V ^8„  d   QhRS[ RS[ R,          RR/# )r%   Úbar_idÚtotalNrM   ©r˜   )r)   r*   s   "€r!   r+   r²   d  s&   ø€ ÷ ñ ™sð ©3°­:ð ¸$ñ r    c                óH   € ^ V3V P                   V&   V P                  4        R# )z6Register a new download bar with its total byte count.N©r´   Ú_emit)r0   r¸   r¹   s   &&&r!   ÚregisterÚDownloadAggregator.registerd  s   € à ˜Jˆ�	‰	�&ÑØ�
‰
Žr    c                ó>   <€ V ^8„  d   QhRS[ RS[ RS[ R,          RR/# )r%   r¸   Úcurrentr¹   NrM   rº   )r)   r*   s   "€r!   r+   r²   i  s-   ø€ ÷ ñ ™Sð ©3ð ±s¸Tµzð Àdñ r    c                óF   € W#3V P                   V&   V P                  4        R# )z>Update a bar's current byte count and emit aggregate progress.Nr¼   )r0   r¸   rÁ   r¹   s   &&&&r!   ÚupdateÚDownloadAggregator.updatei  s   € à$Ð,ˆ�	‰	�&ÑØ�
‰
Žr    c                ó$   <€ V ^8„  d   QhRS[ RR/# )r%   r¸   rM   Nrº   )r)   r*   s   "€r!   r+   r²   n  s   ø€ ÷ ñ ™Cð  Dñ r    c                ó   € R # r.   r   )r0   r¸   s   &&r!   ÚcloseÚDownloadAggregator.closen  ó   € Ùr    c                ó   <€ V ^8„  d   QhRR/# ©r%   rM   Nr   )r)   r*   s   "€r!   r+   r²   q  s   ø€ ÷ 
ñ 
�tñ 
r    c                ót  € \        R  V P                  P                  4        4       4      pWP                  8X  d   R# Wn        V P                  P                  4        UUu. uF  w  r#Vf   K  VNK  	  pppV'       d   \        V4      MRpV P	                  RRRV P
                  RRRRVR	V//4       R# u uppi )
c              3   ó*   "  € T F	  w  rVx € K  	  R # 5ir.   r   )rW   ÚcÚ_s   &  r!   r[   Ú+DownloadAggregator._emit.<locals>.<genexpr>r  s   é € Ð;Ñ(:¡ �!Ó(:ùs   ‚NÚstatusÚloadingrL   ÚstageÚdownloadÚprogressrÁ   r¹   )Úsumr´   Úvaluesrµ   r°   rL   )r0   Úagg_currentrÏ   rŽ   ÚtotalsÚ	agg_totals   &     r!   r½   ÚDownloadAggregator._emitq  s¢   € ÜÑ;¨¯	©	×(8Ñ(8Ô(:Ó;Ó;ˆØ×3Ñ3Ô3ÙØ$/Ô!Ø $§	¡	× 0Ñ 0Ô 2ÔDÑ 2™˜°a—!�!Ñ 2ˆÑDß#)”C˜”K¨tˆ	Ø�‰à˜)Ø˜Ÿ™Ø˜Ø˜Y¨°W¸iÐHð	ö	
ùó Es   Á B4Á-B4)r´   r°   rµ   rL   N)r   r   r   r   r3   r1   r¾   rÃ   rÇ   r½   r   r4   r5   s   @r!   r®   r®   W  s<   ø‡ € ñ÷5ð 5÷ð ÷
ð ÷
ð ÷
ö 
r    r®   c                ó<   € V ^8„  d   QhR\         R\        R\        /# )r%   Úcallbackr±   rM   )r   r(   rA   )r)   s   "r!   r+   r+   ‚  s&   € ÷ Dñ D¤xð D¼3ð DÄ4ñ Dr    c                óP   a aa€ ^ RI Hp \        S S4      o ! V VV3R lRV4      pV# )uè  Create a tqdm subclass that routes progress to a callback.

Bars with ``unit="B"`` are download bars â€” aggregated via ``DownloadAggregator``.
Other bars (e.g. "Loading weights") emit ``weights`` stage events.

Args:
    callback (`callable`): Called with a dict payload
        ``{"status": "loading", "model": ..., "stage": ..., "progress": ...}``.
    model_id (`str`): The model ID (included in progress payloads).

Returns:
    A tqdm subclass that forwards progress to *callback*.
)Útqdmc                   óh   <a a€ ] tR tRt oV V3R ltRVVV3R lltVVV3R ltV V3R ltRtVt	V ;t
# )Ú.make_progress_tqdm_class.<locals>.ProgressTqdmi”  c                ó   <€ VP                  R 4      ;'       g    RV n        RVR&   \        SV `  ! V/ VB  ^ V n        RV n        V P                  R8X  d9   \        V 4      V n        SP                  V P                  V P                  4       R# R# )ÚunitÚitTÚdisableÚBNéÿÿÿÿ)
rl   Ússe_unitÚsuperr1   r    Úlast_emittedÚidÚ_bar_idr¾   r¹   )r0   ÚargsÚkwargsÚ	__class__Údownload_aggregators   &*,€€r!   r1   Ú7make_progress_tqdm_class.<locals>.ProgressTqdm.__init__•  sz   ø€ Ø"ŸJ™J vÓ.×6Ð6°$ˆDŒMØ $ˆF�9ÑÜ‰GÒ˜dÐ- fÒ-ØˆDŒFØ "ˆDÔØ�}‰} Ô#Ü! $›x�”Ø#×,Ñ,¨T¯\©\¸4¿:¹:ÖFñ $r    c                óz  <€ Vf   ^pV ;P                   V,          un         V P                  R8X  d4   SP                  V P                  V P                   V P                  4       R # V P                   V P
                  8w  d<   V P                   V n        S! RRRSRRRRV P                   R	V P                  //4       R # R # )
Nræ   rÑ   rÒ   rL   rÓ   ÚweightsrÕ   rÁ   r¹   )r    rè   rÃ   rì   r¹   rê   )r0   r    rÝ   rð   r±   s   &&€€€r!   rÃ   Ú5make_progress_tqdm_class.<locals>.ProgressTqdm.updateŸ  s™   ø€ ØŠyØ�Ø�FŠF�a�K�FØ�}‰} Ô#Ø#×*Ñ*¨4¯<©<¸¿¹ÀÇÁÖLØ—‘˜4×,Ñ,Ô,Ø$(§F¡F�Ô!Ùà  )Ø Ø Ø" Y°·±¸ÀÇÁÐ$Lð	öñ -r    c              3  óž  <"  € V P                    F·  pV ;P                  ^,          un        V P                  R8X  d3   SP                  V P                  V P                  V P
                  4       MTV P                  V P                  8w  d:   V P                  V n        S! RRRSRRRRV P                  R	V P
                  //4       Vx € K¹  	  R
# 5i)r�   ræ   rÑ   rÒ   rL   rÓ   ró   rÕ   rÁ   r¹   N)Úiterabler    rè   rÃ   rì   r¹   rê   )r0   ÚitemrÝ   rð   r±   s   & €€€r!   Ú__iter__Ú7make_progress_tqdm_class.<locals>.ProgressTqdm.__iter__°  s¢   øé € ØŸœ�Ø—’˜!••Ø—=‘= CÔ'Ø'×.Ñ.¨t¯|©|¸T¿V¹VÀTÇZÁZÕPØ—V‘V˜t×0Ñ0Ô0Ø(,¯©�DÔ%Ùà$ iØ# XØ# YØ&¨°D·F±F¸GÀTÇZÁZÐ(Pð	ôð ”
ó &ùs   ƒC
Cc                ó|   <€ V P                   R 8X  d   SP                  V P                  4       \        SV `  4        R# )ræ   N)rè   rÇ   rì   ré   )r0   rï   rð   s   &€€r!   rÇ   Ú4make_progress_tqdm_class.<locals>.ProgressTqdm.closeÁ  s*   ø€ Ø�}‰} Ô#Ø#×)Ñ)¨$¯,©,Ô7Ü‰G‰MŽOr    )rì   rê   r    rè   )r�   )r   r   r   r   r1   rÃ   rø   rÇ   r   r4   Ú__classcell__)rï   r*   rÝ   rð   r±   s   @@€€€r!   ÚProgressTqdmrá   ”  s$   ú‡ € ö	G÷	ñ 	÷"	÷"	ö 	r    rý   )Ú	tqdm.autorß   r®   )rÝ   r±   Ú	base_tqdmrý   rð   s   ff  @r!   Úmake_progress_tqdm_classr   ‚  s/   ú€ õ ,ä,¨X°xÓ@Ð÷0ñ 0�yô 0ðd Ðr    c                   óp   a € ] tR tRt o RtRV 3R lR lltV 3R lR ltV 3R lR	 ltV 3R
 lR ltRt	V t
R# )ÚDirectStreameriÉ  ar  Streamer for ``model.generate()`` (used by :class:`GenerateManager`).

Implements the ``put``/``end`` protocol that ``model.generate()`` expects:
generate calls ``put(token_tensor)`` after each decode step, and ``end()``
when generation is complete. Tokens are decoded incrementally via the Rust
``DecodeStream`` (O(1) per token) and pushed as text to an asyncio.Queue.
Nc                ó€   <€ V ^8„  d   QhRRRS[ P                  RS[ P                  RS[RS[R,          RS[R,          /# )	r%   rQ   útokenizers.TokenizerÚloopÚqueueÚskip_special_tokensÚtool_configNr‘   )ÚasyncioÚAbstractEventLoopr   r™   rO   )r)   r*   s   "€r!   r+   ÚDirectStreamer.__annotate__Ò  sY   ø€ ÷ '1ñ '1à)ð'1ñ ×'Ñ'ð'1ñ �}‰}ð	'1ñ
 "ð'1ñ ˜D•[ð'1ñ  �+ñ'1r    c                óÞ  € ^ RI Hp Wn        W n        W0n        V! . V4      V n        V'       d
   VR,          MRV n        V'       d
   VR,          MRV n        RV n        V'       d
   VR,          MRV n	        V'       d
   VR,          MRV n
        \        T;'       d    VP                  R4      4      V n        . V n        R	V n        \         P"                  ! 4       V n        ^ V n        . V n        R# )
a:  
Args:
    tokenizer: The Rust tokenizer (``tokenizer._tokenizer``).
    loop (`asyncio.AbstractEventLoop`): The event loop to push decoded text to.
    queue (`asyncio.Queue`): The queue that receives decoded text chunks.
    skip_special_tokens (`bool`, *optional*, defaults to `True`):
        Whether to strip special tokens during decoding.
    tool_config (`dict`, *optional*): Tool call config from ``get_tool_call_config``.
        When set, tokens between stc/etc delimiters (inclusive) are suppressed
        from the queue so tool call markup is never streamed to the client.
    reasoning_config (`dict`, *optional*): Thinking config from ``get_reasoning_config``.
        When set, tokens between start/end delimiters are wrapped as
        :class:`ReasoningText` so handlers route them to ``reasoning_content``.
©ÚDecodeStreamr^   Nr_   Fr„   r…   r†   T)Útokenizers.decodersr  Ú
_tokenizerÚ_loopÚ_queueÚ_decode_streamÚ_stc_idÚ_etc_idÚ_inside_tool_callr¥   r§   r™   rl   r¦   r¨   Ú_firstÚ	threadingÚEventÚ
_cancelledÚtotal_tokensÚgenerated_token_ids)r0   rQ   r  r  r  r  r‘   r  s   &&&&&&& r!   r1   ÚDirectStreamer.__init__Ò  sÄ   € õ. 	5à#ŒØŒ
ØŒÙ*¨2Ð/BÓCˆÔß0;�{ 8Ö,ÀˆŒß0;�{ 8Ö,ÀˆŒØ!&ˆÔßDTÐ#3°KÖ#@ÐZ^ˆÔ ß>NÐ 0°Ö :ÐTXˆÔÜ $Ð%5×%cÐ%cÐ:J×:NÑ:NÐObÓ:cÓ dˆÔØ+-ˆÔØˆŒÜ#Ÿ/š/Ó+ˆŒØˆÔØ.0ˆÖ r    c                ó"   <€ V ^8„  d   QhRRRR/# )r%   Úvalueútorch.TensorrM   Nr   )r)   r*   s   "€r!   r+   r  û  s   ø€ ÷ Jñ J˜ð J¨Dñ Jr    c                óä  € V P                   P                  4       '       d   \        4       hV P                  '       d
   RV n        R# VP	                  4        EF  pV ;P
                  ^,          un        V P                  P                  V4       W P                  8X  d	   RV n	        MW P                  8X  d   RV n	        \        W4      pV P                  P                  V P                  V4      pVe+   V P                  '       g   W P                  8X  g	   V'       d   KÈ  V P                  '       d   \!        V4      pV P"                  P%                  V P&                  P(                  V4       EK  	  R# )zHCalled by ``model.generate()`` after each decode step with new token(s).FNT)r  Úis_setr7   r  r›   r  r  r©   r  r  r  r¬   r  Ústepr  r¦   r:   r  Úcall_soon_threadsafer  Ú
put_nowait)r0   r  r£   Úis_start_or_end_tokenÚtexts   &&   r!   ÚputÚDirectStreamer.putû  s  € à�?‰?×!Ñ!×#Ò#Ü&Ó(Ð(à�;�;ˆ;ØˆDŒKÙØŸ™ŸˆHØ×Ò Õ"ÕØ×$Ñ$×+Ñ+¨HÔ5àŸ<™<Ô'Ø)-�Õ&ØŸ\™\Ô)Ø).�Ô&ä$;¸DÓ$KÐ!à×&Ñ&×+Ñ+¨D¯O©O¸XÓFˆDØŠ|˜t×5×5Ð5¸Ç\Á\Ô9Q×UjÙØ×$×$Ð$Ü$ TÓ*�Ø�J‰J×+Ñ+¨D¯K©K×,BÑ,BÀD×Ió! 'r    c                ó   <€ V ^8„  d   QhRR/# rË   r   )r)   r*   s   "€r!   r+   r    s   ø€ ÷ Fñ F�Tñ Fr    c                óf   € V P                   P                  V P                  P                  R4       R# )z;Called by ``model.generate()`` when generation is complete.N)r  r$  r  r%  ©r0   s   &r!   rx   ÚDirectStreamer.end  s    € à�
‰
×'Ñ'¨¯©×(>Ñ(>ÀÖEr    c                ó   <€ V ^8„  d   QhRR/# rË   r   )r)   r*   s   "€r!   r+   r    s   ø€ ÷ ñ ˜ñ r    c                ó:   € V P                   P                  4        R# )zWSignal cancellation. The next ``put()`` call will raise and abort ``model.generate()``.N)r  Úsetr,  s   &r!   ÚcancelÚDirectStreamer.cancel  s   € à�‰×ÑÖr    )r  r  r  r  r¦   r  r  r  r  r§   r¨   r¥   r  r  r  )TNN©r   r   r   r   r3   r1   r(  rx   r1  r   r4   r5   s   @r!   r  r  É  s7   ø‡ € ñ÷'1ò '1÷RJð J÷4Fð F÷ö r    r  c                   óp   a € ] tR tRt o RtRV 3R lR lltV 3R lR ltV 3R lR	 ltV 3R
 lR ltRt	V t
R# )Ú
CBStreameri  ap  Streamer for continuous batching (used by :class:`CBGenerateManager`).

Same ``put``/``end`` protocol as :class:`DirectStreamer`, but called manually
by :class:`CBGenerateManager` instead of by ``model.generate()``:
``put(output)`` receives a CB ``GenerationOutput``, decodes new tokens, and
pushes text to the asyncio.Queue. ``end()`` signals the stream is complete.
Nc                ó„   <€ V ^8„  d   QhRRRS[ RRRS[P                  RS[P                  RS[R	,          R
S[R	,          /# )r%   Ú
cb_managerr   Ú
request_idrQ   r  r  r  r  Nr‘   )r(   r	  r
  r   rO   )r)   r*   s   "€r!   r+   ÚCBStreamer.__annotate__'  sc   ø€ ÷ %1ñ %1à/ð%1ñ ð%1ð *ð	%1ñ
 ×'Ñ'ð%1ñ �}‰}ð%1ñ ˜D•[ð%1ñ  �+ñ%1r    c                óÂ  € ^ RI Hp Wn        W n        W@n        WPn        W0n        V! . R4      V n        V'       d
   VR,          MRV n        V'       d
   VR,          MRV n	        RV n
        V'       d
   VR,          MRV n        V'       d
   VR,          MRV n        \        T;'       d    VP                  R	4      4      V n        . V n        ^ V n        ^ V n        . V n        R# )
a  
Args:
    cb_manager (`ContinuousBatchingManager`): The CB manager instance.
    request_id (`str`): The request ID to track in the CB scheduler.
    tokenizer: The Rust tokenizer (``tokenizer._tokenizer``).
    loop (`asyncio.AbstractEventLoop`): The event loop to push decoded text to.
    queue (`asyncio.Queue`): The queue that receives decoded text chunks.
    tool_config (`dict`, *optional*): Tool call config (see ``DirectStreamer``).
    reasoning_config (`dict`, *optional*): Thinking config (see ``DirectStreamer``).
r  Tr^   Nr_   Fr„   r…   r†   )r  r  Ú_cbÚ_request_idr  r  r  r  r  r  r  r¥   r§   r™   rl   r¦   r¨   Ú	_prev_lenr  r  )	r0   r7  r8  rQ   r  r  r  r‘   r  s	   &&&&&&&& r!   r1   ÚCBStreamer.__init__'  sÀ   € õ( 	5àŒØ%ÔØŒ
ØŒØ#ŒÙ*¨2¨tÓ4ˆÔß0;�{ 8Ö,ÀˆŒß0;�{ 8Ö,ÀˆŒØ!&ˆÔßDTÐ#3°KÖ#@ÐZ^ˆÔ ß>NÐ 0°Ö :ÐTXˆÔÜ $Ð%5×%cÐ%cÐ:J×:NÑ:NÐObÓ:cÓ dˆÔØ+-ˆÔØˆŒØˆÔØ.0ˆÖ r    c                ó"   <€ V ^8„  d   QhRRRR/# )r%   Úoutputr   rM   Nr   )r)   r*   s   "€r!   r+   r9  N  s   ø€ ÷ )ñ )Ð,ð )°ñ )r    c                óz  € VP                   V P                  R p\        VP                   4      V n        V EF   pV ;P                  ^,          un        V P                  P                  V4       W0P                  8X  d	   RV n        MW0P                  8X  d   RV n        \        W4      pV P                  P                  V P                  V4      pVe+   V P                  '       g   W0P                  8X  g	   V'       d   KÈ  V P                  '       d   \        V4      pV P                  P!                  V4       EK  	  R# )zLDecode new tokens from a CB ``GenerationOutput`` and push text to the queue.NTF)Úgenerated_tokensr=  rŸ   r  r  r©   r  r  r  r¬   r  r#  r  r¦   r:   r  r%  )r0   r@  Ú
new_tokensr£   r&  r'  s   &&    r!   r(  ÚCBStreamer.putN  sì   € à×,Ñ,¨T¯^©^Ð-=Ð>ˆ
Ü˜V×4Ñ4Ó5ˆŒÜ"ˆHØ×Ò Õ"ÕØ×$Ñ$×+Ñ+¨HÔ5àŸ<™<Ô'Ø)-�Õ&ØŸ\™\Ô)Ø).�Ô&ä$;¸DÓ$KÐ!à×&Ñ&×+Ñ+¨D¯O©O¸XÓFˆDØŠ|˜t×5×5Ð5¸Ç\Á\Ô9Q×UjÙØ×$×$Ð$Ü$ TÓ*�Ø�K‰K×"Ñ" 4×(ó! #r    c                ó   <€ V ^8„  d   QhRR/# rË   r   )r)   r*   s   "€r!   r+   r9  d  s   ø€ ÷ %ñ %�Tñ %r    c                ó<   € V P                   P                  R4       R# )zSignal end of stream.N)r  r%  r,  s   &r!   rx   ÚCBStreamer.endd  s   € à�‰×Ñ˜tÖ$r    c                ó   <€ V ^8„  d   QhRR/# rË   r   )r)   r*   s   "€r!   r+   r9  h  s   ø€ ÷ 2ñ 2˜ñ 2r    c                óP   € V P                   P                  V P                  4       R# )zCancel the CB request.N)r;  Úcancel_requestr<  r,  s   &r!   r1  ÚCBStreamer.cancelh  s   € à�‰×Ñ × 0Ñ 0Ö1r    )r;  r  r  r¦   r  r  r=  r  r<  r  r§   r¨   r¥   r  r  r  ©NNr3  r5   s   @r!   r5  r5    s3   ø‡ € ñ÷%1ò %1÷N)ð )÷,%ð %÷2ö 2r    r5  c                ó(   € V ^8„  d   QhR\         RR/# )r%   ÚseedrM   Nrº   )r)   s   "r!   r+   r+   m  s   € ÷ ñ œð  ñ r    c                ó2   € ^ RI pVP                  ! V 4       R# )z8Set the PyTorch random seed for reproducible generation.N)ÚtorchÚmanual_seed)rN  rP  s   & r!   Úset_torch_seedrR  m  s   € ãà	×Ò�dÖr    c                ó   € V ^8„  d   QhRR/# rË   r   )r)   s   "r!   r+   r+   t  s   € ÷ !ñ !˜4ñ !r    c                 ó†   € ^ RI p V P                  P                  4       '       d   V P                  P                  4        R# R# )z+Empty the CUDA cache if a GPU is available.N)rP  ÚcudaÚis_availableÚempty_cache)rP  s    r!   Úreset_torch_cacherX  t  s-   € ãà‡z�z×Ñ× Ò Ø�
‰
×ÑÖ ñ !r    c                   ó`   a € ] tR tRt o RtR tV 3R lR ltV 3R lR ltV 3R lR	 ltR
t	V t
R# )ÚInferenceThreadi|  zÆPersistent thread for ``model.generate()`` calls.

``torch.compile`` with CUDA graphs stores state in thread-local storage.
All inference must run on the same thread to avoid corrupted graph state.
c                ó¦   € \        4       V n        \        P                  ! V P                  R R7      V n        V P
                  P                  4        R# )T)ÚtargetÚdaemonN)r   r  r  ÚThreadÚ_runÚ_threadrw   r,  s   &r!   r1   ÚInferenceThread.__init__ƒ  s3   € Ü"›WˆŒÜ ×'Ò'¨t¯y©yÀÔFˆŒØ�‰×ÑÖr    c                ó   <€ V ^8„  d   QhRR/# rË   r   )r)   r*   s   "€r!   r+   ÚInferenceThread.__annotate__ˆ  s   ø€ ÷ ,ñ ,�dñ ,r    c                ó\  €  V P                   P                  4       w  rr4p V! V/ VB pVe   VP                  VP                  V4       KJ  VP                  V4       K]    \         dC   pTe#   TP                  TP
                  T4        Rp?KŽ  TP                  T4        Rp?K¥  Rp?ii ; i)TN)r  rl   r$  Ú
set_resultÚ	ExceptionÚset_exception)r0   Úfnrí   rî   Úfuturer  ÚresultÚes   &       r!   r_  ÚInferenceThread._runˆ  sš   € ØØ-1¯[©[¯_©_Ó->Ñ*ˆB�f dð
,Ù˜TÐ, VÑ,�ØÒ#Ø×-Ñ-¨f×.?Ñ.?ÀÖHà×%Ñ% fÖ-øÜô ,ØÒ#Ø×-Ñ-¨f×.BÑ.BÀA×FÒFà×(Ñ(¨×+Ò+ûð	,ús#   ¡(A ÁA ÁB+Á) B&ÂB&Â&B+c                ó    <€ V ^8„  d   QhRS[ /# ©r%   rM   r   )r)   r*   s   "€r!   r+   rc  —  s   ø€ ÷ ñ ©Vñ r    c                óV   € \        4       pV P                  P                  WW4R34       V# )úESubmit a callable to the inference thread. Returns a blocking Future.N)r   r  r(  )r0   rh  rí   rî   ri  s   &&*, r!   ÚsubmitÚInferenceThread.submit—  s%   € ä›ˆØ�‰�‰˜ 6°4Ð8Ô9Øˆr    c                ó4   <€ V ^8„  d   QhRS[ P                  /# rn  )r	  r   )r)   r*   s   "€r!   r+   rc  �  s   ø€ ÷ ñ ±7·>±>ñ r    c                óŒ   € \         P                  ! 4       pVP                  4       pV P                  P	                  WW5V34       V# ©zOSubmit a callable to the inference thread. Returns an awaitable asyncio.Future.)r	  Úget_running_loopÚcreate_futurer  r(  )r0   rh  rí   rî   r  ri  s   &&*,  r!   Úasync_submitÚInferenceThread.async_submit�  s:   € ä×'Ò'Ó)ˆØ×#Ñ#Ó%ˆØ�‰�‰˜ 6°4Ð8Ô9Øˆr    )r  r`  N)r   r   r   r   r3   r1   r_  rq  rx  r   r4   r5   s   @r!   rZ  rZ  |  s-   ø‡ € ñò÷
,ð ,÷ð ÷ö r    rZ  c                   óŽ   a € ] tR tRt o RtV 3R lR lt]RV 3R lR ll4       t]V 3R lR	 l4       t]V 3R
 lR l4       t	Rt
V tR# )ÚBaseGenerateManageri¥  uÓ   Base class for generation managers.

Subclasses:
- :class:`GenerateManager` â€” sequential ``model.generate()`` on a persistent thread.
- :class:`CBGenerateManager` â€” continuous batching with paged attention.
c                ó&   <€ V ^8„  d   QhRRRRRR/# ©r%   rL   r   Ú
gen_configr   rM   Nr   )r)   r*   s   "€r!   r+   Ú BaseGenerateManager.__annotate__­  s*   ø€ ÷ Iñ IÐ.ð IÐ<Nð IÐSWñ Ir    c                ó   € R# )z:Initialize continuous batching. No-op for non-CB managers.Nr   ©r0   rL   r~  s   &&&r!   Úinit_cbÚBaseGenerateManager.init_cb­  ó   ‚ r    Nc                óˆ   <€ V ^8„  d   QhRRRRRS[ RRRS[R	S[ R
,          RS[ R
,          RS[S[P                  R3,          /# )r%   rL   r   re   ú(ProcessorMixin | PreTrainedTokenizerFastÚinputsr~  r   r8  r  Nr‘   rM   zDirectStreamer | CBStreamer)rO   r(   r’   r	  r   )r)   r*   s   "€r!   r+   r  ±  sr   ø€ ÷ ñ à ðð >ðñ ð	ð
 'ðñ ðñ ˜D•[ðñ  �+ðñ 
‰w�}‰}Ð;Ð;Õ	<ñr    c                ó   € R# )aj  Start streaming generation.

Args:
    model (`PreTrainedModel`): The loaded model.
    processor: The processor or tokenizer for decoding.
    inputs (`dict`): Tokenized inputs (tensors for sequential, lists for CB).
    gen_config (`GenerationConfig`): Generation parameters.
    request_id (`str`): Unique request identifier.
    tool_config (`dict`, *optional*): Tool call config from ``get_tool_call_config``.
        When set, tool call tokens (between stc/etc) are suppressed from output.
    reasoning_config (`dict`, *optional*): Thinking config from ``get_reasoning_config``.
        When set, thinking tokens are wrapped as :class:`ReasoningText`.

Returns:
    `tuple[asyncio.Queue, DirectStreamer | CBStreamer]`: A ``(queue, streamer)`` pair
    where *queue* yields ``str | _StreamError | None`` and *streamer* exposes
    ``.total_tokens`` and ``.cancel()``.
Nr   )r0   rL   re   r‡  r~  r8  r  r‘   s   &&&&&&&&r!   Úgenerate_streamingÚ&BaseGenerateManager.generate_streaming°  r„  r    c                ób   <€ V ^8„  d   QhRRRRRS[ RRRS[R	S[S[S[S[S[,          3,          /# ©
r%   rL   r   re   r†  r‡  r~  r   r8  rM   ©rO   r(   r’   r˜   rq   )r)   r*   s   "€r!   r+   r  Ï  sW   ø€ ÷ ñ à ðð >ðñ ð	ð
 'ðñ ðñ 
‰s‘C™™c�Ð"Õ	#ñr    c              ƒ  ó   "  € R# 5i)a�  Run generation to completion.

Args:
    model (`PreTrainedModel`): The loaded model.
    processor: The processor or tokenizer for decoding.
    inputs (`dict`): Tokenized inputs (tensors for sequential, lists for CB).
    gen_config (`GenerationConfig`): Generation parameters.
    request_id (`str`): Unique request identifier.

Returns:
    `tuple[str, int, list[int]]`: ``(text, input_len, generated_ids)``.
Nr   )r0   rL   re   r‡  r~  r8  s   &&&&&&r!   Úgenerate_non_streamingÚ*BaseGenerateManager.generate_non_streamingÎ  s   é ‚ ùs   ‚c                ó   <€ V ^8„  d   QhRR/# rË   r   )r)   r*   s   "€r!   r+   r  å  s   ø€ ÷ >ñ >�dñ >r    c                ó   € R# )z/Stop the generation manager and free resources.Nr   r,  s   &r!   ÚstopÚBaseGenerateManager.stopä  r„  r    r   rL  )r   r   r   r   r3   r‚  r   r‰  r�  r“  r   r4   r5   s   @r!   r{  r{  ¥  sW   ø‡ € ñ÷Ið Ið ÷ñ ó ðð: ÷ó ðð* ÷>ó ö>r    r{  c                   óˆ   a € ] tR tRt o RtR tRV 3R lR lltV 3R lR ltV 3R	 lR
 ltV 3R lR lt	V 3R lR lt
RtV tR# )ÚGenerateManagerié  zFSequential generation via ``model.generate()`` on a persistent thread.c                ó$   € \        4       V n        R # r.   )rZ  r`  r,  s   &r!   r1   ÚGenerateManager.__init__ì  s   € Ü&Ó(ˆŽr    Nc                óŠ   <€ V ^8„  d   QhRRRRRS[ RRRS[R	S[ R
,          RS[ R
,          RS[S[P                  S[3,          /# ©r%   rL   r   re   r†  r‡  r~  r   r8  r  Nr‘   rM   )rO   r(   r’   r	  r   r  )r)   r*   s   "€r!   r+   ÚGenerateManager.__annotate__ï  sq   ø€ ÷ ñ à ðð >ðñ ð	ð
 'ðñ ðñ ˜D•[ðñ  �+ðñ 
‰w�}‰}™nÐ,Õ	-ñr    c                ó2  aaaa€ \         P                  ! 4       o\         P                  ! 4       o\        VRV4      P                  p\        VSSWgR7      p	/ VCRV	RVRV/Co\        SR4      '       d   RSR&   R VVVV3R	 llp
V P                  V
4       SV	3# )
zLStart streaming generation via ``model.generate()`` on the inference thread.rQ   ©r  r‘   rª   Úgeneration_configÚ
has_talkerr'  Úgeneration_modec                ó   € V ^8„  d   QhRR/# rË   r   )r)   s   "r!   r+   Ú8GenerateManager.generate_streaming.<locals>.__annotate__  s   € ÷ 	Rñ 	R�dñ 	Rr    c            	      ó  <€  SP                   ! R/ SB  R #   \         d!    SP                  SP                  R 4        R # \         d:   p SP                  SP                  \        \        T 4      4      4        R p ? R # R p ? ii ; i)Nr   )Úgenerater7   r$  r%  rf  r#   r(   )rk  Ú
gen_kwargsr  rL   r  s    €€€€r!   r_  Ú0GenerateManager.generate_streaming.<locals>._run  sk   ø€ ðRØ—’Ñ, Ô,øÜ'ô BØ×)Ñ)¨%×*:Ñ*:¸D×AÜô RØ×)Ñ)¨%×*:Ñ*:¼LÌÈQËÓ<P×QÒQûðRús!   ƒ —'BÁBÁ
BÁ.A?Á?B)r	  rv  r   r`   r  r  rž   rq  )r0   rL   re   r‡  r~  r8  r  r‘   Úrust_tokenizerrª   r_  r¥  r  r  s   &f&&&&&&   @@@r!   r‰  Ú"GenerateManager.generate_streamingï  s    û€ ô ×'Ò'Ó)ˆÜ&Ÿ}š}›ˆä  ¨K¸ÓC×NÑNˆÜ!Ø˜D %°[ô
ˆð o˜Ðn 
¨HÐ6IÈ:ÐWbÐdmÑnˆ
Ü�5˜,×'Ò'Ø,2ˆJÐ(Ñ)÷	Ró 	Rð 	�‰�DÔØ�hˆÐr    c                óP   <€ V ^8„  d   QhRRRRRS[ RRRS[R	S[S[S[R
3,          /# )r%   rL   r   re   r†  r‡  r~  r   r8  rM   r   )rO   r(   r’   r˜   )r)   r*   s   "€r!   r+   r›    sS   ø€ ÷ .ñ .à ð.ð >ð.ñ ð	.ð
 'ð.ñ ð.ñ 
‰s‘C˜Ð'Õ	(ñ.r    c              ƒ  ó  "  € / VCRVRV/Cp\        VR4      '       d   RVR&   V P                  ! VP                  3/ VB G Rj  x€L
 pVR,          P                  R
,          pV^ VR13,          p	VP	                  V	RR	7      p
W¨V	3#  LB5i)zNRun generation to completion via ``model.generate()`` on the inference thread.rž  rQ   rŸ  r'  r   NrŒ   T©r  rç   )rž   rx  r¤  ÚshapeÚdecode)r0   rL   re   r‡  r~  r8  Úgenerate_kwargsÚ	sequencesÚ	input_lenrt   r'  s   &&&&&&     r!   r�  Ú&GenerateManager.generate_non_streaming  sž   é € ð ^˜VÐ]Ð%8¸*ÀkÐS\Ñ]ˆÜ�5˜,×'Ò'Ø17ˆOÐ-Ñ.Ø×+Ò+¨E¯N©NÑN¸oÑN×Nˆ	Ø˜;Õ'×-Ñ-¨bÕ1ˆ	Ø! ! Y¡Z -Õ0ˆØ×Ñ À4ÐÓHˆØ Ð-Ð-ñ	 Oùs   ‚AB	ÁBÁAB	c                ó&   <€ V ^8„  d   QhRS[ RS[/# ©r%   rh  rM   )r   r   )r)   r*   s   "€r!   r+   r›  $  s   ø€ ÷ 8ñ 8™ð 8±vñ 8r    c                óB   € V P                   P                  ! V.VO5/ VB # )rp  )r`  rq  ©r0   rh  rí   rî   s   &&*,r!   rq  ÚGenerateManager.submit$  s!   € à�|‰|×"Ò" 2Ð7¨Ò7°Ñ7Ð7r    c                ó:   <€ V ^8„  d   QhRS[ RS[P                  /# r³  )r   r	  r   )r)   r*   s   "€r!   r+   r›  (  s   ø€ ÷ >ñ >™xð >¹W¿^¹^ñ >r    c                óB   € V P                   P                  ! V.VO5/ VB # ru  )r`  rx  rµ  s   &&*,r!   rx  ÚGenerateManager.async_submit(  s!   € à�|‰|×(Ò(¨Ð=¨dÒ=°fÑ=Ð=r    c                ó   <€ V ^8„  d   QhRR/# rË   r   )r)   r*   s   "€r!   r+   r›  ,  s   ø€ ÷ ñ �dñ r    c                ó   € R # r.   r   r,  s   &r!   r“  ÚGenerateManager.stop,  rÉ   r    )r`  rL  )r   r   r   r   r3   r1   r‰  r�  rq  rx  r“  r   r4   r5   s   @r!   r–  r–  é  s@   ø‡ € ÙPò)÷ò ÷B.ð .÷(8ð 8÷>ð >÷ö r    r–  c                   óÆ   a € ] tR tRt o RtRV 3R lR lltV 3R lR ltV 3R lR	 ltV 3R
 lR ltRV 3R lR llt	V 3R lR lt
]V 3R lR l4       tV 3R lR ltRtV tR# )ÚCBGenerateManageri0  aQ  Continuous batching generation via paged attention.

Translates between the handler's text-level asyncio.Queue and CB's
token-level interface. Per-request: ``max_new_tokens``, ``eos_token_id``.

The CB manager is initialized lazily on the first request via
:meth:`ensure_initialized`, using that request's ``gen_config`` for shared
sampling params (temperature, top_p, do_sample).

.. todo:: Remove :meth:`init_cb` when CB supports per-request
   generation config. At that point, ``gen_config`` can be passed directly
   to ``add_request`` and the CB manager no longer needs a shared config.
Nc                ó   <€ V ^8„  d   QhRR/# )r%   Ú	cb_configúContinuousBatchingConfig | Noner   )r)   r*   s   "€r!   r+   ÚCBGenerateManager.__annotate__?  s   ø€ ÷ $ñ $Ð"Cñ $r    c                ó    € R V n         Wn        R # r.   ©r;  Ú
_cb_config)r0   rÀ  s   &&r!   r1   ÚCBGenerateManager.__init__?  s   € Ø59ˆŒØ#Žr    c                ó&   <€ V ^8„  d   QhRRRRRR/# r}  r   )r)   r*   s   "€r!   r+   rÂ  C  s%   ø€ ÷ ñ Ð.ð Ð<Nð ÐSWñ r    c                óœ   € V P                   e   R# VP                  W P                  R7      V n         V P                   P                  4        R# )aL  Initialize the CB manager on first call with the request's generation config.

.. todo:: Remove when CB supports per-request generation config.

Args:
    model (`PreTrainedModel`): The loaded model (must support ``init_continuous_batching``).
    gen_config (`GenerationConfig`): Generation config used for shared sampling params.
N)rž  Úcontinuous_batching_config)r;  Úinit_continuous_batchingrÅ  rw   r�  s   &&&r!   r‚  ÚCBGenerateManager.init_cbC  s?   € ð �8‰8ÒÙà×1Ñ1Ø(Ç_Á_ð 2ó 
ˆŒð 	�‰�‰Ör    c                ó    <€ V ^8„  d   QhRS[ /# rn  ©r™   )r)   r*   s   "€r!   r+   rÂ  T  s   ø€ ÷ @ñ @™$ñ @r    c                ó^   € V P                   RJ ;'       g    V P                   P                  RJ # )zJWhether the CB worker is healthy. ``True`` before ``init_cb()`` is called.N)r;  Úfatal_errorr,  s   &r!   Úis_aliveÚCBGenerateManager.is_aliveT  s(   € à�x‰x˜4Ð×?Ð? 4§8¡8×#7Ñ#7¸4Ð#?Ð?r    c                ó$   <€ V ^8„  d   QhRS[ RR/# )r%   r8  rM   Nr'   )r)   r*   s   "€r!   r+   rÂ  X  s   ø€ ÷ 	ñ 	¡sð 	¨tñ 	r    c                ó    € V P                   e@   V P                   P                  e&   \        RV RV P                   P                   24      hR# R# )uÑ   Raise :class:`CBWorkerDeadError` if the CB worker has died.

Called at request entry to fail fast â€” submitting to a dead worker would otherwise
enqueue the request into a void where it never gets processed.
Nz,CB worker is dead and cannot accept request ú: )r;  rÏ  r=   )r0   r8  s   &&r!   Ú_check_aliveÚCBGenerateManager._check_aliveX  sN   € ð �8‰8Ò D§H¡H×$8Ñ$8Ò$DÜ#Ø>¸z¸lÈ"ÈTÏXÉX×MaÑMaÐLbÐcóð ñ %EÑr    c                óŠ   <€ V ^8„  d   QhRRRRRS[ RRRS[R	S[ R
,          RS[ R
,          RS[S[P                  S[3,          /# rš  )rO   r(   r’   r	  r   r5  )r)   r*   s   "€r!   r+   rÂ  c  sq   ø€ ÷ 8$ñ 8$à ð8$ð >ð8$ñ ð	8$ð
 'ð8$ñ ð8$ñ ˜D•[ð8$ñ  �+ð8$ñ 
‰w�}‰}™jÐ(Õ	)ñ8$r    c           
     ó¶  aa€ V P                   pVf   \        R4      hV P                  V4       \        P                  ! 4       p	\        P
                  ! 4       oVR,          p
VP                  V
VRVP                  VP                  R7      p\        VRV4      P                  p\        V P                   VVV	SVVR7      oVV3R lpVP                  W\4       SS3# )zFStart streaming CB generation. Registers a per-request output handler.ú3CB manager not initialized. Call `init_cb()` first.rŒ   T)r8  Ú	streamingÚmax_new_tokensÚeos_token_idrQ   r�  c                 ó|  <€  SP                  V 4       V P                  e7   SP                  \        V P                  4      4       SP	                  4        R # V P                  4       '       d   SP	                  4        R # R #   \         d/   pSP                  \        \        T4      4      4        R p?R # R p?ii ; ir.   )r(  Úerrorr%  r#   rx   Úis_finishedrf  r(   )r@  rk  rª   Ú
text_queues   & €€r!   Ú
_on_outputÚ8CBGenerateManager.generate_streaming.<locals>._on_outputŒ  s‡   ø€ ð<Ø—‘˜VÔ$ð —<‘<Ò+Ø×)Ñ)¬,°v·|±|Ó*DÔEØ—L‘L–NØ×'Ñ'×)Ò)Ø—L‘L–Nñ *øäô <Ø×%Ñ%¤l´3°q³6Ó&:×;Ò;ûð<ús$   ƒAB ÁB Á.B ÂB;Â#B6Â6B;)r;  ÚRuntimeErrorrÕ  r	  rv  r   Úadd_requestrÛ  rÜ  r`   r  r5  Úregister_result_handler)r0   rL   re   r‡  r~  r8  r  r‘   Úcbr  rŒ   r§  rá  rª   rà  s   &&&&&&&&     @@r!   r‰  Ú$CBGenerateManager.generate_streamingc  s×   ù€ ð �X‰XˆØŠ:ÜÐTÓUÐUØ×Ñ˜*Ô%ä×'Ò'Ó)ˆÜ$+§M¢M£Oˆ
à˜;Õ'ˆ	Ø—^‘^ØØ!ØØ%×4Ñ4Ø#×0Ñ0ð $ó 
ˆ
ô ! ¨K¸ÓC×NÑNˆÜØ�H‰HØØØØØ#Ø-ô
ˆö	<ð 	×"Ñ" :Ô:Ø˜8Ð#Ð#r    c                ób   <€ V ^8„  d   QhRRRRRS[ RRRS[R	S[S[S[S[S[,          3,          /# rŒ  r�  )r)   r*   s   "€r!   r+   rÂ  �  sW   ø€ ÷ /.ñ /.à ð/.ð >ð/.ñ ð	/.ð
 'ð/.ñ ð/.ñ 
‰s‘C™™c�Ð"Õ	#ñ/.r    c              ƒ  óZ  a"  € V P                   pVf   \        R4      hV P                  V4       VR,          p\        V4      p\        P
                  ! 4       p	V	P                  4       oV3R lp
VP                  WZ4       VP                  VVVP                  RVP                  R7       SG Rj  x€L
 pVP                  eE   VP                  e   \        RV RVP                   24      h\        R	V RVP                   24      hVP                  pVP                  VR
R7      pWØV3#  Ly5i)zcRun non-streaming CB generation. Registers a handler that resolves an asyncio.Future on completion.NrÙ  rŒ   c                 óZ   <€ SP                  4       '       g   SP                  V 4       R # R # r.   )Údonere  )rj  ri  s   &€r!   Ú
_on_resultÚ<CBGenerateManager.generate_non_streaming.<locals>._on_result²  s!   ø€ Ø—;‘;—=’=Ø×!Ñ! &Ö)ñ !r    F)r8  rÛ  rÚ  rÜ  zCB worker died during request rÔ  zCB generation failed for Tr«  )r;  rã  rÕ  rŸ   r	  rv  rw  rå  rä  rÛ  rÜ  rÞ  rÏ  r=   rB  r­  )r0   rL   re   r‡  r~  r8  ræ  rŒ   r°  r  rì  rj  rt   r'  ri  s   &&&&&&        @r!   r�  Ú(CBGenerateManager.generate_non_streaming�  s0  øé € ð �X‰XˆØŠ:ÜÐTÓUÐUØ×Ñ˜*Ô%à˜;Õ'ˆ	Ü˜	“Nˆ	ô ×'Ò'Ó)ˆØ×#Ñ#Ó%ˆõ	*ð 	×"Ñ" :Ô:à
�‰ØØ!Ø%×4Ñ4ØØ#×0Ñ0ð 	ô 	
ð —ˆð �<‰<Ò#Ø�~‰~Ò)Ü'Ð*HÈÈÐTVÐW]×WcÑWcÐVdÐ(eÓfÐfÜÐ!:¸:¸,ÀbÈÏÉÈÐWÓXÐXØ×/Ñ/ˆØ×Ñ À4ÐÓHˆØ Ð-Ð-ñ ùs   ƒB,D+Â/D)Â0A:D+c                ó   <€ V ^8„  d   QhRR/# )r%   rM   r   r   )r)   r*   s   "€r!   r+   rÂ  Ï  s   ø€ ÷ 2ñ 2˜;ñ 2r    c                ó¤   € V P                   e   V P                   P                  f   \        R4      hV P                   P                  P                  # )z*The CB scheduler (for testing/monitoring).z.Continuous batching processor not initialized.)r;  Úbatch_processorrã  Ú	schedulerr,  s   &r!   rò  ÚCBGenerateManager.schedulerÎ  s?   € ð �8‰8Ò˜tŸx™x×7Ñ7Ò?ÜÐOÓPÐPØ�x‰x×'Ñ'×1Ñ1Ð1r    c                ó   <€ V ^8„  d   QhRR/# rË   r   )r)   r*   s   "€r!   r+   rÂ  Õ  s   ø€ ÷ 1ñ 1�dñ 1r    c                ó`   € V P                   e    V P                   P                  R^R7       R # R # )NT)ÚblockÚtimeout)r;  r“  r,  s   &r!   r“  ÚCBGenerateManager.stopÕ  s%   € Ø�8‰8ÒØ�H‰H�M‰M ¨aˆMÖ0ñ  r    rÄ  r.   rL  )r   r   r   r   r3   r1   r‚  rÐ  rÕ  r‰  r�  Úpropertyrò  r“  r   r4   r5   s   @r!   r¾  r¾  0  sh   ø‡ € ñ÷$ò $÷ð ÷"@ð @÷	ð 	÷8$ò 8$÷t/.ð /.ðb ÷2ó ð2÷1ö 1r    r¾  c                   ó†   a € ] tR tRt o RtRV 3R lR lltV 3R lR ltRV 3R lR	 lltV 3R
 lR ltV 3R lR lt	Rt
V tR# )ÚGenerationStateiÚ  a  Shared generation state across all handlers.

Manages per-model :class:`GenerateManager` instances (each with its own
:class:`InferenceThread` so different models can run concurrently while
``torch.compile`` / CUDA graphs require same-model-same-thread) and a
single :class:`CBGenerateManager` for continuous batching.

Args:
    continuous_batching (`bool`, *optional*, defaults to `False`):
        Whether to use continuous batching with paged attention instead of
        sequential ``model.generate()`` calls.
Nc                ó*   <€ V ^8„  d   QhRS[ RS[ RR/# )r%   Úcontinuous_batchingÚcompilerÀ  rÁ  rÍ  )r)   r*   s   "€r!   r+   ÚGenerationState.__annotate__è  s)   ø€ ÷ -ñ -á!ð-ñ ð-ð 5ñ	-r    c                óT   € Wn         W n        W0n        / V n        R V n        R V n        R # r.   )Ú_continuous_batchingÚ_compilerÅ  Ú_generate_managersÚ_cb_managerÚ_cb_model_id)r0   rý  rþ  rÀ  s   &&&&r!   r1   ÚGenerationState.__init__è  s,   € ð %8Ô!ØŒØ#ŒØ>@ˆÔØ59ˆÔØ(,ˆÖr    c                ó*   <€ V ^8„  d   QhRRRS[ RS[/# )r%   rL   r   ÚmodalityrM   )r   r™   )r)   r*   s   "€r!   r+   rÿ  õ  s$   ø€ ÷ ñ Ð->ð É(ð ÑW[ñ r    c                óä   € V P                   '       g   R# \        VR4      ;'       d    V\        P                  8H  pV'       g-   \        P                  VP                  P                   R24       V# )a'  Check if continuous batching can be used for this model and modality.

Args:
    model (`PreTrainedModel`): The loaded model.
    modality (`Modality`): The detected model modality (LLM, VLM, etc.).

Returns:
    `bool`: ``True`` if CB is enabled and the model supports it, ``False`` otherwise.
FrÊ  zM does not support continuous batching. Falling back to sequential generation.)r  rž   r   r   ÚloggerÚwarning_oncerï   r   )r0   rL   r  Úcans   &&& r!   Úuse_continuous_batchingÚ'GenerationState.use_continuous_batchingõ  sc   € ð ×(×(Ð(ÙÜ�eÐ7Ó8×UÐU¸XÌÏÉÑ=UˆßÜ×ÑØ—?‘?×+Ñ+Ð,ð -9ð 9ôð ˆ
r    c                ó,   <€ V ^8„  d   QhRS[ RS[RS[/# )r%   r±   Úuse_cbrM   )r(   r™   r{  )r)   r*   s   "€r!   r+   rÿ  	  s#   ø€ ÷ 1ñ 1¡Cð 1±ð 1ÑBUñ 1r    c                ó|  € V'       d|   V P                   V8w  d0   V P                  e"   V P                  P                  4        RV n        V P                  f"   \        V P                  R7      V n        Wn         V P                  # WP
                  9  d   \        4       V P
                  V&   V P
                  V,          # )a6  Return a per-model generation manager, lazily created on first request.

Args:
    model_id (`str`): The model ID in ``'model_id@revision'`` format.
    use_cb (`bool`): Whether to return a CB manager or a sequential one.

Returns:
    `BaseGenerateManager`: Either a `GenerateManager` or `CBGenerateManager`.
N)rÀ  )r  r  r“  r¾  rÅ  r  r–  )r0   r±   r  s   &&&r!   Úget_managerÚGenerationState.get_manager	  sš   € ÷ Ø× Ñ  HÔ,Ø×#Ñ#Ò/Ø×$Ñ$×)Ñ)Ô+Ø'+�DÔ$Ø×ÑÒ'Ü#4¸t¿¹Ô#O�Ô Ø$,Ô!Ø×#Ñ#Ð#Ø×2Ñ2Ô2Ü0?Ó0AˆD×#Ñ# HÑ-Ø×&Ñ& xÕ0Ð0r    c                ó   <€ V ^8„  d   QhRR/# rË   r   )r)   r*   s   "€r!   r+   rÿ     s   ø€ ÷ $ñ $˜$ñ $r    c                óh   € V P                   e$   V P                   P                  4        RV n         R# R# )z$Stop any active generation managers.N)r  r“  r,  s   &r!   ÚshutdownÚGenerationState.shutdown   s-   € à×ÑÒ'Ø×Ñ×!Ñ!Ô#Ø#ˆDÖñ (r    c                ó    <€ V ^8„  d   QhRS[ /# rn  rÍ  )r)   r*   s   "€r!   r+   rÿ  &  s   ø€ ÷ Gñ G™Tñ Gr    c                ób   € V P                   RJ ;'       g    V P                   P                  4       # )zTWhether the CB worker is healthy. ``True`` if CB is disabled or not yet initialized.N)r  rÐ  r,  s   &r!   Úis_cb_aliveÚGenerationState.is_cb_alive&  s*   € à×Ñ 4Ð'×FÐF¨4×+;Ñ+;×+DÑ+DÓ+FÐFr    )rÅ  r  r  r  r  r  )FFN©F)r   r   r   r   r3   r1   r  r  r  r  r   r4   r5   s   @r!   rû  rû  Ú  s>   ø‡ € ñ÷-ò -÷ð ÷(1ò 1÷.$ð $÷Gö Gr    rû  c                   óÊ   a € ] tR tRt o RtRt]! 4       tRV 3R lR lltV 3R lR lt	]
V 3R lR	 l4       tV 3R
 lR ltRV 3R lR llt]
V 3R lR l4       tV 3R ltRtV tR# )ÚBaseHandleri+  aŽ  Shared logic for chat completion and responses handlers.

Provides model resolution, generation config building, and SSE formatting.
Generation is delegated to the shared :class:`GenerationState`.

Args:
    model_manager (`ModelManager`):
        Handles model loading, caching, and lifecycle.
    generation_state (`GenerationState`):
        Shared state managing per-model generation managers.
Nc                ó8   <€ V ^8„  d   QhRRRS[ RS[R,          /# )r%   Úmodel_managerr   Úgeneration_stateÚchat_template_kwargsN)rû  rO   )r)   r*   s   "€r!   r+   ÚBaseHandler.__annotate__;  s-   ø€ ÷ ?ñ ?à%ð?ñ *ð?ñ # T�kñ	?r    c                ó@   € Wn         W n        T;'       g    / V n        R # r.   )r   r!  r"  )r0   r   r!  r"  s   &&&&r!   r1   ÚBaseHandler.__init__;  s    € ð +ÔØ 0ÔØ$8×$>Ð$>¸BˆÖ!r    c                ó$   <€ V ^8„  d   QhRS[ RR/# )r%   ÚbodyrM   NrN   )r)   r*   s   "€r!   r+   r#  E  s   ø€ ÷ Yñ Y¡dð Y¨tñ Yr    c                ó>  € ^ RI Hp \        VP                  4       4      pV P                  e<   V\        V P                  R\        4       4      ,
          pV'       d   V! RRV 2R7      hW0P                  ,          pV'       d   \        P                  RV 24       R# R# )zMValidate request fields against the handler's params class and unused fields.©ÚHTTPExceptionNÚ__mutable_keys__i¦  z"Unexpected fields in the request: ©Ústatus_codeÚdetailz,Ignoring unsupported fields in the request: )	Úfastapir*  r0  ÚkeysÚ_valid_params_classr`   Ú_unused_fieldsr
  r  )r0   r'  r*  Ú
input_keysÚ
unexpectedÚunuseds   &&    r!   Ú_validate_requestÚBaseHandler._validate_requestE  s…   € å)ä˜Ÿ™›Ó%ˆ
Ø×#Ñ#Ò/Ø#¤g¨d×.FÑ.FÐHZÔ\_Ó\aÓ&bÕbˆJßÙ#°Ð>`ÐakÐ`lÐ<mÔnÐnØ×1Ñ1Õ1ˆßÜ×ÑÐ"NÈvÈhÐ WÖXñ r    c                ó$   <€ V ^8„  d   QhRRRS[ /# )r%   Úchunkzstr | pydantic.BaseModelrM   r'   )r)   r*   s   "€r!   r+   r#  S  s    ø€ ÷ Gñ GÐ6ð G¹3ñ Gr    c                ó˜   € \        V \        4      '       d    V P                  R4      '       d   V # RV  R2# RV P                  RR7       R2# )z;Format a pydantic model or string as an SSE ``data:`` line.zdata: z

T)Úexclude_none)rm   r(   Ú
startswithÚmodel_dump_json)r9  s   &r!   Úchunk_to_sseÚBaseHandler.chunk_to_sseR  sS   € ô �eœS×!Ò!Ø!×,Ñ,¨X×6Ò6�5ÐP¸fÀUÀGÈ4Ð<PÐPØ˜×-Ñ-¸4Ð-Ó@ÐAÀÐFÐFr    c                ó<   <€ V ^8„  d   QhRS[ RS[S[RR3,          /# )r%   r'  rM   r   r†  )rO   r’   r(   )r)   r*   s   "€r!   r+   r#  Y  s)   ø€ ÷ *ñ *¡4ð *©E±#Ð7HÐJtÐ2tÕ,uñ *r    c                óž  € ^ RI Hp V P                  P                  en   VP	                  R4      pVe@   W0P                  P                  8w  d&   V! RRV P                  P                   RV R2R7      hV P                  P                  VR&   V P                  P                  VR,          4      pV P                  P                  V4      w  rVWEV3# )zVApply force_model, load model + processor.

Returns ``(model_id, model, processor)``.
r)  rL   i�  zServer is pinned to 'z'; requested 'z'.r,  )r/  r*  r   Úforce_modelrl   Úprocess_model_nameÚload_model_and_processor)r0   r'  r*  Ú	requestedr±   rL   re   s   &&     r!   Ú_resolve_modelÚBaseHandler._resolve_modelY  sÆ   € õ
 	*à×Ñ×)Ñ)Ò5ØŸ™ Ó)ˆIØÒ$¨×6HÑ6H×6TÑ6TÔ)TÙ#Ø #Ø3°D×4FÑ4F×4RÑ4RÐ3SÐSaÐbkÐalÐlnÐoôð ð !×.Ñ.×:Ñ:ˆD�‰Mà×%Ñ%×8Ñ8¸¸g½ÓGˆØ×-Ñ-×FÑFÀxÓPÑˆà 	Ð)Ð)r    c                ó.   <€ V ^8„  d   QhRS[ RRRS[RR/# )r%   r'  Úmodel_generation_configr   r  rM   )rO   r™   )r)   r*   s   "€r!   r+   r#  n  s-   ø€ ÷ 0!ñ 0!Ùð0!Ø3Eð0!ÙOSð0!à	ñ0!r    c                ó¦  € ^ RI Hp VP                  R4      e%   V! R
/ \        P                  ! VR,          4      B pM<\
        P                  ! V4      pVP                  e   VP                  R8  d   RVn        VP                  R4      e6   \        VR,          4      Vn	        \        VR,          4      R8X  d   RVn
        VP                  R4      e   \        VR,          4      Vn        VP                  R4      e   \        VR,          4       V P                  P                  '       d   VP                  f   R	Vn        V'       d   RVn        V# )a¿  Build a GenerationConfig from shared params (temperature, top_p, seed, generation_config JSON).

Subclasses should call ``super()._build_generation_config(...)`` then apply
endpoint-specific params (``max_tokens``, ``max_output_tokens``, etc.).

Args:
    body (`dict`):
        The raw request body.
    model_generation_config (`GenerationConfig`):
        The model's default generation config (will be deep-copied).
    use_cb (`bool`, *optional*, defaults to `False`):
        Whether continuous batching is active. If ``True``, disables the model's
        internal KV cache (CB manages its own paged cache).

Returns:
    `GenerationConfig`: A new config with request-specific overrides applied.
)r   rž  i   Útemperatureg        FÚtop_prN  Ústaticr   )Útransformersr   rl   rE   ÚloadsÚcopyÚdeepcopyrÛ  ÚfloatrK  Ú	do_samplerL  rR  r!  r  Úcache_implementationÚ	use_cache)r0   r'  rI  r  r   rž  s   &&&&  r!   Ú_build_generation_configÚ$BaseHandler._build_generation_confign  s   € õ( 	2à�8‰8Ð'Ó(Ò4Ù 0Ñ Y´4·:²:¸dÐCVÕ>WÓ3XÑ YÑä $§¢Ð.EÓ FÐØ ×/Ñ/Ò7Ð;L×;[Ñ;[Ð^bÔ;bØ37Ð!Ô0à�8‰8�MÓ"Ò.Ü,1°$°}Õ2EÓ,FÐÔ)Ü�T˜-Õ(Ó)¨SÔ0Ø.3Ð!Ô+Ø�8‰8�GÓÒ(Ü&+¨D°­MÓ&:ÐÔ#Ø�8‰8�FÓÒ'Ü˜4 �<Ô(ð × Ñ ×)×)Ð)Ð.?×.TÑ.TÒ.\Ø5=ÐÔ2÷ Ø*/ÐÔ'ð !Ð r    c                óL   <€ V ^8„  d   QhRS[ S[,          RS[RS[ S[,          /# )r%   Úmessagesr  rM   )rq   rO   r   )r)   r*   s   "€r!   r+   r#  ¡  s2   ø€ ÷ C ñ C ±T¹$µZð C É8ð C ÑX\Ñ]aÕXbñ C r    c                óR  € . pV  EF  pRVR,          R. /pRV9   d–   . pVR,           F‚  p\         P                  ! V4      pVP                  R4      ;'       g    Tp\        VR,          \        4      '       d!   \
        P                  ! VR,          4      VR&   VP                  V4       K„  	  WTR&   RV9   d   VR,          VR&   RV9   d   . MVP                  R4      ;'       g    . p\        V\        4      '       d   RRRV/.pV EFÍ  p	V	R,          p
V
R9   d&   VR,          P                  RRRV	R,          /4       K9  V
R9   dl   V\        P                  \        P                  39   dG   V	R	,          p\        V\        4      '       d
   VR
,          pVR,          P                  RRR
V/4       K«  V
R8X  dw   V\        P                  8X  db   V	R,          p\        V\        4      '       d   VP                  RR4      MRpVR,          pVR,          P                  RRR
RV RV 2/4       EK(  V
R8X  dS   V\        P                  \        P                  39   d.   VR,          P                  RRR
V	R,          R
,          /4       EK�  V
R8X  g   EK‹  V\        P                  8X  g   EK£  VR,          P                  RRR
V	R,          R
,          /4       EKÐ  	  V\        P                  8X  d#   RP                  R VR,           4       4      VR&   VP                  V4       EK   	  V# )aõ  Convert OpenAI-format messages to the format expected by HF processors.

All modalities extract text. VLM additionally handles ``image_url`` and ``video_url``.
MULTIMODAL handles all of the above plus ``input_audio`` and ``audio_url``.
For LLMs, the content parts are collapsed into a plain text string.

Args:
    messages (`list[dict]`): OpenAI-format chat messages.
    modality (`Modality`): The model modality (LLM, VLM, or MULTIMODAL).

Returns:
    `list[dict]`: Processor-compatible messages.
Úrolerz   rU   rk   rI   Útool_call_idrA   r'  Ú	image_urlÚurlÚimageÚinput_audior)   ÚwavÚdataÚaudiozdata:audio/z;base64,Ú	video_urlÚvideoÚ	audio_urlÚ c              3   ó2   "  € T F  qR ,          x € K  	  R# 5i)r'  Nr   )rW   rÎ   s   & r!   r[   ÚABaseHandler.get_processor_inputs_from_messages.<locals>.<genexpr>á  s   é € Ð,RÑ@Q¸1¨v¯YªYÓ@Qùs   ‚)r'  Ú
input_textÚoutput_text)r]  Úinput_image)rP  rQ  rl   rm   r(   rE   rO  r©   r   r   r   rO   r   Újoin)rY  r  Úprocessor_inputsÚmessageru   rU   Útcrh  Úraw_contentrz   Úcontent_typer^  r`  ÚfmtÚ	audio_b64s   &&             r!   Ú"get_processor_inputs_from_messagesÚ.BaseHandler.get_processor_inputs_from_messages   sÙ  € ð ÐäˆGØ˜g f�o¨y¸"Ð=ˆFð ˜wÔ&Ø�
Ø! ,×/Ð/�BÜŸš rÓ*�BØŸ™ 
Ó+×1Ð1¨r�BÜ! " [¥/´3×7Ò7Ü*.¯*ª*°R¸µ_Ó*E˜˜;™Ø×%Ñ% bÖ)ñ 0ð (2�|Ñ$Ø Ô(Ø)0°Õ)@��~Ñ&ð !-°Ô 7™"¸g¿k¹kÈ)Ó>T×>ZÐ>ZÐXZˆKÜ˜+¤s×+Ò+Ø &¨°¸ÐDÐE�ä&�Ø& v��àÐ#HÔHØ˜9Õ%×,Ñ,¨f°f¸fÀgÈfÅoÐ-VÖWà!Ð%AÔAÀhÔS[×S_ÑS_Ôai×atÑatÐRuÔFuà! +Õ.�CÜ! #¤t×,Ò,Ø! %�j˜Ø˜9Õ%×,Ñ,¨f°g¸uÀcÐ-JÖKà! ]Ô2°xÄ8×CVÑCVÔ7VØ")¨-Õ"8�KÜ>HÈÔVZ×>[Ò>[˜+Ÿ/™/¨(°EÔ:Ðaf�CØ +¨FÕ 3�IØ˜9Õ%×,Ñ,¨f°g¸uÈÐTWÐSXÐX`ÐajÐ`kÐFlÐ-m×nà! [Ô0°XÄ(Ç,Á,ÔPX×PcÑPcÐAdÔ5dØ˜9Õ%×,Ñ,¨f°g¸uÀgÈkÕFZÐ[`ÕFaÐ-b×cØ! [×0°XÄ×ATÑAT×5TØ˜9Õ%×,Ñ,¨f°g¸uÀgÈkÕFZÐ[`ÕFaÐ-b×cñ- 'ð2 œ8Ÿ<™<Ô'Ø$'§H¡HÑ,RÀÀyÖ@QÓ,RÓ$R��yÑ!à×#Ñ# F×+ñe  ðf  Ðr    c                óP   <€ V ^8„  d   Qh/ S[ R,          ;R&   S[S[,          ;R&   # )r%   Nr1  r2  )rA   r0  r(   )r)   r*   s   "€r!   r+   r#  +  s'   ø‡ ‚ ñ  �Ñ+ñ ñ ™•HÑ$ò r    )r"  r!  r   r.   r  )r   r   r   r   r3   r1  r0  r2  r1   r6  Ústaticmethodr>  rF  rV  ru  Ú__annotate_func__r   r4   r5   s   @r!   r  r  +  sx   ø‡ € ñ
ð (,ÐÙ"›u€N÷?ò ?÷Yð Yð ÷Gó ðG÷*ð *÷*0!ò 0!ðd ÷C ó ðC ÷m ƒ r    r  )	Úqwen2Ú	qwen2_moeÚqwen2_vlÚ
qwen2_5_vlÚqwen3Ú	qwen3_moeÚ
qwen3_nextÚqwen3_vlÚqwen3_vl_moe)Úqwen3_5Úqwen3_5_moe)z
<|channel>ÚthoughtÚ
r.   )Gr3   r	  rP  ÚenumrE   r  Úabcr   r   Úcollections.abcr   Úconcurrent.futuresr   r  r   Útypingr   Útransformers.utilsr	   ÚpydanticÚ
tokenizersrP  rN  r
   r   r   r   r   Ú:transformers.generation.continuous_batching.continuous_apir   Ú4transformers.generation.continuous_batching.requestsr   Ú5transformers.generation.continuous_batching.schedulerr   r   r   Ú
get_loggerr   r
  ÚX_REQUEST_IDÚEnumr   r#   rf  r7   r(   r:   rã  r=   rc   rg   ro   rv   r‰   rˆ   r�   r–   r‹   r¬   r®   r   r  r5  rR  rX  rZ  r{  r–  r¾  rû  r  r   r    r!   Ú<module>r•     sV  ðñó Û Û Û Û ß #Ý $Ý %Ý Ý  å &÷ ÛÛÛ÷õ õ eÝUÝOå+ð 
×	Ò	˜HÓ	%€ð €ôˆt�y‰yô ÷ñ ôP˜9ô Pô�Cô ô˜ô ð
ð 	ˆ}Øˆ~ØØÐ ?Ø�GØ�f˜h¨
°FÐ;ð
ðð Øˆ}Øˆ~ØØÐ \Ø�GØØ˜ØØ˜V XÐ.ØØ Ø+Ð-lð"ðð	ð
ð!ð1*Ð õZBõ<õ".ð0 ˆiˆ[Ø	ˆ:ØØ�ØØ˜ Ð*Ø˜ Ð)ð
ð 	Ðað
ðÐ ð, ˆwÒ7¸ÀÐMð	Ð ÷õDõ(õ8÷2(
ñ (
õVD÷NRñ R÷jL2ñ L2õ^õ!÷&ñ &ôRA>˜#ô A>ôHDÐ)ô DôNg1Ð+ô g1÷TNGñ NG÷by ó y r    