Ë
    g^(hm  ã                   ón  — d dl mZ d dlmZ d dlmZ d dlmZmZ d dl	Z	e	j                  j                  ZdZdZdZd	Zd
Zg d¢Zedefd„«       Zdededefd„Z	 ddedededededee   fd„Zdedededeegef   dedeee      fd„Zdededefd„Z	 	 	 	 ddedededee   deeegef      deddfd„Zy)é    )ÚIterable)Úcontextmanager)Ú	timedelta)ÚCallableÚOptionalNz/num_membersz/last_memberz/TRACEz/TRACING_GATEé   )Ústore_timeoutÚget_allÚsynchronizeÚbarrierÚtimeoutc              #   ó„   K  — | j                   }| j                  t        |¬«      «       d–— | j                  |«       y­w)zÃ
    This sets the timeout and then restores the old timeout when the context
    manager exits.

    Args:
        store: the store to set the timeout on
        timeout: the timeout to set
    ©ÚsecondsN)r   Úset_timeoutr   )Ústorer   Úold_timeouts      úc/var/www/skyplay_api_hub/venv/lib/python3.12/site-packages/torch/distributed/elastic/utils/store.pyr	   r	      s5   è ø€ ð —-‘-€KØ	×Ñ”i¨Ô0Ô1Û	Ø	×Ñ�kÕ"ùs   ‚>A ÚrankÚprefixÚ
world_sizec                 ó¸   — | j                  t        |«      D �cg c]  }|› |› �‘Œ
 c}«      }t        | ||› d�¬«      }|dk(  r| j                  |g«       |S c c}w )ad  
    Given a store and a prefix, the method goes through the array of keys
    of the following format: ``{prefix}{idx}``, where idx is in a range
    from 0 to size, and tries to retrieve the data.

    The Rank0 process waits at the end to make sure all other processes
    finished the procedure before exiting.

    Usage

    ::

     values = get_all(store, "torchelastic/data", 3)
     value1 = values[0]  # retrieves the data for key torchelastic/data0
     value2 = values[1]  # retrieves the data for key torchelastic/data1
     value3 = values[2]  # retrieves the data for key torchelastic/data2

    z	/finished©r   r   Ú
key_prefixr   )Ú	multi_getÚrangeÚ_barrier_nonblockingÚwait)r   r   r   r   ÚidxÚdata_arrÚbarrier_keys          r   r
   r
   /   si   € ð& �‰¼EÀ*Ó<MÖN°S 6 (¨3¨%Ò 0ÒNÓO€Hä&ØØØ�X˜YÐ'ô€Kð
 ˆq‚yð 	�
‰
�K�=Ô!à€Oùò  Os   ™AÚdatar   Úreturnc                 ó’   — t        | |«      5  | j                  |› |› �|«       t        | |||«      }|cddd«       S # 1 sw Y   yxY w)aT  
    Synchronizes ``world_size`` agents between each other using the underlying c10d store.
    The ``data`` will be available on each of the agents.

    Note: The data on the path is not deleted, as a result there can be stale data if
        you use the same key_prefix twice.

    Time complexity: O(N) per worker, O(N^2) globally.
    N)r	   Úsetr
   )r   r"   r   r   r   r   Ú
agent_datas          r   r   r   R   sM   € ô" 
�u˜gÓ	&ñ Ø�	‰	�Z�L  Ð'¨Ô.Ü˜U D¨*°jÓAˆ
Ø÷÷ ò ús	   �&=½AÚrank_decoderÚtrace_timeoutc                 óÂ   ‡ ‡‡‡‡— ‰ j                  ‰› |› t        › �d«       ˆˆˆ ˆˆfd„}ˆˆˆ fd„}|dk(  r# |«       }‰ j                  ‰› t        › �d«       |S  |«       S )Nú<val_ignored>c                  óZ  •— t        «       } d}t        d‰«      D ]c  }|t        k\  r | S 	 |dk(  r(‰j                  ‰› |› t        › �gt        ‰¬«      «       n'‰j                  ‰› |› t        › �gt        d¬«      «       Œe | S # t        $ r |dz  }| j                   ‰|«      «       Y Œ�w xY w)Nr   é   r   )Úmilliseconds)r%   r   Ú_MAX_TRACE_MISSING_RANKSr   Ú_TRACEr   ÚDistStoreErrorÚadd)Úmissing_rank_infoÚranks_missingÚir   r'   r   r(   r   s      €€€€€r   Ú_find_missing_ranksz9_try_detecting_missing_ranks.<locals>._find_missing_rankss   sË   ø€ Ü›EÐØˆÜ�q˜*Ó%ò 	7ˆAð Ô 8Ò8Øð !Ð ð
7Ø  AÒ%Ø—J‘JØ&˜<¨ s¬6¨(Ð3Ð4´iÈÔ6Võð
 —J‘J : ,¨q¨c´&°Ð :Ð;¼YÐTUÔ=VÔWøð	7ð  !Ð øô "ò 7Ø Ñ"�Ø!×%Ñ%¡l°1£oÖ6ð7ús   ªABÂ%B*Â)B*c                  ór   •— 	 ‰j                  ‰ › t        › �g«       d ‰d«      › d�gS # t        $ r Y y w xY w)Nz[<check rank 0 (r   z) for missing rank info>])r   Ú_TRACING_GATEr0   )r   r'   r   s   €€€r   Ú_checkinz._try_detecting_missing_ranks.<locals>._checkinˆ   sJ   ø€ ð	Ø�J‰J˜:˜,¤} oÐ6Ð7Ô8Ø&¡|°A£Ð&7Ð7PÐQÐRÐRøÜò 	áð	ús   ƒ&* ª	6µ6r   )r%   r/   r7   )	r   r   r   r   r'   r(   r5   r8   r2   s	   ``` ``   r   Ú_try_detecting_missing_ranksr9   i   sf   ü€ ð 
‡I�I��˜T˜F¤6 (Ð+¨_Ô=÷!ð !ö*ð ˆq‚yÙ/Ó1ÐØ�	‰	�Z�L¤ Ð0°/ÔBØ Ð á‹zÐó    c                 ó|   — |t         z   }|t        z   }| j                  |d«      }||k(  r| j                  |d«       |S )zq
    Does all the non-blocking operations for a barrier and returns the final key
    that can be waited on.
    r,   r*   )Ú_NUM_MEMBERSÚ_LAST_MEMBER_CHECKINr1   r%   )r   r   r   Únum_members_keyÚlast_member_keyr   s         r   r   r   ˜   sE   € ð
 !¤<Ñ/€OØ Ô#7Ñ7€Oà
