§
    ‚ŠtjxÔ  ã            
       óô  — d Z ddlZddlZddlZddlZddlZddlmZmZ ddl	m
Z
 ddlmZ ddlmZ ddlmZ ddlmZ er2ddlZddlZddlZdd	lmZmZmZmZmZ dd
lmZ ddlmZ ddl m!Z! ddl"m#Z#  ej$        e%¦  «        Z&dZ' G d„ dej(        ¦  «        Z) G d„ d¦  «        Z* G d„ de+¦  «        Z, G d„ de-¦  «        Z. G d„ de/¦  «        Z0ddi ddddddd œid!œd"œddd#d$d%d&d'id%d(d)œd*œd+œd,œd"œd-œZ1d.d/d0e2dz  fd1„Z3d2e2d0e2fd3„Z4d4e2d0e5e2         dz  fd5„Z6d6gd7d%d&d'id&d'id8œd9d:œd;œZ7d<g d=¢d>d?œiZ8d`d.d/d0e2dz  fd@„Z9dAe-dBe2d0e:e-e-dz  f         fdC„Z;dDe5e<         d0e=fdE„Z>dFe<d0e=fdG„Z? G dH„ dI¦  «        Z@dJe
dKe-d0eAfdL„ZB G dM„ dN¦  «        ZC G dO„ dP¦  «        ZDdQe<d0dfdR„ZEdadS„ZF G dT„ dU¦  «        ZG G dV„ dWe¦  «        ZH G dX„ dYeH¦  «        ZI G dZ„ d[eH¦  «        ZJ G d\„ d]¦  «        ZK G d^„ d_¦  «        ZLdS )bz?
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                   ó"   — e Zd ZdZdZdZdZdZdS )ÚModalityÚLLMÚVLMÚ
MULTIMODALÚSTTÚTTSN)Ú__name__Ú
__module__Ú__qualname__r   r   r   r   r   © ó    ú\/var/www/html/CA-Chatbot/venv/lib/python3.11/site-packages/transformers/cli/serving/utils.pyr   r   9   s'   € € € € € Ø
€CØ
€CØ€JØ
€CØ
€C€C€Cr   r   c                   ó   — e Zd ZdZdefd„ZdS )Ú_StreamErrorz5Sentinel to signal an error from the generate thread.Úmsgc                 ó   — || _         d S ©N)r#   )Úselfr#   s     r    Ú__init__z_StreamError.__init__D   s   € ØˆŒˆˆr   N)r   r   r   Ú__doc__Ústrr'   r   r   r    r"   r"   A   s5   € € € € € Ø?Ð?ð˜Cð ð ð ð ð ð r   r"   c                   ó   — e Zd ZdZdS )Ú_GenerationCancelledzERaised inside ``DirectStreamer.put()`` to abort ``model.generate()``.N©r   r   r   r(   r   r   r    r+   r+   H   s   € € € € € ØOÐOÐOÐOr   r+   c                   ó   — e Zd ZdZdS )Ú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``.
    Nr,   r   r   r    r.   r.   L   ó   € € € € € ðð ð ð r   r.   c                   ó   — e Zd ZdZdS )Ú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.
    Nr,   r   r   r    r1   r1   T   r/   r   r1   z<tool_call>z</tool_call>z<|im_start|>assistant
Ú
tool_callsTÚjson)ÚopenÚcloseÚrepeatsÚcontent)ÚdefaultsÚstart_anchorÚfields)ÚstcÚetcÚschemaz9<function=(?P<name>[^>\n]+)>(?P<arguments>.*?)</function>ÚarrayÚobjectÚtypeÚstringz<<parameter=(?P<key>[^>\n]+)>\s*(?P<value>.*?)\s*</parameter>)r@   zx-regex-key-value©ÚnameÚ	arguments)r@   Ú
properties)zx-regex-iteratorr@   Úitems))	Úqwen2Ú	qwen2_moeÚqwen2_vlÚ
qwen2_5_vlÚqwen3Ú	qwen3_moeÚ
qwen3_nextÚqwen3_vlÚqwen3_vl_moe)Úqwen3_5Úqwen3_5_moeÚmodelr   Úreturnc                 ó‚  ‡— t          | d| ¦  «        }t          |dd¦  «        }t          |dd¦  «        }t          |dd¦  «        }t          |dd¦  «        }d}|rF|rD|rBd|                     di ¦  «        v r*i d|d         d         id	œ}d
D ]}||v r||         ||<    nŒnp|r|r|r|d         d         }n[|j        j        Št	          ˆfd„t
                               ¦   «         D ¦   «         d¦  «        }	|	€dS |	d         |	d         |	d         }}}|                     |¦  «        }
|                     |¦  «        }||
|dœS )af  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_templateÚresponse_schemar2   r:   )r8   r:   )r9   Ústart_anchor_patternrE   c              3   ó*   •K  — | ]\  }}‰|v ¯	|V — Œd S r%   r   )Ú.0ÚtypesÚvÚ
model_types      €r    ú	<genexpr>z'get_tool_call_config.<locals>.<genexpr>´   s2   øè è € Ð_Ð_™x˜u aÈ:ÐY^ÐK^ÐK^˜ÐK^ÐK^ÐK^ÐK^Ð_Ð_r   r;   r<   r=   )r=   Ústc_idÚetc_id)ÚgetattrÚgetÚconfigr_   ÚnextÚ_TOOL_CALL_FALLBACKSrF   Úconvert_tokens_to_ids)Ú	processorrR   rU   r;   r<   rX   rY   r=   Ú
anchor_keyÚfallbackra   rb   r_   s               @r    Úget_tool_call_configrl   “   s¸  ø€ õ ˜	 ;°	Ñ:Ô:€IÝ
�)˜[¨$Ñ
/Ô
/€CÝ
�)˜[¨$Ñ
/Ô
/€CÝ 	Ð+>ÀÑEÔEÐÝ˜iÐ):¸DÑAÔA€Oà€Fà
ð Pˆsð PÐ(ð P¨\Ð=N×=RÒ=RÐS[Ð]_Ñ=`Ô=`Ð-`Ð-`àØ#Ð%6°xÔ%@ÀÔ%NÐOð
ð 
ˆð
 Cð 	ð 	ˆJØÐ.Ð.Ð.Ø%6°zÔ%B��zÑ"Ø�ð /øð 
ð 	P�ð 	P˜ð 	PØ  Ô.¨|Ô<ˆˆð ”\Ô,ˆ
ÝÐ_Ð_Ð_Ð_Õ+?×+EÒ+EÑ+GÔ+GÐ_Ñ_Ô_ÐaeÑfÔfˆØÐØ�4Ø# Eœ?¨H°U¬O¸XÀhÔ=O�&ˆSˆà×,Ò,¨SÑ1Ô1€FØ×,Ò,¨SÑ1Ô1€FØ¨¸&ÐAÐAÐAr   Ú	tool_callc                 óÂ   — |                       d| ¦  «        }|                      di ¦  «        }|d         t          |t          ¦  «        st          j        |¦  «        n|dœS )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.
    ÚfunctionrD   rC   rB   )rd   Ú
isinstancer)   r3   Údumps)rm   ro   rD   s      r    Ú_normalize_tool_callrr   ¾   sc   € ð �}Š}˜Z¨Ñ3Ô3€HØ—’˜[¨"Ñ-Ô-€Ià˜Ô Ý2<¸YÍÑ2LÔ2LÐ[•T”Z 	Ñ*Ô*Ð*ÐR[ðð ð r   r=   c                 óÐ   — |                       ||d¬¦  «        }t          |t          ¦  «        rd|v r|d         }|sdS t          |t          ¦  «        s|g}d„ |D ¦   «         }|r|ndS )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.
    Ú )Úprefixr2   Nc                 ó,   — g | ]}t          |¦  «        ‘ŒS r   )rr   )r\   rm   s     r    ú
<listcomp>z$parse_tool_calls.<locals>.<listcomp>ã   s!   € ÐJÐJÐJ°iÕ& yÑ1Ô1ÐJÐJÐJr   )Úparse_responserp   ÚdictÚlist)ri   Úgenerated_idsr=   Úparsedr2   s        r    Úparse_tool_callsr}   Ï   sŒ   € ð ×%Ò% m°VÀBÐ%ÑGÔG€Få�&�$ÑÔð & L°FÐ$:Ð$:Ø˜Ô%ˆØð ØˆtÝ�f�dÑ#Ô#ð Ø�ˆØJÐJÀ6ÐJÑJÔJ€JØ#Ð-ˆ:ˆ:¨Ð-r   z<think>z</think>)Úthinkingr7   zK(?:<think>)?(?P<thinking>.*?)</think>(?P<content>.*?)(?:<\|[^|<>\s]+\|>)?\Z)r@   rE   zx-regex)ÚstartÚendr=   Úgemma4)z
<|channel>Úthoughtú
z
<channel|>)r   r€   c                 ó  ‡‡	— t          | d| ¦  «        Š	|j        j                             ¦   «         Št	          ˆfd„t
                               ¦   «         D ¦   «         t          ¦  «        }ˆ	fd„|d         D ¦   «         }‰	                     |d         ¦  «        }t          ˆ	fd„|D ¦   «         ¦  «        s|d‰	j
        fv rdS t          ‰	dd¦  «        }|r
