o
    wvXj#6  ã                   @   s   d Z ddlmZm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mZmZmZmZ dd	lmZmZ d
dlmZmZ d
dlmZmZ d
dlmZ d
dlm Z! d
dl"mZ zddl#m$Z$ W n! e%y�   zddl&m$Z$ W n e%yŠ   ddl'm$Z$ Y nw Y nw dZ(dZ)ee*ƒZ+da,dd„ Z-dd„ Z.d=dd„Z/G dd„ deƒZ0dd„ Z1d>dd „Z2d!d"„ Z3d#d$„ Z4d%d&„ Z5d?d'd(„Z6		d?d)d*„Z7d@d+d,„Z8	dAd-d.„Z9d/d0„ Z:d1d2„ Z;e
d3d4„ ƒZ<dBd5d6„Z=dBd7d8„Z>dCd9d:„Z?G d;d<„ d<e@ƒZAdS )DzCommon Utilities.é    )Úabsolute_importÚunicode_literalsN)Údeque)Úcontextmanager)Úpartial)Úcount)Úuuid5Úuuid4Úuuid3ÚNAMESPACE_OID)ÚChannelErrorÚRecoverableConnectionErroré   )ÚExchangeÚQueue)Úbytes_if_py2Úrange)Ú
get_logger)Úregistry)Úuuid)Ú	get_ident)	Ú	BroadcastÚmaybe_declarer   ÚitermessagesÚ
send_replyÚcollect_repliesÚinsuredÚdrain_consumerÚ	eventloopiÿÿ  c                   C   s   t d u rtƒ ja t S ©N)Ú_node_idr	   Úint© r"   r"   úI/var/www/html/myproject/venv/lib/python3.10/site-packages/kombu/common.pyÚget_node_id+   s   r$   c                 C   sP   t d| ||t|ƒf ƒ}z
ttt|ƒƒ}W |S  ty'   ttt|ƒƒ}Y |S w )Nz%x-%x-%x-%x)r   ÚidÚstrr
   r   Ú
ValueErrorr   )Únode_idÚ
process_idÚ	thread_idÚinstanceÚentÚretr"   r"   r#   Úgenerate_oid2   s   ÿþþr.   Tc                 C   s"   t tƒ t ¡ |rtƒ | ƒS d| ƒS ©Nr   )r.   r$   ÚosÚgetpidr   )r+   Úthreadsr"   r"   r#   Úoid_from<   s   üür3   c                       s8   e Zd ZdZejd Z						d‡ fdd„	Z‡  ZS )	r   aˆ  Broadcast queue.

    Convenience class used to define broadcast queues.

    Every queue instance will have a unique name,
    and both the queue and exchange is configured with auto deletion.

    Arguments:
        name (str): This is used as the name of the exchange.
        queue (str): By default a unique id is used for the queue
            name for every consumer.  You can specify a custom
            queue name here.
        unique (bool): Always create a unique queue
            even if a queue name is supplied.
        **kwargs (Any): See :class:`~kombu.Queue` for a list
            of additional keyword arguments supported.
    ))ÚqueueNNFTc              
      sf   |rd  |pdtƒ ¡}n|pd  tƒ ¡}tt| ƒjd|p|||||d ur&|nt|dd�dœ|¤Ž d S )Nz{0}.{1}Úbcastz	bcast.{0}Úfanout)Útype)Úaliasr4   ÚnameÚauto_deleteÚexchanger"   )Úformatr   Úsuperr   Ú__init__r   )Úselfr9   r4   Úuniquer:   r;   r8   Úkwargs©Ú	__class__r"   r#   r>   Z   s   
