Ë
    GêñiÔœ  ã                   ó¦  — 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jH                  e%«      Z&dZ' G d„ dejP                  «      Z) G d„ d«      Z* G d„ de+«      Z,ddddddddœdœdœiZ-d d!d"e.dz  fd#„Z/d$e.d"e.fd%„Z0d&e.d"e1e.   dz  fd'„Z2 G d(„ d)«      Z3d*e
d+e4d"e5fd,„Z6 G d-„ d.«      Z7 G d/„ d0«      Z8d1e9d"dfd2„Z:d@d3„Z; G d4„ d5«      Z< G d6„ d7e«      Z= G d8„ d9e=«      Z> G d:„ d;e=«      Z? G d<„ d=«      Z@ G d>„ d?«      ZAy)Az?
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y)ÚModalityÚLLMÚVLMÚ
MULTIMODALÚSTTÚTTSN)Ú__name__Ú
__module__Ú__qualname__r   r   r   r   r   © ó    ú`/var/www/pod-logistic/pod-ai/venv/lib/python3.12/site-packages/transformers/cli/serving/utils.pyr   r   9   s   „ Ø
€CØ
€CØ€JØ
€CØ
�Cr   r   c                   ó   — e Zd ZdZdefd„Zy)Ú_StreamErrorz5Sentinel to signal an error from the generate thread.Úmsgc                 ó   — || _         y ©N)r#   )Úselfr#   s     r    Ú__init__z_StreamError.__init__D   s	   € Øˆ�r   N)r   r   r   Ú__doc__Ústrr'   r   r   r    r"   r"   A   s   „ Ù?ð˜Cô r   r"   c                   ó   — e Zd ZdZy)Ú_GenerationCancelledzERaised inside ``DirectStreamer.put()`` to abort ``model.generate()``.N)r   r   r   r(   r   r   r    r+   r+   H   s   „ ÚOr   r+   Úqwenz<tool_call>z</tool_call>z<tool_call>(.*?)</tool_call>ÚarrayÚobjectÚjson)Útypezx-parser)zx-regex-iteratorr0   Úitems)ÚstcÚetcÚschemaÚmodelr   Úreturnc                 óL  ‡— t        | d| «      }t        |dd«      }t        |dd«      }t        |dd«      }|r|r|r	|d   d   }n9t        ˆfd„t        j                  «       D «       d«      }|€y|d	   |d
   |d   }}}|j	                  |«      }|j	                  |«      }	|||	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_schemaÚ
propertiesÚ
tool_callsc              3   óZ   •K  — | ]"  \  }}|‰j                   j                  v sŒ|–— Œ$ y ­wr%   )ÚconfigÚ
model_type)Ú.0ÚkÚvr5   s      €r    ú	<genexpr>z'get_tool_call_config.<locals>.<genexpr>o   s&   øè ø€ Òd™t˜q !ÀqÈEÏLÉL×LcÑLcÒGcœÑdùs   ƒ +¤+r2   r3   r4   )r4   Ústc_idÚetc_id)ÚgetattrÚnextÚ_TOOL_CALL_FALLBACKSr1   Úconvert_tokens_to_ids)
Ú	processorr5   r8   r2   r3   r;   r4   ÚfallbackrE   rF   s
    `        r    Úget_tool_call_configrM   ]   sÅ   ø€ ô ˜	 ;°	Ó:€IÜ
�)˜[¨$Ó
/€CÜ
�)˜[¨$Ó
/€CÜ˜iÐ):¸DÓA€Oñ ‰s‘Ø  Ñ.¨|Ñ<‰ô ÓdÔ';×'AÑ'AÓ'CÔdÐfjÓkˆØÐØØ# E™?¨H°U©O¸XÀhÑ=O�&ˆSˆà×,Ñ,¨SÓ1€FØ×,Ñ,¨SÓ1€FØ¨¸&ÑAÐAr   Ú	tool_callc                 ó¨   — | j                  d| «      }|j                  di «      }|d   t        |t        «      st        j                  |«      dœS |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.
    ÚfunctionÚ	argumentsÚname)rR   rQ   )ÚgetÚ