d	|d
         v st          d         }|||dœ}|�t          ||¦  «        |d<   |S )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.
    rU   c              3   ó.   •K  — | ]\  }}|‰k    ¯|V — Œd S r%   r   )r\   Úkr^   r_   s      €r    r`   z'get_reasoning_config.<locals>.<genexpr>  s+   øè è € ÐCÐC‰tˆq�!°1¸
²?°?ˆ°?°?°?°?ÐCÐCr   c                 ó:   •— g | ]}‰                      |¦  «        ‘ŒS r   )rh   )r\   ÚtrU   s     €r    rw   z(get_reasoning_config.<locals>.<listcomp>  s'   ø€ ÐVÐVÐV¸�×0Ò0°Ñ3Ô3ÐVÐVÐVr   r   r€   c              3   ó.   •K  — | ]}|d ‰j         fv V — Œd S r%   )Úunk_token_id)r\   ÚtidrU   s     €r    r`   z'get_reasoning_config.<locals>.<genexpr>  s0   øè è € Ð
FÐ
F°Sˆ3�4˜Ô/Ð0Ð0Ð
FÐ
FÐ
FÐ
FÐ
FÐ
Fr   NrY   r~   rE   r=   )Ú	start_idsÚend_idr=   Ústart_in_thinking)rc   re   r_   Úlowerrf   Ú_THINKING_TOKENSrF   Ú_DEFAULT_THINKING_TOKENSrh   ÚanyrŠ   Ú_starts_in_thinking)
ri   rR   Ú	input_idsÚthinking_tokensrŒ   r�   r=   re   r_   rU   s
           @@r    Úget_reasoning_configr–     s>  øø€ õ ˜	 ;°	Ñ:Ô:€IØ”Ô(×.Ò.Ñ0Ô0€JÝØCÐCÐCÐCÕ'×-Ò-Ñ/Ô/ÐCÑCÔCÝ ñô €Oð WÐVÐVÐV¸_ÈWÔ=UÐVÑVÔV€IØ×,Ò,¨_¸UÔ-CÑDÔD€FÝ
Ð
FÐ
FÐ
FÐ
F¸IÐ
FÑ
FÔ
FÑFÔFð È&ÐUYÐ[dÔ[qÐTrÐJrÐJrØˆtõ �YÐ 1°4Ñ8Ô8€FØð 4�z V¨LÔ%9Ð9Ð9Ý)¨(Ô3ˆØ!*°fÈÐOÐO€FØÐÝ&9¸)ÀYÑ&OÔ&OˆÐ"Ñ#Ø€Mr   r7   Úreasoning_configc                 óØ   — |                       ||d         ¦  «        }|r0|                     dd¦  «        }|r|                     dd¦  «        |fS |                     d¦  «        rd|fS |dfS )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=   r~   rt   r7   rŽ   N)rx   rd   )ri   r{   r7   r—   r|   Ú	reasonings         r    Úparse_reasoningrš   %  sˆ   € ð ×%Ò% mÐ5EÀhÔ5OÑPÔP€FØð 8Ø—J’J˜z¨2Ñ.Ô.ˆ	Øð 	8Ø—:’:˜i¨Ñ,Ô,¨iÐ7Ð7ð ×ÒÐ/Ñ0Ô0ð Ø�7ˆ{ÐØ�Dˆ=Ðr   rŒ   c                 ób  — t          | d¦  «        r|                      ¦   «         } | r8t          | d         t          ¦  «        rt	          | ¦  «        dk    rdS | d         } t	          |¦  «        }dD ]>}t	          | ¦  «        ||z   k    r&t	          | ¦  «        |z
  }| ||z
  |…         |k    r dS Œ?dS )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.
    Útolistr   r   F)r   r   T)Úhasattrrœ   rp   rz   Úlen)r”   rŒ   ÚnÚtrailingr€   s        r    r“   r“   9  sÇ   € õ ˆy˜(Ñ#Ô#ð 'Ø×$Ò$Ñ&Ô&ˆ	Øð !•Z 	¨!¤­dÑ3Ô3ð !Ýˆy‰>Œ>˜QÒÐØ�5Ø˜a”Lˆ	ÝˆI‰Œ€Aàð ð ˆÝˆy‰>Œ>˜Q ™\Ò)Ð)Ý�i‘.”. 8Ñ+ˆCØ˜˜q™ 3˜Ô'¨9Ò4Ð4Ø�t�tøØˆ5r   Útoken_idc                 óR  — | j         €dS | j        r|| j        k    r	d| _        dS dS | j         t          | j        ¦  «                 }||k    r	g | _        dS | j                             |¦  «         t          | j        ¦  «        t          | j         ¦  «        k    rd| _        g | _        dS )uU  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.
    NFT)Ú_thinking_start_idsÚ_inside_thinkingÚ_thinking_end_idrž   Ú_thinking_prefixÚappend)Ústreamerr¡   Úexpecteds      r    Ú_advance_thinking_staterª   U  sº   € ð Ô#Ð+ØˆuØÔ ð Ø�xÔ0Ò0Ð0Ø(-ˆHÔ%Ø�4ØˆuØÔ+­C°Ô0IÑ,JÔ,JÔK€HØ�8ÒÐØ$&ˆÔ!ØˆuØÔ×$Ò$ XÑ.Ô.Ð.Ý
ˆ8Ô$Ñ%Ô%­¨XÔ-IÑ)JÔ)JÒJÐJØ$(ˆÔ!Ø$&ˆÔ!Øˆ4r   c                   ór   — e Zd ZdZdedefd„Zdededz  ddfd	„Zded
ededz  ddfd„Z	deddfd„Z
dd„ZdS )ÚDownloadAggregatora	  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.
    ÚenqueueÚmodel_idc                 ó>   — || _         || _        i | _        d | _        d S r%   )r­   rR   ÚbarsÚlast_emitted_current)r&   r­   r®   s      r    r'   zDownloadAggregator.__init__u  s%   € ØˆŒØˆŒ
