o
    wvXjÒ/  ã                   @   sØ   d Z ddlmZm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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.
é    )Úabsolute_importÚunicode_literalsN)Úmaybe_declare)Úregister_after_fork)Úcached_property)Ústates)Úcurrent_taskÚtask_join_will_block)ÚitemsÚrangeé   )Ú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   úP/var/www/html/myproject/venv/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 t| ƒj|i |¤Ž | jj| _d S r   )Úsuperr   Ú__init__r   Ú_create_binding©ÚselfÚargsÚkwargs©Ú	__class__r   r   r   0   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   Ústart4   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"   r5   r   r   r   r6   =   s
   ÿzResultConsumer.drain_eventsc                 C   s(   z| j  ¡  W | j ¡  d S | j ¡  w r   )r0   Úcancelr,   Úclose©r"   r   r   r   ÚstopC   s   zResultConsumer.stopc                 C   s(   d | _ | jd ur| j ¡  d | _d S d S r   )r0   r,   Úcollectr;   r   r   r   Úon_after_forkI   s
   


þzResultConsumer.on_after_forkc                 C   sH   | j d u r
|  |¡S |  |¡}| j  |¡s"| j  |¡ | j  ¡  d S d S r   )r0   r4   r    Úconsuming_fromÚ	add_queuer1   )r"   Útask_idÚqueuer   r   r   Úconsume_fromO   s   


þzResultConsumer.consume_fromc                 C   s"   | j r| j  |  |¡j¡ d S d S r   )r0   Úcancel_by_queuer    Úname©r"   rA   r   r   r   Ú
cancel_forW   s   ÿzResultConsumer.cancel_for©Tr   )r   r   r   Úkombur-   r,   r0   r   r4   r6   r<   r>   rC   rG   Ú__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-   q   ó    r-   c                   @   rP   )zRPCBackend.Queuez$Queue that never caches declaration.FN)r   r   r   r   Úcan_cache_declarationr   r   r   r   ÚQueuev   rR   rT   Nc           
         s¶   t t| ƒj|fi |¤Ž | jj}	|| _i | _|  |¡| _| jr!dnd| _	|p(|	j
}|p-|	j}|  ||| j	¡| _|p;|	j| _|| _|  | | j| j| j| j¡| _td urYt| tƒ d S d S )Né   r   )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+   r^   Úexchange_typerY   r`   ra   r$   rV   r%   r   r   r   {   s(   

ÿ
þÿzRPCBackend.__init__c                 C   s   | j  ¡  | j ¡  d S r   )rb   Úclearrd   r   r;   r   r   r   r   ‘   s   
zRPCBackend._after_forkÚdirectrU   c                 C   s
   |   d ¡S r   )ÚExchange)r"   rE   ÚtyperZ   r   r   r   r]   –   s   
zRPCBackend._create_exchangec                 C   s   | j S )z$Create new binding for task with id.)ÚbindingrF   r   r   r   r    š   s   zRPCBackend._create_bindingc                 C   s   t t ¡ ƒ‚r   )ÚNotImplementedErrorÚE_NO_CHORD_SUPPORTÚstripr;   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   rj   Úchannel)r"   ÚproducerrA   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 {0!r})r   ÚrequestÚAttributeErrorÚRuntimeErrorÚformatÚreply_toÚcorrelation_id)r"   rA   rs   r   r   r   Údestination_forª   s   ÿÿzRPCBackend.destination_forc                 C   ó   d S r   r   rF   r   r   r   Úon_reply_declare¹   s   zRPCBackend.on_reply_declarec                 C   rz   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)r^   Úrouting_keyrx   r`   ro   Úretry_policyÚdeclarerZ   )ry   r*   ÚamqpÚproducer_poolÚacquireÚpublishÚ
_to_resultr^   r`   r„   r{   rZ   )
r"   rA   r|   ÚstateÚ	tracebackrs   r$   rƒ   rx   rq   r   r   r   Ústore_resultÇ   s$   ø
ÿõzRPCBackend.store_resultc                 C   s   |||   ||¡||  |¡dœS )N)rA   Ústatusr|   rŒ   Úchildren)Úencode_resultÚcurrent_task_children)r"   rA   r‹   r|   rŒ   rs   r   r   r   rŠ   Ú   s   
ûzRPCBackend._to_resultc                 C   s    | j r	| j  |¡ || j|< d S r   )rd   Úon_out_of_band_resultrW   )r"   rA   Ú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 ¡}t|ƒD ]
\}}	|  	||	¡ q?|rV| 
¡  |  ||¡S z| j| W S  tyk   tjd dœ Y S w )N)rŽ   r|   )rW   ÚpopÚ_set_cache_by_messageÚ_slurp_from_queuer)   Ú_get_message_task_idÚgetÚackr
   r’   ÚrequeueÚ_cacheÚKeyErrorr   ÚPENDING)
r"   rA   Ú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"   rA   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…   r   r™   r   )	r"   rA   r)   Úlimitr(   Ú_rp   rj   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 )Nrx   rA   )Ú
propertiesrt   r�   r©   )r"   r“   r   r   r   r˜      s
   þzRPCBackend._get_message_task_idc                 C   rz   r   r   )r"   rp   r   r   r   Úrevive)  r€   zRPCBackend.revivec                 C   ó   t dƒ‚)Nz4reload_task_result is not supported by this backend.©rk   rF   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±   rF   r   r   r   Úreload_group_result0  s   ÿzRPCBackend.reload_group_resultc                 C   r°   )Nz,save_group is not supported by this backend.r±   )r"   Úgroup_idr|   r   r   r   Ú
save_group5  r³   zRPCBackend.save_groupc                 C   r°   )Nz/restore_group is not supported by this backend.r±   )r"   rµ   Úcacher   r   r   Úrestore_group9  r³   zRPCBackend.restore_groupc                 C   r°   )Nz.delete_group is not supported by this backend.r±   )r"   rµ   r   r   r   Údelete_group=  r³   zRPCBackend.delete_groupr   c                    sD   |si n|}t t| ƒ |t|| j| jj| jj| j| j	| j
| jd�¡S )N)r+   r^   re   rY   r`   ra   Úexpires)r   r   Ú
__reduce__Údictr,   r^   rE   ri   rY   r`   ra   rº   r!   r%   r   r   r»   A  s   øzRPCBackend.__reduce__c                 C   s   | j | j| j| jdd| jd�S )NFT)Údurablera   rº   )rT   Úoidr^   rº   r;   r   r   r   rj   N  s   üzRPCBackend.bindingc                 C   s   | j jS r   )r*   r¾   r;   r   r   r   r¾   W  s   zRPCBackend.oid)NNNNNT)rg   rU   rH   )NN)r”   )r”   F)r   N)-r   r   r   r   rI   rh   ÚProducerr   r   rY   Úsupports_autoexpireÚsupports_native_joinr„   r-   rT   r   r   r]   r    rn   rr   ry   r{   r}   r   r�   rŠ   r’   r§   Úpollr–   r—   r˜   r¯   r²   r´   r¶   r¸   r¹   r»   Úpropertyrj   r   r¾   rJ   r   r   r%   r   r   \   sb    üÿ


ÿ	
	
ÿ	

r   )!r   Ú
__future__r   r   r7   rI   Úkombu.commonr   Úkombu.utils.compatr   Úkombu.utils.objectsr   Úceleryr   Úcelery._stater   r	   Úcelery.fiver
   r   Ú r   Úasynchronousr   r   Ú__all__rl   Ú	Exceptionr   r   r   ÚBackendr   r   r   r   r   Ú<module>   s$   
2