Ë
    UV.j  ã                  ó>  — d Z ddlmZ ddlZddlZddlm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 erddlmZ dZdefdefdefdefdefdœZd„ Zd„ Zd„ Z G d„ d«      Z G d„ d«      Z G d„ de«      Z ed e g d¢«      d¬«      Z! G d„ d «      Z"y)!zBase transport interface.é    )ÚannotationsN)ÚTYPE_CHECKING)ÚRecoverableConnectionError)ÚChannelErrorÚConnectionError)ÚMessage)Ú
dictfilter)Úcached_property)Úmaybe_s_to_ms)ÚTracebackType)r   Ú
StdChannelÚ
ManagementÚ	Transportz	x-expireszx-message-ttlzx-max-lengthzx-max-length-byteszx-max-priority)ÚexpiresÚmessage_ttlÚ
max_lengthÚmax_length_bytesÚmax_priorityc                ót   — t        t        d„ |j                  «       D «       «      «      }|rt        | fi |¤ŽS | S )a!  Convert queue arguments to RabbitMQ queue arguments.

    This is the implementation for Channel.prepare_queue_arguments
    for AMQP-based transports.  It's used by both the pyamqp and librabbitmq
    transports.

    Arguments:
        arguments (Mapping):
            User-supplied arguments (``Queue.queue_arguments``).

    Keyword Arguments:
        expires (float): Queue expiry time in seconds.
            This will be converted to ``x-expires`` in int milliseconds.
        message_ttl (float): Message TTL in seconds.
            This will be converted to ``x-message-ttl`` in int milliseconds.
        max_length (int): Max queue length (in number of messages).
            This will be converted to ``x-max-length`` int.
        max_length_bytes (int): Max queue size in bytes.
            This will be converted to ``x-max-length-bytes`` int.
        max_priority (int): Max priority steps for queue.
            This will be converted to ``x-max-priority`` int.

    Returns
    -------
        Dict: RabbitMQ compatible queue arguments.
    c              3  ó:   K  — | ]  \  }}t        ||«      –— Œ y ­w©N)Ú_to_rabbitmq_queue_argument)Ú.0ÚkeyÚvalues      úVC:\xampp\htdocs\tradingbinance\backend\.venv\Lib\site-packages\kombu/transport/base.pyÚ	<genexpr>z.to_rabbitmq_queue_arguments.<locals>.<genexpr>=   s#   è ø€ ð á)‰JˆC�ô 	$ C¨×/Ù)ùs   ‚)r	   ÚdictÚitems)Ú	argumentsÚoptionsÚprepareds      r   Úto_rabbitmq_queue_argumentsr#   !   sA   € ô8 œ$ñ à!Ÿ-™-œ/óó ó €Hñ +3Œ4�	Ñ&˜XÑ&ÐA¸	ÐAó    c                ó8   — t         |    \  }}||�	 ||«      fS |fS r   )ÚRABBITMQ_QUEUE_ARGUMENTS)r   r   ÚoptÚtyps       r   r   r   D   s,   € ä'¨Ñ,�H€CˆØ˜eÐ/‘�E“
Ð:Ð:°UÐ:Ð:r$   c                óL   — t        dj                  | j                  |«      «      S )Nz<Transport {0.__module__}.{0.__name__} does not implement {1})ÚNotImplementedErrorÚformatÚ	__class__)ÚobjÚmethods     r   Ú
_LeftBlankr/   J   s&   € ÜØF×MÑMØ�M‰M˜6ó	#ó$ð $r$   c                  óP   — e Zd ZdZdZd„ Zd„ Zd„ Zd„ Zd„ Z	d„ Z
	 	 	 	 	 	 	 	 d
d	„Zy)r   zStandard channel base class.Nc                ó&   — ddl m}  || g|¢­i |¤ŽS )Nr   )ÚConsumer)Úkombu.messagingr2   )ÚselfÚargsÚkwargsr2   s       r   r2   zStdChannel.ConsumerU   ó   € Ý,Ù˜Ð.˜tÒ. vÑ.Ð.r$   c                ó&   — ddl m}  || g|¢­i |¤ŽS )Nr   )ÚProducer)r3   r9   )r4   r5   r6   r9   s       r   r9   zStdChannel.ProducerY   r7   r$   c                ó   — t        | d«      ‚©NÚget_bindings©r/   ©r4   s    r   r<   zStdChannel.get_bindings]   ó   € Ü˜˜~Ó.Ð.r$   c                 ó   — y)zÄCallback called after RPC reply received.

        Notes
        -----
           Reply queue semantics: can be used to delete the queue
           after transient reply message received.
        N© )r4   Úqueues     r   Úafter_reply_message_receivedz'StdChannel.after_reply_message_received`   s   � r$   c                ó   — |S r   rA   )r4   r    r6   s      r   Úprepare_queue_argumentsz"StdChannel.prepare_queue_argumentsi   s   € ØÐr$   c                ó   — | S r   rA   r>   s    r   Ú	__enter__zStdChannel.__enter__l   s   € Øˆr$   c                ó$   — | j                  «        y r   )Úclose)r4   Úexc_typeÚexc_valÚexc_tbs       r   Ú__exit__zStdChannel.__exit__o   s   € ð 	�
‰
�r$   )rJ   ztype[BaseException] | NonerK   zBaseException | NonerL   zTracebackType | NoneÚreturnÚNone)Ú__name__Ú
__module__Ú__qualname__Ú__doc__Úno_ack_consumersr2   r9   r<   rC   rE   rG   rM   rA   r$   r   r   r   P   sT   „ Ù&àÐò/ò/ò/òòòðà,ðð &ðð %ð	ð
 
ôr$   r   c                  ó   — e Zd ZdZd„ Zd„ Zy)r   z!AMQP Management API (incomplete).c                ó   — || _         y r   )Ú	transport)r4   rW   s     r   Ú__init__zManagement.__init__{   s	   € Ø"ˆ�r$   c                ó   — t        | d«      ‚r;   r=   r>   s    r   r<   zManagement.get_bindings~   r?   r$   N)rP   rQ   rR   rS   rX   r<   rA   r$   r   r   r   x   s   „ Ù+ò#ó/r$   r   c                  ó"   — e Zd ZdZd„ Zd„ Zd„ Zy)Ú
Implementsz/Helper class used to define transport features.c                ó>   — 	 | |   S # t         $ r t        |«      ‚w xY wr   )ÚKeyErrorÚAttributeError)r4   r   s     r   Ú__getattr__zImplements.__getattr__…   s+   € ð	&Ø˜‘9ÐøÜò 	&Ü  Ó%Ð%ð	&ús   ‚ ‡c                ó   — || |<   y r   rA   )r4   r   r   s      r   Ú__setattr__zImplements.__setattr__‹   s   € ØˆˆSŠ	r$   c                ó(   —  | j                   | fi |¤ŽS r   )r,   )r4   r6   s     r   ÚextendzImplements.extendŽ   s   € Øˆt�~‰~˜dÑ- fÑ-Ð-r$   N)rP   rQ   rR   rS   r_   ra   rc   rA   r$   r   r[   r[   ‚   s   „ Ù9ò&òó.r$   r[   F)ÚdirectÚtopicÚfanoutÚheaders)ÚasynchronousÚexchange_typeÚ
heartbeatsc                  ó`  — e Zd ZdZeZdZdZdZefZ	e
fZdZdZdZej!                  «       Zd„ Zd„ Zd„ Zd„ Zd	„ Zd
„ Zdd„Zd„ Zd„ Zd„ Zd„ Zd„ Zej>                  ej@                  e!jD                  e!jF                  ffd„Z$d„ Z%d„ Z&ddd„Z'e(d„ «       Z)d„ Z*e+d„ «       Z,e(d„ «       Z-e(d„ «       Z.y)r   zBase class for transports.NFúN/Ac                ó   — || _         y r   )Úclient)r4   rn   r6   s      r   rX   zTransport.__init__º   s	   € Øˆ�r$   c                ó   — t        | d«      ‚)NÚestablish_connectionr=   r>   s    r   rp   zTransport.establish_connection½   s   € Ü˜Ð5Ó6Ð6r$   c                ó   — t        | d«      ‚)NÚclose_connectionr=   ©r4   Ú
connections     r   rr   zTransport.close_connectionÀ   s   € Ü˜Ð1Ó2Ð2r$   c                ó   — t        | d«      ‚)NÚcreate_channelr=   rs   s     r   rv   zTransport.create_channelÃ   s   € Ü˜Ð/Ó0Ð0r$   c                ó   — t        | d«      ‚)NÚclose_channelr=   rs   s     r   rx   zTransport.close_channelÆ   s   € Ü˜˜Ó/Ð/r$   c                ó   — t        | d«      ‚)NÚdrain_eventsr=   )r4   rt   r6   s      r   rz   zTransport.drain_eventsÉ   r?   r$   c                 ó   — y r   rA   )r4   rt   Úrates      r   Úheartbeat_checkzTransport.heartbeat_checkÌ   ó   € Ør$   c                 ó   — y)Nrl   rA   r>   s    r   Údriver_versionzTransport.driver_versionÏ   s   € Ør$   c                 ó   — y)Nr   rA   rs   s     r   Úget_heartbeat_intervalz Transport.get_heartbeat_intervalÒ   s   € Ør$   c                 ó   — y r   rA   ©r4   rt   Úloops      r   Úregister_with_event_loopz"Transport.register_with_event_loopÕ   r~   r$   c                 ó   — y r   rA   r„   s      r   Úunregister_from_event_loopz$Transport.unregister_from_event_loopØ   r~   r$   c                 ó   — y©NTrA   rs   s     r   Úverify_connectionzTransport.verify_connectionÛ   ó   € Ør$   c                ó>   ‡‡‡‡‡‡— ‰j                   Šˆˆˆˆˆˆfd„Š‰S )Nc                óº   •— ‰j                   st        d«      ‚	  ‰d¬«       | j                  ‰| «       y # ‰$ r Y y ‰$ r}|j                  ‰v rY d }~y ‚ d }~ww xY w)NzSocket was disconnectedr   )Útimeout)Ú	connectedr   ÚerrnoÚ	call_soon)r…   ÚexcÚ_readÚ_unavailrt   rz   Úerrorr�   s     €€€€€€r   r”   z%Transport._make_reader.<locals>._readâ   sd   ø€ Ø×'Ò'Ü0Ð1JÓKÐKðÙ QÕ'ð �N‰N˜5 $Õ'øð ò ÙØò Ø—9‘9 Ñ(ÜØûðús    š	6 ¶A½AÁAÁAÁA)rz   )r4   rt   r�   r–   r•   r”   rz   s    ````@@r   Ú_make_readerzTransport._make_readerÞ   s   ý€ à!×.Ñ.ˆ÷	(ñ 	(ð ˆr$   c                 ó   — yrŠ   rA   rs   s     r   Úqos_semantics_matches_specz$Transport.qos_semantics_matches_specñ   rŒ   r$   c                ó`   — | j                   }|€| j                  |«      x}| _          ||«       y r   )Ú_Transport__readerr—   )r4   rt   r…   Úreaders       r   Úon_readablezTransport.on_readableô   s.   € Ø—‘ˆØˆ>Ø%)×%6Ñ%6°zÓ%BÐBˆF�T”]Ùˆt�r$   c                ó   — t        «       ‚)z(Customise the display format of the URI.)r*   )r4   ÚuriÚinclude_passwordÚmasks       r   Úas_urizTransport.as_uriú   s   € ä!Ó#Ð#r$   c                ó   — i S r   rA   r>   s    r   Údefault_connection_paramsz#Transport.default_connection_paramsþ   s   € àˆ	r$   c                ó$   — | j                  | «      S r   )r   )r4   r5   r6   s      r   Úget_managerzTransport.get_manager  s   € Ø�‰˜tÓ$Ð$r$   c                ó"   — | j                  «       S r   )r¦   r>   s    r   ÚmanagerzTransport.manager  s   € à×ÑÓ!Ð!r$   c                ó.   — | j                   j                  S r   )Ú
implementsrj   r>   s    r   Úsupports_heartbeatszTransport.supports_heartbeats	  s   € à�‰×)Ñ)Ð)r$   c                ó.   — | j                   j                  S r   )rª   rh   r>   s    r   Úsupports_evzTransport.supports_ev  s   € à�‰×+Ñ+Ð+r$   )é   )Fz**)rŸ   ÚstrrN   r¯   )/rP   rQ   rR   rS   r   rn   Úcan_parse_urlÚdefault_portr   Úconnection_errorsr   Úchannel_errorsÚdriver_typeÚdriver_namer›   Údefault_transport_capabilitiesrc   rª   rX   rp   rr   rv   rx   rz   r}   r€   r‚   r†   rˆ   r‹   Úsocketr�   r–   r‘   ÚEAGAINÚEINTRr—   r™   r�   r¢   Úpropertyr¤   r¦   r
   r¨   r«   r­   rA   r$   r   r   r   ™   s  „ Ù$à€Jð €Fð €Mð €Lð )Ð*Ðð #�_€Nð
 €Kð €Kà€Hà/×6Ñ6Ó8€Jòò7ò3ò1ò0ò/óòòòòòð 06¯~©~Ø!Ÿ<™<°5·<±<ÀÇÁÐ2Móò&òô$ð ñó ðò%ð ñ"ó ð"ð ñ*ó ð*ð ñ,ó ñ,r$   r   )#rS   Ú
__future__r   r‘   r·   Útypingr   Úamqp.exceptionsr   Úkombu.exceptionsr   r   Úkombu.messager   Úkombu.utils.functionalr	   Úkombu.utils.objectsr
   Úkombu.utils.timer   Útypesr   Ú__all__Úintr&   r#   r   r/   r   r   r   r[   Ú	frozensetr¶   r   rA   r$   r   Ú<module>rÇ      s¿   ðÙ õ #ã Û Ý  å 6ç :Ý !Ý -Ý /Ý *áÝ#à
>€ð ˜]Ð+Ø# ]Ð3Ø! 3Ð'Ø-¨sÐ3Ø% sÐ+ñÐ ò BòF;ò$÷%ñ %÷P/ñ /ô.�ô .ñ  ",ØÙÒDÓEØô"Ð ÷v,ò v,r$   