Ø79ˆŒ	Ø04ˆÔ!Ð!Ð!r   Úbar_idÚtotalNrS   c                 óF   — d|f| j         |<   |                      ¦   «          dS )z6Register a new download bar with its total byte count.r   N©r°   Ú_emit)r&   r²   r³   s      r    ÚregisterzDownloadAggregator.register{  s#   € à ˜JˆŒ	�&ÑØ�
Š
‰Œˆˆˆr   Úcurrentc                 óF   — ||f| j         |<   |                      ¦   «          dS )z>Update a bar's current byte count and emit aggregate progress.Nrµ   )r&   r²   r¸   r³   s       r    ÚupdatezDownloadAggregator.update€  s$   € à$ eÐ,ˆŒ	�&ÑØ�
Š
‰Œˆˆˆr   c                 ó   — d S r%   r   )r&   r²   s     r    r5   zDownloadAggregator.close…  ó   € Øˆr   c                 ó>  — t          d„ | j                             ¦   «         D ¦   «         ¦  «        }|| j        k    rd S || _        d„ | j                             ¦   «         D ¦   «         }|rt          |¦  «        nd }|                      d| j        d||dœdœ¦  «         d S )Nc              3   ó    K  — | ]	\  }}|V — Œ
d S r%   r   )r\   ÚcÚ_s      r    r`   z+DownloadAggregator._emit.<locals>.<genexpr>‰  s&   è è € Ð;Ð;¡  1˜!Ð;Ð;Ð;Ð;Ð;Ð;r   c                 ó   — g | ]	\  }}|®|‘Œ
S r%   r   )r\   rÀ   rˆ   s      r    rw   z,DownloadAggregator._emit.<locals>.<listcomp>�  s   € ÐDÐDÐD™˜˜1°a°m�!°m°m°mr   ÚloadingÚdownload©r¸   r³   ©ÚstatusrR   ÚstageÚprogress)Úsumr°   Úvaluesr±   r­   rR   )r&   Úagg_currentÚtotalsÚ	agg_totals       r    r¶   zDownloadAggregator._emitˆ  s¼   € ÝÐ;Ð;¨¬	×(8Ò(8Ñ(:Ô(:Ð;Ñ;Ô;Ñ;Ô;ˆØ˜$Ô3Ò3Ð3ØˆFØ$/ˆÔ!ØDÐD ¤	× 0Ò 0Ñ 2Ô 2ÐDÑDÔDˆØ#)Ð3•C˜‘K”K�K¨tˆ	Ø�Šà#ØœØ#Ø(3¸iÐHÐHð	ð ñ	
ô 	
ð 	
ð 	
ð 	
r   ©rS   N)r   r   r   r(   r   r)   r'   Úintr·   rº   r5   r¶   r   r   r    r¬   r¬   n  s×   € € € € € ðð ð5 ð 5°Cð 5ð 5ð 5ð 5ð˜sð ¨3°©:ð ¸$ð ð ð ð ð
˜Sð ¨3ð °s¸T±zð Àdð ð ð ð ð
˜Cð  Dð ð ð ð ð
ð 
ð 
ð 
ð 
ð 
r   r¬   Úcallbackr®   c                 ó\   ‡ ‡‡— ddl m} t          ‰ ‰¦  «        Š G ˆ ˆˆfd„d|¦  «        }|S )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*.
    r   )Útqdmc                   óL   •‡ — e Zd Zˆ ˆfd„Zdˆˆˆfd„	Zˆˆˆfd„Zˆ ˆfd„Zˆ xZS )ú.make_progress_tqdm_class.<locals>.ProgressTqdmc                 ó  •— |                      d¦  «        pd| _        d|d<    t          ¦   «         j        |i |¤Ž d| _        d| _        | j        dk    r6t          | ¦  «        | _        ‰                     | j        | j	        ¦  «         d S d S )NÚunitÚitTÚdisabler   éÿÿÿÿÚB)
rd   Ússe_unitÚsuperr'   rŸ   Úlast_emittedÚidÚ_bar_idr·   r³   )r&   ÚargsÚkwargsÚ	__class__Údownload_aggregators      €€r    r'   z7make_progress_tqdm_class.<locals>.ProgressTqdm.__init__¬  s’   ø€ Ø"ŸJšJ vÑ.Ô.Ð6°$ˆDŒMØ $ˆF�9ÑØ�E‰GŒGÔ˜dÐ- fÐ-Ð-Ð-ØˆDŒFØ "ˆDÔØŒ} Ò#Ð#Ý! $™xœx�”Ø#×,Ò,¨T¬\¸4¼:ÑFÔFÐFÐFÐFð $Ð#r   r   c                 ó  •— |€d}| xj         |z  c_         | j        dk    r(‰                     | j        | j         | j        ¦  «         d S | j         | j        k    r+| j         | _         ‰d‰d| j         | j        dœdœ¦  «         d S d S ©Nr   rÚ   rÂ   ÚweightsrÄ   rÅ   )rŸ   rÛ   rº   rß   r³   rÝ   )r&   rŸ   rÐ   rã   r®   s     €€€r    rº   z5make_progress_tqdm_class.<locals>.ProgressTqdm.update¶  s¯   ø€ ØˆyØ�ØˆFŒF�a‰KˆFŒFØŒ} Ò#Ð#Ø#×*Ò*¨4¬<¸¼ÀÄÑLÔLÐLÐLÐLØ”˜4Ô,Ò,Ð,Ø$(¤F�Ô!Ø�à"+Ø!)Ø!*Ø04´ÀÄÐ$LÐ$Lð	ð ñô ð ð ð ð -Ð,r   c           	   3   ó  •K  — | j         D ]�}| xj        dz  c_        | j        dk    r'‰                     | j        | j        | j        ¦  «         n9| j        | j        k    r)| j        | _         ‰d‰d| j        | j        dœdœ¦  «         |V — Œ‚d S rå   )ÚiterablerŸ   rÛ   rº   rß   r³   rÝ   )r&   ÚitemrÐ   rã   r®   s     €€€r    Ú__iter__z7make_progress_tqdm_class.<locals>.ProgressTqdm.__iter__Ç  s¼   øè è € Øœð ð �Ø�”˜!‘�”Ø”= CÒ'Ð'Ø'×.Ò.¨t¬|¸T¼VÀTÄZÑPÔPÐPÐPØ”V˜tÔ0Ò0Ð0Ø(,¬�DÔ%Ø�Hà&/Ø%-Ø%.Ø48´FÀTÄZÐ(PÐ(Pð	ð ñô ð ð �
�
�
�
ðð r   c                 ó’   •— | j         dk    r‰                     | j        ¦  «         t          ¦   «                              ¦   «          d S )NrÚ   )rÛ   r5   rß   rÜ   )r&   râ   rã   s    €€r    r5   z4make_progress_tqdm_class.<locals>.ProgressTqdm.closeØ  s;   ø€ ØŒ} Ò#Ð#Ø#×)Ò)¨$¬,Ñ7Ô7Ð7Ý‰GŒG�MŠM‰OŒOˆOˆOˆOr   )r   )r   r   r   r'   rº   rê   r5   Ú__classcell__)râ   rÐ   rã   r®   s   @€€€r    ÚProgressTqdmrÔ   «  s­   øø€ € € € € ð	Gð 	Gð 	Gð 	Gð 	Gð 	Gð	ð 	ð 	ð 	ð 	ð 	ð 	ð 	ð"	ð 	ð 	ð 	ð 	ð 	ð 	ð"	ð 	ð 	ð 	ð 	ð 	ð 	ð 	ð 	ð 	r   rí   )Ú	tqdm.autorÒ   r¬   )rÐ   r®   Ú	base_tqdmrí   rã   s   ``  @r    Úmake_progress_tqdm_classrð   ™  sp   øøø€ ð ,Ð+Ð+Ð+Ð+Ð+å,¨X°xÑ@Ô@Ðð0ð 0ð 0ð 0ð 0ð 0ð 0ð 0ð 0�yñ 0ô 0ð 0ðd Ðr   c                   ór   — e Zd ZdZ	 	 	 ddddej        dej        ded	edz  d
edz  fd„Z	dd„Z
dd„Zdd„ZdS )ÚDirectStreamera†  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.
    TNrU   útokenizers.TokenizerÚloopÚqueueÚskip_special_tokensÚtool_configr—   c                 ó®  — ddl m} || _        || _        || _         |g |¦  «        | _        |r|d         nd| _        |r|d         nd| _        d| _        |r|d         nd| _	        |r|d         nd| _
        t          |o|                     d	¦  «        ¦  «        | _        g | _        d
| _        t!          j        ¦   «         | _        d| _        g | _        dS )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``.
        r   ©ÚDecodeStreamra   Nrb   FrŒ   r�   rŽ   T)Útokenizers.decodersrú   Ú
_tokenizerÚ_loopÚ_queueÚ_decode_streamÚ_stc_idÚ_etc_idÚ_inside_tool_callr£   r¥   Úboolrd   r¤   r¦   Ú_firstÚ	threadingÚEventÚ
_cancelledÚtotal_tokensÚgenerated_token_ids)r&   rU   rô   rõ   rö   r÷   r—   rú   s           r    r'   zDirectStreamer.__init__é  s   € ð. 	5Ð4Ð4Ð4Ð4Ð4à#ˆŒØˆŒ
ØˆŒØ*˜l¨2Ð/BÑCÔCˆÔØ0;ÐE�{ 8Ô,Ð,ÀˆŒØ0;ÐE�{ 8Ô,Ð,ÀˆŒØ!&ˆÔØDTÐ#^Ð#3°KÔ#@Ð#@ÐZ^ˆÔ Ø>NÐ XÐ 0°Ô :Ð :ÐTXˆÔÝ $Ð%5Ð%cÐ:J×:NÒ:NÐObÑ:cÔ:cÑ dÔ dˆÔØ+-ˆÔØˆŒÝ#œ/Ñ+Ô+ˆŒØˆÔØ.0ˆÔ Ð Ð r   Úvalueútorch.TensorrS   c                 óD  — | j                              ¦   «         rt          ¦   «         ‚| j        r	d| _        dS |                     ¦   «         D ]Ó}| xj        dz  c_        | j                             |¦  «         || j        k    rd| _	        n|| j
        k    rd| _	        t          | |¦  «        }| j                             | j        |¦  «        }|�| j	        s|| j
        k    s|rŒ˜| j        rt!          |¦  «        }| j                             | j        j        |¦  «         ŒÔdS )zHCalled by ``model.generate()`` after each decode step with new token(s).FNr   T)r  Úis_setr+   r  rœ   r  r	  r§   r   r  r  rª   rÿ   Ústeprü   r¤   r.   rý   Úcall_soon_threadsaferþ   Ú
put_nowait)r&   r
  r¡   Úis_start_or_end_tokenÚtexts        r    ÚputzDirectStreamer.put  s:  € àŒ?×!Ò!Ñ#Ô#ð 	)Ý&Ñ(Ô(Ð(àŒ;ð 	ØˆDŒKØˆFØŸš™œð 	Jð 	JˆHØÐÔ Ñ"ÐÔØÔ$×+Ò+¨HÑ5Ô5Ð5à˜4œ<Ò'Ð'Ø)-�Ô&Ð&Ø˜Tœ\Ò)Ð)Ø).�Ô&å$;¸DÀ(Ñ$KÔ$KÐ!àÔ&×+Ò+¨D¬O¸XÑFÔFˆDØˆ|˜tÔ5ˆ|¸ÀTÄ\Ò9QÐ9QÐUjÐ9QØØÔ$ð +Ý$ TÑ*Ô*�ØŒJ×+Ò+¨D¬KÔ,BÀDÑIÔIÐIÐIð!	Jð 	Jr   c                 óP   — | j                              | j        j        d¦  «         dS )z;Called by ``model.generate()`` when generation is complete.N)rý   r  rþ   r  ©r&   s    r    r€   zDirectStreamer.end,  s%   € àŒ
×'Ò'¨¬Ô(>ÀÑEÔEÐEÐEÐEr   c                 ó8   — | j                              ¦   «          dS )zWSignal cancellation. The next ``put()`` call will raise and abort ``model.generate()``.N)r  Úsetr  s    r    ÚcancelzDirectStreamer.cancel0  s   € àŒ×ÒÑÔÐÐÐr   )TNN)r
  r  rS   NrÎ   )r   r   r   r(   ÚasyncioÚAbstractEventLoopr   r  ry   r'   r  r€   r  r   r   r    rò   rò   à  sÎ   € € € € € ðð ð %)Ø#'Ø(,ð'1ð '1à)ð'1ð Ô'ð'1ð Œ}ð	'1ð
 "ð'1ð ˜D‘[ð'1ð  ™+ð'1ð '1ð '1ð '1ðRJð Jð Jð Jð4Fð Fð Fð Fðð ð ð ð ð r   rò   c                   ót   — e Zd ZdZ	 	 ddddedddej        d	ej        d
edz  dedz  fd„Z	dd„Z
dd„Zdd„ZdS )Ú
CBStreamera„  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.
    NÚ
cb_managerr   Ú
request_idrU   ró   rô   rõ   r÷   r—   c                 óš  — ddl m} || _        || _        || _        || _        || _         |g d¦  «        | _        |r|d         nd| _        |r|d         nd| _	        d| _
        |r|d         nd| _        |r|d	         nd| _        t          |o|                     d
¦  «        ¦  «        | _        g | _        d| _        d| _        g | _        dS )aY  
        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   rù   Tra   Nrb   FrŒ   r�   rŽ   )rû   rú   Ú_cbÚ_request_idrý   rþ   rü   rÿ   r   r  r  r£   r¥   r  rd   r¤   r¦   Ú	_prev_lenr  r	  )	r&   r  r  rU   rô   rõ   r÷   r—   rú   s	            r    r'   zCBStreamer.__init__>  sÿ   € ð( 	5Ð4Ð4Ð4Ð4Ð4àˆŒØ%ˆÔØˆŒ
ØˆŒØ#ˆŒØ*˜l¨2¨tÑ4Ô4ˆÔØ0;ÐE�{ 8Ô,Ð,ÀˆŒØ0;ÐE�{ 8Ô,Ð,ÀˆŒØ!&ˆÔØDTÐ#^Ð#3°KÔ#@Ð#@ÐZ^ˆÔ Ø>NÐ XÐ 0°Ô :Ð :ÐTXˆÔÝ $Ð%5Ð%cÐ:J×:NÒ:NÐObÑ:cÔ:cÑ dÔ dˆÔØ+-ˆÔØˆŒØˆÔØ.0ˆÔ Ð Ð r   Úoutputr   rS   c                 óö  — |j         | j        d…         }t          |j         ¦  «        | _        |D ]È}| xj        dz  c_        | j                             |¦  «         || j        k    rd| _        n|| j        k    rd| _        t          | |¦  «        }| j
                             | j        |¦  «        }|�| j        s|| j        k    s|rŒ˜| j        rt          |¦  «        }| j                             |¦  «         ŒÉdS )zLDecode new tokens from a CB ``GenerationOutput`` and push text to the queue.Nr   TF)Úgenerated_tokensr"  rž   r  r	  r§   r   r  r  rª   rÿ   r  rü   r¤   r.   rþ   r  )r&   r#  Ú
new_tokensr¡   r  r  s         r    r  zCBStreamer.pute  s  € àÔ,¨T¬^Ð-=Ð-=Ô>ˆ
Ý˜VÔ4Ñ5Ô5ˆŒØ"ð 	)ð 	)ˆHØÐÔ Ñ"ÐÔØÔ$×+Ò+¨HÑ5Ô5Ð5à˜4œ<Ò'Ð'Ø)-�Ô&Ð&Ø˜Tœ\Ò)Ð)Ø).�Ô&å$;¸DÀ(Ñ$KÔ$KÐ!àÔ&×+Ò+¨D¬O¸XÑFÔFˆDØˆ|˜tÔ5ˆ|¸ÀTÄ\Ò9QÐ9QÐUjÐ9QØØÔ$ð +Ý$ TÑ*Ô*�ØŒK×"Ò" 4Ñ(Ô(Ð(Ð(ð!	)ð 	)r   c                 ó:   — | j                              d¦  «         dS )zSignal end of stream.N)rþ   r  r  s    r    r€   zCBStreamer.end{  s   € àŒ×Ò˜tÑ$Ô$Ð$Ð$Ð$r   c                 óD   — | j                              | j        ¦  «         dS )zCancel the CB request.N)r   Úcancel_requestr!  r  s    r    r  zCBStreamer.cancel  s!   € àŒ×Ò Ô 0Ñ1Ô1Ð1Ð1Ð1r   ©NN)r#  r   rS   NrÎ   )r   r   r   r(   r)   r  r  r   ry   r'   r  r€   r  r   r   r    r  r  5  sÍ   € € € € € ðð ð $(Ø(,ð%1ð %1à/ð%1ð ð%1ð *ð	%1ð
 Ô'ð%1ð Œ}ð%1ð ˜D‘[ð%1ð  ™+ð%1ð %1ð %1ð %1ðN)ð )ð )ð )ð,%ð %ð %ð %ð2ð 2ð 2ð 2ð 2ð 2r   r  Úseedc                 ó.   — ddl } |j        | ¦  «         dS )z8Set the PyTorch random seed for reproducible generation.r   N)ÚtorchÚmanual_seed)r+  r-  s     r    Úset_torch_seedr/  „  s$   € à€L€L€Là€EÔ�dÑÔÐÐÐr   c                  óv   — ddl } | j                             ¦   «         r| j                             ¦   «          dS dS )z+Empty the CUDA cache if a GPU is available.r   N)r-  ÚcudaÚis_availableÚempty_cache)r-  s    r    Úreset_torch_cacher4  ‹  sE   € à€L€L€Là„z×ÒÑ Ô ð !ØŒ
×ÒÑ Ô Ð Ð Ð ð!ð !r   c                   óB   — e Zd ZdZd„ Zdd„Zdefd„Zdej        fd„Z	dS )	ÚInferenceThreadzÒ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                 óž   — t          ¦   «         | _        t          j        | j        d¬¦  «        | _        | j                             ¦   «          d S )NT)ÚtargetÚdaemon)r   rþ   r  ÚThreadÚ_runÚ_threadr   r  s    r    r'   zInferenceThread.__init__š  s@   € Ý"™WœWˆŒÝ Ô'¨t¬yÀÐFÑFÔFˆŒØŒ×ÒÑÔÐÐÐr   rS   Nc                 óR  — 	 | j                              ¦   «         \  }}}}}	  ||i |¤Ž}|�|                     |j        |¦  «         n|                     |¦  «         nJ# t          $ r=}|�|                     |j        |¦  «         n|                     |¦  «         Y d }~nd }~ww xY wŒ§r%   )rþ   rd   r  Ú
set_resultÚ	ExceptionÚset_exception)r&   Úfnrà   rá   Úfuturerô   ÚresultÚes           r    r;  zInferenceThread._runŸ  sÒ   € ð	,Ø-1¬[¯_ª_Ñ->Ô->Ñ*ˆB��f˜f dð
,Ø˜˜TÐ, VÐ,Ð,�ØÐ#Ø×-Ò-¨fÔ.?ÀÑHÔHÐHÐHà×%Ò% fÑ-Ô-Ð-øøÝð ,ð ,ð ,ØÐ#Ø×-Ò-¨fÔ.BÀAÑFÔFÐFÐFà×(Ò(¨Ñ+Ô+Ð+øøøøøøøøøð	,øøøð	,s   ¢;A Á
B%Á(3B Â B%c                 ó`   — t          ¦   «         }| j                             ||||df¦  «         |S )úESubmit a callable to the inference thread. Returns a blocking Future.N)r   rþ   r  )r&   rA  rà   rá   rB  s        r    ÚsubmitzInferenceThread.submit®  s/   € å™œˆØŒ�Š˜˜T 6¨6°4Ð8Ñ9Ô9Ð9Øˆr   c                 ó’   — t          j        ¦   «         }|                     ¦   «         }| j                             |||||f¦  «         |S ©zOSubmit a callable to the inference thread. Returns an awaitable asyncio.Future.)r  Úget_running_loopÚcreate_futurerþ   r  )r&   rA  rà   rá   rô   rB  s         r    Úasync_submitzInferenceThread.async_submit´  sE   € åÔ'Ñ)Ô)ˆØ×#Ò#Ñ%Ô%ˆØŒ�Š˜˜T 6¨6°4Ð8Ñ9Ô9Ð9Øˆr   rÎ   )
r   r   r   r(   r'   r;  r   rG  r  rL  r   r   r    r6  r6  “  sy   € € € € € ðð ðð ð ð
,ð ,ð ,ð ,ð¨Vð ð ð ð ð°7´>ð ð ð ð ð ð r   r6  c                   óä   — e Zd ZdZdd„Ze	 	 dddd	d
dedddededz  dedz  dee	j
        df         fd„¦   «         Zeddd	d
dedddedeeeee         f         fd„¦   «         Zedd„¦   «         ZdS )ÚBaseGenerateManageruã   Base class for generation managers.

    Subclasses:
    - :class:`GenerateManager` â€” sequential ``model.generate()`` on a persistent thread.
    - :class:`CBGenerateManager` â€” continuous batching with paged attention.
    rR   r   Ú
