o
    wvXjy#  ã                   @   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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 d
dlmZmZmZ dZG dd„ deƒZdS )zEvent dispatcher sends events.é    )Úabsolute_importÚunicode_literalsN)ÚdefaultdictÚdeque)ÚProducer)Úapp_or_default)Úitems)Úanon_nodename)Ú	utcoffseté   )ÚEventÚget_exchangeÚ
group_from)ÚEventDispatcherc                   @   sº   e Zd ZdZdhZdZdZdZ				d"dd„Zd	d
„ Z	dd„ Z
dd„ Zdd„ Zdefdd„Zddefdd„Zdeddefdd„Zd#dd„Zdd„ Zdd„ Zdd„ Zd d!„ ZeeeƒZdS )$r   a0  Dispatches event messages.

    Arguments:
        connection (kombu.Connection): Connection to the broker.

        hostname (str): Hostname to identify ourselves as,
            by default uses the hostname returned by
            :func:`~celery.utils.anon_nodename`.

        groups (Sequence[str]): List of groups to send events for.
            :meth:`send` will ignore send requests to groups not in this list.
            If this is :const:`None`, all events will be sent.
            Example groups include ``"task"`` and ``"worker"``.

        enabled (bool): Set to :const:`False` to not actually publish any
            events, making :meth:`send` a no-op.

        channel (kombu.Channel): Can be used instead of `connection` to specify
            an exact channel to use when sending events.

        buffer_while_offline (bool): If enabled events will be buffered
            while the connection is down. :meth:`flush` must be called
            as soon as the connection is re-established.

    Note:
        You need to :meth:`close` this after use.
    ÚsqlNTr   é   c                 C   s0  t |p| jƒ| _|| _|| _|ptƒ | _|| _|
ptƒ | _|| _	|| _
ttƒ| _t ¡ | _d | _tƒ | _|p:| jjj| _tƒ | _tƒ | _t|pHg ƒ| _tj tj g| _| jj| _|	| _ |se|re|jj!| _|| _"| jpo| j #¡ }t$|| jjj%d�| _&|j'j(| j)v r„d| _"| j"r‹|  *¡  d| ji| _+t, -¡ | _.d S )N)ÚnameFÚhostname)/r   ÚappÚ
connectionÚchannelr	   r   Úbuffer_while_offlineÚ	frozensetÚbuffer_groupÚbuffer_limitÚon_send_bufferedr   ÚlistÚ_group_bufferÚ	threadingÚLockÚmutexÚproducerr   Ú_outbound_bufferÚconfÚevent_serializerÚ
serializerÚsetÚ
on_enabledÚon_disabledÚgroupsÚtimeÚtimezoneÚaltzoneÚtzoffsetÚclockÚdelivery_modeÚclientÚenabledÚconnection_for_writer   Úevent_exchangeÚexchangeÚ	transportÚdriver_typeÚDISABLED_TRANSPORTSÚenableÚheadersÚosÚgetpidÚpid)Úselfr   r   r1   r   r   r   r%   r)   r/   r   r   r   Úconninfo© r?   úU/var/www/html/myproject/venv/lib/python3.10/site-packages/celery/events/dispatcher.pyÚ__init__<   s@   



