ó
    pyüi@  ã                   óx   • S 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JrJrJr  SS
KJr   " S S5      rg)u©  Centralized offloading logic for continuous batching.

Handles two offloading strategies when the GPU KV cache is full:
  1. CPU offloading: copy the KV cache to a pre-allocated pinned CPU buffer, preserving exact request state.
  2. Soft reset: discard the KV cache and re-prefill from scratch when the request is re-scheduled. This incurs no data
    transfer overhead, but we need to re-run prefill over all intial + generated tokens (so more compute overhead).

The CPU swap pool is a static set of pinned tensors allocated once at init (like vLLM/SGLang). Blocks are tracked
with a simple free set â€” no dynamic allocation or deallocation of tensors ever happens at runtime.
é    )Údeque)ÚnullcontextNé   )Úis_psutil_availableé   )ÚPagedAttentionCache)ÚFutureRequestStateÚRequestStateÚRequestStatusÚlogger)Ú	Schedulerc                   ó  • \ rS rSrSrS\S\S\S-  S\S\R                  R                  S-  S	S4S
 jrS\S-  S\S	\4S jrS rSS jrS\\   S	S4S jrS\S	S4S jrSS jrSS jrS\S\S	\4S jrS\S	\\\   \\   4   4S jrSrg)ÚOffloadingManageré$   u#  Manages request offloading and restoration for continuous batching.

