o
    ›¨Êh-/  ã                   @   s¸   d 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 dd	lmZ dd
lmZmZ dZdZG dd„ deƒZdd„ ZG dd„ deƒZG dd„ dejeƒZdS )zqThe ``RPC`` result backend for AMQP brokers.

RPC-style result backend, using reply-to and one queue per client.
é    N)Úmaybe_declare)Úregister_after_fork)Úcached_property)Ústates)Úcurrent_taskÚtask_join_will_blocké   )Úbase)ÚAsyncBackendMixinÚBaseResultConsumer)ÚBacklogLimitExceededÚ
RPCBackendzñ
The "rpc" result backend does not support chords!

Note that a group chained with a task is also upgraded to be a chord,
as this pattern requires synchronization.

Result backends that supports chords: Redis, Database, Memcached, and more.
c                   @   s   e Zd ZdZdS )r   z'Too much state history to fast-forward.N)Ú__name__Ú
__module__Ú__qualname__Ú__doc__© r   r   úE/var/www/html/env/lib/python3.10/site-packages/celery/backends/rpc.pyr      s    r   c                 C   s   |   ¡  d S ©N)Ú_after_fork)Úbackendr   r   r   Ú_on_after_fork_cleanup_backend"   ó   r   c                       s^   e Zd ZejZdZdZ‡ fdd„Zddd„Zddd„Z	d	d
„ Z
dd„ Zdd„ Zdd„ Z‡  ZS )ÚResultConsumerNc                    s    t ƒ j|i |¤Ž | jj| _d S r   )ÚsuperÚ__init__r   Ú_create_binding©ÚselfÚargsÚkwargs©Ú	__class__r   r   r   ,   s   zResultConsumer.__init__Tc                 K   sF   | j  ¡ | _|  |¡}| j| jj|g| jg|| jd�| _| j 	¡  d S )N)Ú	callbacksÚno_ackÚaccept)
ÚappÚ
connectionÚ_connectionr   ÚConsumerÚdefault_channelÚon_state_changer%   Ú	_consumerÚconsume)r   Úinitial_task_idr$   r    Úinitial_queuer   r   r   Ústart0   s   

ýzResultConsumer.startc                 C   s*   | j r
| j j|d�S |rt |¡ d S d S )N)Útimeout)r(   Údrain_eventsÚtimeÚsleep)r   r1   r   r   r   r2   9   s
   ÿzResultConsumer.drain_eventsc                 C   s(   z| j  ¡  W | j ¡  d S | j ¡  w r   )r,   Úcancelr(   Úclose©r   r   r   r   Ústop?   s   zResultConsumer.stopc                 C   s(   d | _ | jd ur| j ¡  d | _d S d S r   )r,   r(   Úcollectr7   r   r   r   Úon_after_forkE   s
   


þzResultConsumer.on_after_forkc                 C   sH   | j d u r
|  |¡S |  |¡}| j  |¡s"| j  |¡ | j  ¡  d S d S r   )r,   r0   r   Úconsuming_fromÚ	add_queuer-   )r   Útask_idÚqueuer   r   r   Úconsume_fromK   s   


þzResultConsumer.consume_fromc                 C   s"   | j r| j  |  |¡j¡ d S d S r   )r,   Úcancel_by_queuer   Úname©r   r=   r   r   r   Ú
cancel_forS   s   ÿzResultConsumer.cancel_for©Tr   )r   r   r   Úkombur)   r(   r,   r   r0   r2   r8   r:   r?   rC   Ú__classcell__r   r   r!   r   r   &   s    

	r   c                       sb  e Zd ZdZejZejZeZeZdZ	dZ
dZdddddœZG dd	„ d	ejƒZG d
d„ dejƒZ		dE‡ fdd„	Zdd„ ZdFdd„Zdd„ Zdd„ Zdd„ Zdd„ Zdd„ Zdd „ ZdGd!d"„Z	dHd#d$„Zd%d&„ Zd'd(„ ZdId*d+„ZeZd,d-„ Z	dJd.d/„Zd0d1„ Z d2d3„ Z!d4d5„ Z"d6d7„ Z#d8d9„ Z$dGd:d;„Z%d<d=„ Z&dK‡ fd?d@„	Z'e(dAdB„ ƒZ)e*dCdD„ ƒZ+‡  Z,S )Lr   z&Base class for the RPC result backend.FTé   r   r   )Úmax_retriesÚinterval_startÚinterval_stepÚinterval_maxc                   @   ó   e Zd ZdZdZdS )zRPCBackend.Consumerz4Consumer that requires manual declaration of queues.FN)r   r   r   r   Úauto_declarer   r   r   r   r)   m   ó    r)   c                   @   rL   )zRPCBackend.Queuez$Queue that never caches declaration.FN)r   r   r   r   Úcan_cache_declarationr   r   r   r   ÚQueuer   rN   rP   Nc           
         s²   t ƒ j|fi |¤Ž | jj}	|| _i | _|  |¡| _| jrdnd| _|p&|	j	}|p+|	j
}|  ||| j¡| _|p9|	j| _|| _|  | | j| j| j| j¡| _td urWt| tƒ d S d S )Né   r   )r   r   r&   Úconfr(   Ú_out_of_bandÚprepare_persistentÚ
persistentÚdelivery_modeÚresult_exchangeÚresult_exchange_typeÚ_create_exchangeÚexchangeÚresult_serializerÚ
serializerÚauto_deleter   r%   Ú_pending_resultsÚ_pending_messagesÚresult_consumerr   r   )
r   r&   r'   rZ   Úexchange_typerU   r\   r]   r    rR   r!   r   r   r   w   s(   

