ó
    pyüiÿ	 ã                   ó   • S SK r S SKrS SKrS SKrS SKJr  S SKJrJr  S SK	J
r
Jr  S SKJr  S SKJr  S SKJr  S SKrS SKJr  S S	KJr  S S
KJr  SSKJr  SSKJrJr  SSKJr  SSKJr  SSK J!r!  SSK"J#r#J$r$J%r%  SSK&J'r'  SSK(J)r)  SSK*J+r+  SSK,J-r-J.r.  SSK/J0r0  SSK1J2r2J3r3J4r4J5r5  SSK6J7r7J8r8J9r9  SSK:J;r;J<r<J=r=J>r>J?r?    " S S\R€                  5      rA " S S5      rB\$" 5        " S  S!5      5       rC\$" 5        " S" S#5      5       rD " S$ S%5      rEg)&é    N)Úabstractmethod)ÚCallableÚ	Generator)ÚcontextmanagerÚnullcontext)Úceil)Úperf_counter)ÚAny)Únn)Útqdm)Úlogging_redirect_tqdmé   )ÚPretrainedConfig)ÚContinuousBatchingConfigÚGenerationConfig)Ú!lazy_import_paged_flash_attention)Úis_flash_attention_requested)Úlogging)ÚContinuousBatchProcessorMetricsÚattach_tracerÚtracedé   )ÚLogitsProcessorListé   )ÚPagedAttentionCache)Ú%ContinuousBatchingLogitsProcessorList)ÚContinuousBatchingAsyncIOsÚContinuousBatchingIOs)ÚOffloadingManager)ÚGenerationOutputÚRequestStateÚRequestStatusÚlogger)ÚSCHEDULER_MAPPINGÚFIFOSchedulerÚ	Scheduler)ÚWorkloadHintsÚattn_mask_is_neededÚcreate_warmup_future_statesÚpad_to_intervalÚpad_to_pow2c                   ó”   • \ rS rSr% \\S'   \R                  \S'   \R                  \S'   \	S\
SS4S j5       r\	S	\S\4S
 j5       rSrg)ÚProtoPretrainedModeléF   ÚconfigÚdtypeÚdeviceÚattn_implementationÚreturnNc                 ó   • g ©N© )Úselfr2   s     Úw/home/mande/repo/quber/.venv/lib/python3.13/site-packages/transformers/generation/continuous_batching/continuous_api.pyÚset_attn_implementationÚ,ProtoPretrainedModel.set_attn_implementationK   ó   € àó    Úgeneration_configc                 ó   • g r5   r6   )r7   r=   s     r8   Ú_get_logits_processorÚ*ProtoPretrainedModel._get_logits_processorO   r;   r<   r6   )Ú__name__Ú
__module__Ú__qualname__Ú__firstlineno__r   Ú__annotations__Útorchr0   r1   r   Ústrr9   r   r   r?   Ú__static_attributes__r6   r<   r8   r-   r-   F   s`   ‡ ØÓØ�;‰;ÓØ�L‰LÓàð¸3ð À4ó ó ðð ðÐ7Gð ÐL_ó ó ór<   r-   c                   óL   • \ rS rSrSrSS jrS\SS4S jrS\\   SS4S	 jr	S
r
g)ÚOutputRouteréT   a  Dedicated object for routing generation outputs to the right destination.

When an async handler is registered for a request, the output is forwarded
to that handler via ``call_soon_threadsafe``. Otherwise the output is placed
on the shared ``output_queue``.
r3   Nc                 óz   • [         R                  " 5       U l        0 U l        [        R
                  " 5       U l        g r5   )ÚqueueÚQueueÚoutput_queueÚresult_handlersÚ	threadingÚLockÚ_lock©r7   s    r8   Ú__init__ÚOutputRouter.__init__\   s&   € Ü!ŸKšK›MˆÔØVXˆÔÜ—^’^Ó%ˆ�
r<   Úoutputc                 ó  • U R                      U R                  R                  UR                  5      nSSS5        Wb  Uu  p4UR	                  X15        gU R
                  R                  U5        g! , (       d  f       NC= f)zDRoute a single output to its registered handler or the output_queue.N)rS   rP   ÚgetÚ
request_idÚcall_soon_threadsaferO   Úput)r7   rW   ÚentryÚcallbackÚloops        r8   ÚdeliverÚOutputRouter.delivera   sa   € à�Z‹ZØ×(Ñ(×,Ñ,¨V×->Ñ->Ó?ˆE÷ àÑØ"‰NˆHØ×%Ñ% hÕ7à×Ñ×!Ñ! &Õ)÷ �Zús   �&A0Á0
A>Úoutputsc                 óf  • / nSnU R                      U H^  nU R                  R                  UR                  5      nUb  Uu  pcUR	                  Xd45        MC  U R
                  R                  U5        M`     SSS5        U(       a  Ub  U4S jnUR                  U5        ggg! , (       d  f       N2= f)zµRoute a batch of outputs, using a single ``call_soon_threadsafe`` to minimize cross-thread overhead.

Outputs without a registered handler fall back to the shared ``output_queue``.
Nc                 ó*   • U  H  u  pU" U5        M     g r5   r6   )ÚbatchÚcbÚouts      r8   Ú
_run_batchÚ.OutputRouter.deliver_batch.<locals>._run_batch|   s   € Û$‘G�BÙ�s–Gò  %r<   )rS   rP   rY   rZ   ÚappendrO   r\   r[   )r7   rb   Ú	callbacksr_   rW   r]   r^   rh   s           r8   Údeliver_batchÚOutputRouter.deliver_batchk   s¤   € ð
 >@ˆ	ØˆØ�Z‹ZÛ!�Ø×,Ñ,×0Ñ0°×1BÑ1BÓC�ØÑ$Ø%*‘N�HØ×$Ñ$ hÐ%7Ö8à×%Ñ%×)Ñ)¨&Ö1ñ "÷ ö ˜Ñ)à!*ô ð ×%Ñ% jÕ1ð *ˆ9÷ �Zús   ‘A%B"Â"
