§
    �Ÿj2  ã                   ó  — d 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 ddlZddlmZ ddlmZmZ dd	lmZ dd
lmZ ddlmZ ddlmZ ddlmZ  G d„ d¦  «        Z G d„ de	¦  «        Z G d„ de¦  «        ZdS )zbDefines a KernelClient that provides thread-safe sockets with async callbacks on message
replies.
é    N)ÚFuture)Úpartial)ÚThread)ÚAny)ÚIOLoop)ÚInstanceÚType)Ú
get_logger)Ú	zmqstreamé   )Ú	HBChannel)ÚKernelClient)ÚSessionc                   óþ   ‡ — e Zd ZdZdZdZdZdZdZde	j
        dz  dedz  dedz  ddfˆ fd„ZdZdefd	„Zdd
„Zdd„Zdd„Zdeeef         ddfd„Zdeddfd„Zdeeef         ddfd„Zdd„Zddeddfd„Zdd„Zˆ xZS )ÚThreadedZMQSocketChannelz.A ZMQ socket invoking a callback in the ioloopNÚsocketÚsessionÚloopÚreturnc                 ó  •‡ ‡— t          ¦   «                              ¦   «          |‰ _        |‰ _        |‰ _        t          ¦   «         Šdˆˆ fd„}‰ j        €J ‚‰ j                             |¦  «         ‰                     d¬¦  «         dS )a)  Create a channel.

        Parameters
        ----------
        socket : :class:`zmq.Socket`
            The ZMQ socket to use.
        session : :class:`session.Session`
            The session to use.
        loop
            A tornado ioloop to connect the socket to using a ZMQStream
        r   Nc                  ó&  •— 	 ‰j         €J ‚t          j        ‰j         ‰j        ¦  «        ‰_        ‰j                             ‰j        ¦  «         ‰                     d ¦  «         d S # t          $ r } ‰ 	                    | ¦  «         Y d } ~ d S d } ~ ww xY w©N)
r   r   Ú	ZMQStreamÚioloopÚstreamÚon_recvÚ_handle_recvÚ
set_resultÚ	ExceptionÚset_exception©ÚeÚfÚselfs    €€úd/var/www/finuniver-perm.ru/html/student/venv/lib/python3.11/site-packages/jupyter_client/threaded.pyÚsetup_streamz7ThreadedZMQSocketChannel.__init__.<locals>.setup_stream=   s¦   ø€ ð#Ø”{Ð.Ð.Ð.Ý'Ô1°$´+¸t¼{ÑKÔK�”Ø”×#Ò# DÔ$5Ñ6Ô6Ð6ð —’˜TÑ"Ô"Ð"Ð"Ð"øõ ð #ð #ð #Ø—’ Ñ"Ô"Ð"Ð"Ð"Ð"Ð"Ð"Ð"øøøøð#øøøs   ƒAA& Á&
BÁ0BÂBé
   ©Útimeout©r   N)ÚsuperÚ__init__r   r   r   r   Úadd_callbackÚresult)r$   r   r   r   r&   r#   Ú	__class__s   `    @€r%   r,   z!ThreadedZMQSocketChannel.__init__%   sž   øøø€ õ" 	‰Œ×ÒÑÔÐàˆŒØˆŒØˆŒÝ‘H”Hˆð	#ð 	#ð 	#ð 	#ð 	#ð 	#ð 	#ð Œ{Ð&Ð&Ð&ØŒ× Ò  Ñ.Ô.Ð.à	�Š˜ˆÑÔÐÐÐó    Fc                 ó   — | j         S )zWhether the channel is alive.©Ú	_is_alive©r$   s    r%   Úis_alivez!ThreadedZMQSocketChannel.is_aliveN   s
   € àŒ~Ðr0   c                 ó   — d| _         dS )zStart the channel.TNr2   r4   s    r%   ÚstartzThreadedZMQSocketChannel.startR   s   € àˆŒˆˆr0   c                 ó   — d| _         dS )zStop the channel.FNr2   r4   s    r%   ÚstopzThreadedZMQSocketChannel.stopV   s   € àˆŒˆˆr0   c                 óÌ  ‡ ‡— ‰ j         ��‰ j        �–t          ¦   «         Šdˆˆ fd„}‰ j                             |¦  «         	 ‰                     d¬¦  «         nO# t
          $ rB}t          ¦   «         }d‰ j         › d|› �}|                     |t          d¬	¦  «         Y d}~nd}~ww xY w‰ j	        �6	 ‰ j	         
                    d
