Ë
    g^(hh+  ã                   ór  — d dl Z d dlZ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	Z	d dl
mZ d dlmZ g d¢Z ej                   «       aej$                  Zej(                  Z G d„ de«      Z G d„ d«      Z e«       Zd	„ Zd
„ Zd„ Zd„ Zd„ Zd„ Z e j>                  dg d¢«      Z  e j>                  dddg«      Z!y)é    N)ÚEnum)Ú_get_current_rpc_agent)ÚRPCExecModeÚ	serializeÚdeserializeÚ	PythonUDFÚRemoteExceptionc                   ó   — e Zd ZdZdZdZdZy)r   ÚsyncÚasyncÚ	async_jitÚremoteN)Ú__name__Ú
__module__Ú__qualname__ÚSYNCÚASYNCÚ	ASYNC_JITÚREMOTE© ó    ú\/var/www/skyplay_api_hub/venv/lib/python3.12/site-packages/torch/distributed/rpc/internal.pyr   r      s   „ Ø€DØ€EØ€IØ�Fr   r   c                   óp   — e Zd ZdZd„ Zd„ Zed„ «       Zd„ Zed„ «       Z	d„ Z
d„ Zed	„ «       Zd
„ Zd„ Zd„ Zy)Ú_InternalRPCPicklera	  
    This class provides serialize() and deserialize() interfaces to serialize
    data to be "binary string + tensor table" format
    So for RPC python UDF function and args, non tensor data will be serialized
    into regular binary string, tensor data will be put into thread local tensor
    tables, this serialization format is consistent with builtin operator and args
    using JIT pickler. This format will make tensor handling in C++ much easier,
    e.g. attach tensor to distributed autograd graph in C++
    c                 ó¦   — t         j                  j                  «       | _        | j                  | j                  t
        j                  <   i | _        y ©N)ÚcopyregÚdispatch_tableÚcopyÚ_dispatch_tableÚ_tensor_reducerÚtorchÚTensorÚ_class_reducer_dict)Úselfs    r   Ú__init__z_InternalRPCPickler.__init__+   s;   € ä&×5Ñ5×:Ñ:Ó<ˆÔØ-1×-AÑ-Aˆ×ÑœUŸ\™\Ñ*à#%ˆÕ r   c                 ó@   — || j                   vr|| j                   |<   y y r   )r$   )r%   Ú	obj_classÚreducers      r   Ú_register_reducerz%_InternalRPCPickler._register_reducer2   s%   € à˜D×4Ñ4Ñ4Ø29ˆD×$Ñ$ YÒ/ð 5r   c                 ó(   — t         j                  |   S r   )Ú_thread_local_tensor_tablesÚrecv_tables)ÚclsÚtensor_indexs     r   Ú_tensor_receiverz$_InternalRPCPickler._tensor_receiver7   s   € ô +×6Ñ6°|ÑDÐDr   c                 óž   — t         j                  j                  |«       t        t         j                  «      dz
  }t        j
                  |ffS )Né   )r,   Úsend_tablesÚappendÚlenr   r0   )r%   Útensorr/   s      r   r!   z#_InternalRPCPickler._tensor_reducer<   s?   € ä#×/Ñ/×6Ñ6°vÔ>ÜÔ6×BÑBÓCÀaÑGˆÜ#×4Ñ4°|°oÐFÐFr   c                 óT   — t         j                  j                  j                  |«      S r   )ÚdistÚrpcÚPyRRefÚ_deserialize)r.   Úrref_fork_datas     r   Ú_py_rref_receiverz%_InternalRPCPickler._py_rref_receiverB   s   € ä�x‰x�‰×+Ñ+¨NÓ;Ð;r   c                 óH   — |j                  «       }t        j                  |ffS r   )Ú
_serializer   r=   )r%   Úpy_rrefr<   s      r   Ú_py_rref_reducerz$_InternalRPCPickler._py_rref_reducerF   s$   € Ø ×+Ñ+Ó-ˆÜ#×5Ñ5¸Ð7HÐIÐIr   c                 ó$   — | j                  |«      S r   )rA   )r%   Úrrefs     r   Ú_rref_reducerz!_InternalRPCPickler._rref_reducerJ   s   € Ø×$Ñ$ TÓ*Ð*r   c                 ón   — t        j                  |«      }t        j                  j	                  |«      }|S )zŽ
        Given a serialized representation of a ScriptModule created with torch.jit.save,
        loads and returns the ScriptModule.
        )ÚioÚBytesIOr"   ÚjitÚload)r.   Úscript_module_serializedÚfÚms       r   Ú_script_module_receiverz+_InternalRPCPickler._script_module_receiverM   s*   € ô �J‰JÐ/Ó0ˆÜ�I‰I�N‰N˜1ÓˆØˆr   c                 ó¬   — t        j                  «       }t        j                  j	                  ||«       t
        j                  |j                  «       ffS )z,
        Serializes a ScriptModule.
        )rF   rG   r"   rH   Úsaver   rM   Úgetvalue)r%   Úscript_modulerK   s      r   Ú_script_module_reducerz*_InternalRPCPickler._script_module_reducerW   s:   € ô �J‰J‹LˆÜ�	‰	�‰�} aÔ(Ü#×;Ñ;¸a¿j¹j»l¸_ÐMÐMr   c                 ó  — t        j                  «       }t        |«      }| j                  |_        | j
                  |j                  t        j                  j                  <   | j                  |j                  t        j                  j                  <   t        |t        j                  j                  «      r#| j                  |j                  |j                   <   | j"                  j%                  «       D ]  }| j"                  |   |j                  |<   Œ  t'        t(        d«      rt(        j*                  }nd}g t(        _        |j-                  |«       t(        j*                  }|�|t(        _        nt(        `|j/                  «       |fS )ze
        Serialize non tensor data into binary string, tensor data into
        tensor table
        r3   N)rF   rG   Ú_picklerr    r   rA   r8   r9   r:   rD   ÚRRefÚ
