Ë
    šQjû7  ã                  óŒ  — d Z ddl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 d
dlmZ d
dlmZ dZdZ  ee!«      Z"da#d„ Z$d„ Z%d'd„Z& G d„ de«      Z'd„ Z(d(d„Z)d„ Z*d„ Z+d„ Z,d)d„Z-	 	 d)d„Z.d*d„Z/	 d+d„Z0d„ Z1d „ Z2e	d!„ «       Z3d,d"„Z4d,d#„Z5d-d$„Z6 G d%„ d&«      Z7y).zCommon Utilities.é    )ÚannotationsN)Údeque)Úcontextmanager)Úpartial)Úcount)ÚNAMESPACE_OIDÚuuid3Úuuid4Úuuid5)ÚChannelErrorÚRecoverableConnectionErroré   )ÚExchangeÚQueue)Ú
get_logger)Úregistry)Úuuid)	Ú	BroadcastÚmaybe_declarer   ÚitermessagesÚ
send_replyÚcollect_repliesÚinsuredÚdrain_consumerÚ	eventloopiÿÿ  c                 óB   — t         €t        «       j                  a t         S ©N)Ú_node_idr
   Úint© ó    úH/var/www/html/truck-me/venv/lib/python3.12/site-packages/kombu/common.pyÚget_node_idr#   "   s   € äÐÜ“7—;‘;ˆÜ€Or!   c                óÆ   — dj                  | ||t        |«      «      }	 t        t        t        |«      «      }|S # t
        $ r t        t        t        |«      «      }Y |S w xY w)Nz{:x}-{:x}-{:x}-{:x})ÚformatÚidÚstrr	   r   Ú
ValueErrorr   )Únode_idÚ
process_idÚ	thread_idÚinstanceÚentÚrets         r"   Úgenerate_oidr/   )   sc   € Ø
×
&Ñ
&Ø�˜Y¬¨8«ó6€Cð-Ü”%œ sÓ+Ó,ˆð €Jøô ò -Ü”%œ sÓ+Ó,‰Ø€Jð-ús   Ÿ: º"A ÁA c                óˆ   — t        t        «       t        j                  «       |rt	        j
                  «       | «      S d| «      S ©Nr   )r/   r#   ÚosÚgetpidÚ	threadingÚ	get_ident)r,   Úthreadss     r"   Úoid_fromr7   3   s@   € ÜÜ‹Ü
�	‰	‹Ù!(Œ	×ÑÓØó	ð ð /0Øó	ð r!   c                  óN   ‡ — e Zd ZdZej
                  dz   Z	 	 	 	 	 	 dˆ fd„	Zˆ x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.
    ))ÚqueueNc                óº   •— |rdj                  |xs dt        «       «      }n|xs dt        «       › �}t        ‰| �  d|xs |||||�|nt	        |d¬«      dœ|¤Ž y )Nz{}.{}Úbcastzbcast.Úfanout)Útype)Úaliasr9   ÚnameÚauto_deleteÚexchanger    )r%   r   ÚsuperÚ__init__r   )	Úselfr?   r9   Úuniquer@   rA   r>   ÚkwargsÚ	__class__s	           €r"   rC   zBroadcast.__init__R   sp   ø€ ñ Ø—N‘N 5Ò#3¨G´T³VÓ<‰EàÒ.˜v¤d£f XÐ.ˆEÜ‰Ñð 	