¬¦  «         n# t
          $ r Y nw xY wd‰ _	        dS dS )zClose the channel.Nr   c                  óà   •— 	 ‰j         �"‰j                              d¬¦  «         d ‰_         ‰                     d ¦  «         d S # t          $ r } ‰                     | ¦  «         Y d } ~ d S d } ~ ww xY w)Nr   ©Úlinger)r   Úcloser   r   r    r!   s    €€r%   Úclose_streamz4ThreadedZMQSocketChannel.close.<locals>.close_stream`   s�   ø€ ð'Ø”{Ð.Øœ×)Ò)°Ð)Ñ3Ô3Ð3Ø&*˜œð —L’L Ñ&Ô&Ð&Ð&Ð&øõ !ð 'ð 'ð 'Ø—O’O AÑ&Ô&Ð&Ð&Ð&Ð&Ð&Ð&Ð&øøøøð'øøøs   ƒ)A Á
A-ÁA(Á(A-é   r(   zError closing stream z: é   )Ú
stacklevelr   r<   r*   )r   r   r   r-   r.   r   r
   ÚwarningÚRuntimeWarningr   r>   )r$   r?   r"   ÚlogÚmsgr#   s   `    @r%   r>   zThreadedZMQSocketChannel.closeZ   s7  øø€ àŒ;Ð" t¤{Ð'>å™œˆAð'ð 'ð 'ð 'ð 'ð 'ð 'ð ŒK×$Ò$ \Ñ2Ô2Ð2ð?Ø—’ �Ñ#Ô#Ð#Ð#øÝð ?ð ?ð ?Ý ‘l”l�Ø@¨d¬kÐ@Ð@¸QÐ@Ð@�Ø—’˜C¥¸A�Ñ>Ô>Ð>Ð>Ð>Ð>Ð>Ð>øøøøð?øøøð
 Œ;Ð"ðØ”×!Ò!¨Ð!Ñ+Ô+Ð+Ð+øÝð ð ð Ø�ðøøøàˆDŒKˆKˆKð #Ð"s*   ÁA Á
B$Á"8BÂB$Â/C Ã
CÃCrF   c                 ó^   ‡ ‡— dˆˆ fd„}‰ j         €J ‚‰ j                              |¦  «         dS )z÷Queue a message to be sent from the IOLoop's thread.

        Parameters
        ----------
        msg : message to send

        This is threadsafe, as it uses IOLoop.add_callback to give the loop's
        thread control of the action.
        r   Nc                  óZ   •— ‰j         €J ‚‰j                              ‰j        ‰ ¦  «         d S r   )r   Úsendr   )rF   r$   s   €€r%   Úthread_sendz2ThreadedZMQSocketChannel.send.<locals>.thread_send…   s1   ø€ Ø”<Ð+Ð+Ð+ØŒL×Ò˜dœk¨3Ñ/Ô/Ð/Ð/Ð/r0   r*   )r   r-   )r$   rF   rJ   s   `` r%   rI   zThreadedZMQSocketChannel.sendz   sS   øø€ ð	0ð 	0ð 	0ð 	0ð 	0ð 	0ð 	0ð Œ{Ð&Ð&Ð&ØŒ× Ò  Ñ-Ô-Ð-Ð-Ð-r0   Úmsg_listc                 óú   — | j         €J ‚| j        €J ‚| j                             |¦  «        \  }}| j                             |¦  «        }| j        r|                      |¦  «         |                      |¦  «         dS )z[Callback for stream.on_recv.

        Unpacks message, and calls handlers with it.
        N)r   r   Úfeed_identitiesÚdeserializeÚ_inspectÚcall_handlers)r$   rK   Ú_identÚsmsgrF   s        r%   r   z%ThreadedZMQSocketChannel._handle_recvŒ   s„   € ð
 Œ{Ð&Ð&Ð&ØŒ|Ð'Ð'Ð'Ø”|×3Ò3°HÑ=Ô=‰ˆ�ØŒl×&Ò& tÑ,Ô,ˆàŒ=ð 	Ø�MŠM˜#ÑÔÐØ×Ò˜3ÑÔÐÐÐr0   c                 ó   — dS )ai  This method is called in the ioloop thread when a message arrives.

        Subclasses should override this method to handle incoming messages.
        It is important to remember that this method is called in the thread
        so that some logic must be done to ensure that the application level
        handlers are called in the application thread.
        N© ©r$   rF   s     r%   rP   z&ThreadedZMQSocketChannel.call_handlersš   s	   € ð 	ˆr0   c                 ó   — dS )zaSubclasses should override this with a method
        processing any pending GUI events.
        NrT   r4   s    r%   Úprocess_eventsz'ThreadedZMQSocketChannel.process_events¤   s	   € ð 	ˆr0   ç      ð?r)   c                 ó2  ‡ — t          j        ¦   «         |z   }‰ j        €J ‚‰ j        �‰ j                             ¦   «         rd}t          |¦  «        ‚dt          ddfˆ fd„}t          d¦  «        D ]¦}t          ¦   «         }‰ j         	                    t          ||¦  «        ¦  «         t          |t          j        ¦   «         z
  d¦  «        }	 |                     t          |t          j        ¦   «         z
  d¦  «        ¦  «         Œ•# t          $ r Y  dS w xY wdS )a  Immediately processes all pending messages on this channel.

        This is only used for the IOPub channel.

        Callers should use this method to ensure that :meth:`call_handlers`
        has been called for all messages that have been received on the
        0MQ SUB socket of this channel.

        This method is thread safe.

        Parameters
        ----------
        timeout : float, optional
            The maximum amount of time to spend flushing, in seconds. The
            default is one second.
        NzAttempt to flush closed streamr#   r   c                 ó¶   •— 	 ‰                      ¦   «          |                      d ¦  «         d S # t          $ r }|                      |¦  «         Y d }~d S d }~ww xY wr   )Ú_flushr   r   r    )r#   r"   r$   s     €r%   Úflushz-ThreadedZMQSocketChannel.flush.<locals>.flushÄ   st   ø€ ð#Ø—’‘”�ð —’˜TÑ"Ô"Ð"Ð"Ð"øõ ð #ð #ð #Ø—’ Ñ"Ô"Ð"Ð"Ð"Ð"Ð"Ð"Ð"øøøøð#øøøs   ƒ. ®