gen_configr   rS   Nc                 ó   — dS )z:Initialize continuous batching. No-op for non-CB managers.Nr   ©r&   rR   rO  s      r    Úinit_cbzBaseGenerateManager.init_cbÄ  ó   € € € r   ri   ú(ProcessorMixin | PreTrainedTokenizerFastÚinputsr  r÷   r—   zDirectStreamer | CBStreamerc                 ó   — dS )aâ  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   )r&   rR   ri   rU  rO  r  r÷   r—   s           r    Úgenerate_streamingz&BaseGenerateManager.generate_streamingÇ  rS  r   c              ƒ   ó
   K  — dS )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   )r&   rR   ri   rU  rO  r  s         r    Úgenerate_non_streamingz*BaseGenerateManager.generate_non_streamingå  s
   è è € € € r   c                 ó   — dS )z/Stop the generation manager and free resources.Nr   r  s    r    ÚstopzBaseGenerateManager.stopû  rS  r   ©rR   r   rO  r   rS   Nr*  rÎ   )r   r   r   r(   rR  r   ry   r)   Útupler  r   rW  rÏ   rz   rY  r[  r   r   r    rN  rN  ¼  sP  € € € € € ðð ðIð Ið Ið Ið ð $(Ø(,ðð à ðð >ðð ð	ð
 'ðð ðð ˜D‘[ðð  ™+ðð 
ˆwŒ}Ð;Ð;Ô	<ðð ð ñ „^ðð: ðà ðð >ðð ð	ð
 'ðð ðð 
ˆs�C˜˜cœÐ"Ô	#ðð ð ñ „^ðð* ð>ð >ð >ñ „^ð>ð >ð >r   rN  c                   óÐ   — e Zd ZdZd„ Z	 	 dddddded	d
dededz  dedz  deej	        e
f         fd„Zddddded	d
dedeeedf         fd„Zdedefd„Zdedej        fd„Zdd„ZdS )ÚGenerateManagerzFSequential generation via ``model.generate()`` on a persistent thread.c                 ó,   — t          ¦   «         | _        d S r%   )r6  r<  r  s    r    r'   zGenerateManager.__init__  s   € Ý&Ñ(Ô(ˆŒˆˆr   NrR   r   ri   rT  rU  rO  r   r  r÷   r—   rS   c                 ó,  ‡‡‡‡— t          j        ¦   «         Št          j        ¦   «         Št          |d|¦  «        j        }t          |‰‰||¬¦  «        }	i |¥|	||dœ¥Št          ‰d¦  «        rd‰d<   d
ˆˆˆˆfd	„}
|                      |
¦  «         ‰|	fS )zLStart streaming generation via ``model.generate()`` on the inference thread.rU   ©r÷   r—   )r¨   Úgeneration_configrU   Ú
has_talkerr  Úgeneration_moderS   Nc            	      ó  •— 	  ‰j         di ‰¤Ž d S # t          $ r ‰                     ‰j        d ¦  «         Y d S t          $ r@} ‰                     ‰j        t          t          | ¦  «        ¦  «        ¦  «         Y d } ~ d S d } ~ ww xY w)Nr   )Úgenerater+   r  r  r?  r"   r)   )rD  Ú
gen_kwargsrô   rR   rõ   s    €€€€r    r;  z0GenerateManager.generate_streaming.<locals>._run  sº   ø€ ðRØ�”Ð,Ð, Ð,Ð,Ð,Ð,Ð,øÝ'ð Bð Bð BØ×)Ò)¨%Ô*:¸DÑAÔAÐAÐAÐAÐAÝð Rð Rð RØ×)Ò)¨%Ô*:½LÍÈQÉÌÑ<PÔ<PÑQÔQÐQÐQÐQÐQÐQÐQÐQøøøøðRøøøs   ƒ ’%Bº	BÁ5A>Á>BrÎ   )r  rJ  r   rc   rü   rò   r�   rG  )r&   rR   ri   rU  rO  r  r÷   r—   Úrust_tokenizerr¨   r;  rh  rô   rõ   s    `         @@@r    rW  z"GenerateManager.generate_streaming  sã   øøøø€ õ Ô'Ñ)Ô)ˆÝ&œ}™œˆå  ¨K¸ÑCÔCÔNˆÝ!Ø˜D %°[ÐScð
ñ 
ô 
ˆð o˜Ðn¨HÈ:ÐdmÐnÐnÐnˆ
Ý�5˜,Ñ'Ô'ð 	3Ø,2ˆJÐ(Ñ)ð	Rð 	Rð 	Rð 	Rð 	Rð 	Rð 	Rð 	Rð 	Rð 	�Š�DÑÔÐØ�hˆÐr   r  c              ƒ   óê   K  — i |¥||dœ¥}t          |d¦  «        rd|d<    | j        |j        fi |¤Žƒ d{V —†}|d         j        d         }|d|d…f         }	|                     |	d	¬
¦  «        }
|
||	fS )zNRun generation to completion via ``model.generate()`` on the inference thread.)rc  rU   rd  r  re  Nr”   rÙ   r   T©rö   )r�   rL  rg  ÚshapeÚdecode)r&   rR   ri   rU  rO  r  Úgenerate_kwargsÚ	sequencesÚ	input_lenr{   r  s              r    rY  z&GenerateManager.generate_non_streaming'  sµ   è è € ð ^˜VÐ]¸*ÐS\Ð]Ð]Ð]ˆÝ�5˜,Ñ'Ô'ð 	8Ø17ˆOÐ-Ñ.Ø+˜$Ô+¨E¬NÐNÐN¸oÐNÐNÐNÐNÐNÐNÐNÐNˆ	Ø˜;Ô'Ô-¨bÔ1ˆ	Ø! ! Y Z Z -Ô0ˆØ×Ò À4ÐÑHÔHˆØ�Y Ð-Ð-r   rA  c                 ó.   —  | j         j        |g|¢R i |¤ŽS )rF  )r<  rG  ©r&   rA  rà   rá   s       r    rG  zGenerateManager.submit;  s'   € à"ˆtŒ|Ô" 2Ð7¨Ð7Ð7Ð7°Ð7Ð7Ð7r   c                 ó.   —  | j         j        |g|¢R i |¤ŽS rI  )r<  rL  rr  s       r    rL  zGenerateManager.async_submit?  s'   € à(ˆtŒ|Ô(¨Ð=¨dÐ=Ð=Ð=°fÐ=Ð=Ð=r   c                 ó   — d S r%   r   r  s    r    r[  zGenerateManager.stopC  r¼   r   r*  rÎ   )r   r   r   r(   r'   ry   r)   r]  r  r   rò   rW  rÏ   rY  r   r   rG  rL  r[  r   r   r    r_  r_     sa  € € € € € ØPÐPð)ð )ð )ð $(Ø(,ðð à ðð >ðð ð	ð
 'ðð ðð ˜D‘[ðð  ™+ðð 
ˆwŒ}˜nÐ,Ô	-ðð ð ð ðB.à ð.ð >ð.ð ð	.ð
 'ð.ð ð.ð 
ˆs�C˜Ð'Ô	(ð.ð .ð .ð .ð(8˜ð 8°vð 8ð 8ð 8ð 8ð>˜xð >¸W¼^ð >ð >ð >ð >ðð ð ð ð ð r   r_  c                   óò   — e Zd ZdZddd„Zdd„Zd
efd„Zded
dfd„Z		 	 dddddde
dd	dede
dz  de
dz  d
eej        ef         fd„Zddddde
dd	ded
eeeee         f         fd„Zedd„¦   «         Zdd„ZdS )ÚCBGenerateManagerau  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.
    NÚ	cb_configúContinuousBatchingConfig | Nonec                 ó"   — d | _         || _        d S r%   )r   Ú
_cb_config)r&   rw  s     r    r'   zCBGenerateManager.__init__V  s   € Ø59ˆŒØ#ˆŒˆˆr   rR   r   rO  r   rS   c                 óŒ   — | j         �dS |                     || j        ¬¦  «        | _         | j                              ¦   «          dS )at  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)rc  Úcontinuous_batching_config)r   Úinit_continuous_batchingrz  r   rQ  s      r    rR  zCBGenerateManager.init_cbZ  sN   € ð Œ8ÐØˆFà×1Ò1Ø(ÀTÄ_ð 2ñ 
ô 
ˆŒð 	Œ�ŠÑÔÐÐÐr   c                 ó0   — | j         du p| j         j        du S )zJWhether the CB worker is healthy. ``True`` before ``init_cb()`` is called.N)r   Úfatal_errorr  s    r    Úis_alivezCBGenerateManager.is_alivek  s   € àŒx˜4ÐÐ? 4¤8Ô#7¸4Ð#?Ð?r   r  c                 ón   — | j         �+| j         j        �!t          d|› d| j         j        › �¦  «        ‚dS dS )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  r1   )r&   r  s     r    Ú_check_alivezCBGenerateManager._check_aliveo  sN   € ð Œ8Ð D¤HÔ$8Ð$DÝ#Øc¸zÐcÐcÈTÌXÔMaÐcÐcñô ð ð  ÐÐ$DÐ$Dr   ri   rT  rU  r÷   r—   c           	      ó¦  ‡‡— | j         }|€t          d¦  «        ‚|                      |¦  «         t          j        ¦   «         }	t          j        ¦   «         Š|d         }
|                     |
|d|j        |j        ¬¦  «        }t          |d|¦  «        j
        }t          | j         |||	‰||¬¦  «        Šˆˆfd„}|                     ||¦  «         ‰‰fS )	zFStart streaming CB generation. Registers a per-request output handler.Nú3CB manager not initialized. Call `init_cb()` first.r”   T)r  Ú	streamingÚmax_new_tokensÚeos_token_idrU   rb  c                 óž  •— 	 ‰                      | ¦  «         | j        �=‰                     t          | j        ¦  «        ¦  «         ‰                     ¦   «          d S |                      ¦   «         r‰                     ¦   «          d S d S # t          $ r:}‰                     t          t          |¦  «        ¦  «        ¦  «         Y d }~d S d }~ww xY wr%   )r  Úerrorr  r"   r€   Úis_finishedr?  r)   )r#  rD  r¨   Ú