isinstancer)   r/   Údumps)rN   rP   rQ   s      r    Ú_normalize_tool_callrV   y   sX   € ð �}‰}˜Z¨Ó3€HØ—‘˜[¨"Ó-€Ià˜Ñ Ü2<¸YÌÔ2L”T—Z‘Z 	Ó*ñð àR[ñð r   r4   c                 ó˜   — | j                  ||«      }|syt        |t        «      s|g}|D �cg c]  }t        |«      ‘Œ }}|r|S dS c c}w )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_responserT   ÚlistrV   )rK   Úgenerated_idsr4   ÚparsedrN   r=   s         r    Úparse_tool_callsr\   Š   sY   € ð ×%Ñ% m°VÓ<€FÙØÜ�fœdÔ#Ø�ˆØCIÖJ°iÔ& yÕ1ÐJ€JÐJÙ#ˆ:Ð-¨Ð-ùò Ks   ­Ac                   óp   — 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y)Ú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 | _        y r%   )r_   r5   ÚbarsÚlast_emitted_current)r&   r_   r`   s      r    r'   zDownloadAggregator.__init__¦   s   € ØˆŒØˆŒ
Ø79ˆŒ	Ø04ˆÕ!r   Úbar_idÚtotalNr6   c                 óF   — d|f| j                   |<   | j                  «        y)z6Register a new download bar with its total byte count.r   N©rb   Ú_emit)r&   rd   re   s      r    ÚregisterzDownloadAggregator.register¬   s   € à ˜Jˆ�	‰	�&ÑØ�
‰
�r   Úcurrentc                 óF   — ||f| j                   |<   | j                  «        y)z>Update a bar's current byte count and emit aggregate progress.Nrg   )r&   rd   rj   re   s       r    ÚupdatezDownloadAggregator.update±   s   € à$ eÐ,ˆ�	‰	�&ÑØ�
‰
�r   c                  ó   — y r%   r   )r&   rd   s     r    ÚclosezDownloadAggregator.close¶   ó   € Ør   c                 óT  — t        d„ | j                  j                  «       D «       «      }|| j                  k(  ry || _        | j                  j                  «       D ��cg c]
  \  }}|€Œ	|‘Œ }}}|rt        |«      nd }| j	                  d| j
                  d||dœdœ«       y c c}}w )Nc              3   ó&   K  — | ]	  \  }}|–— Œ y ­wr%   r   )rA   ÚcÚ_s      r    rD   z+DownloadAggregator._emit.<locals>.<genexpr>º   s   è ø€ Ò;¡  1œ!Ñ;ùó   ‚ÚloadingÚdownload©rj   re   ©Ústatusr5   ÚstageÚprogress)Úsumrb   Úvaluesrc   r_   r5   )r&   Úagg_currentrs   ÚtÚtotalsÚ	agg_totals         r    rh   zDownloadAggregator._emit¹   s™   € ÜÑ;¨¯	©	×(8Ñ(8Ó(:Ô;Ó;ˆØ˜$×3Ñ3Ò3ØØ$/ˆÔ!Ø $§	¡	× 0Ñ 0Ó 2×D™˜˜1°a±m’!ÐDˆÑDÙ#)”C˜”K¨tˆ	Ø�‰à#ØŸ™Ø#Ø(3¸iÑHñ	õ	
ùó Es   Á
B$Á*B$©r6   N)r   r   r   r(   r   r)   r'   Úintri   rl   rn   rh   r   r   r    r^   r^   Ÿ   su   „ ñð5 ð 5°Có 5ð˜sð ¨3°©:ð ¸$ó ð
˜Sð ¨3ð °s¸T±zð Àdó ð
˜Cð  Dó ô
r   r^   Úcallbackr`   c                 óN   ‡ ‡‡— 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                 ó
  •— |j                  d«      xs d| _        d|d<   t        ‰| �  |i |¤Ž d| _        d| _        | j                  dk(  r7t        | «      | _        ‰j                  | j                  | j                  «       y y )NÚunitÚitTÚdisabler   éÿÿÿÿÚB)
rS   Ússe_unitÚsuperr'   ÚnÚlast_emittedÚidÚ_bar_idri   re   )r&   ÚargsÚkwargsÚ	__class__Údownload_aggregators      €€r    r'   z7make_progress_tqdm_class.<locals>.ProgressTqdm.__init__Ý   sw   ø€ Ø"ŸJ™J vÓ.Ò6°$ˆDŒMØ $ˆF�9ÑÜ‰GÑ˜dÐ- fÒ-ØˆDŒFØ "ˆDÔØ�}‰} Ò#Ü! $›x�”Ø#×,Ñ,¨T¯\©\¸4¿:¹:ÕFð $r   c                 óX  •— |€d}| xj                   |z  c_         | j                  dk(  r2‰j                  | j                  | j                   | j                  «       y | j                   | j
                  k7  r6| j                   | _         ‰d‰d| j                   | j                  dœdœ«       y y ©Nr   rŽ   ru   Úweightsrw   rx   )r‘   r�   rl   r”   re   r’   )r&   r‘   r„   r˜   r`   s     €€€r    rl   z5make_progress_tqdm_class.<locals>.ProgressTqdm.updateç   s�   ø€ ØˆyØ�Ø�FŠF�a‰K�FØ�}‰} Ò#Ø#×*Ñ*¨4¯<©<¸¿¹ÀÇÁÕLØ—‘˜4×,Ñ,Ò,Ø$(§F¡F�Ô!Ùà"+Ø!)Ø!*Ø04·±ÀÇÁÑ$Lñ	õð -r   c           	   3   ó€  •K  — | j                   D ]ª  }| xj                  dz  c_        | j                  dk(  r2‰j                  | j                  | j                  | j
                  «       nN| j                  | j                  k7  r5| j                  | _         ‰d‰d| j                  | j
                  dœdœ«       |–— Œ¬ y ­wrš   )Úiterabler‘   r�   rl   r”   re   r’   )r&   Úitemr„   r˜   r`   s     €€€r    Ú__iter__z7make_progress_tqdm_class.<locals>.ProgressTqdm.__iter__ø   sœ   øè ø€ ØŸ™ò �Ø—’˜!‘•Ø—=‘= CÒ'Ø'×.Ñ.¨t¯|©|¸T¿V¹VÀTÇZÁZÕPØ—V‘V˜t×0Ñ0Ò0Ø(,¯©�DÔ%Ùà&/Ø%-Ø%.Ø48·F±FÀTÇZÁZÑ(Pñ	ôð “
ñùs   ƒB;B>c                 óv   •— | j                   dk(  r‰j                  | j                  «       t        ‰| �  «        y )NrŽ   )r�   rn   r”   r�   )r&   r—   r˜   s    €€r    rn   z4make_progress_tqdm_class.<locals>.ProgressTqdm.close	  s*   ø€ Ø�}‰} Ò#Ø#×)Ñ)¨$¯,©,Ô7Ü‰G‰M�Or   )r   )r   r   r   r'   rl   rŸ   rn   Ú__classcell__)r—   r„   r˜   r`   s   @€€€r    ÚProgressTqdmrˆ   Ü   s   ù„ õ	G÷	ö"	÷"	ñ 	r   r¢   )Ú	tqdm.autor†   r^   )r„   r`   Ú	base_tqdmr¢   r˜   s   ``  @r    Úmake_progress_tqdm_classr¥   Ê   s/   ú€ õ ,ä,¨X°xÓ@Ð÷0ð 0�yô 0ðd Ðr   c                   óx   — e Zd ZdZ	 	 ddddej
                  dej                  dededz  f
d	„Z	dd
„Z
dd„Zdd„Zy)Ú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.
    Nr8   útokenizers.TokenizerÚloopÚqueueÚskip_special_tokensÚtool_configc                 óø   — ddl m} || _        || _        || _         |g |«      | _        |r|d   nd| _        |r|d   nd| _        d| _        d| _	        t        j                  «       | _        d| _        g | _        y)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.
        r   ©ÚDecodeStreamrE   NrF   FT)Útokenizers.decodersr¯   Ú
_tokenizerÚ_loopÚ_queueÚ_decode_streamÚ_stc_idÚ_etc_idÚ_inside_tool_callÚ_firstÚ	threadingÚEventÚ
_cancelledÚtotal_tokensÚgenerated_token_ids)r&   r8   r©   rª   r«   r¬   r¯   s          r    r'   zDirectStreamer.__init__  sy   € õ& 	5à#ˆŒØˆŒ
ØˆŒÙ*¨2Ð/BÓCˆÔÙ0;�{ 8Ò,ÀˆŒÙ0;�{ 8Ò,ÀˆŒØ!&ˆÔØˆŒÜ#Ÿ/™/Ó+ˆŒØˆÔØ.0ˆÕ r   c                 óD  — | j                   j                  «       r
t        «       ‚| j                  rd| _        y|j	                  «       D ]Õ  }| xj
                  dz  c_        | j                  j                  |«       || j                  k(  rd| _	        n|| j                  k(  rd| _	        | j                  j                  | j                  |«      }|€Œ‰| j                  rŒ–|| j                  k7  sŒ¦| j                  j                  | j                   j"                  |«       Œ× y)zHCalled by ``model.generate()`` after each decode step with new token(s).FNr   T)r»   Úis_setr+   r¸   Útolistr¼   r½   Úappendrµ   r·   r¶   r´   Ústepr±   r²   Úcall_soon_threadsafer³   Ú
put_nowait)r&   ÚvalueÚtoken_idÚtexts       r    ÚputzDirectStreamer.put;  sã   € à�?‰?×!Ñ!Ô#Ü&Ó(Ð(à�;Š;ØˆDŒKØØŸ™›ò 	NˆHØ×Ò Ñ"ÕØ×$Ñ$×+Ñ+¨HÔ5à˜4Ÿ<™<Ò'Ø)-�Õ&Ø˜TŸ\™\Ò)Ø).�Ô&à×&Ñ&×+Ñ+¨D¯O©O¸XÓFˆDØÑ¨×(>Ó(>À8ÈtÏ|É|ÓC[Ø—
‘
×/Ñ/°·±×0FÑ0FÈÕMñ	Nr   c                 ód   — | j                   j                  | j                  j                  d«       y)z;Called by ``model.generate()`` when generation is complete.N)r²   rÃ   r³   rÄ   ©r&   s    r    ÚendzDirectStreamer.endP  s    € à�
‰
×'Ñ'¨¯©×(>Ñ(>ÀÕEr   c                 ó8   — | j                   j                  «        y)zWSignal cancellation. The next ``put()`` call will raise and abort ``model.generate()``.N)r»   ÚsetrÊ   s    r    ÚcancelzDirectStreamer.cancelT  s   € à�‰×ÑÕr   )TN)rÅ   útorch.Tensorr6   Nr‚   )r   r   r   r(   ÚasyncioÚAbstractEventLoopr   ÚboolÚdictr'   rÈ   rË   rÎ   r   r   r    r§   r§     sd   „ ñð %)Ø#'ñ1à)ð1ð ×'Ñ'ð1ð �}‰}ð	1ð
 "ð1ð ˜D‘[ó1óBNó*Fôr   r§   c                   óz   — e Zd ZdZ	 ddddedddej                  d	ej                  d
edz  fd„Z	dd„Z
dd„Zdd„Zy)Ú
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_idr8   r¨   r©   rª   r¬   c                 óâ   — ddl m} || _        || _        || _        || _        || _         |g d«      | _        |r|d   nd| _        |r|d   nd| _	        d| _
        d| _        d| _        g | _        y)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``).
        r   r®   TrE   NrF   F)r°   r¯   Ú_cbÚ_request_idr²   r³   r±   r´   rµ   r¶   r·   Ú	_prev_lenr¼   r½   )r&   rÖ   r×   r8   r©   rª   r¬   r¯   s           r    r'   zCBStreamer.__init__b  sy   € õ$ 	5àˆŒØ%ˆÔØˆŒ
ØˆŒØ#ˆŒÙ*¨2¨tÓ4ˆÔÙ0;�{ 8Ò,ÀˆŒÙ0;�{ 8Ò,ÀˆŒØ!&ˆÔØˆŒØˆÔØ.0ˆÕ r   c                 óô  — |j                   | j                  d }t        |j                   «      | _        |D ]À  }| xj                  dz  c_        | j                  j                  |«       || j                  k(  rd| _        n|| j                  k(  rd| _        | j                  j                  | j                  |«      }|€Œ‰| j                  rŒ–|| j                  k7  sŒ¦| j                  j                  |«       ŒÂ y)zLDecode new tokens from a CB ``GenerationOutput`` and push text to the queue.Nr   TF)Úgenerated_tokensrÛ   Úlenr¼   r½   rÁ   rµ   r·   r¶   r´   rÂ   r±   r³   rÄ   )r&   ÚoutputÚ
new_tokensrÆ   rÇ   s        r    rÈ   zCBStreamer.putƒ  sÎ   € à×,Ñ,¨T¯^©^Ð-=Ð>ˆ
Ü˜V×4Ñ4Ó5ˆŒØ"ò 	-ˆHØ×Ò Ñ"ÕØ×$Ñ$×+Ñ+¨HÔ5à˜4Ÿ<™<Ò'Ø)-�Õ&Ø˜TŸ\™\Ò)Ø).�Ô&à×&Ñ&×+Ñ+¨D¯O©O¸XÓFˆDØÑ¨×(>Ó(>À8ÈtÏ|É|ÓC[Ø—‘×&Ñ& tÕ,ñ	-r   c                 ó:   — | j                   j                  d«       y)zSignal end of stream.N)r³   rÄ   rÊ   s    r    rË   zCBStreamer.end”  s   € à�‰×Ñ˜tÕ$r   c                 óN   — | j                   j                  | j                  «       y)zCancel the CB request.N)rÙ   Úcancel_requestrÚ   rÊ   s    r    rÎ   zCBStreamer.cancel˜  s   € à�‰×Ñ × 0Ñ 0Õ1r   r%   )rß   r   r6   Nr‚   )r   r   r   r(   r)   rÐ   rÑ   r   rÓ   r'   rÈ   rË   rÎ   r   r   r    rÕ   rÕ   Y  si   „ ñð $(ñ1à/ð1ð ð1ð *ð	1ð
 ×'Ñ'ð1ð �}‰}ð1ð ˜D‘[ó1óB-ó"%ô2r   rÕ   Úseedc                 ó0   — ddl } |j                  | «       y)z8Set the PyTorch random seed for reproducible generation.r   N)ÚtorchÚmanual_seed)rä   ræ   s     r    Úset_torch_seedrè   �  s   € ãà€E×Ñ�dÕr   c                  óv   — ddl } | j                  j                  «       r| j                  j                  «        yy)z+Empty the CUDA cache if a GPU is available.r   N)ræ   ÚcudaÚis_availableÚempty_cache)ræ   s    r    Úreset_torch_cacherí   ¤  s*   € ãà‡z�z×ÑÔ Ø�
‰
×ÑÕ ð !r   c                   óJ   — e Zd ZdZd„ Zdd„Zdefd„Zdej                  fd„Z	y)	Ú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
                  j                  «        y )NT)ÚtargetÚdaemon)r   r³   r¹   ÚThreadÚ_runÚ_threadÚstartrÊ   s    r    r'   zInferenceThread.__init__³  s3   € Ü"›WˆŒÜ ×'Ñ'¨t¯y©yÀÔFˆŒØ�‰×ÑÕr   r6   Nc                 óD  — 	 | j                   j                  «       \  }}}}}	  ||i |¤Ž}|�|j                  |j                  |«       n|j                  |«       ŒZ# t        $ r:}|�|j                  |j
                  |«       n|j                  |«       Y d }~Œ?d }~ww xY wr%   )r³   rS   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à×%Ñ% fÔ-ð øô ò ,ØÐ#Ø×-Ñ-¨f×.BÑ.BÀAÕFà×(Ñ(¨Ô+ÿøð	,ús   £8A Á	BÁ%0BÂBc                 óZ   — t        «       }| j                  j                  ||||df«       |S )úESubmit a callable to the inference thread. Returns a blocking Future.N)r   r³   rÈ   )r&   rû   r•   r–   rü   s        r    ÚsubmitzInferenceThread.submitÇ  s)   € ä›ˆØ�‰�‰˜˜T 6¨6°4Ð8Ô9Øˆr   c                 óŽ   — t        j                  «       }|j                  «       }| j                  j	                  |||||f«       |S ©zOSubmit a callable to the inference thread. Returns an awaitable asyncio.Future.)rÐ   Úget_running_loopÚcreate_futurer³   rÈ   )r&   rû   r•   r–   r©   rü   s         r    Úasync_submitzInferenceThread.async_submitÍ  s>   € ä×'Ñ'Ó)ˆØ×#Ñ#Ó%ˆØ�‰�‰˜˜T 6¨6°4Ð8Ô9Øˆr   r‚   )
r   r   r   r(   r'   rô   r   r  rÐ   r  r   r   r    rï   rï   ¬  s-   „ ñòó
,ð¨Vó ð°7·>±>ô r   rï   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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y)ÚBaseGenerateManageruã   Base class for generation managers.

    Subclasses:
    - :class:`GenerateManager` â€” sequential ``model.generate()`` on a persistent thread.
    - :class:`CBGenerateManager` â€” continuous batching with paged attention.
    r5   r   Ú
gen_configr   r6   Nc                  ó   — y)z:Initialize continuous batching. No-op for non-CB managers.Nr   ©r&   r5   r	  s      r    Úinit_cbzBaseGenerateManager.init_cbÝ  ó   � r   rK   ú(ProcessorMixin | PreTrainedTokenizerFastÚinputsr×   r¬   zDirectStreamer | CBStreamerc                  ó   — y)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.

        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&   r5   rK   r  r	  r×   r¬   s          r    Úgenerate_streamingz&BaseGenerateManager.generate_streamingà  r  r   c              ƒ   ó   K  — y­w)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&   r5   rK   r  r	  r×   s         r    Úgenerate_non_streamingz*BaseGenerateManager.generate_non_streamingû  s   è ø� ùs   ‚c                  ó   — y)z/Stop the generation manager and free resources.Nr   rÊ   s    r    ÚstopzBaseGenerateManager.stop  r  r   ©r5   r   r	  r   r6   Nr%   r‚   )r   r   r   r(   r  r   rÓ   r)   ÚtuplerÐ   r   r  rƒ   rY   r  r  r   r   r    r  r  Õ  sï   „ ñóIð ð $(ñà ðð >ðð ð	ð
 'ðð ðð ˜D‘[ðð 
ˆw�}‰}Ð;Ð;Ñ	<òó ðð4 ðà ðð >ðð ð	ð
 'ðð ðð 
ˆs�C˜˜c™Ð"Ñ	#òó ðð* ò>ó ñ>r   r  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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y)ÚGenerateManagerzFSequential generation via ``model.generate()`` on a persistent thread.c                 ó"   — t        «       | _        y r%   )rï   rõ   rÊ   s    r    r'   zGenerateManager.__init__  s   € Ü&Ó(ˆ�r   Nr5   r   rK   r  r  r	  r   r×   r¬   r6   c                 ó  ‡‡
‡‡— t        j                  «       Št        j                  «       Št        |d|«      j                  }t        |‰‰|¬«      }i |¥|||dœ¥Š
t        ‰d«      rd‰
d<   dˆ
ˆˆˆfd„}	| j                  |	«       ‰|fS )	zLStart streaming generation via ``model.generate()`` on the inference thread.r8   ©r¬   )ÚstreamerÚgeneration_configr8   Ú
has_talkerrÇ   Úgeneration_modec            	      ó   •— 	  ‰j                   di ‰¤Ž y # t        $ r ‰j                  ‰j                  d «       Y y t        $ r8} ‰j                  ‰j                  t        t        | «      «      «       Y d } ~ y d } ~ ww xY w)Nr   )Úgenerater+   rÃ   rÄ   rù   r"   r)   )rþ   Ú
gen_kwargsr©   r5   rª   s    €€€€r    rô   z0GenerateManager.generate_streaming.<locals>._run/  sm   ø€ ðRØ�—‘Ñ, Ó,øÜ'ò BØ×)Ñ)¨%×*:Ñ*:¸DÖAÜò RØ×)Ñ)¨%×*:Ñ*:¼LÌÈQËÓ<P×QÑQûðRús   ƒ –%A=½A=Á.A8Á8A=r‚   )rÐ   r  r   rG   r±   r§   Úhasattrr  )r&   r5   rK   r  r	  r×   r¬   Úrust_tokenizerr  rô   r#  r©   rª   s    `        @@@r    r  z"GenerateManager.generate_streaming  s�   û€ ô ×'Ñ'Ó)ˆÜ&Ÿ}™}›ˆä  ¨K¸ÓC×NÑNˆÜ! .°$¸È;ÔWˆØn˜Ðn¨HÈ:ÐdmÒnˆ
Ü�5˜,Ô'Ø,2ˆJÐ(Ñ)÷	Rð 	Rð 	�‰�DÔØ�hˆÐr   rÏ   c              ƒ   óò   K  — i |¥||dœ¥}t        |d«      rd|d<    | j                  |j                  fi |¤Žƒ d{  –—† }|d   j                  d   }|d|d…f   }	|j	                  |	d	¬
«      }
|
||	fS 7 Œ7­w)zNRun generation to completion via ``model.generate()`` on the inference thread.)r  r8   r  rÇ   r   NÚ	input_idsr�   r   T©r«   )r$  r  r"  ÚshapeÚdecode)r&   r5   rK   r  r	  r×   Úgenerate_kwargsÚ	sequencesÚ	input_lenrZ   rÇ   s              r    r  z&GenerateManager.generate_non_streaming:  s›   è ø€ ð ^˜VÐ]¸*ÐS\Ò]ˆÜ�5˜,Ô'Ø17ˆOÐ-Ñ.Ø+˜$×+Ñ+¨E¯N©NÑN¸oÑN×Nˆ	Ø˜;Ñ'×-Ñ-¨bÑ1ˆ	Ø! ! Y¡Z -Ñ0ˆØ×Ñ À4ÐÓHˆØ�Y Ð-Ð-ð	 Oús   ‚;A7½A5¾8A7rû   c                 óB   —  | j                   j                  |g|¢­i |¤ŽS )r   )rõ   r  ©r&   rû   r•   r–   s       r    r  zGenerateManager.submitN  s#   € à"ˆt�|‰|×"Ñ" 2Ð7¨Ò7°Ñ7Ð7r   c                 óB   —  | j                   j                  |g|¢­i |¤ŽS r  )rõ   r  r/  s       r    r  zGenerateManager.async_submitR  s#   € à(ˆt�|‰|×(Ñ(¨Ð=¨dÒ=°fÑ=Ð=r   c                  ó   — y r%   r   rÊ   s    r    r  zGenerateManager.stopV  ro   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    sä   „ ÙPò)ð $(ñà ðð >ðð ð	ð
 'ðð ðð ˜D‘[ðð 