Ø’-˜4ØØØ#Ø"*Ð"6‘hÜ# D¨xÔ8ñ	
ð ó	
r!   )NNFTNN)Ú__name__Ú
__module__Ú__qualname__Ú__doc__r   ÚattrsrC   Ú__classcell__)rG   s   @r"   r   r   <   s7   ø„ ñð& �K‰KÐ,Ñ,€Eð ØØØ!ØØ÷
ñ 
r!   r   c                óF   — | |j                   j                  j                  v S r   )Ú
connectionÚclientÚdeclared_entities)ÚentityÚchannels     r"   Údeclaration_cachedrT   i   s   € Ø�W×'Ñ'×.Ñ.×@Ñ@Ð@Ð@r!   c                ó8   — |rt        | |fi |¤ŽS t        | |«      S )zDeclare entity (cached).)Ú_imaybe_declareÚ_maybe_declare)rR   rS   ÚretryÚretry_policys       r"   r   r   m   s$   € áÜ˜v wÑ?°,Ñ?Ð?Ü˜& 'Ó*Ð*r!   c                ój   — | j                   }|s$|st        d|› d| › �«      ‚| j                  |«      } | 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 z to entity )Úis_boundr   Úbind)rR   rS   r[   s      r"   Ú_ensure_channel_is_boundr]   t   sE   € ð �‰€HÙÙÜØ& w i¨{¸6¸(ÐCóEð Eà—‘˜WÓ%ˆØ€Mr!   c                óÄ  — | }t        | |«       |�|j                  €'| j                  st        d| › d�«      ‚| j                  }d x}}|j                  r<| j
                  r0|j                  j                  j                  }t        | «      }||v ry|j                  st        d«      ‚| j                  |¬«       |�|r|j                  |«       |�| j                  |_        y)Nzchannel is None and entity z not bound.Fúchannel disconnected)rS   T)r]   rO   r[   r   rS   Úcan_cache_declarationrP   rQ   Úhashr   ÚdeclareÚaddr?   )rR   rS   ÚorigÚdeclaredÚidents        r"   rW   rW   „   sß   € à€Dä˜V WÔ-à€˜'×,Ñ,Ð4ð �ŠÜØ-¨f¨X°[ÐAóCð Cà—.‘.ˆàÐ€HˆuØ×Ò˜f×:Ò:Ø×%Ñ%×,Ñ,×>Ñ>ˆÜ�V“ˆØ�HÑØà×ÒÜ(Ð)?Ó@Ð@Ø
‡N�N˜7€NÔ#ØÐ¡Ø�‰�UÔØÐØ—K‘KˆŒ	Ør!   c                óÖ   — t        | |«      } | j                  j                  st        d«      ‚  | j                  j                  j                  j
                  | t        fi |¤Ž| |«      S )Nr_   )r]   rS   rO   r   rP   ÚensurerW   )rR   rS   rY   s      r"   rV   rV   £   sj   € Ü% f¨gÓ6€Fà�>‰>×$Ò$Ü(Ð)?Ó@Ð@ð0Ð2ˆ6�>‰>×$Ñ$×+Ñ+×2Ñ2Ø”ñ0Ø".ñ0Ø06¸óAð Ar!   c              #  ó"  ‡K  — t        «       Šˆfd„}|g|xs g z   | _        | 5  t        | j                  j                  j
                  ||d¬«      D ]  }	 ‰j                  «       –— Œ 	 ddd«       y# t        $ r Y Œ-w xY w# 1 sw Y   yxY w­w)z&Drain messages from consumer instance.c                ó,   •— ‰j                  | |f«       y r   )Úappend)ÚbodyÚmessageÚaccs     €r"   Ú
on_messagez"drain_consumer.<locals>.on_message±   s   ø€ Ø�
‰
�D˜'�?Õ#r!   T)ÚlimitÚtimeoutÚignore_timeoutsN)r   Ú	callbacksr   rS   rO   rP   ÚpopleftÚ
IndexError)Úconsumerrp   rq   rs   ro   Ú_rn   s         @r"   r   r   ­   s™   øè ø€ ä
‹'€Cô$ð %˜¨ª°bÑ9€HÔà	ñ Ü˜8×+Ñ+×6Ñ6×=Ñ=Ø!&°ÈôOò 	ˆAðØ—k‘k“mÓ#ñ	÷ð øô
 ò Ùðú÷ð üs@   ƒ!B¤1BÁA4Á(BÁ+	BÁ4	B Á=BÁ?B Â BÂBÂBc                óH   — t         | j                  d|g|dœ|¤Ž|||¬«      S )zIterator over messages.)ÚqueuesrS   )rp   rq   rs   r    )r   ÚConsumer)ÚconnrS   r9   rp   rq   rs   rF   s          r"   r   r   ¿   s2   € ô Øˆ�‰Ð@˜e˜W¨gÑ@¸Ñ@Ø˜W°	ôð r!   c              #  ó²   K  — |xr t        |«      xs
 t        «       D ]  }	 | j                  |¬«      –— Œ y# t        j                  $ r |r|s‚ Y Œ5w xY w­w)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.
    )rq   N)Úranger   Údrain_eventsÚsocketrq   )r{   rp   rq   rr   Úis        r"   r   r   È   s]   è ø€ ð> Ò#”u˜U“|Ò.¤u£wò ˆð	Ø×#Ñ#¨GÐ#Ó4Ó4ñøô �~‰~ò 	Ù™Øùð	üs%   ‚A¢9¶A¹AÁAÁAÁAc                óä   —  |j                   |f| ||dœt        |j                  d   |j                  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.
    )rA   rX   rY   Úreply_toÚcorrelation_id)Úrouting_keyrƒ   Ú
