o
    Ù­jT'  ã                   @   sâ  U d 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m	Z	 ddl
mZmZmZmZmZ ddlZddlmZ ddlmZmZmZmZmZmZ e d	¡Zeeeeef f Zee d
< eeeef  Z!ee d< eee!f Z"ee d< eG dd„ dƒƒZ#de!ddfdd„Z$d2dd„Z%defdd„Z&defdd„Z'de(fdd„Z)deddfdd„Z*defdd„Z+d ed!edefd"d#„Z,d$ej-defd%d&„Z.e	G d'd(„ d(eƒƒZ/d ej0d)e/dej0fd*d+„Z1d2d,d-„Z2dee fd.d/„Z3G d0d1„ d1ƒZ4dS )3z-XGBoost collective communication related API.é    N)Ú	dataclass)ÚIntEnumÚunique)ÚAnyÚDictÚOptionalÚ	TypeAliasÚUnioné   )Ú_T)Ú_LIBÚ_check_callÚ
build_infoÚc_strÚmake_jcargsÚpy_strz[xgboost.collective]Ú_ConfÚ_ArgValsÚ_Argsc                   @   st   e Zd ZU dZdZee ed< dZee ed< dZ	ee
 ed< dZee ed< dZee ed< ded	efd
d„ZdS )ÚConfiga§  User configuration for the communicator context. This is used for easier
    integration with distributed frameworks. Users of the collective module can pass the
    parameters directly into tracker and the communicator.

    .. versionadded:: 3.0

    Attributes
    ----------
    retry : See `dmlc_retry` in :py:meth:`init`.

    timeout :
        See `dmlc_timeout` in :py:meth:`init`. This is only used for communicators, not
        the tracker. They are different parameters since the timeout for tracker limits
        only the time for starting and finalizing the communication group, whereas the
        timeout for communicators limits the time used for collective operations, like
        :py:meth:`allreduce`.

    tracker_host_ip : See :py:class:`~xgboost.tracker.RabitTracker`.

    tracker_port : See :py:class:`~xgboost.tracker.RabitTracker`.

    tracker_timeout : See :py:class:`~xgboost.tracker.RabitTracker`.

    NÚretryÚtimeoutÚtracker_host_ipÚtracker_portÚtracker_timeoutÚargsÚreturnc                 C   s,   | j dur
| j |d< | jdur| j|d< |S )z*Update the arguments for the communicator.NÚ
dmlc_retryÚdmlc_timeout)r   r   ©Úselfr   © r!   úO/var/www/html/CropPilot/venv/lib/python3.10/site-packages/xgboost/collective.pyÚget_comm_config:   s
   