ˆw�}‰}˜nÐ,Ñ	-óð<.à ð.ð >ð.ð ð	.ð
 'ð.ð ð.ð 
ˆs�C˜Ð'Ñ	(ó.ð(8˜ð 8°vó 8ð>˜xð >¸W¿^¹^ó >ôr   r  c                   óº   — e Zd ZdZddd„Zdd	„Z	 dddd
ddeddde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y)Ú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.
    Nc                 ó    — d | _         || _        y r%   )rÙ   Ú
_cb_config)r&   Ú	cb_configs     r    r'   zCBGenerateManager.__init__i  s   € ØˆŒØ#ˆ�r   r5   r   r	  r   r6   c                 ó–   — | j                   �y|j                  || j                  ¬«      | _         | j                   j                  «        y)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)r  Úcontinuous_batching_config)rÙ   Úinit_continuous_batchingr5  rö   r  s      r    r  zCBGenerateManager.init_cbm  sA   € ð �8‰8ÐØà×1Ñ1Ø(ÀTÇ_Á_ð 2ó 
ˆŒð 	�‰�‰Õr   rK   r  r  r×   r¬   c                 ó‚  ‡‡— | j                   }|€t        d«      ‚t        j                  «       }t        j                  «       Š|d   }	|j                  |	|d|j                  |j                  ¬«      }t        |d|«      j                  }
t        | j                   ||
|‰|¬«      Šˆˆfd„}|j                  ||«       ‰‰fS )zFStart streaming CB generation. Registers a per-request output handler.ú3CB manager not initialized. Call `init_cb()` first.r'  T)r×   Ú	streamingÚmax_new_tokensÚeos_token_idr8   r  c                 óÞ   •— 	 ‰j                  | «       | j                  «       r‰j                  «        y y # t        $ r-}‰j	                  t        t        |«      «      «       Y d }~y d }~ww xY wr%   )rÈ   Úis_finishedrË   rù   rÄ   r"   r)   )rß   rþ   r  Ú
text_queues     €€r    Ú
_on_outputz8CBGenerateManager.generate_streaming.<locals>._on_output�  sX   ø€ ð<Ø—‘˜VÔ$Ø×%Ñ%Ô'Ø—L‘L•Nð (øäò <Ø×%Ñ%¤l´3°q³6Ó&:×;Ñ;ûð<ús   ƒ16 ¶	A,¿#A'Á'A,)rÙ   ÚRuntimeErrorrÐ   r  r   Úadd_requestr=  r>  rG   r±   rÕ   Úregister_result_handler)r&   r5   rK   r  r	  r×   r¬   Úcbr©   r'  r%  rB  r  rA  s               @@r    r  z$CBGenerateManager.generate_streaming~  sÁ   ù€ ð �X‰XˆØˆ:ÜÐTÓUÐUä×'Ñ'Ó)ˆÜ$+§M¡M£Oˆ
à˜;Ñ'ˆ	Ø—^‘^ØØ!ØØ%×4Ñ4Ø#×0Ñ0ð $ó 
ˆ
ô ! ¨K¸ÓC×NÑNˆÜ˜dŸh™h¨
°NÀDÈ*ÐbmÔnˆõ	<ð 	×"Ñ" :¨zÔ:Ø˜8Ð#Ð#r   c              ƒ   ó¨  ‡K  — | j                   }|€t        d«      ‚|d   }t        |«      }t        j                  «       }	|	j                  «       Šˆfd„}
|j                  ||
«       |j                  |||j                  d|j                  ¬«       ‰ƒ d{  –—† }|€t        d|› �«      ‚|j                  }|j                  |d¬	«      }|||fS 7 Œ8­w)
zcRun non-streaming CB generation. Registers a handler that resolves an asyncio.Future on completion.Nr;  r'  c                 óJ   •— ‰j                  «       s‰j                  | «       y y r%   )Údonerø   )rý   rü   s    €r    Ú
_on_resultz<CBGenerateManager.generate_non_streaming.<locals>._on_result¼  s   ø€ Ø—;‘;”=Ø×!Ñ! &Õ)ð !r   F)r×   r=  r<  r>  z1CB manager stopped before producing a result for Tr(  )rÙ   rC  rÞ   rÐ   r  r  rE  rD  r=  r>  rÝ   r*  )r&   r5   rK   r  r	  r×   rF  r'  r-  r©   rJ  rý   rZ   rÇ   rü   s                 @r    r  z(CBGenerateManager.generate_non_streaming¨  sì   øè ø€ ð �X‰XˆØˆ:ÜÐTÓUÐUà˜;Ñ'ˆ	Ü˜	“Nˆ	ô ×'Ñ'Ó)ˆØ×#Ñ#Ó%ˆô	*ð 	×"Ñ" :¨zÔ:à
�‰ØØ!Ø%×4Ñ4ØØ#×0Ñ0ð 	ô 	
ð —ˆØˆ>ÜÐ!RÐS]ÐR^Ð_Ó`Ð`Ø×/Ñ/ˆØ×Ñ À4ÐÓHˆØ�Y Ð-Ð-ð ús   ƒBCÂCÂ9Cc                 óp   — | j                   €t        d«      ‚| j                   j                  j                  S )z*The CB scheduler (for testing/monitoring).zCB manager not initialized.)rÙ   rC  Úbatch_processorÚ	schedulerrÊ   s    r    rM  zCBGenerateManager.schedulerÐ  s0   € ð �8‰8ÐÜÐ<Ó=Ð=Ø�x‰x×'Ñ'×1Ñ1Ð1r   c                 óX   — | j                   �| j                   j                  dd¬«       y y )NTé   )ÚblockÚtimeout)rÙ   r  rÊ   s    r    r  zCBGenerateManager.stop×  s%   € Ø�8‰8ÐØ�H‰H�M‰M ¨aˆMÕ0ð  r   r%   )r6  úContinuousBatchingConfig | Noner  )r6   r   r‚   )r   r   r   r(   r'   r  rÓ   r)   r  rÐ   r   rÕ   r  rƒ   rY   r  ÚpropertyrM  r  r   r   r    r3  r3  Z  sÛ   „ ñô$óð0 $(ñ($à ð($ð >ð($ð ð	($ð
 'ð($ð ð($ð ˜D‘[ð($ð 
ˆw�}‰}˜jÐ(Ñ	)ó($ðT&.à ð&.ð >ð&.ð ð	&.ð
 'ð&.ð ð&.ð 
ˆs�C˜˜c™Ð"Ñ	#ó&.ðP ò2ó ð2ô1r   r3  c                   ó^   — 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y)Ú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.
    NÚcontinuous_batchingÚcompiler6  rR  c                 óX   — || _         || _        || _        i | _        d | _        d | _        y r%   )Ú_continuous_batchingÚ_compiler5  Ú_generate_managersÚ_cb_managerÚ_cb_model_id)r&   rV  rW  r6  s       r    r'   zGenerationState.__init__ê  s2   € ð %8ˆÔ!ØˆŒØ#ˆŒØ>@ˆÔØ59ˆÔØ(,ˆÕr   r5   r   Úmodalityr6   c                 ó¾   — | j                   syt        |d«      xr |t        j                  k(  }|s,t        j                  |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.
        Fr9  zM does not support continuous batching. Falling back to sequential generation.)rY  r$  r   r   ÚloggerÚwarning_oncer—   r   )r&   r5   r^  Úcans       r    Úuse_continuous_batchingz'GenerationState.use_continuous_batching÷  s]   € ð ×(Ò(ØÜ�eÐ7Ó8ÒU¸XÌÏÉÑ=UˆÙÜ×ÑØ—?‘?×+Ñ+Ð,ð -9ð 9ôð ˆ