serializerÚcontent_encoding)ÚpublishÚdictÚ
propertiesÚgetÚserializersÚtype_to_nameÚcontent_typer†   )rA   ÚreqÚmsgÚproducerrX   rY   Úpropss          r"   r   r   ï   s…   € ð ˆ8×ÑØðØØ ,ñô ˜sŸ~™~¨jÑ9Ø"%§.¡.×"4Ñ"4Ð5EÓ"FÜ)×6Ñ6°s×7GÑ7GÑHØ$'×$8Ñ$8ñ:ñ Dð >CñDñð r!   c              /  ó  K  — |j                  dd«      }d}	 t        | ||g|¢­i |¤ŽD ]  \  }}|s|j                  «        d}|–— Œ 	 |r|j                  |j                  «       yy# |r|j                  |j                  «       w w xY w­w)z,Generator collecting replies from ``queue``.Úno_ackTFN)Ú
setdefaultr   ÚackÚafter_reply_message_receivedr?   )	r{   rS   r9   ÚargsrF   r“   Úreceivedrl   rm   s	            r"   r   r     sœ   è ø€ à×Ñ˜x¨Ó.€FØ€Hð	=Ü)¨$°¸ð ;Ø+/ò;Ø39ñ;ò 	‰MˆD�'áØ—‘”ØˆHØ‹Jñ	ñ Ø×0Ñ0°·±Õ<ð ø‰8Ø×0Ñ0°·±Õ<ð üs   ‚B˜1A) Á
BÁ) B	Â	Bc                ó6   — t         j                  d| |d¬«       y )Nz#Connection error: %r. Retry in %ss
T)Úexc_info)ÚloggerÚerror)ÚexcÚintervals     r"   Ú_ensure_errbackrŸ     s   € Ü
‡L�LØ.°°XØð õ r!   c              #  óZ   K  — 	 d –— y # | j                   | j                  z   $ r Y y w xY w­wr   )Úconnection_errorsÚchannel_errors)r{   s    r"   Ú_ignore_errorsr£     s0   è ø€ ðÜøØ×!Ñ! D×$7Ñ$7Ñ7ò Ùðüs   ‚+„	 ˆ+‰(¥+§(¨+c                ó‚   — |rt        | «      5   ||i |¤Žcddd«       S t        | «      S # 1 sw Y   t        | «      S xY w)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£   )r{   Úfunr—   rF   s       r"   Úignore_errorsr¦   '  sG   € ñ8 Ü˜DÓ!ñ 	(Ù˜Ð' Ñ'÷	(ñ 	(ä˜$ÓÐ÷	(ä˜$ÓÐús   Ž+«>c                ó   — |r	 ||«       y y r   r    )rO   rS   Ú	on_revives      r"   Úrevive_connectionr©   I  s   € ÙÙ�'Õð r!   c           	     ó$  — |xs t         }| j                  d¬«      5 }|j                  |¬«       |j                  }t	        t
        ||¬«      }	 |j                  ||f||	dœ|¤Ž}
 |
|i t        ||¬«      ¤Ž\  }}|cddd«       S # 1 sw Y   yxY w)z›Function wrapper to handle connection errors.

    Ensures function performing broker commands completes
    despite intermittent connection failures.
    T)Úblock)Úerrback)r¨   )r¬   r¨   )rO   N)rŸ   ÚacquireÚensure_connectionÚdefault_channelr   r©   Ú	autoretryrˆ   )Úpoolr¥   r—   rF   r¬   r¨   Úoptsr{   rS   Úreviver   Úretvalrw   s                r"   r   r   N  s£   € ð Ò(œ€Gà	�‰˜DˆÓ	!ð 	 TØ×Ñ wÐÔ/ð ×&Ñ&ˆÜÔ*¨D¸IÔFˆØ �$—.‘.  gð ;°wØ+1ñ;Ø59ñ;ˆá˜TÐC¤T¨&¸TÔ%BÑC‰	ˆ�Ø÷	÷ 	ò 	ús   �ABÂBc                  ó8   — e Zd ZdZdZdd„Zd	d„Zd	d„Zd„ Zd„ Z	y)
Ú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..
        max_prefetch (int or None): Maximum allowed prefetch count. If specified
            as an integer, increment_eventually will not allow the value to exceed this limit.
            If None (the default), there is no upper limit on the prefetch count.

    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                óh   — || _         t        j                  «       | _        |xs d| _        || _        y r1   )Úcallbackr4   ÚRLockÚ_mutexÚvalueÚmax_prefetch)rD   r¸   Úinitial_valuer¼   s       r"   rC   zQoS.__init__’  s+   € Ø ˆŒÜ—o‘oÓ'ˆŒØ"Ò' aˆŒ
