o
    ›¨Êh†  ã                   @   s’   d 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 ddlmZ dd	lmZ d
dlmZ dZdZeddƒZG dd„ de	ƒZdS )zEvent receiver implementation.é    N)Ú
itemgetter)ÚQueue)Úmaybe_channel)ÚConsumerMixin)Úuuid)Úapp_or_default)Úadjust_timestampé   )Úget_exchange)ÚEventReceiveréÿÿÿÿÚ	utcoffsetÚ	timestampc                   @   sŽ   e Zd ZdZdZ			ddd„Zdd„ Zdd	„ Z	
ddd„Zddd„Z	ddd„Z
ddd„Zd
ejeeefdd„Zeefdd„Zedd„ ƒZdS )r   a?  Capture events.

    Arguments:
        connection (kombu.Connection): Connection to the broker.
        handlers (Mapping[Callable]): Event handlers.
            This is  a map of event type names and their handlers.
            The special handler `"*"` captures all events that don't have a
            handler.
    Nú#c
           
   	   C   sú   t |p| jƒ| _t|ƒ| _|d u ri n|| _|| _|ptƒ | _|p%| jjj	| _
t| jp/| j ¡ | jjjd�| _|d u r@| jjj}|	d u rI| jjj}	td | j
| jg¡| j| jdd||	d�| _| jj| _| jj| _| jj| _|d u rx| jjjdh}|| _d S )N)ÚnameÚ.TF)ÚexchangeÚrouting_keyÚauto_deleteÚdurableÚmessage_ttlÚexpiresÚjson)r   Úappr   ÚchannelÚhandlersr   r   Únode_idÚconfÚevent_queue_prefixÚqueue_prefixr
   Ú
connectionÚconnection_for_writeÚevent_exchanger   Úevent_queue_ttlÚevent_queue_expiresr   ÚjoinÚqueueÚclockÚadjustÚadjust_clockÚforwardÚforward_clockÚevent_serializerÚaccept)
Úselfr   r   r   r   r   r   r-   Ú	queue_ttlÚqueue_expires© r1   úH/var/www/html/env/lib/python3.10/site-packages/celery/events/receiver.pyÚ__init__#   s8   
þ

ú



zEventReceiver.__init__c                 C   s.   | j  |¡p| j  d¡}|o||ƒ dS  dS )z3Process event by dispatching to configured handler.Ú*N)r   Úget)r.   ÚtypeÚeventÚhandlerr1   r1   r2   ÚprocessB   s   zEventReceiver.processc                 C   s   || j g| jgd| jd�gS )NT)ÚqueuesÚ	callbacksÚno_ackr-   )r&   Ú_receiver-   )r.   ÚConsumerr   r1   r1   r2   Úget_consumersG   s   þzEventReceiver.get_consumersTc                 K   s   |r
| j |d� d S d S )N)r   )Úwakeup_workers)r.   r    r   Ú	consumersÚwakeupÚkwargsr1   r1   r2   Úon_consume_readyL   s   ÿzEventReceiver.on_consume_readyc                 C   s   | j |||d�S )N©ÚlimitÚtimeoutrB   ©Úconsume)r.   rF   rG   rB   r1   r1   r2   ÚitercaptureQ   s   zEventReceiver.itercapturec                 C   s   | j |||d�D ]}qdS )zúOpen up a consumer capturing events.

        This has to run in the main process, and it will never stop
        unless :attr:`EventDispatcher.should_stop` is set to True, or
        forced via :exc:`KeyboardInterrupt` or :exc:`SystemExit`.
        rE   NrH   )r.   rF   rG   rB   Ú_r1   r1   r2   ÚcaptureT   s   ÿzEventReceiver.capturec                 C   s   | j jjd| j|d� d S )NÚ	heartbeat)r    r   )r   ÚcontrolÚ	broadcastr    )r.   r   r1   r1   r2   r@   ^   s   

þzEventReceiver.wakeup_workersc                 C   s²   |d }|dkr| j jpd|  }|d< |  |¡ nz|d }	W n ty/   |  ¡ |d< Y nw |  |	¡ |rPz||ƒ\}
}W n	 tyH   Y nw |||
ƒ|d< |ƒ |d< ||fS )Nr6   z	task-sentr	   r'   r   Úlocal_received)r'   Úvaluer)   ÚKeyErrorr+   )r.   ÚbodyÚlocalizeÚnowÚtzfieldsr   ÚCLIENT_CLOCK_SKEWr6   Ú_cr'   Úoffsetr   r1   r1   r2   Úevent_from_messagec   s&   ÿ
ÿ
z EventReceiver.event_from_messagec                    sD   |||ƒr| j | j‰‰ ‡ ‡fdd„|D ƒ d S | j |  |¡Ž  d S )Nc                    s   g | ]}ˆˆ |ƒŽ ‘qS r1   r1   )Ú.0r7   ©Úfrom_messager9   r1   r2   Ú
<listcomp>�   s    z*EventReceiver._receive.<locals>.<listcomp>)r9   rZ   )r.   rS   ÚmessageÚlistÚ
isinstancer1   r\   r2   r=   ~   s   
zEventReceiver._receivec                 C   s   | j r| j jjS d S ©N)r   r    Úclient)r.   r1   r1   r2   r    …   s   zEventReceiver.connection)Nr   NNNNNN)T)NNTrb   )Ú__name__Ú
__module__Ú__qualname__Ú__doc__r   r3   r9   r?   rD   rJ   rL   r@   ÚtimeÚ	_TZGETTERr   rW   rZ   r`   ra   r=   Úpropertyr    r1   r1   r1   r2   r      s,    

þ
ÿ




ýr   )rg   rh   Úoperatorr   Úkombur   Úkombu.connectionr   Úkombu.mixinsr   Úceleryr   Ú
celery.appr   Úcelery.utils.timer   r7   r
   Ú__all__rW   ri   r   r1   r1   r1   r2   Ú<module>   s    
