o
    —¨Êh  ã                   @  s  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rDddlmZ dZdefdefdefdefdefdœZdd„ Zdd„ Zdd„ ZG dd„ dƒZG dd„ dƒZG dd„ deƒZede g d ¢ƒdd!�Z!G d"d#„ d#ƒZ"dS )$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                 K  s2   t tdd„ | ¡ 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                 s  s   � | ]
\}}t ||ƒV  qd S ©N)Ú_to_rabbitmq_queue_argument)Ú.0ÚkeyÚvalue© r   úF/var/www/html/env/lib/python3.10/site-packages/kombu/transport/base.pyÚ	<genexpr>=   s
   € ÿ
ÿz.to_rabbitmq_queue_arguments.<locals>.<genexpr>)r   ÚdictÚitems)Ú	argumentsÚoptionsÚpreparedr   r   r   Úto_rabbitmq_queue_arguments!   s   

þr!   c                 C  s&   t |  \}}||d ur||ƒfS |fS r   )ÚRABBITMQ_QUEUE_ARGUMENTS)r   r   ÚoptÚtypr   r   r   r   D   s   r   c                 C  s   t d | j|¡ƒS )Nz<Transport {0.__module__}.{0.__name__} does not implement {1})ÚNotImplementedErrorÚformatÚ	__class__)ÚobjÚmethodr   r   r   Ú
_LeftBlankJ   s
   ÿÿr*   c                   @  sN   e Zd ZdZdZdd„ Zdd„ Zdd„ Zd	d
„ Zdd„ Z	dd„ Z
ddd„ZdS )r   zStandard channel base class.Nc                 O  ó"   ddl m} || g|¢R i |¤ŽS )Nr   )ÚConsumer)Úkombu.messagingr,   )ÚselfÚargsÚkwargsr,   r   r   r   r,   U   ó   zStdChannel.Consumerc                 O  r+   )Nr   )ÚProducer)r-   r2   )r.   r/   r0   r2   r   r   r   r2   Y   r1   zStdChannel.Producerc                 C  ó
   t | dƒ‚©NÚget_bindings©r*   ©r.   r   r   r   r5   ]   ó   
zStdChannel.get_bindingsc                 C  ó   dS )zÄCallback called after RPC reply received.

        Notes
        -----
           Reply queue semantics: can be used to delete the queue
           after transient reply message received.
        Nr   )r.   Úqueuer   r   r   Úafter_reply_message_received`   s    z'StdChannel.after_reply_message_receivedc                 K  s   |S r   r   )r.   r   r0   r   r   r   Úprepare_queue_argumentsi   ó   z"StdChannel.prepare_queue_argumentsc                 C  s   | S r   r   r7   r   r   r   Ú	__enter__l   r=   zStdChannel.__enter__Úexc_typeútype[BaseException] | NoneÚexc_valúBaseException | NoneÚexc_tbúTracebackType | NoneÚreturnÚNonec                 C  s   |   ¡  d S r   )Úclose)r.   r?   rA   rC   r   r   r   Ú__exit__o   s   zStdChannel.__exit__)r?   r@   rA   rB   rC   rD   rE   rF   )Ú__name__Ú
__module__Ú__qualname__Ú__doc__Úno_ack_consumersr,   r2   r5   r;   r<   r>   rH   r   r   r   r   r   P   s    	r   c                   @  s    e Zd ZdZdd„ Zdd„ ZdS )r   z!AMQP Management API (incomplete).c                 C  ó
   || _ d S r   )Ú	transport)r.   rO   r   r   r   Ú__init__{   r8   zManagement.__init__c                 C  r3   r4   r6   r7   r   r   r   r5   ~   r8   zManagement.get_bindingsN)rI   rJ   rK   rL   rP   r5   r   r   r   r   r   x   s    r   c                   @  s(   e Zd ZdZdd„ Zdd„ Zdd„ ZdS )	Ú
Implementsz/Helper class used to define transport features.c                 C  s"   z| | W S  t y   t|ƒ‚w r   )ÚKeyErrorÚAttributeError)r.   r   r   r   r   Ú__getattr__…   s
   
ÿzImplements.__getattr__c                 C  s   || |< d S r   r   )r.   r   r   r   r   r   Ú__setattr__‹   s   zImplements.__setattr__c                 K  s   | j | fi |¤ŽS r   )r'   )r.   r0   r   r   r   ÚextendŽ   s   zImplements.extendN)rI   rJ   rK   rL   rT   rU   rV   r   r   r   r   rQ   ‚   s
    rQ   F)ÚdirectÚtopicÚfanoutÚheaders)ÚasynchronousÚexchange_typeÚ
heartbeatsc                   @  s  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 ¡ Zdd„ Zdd„ Zd	d
„ Zdd„ Zdd„ Zdd„ Zd4dd„Zdd„ Zdd„ Zdd„ Zdd„ Zdd„ Zejej e!j"e!j#ffdd„Z$d d!„ Z%d"d#„ Z&d5d6d(d)„Z'e(d*d+„ ƒZ)d,d-„ Z*e+d.d/„ ƒZ,e(d0d1„ ƒZ-e(d2d3„ ƒZ.dS )7r   zBase class for transports.NFúN/Ac                 K  rN   r   )Úclient)r.   r_   r0   r   r   r   rP   º   r8   zTransport.__init__c                 C  r3   )NÚestablish_connectionr6   r7   r   r   r   r`   ½   r8   zTransport.establish_connectionc                 C  r3   )NÚclose_connectionr6   ©r.   Ú
connectionr   r   r   ra   À   r8   zTransport.close_connectionc                 C  r3   )NÚcreate_channelr6   rb   r   r   r   rd   Ã   r8   zTransport.create_channelc                 C  r3   )NÚclose_channelr6   rb   r   r   r   re   Æ   r8   zTransport.close_channelc                 K  r3   )NÚdrain_eventsr6   )r.   rc   r0   r   r   r   rf   É   r8   zTransport.drain_eventsé   c                 C  ó   d S r   r   )r.   rc   Úrater   r   r   Úheartbeat_checkÌ   r=   zTransport.heartbeat_checkc                 C  r9   )Nr^   r   r7   r   r   r   Údriver_versionÏ   r=   zTransport.driver_versionc                 C  r9   )Nr   r   rb   r   r   r   Úget_heartbeat_intervalÒ   r=   z Transport.get_heartbeat_intervalc                 C  rh   r   r   ©r.   rc   Úloopr   r   r   Úregister_with_event_loopÕ   r=   z"Transport.register_with_event_loopc                 C  rh   r   r   rm   r   r   r   Úunregister_from_event_loopØ   r=   z$Transport.unregister_from_event_loopc                 C  r9   ©NTr   rb   r   r   r   Úverify_connectionÛ   r=   zTransport.verify_connectionc                   s    ˆj ‰‡ ‡‡‡‡‡fdd„‰ ˆ S )Nc              
     sr   ˆj stdƒ‚zˆdd� W n" ˆy   Y d S  ˆy0 } z|jˆv r+W Y d }~d S ‚ d }~ww |  ˆ | ¡ d S )NzSocket was disconnectedr   )Útimeout)Ú	connectedr   ÚerrnoÚ	call_soon)rn   Úexc©Ú_readÚ_unavailrc   rf   Úerrorrs   r   r   ry   â   s   