Ø(ˆÕr!   c                ó  — | j                   5  | j                  rG| j                  t        |d«      z   }| j                  �|| j                  kD  r| j                  }|| _        ddd«       | j                  S # 1 sw Y   | j                  S xY w)a  Increment the value, but do not update the channels QoS.

        Note:
        ----
            The MainThread will be responsible for calling :meth:`update`
            when necessary. If max_prefetch is set, the value will not
            exceed this limit.
        r   N)rº   r»   Úmaxr¼   )rD   ÚnÚ	new_values      r"   Úincrement_eventuallyzQoS.increment_eventually˜  sy   € ð �[‰[ñ 	'Ø�zŠzØ ŸJ™J¬¨Q°«Ñ2�	Ø×$Ñ$Ð0°YÀ×ARÑARÒ5RØ $× 1Ñ 1�IØ&�”
÷	'ð �z‰zÐ÷	'ð �z‰zÐús   �AA5Á5B	c                óà   — | j                   5  | j                  r+| xj                  |z  c_        | j                  dk  rd| _        ddd«       | j                  S # 1 sw Y   | j                  S xY w)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»   )rD   rÀ   s     r"   Údecrement_eventuallyzQoS.decrement_eventually©  sY   € ð �[‰[ñ 	#Ø�zŠzØ—
’
˜a‘•
Ø—:‘: ’>Ø!"�D”J÷		#ð
 �z‰zÐ÷	#ð
 �z‰zÐús   �8AÁA-c                óÐ   — || j                   k7  rV|}|t        kD  rt        j                  dt        «       d}t        j	                  d|«       | j                  |¬«       || _         |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¸   )rD   ÚpcountrÁ   s      r"   ÚsetzQoS.set¸  s\   € à�T—Y‘YÒØˆIØÔ*Ò*Ü—‘ÐIÜ1ô3à�	Ü�L‰LÐ8¸)ÔDØ�M‰M¨ˆMÔ3ØˆDŒIØˆr!   c                ó|   — | j                   5  | j                  | j                  «      cddd«       S # 1 sw Y   yxY w)z)Update prefetch count with current value.N)rº   rÌ   r»   )rD   s    r"   Úupdatez
QoS.updateÅ  s.   € à�[‰[ñ 	(Ø—8‘8˜DŸJ™JÓ'÷	(÷ 	(ò 	(ús   �2²;r   )r   )
rH   rI   rJ   rK   rÇ   rC   rÂ   rÄ   rÌ   rÎ   r    r!   r"   r¶   r¶   b  s(   „ ñ+ðZ €Dó)óó"òó(r!   r¶   )T)NF)r   NN)NNF)NFNr   )NN)8rK   Ú
__future__r   r2   r   r4   Úcollectionsr   Ú
contextlibr   Ú	functoolsr   Ú	itertoolsr   r   r   r	   r
   r   Úamqpr   r   rR   r   r   Úlogr   Úserializationr   r‹   Ú
utils.uuidÚ__all__rÈ   rH   r›   r   r#   r/   r7   r   rT   r   r]   rW   rV   r   r   r   r   r   rŸ   r£   r¦   r©   r   r¶   r    r!   r"   ú<module>rÙ      sæ   ðÙ å "ã 	Û Û Ý Ý %Ý Ý ß 3Ó 3ç 9ç #Ý Ý 2Ý ð€ð Ð á	�HÓ	€à€òòóô*
�ô *
òZAó+òò ò>Aóð$ 9=Øóó$ðP 9=óò2=ò ð ñó ðó óDó
÷(f(ò f(r!   