Ë
    ˆ\;j™3  ã                   óô   — 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	 d dl
mZmZ d dlmZ  edg d¢«      Zd	Zd
ZdZdad ad„ Zd„ Zd„ Zd„ Zd„ Zdd„Zddefd„Zddefd„Zd„ Zd„ Zd„ Zd„ Z d„ Z!d„ Z"y)é    N)Ú
namedtuple)Úcore)ÚNode)Ú
PythonFuncÚ
_serialize)ÚloggerÚ
WorkerInfo)ÚnameÚrankÚipÚportéÿÿÿÿiÿÿÿiÿàõc                 ó   — | a y ©N©Ú_barrier_store)Ústores    úcG:\00. PROJECTS\API\Inventory\templateJSON\kerjaOCR\Lib\site-packages\paddle/distributed/rpc/rpc.pyÚ_set_barrier_storer   &   s   € à�Nó    c                   ó   — b y r   r   © r   r   Ú_del_barrier_storer   +   s   € ár   c                 ó„   — t        j                  t        | |||«      «      }t        j	                  t        |«      |«       y r   )ÚpickleÚdumpsr	   r   ÚsetÚstr)r
   r   r   r   Ú	self_infos        r   Ú_set_self_infor    0   s/   € Ü—‘œZ¨¨d°B¸Ó=Ó>€IÜ×Ñ”s˜4“y )Õ,r   c                 ó"  — g }t        «       }t        | «      D ]t  }t        j                  t        j                  t        |«      «      «      }|j                  |vsJ d«       ‚|j                  |j                  «       |j                  |«       Œv |S )Nz:The Worker name must be unique, but name `{}` is repeated.)
r   Úranger   Úloadsr   Úgetr   r
   ÚaddÚappend)Ú
world_sizeÚ	all_infosÚsr   Úinfos        r   Ú_exchange_all_service_infosr+   5   s~   € Ø€IÜ‹€AÜ�jÖ!ˆÜ�|‰|œN×.Ñ.¬s°4«yÓ9Ó:ˆà�I‰I˜QÑð	HàGó	HØà	�‰ˆd�i‰iÔØ×Ñ˜Õð "ð Ðr   c                  ód   — t        «       } | j                  «       }| j                  «       }|› d|› �S )NÚ:)r   Úget_host_ipÚget_free_port)Únoder   Ú	free_ports      r   Ú_gen_endpointr2   B   s6   € Ü‹6€DØ	×	Ñ	Ó	€BØ×"Ñ"Ó$€IØˆT��9�+ÐÐr   c                 óâ  — |€t        t        j                  d   «      n|}|€t        t        j                  d   «      n|}t        j                  dd«      }|€
t	        «       }t        j                  d|› d|› �«       |�|nt        j                  d   }|j                  d«      \  }}t        |«      }t        t        j                  d	d
«      «      }t        j                  |||dk(  ||¬«      }t        |«       |j                  d«      \  }	}
t        |
«      }
t        | ||	|
«       t        |«      }g }|D ]S  }t        j                  |j                  |j                  |j                   |j"                  «      }|j%                  |«       ŒU t        j&                  | |«       t        j(                  «        t+        ||«       t        j,                  «        t        j                  d|› d�«       y)aÅ  
    init rpc.

    Args:
        name (str): worker name.
        rank (int, optional): worker id, default is None.
        world_size (int, optional): number of workers, default is None.
        master_endpoint (str, optional): id address of master, other nodes communicate with the master to
            get the information of all worker nodes, default is None.

    Returns:
        None.

    Examples:
        .. code-block:: python

            >>> # doctest: +REQUIRES(env:DISTRIBUTED)
            >>> import paddle.distributed.rpc as rpc

            >>> rpc.init_rpc("worker0", rank=0, world_size=1,
            ...             master_endpoint="127.0.0.1:8001")

            >>> rpc.shutdown()

    NÚPADDLE_TRAINER_IDÚPADDLE_TRAINERS_NUMÚPADDLE_WORKER_ENDPOINTúTrainer z: worker endpoint: ÚPADDLE_MASTER_ENDPOINTr-   ÚFLAGS_stop_check_timeoutÚ900r   )Útimeoutz: Init RPC done!)ÚintÚosÚenvironÚgetenvr2   r   r*   Úsplitr   ÚTCPStorer   r    r+   r	   r
   r   r   r   r&   Úinit_and_set_agent_instanceÚrpc_start_workerÚ_barrier_never_timeoutÚrpc_start_client)r
   r   r'   Úmaster_endpointÚworker_endpointÚmaster_addrÚmaster_portÚstop_check_timeoutr   r   r   r(   Úc_infosÚ	node_infor*   s                  r   Úinit_rpcrM   I   sÂ  € ð4 48°<Œ3Œr�z‰zÐ-Ñ.Ô/ÀT€Dð Ðô 	ŒB�J‰JÐ,Ñ-Ô.àð ô
 —i‘iÐ 8¸$Ó?€OØÐÜ'›/ˆÜ