B0)rS   rO   rP   ©r3   N)rA   rB   rC   rD   Ú__doc__rU   r    r`   Úlistrl   rH   r6   r<   r8   rJ   rJ   T   s9   † ñô&ð
*Ð.ð *°4ô *ð2 TÐ*:Ñ%;ð 2À÷ 2r<   rJ   c                   ó  • \ rS rSr% \\-  \S'   \\S'   S\S\	S\
S\S\S	\R                  S
\S\R"                  S\R&                  S\R(                  S\SS4S jrS\4S jrS8S jrS8S jrS8S jr\S8S j5       r\S\S\SS4S j5       rS\ S\ S\!S\"\ \ 4   4S jr#\S\!4S j5       r$\S8S j5       r%\S\!4S j5       r&\S  5       r'\S\SS4S! j5       r(\\RR                  " 5       S"\*RV                  SS4S# j5       5       r,S$\-S%\R\                  R^                  SS4S& jr0\S"\*RV                  S'\1S(\Rd                  S)\Rd                  S*\Rd                  SS4S+ j5       r3\" S,S-9S"\*RV                  S'\1S\Rd                  4S. j5       r4\" S/S-9S'\1S0\Rd                  S\Rd                  4S1 j5       r5\" S2S-9S3\Rd                  S4\Rd                  S*\Rd                  SS4S5 j5       r6\Rn                  " 5       S"\*RV                  SS4S6 j5       r8S7r9g)9ÚContinuousBatchProcessoré„   Úinputs_and_outputsÚ	schedulerÚcacher/   r=   Úcontinuous_batching_configÚlogit_processorÚinput_queueÚoutput_routerÚ
stop_eventÚmodel_deviceÚmodel_dtyper3   Nc                 ó´  • Xl         X l        X@l        XPl        X`l        Xpl        X€l        X�l        X l        X°l	        [        USS5      U l        UR                  U l        [        USS5      c  SOUR                  U l        U R                  R                  U l        U R                  R                  U l        U R                  R!                  5       u  U l        U l        UR&                  U l        [)        UR&                  5      U l        U R-                  5         U R                  R/                  [        USS5      [1        US9U R                   R2                  S:„  S	9  U R                  R4                  U R                  R6                  pÜSU l        Ub4  [:        R<                  " U R>                  40 URA                  5       D6U l        SU l!        Ub4  [:        R<                  " U R>                  40 URA                  5       D6U l!        U R"                  =(       d    U R$                  nU=(       d    USL=(       d    USLU l"        U(       a  [:        RF                  RI                  5       OSU l%        UUU	U
U R                  U R                  U R"                  S
.nU R                  RL                  U l&        U RL                  (       a6  [O        U R                  RP                  S-  5      US'   [S        S0 UD6U l*        O)U R                  RP                  US'   [W        S0 UD6U l*        [Y        UUURZ                  UR\                  U RT                  R^                  S9U l0        g)a·  Initialize the continuous batch processor.

Args:
    cache: A [`PagedAttentionCache`] object
    config: The model configuration
    generation_config: The generation configuration
    continuous_batching_config: The continuous batching configuration
    logit_processor: The [`ContinuousBatchingLogitsProcessorList`] object used to process the logits.
    input_queue: Queue for incoming requests
    output_router: An [`OutputRouter`] object that routes outputs to handlers or the output queue.
    stop_event: Event to signal processing should stop
    model_device: Device for model inputs/outputs
    model_dtype: Data type for model inputs/outputs
    scheduler: The [`Scheduler`] to use
Ú	do_sampleTÚsliding_windowNr   Úcompile_config)r/   r   )Úfallback_compile_configÚis_flash_attnÚdecode_fast_path_available)rv   r/   r1   r}   Úreturn_logprobsrx   Úuse_cuda_graph_varlenr   Ú
max_graphs)rv   ru   Úcpu_offload_space_gibÚsafety_thresholdÚcompute_streamr6   )1rv   r/   Ú	cb_configrx   ry   rz   r{   r|   r}   ru   Úgetattrr   r…   r€   Úq_padding_interval_sizeÚkv_padding_interval_sizeÚget_cuda_graph_booleansr†   Úuse_cuda_graph_decodeÚmax_batch_tokensr   ÚmetricsÚ%_ensure_decode_fast_path_is_availableÚresolve_compile_configsr   Úmax_blocks_per_requestÚvarlen_compile_configÚdecode_compile_configÚ_compiled_varlenrF   ÚcompileÚ_forward_process_and_sampleÚto_dictÚ_compiled_decodeÚ_pad_inputsÚcudaÚgraph_pool_handleÚ
graph_poolÚuse_async_batchingr   Úmax_cached_graphsr   rt   r   r   Úcpu_offload_spaceÚ"cpu_offload_space_safety_thresholdrŠ   Úoffloading_manager)r7   rv   r/   r=   rw   rx   ry   rz   r{   r|   r}   ru   Úvarlen_configÚdecode_configÚuse_cuda_graphsÚ	io_kwargss                   r8   rU   Ú!ContinuousBatchProcessor.__init__‰   sÝ  € ð: Œ
ØŒØ3ŒØ.ÔØ&ÔØ*ÔØ$ŒØ(ÔØ&ÔØ"Œô !Ð!2°KÀÓFˆŒØ9×IÑIˆÔô $+¨6Ð3CÀTÓ#JÑ#R™aÐX^×XmÑXmˆÔà'+§~¡~×'MÑ'MˆÔ$Ø(,¯©×(OÑ(OˆÔ%ØAEÇÁ×AgÑAgÓAiÑ>ˆÔ" DÔ$>ð !&× 6Ñ 6ˆÔÜ6°u×7MÑ7MÓNˆŒð 	×2Ñ2Ô4ð 	�‰×.Ñ.Ü$+Ð,=Ð?OÐQUÓ$VÜ6¸fÑEØ'+§z¡z×'HÑ'HÈ1Ñ'Lð 	/ñ 	
ð
 (,§~¡~×'KÑ'KÈTÏ^É^×MqÑMq�}ð !%ˆÔØÑ$Ü$)§M¢M°$×2RÑ2RÑ$nÐVc×VkÑVkÓVmÑ$nˆDÔ!ð !%ˆÔØÑ$Ü$)§M¢M°$×2RÑ2RÑ$nÐVc×VkÑVkÓVmÑ$nˆDÔ!ð ×4Ñ4×R¸×8RÑ8RˆØ*×f¨}ÀDÐ/H×/eÈMÐaeÐLeˆÔæ<Kœ%Ÿ*™*×6Ñ6Ô8ÐQUˆŒð ØØ"Ø&Ø#×3Ñ3Ø#×3Ñ3Ø%)×%?Ñ%?ñ
ˆ	ð #'§.¡.×"CÑ"CˆÔà×"×"ä&*¨4¯>©>×+KÑ+KÈaÑ+OÓ&PˆI�lÑ#Ü&@Ñ&MÀ9Ñ&MˆDÕ#à&*§n¡n×&FÑ&FˆI�lÑ#Ü&;Ñ&H¸iÑ&HˆDÔ#ô #4ØØØ"<×"NÑ"NØ7×ZÑZØ×2Ñ2×AÑAñ#
ˆÕr<   c                 óÐ   • SU R                    SU R                  R                   SU R                  R                   S3U R                  R                  5       R                  5       -   $ )Nz%ContinuousBatchProcessor(input_queue=z, active_requests=z, waiting_requests=Ú))ry   ru   Úactive_requestsÚwaiting_requestsrt   Úget_model_kwargsÚ__repr__rT   s    r8   r°   Ú!ContinuousBatchProcessor.__repr__ø   sj   € à3°D×4DÑ4DÐ3Eð FØ#Ÿ~™~×=Ñ=Ð>Ð>QÐRV×R`ÑR`×RqÑRqÐQrÐrsðuà×%Ñ%×6Ñ6Ó8×AÑAÓCñDð	
r<   c                 óÀ   • S U l         [        R                  " 5         [        R                  R                  5       (       a  [        R                  R                  5         g g r5   )rt   ÚgcÚcollectrF   rž   Úis_availableÚempty_cacherT   s    r8   Ú__del__Ú ContinuousBatchProcessor.__del__ÿ   s;   € Ø"&ˆÔÜ
�
Š
ŒÜ�:‰:×"Ñ"×$Ñ$Ü�J‰J×"Ñ"Õ$ð %r<   c                 óø  • U R                   R                  SLnU(       d)  U R                   R                  U R                  l        S nO[        R
                  nU R                  R                  S:w  Ga  [        U R                  SS9(       a£  [        U R                  R                  5      S   nU R                  R                  S:H  [        R                  R                  5       USL/n[        U5      (       d6  U" SU R                  R                  < SU S	35        SU R                  l        ggU" SU R                  R                  < S
U R                  R                  < S	35        SU R                  l        gg)z¾Ensures the decode fast path is available. If it is not, set the max blocks per request to 0. If it is
available, and no user-provided max blocks per request, set it to the fallback default.Nc                 ó   • U $ r5   r6   ©Úxs    r8   Ú<lambda>ÚPContinuousBatchProcessor._ensure_decode_fast_path_is_available.<locals>.<lambda>  s   € ¡qr<   r   r   )Úversionr   z-Although self.cache.max_blocks_per_request = zN, the decode fast path is not available because the one condition is not met: Ú.z€, the decode fast path is not available because the attention implementation is not FA3. Got self.config._attn_implementation = )r‹   r•   Úfallback_max_blocks_per_requestrv   r#   Úwarningr   r/   r   Ú_attn_implementationÚnum_sliding_attention_groupsrF   rž   rµ   Úall)r7   Úuser_requestedÚlogger_warningÚflash_attn_with_kvcacheÚ
conditionss        r8   r“   Ú>ContinuousBatchProcessor._ensure_decode_fast_path_is_available  sH  € ð Ÿ™×>Ñ>ÀdÐJˆÞØ04·±×0^Ñ0^ˆD�J‰JÔ-Ù(‰Nä#Ÿ^™^ˆNð �:‰:×,Ñ,°Ô1ä+¨D¯K©KÀ×CÜ*KÈDÏKÉK×LlÑLlÓ*mÐnoÑ*pÐ'à—J‘J×;Ñ;¸qÑ@Ü—J‘J×+Ñ+Ó-Ø+°4Ð7ð�
ô ˜:—‘Ù"ØH D§J¡J×$EÑ$EÑ#Ið JAØAKÀÈAðOôð 9:�D—J‘JÕ5ð 'ñ ØD §
¡
× AÑ AÑEð FpØLPÏKÉK×LlÑLlÑKpÐpqðsôð 56�—
‘
Õ1ð- 2r<   c                 ó  • U R                   R                  5         U R                  R                  5         U R                  R                  5         U R                  R                  5         [        U R                  R                  5      U l        g)z4Reset the batch processor for a new generation loop.N)	r¥   Úresetru   rt   rv   Úfree_all_requestsr   r‘   r’   rT   s    r8   rÌ   ÚContinuousBatchProcessor.reset)  s\   € à×Ñ×%Ñ%Ô'Ø�‰×ÑÔØ×Ñ×%Ñ%Ô'Ø�
‰
×$Ñ$Ô&Ü6°t·z±z×7RÑ7RÓSˆ�r<   c                 ó(  • U R                   R                  5       (       d‚   U R                   R                  5       nUc  M?  U R                  R	                  UR
                  5        U R                  R                  U5        U R                   R                  5       (       d  M�  gg! [        R                   a     g[         aO  n[        R                  " SU 3SS9  [        5       R                  S5      nUb  U R                  X!5         SnANŒSnAff = f)z?Pull new requests from the input queue and add to waiting list.NzError processing new request: T©Úexc_infoÚstate)ry   ÚemptyÚ
get_nowaitrx   Úcheck_kwargsÚlogit_processor_kwargsru   Úadd_waiting_requestrM   ÚEmptyÚ	Exceptionr#   ÚerrorÚlocalsrY   Ú_handle_request_error)r7   rÒ   Úes      r8   Ú_get_new_requestsÚ*ContinuousBatchProcessor._get_new_requests1  sÙ   € ð ×"Ñ"×(Ñ(×*Ñ*ð9Ø×(Ñ(×3Ñ3Ó5�Ø‘=ÙØ×$Ñ$×1Ñ1°%×2NÑ2NÔOØ—‘×2Ñ2°5Ô9ð ×"Ñ"×(Ñ(×*Õ*øô —;‘;ó ÙÜó 9Ü—’Ð=¸a¸SÐAÈDÒQÜ&,£h§l¡l°7Ó&;�ØÑ$Ø×.Ñ.¨qÔ8ÿøð	9ús%   ¡B# Á A B# Â#DÂ9	DÃADÄDrÚ   rÒ   c                 ó¢  • [         R                  Ul        [        U5      Ul        [        UR                  [        5      (       a+  U R                  R                  UR                  5      Ul	        O/ Ul	        U R                  R                  UR                  UR                  5        U R                  R                  UR                  5       5        g)z(Handle general request processing error.N)r"   ÚFAILEDÚstatusrG   rÚ   Ú
isinstancerZ   ru   Ú!get_active_request_static_outputsÚgenerated_tokensr’   Úrecord_request_completionÚcreated_timerz   r`   Úto_generation_output)r7   rÚ   rÒ   s      r8   rÜ   Ú.ContinuousBatchProcessor._handle_request_errorD  s”   € ô %×+Ñ+ˆŒÜ˜%“jˆŒô �e×&Ñ&¬×,Ñ,Ø%)§^¡^×%UÑ%UÐV[×VfÑVfÓ%gˆEÕ"à%'ˆEÔ"à�‰×.Ñ.¨u×/AÑ/AÀ5×CSÑCSÔTØ×Ñ×"Ñ" 5×#=Ñ#=Ó#?Õ@r<   Únum_q_tokensÚmax_kv_readÚuse_decode_fast_pathc                 ó   • U R                   (       ak  U(       dM  [        XR                  U R                  5      n[        X R                  U R
                  R                  5      nX4$ [        XR                  5      nSnX4$ )z[Pads the inputs sizes for the next batch if it is needed. Often it is, for max performance.r   )r�   r*   r�   r‘   rŽ   rv   Ú	num_pagesr+   )r7   rê   rë   rì   s       r8   Úmaybe_pad_inputsÚ)ContinuousBatchProcessor.maybe_pad_inputsS  sq   € à××æ'Ü.¨|×=YÑ=YÐ[_×[pÑ[pÓq�Ü-¨k×;XÑ;XÐZ^×ZdÑZd×ZnÑZnÓo�ð
 Ð(Ð(ô  +¨<×9NÑ9NÓO�Ø�ØÐ(Ð(r<   c                 óà  • U R                  5         U R                  R                  5       nU H  nU R                  R	                  U5        M      U R                  R                  5       (       d  gU R                  R                  [        U R                  R                  5      [        U R                  R                  5      5        U R                  R                  U R                  U R                  R                  5      u  p4pVUcI  [        U R                  R                  5      S:”  a  U R                  R                  5         g[!        S5      eU(       d  gU R                  R#                  U5        U R                  R%                  U5        [&        R(                  " S[        U5       S[        U R                  R                  5       S[        U R                  R                  5       SU SU S	U R                  R+                  5        35        U R-                  XVU5      u  pVU R.                  R1                  X0R2                  XEU5        U R                  R5                  U R                  5        g
)z}Prepare tensors and metadata for the next model forward pass. Returns True if there are requests to process,
False otherwise.Fr   z=No requests can be scheduled and no request can be offloaded.zScheduled: z, Waiting: z
, Active: z	. cum Q: z
. cum KV: z, free blocks: T)rÞ   ru   Úclear_cancelled_requestsr¥   Úfree_request_cpu_cacheÚhas_pending_requestsr’   Úrecord_queue_metricsÚlenr­   r®   Úschedule_batchr‘   rv   rî   Úoffload_one_requestÚRuntimeErrorÚrestore_scheduled_requestsÚrecord_batch_metricsr#   ÚdebugÚget_num_free_blocksrï   rt   Úprepare_batch_tensorsrx   Úrecord_kv_cache_memory_metrics)r7   Úcancelled_statesrÒ   Úrequests_in_batchrì   rê   rë   s          r8   Úprepare_next_batchÚ+ContinuousBatchProcessor.prepare_next_batch`  sþ  € ð 	×ÑÔ ØŸ>™>×BÑBÓDÐã%ˆEØ×#Ñ#×:Ñ:¸5ÖAñ &à�~‰~×2Ñ2×4Ñ4ØØ�‰×)Ñ)¬#¨d¯n©n×.LÑ.LÓ*MÌsÐSW×SaÑSa×SrÑSrÓOsÔtð NRÏ^É^×MjÑMjØ×!Ñ! 4§:¡:×#7Ñ#7óN
ÑJÐ°ð Ñ$Ü�4—>‘>×1Ñ1Ó2°QÓ6Ø×'Ñ'×;Ñ;Ô=Øä"Ð#bÓcÐcæ Øð 	×Ñ×:Ñ:Ð;LÔMð 	�‰×)Ñ)Ð*;Ô<Ü�ŠØœ#Ð/Ó0Ð1°¼SÀÇÁ×A`ÑA`Ó=aÐ<bð cÜ˜4Ÿ>™>×9Ñ9Ó:Ð;¸9À\ÀNð SØ"�m ?°4·:±:×3QÑ3QÓ3SÐ2TðVô	
ð %)×$9Ñ$9¸,ÐUiÓ$jÑ!ˆà×Ñ×5Ñ5Ø×3Ñ3Ð5IÐYdô	
ð 	�‰×3Ñ3°D·J±JÔ?Ør<   c                 ót  • U R                   R                  5       u  pnSn/ nU GHQ  nUR                  nUR                  [        R
                  [        R                  4;   aY  U R                  (       a  UR                  (       a  US-  nMg  [        SUR                  R                   SUR                   S35      eUR                  (       Gab  UR                  5       S:X  aE  U R                  R                  UR                  UR                  5        [        R                   Ul        X$   nUb  X4   OSn	US-  nUR#                  X‰5      n
U R$                  R'                  XvR(                  5        U
(       af  U R                  R+                  UR                  UR                  5        U R,                  R/                  UR                  5        SU R,                  l        UR2                  (       d  UR                  [        R
                  :X  a"  UR5                  UR7                  5       5        GM  GM  UR                  [        R8                  :X  d  GM,  U R$                  R'                  XvR(                  5        GMT     U(       a  U R:                  R=                  U5        / / pËU R,                  R>                  (       aù  U R,                  R>                  RA                  5       nURB                  nSUl!        [E        U5       Vs/ s H  oýR                   SU 3PM     nnU H+  nURG                  U5      U R,                  RH                  U'   M-     U R$                  RK                  UR                  U5      u  nnURM                  U5        URM                  U5        U R,                  R>                  (       a  Mù  U(       ai  U R                   RN                  nUb  [P        RR                  RU                  U5      O	[W        5       nU   U R$                  RY                  X¼5        SSS5        ggs  snf ! , (       d  f       g= f)	z0Update request states based on generated tokens.r   r   zTried to update z	 request z in sync mode.NFz__child#)-rt   Úprepare_batch_updaterÒ   râ   r"   ÚFINISHEDÚPENDINGr¡   Úhas_new_tokenrù   ÚnamerZ   Úgenerated_lenr’   Úrecord_ttft_metricrç   ÚDECODINGÚupdate_and_check_completionrv   Ú!mark_shareable_blocks_as_completeÚcomplete_blocksræ   ru   Úfinish_requestÚblock_new_requestsÚ	streamingrj   rè   Ú
PREFILLINGrz   rl   Ú_requests_to_forkÚpopÚnum_childrenÚrangeÚforkr­   Úfork_requestÚextendrŠ   rF   rž   Ústreamr   Ú
copy_cache)r7   r  Ú
new_tokensÚlogprobsÚcurrent_logits_indexÚpending_outputsÚfuture_staterÒ   ÚtokenÚlogprobÚis_finishedÚcopy_sourceÚcopy_destinationÚstate_to_forkr  ÚiÚnew_request_idsÚnew_request_idÚcopy_srcÚcopy_dstrŠ   Úmaybe_streams                         r8   Úupdate_batchÚ%ContinuousBatchProcessor.update_batch’  sT  € ð 37×2IÑ2I×2^Ñ2^Ó2`Ñ/Ð xØ ÐØˆÜ-ˆLØ ×&Ñ&ˆEà�|‰|¤× 6Ñ 6¼×8MÑ8MÐNÓNØ×*×*à#×1×1Ø,°Ñ1Ð,ÙÜ"Ð%5°e·l±l×6GÑ6GÐ5HÈ	ÐRW×RbÑRbÐQcÐcqÐ#rÓsÐsà×)×)Ð)à×&Ñ&Ó(¨AÓ-Ø—L‘L×3Ñ3°E×4FÑ4FÈ×HXÑHXÔYÜ#0×#9Ñ#9�E”Là"Ñ8�Ø<DÑ<P˜(Ò8ÐVZ�Ø$¨Ñ)Ð$ð $×?Ñ?ÀÓO�à—
‘
×<Ñ<¸U×D`ÑD`ÔaÞØ—L‘L×:Ñ:¸5×;MÑ;MÈu×O_ÑO_Ô`Ø—N‘N×1Ñ1°%×2BÑ2BÔCØ8=�D—N‘NÔ5Ø—?—? e§l¡l´m×6LÑ6LÓ&LØ#×*Ñ*¨5×+EÑ+EÓ+G×Hò 'Mð —‘¤×!9Ñ!9Ö9Ø—
‘
×<Ñ<¸U×D`ÑD`×añC .öF Ø×Ñ×,Ñ,¨_Ô=ð )+¨BÐ%Ø�n‰n×.×.à ŸN™N×<Ñ<×@Ñ@ÓBˆMØ(×5Ñ5ˆLØ)*ˆMÔ&äQVÐWcÔQdÓeÒQdÈA×":Ñ":Ð!;¸8ÀAÀ3ÓGÑQdˆOÐeÛ"1�ØAN×ASÑASÐTbÓAc�—‘×.Ñ.¨~Ó>ñ #2ð "&§¡×!8Ñ!8¸×9QÑ9QÐSbÓ!cÑˆH�hØ×Ñ˜xÔ(Ø×#Ñ# HÔ-ð �n‰n×.×.Ñ.ö  ð "×4Ñ4×CÑCˆNØ@NÑ@Zœ5Ÿ:™:×,Ñ,¨^Ô<Ô`kÓ`mˆLÚØ—
‘
×%Ñ% kÔD÷ �ð ùò f÷ •ús   Ë>P$Ï>P)Ð)
P7c                 ó6   • U R                   R                  5       $ )z2Check if there are any active or waiting requests.)ru   rô   rT   s    r8   rô   Ú-ContinuousBatchProcessor.has_pending_requestsØ  s   € ð �~‰~×2Ñ2Ó4Ð4r<   c                 óä   • U R                   R                  5       S   nU HM  nU R                  XR                  5        U R                  R                  UR                  R                  5        MO     g)z&Handle errors during batch processing.r   N)rt   r  rÜ   rÒ   ru   r  rZ   )r7   rÚ   Úfailed_future_statesr!  s       r8   Úhandle_batch_errorÚ+ContinuousBatchProcessor.handle_batch_errorÝ  sZ   € ð  $×6Ñ6×KÑKÓMÈaÑPÐÛ0ˆLØ×&Ñ& u×.@Ñ.@ÔAØ�N‰N×)Ñ)¨,×*<Ñ*<×*GÑ*GÖHò 1r<   c                 ó,  • [        U R                  R                  R                  5       5      nU H9  nU R	                  X5        U R                  R                  UR                  5        M;     U R                  R                  5         [        U R                  R                  R                  5       5       H9  nU R                  R                  R                  U5      nU R	                  X5        M;     U R                  R                  R                  5         g)z.Fail all active requests with the given error.N)rp   ru   r­   ÚvaluesrÜ   r  rZ   r¥   Úfree_all_waiting_cpu_cachesr®   Úkeysr  Úwaiting_requests_orderÚclear)r7   rÚ   ÚrequestsrÒ   Úreq_ids        r8   Úfail_all_requestsÚ*ContinuousBatchProcessor.fail_all_requestså  sÈ   € ô ˜Ÿ™×6Ñ6×=Ñ=Ó?Ó@ˆÛˆEØ×&Ñ& uÔ4Ø�N‰N×)Ñ)¨%×*:Ñ*:Ö;ñ ð
 	×Ñ×;Ñ;Ô=Ü˜4Ÿ>™>×:Ñ:×?Ñ?ÓAÖBˆFØ—N‘N×3Ñ3×7Ñ7¸Ó?ˆEØ×&Ñ& uÖ4ñ Cð
 	�‰×-Ñ-×3Ñ3Õ5r<   Úmodelc                 ó’  • U R                   R                  U R                  S9nU R                   R                  5       u  p4nU R                   R                  nU R                   R
                  (       a2  U R                  c  U R                  OU R                  nU R                  nO1U R                  c  U R                  OU R                  nU R                  nU(       dB  Ub  [        R                  R                  U5      O	[        5       n	U	   U" XX4U5        SSS5        OnU R                   R                  5       n
U
b9  [        R                  R                  U5         U
R!                  5         SSS5        OXX4U4nU R"                  " Xv/UQ76   U R                   R%                  5         g! , (       d  f       N)= f! , (       d  f       N:= f)z!Perform a single generation step.©Úuse_paddingN)rt   r¯   r�   Úget_cb_kwargsrŠ   Úuse_block_tablerœ   rš   r�   r˜   r†   rF   rž   r  r   Ú	get_graphÚreplayÚcapture_graphÚretrieve_device_outputs)r7   r@  Ú
batch_dataÚcarry_over_idsÚprev_output_idsÚ
output_idsrŠ   Ú
forward_fnÚuse_cuda_graphr-  ÚgraphÚargss               r8   Ú_generation_stepÚ)ContinuousBatchProcessor._generation_step÷  sq  € ð ×,Ñ,×=Ñ=È$×JZÑJZÐ=Ð[ˆ
Ø6:×6MÑ6M×6[Ñ6[Ó6]Ñ3ˆ¨Ø×0Ñ0×?Ñ?ˆð ×"Ñ"×2×2Ø=A×=RÑ=RÑ=Z˜×9Ò9Ð`d×`uÑ`uˆJØ!×7Ñ7‰Nà=A×=RÑ=RÑ=Z˜×9Ò9Ð`d×`uÑ`uˆJØ!×7Ñ7ˆNö Ø@NÑ@Zœ5Ÿ:™:×,Ñ,¨^Ô<Ô`kÓ`mˆLÚÙ˜5¨nÈzÔZ÷ �ð
 ×+Ñ+×5Ñ5Ó7ˆEàÑ Ü—Z‘Z×&Ñ& ~Õ6Ø—L‘L”N÷ 7Ð6ð ¨>ÈJÐW�Ø×"Ò" :ÐEÀÓEð 	×Ñ×7Ñ7Õ9÷! •ú÷ 7Õ6ús   Ä
F'ÅF8Æ'
F5Æ8
GrN  rŠ   c                 ó‚  • [         R                  R                  U5         U" U6   S S S 5        [         R                  R                  5       n[         R                  R	                  XBU R
                  SS9   U" U6   S S S 5        U R                  R                  U5        g ! , (       d  f       N= f! , (       d  f       N;= f)NÚthread_local)r  ÚpoolÚcapture_error_mode)rF   rž   r  Ú	CUDAGraphrP  r    rt   Ú	set_graph)r7   rN  rŠ   rQ  rP  s        r8   rH  Ú&ContinuousBatchProcessor.capture_graph  sŒ   € ä�Z‰Z×Ñ˜~Õ.Ù˜Ñ÷ /ô —
‘
×$Ñ$Ó&ˆô �Z‰Z×Ñ˜eÀÇÁÐesÐÒtÙ˜Ñ÷ uð 	×Ñ×)Ñ)¨%Õ0÷ /Õ.ú÷ uÕtús    BÁ5B0Â
B-Â0
B>rJ  rK  rL  rM  c                 ó  • U R                   R                  US   X45        U R                  X5      R                  5       nU R                  R
                  (       a  U R                  X&5      OUnU R                  XrS   U5        g)zŸThis function performs the forward pass, logits processing, and sampling; which are broken down into smaller
function to be easier to trace with OpenTelemetry.Ú	input_idsÚlogits_indicesN)rt   Úcarry_over_tokensÚ_model_forwardÚfloatrx   Údo_processingÚ_process_logitÚ_sample)r7   r@  rJ  rK  rL  rM  ÚlogitsÚscoress           r8   rš   Ú4ContinuousBatchProcessor._forward_process_and_sample,  sm   € ð 	×Ñ×1Ñ1°*¸[Ñ2IÈ>ÔkØ×$Ñ$ UÓ7×=Ñ=Ó?ˆØ<@×<PÑ<P×<^×<^�×$Ñ$ ZÔ8ÐdjˆØ�‰�VÐ(8Ñ9¸:ÕFr<   Úmodel_forward©Ú	span_namec                 ó&   • U" S0 UD6R                   $ )Nr6   )rd  )r7   r@  rJ  s      r8   r_  Ú'ContinuousBatchProcessor._model_forward<  s   € áÑ"�zÑ"×)Ñ)Ð)r<   Úlogit_processingrd  c                 óÂ   • UR                   u  p4nUR                  X4-  U5      nUS   R                  X4-  5      nU R                  XvUS   5      nUR                  X4U5      $ )Nr\  Úlogits_processor_args)ÚshapeÚviewrx   )	r7   rJ  rd  Ú
batch_sizeÚseq_lenÚ
vocab_sizeÚ	logits_2dÚinput_ids_2dÚprocessed_logits_2ds	            r8   rb  Ú'ContinuousBatchProcessor._process_logit@  si   € ð +1¯,©,Ñ'ˆ
˜ZØ—K‘K 
Ñ 4°jÓAˆ	Ø! +Ñ.×3Ñ3°JÑ4HÓIˆà"×2Ñ2°<ÈJÐWnÑLoÓpÐà"×'Ñ'¨
¸ZÓHÐHr<   Úsamplingre  r]  c                 óÂ  • U R                   (       d  U R                  (       a"  [        R                  R	                  US   SS9nOUR                  S5      nU R                   (       a  [        R                  " USS9nO[        R                  " USSS9nU R                  (       a/  UR                  SUS9R                  S5      nUR                  5       nUR                  S5      nUR                  S5      nUS U n	XY   nUSS U24   R                  U5        U R                  (       a9  WU	   nUSS U24   R                  UR                  [        R                  S	95        g g )
Nr   éÿÿÿÿ)Údimr   )Únum_samplesT)r{  Úkeepdim)r{  Úindex)r0   )r   r…   r   Ú
functionalÚsoftmaxÚsqueezerF   ÚmultinomialÚargmaxÚgatherÚlogÚsizeÚcopy_rp  Úint32)
r7   re  r]  rM  ÚprobsÚnext_tokensÚper_token_probsr  ÚtokensÚindicess
             r8   rc  Ú ContinuousBatchProcessor._sampleL  s8  € ð �>�>˜T×1×1Ü—M‘M×)Ñ)¨&°©)¸Ð)Ð<‰Eð —N‘N 1Ó%ˆEð �>�>Ü×+Ò+¨E¸qÑA‰KäŸ,š, u°"¸dÑCˆKð ××Ø#Ÿl™l¨q¸˜lÐD×LÑLÈRÓPˆOØ&×*Ñ*Ó,ˆHà!×)Ñ)¨"Ó-ˆð ×!Ñ! !Ó$ˆà   &Ð)ˆØ!Ñ*ˆà�1�g�v�g�:Ñ×$Ñ$ [Ô1Ø××à Ñ(ˆHð �q˜'˜6˜'�zÑ"×(Ñ(¨¯©¼U¿[¹[¨Ð)IÕJð  r<   c           	      óL  • U R                   (       d  [        R                  " S5        gU R                  nU R                  R
                  U R                  R                  -  nX2-
  nU R                  R                  nU R                  (       a  SOSn[        U5       GHÝ  nU R                  (       a-  XpR                  l        [        R                  " SUS-    S35        U R                  UXB-   SS9u  p‰[        R                  " S	U S
U	 S35        [        S[        R                  X$U R                  5      n
 [!        5       nU R                  R#                  X R$                  SX‰U-
  5        U R                  R'                  SS9nU R                  R)                  5       u  pÞnU R*                  =(       d    U R,                  nXXÞU4nU R.                  (       a  U R0                  " UU/UQ76   O-[2        R4                  R7                  U5         U" U6   SSS5        [        R                  " S[!        5       U-
  S S35        U
 H2  nU R                  R=                  UR>                  R@                  5        M4     U R                  RB                  S:X  a  GM÷  [        R                  " S5        Sn[!        5       nSn [        U[        RD                  SU R                  R                  U R                  5      n
U
(       d  GOT U R                  USSS9u  nnU R                  R#                  X R$                  SUS5        U R                  R'                  SS9nU R                  R)                  5       u  pÞnU RF                  =(       d    U R,                  nXXÞU4nU RH                  (       a  U R0                  " UU/UQ76   O-[2        R4                  R7                  U5         U" U6   SSS5        US-  nU
 H2  nU R                  R=                  UR>                  R@                  5        M4     UU R                  :¼  a  O[K        SU-  U R                  5      nGM˜  [        R                  " SU S[!        5       U-
  S S35        GMà     U R                  (       a  SU R                  l        gg! , (       d  f       GN™= f! [8         a%  n[        R:                  " SU S35         SnAGN SnAff = f! U
 H2  nU R                  R=                  UR>                  R@                  5        M4     f = f! , (       d  f       GNL= f! [8         a'  n[        R:                  " SU SU 35         SnAGNvSnAff = f! U
 H2  nU R                  R=                  UR>                  R@                  5        M4     f = f)a7  Pre-capture CUDA graphs (or trigger compile warmup) for varlen and decode paths. In async mode, both IO
pairs are warmed up since each has its own graph buffer and static tensors. The varlen path is warmed up at
the largest possible `(q, kv)` sizes so subsequent captures fit inside it without growing the pool.z6CUDA graphs and compile are disabled, skipping warmup.Nr   r   zWarming up IO pair z/2...F)rê   rë   rì   zWarming up varlen path (z Q tokens, z KV tokens)...TrB  zVarlen warmup completed in ú.2fÚszFailed to warm up varlen path: z-. Graph pool may fragment and OOM under load.r   zWarming up decode fast path...z"Failed to warm up decode path for z requests: zDecode warmup completed (z graphs) in ús.)&r�   r#   Úinfor‘   rv   Ú
num_blocksÚ
block_sizert   rŠ   r¡   r  Úcurrent_pairrï   r)   r"   r  r	   rþ   rx   r¯   rD  r˜   rš   r†   rH  rF   rž   r  rÙ   rÂ   Úfree_blocksrÒ   rZ   r•   r  rœ   r�   Úmin)r7   r@  Únum_query_tokensrî   Únum_cache_tokensrŠ   Únum_io_pairsÚpair_idxÚpadded_qÚ	padded_kvÚfuture_statesÚstartrJ  rK  rL  rM  rN  Úforward_fn_argsrÝ   ÚfsÚdecode_graphsÚnum_requestsÚ_s                          r8   ÚwarmupÚContinuousBatchProcessor.warmupp  sè  € ð ××Ü�KŠKÐPÔQØà×0Ñ0ÐØ—J‘J×)Ñ)¨D¯J©J×,AÑ,AÑAˆ	Ø$Ñ7ÐØ×0Ñ0×?Ñ?ˆð !×3×3‘q¸ˆä˜l×+ˆHØ×&×&Ø7?×'Ñ'Ô4Ü—’Ð1°(¸Q±,°¸uÐEÔFð #'×"7Ñ"7Ø-Ø,Ñ?Ø%*ð #8ð #ÑˆHô
 �KŠKÐ2°8°*¸KÈ	À{ÐR`ÐaÔbä7Ø”=×+Ñ+Ð-=ÐQU×Q[ÑQ[óˆMð@Ü$›�Ø×'Ñ'×=Ñ=Ø!×#7Ñ#7¸ÀÐV^ÑJ^ôð "×4Ñ4×EÑEÐRVÐEÐW�
Ø>B×>UÑ>U×>cÑ>cÓ>eÑ;�°Ø!×2Ñ2×V°d×6VÑ6V�
Ø#(°nÐWaÐ"b�Ø×-×-Ø×&Ò& z°>ÐTÀOÔTäŸ™×*Ñ*¨>Õ:Ù" OÑ4÷ ;ä—’Ð9¼,».È5Ñ:PÐQTÐ9UÐUVÐWÔXó (�BØ—J‘J×*Ñ*¨2¯8©8×+>Ñ+>Ö?ñ (ð �z‰z×0Ñ0°AÓ5Úô �KŠKÐ8Ô9ØˆMÜ “NˆEàˆLØÜ ;Ø ¤-×"8Ñ"8¸!¸T¿Z¹Z×=RÑ=RÐTX×T^ÑT^ó!�ö %ÙðDØ"&×"7Ñ"7Ø%1¸qÐW[ð #8ð #‘K�H˜að ×+Ñ+×AÑAØ%×';Ñ';¸TÀ8ÈQôð "&×!8Ñ!8×!IÑ!IÐVZÐ!IÐ![�JØBF×BYÑBY×BgÑBgÓBiÑ?�N°ZØ!%×!6Ñ!6×!Z¸$×:ZÑ:Z�JØ',¸.Ð[eÐ&f�OØ×1×1Ø×*Ò*¨:°~ÐXÈÔXä"ŸZ™Z×.Ñ.¨~Õ>Ù&¨Ñ8÷ ?à! QÑ&�Mó ,˜ØŸ
™
×.Ñ.¨r¯x©x×/BÑ/BÖCñ ,à 4×#8Ñ#8Ó8ØÜ" 1 |Ñ#3°T×5JÑ5JÓK�ò= ô> �KŠKÐ3°M°?À,Ì|Ë~Ð`eÑOeÐfiÐNjÐjlÐm×nñ] ,ðb ×"×"Ø34ˆD×#Ñ#Õ0ð #÷k ;Ö:ûô ó sÜ—’Ð!@ÀÀÐCpÐq×rÒrûðsûó (�BØ—J‘J×*Ñ*¨2¯8©8×+>Ñ+>Ö?ò (ú÷B ?Ö>ûô !ó fÜ—N’NÐ%GÈÀ~ÐU`ÐabÐ`cÐ#d×eÒeûðfûó ,˜ØŸ
™
×.Ñ.¨r¯x©x×/BÑ/BÖCò ,ús†   Ä7CR4È
R"È.R4Ë=CT5ÏT#ÏT5Ò"
R1	Ò,R4Ò4
S#Ò>SÓS&ÓS#Ó#S&Ó&:T Ô#
T2	Ô-T5Ô5
U&Ô?U!ÕU)Õ!U&Õ&U)Õ):V#)rœ   r˜   r�   rv   r‹   r/   r   r    ry   rt   rŽ   rx   r‘   r’   r|   r}   r¥   rz   r�   r…   ru   r€   r{   r¡   r�   r†   rn   ):rA   rB   rC   rD   r   r   rE   r&   r   r   r   r   r   rM   rN   rJ   rQ   ÚEventrF   r1   r0   rU   rG   r°   r·   r“   rÌ   r   rÞ   rÙ   r!   rÜ   ÚintÚboolÚtuplerï   r  r.  rô   r4  r>  Úno_gradr   ÚModulerR  r
   rž   ÚStreamrH  ÚdictÚTensorrš   r_  rb  rc  Úinference_moder¦  rH   r6   r<   r8   rr   rr   „   st  ‡ à-Ð0JÑJÓJØÓðm
à"ðm
ð !ðm
ð ,ð	m
ð
 %=ðm
ð ?ðm
ð —[‘[ðm
ð $ðm
ð —O‘Oðm
ð —l‘lðm
ð —[‘[ðm
ð ðm
ð 
ôm
ð^
˜#ô 
ô%ô"6ôHTð ó9ó ð9ð$ ðA¨9ð A¸\ð AÈdó Aó ðAð)¨Sð )¸sð )ÐZ^ð )ÐchÐilÐnqÐiqÑcrô )ð ð/ Dó /ó ð/ðb óCEó ðCEðJ ð5 dó 5ó ð5ð ñIó ðIð ð6 yð 6°Tó 6ó ð6ð" Ø
‡]‚]ƒ_ð#: b§i¡ið #:°Dó #:ó ó ð#:ðJ1¨ð 1¸U¿Z¹Z×=NÑ=Nð 1ÐZ^ô 1ð ðGà�y‰yðGð ðGð Ÿ™ð	Gð
 Ÿ™ðGð —L‘LðGð 
óGó ðGñ �oÑ&ð* B§I¡Ið *¸4ð *ÀEÇLÁLó *ó 'ð*ñ Ð(Ñ)ð	I¨ð 	I°u·|±|ð 	IÈÏÉó 	Ió *ð	Iñ �jÑ!ð!K˜eŸl™lð !K¸E¿L¹Lð !KÐV[×VbÑVbð !KÐgkó !Kó "ð!KðF ×ÒÓðc5˜BŸI™Ið c5¨$ó c5ó óc5r<   rr   c                   óL  • \ rS rSrSrS\S\S\SS4S jr\	S,S	 j5       r
S\4S
 jrS,S jrS-S\S\S-  S\SS4S jjrS.S\S\S-  SS4S jjr     S/S\\   S\S-  S\S-  S\S\S\\\   -  S-  S\S\4S jjr   S0S\\\      S\S-  S\S\S\SS4S jjrS\SS4S jrS1S\S-  S\S-  S\S-  4S jjrS rS\S\\   4S jrS\S \SS4S! jr\	S,S" j5       rS\ 4S# jr!\"RF                  " 5       S,S$ j5       r$\	" S%S&9S'\ SS4S( j5       r%\	S)\&S'\ S-  SS4S* j5       r'S+r(g)2ÚContinuousBatchingManageriØ  a¤  Manager for handling continuous batching of generation requests. It provides a user interface for submitting
generation requests, retrieving results, and managing the background generation thread. This class should not be
created directly, but through one of the following entry points (all methods of the `ContinuousMixin` mixin):
- `init_continuous_batching`
- `continuous_batching_context_manager`
- `generate_batch`
r@  r=   rw   r3   Nc                 ó¾  • SUR                   R                  ;  a(  UR                  SUR                   R                   35        UR                  5       U l        X l        X0l        SU l        U R                  R                  U l	        [        R                  " U R                  R                  S9U l        [        R                  " 5       U l        [#        5       U l        [        R                  " 5       U l        SU l        SU l        SU l        [        R.                  " 5       U l        [3        USS5      nUb  UOSU l        [7        U R                  R9                  U5      U R                  R:                  U R                  R<                  S9U l        [A        U R                  R                   5      nU R                  RC                  [3        US	S5      US
9  U R                  RE                  U5      U l#        U R                  RI                  5         U R                  RJ                  U l%        U R                  RL                  U l&        U R                  RN                  U l'        g)zðInitialize the continuous batching manager.

Args:
    model: The language model for generation
    generation_config: Configuration for generation parameters
    continuous_batching_config: Configuration for continuous batching parameters
zpaged|F)ÚmaxsizeNr   Únum_return_sequencesr   )Úlogits_processorÚper_request_processorsÚdrop_unsupported_processorsr�   )r�   Úis_attn_mask_needed)(r/   rÃ   r9   Úevalr@  r=   rw   Ú	warmed_upÚallow_block_sharingÚ_use_prefix_sharingrM   rN   Úmax_queue_sizery   rQ   r¨  Ú_has_new_requestsrJ   rz   r{   Úbatch_processorÚ_generation_threadÚ_request_counterrR   Ú_request_lockrŒ   r¶  r   r?   r¸  r¹  rx   r(   Údecide_use_cuda_graphsÚdecide_use_async_batchingr¡   Úresolve_sentinel_valuesr�   rŽ   r¢   )r7   r@  r=   rw   r¶  rº  s         r8   rU   Ú"ContinuousBatchingManager.__init__â  sç  € ð ˜5Ÿ<™<×<Ñ<Ó<Ø×)Ñ)¨F°5·<±<×3TÑ3TÐ2UÐ*VÔWð —Z‘Z“\ˆŒ
Ø!2ÔØ*DÔ'ØˆŒà#'×#BÑ#B×#VÑ#VˆÔ ä Ÿ;š;¨t×/NÑ/N×/]Ñ/]Ñ^ˆÔÜ!*§¢Ó!2ˆÔÜ)›^ˆÔÜ#Ÿ/š/Ó+ˆŒØ@DˆÔØ"&ˆÔØ !ˆÔÜ&Ÿ^š^Ó-ˆÔô  'Ð'8Ð:PÐRVÓWÐØ<PÑ<\Ñ$8ÐbcˆÔ!äDØ!ŸZ™Z×=Ñ=Ð>OÓPØ#'×#BÑ#B×#YÑ#YØ(,×(GÑ(G×(cÑ(cñ 
ˆÔô 2°$·*±*×2CÑ2CÓDÐØ×'Ñ'×>Ñ>Ü"Ð#4Ð6FÈÓMØ 3ð 	?ñ 	
ð
 #'×"AÑ"A×"[Ñ"[Ð\oÓ"pˆÔð 	×'Ñ'×?Ñ?ÔAØ'+×'FÑ'F×'^Ñ'^ˆÔ$Ø(,×(GÑ(G×(`Ñ(`ˆÔ%Ø!%×!@Ñ!@×!RÑ!RˆÕr<   c                 ó8  • U R                   b6  U R                   R                  5       (       a  [        R                  " S5        gU R                  R                  5         [        R                  " U R                  S9U l         U R                   R                  5         g)z'Start the background generation thread.Nz"Manager thread is already running.)Útarget)
rÂ  Úis_aliver#   rÂ   r{   r;  rQ   ÚThreadÚ_run_generation_loopr   rT   s    r8   r   ÚContinuousBatchingManager.start  so   € ð ×"Ñ"Ñ.°4×3JÑ3J×3SÑ3S×3UÑ3UÜ�NŠNÐ?Ô@ØØ�‰×ÑÔÜ"+×"2Ò"2¸$×:SÑ:SÑ"TˆÔØ×Ñ×%Ñ%Õ'r<   c                 ó`   • U R                   SL=(       a    U R                   R                  5       $ )z5Check if the background generation thread is running.N)rÂ  rË  rT   s    r8   Ú
is_runningÚ$ContinuousBatchingManager.is_running(  s'   € à×&Ñ&¨dÐ2×Y°t×7NÑ7N×7WÑ7WÓ7YÐYr<   c                 ó    • U R                   c  U R                  5       U l         U R                   R                  U R                  5        SU l        g)z‚Pre-capture CUDA graphs for varlen and decode paths by running dummy batches. Initializes the batch
processor if not already done.NT)rÁ  Ú_create_batch_processorr¦  r@  r¼  rT   s    r8   r¦  Ú ContinuousBatchingManager.warmup,  s@   € ð ×ÑÑ'Ø#'×#?Ñ#?Ó#AˆDÔ Ø×Ñ×#Ñ# D§J¡JÔ/Øˆ�r<   ÚblockÚtimeoutÚkeep_for_next_sessionc                 ób  • U R                   c  [        R                  " S5        O\U R                   R                  R                  (       a7  [        R
                  " SU R                   R                  R                   35        U R                  c%  U(       a  SOSn[        R                  " SU-   5        g[        5       nU R                  R                  5       (       d0  U R                  R                  5         [        R
                  " S5        U(       a  U R                  XR5        U(       d  SU l         O&[        R
                  " S5        X R                  l        [        R                   " 5         ["        R$                  R'                  5       (       a  ["        R$                  R)                  5         gg)	zåSignal the background thread to stop.

Args:
    block: Whether to wait for the thread to stop
    timeout: Maximum time to wait for the thread to stop
    keep_for_next_session: Whether to cache this on the model for future use
Nz%
Batch processor was not initialized.z-
Prefix sharing was on. Total prefix length: z? Hence the unstarted manager will not be kept for next session.Ú zManager not started.z'Stopping continuous batching manager...z:Continuous batching manager will be kept for next session.)rÁ  r#   rÂ   rv   Úuse_prefix_sharingr“  Ú_total_prefix_lengthrÂ  r	   r{   Úis_setÚsetÚjoinr@  Ú#_cached_continuous_batching_managerr³   r´   rF   rž   rµ   r¶   )r7   rÕ  rÖ  r×  ÚsuffixÚstop_trigger_times         r8   ÚstopÚContinuousBatchingManager.stop5  s(  € ð ×ÑÑ'Ü�NŠNÐCÕDØ×!Ñ!×'Ñ'×:×:Ü�KŠKØ@À×AUÑAU×A[ÑA[×ApÑApÐ@qÐrôð ×"Ñ"Ñ*ÞZoÑVÐuwˆFÜ�NŠNÐ1°FÑ:Ô;Øä(›NÐØ�‰×%Ñ%×'Ñ'Ø�O‰O×ÑÔ!Ü�KŠKÐAÔBæØ�I‰IÐ'Ô1ö %Ø#'ˆDÕ ô �KŠKÐTÔUØ=A�J‰JÔ:ä
�
Š
ŒÜ�:‰:×"Ñ"×$Ñ$Ü�J‰J×"Ñ"Õ$ð %r<   rá  c                 ó"  • U R                   b‚  U R                   R                  US9  U R                   R                  5       (       a  [        R                  " SU S35        g[        5       n[        R                  " SX1-
  S S35        SU l         gg)zjWait for the background thread to finish.

Args:
    timeout: Maximum time to wait for the thread to stop
N©rÖ  z3Generation thread did not exit after join timeout (z).z*Continuous Batching Manager stopped after r�  r’  )rÂ  rÞ  rË  r#   rÂ   r	   r“  )r7   rá  rÖ  Úends       r8   rÞ  ÚContinuousBatchingManager.join]  s‡   € ð ×"Ñ"Ñ.Ø×#Ñ#×(Ñ(°Ð(Ñ9Ø×&Ñ&×/Ñ/×1Ñ1Ü—’Ð!TÐU\ÐT]Ð]_Ð`Õaä"“n�Ü—’ÐHÈÑI`ÐadÐHeÐegÐhÔiØ*.�Õ'ð /r<   r\  rZ   Úmax_new_tokensr  Úrecord_timestampsÚeos_token_idrÖ   c                 óÂ  • Uc9  U R                      SU R                   3nU =R                  S-  sl        SSS5        Uc  U R                  R                  OUnUc  U R                  R                  OUn[        U[        U5      U R                  S-
  UUUUUS9nU R                  R                  USSS9  U R                  R                  5         U$ ! , (       d  f       N¡= f)a  Add a new generation request to the queue.

Args:
    input_ids: Input token IDs to use as prompt
    request_id: Optional custom request ID (auto-generated if None)
    max_new_tokens: Maximum number of new tokens to generate
    streaming: Whether to stream tokens as they're generated
    record_timestamps: Whether to record timestamps for each generated token
    eos_token_id: End-of-sequence token ID(s)
    logit_processor_kwargs: Keyword arguments for the logits processor.

Returns:
    str: The request ID
NÚreq_r   )rZ   Úinitial_tokensr  ré  rè  rê  r  rÖ   Té
   ©rÕ  rÖ  )rÄ  rÃ  r=   rè  rê  r!   rp   r¶  ry   r\   rÀ  rÝ  )	r7   r\  rZ   rè  r  ré  rê  rÖ   rÒ   s	            r8   Úadd_requestÚ%ContinuousBatchingManager.add_requestl  sã   € ð0 ÑØ×#Ó#Ø# D×$9Ñ$9Ð#:Ð;�
Ø×%Ò%¨Ñ*Õ%÷ $ð CQÑBX˜×/Ñ/×>Ò>Ð^lˆØ>JÑ>R�t×-Ñ-×:Ò:ÐXdˆô Ø!Ü 	›?Ø×2Ñ2°QÑ6Ø/Ø)Ø%ØØ#9ñ	
ˆð 	×Ñ×Ñ˜U¨$¸ÐÑ;Ø×Ñ×"Ñ"Ô$ØÐ÷- $Õ#ús   �%CÃ
CÚinputsc                 óD  • U R                      [        U R                  U R                  [        U5      -   5       Vs/ s H  nSU 3PM
     nnU =R                  [        U5      -  sl        S S S 5        [	        [        WU5      5      nU R                  (       a  [        US SS9nU R                  R                  n	U	c   U R                  R                  R                  OU	n	U	c  SOU	n	U H  u  p«U R                  " SUU
UUUU	S.UD6  M      g s  snf ! , (       d  f       N¬= f)Nrì  c                 ó   • U S   $ )Nr   r6   r»   s    r8   r½   Ú8ContinuousBatchingManager.add_requests.<locals>.<lambda>¬  s   € À!ÀAÂ$r<   T)ÚkeyÚreverserz  )r\  rZ   rè  r  ré  rê  r6   )rÄ  r  rÃ  rö   rp   Úzipr¾  Úsortedr=   rê  r@  r/   rð  )r7   rò  rè  r  ré  rÖ   r(  Úrequest_idsÚids_and_inputsrê  rZ   r\  s               r8   Úadd_requestsÚ&ContinuousBatchingManager.add_requests�  s  € ð ×ÓÜ/4°T×5JÑ5JÈD×LaÑLaÔdgÐhnÓdoÑLoÔ/pÓqÒ/p¨!˜T ! ›:Ñ/pˆKÐqØ×!Ò!¤S¨£[Ñ0Õ!÷  ô œc +¨vÓ6Ó7ˆØ×#×#Ü# N¹ÐPTÑUˆNð ×-Ñ-×:Ñ:ˆØ9EÑ9M�t—z‘z×(Ñ(×5Ò5ÐS_ˆØ)Ñ1‘r°|ˆã%3Ñ!ˆJØ×Òð Ø#Ø%Ø-Ø#Ø"3Ø)ñð )ôò &4ùò r÷  Õús   �/D¼DÁ DÄDÄ
Dc                 ój   • U R                   b&  U R                   R                  R                  U5        gg)zSCancel a request by its ID.

Args:
    request_id: The ID of the request to cancel
N)rÁ  ru   Úset_request_cancellation)r7   rZ   s     r8   Úcancel_requestÚ(ContinuousBatchingManager.cancel_request¾  s/   € ð ×ÑÑ+Ø× Ñ ×*Ñ*×CÑCÀJÕOð ,r<   c                 ód  • U R                   c*  U R                  R                  R                  5       (       a  g U R                  R                  R	                  SUS9nUb6  UR
                  U:w  a&  U R                  R                  R                  U5        gU$ ! [        R                   a     gf = f)a  Retrieve one result from the output queue.

Args:
    request_id: If set, only return results matching this ID (others are requeued).
    timeout: Maximum time to wait for a result.

Returns:
    Optional[GenerationOutput]: The result data or None if timeout.
NTrï  )	rÂ  rz   rO   rÓ   rY   rZ   r\   rM   rØ   )r7   rZ   rÖ  Úresults       r8   Ú
get_resultÚ$ContinuousBatchingManager.get_resultÈ  sž   € ð ×"Ñ"Ñ*¨t×/AÑ/A×/NÑ/N×/TÑ/T×/VÑ/VØð	Ø×'Ñ'×4Ñ4×8Ñ8¸tÈWÐ8ÐUˆFØÑ%¨&×*;Ñ*;¸zÓ*IØ×"Ñ"×/Ñ/×3Ñ3°FÔ;ØØˆMøÜ�{‰{ó 	Ùð	ús   ¹AB ÂB ÂB/Â.B/c              #   óò   #   • U R                   bf  U R                   R                  5       (       aF  U R                  SS9nUb  Uv •  U R                   b"  U R                   R                  5       (       a  MD  gggg7f)z.Iterate over results as they become available.Nçš™™™™™¹?rå  )rÂ  rË  r  )r7   r  s     r8   Ú__iter__Ú"ContinuousBatchingManager.__iter__Ý  sk   é € à×%Ñ%Ñ1°d×6MÑ6M×6VÑ6V×6XÑ6XØ—_‘_¨S�_Ð1ˆFØÑ!Ø’ð ×%Ñ%Ñ1°d×6MÑ6M×6VÑ6V×6XÔ6XÐ1Ð6XÐ1ùs   ‚A/A7Á3A7c              #   ó   #   • U R                   b}  U R                   R                  5       (       a]  U R                  USS9nUb  Uv •  UR                  5       (       a  gU R                   b"  U R                   R                  5       (       a  M[  gggg7f)z·Iterate over results matching a specific request id (blocking).

Uses the shared output queue with requeue. For high-concurrency serving,
use :meth:`register_result_handler` instead.
Nr  )rZ   rÖ  )rÂ  rË  r  r$  )r7   rZ   r  s      r8   Úrequest_id_iterÚ)ContinuousBatchingManager.request_id_iterä  s�   é € ð ×%Ñ%Ñ1°d×6MÑ6M×6VÑ6V×6XÑ6XØ—_‘_°
ÀC�_ÐHˆFØÑ!Ø’Ø×%Ñ%×'Ñ'Øð ×%Ñ%Ñ1°d×6MÑ6M×6VÑ6V×6XÔ6XÐ1Ð6XÐ1ùs   ‚BBÂ
Br^   c                 óØ   ^ ^^• [         R                  " 5       nUUU 4S jnT R                  R                     XC4T R                  R                  T'   SSS5        g! , (       d  f       g= f)aô  Register a callback for result delivery (streaming or non-streaming).

The callback is invoked on the event loop via ``call_soon_threadsafe``
each time a result is produced for this request. For streaming requests,
this happens on every token; for non-streaming, only on completion.

The handler is automatically cleaned up when the request finishes.

Args:
    request_id (`str`): The request ID to receive outputs for.
    callback (`callable`): Called with a ``GenerationOutput`` for each result.
c                 óî   >• T" U 5        U R                  5       (       aF  TR                  R                     TR                  R                  R	                  TS 5        S S S 5        g g ! , (       d  f       g = fr5   )r$  rz   rS   rP   r  )r  r^   rZ   r7   s    €€€r8   Ú_auto_cleanupÚHContinuousBatchingManager.register_result_handler.<locals>._auto_cleanup   sY   ø€ Ù�VÔØ×!Ñ!×#Ñ#Ø×'Ñ'×-Ó-Ø×&Ñ&×6Ñ6×:Ñ:¸:ÀtÔL÷ .Ð-ð $ß-Õ-ús   µ'A&Á&
A4N)ÚasyncioÚget_running_looprz   rS   rP   )r7   rZ   r^   r_   r  s   ```  r8   Úregister_result_handlerÚ1ContinuousBatchingManager.register_result_handlerñ  sO   ú€ ô ×'Ò'Ó)ˆ÷	Mð ×Ñ×%Ó%Ø>KÐ=RˆD×Ñ×.Ñ.¨zÑ:÷ &×%Ö%ús   ·AÁ
A)c                 ó~   • U R                   c  [        S5      eU R                   R                  U R                  5        g)z=Perform a single generation step. This is mostly cuda graphedNzNTried to perform a generation step before the batch processor was initialized.)rÁ  rù   rR  r@  rT   s    r8   rR  Ú*ContinuousBatchingManager._generation_step	  s4   € ð ×ÑÑ'ÜÐoÓpÐpØ×Ñ×-Ñ-¨d¯j©jÕ9r<   c                 ó  • U R                   R                  U R                  R                  5        [	        U R
                  R                  U R                   U R
                  R                  U R
                  R                  [        U R
                  SS 5      S9nUR                  U l        U R                   R                  n[        R                  " US 5      nUc   [        R                   " SU S35        ["        n[%        UU R
                  R                  U R&                  U R                   U R                  U R(                  U R*                  U R,                  U R
                  R                  U R
                  R                  U" U5      S9nU$ )NÚ_tp_size)Útp_sizezScheduler 'z ' not found. Defaulting to FIFO.)rv   r/   r=   rw   rx   ry   rz   r{   r|   r}   ru   )rw   Úresolve_max_memory_percentrx   ra  r   r@  r/   r1   r0   rŒ   rÚ  r¾  Úscheduler_typer$   rY   r#   rÂ   r%   rr   r=   ry   rz   r{   )r7   Úpaged_attention_cacher  ru   rÁ  s        r8   rÓ  Ú1ContinuousBatchingManager._create_batch_processor  s?  € à×'Ñ'×BÑBÀ4×CWÑCW×CeÑCeÔfä 3Ø�J‰J×ÑØ×+Ñ+Ø�J‰J×ÑØ�J‰J×ÑÜ˜DŸJ™J¨
°DÓ9ñ!
Ðð $9×#KÑ#KˆÔ ð ×8Ñ8×GÑGˆÜ%×)Ò)¨.¸$Ó?ˆ	ØÑÜ�NŠN˜[¨Ð(8Ð8XÐYÔZÜ%ˆIô 3Ø'Ø—:‘:×$Ñ$Ø"×4Ñ4Ø'+×'FÑ'FØ ×0Ñ0Ø×(Ñ(Ø×,Ñ,Ø—‘ØŸ™×*Ñ*ØŸ
™
×(Ñ(ÙÐ 5Ó6ñ
ˆð Ðr<   c                 ó  •  [        U SS5      n[        U[        5      (       a  UR                  5         OU R	                  5       nXl        SU l        UR                  (       aE  UR                  5       (       d  [        S5      eU R                  5         U =R                  S-  sl        U R                  R                  5       (       a  UR                  5       (       a^  U R                  U5        U =R                  S-  sl        U R                  R                  5       (       d  MG  UR                  5       (       a  M^  [        UR                  [         5      (       a8  SUR                  R"                  -
  UR                  l        UR%                  5         [(        R.                  " S	5        g! [&         a4  n[(        R*                  " SU 3SS9  U R-                  UW5         SnANPSnAff = f! [(        R.                  " S	5        f = f)
z6Main processing loop running in the background thread.rÁ  Nr   z$Failed to bootstrap the first batch.r   zError in generation loop: TrÐ   zGeneration loop finished.)rŒ   rã   rr   rÌ   rÓ  rÁ  Úcurrent_batchr¡   r  rù   rR  r{   rÜ  rô   Ú_inner_generation_looprt   r   r–  r.  rÙ   r#   rÚ   Ú_handle_critical_errorr“  )r7   rÁ  rÝ   s      r8   rÍ  Ú.ContinuousBatchingManager._run_generation_loop4  s‘  € ð#	5ä% dÐ,=¸tÓDˆOä˜/Ô+C×DÑDØ×%Ñ%Õ'ð #'×">Ñ">Ó"@�ð $3Ô Ø!"ˆDÔð ×1×1Ø&×9Ñ9×;Ñ;Ü&Ð'MÓNÐNØ×%Ñ%Ô'Ø×"Ò" aÑ'Õ"à—‘×-Ñ-×/Ñ/°O×4XÑ4X×4ZÑ4ZØ×+Ñ+¨OÔ<Ø×"Ò" aÑ'Õ"ð —‘×-Ñ-×/Ó/°O×4XÑ4X×4ZÓ4Zô ˜/×<Ñ<Ô>X×YÑYØBCÀo×FhÑFh×FuÑFuÑBu�×2Ñ2Ô?Ø×,Ñ,Ô.ô �KŠKÐ3Õ4øô	 ó 	<Ü�LŠLÐ5°a°SÐ9ÀDÒIØ×'Ñ'¨¨?×;Ñ;ûð	<ûô �KŠKÐ3Õ4ús7   ‚DF( Ä#F( Ä:AF( Æ(
G&Æ2*G!ÇG) Ç!G&Ç&G) Ç)HÚgeneration_looprh  rÁ  c                 óÖ   • UR                  5       (       d4  U R                  R                  SS9  U R                  R                  5         g U R	                  5         UR                  5         g )Nr  rå  )r  rÀ  Úwaitr;  rR  r.  )r7   rÁ  s     r8   r   Ú0ContinuousBatchingManager._inner_generation_loop\  sW   € ð ×1Ñ1×3Ñ3à×"Ñ"×'Ñ'°Ð'Ñ4Ø×"Ñ"×(Ñ(Ô*ØØ×ÑÔØ×$Ñ$Õ&r<   rÚ   c                 óú   • U R                   R                  5           U R                  R                  5       nUb  UR	                  X5        M0  ! [
        R                   a     Of = fUb  UR                  U5        gg)z:Handle critical errors that terminate the generation loop.N)r{   rÝ  ry   rÔ   rÜ   rM   rØ   r>  )r7   rÚ   rÁ  Úreq_datas       r8   r!  Ú0ContinuousBatchingManager._handle_critical_errorg  s|   € ð 	�‰×ÑÔð	ØØ×+Ñ+×6Ñ6Ó8�Ø"Ñ.Ø#×9Ñ9¸%ÔJñ øô �{‰{ó 	Ùð	úð Ñ&Ø×-Ñ-¨eÕ4ð 's   œ1A ÁA$Á#A$)rÂ  rÀ  rÃ  rÄ  r¾  rÁ  rw   r  r=   ry   rŽ   rx   r¢   r@  r¶  rz   r�   r{   r¡   r¼  rn   )TNFr5   )NNFFN)NFF)NN))rA   rB   rC   rD   ro   r-   r   r   rU   r   r   rª  rÐ  r¦  r`  râ  rÞ  rp   r©  rG   r
   rð  rü  r   r    r  r  r   r  r   r  rR  rr   rÓ  rF   r±  rÍ  r   rÙ   r!  rH   r6   r<   r8   r³  r³  Ø  s�  † ñð:Sà#ð:Sð ,ð:Sð %=ð	:Sð
 
ô:Sðx ó(ó ð(ðZ˜Dô Zôñ&%˜$ð &%°¸±ð &%Ð\`ð &%Ðmqõ &%ñP/ eð /°e¸d±lð /Èdõ /ð$ "&Ø%)ØØ"'Ø/3ñ/à˜‘9ð/ð ˜$‘Jð/ð ˜d™
ð	/ð
 ð/ð  ð/ð ˜D ™I‘o¨Ñ,ð/ð #&ð/ð 
õ/ðh &*ØØ"'ñà�T˜#‘Y‘ðð ˜d™
ðð ð	ð
  ðð #&ðð 
õðBP¨ð P°ô Pñ S¨4¡Zð ÀÈÁð ÐYiÐlpÑYpõ ò*ð¨#ð °)Ð<LÑ2Mô ðS°#ð SÀð SÈdô Sð0 ó:ó ð:ð"Ð)Aô "ðH ×ÒÓó%5ó ð%5ñN Ð'Ñ(ð'Ð6Nð 'ÐSWó 'ó )ð'ð ð5¨Ið 5ÐH`ÐcgÑHgð 5Ðlpó 5ó ó5r<   r³  c                   ó¦  • \ rS rSr% Sr\\S'   \R                  " 5          SS\S-  S\	S-  S\
S-  S\4S jj5       rSS	 jr\\R                  " 5              SS\S-  S
\S\S-  S\	S-  S\S\S\
S-  S\\   4S jj5       5       r\\R                  " 5             SS\\\      S\S-  S\	S-  S\S\S\S\S\\\4   4S jj5       5       rSrg)ÚContinuousMixini{  aü  Mixin class for models to add continuous batching capabilities. Continuous batching has three entry points:
- `init_continuous_batching`, which is the actual entry point for continuous batching
- `continuous_batching_context_manager`, which itself is a wrapper around `init_continuous_batching`
- `generate_batch`, which is really a wrapper around `continuous_batching_context_manager`

They are defined in this order. Any change made to any of those three entry points should be reflected in the other
two.
r=   Nrw   Úworkload_hintsr3   c                 óX  • [        U S5      (       a"  [        U S5      (       a  [        U S5      (       d  [        S5      e[        U SS5      n[        U[        5      (       a  [
        R                  " S5        U$ Ub  UOU R                  nUc  [        S5      eUR                  c  [
        R                  " S	5        S
Ul	        Uc7  [        [        USS5      [        5      (       a  UR                  nO
[        5       nUR                  " S0 UD6  Ub  UR                  U5        [	        XUS9$ )aÚ  Initialize a manager for continuous batching inference.

Args:
    generation_config: An optional generation configuration, which may contain a CompileConfig object
    continuous_batching_config: An optional continuous batching configuration
    workload_hints: Optional WorkloadHints to help the continuous batching manager make better decisions for
        default values
    **deprecated_kwargs: Deprecated arguments that are now passed in the continuous_batching_config. Those are:
        max_queue_size, q_padding_interval_size, kv_padding_interval_size, allow_block_sharing,
        use_async_batching, max_cached_graphs
Returns:
    `ContinuousBatchingManager`: The manager instance to add requests and retrieve results.
r/   r1   r0   z;Model must have 'config', 'device', and 'dtype' attributes.rß  NzÄCached continuous batching manager found: it will be re-used instead of creating a new one. If you want to create a new manager, you should call `destroy_cached_continuous_batching_manager` first.z8A GenerationConfig must be provided or set in the model.zE`eos_token_id` not set in GenerationConfig. Setting to -1 (disabled).rz  rw   )r@  r=   rw   r6   )ÚhasattrÚAttributeErrorrŒ   rã   r³  r#   r“  r=   Ú
ValueErrorrê  rÂ   r   rw   Ú#account_for_cb_deprecated_argumentsÚresolve_using_hints)r7   r=   rw   r,  Údeprecated_kwargsÚcached_managerÚ
gen_configs          r8   Úinit_continuous_batchingÚ(ContinuousMixin.init_continuous_batching‡  s)  € ô, �t˜X×&Ñ&¬g°d¸H×.EÑ.EÌWÐUYÐ[b×McÑMcÜ Ð!^Ó_Ð_ô ! Ð'LÈdÓSˆÜ�nÔ&?×@Ñ@Ü�KŠKðuôð "Ð!ð +<Ñ*GÑ&ÈT×McÑMcˆ
ØÑÜÐWÓXÐXà×"Ñ"Ñ*Ü�NŠNÐbÔcØ&(ˆJÔ#ð &Ñ-Üœ' *Ð.JÈDÓQÔSk×lÑlØ-7×-RÑ-RÑ*ä-EÓ-GÐ*Ø"×FÒFÑ[ÐIZÒ[ØÑ%Ø×.Ñ.Ð/IÔJô )ØÐQkñ
ð 	
r<   c                 ó„   • [        U SS5      n[        U[        5      (       a  UR                  SSSS9  [	        U S5        gg)zFDestroy the cached continuous batching manager and free GPU resources.rß  NTF©rÕ  rÖ  r×  )rŒ   rã   r³  râ  Údelattr)r7   r4  s     r8   Ú*destroy_cached_continuous_batching_managerÚ:ContinuousMixin.destroy_cached_continuous_batching_managerÁ  sF   € ä  Ð'LÈdÓSˆÜ�nÔ&?×@Ñ@Ø×Ñ d°DÐPUÐÑVÜ�DÐ?Õ@ð Ar<   rÕ  rÖ  Úpersistent_managerr¦  c              +   óà  #   • U R                   " S	UUUS.UD6n	U(       ag  U	R                  (       dV  [        R                  " S5        [	        5       n
U	R                  5         [        R                  " S[	        5       U
-
  S S35        U	R                  5          U	v •  [        R                  " S5        U	R                  X#US9  g! [        R                  " S5        U	R                  X#US9  f = f7f)
a¤  A context manager to safely use the continuous batching manager. Arguments are similar to the ones of
`init_continuous_batching`, except for:
    - block: whether to block the thread when stopping the manager. Default is True.
    - timeout: maximum time to wait for the thread to stop. Default is None (no timeout).
    - warmup: whether to pre-capture CUDA graphs at the largest sizes before running. Default is True.
)r=   rw   r,  z%Warming up for continuous batching...zWarming up completed in r�  r’  z!Continuous batching loop finishedr9  Nr6   )	r6  r¼  r#   rÂ   r	   r¦  r   rü   râ  )r7   r=   rÕ  rÖ  rw   r=  r¦  r,  r3  Úmanagerr   s              r8   Ú#continuous_batching_context_managerÚ3ContinuousMixin.continuous_batching_context_managerÈ  sÐ   é € ð& ×/Ò/ð 
Ø/Ø'AØ)ñ
ð  ñ	
ˆö ˜'×+×+ä�NŠNÐBÔCÜ “NˆEØ�N‰NÔÜ�NŠNÐ5´l³nÀuÑ6LÈSÐ5QÐQSÐTÔUØ�‰Œð	aØŠMô �LŠLÐ<Ô=Ø�L‰L˜uÐM_ˆLÒ`øô �LŠLÐ<Ô=Ø�L‰L˜uÐM_ˆLÒ`üs   ‚BC.ÂC Â'C.Ã(C+Ã+C.rò  ré  Úprogress_barc                 ón  • U(       d  0 $ [         R                  " 5       [        R                  ::  a  [         R                  " S5        Sn0 n	/ SQn
U
 H  nX¸;   d  M
  UR                  U5      X›'   M     Uc  U R                  OUnUR                  b  UR                  OSn[        U5      U-  nUR                  SS5      nUc  UR                  OUn[        [        S U 5       5      Ub  UOSS	9nU R                  " SUUS
SUUUS.U	D6n[        [         /5      n[        UU(       + SU S3SS9n0 nSnU nU   U n UR                  XUS9  UU:  a’  UR!                  SS9nU(       a=  UR"                  nUR%                  5       (       a  UUU'   US-  nUR'                  S5        O7UR)                  5       (       d"  [         R*                  " S5        [-        S5        OUU:  a  M’  SSS5        SSS5        SSS5        0 n[1        [        U5      5       H>  nUR3                  SU 35      nUb
  UUSU 3'   M$  [         R*                  " SU S35        M@     U$ ! [.         a"  n[         R*                  " SU 3S
S9   SnAN™SnAff = f! , (       d  f       N§= f! , (       d  f       N°= f! , (       d  f       N¹= f)aŒ  Generate sequences for a batch of prompts using continuous batching.

Args:
    inputs: List of input token sequences (prompts)
    generation_config: Optional generation configuration
    continuous_batching_config: Optional continuous batching configuration
    record_timestamps: If set to true, the requests will have a timestamp for each token generated
    progress_bar: If set to true, a progress bar will be displayed
    persistent_manager: whether to persist the manager after the generation is finished. Default is False.
    warmup: whether to pre-capture CUDA graphs before processing requests. Default is True.
    **kwargs: Additional generation parameters. Only max_new_tokens is used, but other deprecated arguments
        are extracted and passed to the continuous_batching_config object.
Returns:
    `dict[str, GenerationOutput]`: a dictionary of request ids to GenerationOutput objects
z=Progress bar is disabled when logger level is less than DEBUGF)r�   rŽ   r½  r¡   r¢   r¿  Nr   rè  c              3   ó8   #   • U  H  n[        U5      v •  M     g 7fr5   )rö   )Ú.0r\  s     r8   Ú	<genexpr>Ú1ContinuousMixin.generate_batch.<locals>.<genexpr>/  s   é € Ð!IÂ&°Y¤# i§. .Â&ùs   ‚r   )Úmax_prompt_lengthÚmax_generated_lengthTé   )r=   rw   rÕ  rÖ  r=  r¦  r,  zSolving z	 requestsÚrequest)ÚtotalÚdisableÚdescÚunit)rò  rè  ré  rå  z*Generation thread terminated unexpectedly.zCReturning results of generate_batch despite unexpected termination.zError during batch generation: rÐ   rì  zRequest req_z not found in results.r6   )r#   ÚgetEffectiveLevelr   ÚDEBUGrÂ   r  r=   r¶  rö   rè  r'   Úmaxr@  r   r   rü  r  rZ   r$  ÚupdaterÐ  rÚ   ÚprintrÙ   r  rY   )r7   rò  r=   rw   ré  rB  r=  r¦  Úkwargsr3  Údeprecated_keysÚdepr_keyÚgen_cfgr¶  r¤  rè  r,  Ú
manager_cmÚ
logging_cmÚpbar_cmÚresultsÚfinished_countr?  Úpbarr  r=  rÝ   Úreordered_resultsr(  s                                r8   Úgenerate_batchÚContinuousMixin.generate_batchð  s¹  € ö: ØˆIô ×#Ò#Ó%¬¯©Ó6Ü�NŠNÐZÔ[Ø ˆLð Ðò
ˆó (ˆHØÕ!Ø.4¯j©j¸Ó.BÐ!Ó+ñ (ð
 ->Ñ,E�$×(Ò(ÐK\ˆØ?F×?[Ñ?[Ñ?g˜w×;Ò;ÐmnÐÜ˜6“{Ð%9Ñ9ˆð  Ÿ™Ð$4°dÓ;ˆØ3AÑ3I˜×/Ò/È~ˆô 'Ü!Ñ!IÁ&Ó!IÓIØ3AÑ3M¡ÐSTñ
ˆð ×=Ò=ð 	
Ø/Ø'AØØØ1ØØ)ñ	
ð  ñ	
ˆ
ô +¬F¨8Ó4ˆ
ÜØØ%Ô%Ø˜L˜>¨Ð3Øñ	
ˆð ˆØˆÙ˜7¢J±¸4ðSØ×$Ñ$¨FÐevÐ$ÑwØ$ |Ó3Ø$×/Ñ/¸Ð/Ð:�FÞØ!'×!2Ñ!2˜Ø!×-Ñ-×/Ñ/Ø.4˜G F™OØ*¨aÑ/˜NØ ŸK™K¨œNøØ$×/Ñ/×1Ñ1ÜŸšÐ%QÔRäÐcÔdØð % |Õ3÷ 18§J�Zð* ÐÜ”s˜6“{Ö#ˆAà—[‘[ 4¨ s Ó,ˆFØÑ!Ø06Ð! D¨¨ *Ó-ä—’˜|¨A¨3Ð.DÐEÖFñ $ð !Ð øô ó SÜ—’Ð>¸q¸cÐBÈT×RûðSú÷# 18µú§J¥Jú�Z�Zúsm   Ä3J&Ä6JÄ9JÄ;BIÇJÇIÇ#JÇ+J&É
J	ÉI<	É7JÉ<J	ÊJÊ
JÊJÊ
J#	ÊJ&Ê&
J4r6   )NNNrn   )NTNNFTN)NNFTFT)rA   rB   rC   rD   ro   r   rE   rF   r±  r   r'   r³  r6  r;  r   rª  r`  r   r@  r   rp   r©  r¯  rG   r    r`  rH   r6   r<   r8   r+  r+  {  sÍ  ‡ ñð (Ó'à
×ÒÓð 6:ØFJØ/3ñ	7
à+¨dÑ2ð7
ð %=¸tÑ$Cð7
ð &¨Ñ,ð	7
ð 
#ô7
ó ð7
ôrAð Ø
×ÒÓð 6:ØØ $ØFJØ#(ØØ/3ñ#aà+¨dÑ2ð#að ð#að ˜‘ð	#að
 %=¸tÑ$Cð#að !ð#að ð#að &¨Ñ,ð#að 
Ð,Ñ	-ô#aó ó ð#aðL Ø
×ÒÓð 6:ØFJØ"'Ø!Ø#(Øñt!à�T˜#‘Y‘ðt!ð ,¨dÑ2ðt!ð %=¸tÑ$Cð	t!ð
  ðt!ð ðt!ð !ðt!ð ðt!ð 
ˆcÐ#Ð#Ñ	$ôt!ó ó ót!r<   r+  )Fr  r³   rM   rQ   Úabcr   Úcollections.abcr   r   Ú
contextlibr   r   Úmathr   Útimer	   Útypingr
   rF   r   r   Útqdm.contrib.loggingr   Úconfiguration_utilsr   Úgeneration.configuration_utilsr   r   Úmodeling_flash_attention_utilsr   Úutils.genericr   Úutils.loggingr   Úutils.metricsr   r   r   Úlogits_processr   rv   r   Úcb_logits_processorsr   Úinput_outputsr   r   r¥   r   r<  r    r!   r"   r#   ru   r$   r%   r&   Úutilsr'   r(   r)   r*   r+   r­  r-   rJ   rr   r³  r+  r6   r<   r8   Ú<module>rs     sÑ   ðó Û 	Û Û Ý ß /ß 2Ý Ý Ý ã Ý Ý Ý 6å 3ß XÝ OÝ 9Ý $ß SÑ SÝ 0Ý &Ý Gß LÝ 1ß KÓ Kß BÑ Bß pÕ pðô.˜2Ÿ9™9ô ÷,2ñ ,2ñ` ƒ÷O	5ð O	5ó ðO	5ñf ƒ÷_5ð _5ó ð_5÷Dk!ò k!r<   