r   r`   Úuse_cbc                 óZ  — |rv| j                   |k7  r-| j                  �!| 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)r6  )r]  r\  r  r3  r5  r[  r  )r&   r`   rd  s      r    Úget_managerzGenerationState.get_manager  sž   € ñ Ø× Ñ  HÒ,Ø×#Ñ#Ð/Ø×$Ñ$×)Ñ)Ô+Ø'+�DÔ$Ø×ÑÐ'Ü#4¸t¿¹Ô#O�Ô Ø$,�Ô!Ø×#Ñ#Ð#Ø˜4×2Ñ2Ñ2Ü0?Ó0AˆD×#Ñ# HÑ-Ø×&Ñ& xÑ0Ð0r   c                 ó`   — | j                   �"| j                   j                  «        d| _         yy)z$Stop any active generation managers.N)r\  r  rÊ   s    r    ÚshutdownzGenerationState.shutdown"  s-   € à×ÑÐ'Ø×Ñ×!Ñ!Ô#Ø#ˆDÕð (r   )FFN©Fr‚   )r   r   r   r(   rÒ   r'   r   rc  r)   r  rf  rh  r   r   r    rU  rU  Ü  so   „ ñð %*ØØ7;ñ	-à!ð-ð ð-ð 5ó	-ðÐ->ð È(ð ÐW[ó ñ(1 Cð 1°ð 1ÐBUó 1ô.$r   rU  c            	       óà   — e Zd ZU dZdZedz  ed<    e«       Zee	   ed<   ddde
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y)Ú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_statec                 ó    — || _         || _        y r%   )rn  ro  )r&   rn  ro  s      r    r'   zBaseHandler.__init__9  s   € ð
 +ˆÔØ 0ˆÕr   Úbodyr6   c                 ó  — ddl m} t        |j                  «       «      }| j                  �1|t        | j                  dt        «       «      z
  }|r |dd|› �¬«      ‚|| j                  z  }|rt        j                  d|› �«       yy)	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: )	Úfastapirt  rÍ   Úkeysrl  rG   rm  r`  ra  )r&   rq  rt  Ú
input_keysÚ
unexpectedÚunuseds         r    Ú_validate_requestzBaseHandler._validate_requestA  s‡   € å)ä˜Ÿ™›Ó%ˆ
Ø×#Ñ#Ð/Ø#¤g¨d×.FÑ.FÐHZÔ\_Ó\aÓ&bÑbˆJÙÙ#°Ð>`ÐakÐ`lÐ<mÔnÐnØ˜d×1Ñ1Ñ1ˆÙÜ×ÑÐ"NÈvÈhÐ WÕXð r   Úchunkzstr | pydantic.BaseModelc                 ó€   — t        | t        «      r| j                  d«      r| S d| › d�S d| j                  d¬«      › d�S )z;Format a pydantic model or string as an SSE ``data:`` line.zdata: z

T)Úexclude_none)rT   r)   Ú
startswithÚmodel_dump_json)r  s    r    Úchunk_to_ssezBaseHandler.chunk_to_sseN  sM   € ô �eœSÔ!Ø!×,Ñ,¨XÔ6�5ÐP¸fÀUÀGÈ4Ð<PÐPØ˜×-Ñ-¸4Ð-Ó@ÐAÀÐFÐFr   r   r  c                 ó�  — ddl m} | j                  j                  �j|j	                  d«      }|�>|| j                  j                  k7  r% |dd| j                  j                  › d|› d�¬«      ‚| j                  j                  |d<   | j                  j                  |d   «      }| j                  j                  |«      \  }}|||fS )	zfApply force_model, load model + processor.

        Returns ``(model_id, model, processor)``.
        r   rs  r5   i�  zServer is pinned to 'z'; requested 'z'.rv  )ry  rt  rn  Úforce_modelrS   Úprocess_model_nameÚload_model_and_processor)r&   rq  rt  Ú	requestedr`   r5   rK   s          r    Ú_resolve_modelzBaseHandler._resolve_modelU  sÌ   € õ
 	*à×Ñ×)Ñ)Ð5ØŸ™ Ó)ˆIØÐ$¨°d×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Ñˆˆyà˜ 	Ð)Ð)r   Úmodel_generation_configr   rd  c                 óB  — ddl m} |j                  d«      � |di t        j                  |d   «      ¤Ž}n7t        j                  |«      }|j                  �|j                  dk  rd|_        |j                  d«      �+t        |d   «      |_	        t        |d   «      dk(  rd|_
        |j                  d«      �t        |d   «      |_        |j                  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   r  i   Útemperatureg        FÚtop_prä   Ústaticr   )Útransformersr   rS   r/   ÚloadsÚcopyÚdeepcopyr=  Úfloatr�  Ú	do_samplerŽ  rè   ro  rZ  Úcache_implementationÚ	use_cache)r&   rq  r‹  rd  r   r  s         r    Ú_build_generation_configz$BaseHandler._build_generation_configj  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   Úmessagesr^  c           	      ó  — g }| D �]þ  }|d   g dœ}d|v r|d   |d<   d|v r|d   |d<   d|v rg n|j                  d«      xs g }t        |t        «      rd|dœg}|D �]b  }|d   }|d	v r|d   j                  d|d   dœ«       Œ(|d
v rT|t        j
                  t        j                  fv r2|d   }t        |t        «      r|d   }|d   j                  d|dœ«       Œ€|dk(  r_|t        j                  k(  rL|d   }	t        |	t        «      r|	j                  dd«      nd}