text_queues     €€r    Ú
_on_outputz8CBGenerateManager.generate_streaming.<locals>._on_output£  sÕ   ø€ ð<Ø—’˜VÑ$Ô$Ð$ð ”<Ð+Ø×)Ò)­,°v´|Ñ*DÔ*DÑEÔEÐEØ—L’L‘N”N�N�N�NØ×'Ò'Ñ)Ô)ð #Ø—L’L‘N”N�N�N�Nð#ð #øåð <ð <ð <Ø×%Ò%¥lµ3°q±6´6Ñ&:Ô&:Ñ;Ô;Ð;Ð;Ð;Ð;Ð;Ð;Ð;øøøøð<øøøs   ƒAB Á(B Â
CÂ/CÃC)r   ÚRuntimeErrorrƒ  r  rJ  r   Úadd_requestr‡  rˆ  rc   rü   r  Úregister_result_handler)r&   rR   ri   rU  rO  r  r÷   r—   Úcbrô   r”   ri  r�  r¨   rŒ  s                @@r    rW  z$CBGenerateManager.generate_streamingz  s  øø€ ð ŒXˆØˆ:ÝÐTÑUÔUÐUØ×Ò˜*Ñ%Ô%Ð%åÔ'Ñ)Ô)ˆÝ$+¤M¡O¤Oˆ
à˜;Ô'ˆ	Ø—^’^ØØ!ØØ%Ô4Ø#Ô0ð $ñ 
ô 
ˆ
õ ! ¨K¸ÑCÔCÔNˆÝØŒHØØØØØ#Ø-ð
ñ 
ô 
ˆð	<ð 	<ð 	<ð 	<ð 	<ð 	<ð 	×"Ò" :¨zÑ:Ô:Ð:Ø˜8Ð#Ð#r   c              ƒ   ó6  ‡K  — | j         }|€t          d¦  «        ‚|                      |¦  «         |d         }t          |¦  «        }t	          j        ¦   «         }	|	                     ¦   «         Šˆfd„}
|                     ||
¦  «         |                     |||j	        d|j
        ¬¦  «         ‰ƒ d{V —†}|j        �;|j        �t          d|› d|j        › �¦  «        ‚t          d	|› d|j        › �¦  «        ‚|j        }|                     |d
