o
    ›¨Êh#  ã                   @   s†   d 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mZmZ d
ZG dd„ dƒZdS )zEvent dispatcher sends events.é    N)ÚdefaultdictÚdeque)ÚProducer)Úapp_or_default)Ú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   r.   r   r   r   r"   r&   r,   r   r   r   Úconninfo© r<   úJ/var/www/html/env/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 r?   )Úclose)r:   Úexc_infor<   r<   r=   Ú__exit__a   s   zEventDispatcher.__exit__c                 C   s:   t | jp| j| j| jdd�| _d| _| jD ]}|ƒ  qd S )NF)r1   r"   Úauto_declareT)r   r   r   r1   r"   r   r.   r$   ©r:   Úcallbackr<   r<   r=   r5   d   s   ý
ÿzEventDispatcher.enablec                 C   s.   | j rd| _ |  ¡  | jD ]}|ƒ  qd S d S )NF)r.   rB   r%   rF   r<   r<   r=   Údisablem   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   r9   r+   Úrouting_keyú-Ú.)r+   Úforwardr   r   r9   r   Ú_publishÚreplace)	r:   ÚtypeÚfieldsr   Úblindr	   Úkwargsr+   Úeventr<   r<   r=   Úpublisht   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)rJ   r1   ÚretryÚretry_policyÚdeclarer"   r6   r,   )
r1   rU   r   r"   r6   r,   Ú	Exceptionr   r   Úappend)	r:   rT   r   rJ   rV   rW   r   r1   Úexcr<   r<   r=   rN   Ž   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.
        NrI   )rR   r	   rV   rW   )r.   r&   r   r   r+   rM   r   r9   r   rZ   Úlenr   Úflushr   rU   r   )r:   rP   rR   r   rV   rW   r	   rQ   r&   Úgroupr+   rT   Ú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�# | 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   rN   r   Úclearr   Úitems)	r:   Úerrorsr&   r_   rT   rJ   Ú_r^   Ú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   r@   r<   r<   r=   rB   Û   s   
ÿ
zEventDispatcher.closec                 C   s   | j S r?   ©r   r@   r<   r<   r=   Ú_get_publisherà   s   zEventDispatcher._get_publisherc                 C   s
   || _ d S r?   rk   )r:   r   r<   r<   r=   Ú_set_publisherã   s   
zEventDispatcher._set_publisher)NNTNTNNNr   Nr   N)TT)Ú__name__Ú
__module__Ú__qualname__Ú__doc__r4   r   r$   r%   r>   rA   rD   r5   rH   r	   rU   r   rN   r`   r]   rh   rB   rl   rm   ÚpropertyÚ	publisherr<   r<   r<   r=   r      s:    
ý$	
ÿ
ÿ
ÿ
%r   )rq   r7   r   r'   Úcollectionsr   r   Úkombur   Ú
celery.appr   Úcelery.utils.nodenamesr   Úcelery.utils.timer   rT   r	   r
   r   Ú__all__r   r<   r<   r<   r=   Ú<module>   s    