A¸AÁArA   r   )ÚtimeÚ	monotonicr   r   ÚclosedÚOSErrorr   Úranger   r-   r   Úmaxr.   ÚTimeoutError)r$   r)   Ú	stop_timeÚ_msgr\   Ú_r#   s   `      r%   r\   zThreadedZMQSocketChannel.flushª   s8  ø€ õ& ”NÑ$Ô$ wÑ.ˆ	ØŒ{Ð&Ð&Ð&ØŒ;Ð $¤+×"4Ò"4Ñ"6Ô"6Ðà3ˆDÝ˜$‘-”-Ðð	#•Sð 	#˜Tð 	#ð 	#ð 	#ð 	#ð 	#ð 	#õ �q‘”ð 		ð 		ˆAÝ™œˆAØŒK×$Ò$¥W¨U°AÑ%6Ô%6Ñ7Ô7Ð7å˜)¥d¤nÑ&6Ô&6Ñ6¸Ñ:Ô:ˆGðØ—’�˜Y­¬Ñ)9Ô)9Ñ9¸1Ñ=Ô=Ñ>Ô>Ð>Ð>øÝð ð ð à���ðøøøð		ð 		s   Ã7DÄ
DÄDc                 óŠ   — | j         �| j                              ¦   «         rdS | j                              ¦   «          d| _        dS )z"Callback for :method:`self.flush`.NT)r   r_   r\   Ú_flushedr4   s    r%   r[   zThreadedZMQSocketChannel._flush×   sC   € ð Œ;Ð $¤+×"4Ò"4Ñ"6Ô"6ÐØˆFØŒ×ÒÑÔÐØˆŒˆˆr0   r*   )rX   ) Ú__name__Ú
__module__Ú__qualname__Ú__doc__r   r   r   r   rO   ÚzmqÚSocketr   r   r,   r3   Úboolr5   r7   r9   r>   ÚdictÚstrr   rI   Úlistr   rP   rW   Úfloatr\   r[   Ú__classcell__©r/   s   @r%   r   r      s·  ø€ € € € € Ø8Ð8à€GØ€FØ€FØ€FØ€Hð%à”
˜TÑ!ð%ð ˜4‘ð%ð �t‰mð	%ð
 
ð%ð %ð %ð %ð %ð %ðN €Ið˜$ð ð ð ð ðð ð ð ðð ð ð ðð ð ð ð@.˜˜S #˜Xœð .¨4ð .ð .ð .ð .ð$  Tð  ¨dð  ð  ð  ð  ð  c¨3 h¤ð °Dð ð ð ð ðð ð ð ð+ð +˜Uð +¨Tð +ð +ð +ð +ðZ
ð 
ð 
ð 
ð 
ð 
ð 
ð 
r0   r   c                   óŽ   ‡ — e Zd ZdZdZdZdˆ fd„Zeej	        dd„¦   «         ¦   «         Z
dd„Zdd„Zdd	„Zdd
„Zdd„Zdd„Zˆ xZS )ÚIOLoopThreadz;Run a pyzmq ioloop in a thread to send and receive messagesFNr   c                 ód   •— t          ¦   «                              ¦   «          d| _        d| _        dS )zInitialize an io loop thread.TFN)r+   r,   ÚdaemonÚ_exiting©r$   r/   s    €r%   r,   zIOLoopThread.__init__ê   s-   ø€ å‰Œ×ÒÑÔÐØˆŒð ˆŒˆˆr0   c                  ó0   — t           �dt           _        d S d S )NT)rw   rz   rT   r0   r%   Ú_notice_exitzIOLoopThread._notice_exitö   s    € õ
 Ð#Ø$(�LÔ!Ð!Ð!ð $Ð#r0   c                 óŠ   — t          ¦   «         | _        t          j        | ¦  «         | j                             d¬¦  «         dS )z{Start the IOLoop thread

        Don't return until self.ioloop is defined,
        which is created in the thread
        r'   r(   N)r   Ú_start_futurer   r7   r.   r4   s    r%   r7   zIOLoopThread.startþ   s@   € õ &,¡X¤XˆÔÝŒ�TÑÔÐàÔ×!Ò!¨"Ð!Ñ-Ô-Ð-Ð-Ð-r0   c                 óà  ‡ — 	 t          j        ¦   «         }t          j        |¦  «         dˆ fd„}|                      |¦   «         ¦  «         ‰ j                             d¦  «         n1# t          $ r$}‰ j                             |¦  «         Y d}~nd}~ww xY w	 |                     ‰                      ¦   «         ¦  «         | 	                    ¦   «          dS # | 	                    ¦   «          w xY w)z0Run my loop, ignoring EINTR events in the pollerr   Nc               “   ó<   •K  — t          j        ¦   «         ‰ _        d S r   )r   Úcurrentr   r4   s   €r%   Úassign_ioloopz'IOLoopThread.run.<locals>.assign_ioloop  s   øè è € Ý$œnÑ.Ô.�”��r0   r*   )