¬¦  «        }|||fS )zcRun non-streaming CB generation. Registers a handler that resolves an asyncio.Future on completion.Nr…  r”   c                 ó^   •— ‰                      ¦   «         s‰                     | ¦  «         d S d S r%   )Údoner>  )rC  rB  s    €r    Ú
_on_resultz<CBGenerateManager.generate_non_streaming.<locals>._on_resultÉ  s7   ø€ Ø—;’;‘=”=ð *Ø×!Ò! &Ñ)Ô)Ð)Ð)Ð)ð*ð *r   F)r  r‡  r†  rˆ  zCB worker died during request r‚  zCB generation failed for Trk  )r   rŽ  rƒ  rž   r  rJ  rK  r�  r�  r‡  rˆ  rŠ  r  r1   r%  rm  )r&   rR   ri   rU  rO  r  r‘  r”   rp  rô   r•  rC  r{   r  rB  s                 @r    rY  z(CBGenerateManager.generate_non_streaming´  se  øè è € ð ŒXˆØˆ:ÝÐTÑUÔUÐUØ×Ò˜*Ñ%Ô%Ð%à˜;Ô'ˆ	Ý˜	‘N”Nˆ	õ Ô'Ñ)Ô)ˆØ×#Ò#Ñ%Ô%ˆð	*ð 	*ð 	*ð 	*ð 	*ð 	×"Ò" :¨zÑ:Ô:Ð:à
�ŠØØ!Ø%Ô4ØØ#Ô0ð 	ñ 	
ô 	
ð 	
ð ������ˆð Œ<Ð#ØŒ~Ð)Ý'Ð(eÈÐ(eÐ(eÐW]ÔWcÐ(eÐ(eÑfÔfÐfÝÐW¸:ÐWÐWÈÌÐWÐWÑXÔXÐXØÔ/ˆØ×Ò À4ÐÑHÔHˆØ�Y Ð-Ð-r   r   c                 óh   — | j         �| j         j        €t          d¦  «        ‚| j         j        j        S )z*The CB scheduler (for testing/monitoring).Nz.Continuous batching processor not initialized.)r   Úbatch_processorrŽ  Ú	schedulerr  s    r    r˜  zCBGenerateManager.schedulerå  s3   € ð Œ8Ð˜tœxÔ7Ð?ÝÐOÑPÔPÐPØŒxÔ'Ô1Ð1r   c                 óP   — | j         �| j                              dd¬¦  «         d S d S )NTé   )ÚblockÚtimeout)r   r[  r  s    r    r[  zCBGenerateManager.stopì  s0   € ØŒ8ÐØŒH�MŠM ¨aˆMÑ0Ô0Ð0Ð0Ð0ð  Ðr   r%   )rw  rx  r\  r*  )rS   r   rÎ   )r   r   r   r(   r'   rR  r  r€  r)   rƒ  ry   r]  r  r   r  rW  rÏ   rz   rY  Úpropertyr˜  r[  r   r   r    rv  rv  G  s¤  € € € € € ðð ð$ð $ð $ð $ð $ðð ð ð ð"@˜$ð @ð @ð @ð @ð	 sð 	¨tð 	ð 	ð 	ð 	ð$ $(Ø(,ð8$ð 8$à ð8$ð >ð8$ð ð	8$ð
 'ð8$ð ð8$ð ˜D‘[ð8$ð  ™+ð8$ð 
ˆwŒ}˜jÐ(Ô	)ð8$ð 8$ð 8$ð 8$ðt/.à ð/.ð >ð/.ð ð	/.ð
 'ð/.ð ð/.ð 
ˆs�C˜˜cœÐ"Ô	#ð/.ð /.ð /.ð /.ðb ð2ð 2ð 2ñ „Xð2ð1ð 1ð 1ð 1ð 1ð 1r   rv  c                   ól   — e Zd ZdZ	 	 	 ddededdfd„Zd	d
dedefd„Zddedede	fd„Z
dd„Zdefd„ZdS )ÚGenerationStatea'  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.
    FNÚcontinuous_batchingÚcompilerw  rx  c                 óZ   — || _         || _        || _        i | _        d | _        d | _        d S r%   )Ú_continuous_batchingÚ_compilerz  Ú_generate_managersÚ_cb_managerÚ_cb_model_id)r&   r   r¡  rw  s       r    r'   zGenerationState.__init__ÿ  s8   € ð %8ˆÔ!ØˆŒØ#ˆŒØ>@ˆÔØ59ˆÔØ(,ˆÔÐÐr   rR   r   ÚmodalityrS   c                 óª   — | j         sdS t          |d¦  «        o|t          j        k    }|s't                               |j        j        › d�¦  «         |S )aW  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   )r&   rR   r¨  Úcans       r    Úuse_continuous_batchingz'GenerationState.use_continuous_batching  sn   € ð Ô(ð 	Ø�5Ý�eÐ7Ñ8Ô8ÐU¸XÍÌÒ=UˆØð 	Ý×ÒØ”?Ô+ð 9ð 9ð 9ñô ð ð ˆ
