Ë
    ‡\;jÌ!  ã                   ó~   — d Z ddlZddlmZ ddlmZ g Z G d„ d«      Z G d„ de«      Z G d	„ d
«      Z	 G d„ d«      Z
y)z£
Communicator is used for async distribute training in distribute_transpiler mode.
It's a wrapper of a cpp class Communicator and should be used inside fleet API.
é    N)ÚDistributedMode)Úcorec                   ód   — e Zd Zdd„Z	 dd„Z	 	 	 dd„Zd„ Zd„ Zd„ Zd„ Z	d	„ Z
d
„ Zd„ Zd„ Zdd„Zy)ÚCommunicatorNc                 óæ  — |€|€qi }nn|t         j                  k(  rdj                  |d   «      |d<   t        |d   «      |d<   t        |d   «      |d<   t        |d   «      |d<   t        |d   «      |d<   d}|t         j                  k(  rd}nA|t         j                  k(  rd	}n+|t         j
                  k(  rd
}n|t         j                  k(  rd}|| _        || _        d| _	        d| _
        d| _        y)a¼  
        Communicator is used for async distribute training in distribute_transpiler mode.
        It's a wrapper of a cpp class Communicator and should be used inside fleet API.

        Args:
            program(Program): the trainers program after transpile of distribute_transpiler.
            It's used by communicator to extract the information to do communication.

        Returns:
            None

        Examples:
            .. code-block:: python

                >>> import paddle

                >>> prog = paddle.static.Program()
                >>> comm = paddle.distributed.communicator.Communicator(prog)
                >>> comm.start()
                >>> comm.stop()
        NÚ,Úpserver_endpointsÚtrainersÚ
trainer_idÚneed_global_stepÚbarrier_table_idÚSYNCÚASYNCÚ
HALF_ASYNCÚGEO)r   r   ÚjoinÚstrr   r   r   ÚmodeÚenvsÚcommunicator_Ú	send_ctx_Ú	recv_ctx_)Úselfr   Úkwargsr   Úmode_strs        úhG:\00. PROJECTS\API\Inventory\templateJSON\kerjaOCR\Lib\site-packages\paddle/distributed/communicator.pyÚ__init__zCommunicator.__init__)   s	  € ð0 ˆ>Øˆ|Ø‘à”×+Ñ+Ò+Ø,/¯H©HØÐ.Ñ/ó-�Ð(Ñ)ô  # 6¨*Ñ#5Ó6ˆD�ÑÜ!$ V¨LÑ%9Ó!:ˆD�ÑÜ'*¨6Ð2DÑ+EÓ'FˆDÐ#Ñ$Ü'*¨6Ð2DÑ+EÓ'FˆDÐ#Ñ$àˆà”?×'Ñ'Ò'Ø‰HØ”_×*Ñ*Ò*Ø‰HØ”_×/Ñ/Ò/Ø#‰HØ”_×(Ñ(Ò(ØˆHàˆŒ	ØˆŒ	Ø!ˆÔØˆŒØˆ�ó    c           	      óÈ   — |€t         j                  j                  «       }t        j                  | j
                  |||||| j                  «      | _        || _        || _	        y ©N)