zConfig.get_comm_config)Ú__name__Ú
__module__Ú__qualname__Ú__doc__r   r   ÚintÚ__annotations__r   r   Ústrr   r   r   r#   r!   r!   r!   r"   r      s   
 r   r   r   c                  K   s   t t tdi | ¤Ž¡ƒ dS )aë  Initialize the collective library with arguments.

    Parameters
    ----------
    args :
        Keyword arguments representing the parameters and their values.

        Accepted parameters:
          - dmlc_communicator: The type of the communicator.
            * rabit: Use Rabit. This is the default if the type is unspecified.
            * federated: Use the gRPC interface for Federated Learning.

        Only applicable to the Rabit communicator:
          - dmlc_tracker_uri: Hostname of the tracker.
          - dmlc_tracker_port: Port number of the tracker.
          - dmlc_task_id: ID of the current task, can be used to obtain deterministic
          - dmlc_retry: The number of retry when handling network errors.
          - dmlc_timeout: Timeout in seconds.
          - dmlc_nccl_path: Path to load (dlopen) nccl for GPU-based communication.

        Only applicable to the Federated communicator:
          - federated_server_address: Address of the federated server.
          - federated_world_size: Number of federated workers.
          - federated_rank: Rank of the current worker.
          - federated_server_cert: Server certificate file path. Only needed for the SSL
            mode.
          - federated_client_key: Client key file path. Only needed for the SSL mode.
          - federated_client_cert: Client certificate file path. Only needed for the SSL
            mode.

        Use upper case for environment variables, use lower case for runtime
        configuration.

    Nr!   )r   r   ÚXGCommunicatorInitr   )r   r!   r!   r"   ÚinitC   s   #r,   c                   C   ó   t t ¡ ƒ dS )zFinalize the communicator.N)r   r   ÚXGCommunicatorFinalizer!   r!   r!   r"   Úfinalizei   ó   r/   c                  C   ó   t  ¡ } | S )zjGet rank of current process.

    Returns
    -------
    rank : int
        Rank of current process.
    )r   ÚXGCommunicatorGetRank©Úretr!   r!   r"   Úget_rankn   ó   r5   c                  C   r1   )z`Get total number workers.

    Returns
    -------
    n :
        Total number of process.
    )r   ÚXGCommunicatorGetWorldSizer3   r!   r!   r"   Úget_world_sizez   r6   r8   c                  C   s   t  ¡ } t| ƒS )z.If the collective communicator is distributed.)r   ÚXGCommunicatorIsDistributedÚbool)Úis_distr!   r!   r"   Úis_distributed†   s   r<   Úmsgc                 C   sP   t | tƒs	t| ƒ} t ¡ }|dkrtt t|  ¡ ƒ¡ƒ dS t|  ¡ dd� dS )zòPrint message to the communicator.

    This function can be used to communicate the information of
    the progress to the communicator.

    Parameters
    ----------
    msg : str
        The message to be printed to the communicator.
    r   T)ÚflushN)	Ú
isinstancer*   r   r9   r   ÚXGCommunicatorPrintr   ÚstripÚprint)r=   r;   r!   r!   r"   Úcommunicator_printŒ   s   
rC   c                  C   s*   t  ¡ } tt t  | ¡¡ƒ | j}t|ƒS )zdGet the processor name.

    Returns
    -------
    name :
        The name of processor(host)
    )ÚctypesÚc_char_pr   r   ÚXGCommunicatorGetProcessorNameÚbyrefÚvaluer   )Úname_strrH   r!   r!   r"   Úget_processor_name    s   rJ   ÚdataÚrootc                 C   sÐ   t ƒ }t ¡ }||kr | dusJ dƒ‚tj| tjd�}t|ƒ|_tt	 
t |¡t tj¡|¡ƒ ||krStj|j ƒ }tt	 
t |tj¡|j|¡ƒ t |j¡} ~| S tt	 
t t |¡tj¡|j|¡ƒ ~| S )aS  Broadcast object from one node to all other nodes.

    Parameters
    ----------
    data : any type that can be pickled
        Input data, if current rank does not equal root, this can be None
    root : int
        Rank of the node to broadcast data from.

    Returns
    -------
    object : int
        the result of broadcast.
    Nz&need to pass in data when broadcasting)Úprotocol)r5   rD   Úc_ulongÚpickleÚdumpsÚHIGHEST_PROTOCOLÚlenrH   r   r   ÚXGCommunicatorBroadcastrG   ÚsizeofÚc_charÚcastÚc_void_pÚloadsÚrawrE   )rK   rL   ÚrankÚlengthÚsÚdptrr!   r!   r"   Ú	broadcast®   s8   
ÿÿÿÿúÿÿr^   Údtypec                 C   s¾   t  d¡dt  d¡dt  d¡dt  d¡dt  d	¡d
t  d¡dt  d¡dt  d¡dt  d¡dt  d¡dt  d¡di}z| t  d¡di¡ W n	 tyN   Y nw | |vr[td| › d�ƒ‚||  S )NÚfloat16r   Úfloat32r
   Úfloat64é   Úint8é   Úint16é   Úint32é   Úint64é   Úuint8é   Úuint16é	   Úuint32é
   Úuint64é   Úfloat128é   z