r   r®   Úuse_cbc                 ó   — |ra| j         |k    r'| j        � | j                             ¦   «          d| _        | j        €!t          | j        ¬¦  «        | _        || _         | j        S || j        vrt          ¦   «         | j        |<   | j        |         S )af  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)rw  )r§  r¦  r[  rv  rz  r¥  r_  )r&   r®   r®  s      r    Úget_managerzGenerationState.get_manager   sž   € ð ð 	$ØÔ  HÒ,Ð,ØÔ#Ð/ØÔ$×)Ò)Ñ+Ô+Ð+Ø'+�DÔ$ØÔÐ'Ý#4¸t¼Ð#OÑ#OÔ#O�Ô Ø$,�Ô!ØÔ#Ð#Ø˜4Ô2Ð2Ð2Ý0?Ñ0AÔ0AˆDÔ# HÑ-ØÔ& xÔ0Ð0r   c                 óX   — | j         �"| j                              ¦   «          d| _         dS dS )z$Stop any active generation managers.N)r¦  r[  r  s    r    ÚshutdownzGenerationState.shutdown7  s6   € àÔÐ'ØÔ×!Ò!Ñ#Ô#Ð#Ø#ˆDÔÐÐð (Ð'r   c                 óF   — | j         du p| j                              ¦   «         S )zTWhether the CB worker is healthy. ``True`` if CB is disabled or not yet initialized.N)r¦  r€  r  s    r    Úis_cb_alivezGenerationState.is_cb_alive=  s$   € àÔ 4Ð'ÐF¨4Ô+;×+DÒ+DÑ+FÔ+FÐFr   )FFN©FrÎ   )r   r   r   r(   r  r'   r   r­  r)   rN  r°  r²  r´  r   r   r    rŸ  rŸ  ñ  sî   € € € € € ðð ð %*ØØ7;ð	-ð -à!ð-ð ð-ð 5ð	-ð -ð -ð -ðÐ->ð È(ð ÐW[ð ð ð ð ð(1ð 1 Cð 1°ð 1ÐBUð 1ð 1ð 1ð 1ð.$ð $ð $ð $ðG˜Tð Gð Gð Gð Gð Gð Gr   rŸ  c            	       ó  — e Zd ZU dZdZedz  ed<    e¦   «         Zee	         ed<   	 dddde
dedz  fd	„Zd
eddfd„Zeddde	fd„¦   «         Zd
edee	ddf         fd„Z	 dd
edddeddfd„Zedee         dedee         fd„¦   «         ZdS )ÚBaseHandlera®  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.
    NÚ_valid_params_classÚ_unused_fieldsÚmodel_managerr   Úgeneration_stateÚchat_template_kwargsc                 ó4   — || _         || _        |pi | _        d S r%   )rº  r»  r¼  )r&   rº  r»  r¼  s       r    r'   zBaseHandler.__init__R  s'   € ð +ˆÔØ 0ˆÔØ$8Ð$>¸BˆÔ!Ð!Ð!r   ÚbodyrS   c                 ó&  — ddl m} t          |                     ¦   «         ¦  «        }| j        �7|t          | j        dt          ¦   «         ¦  «        z
  }|r |dd|› �¬¦  «        ‚|| j        z  }|rt                               d|› �¦  «         dS dS )	zMValidate request fields against the handler's params class and unused fields.r   ©ÚHTTPExceptionNÚ__mutable_keys__i¦  z"Unexpected fields in the request: ©Ústatus_codeÚdetailz,Ignoring unsupported fields in the request: )	ÚfastapirÁ  r  Úkeysr¸  rc   r¹  rª  r«  )r&   r¾  rÁ  Ú
input_keysÚ
unexpectedÚunuseds         r    Ú_validate_requestzBaseHandler._validate_request\  sÀ   € à)Ð)Ð)Ð)Ð)Ð)å˜Ÿš™œÑ%Ô%ˆ
ØÔ#Ð/Ø#¥g¨dÔ.FÐHZÕ\_Ñ\aÔ\aÑ&bÔ&bÑbˆJØð oØ#�m°Ð<mÐakÐ<mÐ<mÐnÑnÔnÐnØ˜dÔ1Ñ1ˆØð 	YÝ×ÒÐ WÈvÐ WÐ WÑXÔXÐXÐXÐXð	Yð 	Yr   Úchunkzstr | pydantic.BaseModelc                 óš   — t          | t          ¦  «        r|                      d¦  «        r| nd| › d�S d|                      d¬¦  «        › d�S )z;Format a pydantic model or string as an SSE ``data:`` line.zdata: z

T)Úexclude_none)rp   r)   Ú
startswithÚmodel_dump_json)rÌ  s    r    Úchunk_to_ssezBaseHandler.chunk_to_ssei  sb   € õ �e�SÑ!Ô!ð 	QØ!×,Ò,¨XÑ6Ô6ÐP�5�5Ð<PÀUÐ<PÐ<PÐ<PÐPØF˜×-Ò-¸4Ð-Ñ@Ô@ÐFÐFÐFÐFr   r   rT  c                 óR  — ddl m} | j        j        �T|                     d¦  «        }|�.|| j        j        k    r |dd| j        j        › d|› d�¬	¦  «        ‚| j        j        |d<   | j                             |d         ¦  «        }| j                             |¦  «        \  }}|||fS )
zfApply force_model, load model + processor.

        Returns ``(model_id, model, processor)``.
        r   rÀ  NrR   i�  zServer is pinned to 'z'; requested 'z'.rÃ  )rÆ  rÁ  rº  Úforce_modelrd   Úprocess_model_nameÚload_model_and_processor)r&   r¾  rÁ  Ú	requestedr®   rR   ri   s          r    Ú_resolve_modelzBaseHandler._resolve_modelp  sÕ   € ð
 	*Ð)Ð)Ð)Ð)Ð)àÔÔ)Ð5ØŸš Ñ)Ô)ˆIØÐ$¨°dÔ6HÔ6TÒ)TÐ)TØ#�mØ #Øo°DÔ4FÔ4RÐoÐoÐbkÐoÐoÐoðñ ô ð ð !Ô.Ô:ˆD�‰MàÔ%×8Ò8¸¸g¼ÑGÔGˆØÔ-×FÒFÀxÑPÔPÑˆˆyà˜ 	Ð)Ð)r   FÚmodel_generation_configr   r®  c                 ón  — ddl m} |                     d¦  «        �! |di t          j        |d         ¦  «        ¤Ž}n-t          j        |¦  «        }|j        �|j        dk     rd|_        |                     d¦  «        �:t          |d         ¦  «        |_	        t          |d         ¦  «        dk    rd|_
        |                     d	¦  «        �t          |d	         ¦  «        |_        |                     d
¦  «        �t          |d
         ¦  «         | j        j        r|j        €d|_        |rd|_        |S )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   rc  Ni   Útemperatureg        FÚtop_pr+  Ústaticr   )Útransformersr   rd   r3   ÚloadsÚcopyÚdeepcopyr‡  ÚfloatrÚ  Ú	do_samplerÛ  r/  r»  r¤  Úcache_implementationÚ	use_cache)r&   r¾  rØ  r®  r   rc  s         r    Ú_build_generation_configz$BaseHandler._build_generation_config…  sV  € ð( 	2Ð1Ð1Ð1Ð1Ð1à�8Š8Ð'Ñ(Ô(Ð4Ø 0Ð 0Ð YÐ Yµ4´:¸dÐCVÔ>WÑ3XÔ3XÐ YÐ YÐÐå $¤Ð.EÑ FÔ FÐØ Ô/Ð7Ð;LÔ;[Ð^bÒ;bÐ;bØ37Ð!Ô0à�8Š8�MÑ"Ô"Ð.Ý,1°$°}Ô2EÑ,FÔ,FÐÔ)Ý�T˜-Ô(Ñ)Ô)¨SÒ0Ð0Ø.3Ð!Ô+Ø�8Š8�GÑÔÐ(Ý&+¨D°¬MÑ&:Ô&:ÐÔ#Ø�8Š8�FÑÔÐ'Ý˜4 œ<Ñ(Ô(Ð(ð Ô Ô)ð 	>Ð.?Ô.TÐ.\Ø5=ÐÔ2ð ð 	0Ø*/ÐÔ'ð !Ð r   Úmessagesr¨  c                 ó¼  — g }| D �]Õ}|d         g dœ}d|v rŠg }|d         D ]z}t          j        |¦  «        }|                     d¦  «        p|}t          |d         t          ¦  «        rt          j        |d         ¦  «        |d<   |                     |¦  «         Œ{||d<   d|v r|d         |d<   d|v rg n|                     d¦  «        pg }t          |t          ¦  «        rd|d	œg}|D �]¡}	|	d
         }
|
dv r%|d                              d|	d         d	œ¦  «         Œ4|
dv r^|t          j	        t          j
        fv rD|	d         }t          |t          ¦  «        r|d         }|d                              d|dœ¦  «         Œ–|