|	d   }|d   j                  dd|
› d|› �dœ«       Œä|dk(  rA|t        j
                  t        j                  fv r|d   j                  d|d   d   dœ«       �Œ*|dk(  s�Œ1|t        j                  k(  s�ŒF|d   j                  d|d   d   dœ«       �Œe |t        j                  k(  rdj                  d„ |d   D «       «      |d<   |j                  |«       �Œ |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›  Úcontentr=   Útool_call_idrœ  rÇ   )r0   rÇ   r0   )rÇ   Ú
input_textÚoutput_text)Ú	image_urlÚinput_imager   ÚurlÚimage)r0   r¢  Úinput_audioÚformatÚwavÚdataÚaudiozdata:audio/z;base64,Ú	video_urlÚvideoÚ	audio_urlú c              3   ó&   K  — | ]	  }|d    –— Œ y­w)rÇ   Nr   )rA   rr   s     r    rD   zABaseHandler.get_processor_inputs_from_messages.<locals>.<genexpr>Ö  s   è ø€ Ò,R¸1¨Q¨v­YÑ,Rùrt   )
rS   rT   r)   rÁ   r   r   r   rÓ   r   Újoin)r™  r^  Úprocessor_inputsÚmessager[   Úraw_contentrœ  Úcontent_typer¢  r¤  ÚfmtÚ	audio_b64s               r    Ú"get_processor_inputs_from_messagesz.BaseHandler.get_processor_inputs_from_messagesœ  sR  € ð Ðàó +	,ˆGØ% f™o¸"Ñ=ˆFð ˜wÑ&Ø'.¨|Ñ'<��|Ñ$Ø Ñ(Ø)0°Ñ)@��~Ñ&ð !-°Ñ 7™"¸g¿k¹kÈ)Ó>TÒ>ZÐXZˆKÜ˜+¤sÔ+Ø(.¸ÑDÐE�à&ó d�Ø& v™�àÐ#HÑHØ˜9Ñ%×,Ñ,°fÀgÈfÁoÑ-VÕWà!Ð%AÑAÀhÔS[×S_ÑS_Ôai×atÑatÐRuÑFuà! +Ñ.�CÜ! #¤tÔ,Ø! %™j˜Ø˜9Ñ%×,Ñ,°gÀcÑ-JÕKà! ]Ò2°xÄ8×CVÑCVÒ7VØ")¨-Ñ"8�KÜ>HÈÔVZÔ>[˜+Ÿ/™/¨(°EÔ:Ðaf�CØ +¨FÑ 3�IØ˜9Ñ%×,Ñ,°gÈÐTWÐSXÐX`ÐajÐ`kÐFlÑ-mÕnà! [Ò0°XÄ(Ç,Á,ÔPX×PcÑPcÐAdÑ5dØ˜9Ñ%×,Ñ,°gÀgÈkÑFZÐ[`ÑFaÑ-bÖcØ! [Ô0°XÄ×ATÑATÔ5TØ˜9Ñ%×,Ñ,°gÀgÈkÑFZÐ[`ÑFaÑ-bÖcð-dð2 œ8Ÿ<™<Ò'Ø$'§H¡HÑ,RÀÀyÑ@QÔ,RÓ$R��yÑ!à×#Ñ# FÖ+ðW+	,ðX  Ðr   ri  )r   r   r   r(   rl  r0   Ú__annotations__rÍ   rm  r)   rU  r'   rÓ   r~  Ústaticmethodr„  r  rŠ  rÒ   r˜  rY   r   rµ  r   r   r    rk  rk  )  sþ   … ñ
ð (,Ð˜ ™Ó+Ù"›u€N�C˜‘HÓ$ð1à%ð1ð *ó1ðY dð Y¨tó Yð ðGÐ6ð G¸3ò Gó ðGð* 4ð *¨E°#Ð7HÐJtÐ2tÑ,uó *ð, W\ñ0!Øð0!Ø3Eð0!ØOSð0!à	ó0!ðd ð< °T¸$±Zð < È8ð < ÐX\Ð]aÑXbò < ó ñ< r   rk  r‚   )Br(   rÐ   r’  Úenumr/   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   rn  r   Ú
get_loggerr   r`  ÚX_REQUEST_IDÚEnumr   r"   rù   r+   rI   rÓ   rM   rV   rY   r\   r^   r)   r0   r¥   r§   rÕ   rƒ   rè   rí   rï   r  r  r3  rU  rk  r   r   r    ú<module>rÆ     s¹  ðñó Û Û Û Û ß #Ý $Ý %Ý Ý  å &ñ ÛÛÛ÷õ õ eÝUÝOå+ð 
ˆ×	Ñ	˜HÓ	%€ð €ôˆt�y‰yô ÷ñ ôP˜9ô Pð ØØà ?ØØ&°FÑ;ñ
ñð
Ð ðBÐ+<ð BÀÈÁó Bð8 Dð ¨Tó ð".°tð .ÀÀTÁ
ÈTÑ@Qó .÷*(
ñ (
ðVD xð D¸3ð DÀ4ó D÷NEñ E÷PA2ñ A2ðH˜ð  ó ó!÷&ñ &ôR>>˜#ô >>ôBAÐ)ô AôH1Ð+ô 1÷DJ$ñ J$÷Zp ò p r   