€ýz%Transport._make_reader.<locals>._read)rf   )r.   rc   rs   r{   rz   r   rx   r   Ú_make_readerÞ   s   zTransport._make_readerc                 C  r9   rq   r   rb   r   r   r   Úqos_semantics_matches_specñ   r=   z$Transport.qos_semantics_matches_specc                 C  s*   | j }|d u r|  |¡ }| _ ||ƒ d S r   )Ú_Transport__readerr|   )r.   rc   rn   Úreaderr   r   r   Úon_readableô   s   zTransport.on_readableú**ÚuriÚstrrE   c                 C  s   t ƒ ‚)z(Customise the display format of the URI.)r%   )r.   r‚   Úinclude_passwordÚmaskr   r   r   Úas_uriú   s   zTransport.as_uric                 C  s   i S r   r   r7   r   r   r   Údefault_connection_paramsþ   s   z#Transport.default_connection_paramsc                 O  s
   |   | ¡S r   )r   )r.   r/   r0   r   r   r   Úget_manager  r8   zTransport.get_managerc                 C  s   |   ¡ S r   )rˆ   r7   r   r   r   Úmanager  ó   zTransport.managerc                 C  ó   | j jS r   )Ú
implementsr]   r7   r   r   r   Úsupports_heartbeats	  rŠ   zTransport.supports_heartbeatsc                 C  r‹   r   )rŒ   r[   r7   r   r   r   Úsupports_ev  rŠ   zTransport.supports_ev)rg   )Fr�   )r‚   rƒ   rE   rƒ   )/rI   rJ   rK   rL   r   r_   Úcan_parse_urlÚdefault_portr   Úconnection_errorsr   Úchannel_errorsÚdriver_typeÚdriver_namer~   Údefault_transport_capabilitiesrV   rŒ   rP   r`   ra   rd   re   rf   rj   rk   rl   ro   rp   rr   Úsocketrs   r{   ru   ÚEAGAINÚEINTRr|   r}   r€   r†   Úpropertyr‡   rˆ   r	   r‰   r�   rŽ   r   r   r   r   r   ™   sN    

ÿ


r   )#rL   Ú
__future__r   ru   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   rQ   Ú	frozensetr•   r   r   r   r   r   Ú<module>   s@    û	#(

ý