o
    wvXj¿  ã                   @   s¢   d Z ddlmZm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.é    )Úabsolute_importÚunicode_literalsN)Ú
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© r3   úS/var/www/html/myproject/venv/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)r0   ÚtypeÚeventÚhandlerr3   r3   r4   ÚprocessD   s   zEventReceiver.processc                 C   s   || j g| jgd| jd�gS )NT)ÚqueuesÚ	callbacksÚno_ackr/   )r(   Ú_receiver/   )r0   ÚConsumerr   r3   r3   r4   Úget_consumersI   s   þzEventReceiver.get_consumersTc                 K   s   |r
| j |d� d S d S )N)r   )Úwakeup_workers)r0   r"   r   Ú	consumersÚwakeupÚkwargsr3   r3   r4   Úon_consume_readyN   s   ÿzEventReceiver.on_consume_readyc                 C   s   | j |||d�S )N©ÚlimitÚtimeoutrD   ©Úconsume)r0   rH   rI   rD   r3   r3   r4   ÚitercaptureS   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`.
        rG   NrJ   )r0   rH   rI   rD   Ú_r3   r3   r4   ÚcaptureV   s   ÿzEventReceiver.capturec                 C   s   | j jjd| j|d� d S )NÚ	heartbeat)r"   r   )r   ÚcontrolÚ	broadcastr"   )r0   r   r3   r3   r4   rB   `   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 )Nr8   z	task-sentr   r)   r   Úlocal_received)r)   Úvaluer+   ÚKeyErrorr-   )r0   ÚbodyÚlocalizeÚnowÚtzfieldsr
   ÚCLIENT_CLOCK_SKEWr8   Ú_cr)   Úoffsetr   r3   r3   r4   Úevent_from_messagee   s&   ÿ
ÿ
z EventReceiver.event_from_messagec                    sD   |||ƒr| j | j‰‰ ‡ ‡fdd„|D ƒ d S | j |  |¡Ž  d S )Nc                    s   g | ]}ˆˆ |ƒŽ ‘qS r3   r3   )Ú.0r9   ©Úfrom_messager;   r3   r4   Ú
<listcomp>ƒ   s    z*EventReceiver._receive.<locals>.<listcomp>)r;   r\   )r0   rU   ÚmessageÚlistÚ
isinstancer3   r^   r4   r?   €   s   
zEventReceiver._receivec                 C   s   | j r| j jjS d S ©N)r   r"   Úclient)r0   r3   r3   r4   r"   ‡   s   zEventReceiver.connection)Nr   NNNNNN)T)NNTrd   )Ú__name__Ú
__module__Ú__qualname__Ú__doc__r   r5   r;   rA   rF   rL   rN   rB   ÚtimeÚ	_TZGETTERr
   rY   r\   rb   rc   r?   Úpropertyr"   r3   r3   r3   r4   r      s,    

þ
ÿ




ýr   )ri   Ú
__future__r   r   rj   Úoperatorr   Úkombur   Úkombu.connectionr   Úkombu.mixinsr   Úceleryr   Ú
celery.appr	   Úcelery.utils.timer
   r9   r   Ú__all__rY   rk   r   r3   r3   r3   r4   Ú<module>   s    