‡K�K�(˜4˜&Ð 3°OÐ3DÐEÔFð Ð&ñ 	ä�Z‰ZÐ0Ñ1ð ð
  /×4Ñ4°SÓ9Ñ€K�Ü�kÓ"€KÜœRŸY™YÐ'AÀ5ÓIÓJÐÜ�M‰MØØØ�‰	ØØ"ô€Eô �uÔØ×$Ñ$ SÓ)�H€BˆÜˆt‹9€DÜ�4˜˜r 4Ô(Ü+¨JÓ7€IØ€GÛˆ	Ü�‰Ø�N‰N˜IŸN™N¨I¯L©L¸)¿.¹.ó
ˆð 	�‰�tÕð	 ô
 	×$Ñ$ T¨7Ô3Ü×ÑÔä˜4 Ô,Ü×ÑÔÜ
‡K�K�(˜4˜&Ð 0Ð1Õ2r   c                 ó@   — t        | ||||«      }|j                  «       S )aë  
    Make a blocking RPC call to run function ``fn`` on worker ``to``. Attention: Users must use this API in a secure network environment.

    Args:
        to (str): name of the destination worker.
        fn (fn): a callable function, such as Python callables.
        args (tuple, optional): the argument tuple for the ``fn`` invocation, default is None.
        kwargs (dict, optional): is a dictionary of keyword arguments for the ``fn``
                       invocation, default is None.
        timeout (int, optional): timeout in seconds to use for this RPC. If
                                   the RPC does not complete in this amount of
                                   time, an exception indicating it has
                                   timed out will be raised. A value less than or equal to 0
                                   indicates an infinite timeout, i.e. a timeout
                                   error will never be raised. The default value is -1.

    Returns:
        Returns the result of running ``fn`` with ``args`` and ``kwargs``.

    Examples:
        .. code-block:: python

            >>> # doctest: +REQUIRES(env:DISTRIBUTED)
            >>> import paddle.distributed.rpc as rpc

            >>> def add(a, b):
            ...     return a + b

            >>> rpc.init_rpc("worker0", rank=0, world_size=1,
            ...         master_endpoint="127.0.0.1:8002")

            >>> ret = rpc.rpc_sync("worker0", add, args=(2, 3))
            >>> rpc.shutdown()

    )Ú_invoke_rpcÚwait)ÚtoÚfnÚargsÚkwargsr;   Úfuts         r   Úrpc_syncrV   �   s#   € ôH �b˜"˜d F¨GÓ
4€CØ�8‰8‹:Ðr   c                 ó    — t        | ||||«      S )a�  
    Make a non-blocking RPC call to run function ``fn`` on worker ``to``. Attention: Users must use this API in a secure network environment.

    Args:
        to (str): name of the destination worker.
        fn (fn): a callable function, such as Python callables.
        args (tuple, optional): the argument tuple for the ``fn`` invocation, default is None.
        kwargs (dict, optional): is a dictionary of keyword arguments for the ``fn``
                       invocation, default is None.
        timeout (int, optional): timeout in seconds to use for this RPC. If
                                   the RPC does not complete in this amount of
                                   time, an exception indicating it has
                                   timed out will be raised. A value less than or equal to 0
                                   indicates an infinite timeout, i.e. a timeout
                                   error will never be raised. The default value is -1.

    Returns:
        Returns a :class:`FutureWrapper` object that can be waited
        on. When completed, the return value of ``fn`` on ``args`` and
        ``kwargs`` can be got by `fut.wait()`.

    Examples:
        .. code-block:: python

            >>> # doctest: +REQUIRES(env:DISTRIBUTED)
            >>> import paddle.distributed.rpc as rpc

            >>> def add(a, b):
            ...     return a + b

            >>> rpc.init_rpc("worker0", rank=0, world_size=1,
            ...         master_endpoint="127.0.0.1:8003")

            >>> fut = rpc.rpc_async("worker0", add, args=(2, 3))
            >>> print(fut.wait())
            5

            >>> rpc.shutdown()

    )rO   )rQ   rR   rS   rT   r;   s        r   Ú	rpc_asyncrX   ·   s   € ôR �r˜2˜t V¨WÓ5Ð5r   c                 óœ   — |r|nd}|r|ni }t        t        |||«      «      }|dz  }|dk  rt        n|}t        j                  | ||«      }|S )Nr   iè  r   )r   r   Ú_MAX_RPC_TIMEOUT_MSr   Ú