ú
ùzBroadcast.__init__)NNFTNN)Ú__name__Ú
__module__Ú__qualname__Ú__doc__r   Úattrsr>   Ú__classcell__r"   r"   rB   r#   r   E   s    
úr   c                 C   s   | |j jjv S r   )Ú
connectionÚclientÚdeclared_entities)ÚentityÚchannelr"   r"   r#   Údeclaration_cachedq   s   rO   Fc                 K   s    |rt | |fi |¤ŽS t| |ƒS )zDeclare entity (cached).)Ú_imaybe_declareÚ_maybe_declare)rM   rN   ÚretryÚretry_policyr"   r"   r#   r   u   s   
r   c                 C   s0   | j }|s|std || ¡ƒ‚|  |¡} | S dS )zÓMake sure the channel is bound to the entity.

    :param entity: generic kombu nomenclature, generally an exchange or queue
    :param channel: channel to bind to the entity
    :return: the updated entity
    z#Cannot bind channel {} to entity {}N)Úis_boundr   r<   Úbind)rM   rN   rT   r"   r"   r#   Ú_ensure_channel_is_bound|   s   
ÿ
ûrV   c                 C   s¦   | }t | |ƒ |d u r| jstd | ¡ƒ‚| j}d  }}|jr1| jr1|jjj}t	| ƒ}||v r1dS |js8t
dƒ‚| j|d� |d urI|rI| |¡ |d urQ| j|_dS )Nz(channel is None and entity {} not bound.Fúchannel disconnected)rN   T)rV   rT   r   r<   rN   rJ   Úcan_cache_declarationrK   rL   Úhashr   ÚdeclareÚaddr9   )rM   rN   ÚorigÚdeclaredÚidentr"   r"   r#   rQ   Œ   s,   
ÿ

rQ   c                 K   s:   t | |ƒ | jjstdƒ‚| jjjj| tfi |¤Ž| |ƒS )NrW   )rV   rN   rJ   r   rK   ÚensurerQ   )rM   rN   rS   r"   r"   r#   rP   ©   s   

ÿÿÿrP   c              
   #   sŠ   � t ƒ ‰ ‡ fdd„}|g|pg  | _| �' t| jjj||dd�D ]}zˆ  ¡ V  W q  ty2   Y q w W d  ƒ dS 1 s>w   Y  dS )z&Drain messages from consumer instance.c                    s   ˆ   | |f¡ d S r   )Úappend)ÚbodyÚmessage©Úaccr"   r#   Ú
on_message·   s   z"drain_consumer.<locals>.on_messageT)ÚlimitÚtimeoutÚignore_timeoutsN)r   Ú	callbacksr   rN   rJ   rK   ÚpopleftÚ
IndexError)Úconsumerrf   rg   ri   re   Ú_r"   rc   r#   r   ³   s   €

ÿÿü"ÿr   c                 K   s$   t | jd|g|dœ|¤Ž|||d�S )zIterator over messages.)ÚqueuesrN   )rf   rg   ri   Nr"   )r   ÚConsumer)ÚconnrN   r4   rf   rg   ri   rA   r"   r"   r#   r   Å   s   þr   c              	   c   sN   � |rt |ƒp	tƒ D ]}z	| j|d�V  W q
 tjy$   |r"|s"‚ Y q
w dS )aè  Best practice generator wrapper around ``Connection.drain_events``.

    Able to drain events forever, with a limit, and optionally ignoring
    timeout errors (a timeout of 1 is often used in environments where
    the socket can get "stuck", and is a best practice for Kombu consumers).

    ``eventloop`` is a generator.

    Examples:
        >>> from kombu.common import eventloop

        >>> def run(conn):
        ...     it = eventloop(conn, timeout=1, ignore_timeouts=True)
        ...     next(it)   # one event consumed, or timed out.
        ...
        ...     for _ in eventloop(conn, timeout=1, ignore_timeouts=True):
        ...         pass  # loop forever.

    It also takes an optional limit parameter, and timeout errors
    are propagated by default::

        for _ in eventloop(connection, limit=1, timeout=1):
            pass

    See Also:
        :func:`itermessages`, which is an event loop bound to one or more
        consumers, that yields any messages received.
    )rg   N)r   r   Údrain_eventsÚsocketrg   )rp   rf   rg   rh   Úir"   r"   r#   r   Î   s   €€þýr   c              	   K   sH   |j |f| ||dœt|jd |j d¡tj|j |jdœfi |¤Ž¤ŽS )a³  Send reply for request.

    Arguments:
        exchange (kombu.Exchange, str): Reply exchange
        req (~kombu.Message): Original request, a message with
            a ``reply_to`` property.
        producer (kombu.Producer): Producer instance
        retry (bool): If true must retry according to
            the ``reply_policy`` argument.
        retry_policy (Dict): Retry settings.
        **props (Any): Extra properties.
    )r;   rR   rS   Úreply_toÚcorrelation_id)Úrouting_keyru   Ú