ÿ
þÿzRPCBackend.__init__c                 C   s   | j  ¡  | j ¡  d S r   )r^   Úclearr`   r   r7   r   r   r   r   �   s   
zRPCBackend._after_forkÚdirectrQ   c                 C   s
   |   d ¡S r   )ÚExchange)r   rA   ÚtyperV   r   r   r   rY   ’   s   
zRPCBackend._create_exchangec                 C   s   | j S )z$Create new binding for task with id.)ÚbindingrB   r   r   r   r   –   s   zRPCBackend._create_bindingc                 C   s   t t ¡ ƒ‚r   )ÚNotImplementedErrorÚE_NO_CHORD_SUPPORTÚstripr7   r   r   r   Úensure_chords_allowed›   r   z RPCBackend.ensure_chords_allowedc                 C   s"   t ƒ st|  |j¡dd� d S d S )NT)Úretry)r   r   rf   Úchannel)r   Úproducerr=   r   r   r   Úon_task_callž   s   ÿzRPCBackend.on_task_callc                 C   s<   z|pt j}W n ty   td|›�ƒ‚w |j|jp|fS )z‹Get the destination for result by task id.

        Returns:
            Tuple[str, str]: tuple of ``(reply_to, correlation_id)``.
        z%RPC backend missing task request for )r   ÚrequestÚAttributeErrorÚRuntimeErrorÚreply_toÚcorrelation_id)r   r=   ro   r   r   r   Údestination_for¦   s   ÿÿzRPCBackend.destination_forc                 C   ó   d S r   r   rB   r   r   r   Úon_reply_declareµ   s   zRPCBackend.on_reply_declarec                 C   ru   r   r   )r   Úresultr   r   r   Úon_result_fulfilled»   s   zRPCBackend.on_result_fulfilledc                 C   s   dS )Nzrpc://r   )r   Úinclude_passwordr   r   r   Úas_uriÀ   ó   zRPCBackend.as_uric           
      K   sˆ   |   ||¡\}}|sdS | jjjjdd��%}	|	j|  |||||¡| j||| jd| j	|  
|¡| jd�	 W d  ƒ |S 1 s=w   Y  |S )z!Send task return value and state.NT©Úblock)rZ   Úrouting_keyrs   r\   rk   Úretry_policyÚdeclarerV   )rt   r&   ÚamqpÚproducer_poolÚacquireÚpublishÚ
_to_resultrZ   r\   r   rv   rV   )
r   r=   rw   ÚstateÚ	tracebackro   r    r~   rs   rm   r   r   r   Ústore_resultÃ   s$   ø
ÿõzRPCBackend.store_resultc                 C   s   |||   ||¡||  |¡dœS )N)r=   Ústatusrw   r‡   Úchildren)Úencode_resultÚcurrent_task_children)r   r=   r†   rw   r‡   ro   r   r   r   r…   Ö   s   
ûzRPCBackend._to_resultc                 C   s    | j r	| j  |¡ || j|< d S r   )r`   Úon_out_of_band_resultrS   )r   r=   Úmessager   r   r   r�   ß   s   z RPCBackend.on_out_of_band_resultéè  c           
      C   sØ   | j  |d ¡}|r|  ||¡S i }d }|  || j|¡D ]}|  |¡}| |¡|}||< |r4| ¡  d }q| |d ¡}| ¡ D ]