invoke_rpc)rQ   rR   rS   rT   r;   Ú
serial_objÚ
timeout_msÚfutures           r   rO   rO   ã   sU   € Ù‰4˜R€DÙ‰V 2€FÜœJ r¨4°Ó8Ó9€JØ˜4‘€JØ(2°aªÕ$¸Z€JÜ�_‰_˜R ¨ZÓ8€FØ€Mr   c                 óº  ‡ ‡— t        j                  t        ¬«      Š|dk  ry dt        t        «      z   dz   }t        dz  a‰ dk(  }ˆ ˆfd„}|rPt        d|«      D �cg c]  }|t        |«      z   ‘Œ }}t        j                  |t        d«      z   d«        ||«       y |t        d«      z   g} ||«       t        j                  |t        ‰ «      z   d«       y c c}w )N)Údaysé   zBarrier/Ú/é   r   c                 óV  •— t        j                   «       }t        | «      dkD  r†t        j                  d«       t        j                   «       |z
  }t        j                  |¬«      ‰kD  rt        dj                  | ‰«      «      ‚t        t        d„ | «      «      } t        | «      dkD  rŒ…y y )Nr   gš™™™™™¹?)Úsecondsz4Keys {} are not ready sinck rank {} is waiting them.c                 óD   — t        t        j                  | «      «      dk7  S )Nrc   )r<   r   r$   )Úkeys    r   Ú<lambda>zC_barrier_never_timeout.<locals>._check_keys_ready.<locals>.<lambda>  s   € ¤3¤~×'9Ñ'9¸#Ó'>Ó#?À1Ò#Dr   )	ÚtimeÚlenÚsleepÚdatetimeÚ	timedeltaÚRuntimeErrorÚformatÚlistÚfilter)Ú	wait_keysÚ
start_timeÚelapse_timeÚglobal_rankr;   s      €€r   Ú_check_keys_readyz1_barrier_never_timeout.<locals>._check_keys_readyù   s�   ø€ Ü—Y‘Y“[ˆ
Ü�)‹n˜qÒ Ü�J‰J�sŒOÜŸ)™)›+¨
Ñ2ˆKÜ×!Ñ!¨+Ô6¸Ò@Ü"ØJ×QÑQØ! ;óóð ô
 ÜÑDÀiÓPóˆIô �)‹n˜qÕ r   )rl   rm   Ú_BARRIER_TIMEOUT_MAX_DAYSr   Ú_barrier_countr"   r   r%   )ru   Úglobal_world_sizeÚbarrier_prefixÚ	is_masterrv   r   rr   r;   s   `      @r   rD   rD   í   sÞ   ù€ ä× Ñ Ô&?Ô@€Gà˜1ÒØð  ¤#¤nÓ"5Ñ5¸Ñ;€NÜ�aÑ€NØ˜qÑ €Iõñ ô 49¸Ð<MÔ3Nó
Ù3N¨4ˆNœS ›YÓ&Ð3Nð 	ð 
ô 	×Ñ˜>¬C°«FÑ2°AÔ6Ù˜)Õ$à#¤c¨!£fÑ,Ð-ˆ	Ù˜)Ô$Ü×Ñ˜>¬C°Ó,<Ñ<¸aÕ@ùò
s   ÁCc                  óÜ   — t        «       } | j                  }t        t        «       «      }t	        ||«       t        j                  «        t        «        t        j                  d|› d�«       y)a+  
    Perform a shutdown of the RPC agent, stop the worker and destroy the agent.
    This will block until all local and remote RPC processes reach this method
    and wait for all outstanding work to complete.

    Returns:
        None.

    Examples:
        .. code-block:: python

            >>> # doctest: +REQUIRES(env:DISTRIBUTED)
            >>> import paddle.distributed.rpc as rpc

            >>> rpc.init_rpc("worker0", rank=0, world_size=1,
            ...             master_endpoint="127.0.0.1:8004")

            >>> rpc.shutdown()

    r7   z: rpc shutdown!N)