isinstancer"   rH   ÚScriptModulerR   Ú	__class__r$   ÚkeysÚhasattrr,   r3   ÚdumprP   )r%   ÚobjrK   ÚpÚ
class_nameÚold_send_tablesÚtensorss          r   r   z_InternalRPCPickler.serialize_   s7  € ô
 �J‰J‹LˆÜ�Q‹KˆØ×/Ñ/ˆÔð -1×,AÑ,Aˆ×ÑœŸ™Ÿ™Ñ)ð +/×*<Ñ*<ˆ×ÑœŸ™Ÿ™Ñ'ô �cœ5Ÿ9™9×1Ñ1Ô2à.2×.IÑ.IˆA×Ñ˜SŸ]™]Ñ+ð ×2Ñ2×7Ñ7Ó9ò 	PˆJØ+/×+CÑ+CÀJÑ+OˆA×Ñ˜ZÒ(ð	Pô
 Ô.°Ô>Ü9×EÑE‰Oà"ˆOØ24Ô#Ô/à	�‰ˆsŒô .×9Ñ9ˆØÐ&Ø6EÔ'Õ3ä+Ð7à—
‘
“˜gÐ&Ð&r   c                 óV  — t        t        d«      rt        j                  }nd}|t        _        	 t        t	        j
                  |«      «      }|j                  «       }|�|t        _        |S t        `|S # t        $ r*}t        |«      dz   }t        |«      }||_	        Y d}~ŒEd}~ww xY w)zJ
        Deserialize binary string + tensor table to original obj
        r-   NzŽ Default RPC pickler does not serialize
            function code. Ensure that UDFs are defined on both caller and
            callee modules.)
rZ   r,   r-   Ú
_unpicklerrF   rG   rI   ÚAttributeErrorÚstrÚ	__cause__)r%   Úbinary_dataÚtensor_tableÚold_recv_tablesÚ	unpicklerÚretÚeÚ
except_strs           r   r   z_InternalRPCPickler.deserialize”   s°   € ô Ô.°Ô>Ü9×EÑE‰Oà"ˆOØ2>Ô#Ô/ð	Ü"¤2§:¡:¨kÓ#:Ó;ˆIØ—.‘.Ó"ˆCð  Ð&Ø6EÔ'Ô3ð ˆ
ô ,Ð7àˆ
øô) ò 	ô �A“ðñð ô ! Ó,ˆCàˆC�M‰Mûð	ús   °.A5 Á5	B(Á> B#Â#B(N)r   r   r   Ú__doc__r&   r*   Úclassmethodr0   r!   r=   rA   rD   rM   rR   r   r   r   r   r   r   r       sq   „ ñò&ò:ð
 ñEó ðEòGð ñ<ó ð<òJò+ð ñó ðòNò3'ój#r   r   c                 ó,   — t         j                  | «      S r   )Ú_internal_rpc_picklerr   )r\   s    r   r   r   ¾   s   € Ü ×*Ñ*¨3Ó/Ð/r   c                 ó.   — t         j                  | |«      S r   )rp   r   )rf   rg   s     r   r   r   Â   s   € Ü ×,Ñ,¨[¸,ÓGÐGr   c                 ó~  — 	 t        | t        «      r| ‚ | j                  | j                  i | j                  ¤Ž}|S # t
        $ rw}dt        «       j                  «       › dt        |«      › dt        j                  «       › �}t        |t        j                  ¬«       t        |t        |«      «      }Y d}~|S d}~ww xY w)zò
    This function is exclusively called from C++.
    See ``torch/csrc/distributed/rpc/python_rpc_handler.cpp``.

    Runs a Python UDF and returns its return value.
    Wraps any exception in ``RemoteException`` if the function raises.
    zOn z:
ú
)ÚfileN)rV   rc   ÚfuncÚargsÚkwargsÚ	Exceptionr   Úget_worker_infoÚreprÚ	tracebackÚ
format_excÚprintÚsysÚstderrr	   Útype)Ú
python_udfÚresultrk   rl   s       r   Ú_run_functionrƒ   Æ   s°   € ð6Ü�j¤.Ô1ØÐØ �—‘ *§/¡/ÐG°Z×5FÑ5FÑGˆð €Møô ò 6ð Ô(Ó*×:Ñ:Ó<Ð=¸SÜ�A‹wˆi�rœ)×.Ñ.Ó0Ð1ð3ð 	ô 	ˆjœsŸz™zÕ*Ü  ¬T°!«WÓ5ŒØ€Mûð6ús   ‚8< ¼	B<ÁA,B7Â7B<c                 ó  — t        | t        «      rC| j                  j                  d«      j	                  d«      }d }	 | j                  |«      }|�|‚y y # t        $ r }t        dt        |«      › d|› �«      |‚d }~ww xY w)Nzutf-8Úunicode_escapez8Failed to create original exception type. Error msg was z' Original exception on remote side was )	rV   r	   ÚmsgÚencodeÚdecodeÚexception_typeÚBaseExceptionÚRuntimeErrorrd   )r‚   Úexception_msgÚexcrk   s       r   Ú_handle_exceptionrŽ   Ý   sš   € Ü�&œ/Ô*ØŸ
™
×)Ñ)¨'Ó2×9Ñ9Ð:JÓKˆð ˆð	Ø×'Ñ'¨Ó6ˆCð ˆ?ØˆIð ð +øô ò 	ÜØJÌ3ÈqË6È(Ø9¸-¸ðJóð ðûð	ús   ¾A Á	A>ÁA9Á9A>c           	      ó8   — d| j                   › d|› d|› d|› d�	}|S )aÓ  
    Builds the key that RPC calls are profiled with using the autograd profiler.
    This will be the name of the corresponding Event recorded in the profiler.

    Args:
        exec_type (RPCExecMode): Type of RPC/RRef call
        func_name (str): Name of function being profiled.
        current_worker_name (str): Name of current worker.
        dst_worker_name (str): Name of the destination worker.

    Returns:
        String representing profiling key
    Úrpc_ú#ú(ú -> ú))Úvalue)Ú	exec_typeÚ	func_nameÚcurrent_worker_nameÚdst_worker_nameÚprofile_keys        r   Ú_build_rpc_profiling_keyr›   ï   s8   € ð" ˆy�‰Ð˜q  ¨1Ð-@Ð,AÀÀoÐEVÐVWÐXð ð Ðr   c           	      ó  — t         j                  j                  «       sJ d«       ‚d| j                  › dt	        |«      › d|› d|› d�	}t         j                  j                  «       }t         j                  j                  ||«       |S )ar  
    This function should be called from RPC/RRef functions to create a
    RecordFunction object for profiling. This function also runs the before
    callbacks that start the profiling, though the user is responsible for
    running the appropriate callbacks when the function to be profiled finishes.

    Args:
        exec_type (RPCExecMode): Type of RPC/RRef call
        func_name (str): Name of function being profiled.
        current_worker_name (str): Name of current worker.
        dest_worker_name (str): Name of the destination worker.

    Returns:
        An instance of `torch.autograd._RecordFunction`.
    z$Autograd profiler should be enabled.r�   r‘   r’   r“   r”   )r"   ÚautogradÚ_profiler_enabledr•   rd   Ú_RecordFunctionÚ_run_before_callbacks)r–   r—   r˜   Údest_worker_namerš   Úrfs         r   Ú_start_record_functionr£     s{   € ô  �>‰>×+Ñ+Ô-ÐUÐ/UÓUÐ-Ø˜Ÿ™Ð)¨¬3¨y«>Ð*:¸!Ð<OÐ;PÐPTÐUeÐTfÐfgÐh€KÜ	�‰×	'Ñ	'Ó	)€BÜ	‡N�N×(Ñ(¨¨[Ô9Ø€Ir   r   )ru   rv   rw   r	   r†   r‰   )"Úcollectionsr   rF   Úpickler~   Ú	threadingr{   Úenumr   r"   Útorch.distributedÚdistributedr8   Útorch._C._distributed_rpcr   Ú__all__Úlocalr,   ÚPicklerrT   Ú	Unpicklerrb   r   r   rp   r   r   rƒ   rŽ   r›   r£   Ú
namedtupler   r	   r   r   r   ú<module>r°      sÈ   ðã Û Û 	Û Û 
Û Û Ý ã Ý  Ý <ò V€ð .˜iŸo™oÓ/Ð Ø�>‰>€Ø×Ñ€
ô�$ô ÷Wñ Wñv ,Ó-Ð ò0òHòò.ò$ò,ð. #ˆK×"Ñ" ;Ò0JÓK€	Ø(�+×(Ñ(Ð):¸UÐDTÐ<UÓV�r   