data type z* is not supported on the current platform.)Únpr_   ÚupdateÚ	TypeError)r_   Ú	dtype_mapr!   r!   r"   Ú
_map_dtypeÞ   s(   










õÿrz   c                   @   s(   e Zd ZdZdZdZdZdZdZdZ	dS )	ÚOpz#Supported operations for allreduce.r   r
   rc   ru   re   rg   N)
r$   r%   r&   r'   ÚMAXÚMINÚSUMÚBITWISE_ANDÚ
BITWISE_ORÚBITWISE_XORr!   r!   r!   r"   r{   ÷   s    r{   Úopc                 C   sN   t | tjƒs
tdƒ‚|  ¡  ¡ }tt |j	 
t	j¡|jt|jƒt|ƒ¡ƒ |S )a'  Perform allreduce, return the result.

    Parameters
    ----------
    data :
        Input data.
    op :
        Reduction operator.

    Returns
    -------
    result :
        The result of allreduce, have same shape as data

    Notes
    -----
    This function is not thread-safe.
    z%allreduce only takes in numpy.ndarray)r?   rv   Úndarrayrx   ÚravelÚcopyr   r   ÚXGCommunicatorAllreducerD   Údata_asrW   Úsizerz   r_   r(   )rK   r‚   Úbufr!   r!   r"   Ú	allreduce  s   üÿrŠ   c                   C   r-   )zKill the process.N)r   r   ÚXGCommunicatorSignalErrorr!   r!   r!   r"   Úsignal_error$  r0   rŒ   c                  C   s¦   ddl m}  | jd urtj | j¡}nt| dƒr%t| jƒdkr%| jd }nd }|s+d S t 	|¡}|s4d S d }|D ]}| 
d¡rC|} nq8|d urQtj ||¡}|S d S )Nr   )ÚlibÚ__path__z
libnccl.so)Únvidia.ncclr�   Ú__file__ÚosÚpathÚdirnameÚhasattrrR   rŽ   ÚlistdirÚ
startswithÚjoin)r�   r“   ÚfilesÚlibnameÚnamer’   r!   r!   r"   Ú
_find_nccl)  s*   


þr›   c                   @   sB   e Zd ZdZdeddfdd„Zdefdd„Zdeddfd	d
„Z	dS )ÚCommunicatorContextzNA context controlling collective communicator initialization and finalization.r   r   Nc                 K   sf   || _ d}| |d ¡d urd S tƒ }|d sd S ztƒ }|r&|| j |< W d S W d S  ty2   Y d S w )NÚdmlc_nccl_pathÚUSE_DLOPEN_NCCL)r   Úgetr   r›   ÚImportError)r    r   ÚkeyÚbinfor’   r!   r!   r"   Ú__init__N  s   ÿÿzCommunicatorContext.__init__c                 C   s*   t di | j¤Ž tƒ sJ ‚t d¡ | jS )Nz8-------------- communicator say hello ------------------r!   )r,   r   r<   ÚLOGGERÚdebug)r    r!   r!   r"   Ú	__enter__`  s   

zCommunicatorContext.__enter__c                 G   s   t ƒ  t d¡ d S )Nz7--------------- communicator say bye ------------------)r/   r¤   r¥   r   r!   r!   r"   Ú__exit__f  s   zCommunicatorContext.__exit__)
r$   r%   r&   r'   r   r£   r   r¦   r   r§   r!   r!   r!   r"   rœ   K  s
    rœ   )r   N)5r'   rD   Úloggingr‘   rO   Údataclassesr   Úenumr   r   Útypingr   r   r   r   r	   Únumpyrv   Ú_typingr   Úcorer   r   r   r   r   r   Ú	getLoggerr¤   r*   r(   r   r)   r   r   r   r,   r/   r5   r8   r:   r<   rC   rJ   r^   r_   rz   r{   rƒ   rŠ   rŒ   r›   rœ   r!   r!   r!   r"   Ú<module>   s@     
*
&0
!"