dk    ry|t          j
        k    ri|	d         }t          |t          ¦  «        r*|d         }|                     d¦  «        }|rd|› d|› �n|}n|}|d                              d|dœ¦  «         �Œ|
dk    rF|t          j	        t          j
        fv r,|d                              d|	d         d         dœ¦  «         �Œa|
dk    r:|t          j
        k    r*|d                              d|	d         d         dœ¦  «         �Œ£|t          j        k    r(d                     d„ |d         D ¦   «         ¦  «        |d<   |                     |¦  «         �Œ×|S )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.
        Úrole)rè  r7   r2   ro   rD   Útool_call_idr7   r  )r@   r  r@   )r  Ú
input_textÚoutput_text)Ú	image_urlÚinput_imagerì  ÚurlÚimage)r@   rî  Úinput_audioÚdataÚformatzdata:audio/z;base64,ÚaudioÚ	video_urlÚvideoÚ	audio_urlú c              3   ó&   K  — | ]}|d          V — ŒdS )r  Nr   )r\   r¿   s     r    r`   zABaseHandler.get_processor_inputs_from_messages.<locals>.<genexpr>ý  s&   è è € Ð,RÐ,R¸1¨Q¨v¬YÐ,RÐ,RÐ,RÐ,RÐ,RÐ,Rr   )rß  rà  rd   rp   r)   r3   rÞ  r§   r   r   r   ry   r   Újoin)ræ  r¨  Úprocessor_inputsÚmessager|   r2   ÚtcrA  Úraw_contentr7   Úcontent_typerî  rð  Ú	audio_b64Úfmts                  r    Ú"get_processor_inputs_from_messagesz.BaseHandler.get_processor_inputs_from_messages·  si  € ð Ðàð 7	,ñ 7	,ˆGØ% fœo¸"Ð=Ð=ˆFð ˜wÐ&Ð&Ø�
Ø! ,Ô/ð *ð *�BÝœ rÑ*Ô*�BØŸš 
Ñ+Ô+Ð1¨r�BÝ! " [¤/µ3Ñ7Ô7ð FÝ*.¬*°R¸´_Ñ*EÔ*E˜˜;™Ø×%Ò% bÑ)Ô)Ð)Ð)Ø'1��|Ñ$Ø Ð(Ð(Ø)0°Ô)@��~Ñ&ð !-°Ð 7Ð 7˜"˜"¸g¿kºkÈ)Ñ>TÔ>TÐ>ZÐXZˆKÝ˜+¥sÑ+Ô+ð FØ(.¸ÐDÐDÐE�à&ð dñ d�Ø& vœ�àÐ#HÐHÐHØ˜9Ô%×,Ò,°fÀgÈfÄoÐ-VÐ-VÑWÔWÐWÐWà!Ð%AÐAÐAÀhÕS[ÔS_ÕaiÔatÐRuÐFuÐFuà! +Ô.�CÝ! #¥tÑ,Ô,ð )Ø! %œj˜Ø˜9Ô%×,Ò,°gÀcÐ-JÐ-JÑKÔKÐKÐKð " ]Ò2Ð2°xÅ8ÔCVÒ7VÐ7VØ")¨-Ô"8�KÝ! +­tÑ4Ô4ð *Ø$/°Ô$7˜	Ø)Ÿošo¨hÑ7Ô7˜ØHKÐZÐD¨CÐDÐD¸ÐDÐDÐDÐQZ˜˜à)˜Ø˜9Ô%×,Ò,°gÀcÐ-JÐ-JÑKÔKÐKÑKà! [Ò0Ð0°XÅ(Ä,ÕPXÔPcÐAdÐ5dÐ5dØ˜9Ô%×,Ò,°gÀgÈkÔFZÐ[`ÔFaÐ-bÐ-bÑcÔcÐcÑcØ! [Ò0Ð0°XÅÔATÒ5TÐ5TØ˜9Ô%×,Ò,°gÀgÈkÔFZÐ[`ÔFaÐ-bÐ-bÑcÔcÐcùð �8œ<Ò'Ð'Ø$'§H¢HÐ,RÐ,RÀÀyÔ@QÐ,RÑ,RÔ,RÑ$RÔ$R��yÑ!à×#Ò# FÑ+Ô+Ð+Ñ+ØÐr   r%   rµ  )r   r   r   r(   r¸  r@   Ú__annotations__r  r¹  r)   rŸ  ry   r'   rË  ÚstaticmethodrÑ  r]  r×  r  rå  rz   r   r  r   r   r    r·  r·  B  s¦  € € € € € € ð
ð 
ð (,Ð˜ ™Ð+Ð+Ñ+Ø"˜s™uœu€N�C˜”HÐ$Ð$Ñ$ð -1ð	?ð ?à%ð?ð *ð?ð # T™kð	?ð ?ð ?ð ?ðY dð Y¨tð Yð Yð Yð Yð ðGÐ6ð G¸3ð Gð Gð Gñ „\ðGð* 4ð *¨E°#Ð7HÐJtÐ2tÔ,uð *ð *ð *ð *ð, W\ð0!ð 0!Øð0!Ø3Eð0!ØOSð0!à	ð0!ð 0!ð 0!ð 0!ðd ðH °T¸$´Zð H È8ð H ÐX\Ð]aÔXbð H ð H ð H ñ „\ðH ð H ð H r   r·  r%   rÎ   )Mr(   r  rß  Úenumr3   r  Úabcr   r   Úcollections.abcr   Úconcurrent.futuresr   rõ   r   Útypingr   Útransformers.utilsr	   ÚpydanticÚ
tokenizersr-  rÝ  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"   r?  r+   r)   r.   rŽ  r1   rg   ry   rl   rr   rz   r}   r‘   r�   r–   r]  rš   rÏ   r  r“   rª   r¬   r@   rð   rò   r  r/  r4  r6  rN  r_  rv  rŸ  r·  r   r   r    ú<module>r     s¼  ððð ð €€€Ø €€€Ø €€€Ø €€€Ø Ð Ð Ð Ø #Ð #Ð #Ð #Ð #Ð #Ð #Ð #Ø $Ð $Ð $Ð $Ð $Ð $Ø %Ð %Ð %Ð %Ð %Ð %Ø Ð Ð Ð Ð Ð Ø  Ð  Ð  Ð  Ð  Ð  à &Ð &Ð &Ð &Ð &Ð &ð ð ,Ø€O€O€OØÐÐÐØ€L€L€Lðð ð ð ð ð ð ð ð ð ð ð ð ð ð eÐdÐdÐdÐdÐdØUÐUÐUÐUÐUÐUØOÐOÐOÐOÐOÐOà+Ð+Ð+Ð+Ð+Ð+ð 
ˆÔ	˜HÑ	%Ô	%€ð €ðð ð ð ð ˆtŒyñ ô ð ðð ð ð ð ñ ô ð ðPð Pð Pð Pð P˜9ñ Pô Pð Pðð ð ð ð �Cñ ô ð ðð ð ð ð ˜ñ ô ð ð0 ØàØ5àØ)Ø+Ø#Ø%ð	ð ðð
ð 
ðð ð( Øà \Øà à# XÐ.à (Ø-lð"ð "ðð ð	ð 	ð
ð 
ð!ð !ð?1ð 1Ð ðh(BÐ+<ð (BÀÈÁð (Bð (Bð (Bð (BðV Dð ¨Tð ð ð ð ð".°tð .ÀÀTÄ
ÈTÑ@Qð .ð .ð .ð .ð6 ˆ[Øàà Ð*Ø Ð)ð
ð 
ð bð
ð 
ðð Ð ð, Ð7Ð7Ð7ÀÐMÐMð	Ð ðð Ð+<ð ÐQUÐX\ÑQ\ð ð ð ð ðD°sð Èdð ÐW\Ð]`ÐbeÐhlÑblÐ]lÔWmð ð ð ð ð(¨d°3¬ið ¸Dð ð ð ð ð8°ð ¸ð ð ð ð ð2(
ð (
ð (
ð (
ð (
ñ (
ô (
ð (
ðVD xð D¸3ð DÀ4ð Dð Dð Dð DðNRð Rð Rð Rð Rñ Rô Rð RðjL2ð L2ð L2ð L2ð L2ñ L2ô L2ð L2ð^˜ð  ð ð ð ð ð!ð !ð !ð !ð&ð &ð &ð &ð &ñ &ô &ð &ðRA>ð A>ð A>ð A>ð A>˜#ñ A>ô A>ð A>ðHDð Dð Dð Dð DÐ)ñ Dô Dð DðNg1ð g1ð g1ð g1ð g1Ð+ñ g1ô g1ð g1ðTNGð NGð NGð NGð NGñ NGô NGð NGðb~ ð ~ ð ~ ð ~ ð ~ ñ ~ ô ~ ð ~ ð ~ ð ~ r   