ÚasyncioÚnew_event_loopÚset_event_loopÚrun_until_completer   r   r   r    Ú
_async_runr>   )r$   r   rƒ   r"   s   `   r%   ÚrunzIOLoopThread.run	  s
  ø€ ð	0ÝÔ)Ñ+Ô+ˆDÝÔ" 4Ñ(Ô(Ð(ð/ð /ð /ð /ð /ð /ð ×#Ò# M M¡O¤OÑ4Ô4Ð4ð Ô×)Ò)¨$Ñ/Ô/Ð/Ð/øõ ð 	0ð 	0ð 	0ØÔ×,Ò,¨QÑ/Ô/Ð/Ð/Ð/Ð/Ð/Ð/øøøøð	0øøøð	Ø×#Ò# D§O¢OÑ$5Ô$5Ñ6Ô6Ð6à�JŠJ‰LŒLˆLˆLˆLøˆD�JŠJ‰LŒLˆLˆLøøøs$   ƒA
A( Á(
BÁ2BÂBÂ'C ÃC-c              ƒ   ó^   K  — | j         s#t          j        d¦  «        ƒ d{V —† | j         ¯!dS dS )z(Run forever (until self._exiting is set)r   N)rz   r„   Úsleepr4   s    r%   rˆ   zIOLoopThread._async_run  sR   è è € à”-ð 	#Ý”- Ñ"Ô"Ð"Ð"Ð"Ð"Ð"Ð"Ð"ð ”-ð 	#ð 	#ð 	#ð 	#ð 	#r0   c                 ór   — d| _         |                      ¦   «          |                      ¦   «          d| _        dS )zÿStop the channel's event loop and join its thread.

        This calls :meth:`~threading.Thread.join` and returns when the thread
        terminates. :class:`RuntimeError` will be raised if
        :meth:`~threading.Thread.start` is called again.
        TN)rz   Újoinr>   r   r4   s    r%   r9   zIOLoopThread.stop!  s0   € ð ˆŒØ�	Š	‰ŒˆØ�
Š
‰ŒˆØˆŒˆˆr0   c                 ó.   — |                       ¦   «          d S r   )r>   r4   s    r%   Ú__del__zIOLoopThread.__del__-  s   € Ø�
Š
‰Œˆˆˆr0   c                 ór   — | j         �/	 | j                              d¬¦  «         dS # t          $ r Y dS w xY wdS )zClose the io loop thread.NT)Úall_fds)r   r>   r   r4   s    r%   r>   zIOLoopThread.close0  sX   € àŒ;Ð"ðØ”×!Ò!¨$Ð!Ñ/Ô/Ð/Ð/Ð/øÝð ð ð Ø��ðøøøð #Ð"s   ‰& ¦
4³4r*   )ri   rj   rk   rl   rz   r   r,   ÚstaticmethodÚatexitÚregisterr}   r7   r‰   rˆ   r9   r�   r>   rt   ru   s   @r%   rw   rw   ä   sñ   ø€ € € € € ØEÐEà€HØ€Fð
ð 
ð 
ð 
ð 
ð 
ð Ø„_ð)ð )ð )ñ „_ñ „\ð)ð	.ð 	.ð 	.ð 	.ðð ð ð ð&#ð #ð #ð #ð

