Ë
    g^(h½  ã                   óÊ   — d dl Z d dlZd dlZd dlZd dlZd dlmZ d dlm	Z	 ddl
mZmZmZmZ g d¢Z e j                   e«      Z G d„ de«      Z G d	„ d
e«      Z G d„ de«      Zy)é    N)ÚEmpty)ÚAnyé   )ÚRequestQueueÚTimerClientÚTimerRequestÚTimerServer)ÚLocalTimerClientÚMultiprocessingRequestQueueÚLocalTimerServerc                   ó.   ‡ — e Zd ZdZˆ fd„Zd„ Zd„ Zˆ xZS )r
   aF  
    Client side of ``LocalTimerServer``. This client is meant to be used
    on the same host that the ``LocalTimerServer`` is running on and uses
    pid to uniquely identify a worker. This is particularly useful in situations
    where one spawns a subprocess (trainer) per GPU on a host with multiple
    GPU devices.
    c                 ó0   •— t         ‰| �  «        || _        y ©N©ÚsuperÚ__init__Ú	_mp_queue©ÚselfÚmp_queueÚ	__class__s     €úi/var/www/skyplay_api_hub/venv/lib/python3.12/site-packages/torch/distributed/elastic/timer/local_timer.pyr   zLocalTimerClient.__init__    ó   ø€ Ü‰ÑÔØ!ˆ�ó    c                 ó|   — t        j                  «       }t        |||«      }| j                  j	                  |«       y r   ©ÚosÚgetpidr   r   Úput)r   Úscope_idÚexpiration_timeÚpidÚacquire_requests        r   ÚacquirezLocalTimerClient.acquire$   s-   € Ü�i‰i‹kˆÜ& s¨H°oÓFˆØ�‰×Ñ˜?Õ+r   c                 ó|   — t        j                  «       }t        ||d«      }| j                  j	                  |«       y )Néÿÿÿÿr   )r   r    r"   Úrelease_requests       r   ÚreleasezLocalTimerClient.release)   s-   € Ü�i‰i‹kˆÜ& s¨H°bÓ9ˆØ�‰×Ñ˜?Õ+r   )Ú__name__Ú
__module__Ú__qualname__Ú__doc__r   r$   r(   Ú__classcell__©r   s   @r   r
   r
      s   ø„ ñô"ò,ö
,r   r
   c                   ó^   ‡ — e Zd ZdZdej
                  fˆ fd„Zdefd„Zde	de
e   fd„Zˆ xZS )r   zG
    A ``RequestQueue`` backed by python ``multiprocessing.Queue``
    r   c                 ó0   •— t         ‰| �  «        || _        y r   r   r   s     €r   r   z$MultiprocessingRequestQueue.__init__4   r   r   Úreturnc                 ó6   — | j                   j                  «       S r   )r   Úqsize)r   s    r   Úsizez MultiprocessingRequestQueue.size8   s   € Ø�~‰~×#Ñ#Ó%Ð%r   Útimeoutc                 ó  — g }|}t        d|«      D ]g  }t        j                  «       }	 | j                  j                  d|¬«      }|j                  |«       |t        j                  «       |z
  z
  }|dk  sŒf |S  |S # t        $ r Y  |S w xY w)Nr   T)Úblockr5   )ÚrangeÚtimer   Úgetr   Úappend)r   r4   r5   ÚrequestsÚwaitÚ_ÚstartÚrs           r   r:   zMultiprocessingRequestQueue.get;   sž   € ØˆØˆÜ�q˜$“ò 	ˆAÜ—I‘I“KˆEðØ—N‘N×&Ñ&¨T¸4Ð&Ó@�ð �O‰O˜AÔØœ4Ÿ9™9›;¨Ñ.Ñ/ˆDØ�q‹yØàˆð	ð ˆøô ò Ùð ˆðús   ©A=Á=	BÂ
B)r)   r*   r+   r,   ÚmpÚQueuer   Úintr4   ÚfloatÚlistr   r:   r-   r.   s   @r   r   r   /   s<   ø„ ñð" §¡õ "ð&�có &ð ð ¨4°Ñ+=÷ r   r   c                   ó¤   ‡ — e Zd ZdZ	 ddej
                  dedefˆ fd„Zde	e
   ddfd	„Zd