ÿzEventDispatcher.__init__c                 C   s   | S ©Nr?   ©r=   r?   r?   r@   Ú	__enter__`   s   zEventDispatcher.__enter__c                 G   s   |   ¡  d S rB   )Úclose)r=   Úexc_infor?   r?   r@   Ú__exit__c   s   zEventDispatcher.__exit__c                 C   s:   t | jp| j| j| jdd�| _d| _| jD ]}|ƒ  qd S )NF)r4   r%   Úauto_declareT)r   r   r   r4   r%   r!   r1   r'   ©r=   Úcallbackr?   r?   r@   r8   f   s   ý
ÿzEventDispatcher.enablec                 C   s.   | j rd| _ |  ¡  | jD ]}|ƒ  qd S d S )NF)r1   rE   r(   rI   r?   r?   r@   Údisableo   s   
üzEventDispatcher.disableFc           	      K   s|   |rdn| j  ¡ }||f| jtƒ | j|dœ|¤Ž}| j� | j||fd| dd¡i|¤ŽW  d  ƒ S 1 s7w   Y  dS )au  Publish event using custom :class:`~kombu.Producer`.

        Arguments:
            type (str): Event type name, with group separated by dash (`-`).
                fields: Dictionary of event fields, must be json serializable.
            producer (kombu.Producer): Producer instance to use:
                only the ``publish`` method will be called.
            retry (bool): Retry in the event of connection failure.
            retry_policy (Mapping): Map of custom retry policy options.
                See :meth:`~kombu.Connection.ensure`.
            blind (bool): Don't set logical clock value (also don't forward
                the internal logical clock).
            Event (Callable): Event type used to create event.
                Defaults to :func:`Event`.
            utcoffset (Callable): Function returning the current
                utc offset in hours.
        N©r   r
   r<   r.   Úrouting_keyú-Ú.)r.   Úforwardr   r
   r<   r    Ú_publishÚreplace)	r=   ÚtypeÚfieldsr!   Úblindr   Úkwargsr.   Úeventr?   r?   r@   Úpublishv   s   ÿÿ
ÿÿ$ÿzEventDispatcher.publishc           	      C   st   | j }z|j|||j|||g| j| j| jd�	 W d S  ty9 } z| js%‚ | j 	|||f¡ W Y d }~d S d }~ww )N)rM   r4   ÚretryÚretry_policyÚdeclarer%   r9   r/   )
r4   rX   r   r%   r9   r/   Ú	Exceptionr   r"   Úappend)	r=   rW   r!   rM   rY   rZ   r
   r4   Úexcr?   r?   r@   rQ   �   s&   ÷ €ýzEventDispatcher._publishc              	   K   s¼   | j r\| jt|ƒ}}	|r|	|vrdS |	| jv rO| j ¡ }
||f| j|ƒ | j|
dœ|¤Ž}| j|	 }| 	|¡ t
|ƒ| jkrD|  ¡  dS | jrM|  ¡  dS dS | j||| j||||d�S dS )aÆ  Send event.

        Arguments:
            type (str): Event type name, with group separated by dash (`-`).
            retry (bool): Retry in the event of connection failure.
            retry_policy (Mapping): Map of custom retry policy options.
                See :meth:`~kombu.Connection.ensure`.
            blind (bool): Don't set logical clock value (also don't forward
                the internal logical clock).
            Event (Callable): Event type used to create event,
                defaults to :func:`Event`.
            utcoffset (Callable): unction returning the current utc offset
                in hours.
            **fields (Any): Event fields -- must be json serializable.
        NrL   )rU   r   rY   rZ   )r1   r)   r   r   r.   rP   r   r<   r   r]   Úlenr   Úflushr   rX   r!   )r=   rS   rU   r
   rY   rZ   r   rT   r)   Úgroupr.   rW   Úbufr?   r?   r@   Úsend¤   s0   


þþ

ÿþðzEventDispatcher.sendc           	      C   sØ   |r8t | jƒ}z*| j� |D ]\}}}|  || j|¡ qW d  ƒ n1 s&w   Y  W | j ¡  n| j ¡  w |rj| j�# t| jƒD ]\}}|  || jd| ¡ g |dd…< qCW d  ƒ dS 1 scw   Y  dS dS )zFlush the outbound buffer.Nz%s.multi)r   r"   r    rQ   r!   Úclearr   r   )	r=   Úerrorsr)   rb   rW   rM   Ú_ra   Úeventsr?   r?   r@   r`   É   s$   
ÿÿ€þ"ÿÿzEventDispatcher.flushc                 C   s   | j  |j ¡ dS )z-Copy the outbound buffer of another instance.N)r"   Úextend)r=   Úotherr?   r?   r@   Úextend_bufferÙ   s   zEventDispatcher.extend_bufferc                 C   s*   | j  ¡ o| j  ¡  d| _dS  d| _dS )zClose the event dispatcher.N)r    ÚlockedÚreleaser!   rC   r?   r?   r@   rE   Ý   s   
ÿ
zEventDispatcher.closec                 C   s   | j S rB   ©r!   rC   r?   r?   r@   Ú_get_publisherâ   s   zEventDispatcher._get_publisherc                 C   s
   || _ d S rB   rm   )r=   r!   r?   r?   r@   Ú_set_publisherå   s   
zEventDispatcher._set_publisher)NNTNTNNNr   Nr   N)TT)Ú__name__Ú
__module__Ú__qualname__Ú__doc__r7   r   r'   r(   rA   rD   rG   r8   rK   r   rX   r
   rQ   rc   r`   rj   rE   rn   ro   ÚpropertyÚ	publisherr?   r?   r?   r@   r      s:    
ý$	
ÿ
ÿ
ÿ
%r   )rs   Ú
__future__r   r   r:   r   r*   Úcollectionsr   r   Úkombur   Ú
celery.appr   Úcelery.fiver   Úcelery.utils.nodenamesr	   Úcelery.utils.timer
   rW   r   r   r   Ú__all__Úobjectr   r?   r?   r?   r@   Ú<module>   s    