ð 
ð 
ð 
ðð ð ð ðð ð ð ð ð ð ð r0   rw   c                   ó*  ‡ — e Zd ZdZededz  fd„¦   «         Z eed¬¦  «        Z		 	 	 	 	 dde
de
d	e
d
e
de
ddfˆ fd„Zdeeef         ddfd„Zdˆ fd„Z ee¦  «        Z ee¦  «        Z ee¦  «        Z ee¦  «        Z ee¦  «        Zde
fd„Zˆ xZS )ÚThreadedKernelClientzYA KernelClient that provides thread-safe sockets with async callbacks on message replies.r   Nc                 ó,   — | j         r| j         j        S d S r   )Úioloop_threadr   r4   s    r%   r   zThreadedKernelClient.ioloop<  s   € àÔð 	-ØÔ%Ô,Ð,Øˆtr0   T)Ú
allow_noneÚshellÚiopubÚstdinÚhbÚcontrolc                 óÐ   •— t          ¦   «         | _        | j                             ¦   «          |r| j        | j        _        t          ¦   «                              |||||¦  «         dS )z!Start the channels on the client.N)rw   r˜   r7   Ú_check_kernel_info_replyÚshell_channelrO   r+   Ústart_channels)r$   rš   r›   rœ   r�   rž   r/   s         €r%   r¢   z#ThreadedKernelClient.start_channelsD  sc   ø€ õ *™^œ^ˆÔØÔ× Ò Ñ"Ô"Ð"àð 	HØ*.Ô*GˆDÔÔ'å‰Œ×Ò˜u e¨U°B¸Ñ@Ô@Ð@Ð@Ð@r0   rF   c                 ód   — |d         dk    r#|                       |¦  «         d| j        _        dS dS )zGThis is run in the ioloop thread when the kernel info reply is receivedÚmsg_typeÚkernel_info_replyN)Ú_handle_kernel_info_replyr¡   rO   rU   s     r%   r    z-ThreadedKernelClient._check_kernel_info_replyU  s?   € àˆzŒ?Ð1Ò1Ð1Ø×*Ò*¨3Ñ/Ô/Ð/Ø*.ˆDÔÔ'Ð'Ð'ð 2Ð1r0   c                 ó  •— | j         r™| j                              ¦   «         r€| j        �| j                             ¦   «          | j        �| j                             ¦   «          | j        �| j                             ¦   «          | j        �| j                             ¦   «          t          ¦   «                              ¦   «          | j         r4| j                              ¦   «         r| j          	                    ¦   «          dS dS dS )z Stop the channels on the client.N)
r˜   r5   Ú_shell_channelr>   Ú_iopub_channelÚ_stdin_channelÚ_control_channelr+   Ústop_channelsr9   r{   s    €r%   r¬   z"ThreadedKernelClient.stop_channels[  s  ø€ ð
 Ôð 	. $Ô"4×"=Ò"=Ñ"?Ô"?ð 	.ØÔ"Ð.ØÔ#×)Ò)Ñ+Ô+Ð+ØÔ"Ð.ØÔ#×)Ò)Ñ+Ô+Ð+ØÔ"Ð.ØÔ#×)Ò)Ñ+Ô+Ð+ØÔ$Ð0ØÔ%×+Ò+Ñ-Ô-Ð-å‰Œ×ÒÑÔÐØÔð 	& $Ô"4×"=Ò"=Ñ"?Ô"?ð 	&ØÔ×#Ò#Ñ%Ô%Ð%Ð%Ð%ð	&ð 	&ð 	&ð 	&r0   c                 óF   — | j         �| j                              ¦   «         S dS )z$Is the kernel process still running?NT)Ú_hb_channelÚ
is_beatingr4   s    r%   r5   zThreadedKernelClient.is_alivet  s)   € àÔÐ'ð Ô#×.Ò.Ñ0Ô0Ð0ð ˆtr0   )TTTTTr*   )ri   rj   rk   rl   Úpropertyr   r   r   rw   r˜   ro   r¢   rp   rq   r   r    r¬   r	   r   Úiopub_channel_classÚshell_channel_classÚstdin_channel_classr   Úhb_channel_classÚcontrol_channel_classr5   rt   ru   s   @r%   r–   r–   9  s–  ø€ € € € € ØcÐcàð˜ ™ð ð ð ñ „Xðð
 �H˜\°dÐ;Ñ;Ô;€Mð ØØØØðAð AàðAð ðAð ð	Að
 ðAð ðAð 