Úget_current_worker_infor   rj   Úget_all_worker_infosrD   r   Úrpc_stop_workerr   r   r*   )r*   r   r'   s      r   Úshutdownr€     sT   € ô* #Ó$€DØ�9‰9€DÜÔ)Ó+Ó,€Jä˜4 Ô,Ü×ÑÔÜÔÜ
‡K�K�(˜4˜& Ð0Õ1r   c                 ó,   — t        j                  | «      S )aÍ  
    Get worker information by worker name.

    Args:
        name (str): name of the worker.

    Returns:
        class `WorkerInfo` with attribute `name`, `rank`, `ip` and `port`.

    Examples:
        .. code-block:: python

            >>> # doctest: +REQUIRES(env:DISTRIBUTED)
            >>> import paddle.distributed.rpc as rpc
            >>> import os

            >>> os.environ["PADDLE_WORKER_ENDPOINT"] = "127.0.0.1:9002"
            >>> rpc.init_rpc("worker0", rank=0, world_size=1,
            ...             master_endpoint="127.0.0.1:8005")

            >>> print(rpc.get_worker_info("worker0"))
            {name: worker0, rank: 0, ip: 127.0.0.1, port: 9002}

            >>> rpc.shutdown()

    )r   Úrpc_get_worker_info)r
   s    r   Úget_worker_inforƒ   5  s   € ô6 ×#Ñ# DÓ)Ð)r   c                  ó*   — t        j                  «       S )aY  
    Get all worker informations.

    Returns:
        List[WorkerInfo].

    Examples:
        .. code-block:: python

            >>> # doctest: +REQUIRES(env:DISTRIBUTED)
            >>> import paddle.distributed.rpc as rpc
            >>> import os

            >>> os.environ["PADDLE_WORKER_ENDPOINT"] = "127.0.0.1:9003"
            >>> rpc.init_rpc("worker0", rank=0, world_size=1,
            ...         master_endpoint="127.0.0.1:8006")

            >>> print(rpc.get_all_worker_infos())
            [{name: worker0, rank: 0, ip: 127.0.0.1, port: 9003}]

            >>> rpc.shutdown()

    )r   Úrpc_get_all_worker_infosr   r   r   r~   r~   S  s   € ô0 ×(Ñ(Ó*Ð*r   c                  ó*   — t        j                  «       S )a’  
    Get current worker information.

    Returns:
        class `WorkerInfo` with attribute `name`, `rank`, `ip` and `port`.

    Examples:
        .. code-block:: python

            >>> # doctest: +REQUIRES(env:DISTRIBUTED)
            >>> import paddle.distributed.rpc as rpc
            >>> import os

            >>> os.environ["PADDLE_WORKER_ENDPOINT"] = "127.0.0.1:9004"
            >>> rpc.init_rpc("worker0", rank=0, world_size=1,
            ...             master_endpoint="127.0.0.1:8007")

            >>> print(rpc.get_current_worker_info())
            {name: worker0, rank: 0, ip: 127.0.0.1, port: 9004}

            >>> rpc.shutdown()

    )r   Úrpc_get_current_worker_infor   r   r   r}   r}   n  s   € ô0 ×+Ñ+Ó-Ð-r   )NNN)#rl   r=   r   ri   Úcollectionsr   Úpaddle.baser   Ú!paddle.distributed.launch.contextr   Úpaddle.distributed.rpc.internalr   r   Ú%paddle.distributed.utils.launch_utilsr   r	   Ú_DEFAULT_RPC_TIMEOUTrZ   rw   r   rx   r   r   r    r+   r2   rM   rV   rX   rO   rD   r€   rƒ   r~   r}   r   r   r   Ú<module>rŽ      s­   ðó Û 	Û Û Ý "å Ý 2ß BÝ 8á˜Ò&DÓE€
àÐ Ø Ð Ø$Ð à€ð €òò
ò
-ò

òóC3ðL  tÐ5Ió %ðP  ¨Ð6Jó )6òXò&AòR2ò>*ò<+ó6.r   