serializerÚcontent_encoding)ÚpublishÚdictÚ
propertiesÚgetÚserializersÚtype_to_nameÚcontent_typerx   )r;   ÚreqÚmsgÚproducerrR   rS   Úpropsr"   r"   r#   r   ó   s   ÿþ


ýýýr   c           	   	   o   s|   � |  dd¡}d}z*t| ||g|¢R i |¤ŽD ]\}}|s!| ¡  d}|V  qW |r2| |j¡ dS dS |r=| |j¡ w w )z,Generator collecting replies from ``queue``.Úno_ackTFN)Ú
setdefaultr   ÚackÚafter_reply_message_receivedr9   )	rp   rN   r4   ÚargsrA   r„   Úreceivedra   rb   r"   r"   r#   r     s&   €
ÿÿûÿÿr   c                 C   s   t jd| |dd� d S )Nz#Connection error: %r. Retry in %ss
T)Úexc_info)ÚloggerÚerror)ÚexcÚintervalr"   r"   r#   Ú_ensure_errback  s   
þr�   c              	   c   s,   � zd V  W d S  | j | j y   Y d S w r   )Úconnection_errorsÚchannel_errors)rp   r"   r"   r#   Ú_ignore_errors"  s   €ÿr’   c                 O   sB   |rt | ƒ� ||i |¤ŽW  d  ƒ S 1 sw   Y  t | ƒS )aÚ  Ignore connection and channel errors.

    The first argument must be a connection object, or any other object
    with ``connection_error`` and ``channel_error`` attributes.

    Can be used as a function:

    .. code-block:: python

        def example(connection):
            ignore_errors(connection, consumer.channel.close)

    or as a context manager:

    .. code-block:: python

        def example(connection):
            with ignore_errors(connection):
                consumer.channel.close()


    Note:
        Connection and channel errors should be properly handled,
        and not ignored.  Using this function is only acceptable in a cleanup
        phase, like when a connection is lost or at shutdown.
    N)r’   )rp   Úfunrˆ   rA   r"   r"   r#   Úignore_errors*  s
   
 ÿr”   c                 C   s   |r||ƒ d S d S r   r"   )rJ   rN   Ú	on_reviver"   r"   r#   Úrevive_connectionK  s   ÿr–   c                 K   s�   |pt }| jdd��4}|j|d� |j}tt||d�}	|j||f||	dœ|¤Ž}
|
|i t||d�¤Ž\}}|W  d  ƒ S 1 sAw   Y  dS )z›Function wrapper to handle connection errors.

    Ensures function performing broker commands completes
    despite intermittent connection failures.
    T)Úblock)Úerrback)r•   )r˜   r•   )rJ   N)r�   ÚacquireÚensure_connectionÚdefault_channelr   r–   Ú	autoretryrz   )Úpoolr“   rˆ   rA   r˜   r•   Úoptsrp   rN   Úreviver   Úretvalrm   r"   r"   r#   r   P  s   ÿÿ$÷r   c                   @   s@   e Zd ZdZdZdd„ Zddd„Zddd	„Zd