\}}	|  	||	¡ q?|rV| 
¡  |  ||¡S z| j| W S  tyk   tjd dœ Y S w )N)r‰   rw   )rS   ÚpopÚ_set_cache_by_messageÚ_slurp_from_queuer%   Ú_get_message_task_idÚgetÚackÚitemsr�   ÚrequeueÚ_cacheÚKeyErrorr   ÚPENDING)
r   r=   Úbacklog_limitÚbufferedÚlatest_by_idÚprevÚaccÚtidÚlatestÚmsgr   r   r   Úget_task_metaè   s.   
€þzRPCBackend.get_task_metac                 C   s   |   |j¡ }| j|< |S r   )Úmeta_from_decodedÚpayloadr˜   )r   r=   rŽ   r¥   r   r   r   r‘   	  s   ÿz RPCBackend._set_cache_by_messagec           	      c   s†   � | j jjdd��0\}}|  |¡|ƒ}| ¡  t|ƒD ]}|j||d�}|s( n	|V  q|  |¡‚W d   ƒ d S 1 s<w   Y  d S )NTr|   )r%   r$   )r&   ÚpoolÚacquire_channelr   r€   Úranger”   r   )	r   r=   r%   Úlimitr$   Ú_rl   rf   r¢   r   r   r   r’     s   €
ý"ùzRPCBackend._slurp_from_queuec              	   C   s.   z|j d W S  ttfy   |jd  Y S w )Nrs   r=   )Ú
propertiesrp   r™   r¥   )r   rŽ   r   r   r   r“     s
   þzRPCBackend._get_message_task_idc                 C   ru   r   r   )r   rl   r   r   r   Úrevive%  r{   zRPCBackend.revivec                 C   ó   t dƒ‚)Nz4reload_task_result is not supported by this backend.©rg   rB   r   r   r   Úreload_task_result(  ó   ÿzRPCBackend.reload_task_resultc                 C   r­   )z<Reload group result, even if it has been previously fetched.z5reload_group_result is not supported by this backend.r®   rB   r   r   r   Úreload_group_result,  s   ÿzRPCBackend.reload_group_resultc                 C   r­   )Nz,save_group is not supported by this backend.r®   )r   Úgroup_idrw   r   r   r   Ú
save_group1  r°   zRPCBackend.save_groupc                 C   r­   )Nz/restore_group is not supported by this backend.r®   )r   r²   Úcacher   r   r   Úrestore_group5  r°   zRPCBackend.restore_groupc                 C   r­   )Nz.delete_group is not supported by this backend.r®   )r   r²   r   r   r   Údelete_group9  r°   zRPCBackend.delete_groupr   c                    s@   |si n|}t ƒ  |t|| j| jj| jj| j| j| j	| j
d�¡S )N)r'   rZ   ra   rU   r\   r]   Úexpires)r   Ú
__reduce__Údictr(   rZ   rA   re   rU   r\   r]   r·   r   r!   r   r   r¸   =  s   
øzRPCBackend.__reduce__c                 C   s   | j | j| j| jdd| jd�S )NFT)Údurabler]   r·   )rP   ÚoidrZ   r·   r7   r   r   r   rf   J  s   üzRPCBackend.bindingc                 C   s   | j jS r   )r&   Ú
thread_oidr7   r   r   r   r»   S  s   zRPCBackend.oid)NNNNNT)rc   rQ   rD   )NN)r�   )r�   F)r   N)-r   r   r   r   rE   rd   ÚProducerr   r   rU   Úsupports_autoexpireÚsupports_native_joinr   r)   rP   r   r   rY   r   rj   rn   rt   rv   rx   rz   rˆ   r…   r�   r£   Úpollr‘   r’   r“   r¬   r¯   r±   r³   rµ   r¶   r¸   Úpropertyrf   r   r»   rF   r   r   r!   r   r   X   sb    üÿ


ÿ	
	
ÿ	

r   )r   r3   rE   Úkombu.commonr   Úkombu.utils.compatr   Úkombu.utils.objectsr   Úceleryr   Úcelery._stater   r   Ú r	   Úasynchronousr
   r   Ú__all__rh   Ú	Exceptionr   r   r   ÚBackendr   r   r   r   r   Ú<module>   s     
2