ÚpaddleÚstaticÚglobal_scoper   ÚDistCommunicatorr   r   r   r   r   )r   Úsend_ctxÚrecv_ctxÚ	proto_txtÚunit64_hostsÚscopes         r   Úinit_with_ctxzCommunicator.init_with_ctx`   s[   € ð ˆ=Ü—M‘M×.Ñ.Ó0ˆEÜ!×2Ñ2Ø�I‰IØØØØØØ�I‰Ió
ˆÔð "ˆŒØ!ˆ�r   c                 ó>   — | j                   j                  |||«       y r    )r   Ú"create_client_to_client_connection)r   Úpserver_timeout_msÚpserver_connect_timeout_msÚ	max_retrys       r   r,   z/Communicator.create_client_to_client_connectionq   s    € ð 	×Ñ×=Ñ=ØÐ :¸Iõ	
r   c                 ó6   — | j                   j                  «       S r    )r   Úget_client_info©r   s    r   r1   zCommunicator.get_client_info{   s   € Ø×!Ñ!×1Ñ1Ó3Ð3r   c                 ó:   — | j                   j                  |«       y r    )r   Úset_clients)r   Ú	host_lists     r   r4   zCommunicator.set_clients~   s   € Ø×Ñ×&Ñ& yÕ1r   c                 óh   — | j                   €t        d«       y| j                   j                  «        y)a‰  
        Start communicator. Should call before training process.

        Returns:
            None

        Examples:
            .. code-block:: python

                >>> import paddle

                >>> prog = paddle.static.Program()
                >>> comm = paddle.distributed.communicator.Communicator(prog)
                >>> comm.start()
                >>> comm.stop()
        Nz;you must call init_with_ctx first to init comm before start)r   ÚprintÚstartr2   s    r   r8   zCommunicator.start�   s.   € ð" ×ÑÐ%ÜÐOÔPØØ×Ñ× Ñ Õ"r   c                 óh   — | j                   €t        d«       y| j                   j                  «        y)a‡  
        Stop communicator. Should call after training process.

        Returns:
            None

        Examples:
            .. code-block:: python

                >>> import paddle

                >>> prog = paddle.static.Program()
                >>> comm = paddle.distributed.communicator.Communicator(prog)
                >>> comm.start()
                >>> comm.stop()
        Nú:you must call init_with_ctx first to init comm before stop)r   r7   Ústopr2   s    r   r;   zCommunicator.stop—   s.   € ð" ×ÑÐ%ÜÐNÔOØØ×Ñ×ÑÕ!r   c                 óh   — | j                   €t        d«       y| j                   j                  «        y)aZ  
        Get communicator is running or stop.

        Returns:
            bool

        Examples:
            .. code-block:: python

                >>> import paddle

                >>> prog = paddle.static.Program()
                >>> comm = paddle.distributed.communicator.Communicator(prog)
                >>> comm.is_running()
        Nr:   )r   r7   Ú
is_runningr2   s    r   r=   zCommunicator.is_running­   s.   € ð  ×ÑÐ%ÜÐNÔOØØ×Ñ×%Ñ%Õ'r   c                 ó8   — | j                   j                  «        y r    )r   Úrecvr2   s    r   r?   zCommunicator.recvÂ   ó   € Ø×Ñ×ÑÕ!r   c                 ó:   — | j                   j                  |«       y r    )r   Úinit_params©r   Úcontexts     r   rB   zCommunicator.init_paramsÅ   s   € Ø×Ñ×&Ñ& wÕ/r   c                 ó:   — | j                   j                  |«       y r    )r   Ú
pull_denserC   s     r   rF   zCommunicator.pull_denseÈ   s   € Ø×Ñ×%Ñ% gÕ.r   c                 ó@  — |€t         j                  j                  «       }| j                  «       st	        d«      ‚t        |t        «      sJ ‚t        |t        «      sJ ‚|dk(  r| j                  |   j                  «       }| j                  j                  |||«       y )NzTCommunicator should init first. Using fleet.init_worker() before push_sparse_param()éÿÿÿÿ)r!   r"   r#   r=   Ú
ValueErrorÚ
isinstancer   Úintr   Útable_idr   Úpush_sparse_param)r   Úvar_namerL   r)   s       r   rM   zCommunicator.push_sparse_paramË   s‹   € Øˆ=Ü—M‘M×.Ñ.Ó0ˆEØ�‰Ô ÜØfóð ô ˜(¤CÔ(Ð(Ð(Ü˜(¤CÔ(Ð(Ð(Ø�rŠ>Ø—~‘~ hÑ/×8Ñ8Ó:ˆHØ×Ñ×,Ñ,¨X°xÀÕGr   )NNr    )i ¡ i'  é   )rH   N)Ú__name__Ú
__module__Ú__qualname__r   r*   r,   r1   r4   r8   r;   r=   r?   rB   rF   rM   © r   r   r   r   (   sR   „ ó5ðp BFó"ð& "Ø#(Øó	
ò4ò2ò#ò,"ò,(ò*"ò0ò/ôHr   r   c                   ó2   ‡ — e Zd Zdˆ fd„	Zd„ Zd„ Zd„ Zˆ xZS )ÚFLCommunicatorc                 ól   •— d }t         ‰| �  ||«       i }i }d}d| _        | j                  ||||«       y )NÚ ÚWITH_COORDINATOR)Úsuperr   r   r*   )r   Úps_hostsr   r   r%   Ú	dense_mapÚprototxtÚ	__class__s          €r   r   zFLCommunicator.__init__Ú   sA   ø€ ØˆÜ‰Ñ˜˜vÔ&ØˆØˆ	ØˆØ&ˆŒ	Ø×Ñ˜8 Y°¸(ÕCr   c                 óV   — | j                   �| j                   j                  ||«       y y r    )r   Ústart_coordinator)r   Úself_endpointÚtrainer_endpointss      r   r_   z FLCommunicator.start_coordinatorã   s-   € Ø×ÑÐ)Ø×Ñ×0Ñ0ØÐ0õð *r   c                 óh   — | j                   �| j                   j                  |«       y t        d«      ‚)Nzself.communicator_ is null)r   Úsave_fl_strategyrI   )r   Úmps     r   rc   zFLCommunicator.save_fl_strategyé   s.   € Ø×ÑÐ)Ø×Ñ×/Ñ/°Õ3äÐ9Ó:Ð:r   c                 óV   — i }| j                   �| j                   j                  «       }|S r    )r   Úquery_fl_clients_info)r   Úinfo_mps     r   rf   z$FLCommunicator.query_fl_clients_infoï   s,   € ØˆØ×ÑÐ)Ø×(Ñ(×>Ñ>Ó@ˆGØˆr   r    )rP   rQ   rR   r   r_   rc   rf   Ú__classcell__)r]   s   @r   rU   rU   Ù   s   ø„ õDòò;ör   rU   c                   ó$   — e Zd Zd„ Zd„ Zd„ Zd„ Zy)ÚLargeScaleKVc                 ó6   — t        j                  «       | _        y r    )r   rj   Úscale_kvr2   s    r   r   zLargeScaleKV.__init__÷   s   € Ü×)Ñ)Ó+ˆ�r   c                 ó<   — | j                   j                  ||«       y r    )rl   Úsave©r   ÚvarnameÚdirnames      r   rn   zLargeScaleKV.saveú   ó   € Ø�‰×Ñ˜7 GÕ,r   c                 ó<   — | j                   j                  ||«       y r    )rl   Úloadro   s      r   rt   zLargeScaleKV.loadý   rr   r   c                 ó8   — | j                   j                  |«      S r    )rl   Úsize)r   rp   s     r   rv   zLargeScaleKV.size   s   € Ø�}‰}×!Ñ! 'Ó*Ð*r   N)rP   rQ   rR   r   rn   rt   rv   rS   r   r   rj   rj   ö   s   „ ò,ò-ò-ó+r   rj   c                   ó   — e Zd Zd„ Zd„ Zy)ÚHeterClientc                 ó<   — t        j                  |||«      | _        y r    )r   rx   Úheter_client_)r   ÚendpointÚprevious_endpointr   s       r   r   zHeterClient.__init__  s   € Ü!×-Ñ-ØÐ'¨ó
ˆÕr   c                 ó8   — | j                   j                  «        y r    )rz   r;   r2   s    r   r;   zHeterClient.stop
  r@   r   N)rP   rQ   rR   r   r;   rS   r   r   rx   rx     s   „ ò
ó
"r   rx   )Ú__doc__r!   Ú"paddle.distributed.ps.utils.publicr   Úpaddle.frameworkr   Ú__all__r   rU   rj   rx   rS   r   r   Ú<module>r‚      sI   ðñ:ó Ý >Ý !à
€÷nHñ nHôb�\ô ÷:+ñ +÷"ò "r   