Owns a static CPU swap pool (pre-allocated pinned tensors mirroring the GPU cache layout), performs GPUâ†”CPU block
copies, decides between CPU offloading and soft reset, and ensures cleanup on cancellation/failure/reset.
ÚcacheÚ	schedulerÚcpu_offload_space_gibNÚsafety_thresholdÚcompute_streamÚreturnc           	      óö  • Xl         X l        XPl        / U l        / U l        / U l        / U l        [        5       U l        0 U l	        0 U l
        U R                  X45      U l        US L=(       a    US:„  nU R                  S:X  a#  U(       a  [        R                  " SUS S35        g U R                  UR                  UR                   UR"                  4nUR$                   Hs  nU R                  R'                  [(        R*                  " XqR,                  SS95        U R                  R'                  [(        R*                  " XqR,                  SS95        Mu     SUR                  UR                   UR"                  4n	[/        UR$                  UR0                  5       HU  u  p«U R
                  R'                  U
R2                  " U	6 5        U R                  R'                  UR2                  " U	6 5        MW     [        [5        U R                  5      5      U l        [(        R*                  " U R                  [(        R6                  SS9U l        [(        R*                  " U R                  [(        R6                  UR:                  S9U l        U R                  S   nS	UR?                  5       -  URA                  5       -  [C        UR$                  5      -  n[        RD                  " S
U R                   SUS-  S S35        g )Nr   úcpu_offload_space=ú.1fz8 GiB is too small for even one block. No CPU offloading.T)ÚdtypeÚ
pin_memoryéÿÿÿÿ)r   Údeviceé   zCPU swap pool initialized: z	 blocks (é   @ú.2fz GiB pinned))#r   r   Ú_compute_streamÚ_cpu_key_cacheÚ_cpu_value_cacheÚ_gpu_key_viewsÚ_gpu_value_viewsr   Ú_free_cpu_blocksÚ_request_id_to_cpu_blocksÚ!_request_id_to_group_block_countsÚ_compute_num_cpu_blocksÚ_num_cpu_blocksr   ÚwarningÚ
block_sizeÚnum_key_value_headsÚhead_dimÚ	key_cacheÚappendÚtorchÚemptyr   ÚzipÚvalue_cacheÚviewÚrangeÚint32Ú_cpu_ids_scratchr   Ú_gpu_ids_scratchÚnumelÚelement_sizeÚlenÚinfo)Úselfr   r   r   r   r   Úoffloading_enabledÚcpu_cache_shapeÚ_Úblock_shapeÚk_cacheÚv_cacheÚcache_tensorÚsize_in_bytess                 Ú{/home/mande/repo/quber/.venv/lib/python3.13/site-packages/transformers/generation/continuous_batching/offloading_manager.pyÚ__init__ÚOffloadingManager.__init__+   s„  € ð Œ
Ø"Œà-Ôð 35ˆÔØ46ˆÔØ24ˆÔØ46ˆÔÜ,1«GˆÔØ?AˆÔ&ØGIˆÔ.ð  $×;Ñ;Ð<QÓdˆÔØ2¸$Ð>×\ÐCXÐ[\ÑC\ÐØ×Ñ 1Ó$Þ!Ü—’Ø(Ð)>¸sÐ(Cð D)ð )ôð ð  ×/Ñ/°×1AÑ1AÀ5×C\ÑC\Ð^c×^lÑ^lÐmˆØ—”ˆAØ×Ñ×&Ñ&¤u§{¢{°?Ï+É+ÐbfÑ'gÔhØ×!Ñ!×(Ñ(¬¯ª°_ÏKÉKÐdhÑ)iÖjñ !ð
 ˜5×+Ñ+¨U×-FÑ-FÈÏÉÐWˆÜ # E§O¡O°U×5FÑ5FÖ GÑˆGØ×Ñ×&Ñ& w§|¢|°[Ð'AÔBØ×!Ñ!×(Ñ(¨¯ª°{Ð)CÖDñ !Hô
 !&¤e¨D×,@Ñ,@Ó&AÓ BˆÔô !&§¢¨D×,@Ñ,@ÌÏÉÐ`dÑ eˆÔÜ %§¢¨D×,@Ñ,@ÌÏÉÐ\a×\hÑ\hÑ iˆÔð ×*Ñ*¨1Ñ-ˆØ˜L×.Ñ.Ó0Ñ0°<×3LÑ3LÓ3NÑNÔQTÐUZ×UdÑUdÓQeÑeˆÜ�ŠØ)¨$×*>Ñ*>Ð)?¸yÈÐZaÑIbÐcfÐHgÐgsÐtõ	
ó    c                 óè  • Ub  [        US-  5      OSn[        5       (       a,  SSKnUR                  5       R                  n[        XR-  5      nOSnUb:  Ub7  X6:”  a1  US-  n[
        R                  " SUS SUS SWS-  S S	US S
3	5        UnOIUb  [
        R                  " S5        O/Ub!  Un[
        R                  " SUS-  S S
35        O[        S5      eS[        U R                  R                  5      -  U R                  R                  -  U R                  R                  -  U R                  R                  -  U R                  R                  R                  -  nUS:X  a  [!        S5      eX8-  $ )z?Returns the number of blocks that can fit in the CPU swap pool.Nr   r   r   r   z GiB exceeds z.0%z of total RAM (z GiB). Clamping to z GiB.u{   psutil is not available â€” cpu_offload_space_safety_threshold cannot be enforced. Install psutil to enable the safety cap.z1Auto-sizing CPU swap pool from safety threshold: r    ztcpu_offload_space=None requires psutil to auto-size the CPU swap pool. Install psutil or pass an explicit GiB value.r   z9The number of bytes per block is 0. This is not possible.)Úintr   ÚpsutilÚvirtual_memoryÚ	availabler   r+   ÚImportErrorr<   r   r/   r,   r-   r.   r   ÚitemsizeÚ
ValueError)	r>   r   r   Úoffload_bytesrM   Ú	total_ramÚ	max_bytesÚclamped_gibÚbytes_per_blocks	            rG   r)   Ú)OffloadingManager._compute_num_cpu_blocksf   s­  € ð CXÑBcœÐ1°WÑ=Ô>Ðimˆô × Ñ Ûà×-Ñ-Ó/×9Ñ9ˆIÜ˜IÑ8Ó9‰IàˆIð Ñ$¨Ñ)>ØÓ(Ø'¨7Ñ3�Ü—’Ø(Ð)>¸sÐ(CÀ=ÐQaÐbeÐPfð gØ! WÑ-¨cÐ2Ð2EÀkÐRUÐEVÐV[ð]ôð !*�øàÑ&Ü�NŠNð;õð
 Ñ"Ø%ˆMÜ�NŠNÐNÈyÐ\cÑOdÐehÐNiÐinÐoÕpô ð&óð ð Ü�$—*‘*×&Ñ&Ó'ñ(à�j‰j×#Ñ#ñ$ð �j‰j×,Ñ,ñ-ð �j‰j×!Ñ!ñ	"ð
 �j‰j×Ñ×'Ñ'ñ(ð 	ð ˜aÓÜÐXÓYÐYØÑ/Ð/rJ   c                 ó‚   • U R                   b)  [        R                  R                  U R                   5      $ [	        5       $ )zdReturns a context manager that runs enclosed ops on the compute stream, or a no-op when none is set.)r!   r1   ÚcudaÚstreamr   ©r>   s    rG   Ú_stream_ctxÚOffloadingManager._stream_ctx›   s1   € à:>×:NÑ:NÑ:ZŒu�z‰z× Ñ  ×!5Ñ!5Ó6ÐmÔ`kÓ`mÐmrJ   c           
      óÎ  • U R                   nUR                  5       u  p#[        R                  " SU S[	        UR
                  5       S[	        UR                  5       S35        U R                  X#5      nU(       a�  SUl        UR                  [        R                  :X  a  UR                  SS Ul        [        R                  Ul	        Un[        R                  " SU S[	        U R                   5       S	35        O?UR#                  5       n[        R$                  Ul	        [        R                  " S
U S35        UR'                  U5        UR)                  U5        SUl        g)z�Offload one active request to make room in the GPU cache. Tries CPU offloading first; if the pool is full,
falls back to the legacy soft reset.zOffloading request ú with z initial tokens and ú generated tokens.r   NzOffloaded request z	 to CPU: z free blocks remaining.zSoft reset request Ú.T)r   Úpop_request_to_evictr   r=   r<   Úinitial_tokensÚgenerated_tokensÚ_offload_to_cpuÚallocated_blocksÚ_statusr   ÚDECODINGÚtokens_to_processÚremaining_prefill_tokensÚPENDINGÚdebugr&   Ú!create_equivalent_initial_requestÚFINISHEDÚfinish_requestÚadd_waiting_requestÚblock_new_requests)r>   r   Ú
request_idÚstateÚoffloaded_to_cpuÚ	new_states         rG   Úoffload_one_requestÚ%OffloadingManager.offload_one_requestŸ   s6  € ð —N‘Nˆ	Ø%×:Ñ:Ó<Ñˆ
Ü�ŠØ! * ¨V´C¸×8LÑ8LÓ4MÐ3NÐNbÜ�5×)Ñ)Ó*Ð+Ð+=ð?ô	
ð  ×/Ñ/°
ÓBÐÞà%&ˆEÔ"ð �}‰}¤× 6Ñ 6Ó6Ø16×1HÑ1HÉÐ1K�Ô.ô *×1Ñ1ˆEŒMØˆIÜ�LŠLÐ-¨j¨\¸Ä3Àt×G\ÑG\ÓC]ÐB^Ð^uÐvÕwà×?Ñ?ÓAˆIÜ)×2Ñ2ˆEŒMÜ�LŠLÐ.¨z¨l¸!Ð<Ô=à× Ñ  Ô,Ø×%Ñ% iÔ0Ø'+ˆ	Õ$rJ   Úrequests_in_batchc                 óö  • U R                   n/ n/ nU GH€  nUR                  nUR                  (       d  M#  U R                  R	                  UR
                  5      nU R                  R	                  UR
                  5      nUR                  U5        Sn	[        U5       HW  u  p«UR                  U
   R                  R                  UR
                  / 5      nUR                  USU 5        [        X›5      n	MY     SUl        X–l        UR                  (       a,  U=R                  UR                   UR"                  -  -  sl        [$        R&                  " SUR
                   S[)        UR*                  5       S[)        UR,                  5       S35        GMƒ     U(       d  gU R.                  S[)        U5       nU R0                  S[)        U5       nUR3                  [4        R6                  " U[4        R8                  S95        U R;                  5          UR3                  [4        R6                  " U[4        R8                  S95        [=        U R>                  U R@                  5       H  u  nnUU   R3                  Xý   5        M     [=        U RB                  U RD                  5       H  u  nnUU   R3                  UU   5        M     SSS5        U RF                  R                  U5        g! , (       d  f       N*= f)	z¸Restore KV caches from CPU for any CPU-offloaded requests in the scheduled batch. Indices are accumulated
per group across all requests, then copied in one batched operation per layer.r   NFzRestored CPU-offloaded request r`   z prefill tokens and ra   ©r   )$r   rt   Úis_cpu_offloadedr'   Úpoprs   r(   ÚextendÚ	enumerateÚgroup_cache_managersÚblock_tableÚgetÚmaxrg   Úallow_block_sharingÚcomplete_blocksÚposition_offsetr,   r   rm   r<   rd   re   r8   r9   Úcopy_r1   Ú	as_tensorr7   r]   r3   r"   r$   r#   r%   r&   )r>   ry   r   Úall_cpu_indicesÚall_gpu_indicesÚfuture_statert   Úcpu_indicesÚgroup_countsÚmax_allocated_blocksÚ	group_idxÚnÚ
gpu_blocksÚcpu_idsÚgpu_idsÚcpu_kÚgpu_kÚcpu_vÚgpu_vs                      rG   Úrestore_scheduled_requestsÚ,OffloadingManager.restore_scheduled_requestsÀ   sy  € ð —
‘
ˆØ%'ˆØ%'ˆä-ˆLà ×&Ñ&ˆEØ×)×)Ùð ×8Ñ8×<Ñ<¸U×=MÑ=MÓNˆKØ×AÑA×EÑEÀe×FVÑFVÓWˆLØ×"Ñ" ;Ô/ð $%Ð Ü )¨,Ö 7‘�	Ø"×7Ñ7¸	ÑB×NÑN×RÑRÐSX×ScÑScÐegÓh�
Ø×&Ñ& z°"°1 ~Ô6Ü'*Ð+?Ó'CÒ$ñ !8ð
 &+ˆEÔ"Ø%9Ô"à×(×(Ø×,Ò,°×0EÑ0EÈ×IYÑIYÑ0YÑYÕ,Ü�LŠLØ1°%×2BÑ2BÐ1CÀ6Ì#Èe×NbÑNbÓJcÐIdð eÜ˜5×1Ñ1Ó2Ð3Ð3EðG÷ñ/ .ö: Øð ×'Ñ'Ð(>¬#¨oÓ*>Ð?ˆØ×'Ñ'Ð(>¬#¨oÓ*>Ð?ˆØ�‰”e—o’o o¼U¿[¹[ÑIÔJØ×ÑÕØ�M‰Mœ%Ÿ/š/¨/ÄÇÁÑMÔNÜ # D×$7Ñ$7¸×9LÑ9LÖ M‘��uØ�g‘×$Ñ$ U¡^Ö4ñ !Nä # D×$9Ñ$9¸4×;PÑ;PÖ Q‘��uØ�g‘×$Ñ$ U¨7¡^Ö4ñ !R÷	  ð 	×Ñ×$Ñ$ _Õ5÷  Õús   ÈB4K*Ë*
K8rt   c                 ól   • UR                   (       a#  U R                  UR                  5        SUl         gg)z=Free CPU blocks for a single request (e.g., on cancellation).FN)r|   Ú_return_cpu_blocksrs   ©r>   rt   s     rG   Úfree_request_cpu_cacheÚ(OffloadingManager.free_request_cpu_cacheó   s,   € à×!×!Ø×#Ñ# E×$4Ñ$4Ô5Ø%*ˆEÕ"ð "rJ   c                 ó|   • U R                   R                  R                  5        H  nU R                  U5        M     g)zPFree all CPU-offloaded caches in the waiting queue (e.g., on fail_all or reset).N)r   Úwaiting_requestsÚvaluesr�   rœ   s     rG   Úfree_all_waiting_cpu_cachesÚ-OffloadingManager.free_all_waiting_cpu_cachesù   s-   € à—^‘^×4Ñ4×;Ñ;Ö=ˆEØ×'Ñ'¨Ö.ò >rJ   c                 óÒ   • U R                  5         U R                  R                  5         U R                  R                  5         [	        [        U R                  5      5      U l        g)z8Reset CPU offloading state for a new generation session.N)r¢   r'   Úclearr(   r   r6   r*   r&   r\   s    rG   ÚresetÚOffloadingManager.resetþ   sJ   € à×(Ñ(Ô*Ø×&Ñ&×,Ñ,Ô.Ø×.Ñ.×4Ñ4Ô6Ü %¤e¨D×,@Ñ,@Ó&AÓ BˆÕrJ   rs   c                 ó.  • / n/ nU R                   R                   HJ  nUR                  R                  U/ 5      nUR	                  U5        UR                  [        U5      5        ML     [        U5      nUS:X  d  [        U R                  5      U:  a  g[        U5       Vs/ s H  o€R                  R                  5       PM     n	nU R                  SU n
U R                  SU nU
R                  [        R                  " U	[        R                  S95        U R!                  5          UR                  [        R                  " U[        R                  S95        [#        U R$                  U R&                  5       H  u  pÍXÊ   R                  XÛ   5        M     [#        U R(                  U R*                  5       H  u  pïXê   R                  Xû   5        M     SSS5        X�R,                  U'   X@R.                  U'   SUl        gs  snf ! , (       d  f       N7= f)zzCopy a request's KV cache blocks from GPU to the static CPU swap pool. Returns True on success, False if
the pool is full.r   FNr{   T)r   r€   r�   r‚   r~   r0   r<   r&   r6   Úpopleftr8   r9   r‡   r1   rˆ   r7   r]   r3   r"   r$   r#   r%   r'   r(   r|   )r>   rs   rt   Úgpu_indicesÚgroup_block_countsÚcmÚblocksÚtotal_gpu_blocksrA   rŒ   r’   r“   Úcpu_key_cacheÚgpu_key_viewÚcpu_value_cacheÚgpu_value_views                   rG   rf   Ú!OffloadingManager._offload_to_cpu  sÇ  € ð
 ˆØÐØ—*‘*×1Ô1ˆBØ—^‘^×'Ñ'¨
°BÓ7ˆFØ×Ñ˜vÔ&Ø×%Ñ%¤c¨&£kÖ2ñ 2ô ˜{Ó+ÐØ˜qÓ ¤C¨×(=Ñ(=Ó$>ÐAQÓ$QØô AFÐFVÔ@WÓXÒ@W¸1×,Ñ,×4Ñ4Ö6Ñ@WˆÐXð ×'Ñ'Ð(9Ð)9Ð:ˆØ×'Ñ'Ð(9Ð)9Ð:ˆØ�‰”e—o’o k¼¿¹ÑEÔFØ×ÑÕØ�M‰Mœ%Ÿ/š/¨+¼U¿[¹[ÑIÔJä/2°4×3FÑ3FÈ×H[ÑH[Ö/\Ñ+�ØÑ&×,Ñ,¨\Ñ-BÖCñ 0]ô 47°t×7LÑ7LÈd×NcÑNcÖ3dÑ/�ØÑ(×.Ñ.¨~Ñ/FÖGñ 4e÷  ð 6A×&Ñ& zÑ2Ø=O×.Ñ.¨zÑ:Ø!%ˆÔØùò+ Y÷  Õús   Â!#HÄ&B/HÈ
Hc                 óª   • U R                   R                  U5      nU R                  R                  U5      nU R                  R	                  U5        X#4$ )z<Return CPU blocks to the free pool without copying anything.)r'   r}   r(   r&   r~   )r>   rs   r’   r�   s       rG   r›   Ú$OffloadingManager._return_cpu_blocks.  sK   € à×0Ñ0×4Ñ4°ZÓ@ˆØ×=Ñ=×AÑAÀ*ÓMˆØ×Ñ×$Ñ$ WÔ-ØÐ$Ð$rJ   )r!   r8   r"   r#   r&   r9   r$   r%   r*   r'   r(   r   r   )r   N)Ú__name__Ú
__module__Ú__qualname__Ú__firstlineno__Ú__doc__r   r   Úfloatr1   rZ   ÚStreamrH   rL   r)   r]   rw   Úlistr	   r˜   r
   r�   r¢   r¦   ÚstrÚboolrf   Útupler›   Ú__static_attributes__© rJ   rG   r   r   $   s   † ñð9
à"ð9
ð ð9
ð  % t™|ð	9
ð
  ð9
ð Ÿ
™
×)Ñ)¨DÑ0ð9
ð 
ô9
ðv30¸UÀT¹\ð 30Ð]bð 30Ðgjô 30òjnô,ðB16¸DÐASÑ<Tð 16ÐY]ô 16ðf+¨Lð +¸Tô +ô/ô
Cð'¨#ð '°lð 'Àtô 'ðR%¨Sð %°U¸4À¹9ÀdÈ3ÁiÐ;OÑ5P÷ %rJ   r   )rº   Úcollectionsr   Ú
contextlibr   r1   Úutilsr   r   r   Úrequestsr	   r
   r   r   r   r   r   rÂ   rJ   rG   Ú<module>rÇ      s0   ðñ	õ Ý "ã å (Ý &ß MÓ MÝ  ÷O%ò O%rJ   