�)‰)�O QÓ
'€CØ
ˆjÒØ�	‰	�/ ?Ô3àÐr:   Úbarrier_timeoutÚrank_tracing_decoderc                 ó`  — |€	|�J d«       ‚t        | |«      5  t        | ||¬«      }	 | j                  |g«       	 ddd«       y# t        $ rT}|€|‚t	        | ||||xs d„ |«      }	|	�2t        dj                  |||ddj                  |	«      › d�|«      «      d‚|‚d}~ww xY w# 1 sw Y   yxY w)	as  
    A global lock between agents. This will pause all workers until at least
    ``world_size`` workers respond.

    This uses a fast incrementing index to assign waiting ranks and a success
    flag set by the last worker.

    Time complexity: O(1) per worker, O(N) globally.

    Optionally, passing rank will enable tracing of missing ranks on timeouts.
    `rank_tracing_decoder` lambda arg can be used to convert rank data
    into a more meaninful information at an app level (e.g. hostname).

    Note: Since the data is not removed from the store, the barrier can be used
        once per unique ``key_prefix``.
    Nz!Tracing requires rank informationr   c                 ó   — t        | «      S )N)Ústr)Úxs    r   ú<lambda>zbarrier.<locals>.<lambda>Ó   s
   € ´s¸1³v€ r:   ziTimed out waiting on barrier on rank {}, for key prefix: {} (world_size={}, missing_ranks={}, timeout={})ú[z, ú])r	   r   r   r0   r9   ÚformatÚjoin)
r   r   r   r@   r   rA   r(   r?   ÚeÚmissing_rankss
             r   r   r   §   s÷   € ð4 €|Ø#Ð+ÐPÐ-PÓPÐ+ä	�u˜oÓ	.ñ Ü.Ø J¸:ô
ˆð	Ø�J‰J˜Ð(Õ)÷ð øô ò 	Øˆ|Ø�ä <ØØØØØ(Ò>Ñ-=Ø!ó!�ð !Ð,Ü(ðdßdjÑdjØ Ø&Ø&Ø §	¡	¨-Ó 8Ð9¸Ð;Ø+óeó	ð  ð	 ð �Gûð1	ú÷ð ús)   ˜B$¨AÁ	B!ÁABÂB!Â!B$Â$B-)é,  )rM   NNé
   )Úcollections.abcr   Ú
contextlibr   Údatetimer   Útypingr   r   ÚtorchÚ_CÚ_DistStoreErrorr0   r<   r=   r/   r7   r.   Ú__all__Úfloatr	   ÚintrD   r
   ÚbytesÚlistr   r9   r   r   © r:   r   ú<module>r\      s¡  ðõ %Ý %Ý ß %ã ð —‘×)Ñ)€à€Ø%Ð Ø	€Ø€ØÐ ò A€ð ð# %ò #ó ð#ð  ˜ð   cð  °só  ðR ñà
ðð ðð ð	ð
 ðð ðð 
ˆ%�[óð.,àð,ð ð,ð ð	,ð
 ˜C˜5 #˜:Ñ&ð,ð ð,ð ˆh�s‰mÑó,ð^¨Cð ¸Sð ÀSó ð& !ØØ;?Øñ;àð;ð ð;ð ð	;ð
 �3‰-ð;ð # 8¨S¨E°3¨JÑ#7Ñ8ð;ð ð;ð 
ô;r:   