ee   ddfd„Zdedeee	e
   f   fd„Zdedefd„Zˆ xZS )r   aK  
    Server that works with ``LocalTimerClient``. Clients are expected to be
    subprocesses to the parent process that is running this server. Each host
    in the job is expected to start its own timer server locally and each
    server instance manages timers for local workers (running on processes
    on the same host).
    r   Úmax_intervalÚdaemonc                 óH   •— t         ‰| �  t        |«      ||«       i | _        y r   )r   r   r   Ú_timers)r   r   rG   rH   r   s       €r   r   zLocalTimerServer.__init__W   s#   ø€ ô 	‰ÑÔ4°XÓ>ÀÈfÔUØ<>ˆ�r   Útimer_requestsr1   Nc                 óÄ   — |D ][  }|j                   }|j                  }|j                  }|dk  r| j                  j	                  ||fd «       ŒK|| j                  ||f<   Œ] y )Nr   )Ú	worker_idr    r!   rJ   Úpop)r   rK   Úrequestr"   r    r!   s         r   Úregister_timersz LocalTimerServer.register_timers]   sf   € Ø%ò 		8ˆGØ×#Ñ#ˆCØ×'Ñ'ˆHØ%×5Ñ5ˆOð  Ò"Ø—‘× Ñ  # x °$Õ7à07�—‘˜c 8˜_Ò-ñ		8r   Ú
worker_idsc                 óž   — t        | j                  j                  «       «      D ]'  \  }}||v sŒ| j                  j                  ||f«       Œ) y r   )rE   rJ   ÚkeysrN   )r   rQ   r"   r    s       r   Úclear_timerszLocalTimerServer.clear_timersi   sE   € Ü! $§,¡,×"3Ñ"3Ó"5Ó6ò 	2‰MˆC�Ø�jÒ Ø—‘× Ñ  # x Õ1ñ	2r   Údeadlinec                 óÂ   — i }| j                   j                  «       D ]?  }|j                  |k  sŒ|j                  |j                  g «      }|j                  |«       ŒA |S r   )rJ   Úvaluesr!   Ú
setdefaultrM   r;   )r   rU   Úexpired_timersrO   Úexpired_scopess        r   Úget_expired_timersz#LocalTimerServer.get_expired_timersn   s_   € à8:ˆØ—|‘|×*Ñ*Ó,ò 	/ˆGØ×&Ñ&¨(Ó2Ø!/×!:Ñ!:¸7×;LÑ;LÈbÓ!Q�Ø×%Ñ% gÕ.ð	/ð Ðr   rM   c                 óØ   — 	 t        j                  |t        j                  «       y# t        $ r t
        j                  d|«       Y yt        $ r t
        j                  d|«       Y yw xY w)NTz,Process with pid=%s does not exist. SkippingzError terminating pid=%sF)	r   ÚkillÚsignalÚSIGKILLÚProcessLookupErrorÚloggerÚinfoÚ	ExceptionÚ	exception)r   rM   s     r   Ú_reap_workerzLocalTimerServer._reap_workerw   s\   € ð	DÜ�G‰G�IœvŸ~™~Ô.ØøÜ!ò 	Ü�K‰KÐFÈ	ÔRÙÜò 	DÜ×ÑÐ7¸ÕCØð	Dús   ‚$' §A)ÁA)Á(A))é<   T)r)   r*   r+   r,   rA   rB   rD   Úboolr   rE   r   rP   ÚsetrC   rT   Údictr   r[   re   r-   r.   s   @r   r   r   N   s”   ø„ ñð LPñ?ØŸ™ð?Ø05ð?ØDHõ?ð
8¨d°<Ñ.@ð 
8ÀTó 
8ð2 s¨3¡xð 2°Dó 2ð
¨5ð °T¸#¸tÀLÑ?QÐ:QÑ5Ró ð	 cð 	¨d÷ 	r   r   )ÚloggingÚmultiprocessingrA   r   r^   r9   Úqueuer   Útypingr   Úapir   r   r   r	   Ú__all__Ú	getLoggerr)   ra   r
   r   r   © r   r   ú<module>rr      s`   ðó Û Û 	Û Û Ý Ý ç EÓ Eò R€à	ˆ×	Ñ	˜8Ó	$€ô,�{ô ,ô0 ,ô ô>2�{õ 2r   