ðAð Að Að Að Að Að"/¨D°°c°¬Nð /¸tð /ð /ð /ð /ð&ð &ð &ð &ð &ð &ð& ˜$Ð7Ñ8Ô8ÐØ˜$Ð7Ñ8Ô8ÐØ˜$Ð7Ñ8Ô8ÐØ�t˜I‘”ÐØ ˜DÐ!9Ñ:Ô:Ðð˜$ð ð ð ð ð ð ð ð r0   r–   )rl   r„   r“   r]   Úconcurrent.futuresr   Ú	functoolsr   Ú	threadingr   Útypingr   rm   Útornado.ioloopr   Ú	traitletsr   r	   Útraitlets.logr
   Úzmq.eventloopr   Úchannelsr   Úclientr   r   r   r   rw   r–   rT   r0   r%   ú<module>rÀ      s£  ððð ð €€€Ø €€€Ø €€€Ø %Ð %Ð %Ð %Ð %Ð %Ø Ð Ð Ð Ð Ð Ø Ð Ð Ð Ð Ð Ø Ð Ð Ð Ð Ð à 
€
€
€
Ø !Ð !Ð !Ð !Ð !Ð !Ø $Ð $Ð $Ð $Ð $Ð $Ð $Ð $Ø $Ð $Ð $Ð $Ð $Ð $Ø #Ð #Ð #Ð #Ð #Ð #à Ð Ð Ð Ð Ð Ø  Ð  Ð  Ð  Ð  Ð  Ø Ð Ð Ð Ð Ð ðEð Eð Eð Eð Eñ Eô Eð EðPRð Rð Rð Rð R�6ñ Rô Rð RðjCð Cð Cð Cð C˜<ñ Cô Cð Cð Cð Cr0   