d„ Zdd„ Z	dS )ÚQoSaÞ  Thread safe increment/decrement of a channels prefetch_count.

    Arguments:
        callback (Callable): Function used to set new prefetch count,
            e.g. ``consumer.qos`` or ``channel.basic_qos``.  Will be called
            with a single ``prefetch_count`` keyword argument.
        initial_value (int): Initial prefetch count value..

    Example:
        >>> from kombu import Consumer, Connection
        >>> connection = Connection('amqp://')
        >>> consumer = Consumer(connection)
        >>> qos = QoS(consumer.qos, initial_prefetch_count=2)
        >>> qos.update()  # set initial

        >>> qos.value
        2

        >>> def in_some_thread():
        ...     qos.increment_eventually()

        >>> def in_some_other_thread():
        ...     qos.decrement_eventually()

        >>> while 1:
        ...    if qos.prev != qos.value:
        ...        qos.update()  # prefetch changed so update.

    It can be used with any function supporting a ``prefetch_count`` keyword
    argument::

        >>> channel = connection.channel()
        >>> QoS(channel.basic_qos, 10)


        >>> def set_qos(prefetch_count):
        ...     print('prefetch count now: %r' % (prefetch_count,))
        >>> QoS(set_qos, 10)
    Nc                 C   s   || _ t ¡ | _|pd| _d S r/   )ÚcallbackÚ	threadingÚRLockÚ_mutexÚvalue)r?   r¢   Úinitial_valuer"   r"   r#   r>   �  s   
zQoS.__init__r   c                 C   sZ   | j � | jr| jt|dƒ | _W d  ƒ | jS W d  ƒ | jS 1 s%w   Y  | jS )z¶Increment the value, but do not update the channels QoS.

        Note:
            The MainThread will be responsible for calling :meth:`update`
            when necessary.
        r   N)r¥   r¦   Úmax©r?   Únr"   r"   r#   Úincrement_eventually”  s   
þþ
ÿýzQoS.increment_eventuallyc                 C   sx   | j �. | jr|  j|8  _| jdk r(d| _W d  ƒ | jS W d  ƒ | jS W d  ƒ | jS 1 s4w   Y  | jS )z¶Decrement the value, but do not update the channels QoS.

        Note:
            The MainThread will be responsible for calling :meth:`update`
            when necessary.
        r   N)r¥   r¦   r©   r"   r"   r#   Údecrement_eventually   s   

üü
ÿþ
ýûzQoS.decrement_eventuallyc                 C   sH   || j kr"|}|tkrt dt¡ d}t d|¡ | j|d� || _ |S )z#Set channel prefetch_count setting.z(QoS: Disabled: prefetch_count exceeds %rr   zbasic.qos: prefetch_count->%s)Úprefetch_count)ÚprevÚPREFETCH_COUNT_MAXr‹   ÚwarningÚdebugr¢   )r?   ÚpcountÚ	new_valuer"   r"   r#   Úset®  s   
ÿzQoS.setc                 C   s6   | j � |  | j¡W  d  ƒ S 1 sw   Y  dS )z)Update prefetch count with current value.N)r¥   r´   r¦   )r?   r"   r"   r#   Úupdate»  s   
$ÿz
QoS.update)r   )
rD   rE   rF   rG   r®   r>   r«   r¬   r´   rµ   r"   r"   r"   r#   r¡   d  s    (

r¡   )T)NF)r   NN)NNF)NFNr   )NN)BrG   Ú
__future__r   r   r0   rr   r£   Úcollectionsr   Ú
contextlibr   Ú	functoolsr   Ú	itertoolsr   r   r   r	   r
   r   Úamqpr   r   rM   r   r   Úfiver   r   Úlogr   Úserializationr   r}   Ú
utils.uuidÚ_threadr   ÚImportErrorÚthreadÚdummy_threadÚ__all__r¯   rD   r‹   r    r$   r.   r3   r   rO   r   rV   rQ   rP   r   r   r   r   r   r�   r’   r”   r–   r   Úobjectr¡   r"   r"   r"   r#   Ú<module>   sl    ÿ€ý

	,



ÿ
	&
ÿ


!
