o
    wvXjÑQ  ã                   @   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
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 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" ddl#m$Z$ ddl%m&Z& ddl'm(Z(m)Z) ddl*m+Z+ zddl,m-Z- W n e.y•   ddl/m-Z- Y nw zddl0Z0ddl1Z0ddl2m3Z3 W n e.y³   dZ0dZ3Y nw zddl0m4Z4 W n e.yÇ   dZ4Y nw dZ5dZ6dZ7dZ8dZ9dZ:dZ;d Z<d!Z=e$e>ƒZ?G d"d#„ d#e)ƒZ@G d$d%„ d%e+e(ƒZAG d&d'„ d'eAƒZBdS )(zRedis result store backend.é    )Úabsolute_importÚunicode_literalsN)Úcontextmanager)Úpartial)Ú	CERT_NONEÚCERT_OPTIONALÚCERT_REQUIRED)Úretry_over_time)Úcached_property)Ú
_parse_url)Ústates)Útask_join_will_block)Úmaybe_signature)Ú
ChordErrorÚImproperlyConfigured)Ústring_tÚtext_t)Ú
deprecated)Ú
dictfilter)Ú
get_logger)Úhumanize_secondsé   )ÚAsyncBackendMixinÚBaseResultConsumer)ÚBaseKeyValueStoreBackend)Úunquote)Úget_redis_error_classes)Úsentinel)ÚRedisBackendÚSentinelBackendzW
You need to install the redis library in order to use the Redis result store backend.
zp
You need to install the redis library with support of sentinel in order to use the Redis result store backend.
zÍ
Setting ssl_cert_reqs=CERT_OPTIONAL when connecting to redis means that celery might not valdate the identity of the redis broker when connecting. This leaves you vulnerable to man in the middle attacks.
zÈ
Setting ssl_cert_reqs=CERT_NONE when connecting to redis means that celery will not valdate the identity of the redis broker when connecting. This leaves you vulnerable to man in the middle attacks.
z”
SSL connection parameters have been provided but the specified URL scheme is redis://. A Redis SSL connection URL should use the scheme rediss://.
zv
A rediss:// URL must have parameter ssl_cert_reqs and this must be set to CERT_REQUIRED, CERT_OPTIONAL, or CERT_NONE
z+Connection to Redis lost: Retry (%s/%s) %s.z„
Retry limit exceeded while trying to reconnect to the Celery redis result store backend. The Celery application must be restarted.
c                       sŽ   e Zd ZdZ‡ fdd„Z‡ fdd„Zdd„ Zedd	„ ƒZd
d„ Z	‡ fdd„Z
dd„ Zdd„ Zdd„ Zddd„Zdd„ Zdd„ Zdd„ Z‡  ZS )ÚResultConsumerNc                    sJ   t t| ƒj|i |¤Ž | jj| _| jj| _| jj| _	| jj
| _tƒ | _d S ©N)Úsuperr    Ú__init__ÚbackendÚget_key_for_taskÚ_get_key_for_taskÚdecode_resultÚ_decode_resultÚensureÚ_ensureÚconnection_errorsÚ_connection_errorsÚsetÚsubscribed_to©ÚselfÚargsÚkwargs©Ú	__class__© úR/var/www/html/myproject/venv/lib/python3.10/site-packages/celery/backends/redis.pyr#   ]   s   



zResultConsumer.__init__c              
      sl   z| j jj ¡  | jd ur| j ¡  W n ty, } zt t	|ƒ¡ W Y d }~nd }~ww t
t| ƒ ¡  d S r!   )r$   ÚclientÚconnection_poolÚresetÚ_pubsubÚcloseÚKeyErrorÚloggerÚwarningr   r"   r    Úon_after_fork)r0   Úer3   r5   r6   r?   e   s   

€€ÿzResultConsumer.on_after_forkc                 C   sr   d | _ | jjj ¡  | jj | j¡}dd„ |D ƒ}|D ]}|  |  |¡d ¡ q| jjj	dd�| _ | j j
| jŽ  d S )Nc                 S   s   g | ]}|r|‘qS r5   r5   )Ú.0Úmetar5   r5   r6   Ú
<listcomp>t   s    z4ResultConsumer._reconnect_pubsub.<locals>.<listcomp>T©Úignore_subscribe_messages)r:   r$   r7   r8   r9   Úmgetr.   Úon_state_changer(   ÚpubsubÚ	subscribe)r0   ÚmetasrB   r5   r5   r6   Ú_reconnect_pubsubn   s   ÿz ResultConsumer._reconnect_pubsubc                 c   sT   � zd V  W d S  | j y)   z|  | jd¡ W Y d S  | j y(   t t¡ ‚ w w ©Nr5   )r,   r*   rK   r=   ÚcriticalÚE_RETRY_LIMIT_EXCEEDED©r0   r5   r5   r6   Úreconnect_on_error|   s   €
þýz!ResultConsumer.reconnect_on_errorc                 C   s$   |d t jv r|  |d ¡ d S d S )NÚstatusÚtask_id)r   ÚREADY_STATESÚ
cancel_for)r0   rB   r5   r5   r6   Ú_maybe_cancel_ready_task‡   s   ÿz'ResultConsumer._maybe_cancel_ready_taskc                    s    t t| ƒ ||¡ |  |¡ d S r!   )r"   r    rG   rU   )r0   rB   Úmessager3   r5   r6   rG   ‹   s   zResultConsumer.on_state_changec                 K   s    | j jjdd�| _|  |¡ d S )NTrD   )r$   r7   rH   r:   Ú_consume_from)r0   Úinitial_task_idr2   r5   r5   r6   Ústart�   s   ÿzResultConsumer.startc                 K   s.   |j di |¤ŽD ]}|d ur|  |d ¡ qd S rL   )Ú
_iter_metarG   )r0   Úresultr2   rB   r5   r5   r6   Úon_wait_for_pending•   s
   €þz"ResultConsumer.on_wait_for_pendingc                 C   s   | j d ur| j  ¡  d S d S r!   )r:   r;   rO   r5   r5   r6   Ústopš   s   
ÿzResultConsumer.stopc                 C   sž   | j rD|  ¡ �3 | j j|d�}|r*|d dkr2|  |  |d ¡|¡ W d   ƒ d S W d   ƒ d S W d   ƒ d S 1 s=w   Y  d S |rMt |¡ d S d S )N)ÚtimeoutÚtyperV   Údata)r:   rP   Úget_messagerG   r(   ÚtimeÚsleep)r0   r^   rV   r5   r5   r6   Údrain_eventsž   s   
ýþ"þÿzResultConsumer.drain_eventsc                 C   s"   | j d u r
|  |¡S |  |¡ d S r!   )r:   rY   rW   ©r0   rR   r5   r5   r6   Úconsume_from§   s   

zResultConsumer.consume_fromc                 C   s^   |   |¡}|| jvr-| j |¡ |  ¡ � | j |¡ W d   ƒ d S 1 s&w   Y  d S d S r!   )r&   r.   ÚaddrP   r:   rI   ©r0   rR   Úkeyr5   r5   r6   rW   ¬   s   


"ÿþzResultConsumer._consume_fromc                 C   sZ   |   |¡}| j |¡ | jr+|  ¡ � | j |¡ W d   ƒ d S 1 s$w   Y  d S d S r!   )r&   r.   Údiscardr:   rP   Úunsubscriberh   r5   r5   r6   rT   ³   s   

"ÿÿzResultConsumer.cancel_forr!   )Ú__name__Ú
__module__Ú__qualname__r:   r#   r?   rK   r   rP   rU   rG   rY   r\   r]   rd   rf   rW   rT   Ú__classcell__r5   r5   r3   r6   r    Z   s     	


	r    c                       s\  e Zd ZdZeZeZdZdZdZ			d=‡ fdd„	Z	dd„ Z
dd	„ Zd
d„ Zdd„ Zdd„ Zdd„ Zdd„ Zdd„ Z‡ fdd„Zdd„ Zdd„ Zdd„ Zdd„ Zejejfd d!„Zd"d#„ Z	d>d$d%„Zd&d'„ Zd(d)„ Zd*d+„ Ze d,d-„ ƒZ!e"d.d/„ ƒZ#d?‡ fd1d2„	Z$e% &d3d4¡d5d6„ ƒZ'e% &d3d4¡d7d8„ ƒZ(e% &d3d4¡d9d:„ ƒZ)e% &d3d4¡d;d<„ ƒZ*‡  Z+S )@r   zyRedis task result store.

    It makes use of the following commands:
    GET, MGET, DEL, INCRBY, EXPIRE, SET, SETEX
    NTc              	      sä  t t| ƒjddti|¤Ž | jjj}	| jd u rtt	 
¡ ƒ‚|r(d|v r(|d }}|p0|	dƒp0| j| _|| _|	dƒ}
|	dƒ}|	dƒ}|	dƒ}|	dƒpJd	|	d
ƒpOd|	dƒpTd|	dƒ| j|
o^t|
ƒ|pad|oft|ƒdœ| _|rq|| jd< |	dƒ}|rƒ| j |¡ tj| jd< |r�|  || j¡| _d| jv rÔ| jd tju rÔd}ttttttdœ}| j d|¡}| ||¡}|| ¡ vr¼ttƒ‚|tkrÆt t¡ n	|tkrÏt t¡ || jd< || _trÜtƒ nd\| _| _|   | | j| j!| j"| j#¡| _$d S )NÚexpires_typez://Úredis_max_connectionsÚredis_socket_timeoutÚredis_socket_connect_timeoutÚredis_retry_on_timeoutÚredis_socket_keepaliveÚ
redis_hostÚ	localhostÚ
redis_portië  Úredis_dbr   Úredis_passwordF)ÚhostÚportÚdbÚpasswordÚmax_connectionsÚsocket_timeoutÚretry_on_timeoutÚsocket_connect_timeoutÚsocket_keepaliveÚredis_backend_use_sslÚconnection_classÚMISSING)r   r   r   ÚrequiredÚoptionalÚnoneÚssl_cert_reqs)r5   r5   r5   )%r"   r   r#   ÚintÚappÚconfÚgetÚredisr   ÚE_REDIS_MISSINGÚstripr   Ú_ConnectionPoolÚfloatÚ
connparamsÚupdateÚSSLConnectionÚ_params_from_urlr   r   r   ÚvaluesÚ
ValueErrorÚ%E_REDIS_SSL_CERT_REQS_MISSING_INVALIDr=   r>   ÚW_REDIS_SSL_CERT_OPTIONALÚW_REDIS_SSL_CERT_NONEÚurlr   r+   Úchannel_errorsr    ÚacceptÚ_pending_resultsÚ_pending_messagesÚresult_consumer)r0   r{   r|   r}   r~   r   r�   r8   r2   Ú_getr€   r‚   r�   rƒ   ÚsslÚssl_cert_reqs_missingÚssl_string_to_constantrŠ   r3   r5   r6   r#   Í   sx   


ÿý




÷

û



þ

þzRedisBackend.__init__c                    sv  t |ƒ\}}}}}}‰t|fi t|||ˆ dd ¡dœƒ¤Ž‰ |dkr@ˆ  | jjd| dœ¡ ˆ  dd ¡ ˆ  dd ¡ ˆ  d¡ n|ˆ d	< g d
¢}	|dkrft‡ fdd„|	D ƒƒsbt‡fdd„|	D ƒƒrftt	ƒ‚|dkr‚tj
ˆ d< |	D ]}
ˆ |
d ¡}|r�t|ƒˆ |
< qqˆ  d	¡pˆd}t|tƒr“| d¡n|}t|ƒˆ d	< ˆ ¡ D ]\}}|tjjv r³tjj| |ƒˆ|< qŸˆ  ˆ¡ ˆ S )NÚvirtual_host)r{   r|   r~   r}   Úsocketú/)r…   Úpathr{   r|   r‚   r}   )Ússl_ca_certsÚssl_certfileÚssl_keyfilerŠ   r�   c                 3   ó   � | ]}|ˆ v V  qd S r!   r5   ©rA   ri   ©r”   r5   r6   Ú	<genexpr>:  ó   € z0RedisBackend._params_from_url.<locals>.<genexpr>c                 3   r®   r!   r5   r¯   )Úqueryr5   r6   r±   ;  r²   Úredissr…   r   )r   Údictr   Úpopr•   r�   ÚUnixDomainSocketConnectionÚanyr™   Ú&E_REDIS_SSL_PARAMS_AND_SCHEME_MISMATCHr–   r   rŽ   Ú
isinstancer   r‘   r‹   ÚitemsÚ
connectionÚURL_QUERY_ARGUMENT_PARSERS)r0   r�   ÚdefaultsÚschemer{   r|   Ú_r~   rª   Ússl_param_keysÚssl_settingÚssl_valr}   ri   Úvaluer5   )r”   r³   r6   r—     sT   ÿ
þÿþÿ
€
ÿ€
zRedisBackend._params_from_urlc                 C   s   t ƒ s| j |¡ d S d S r!   )r   r¢   rf   )r0   ÚproducerrR   r5   r5   r6   Úon_task_callV  s   ÿzRedisBackend.on_task_callc                 C   ó   | j  |¡S r!   )r7   rŽ   ©r0   ri   r5   r5   r6   rŽ   Z  ó   zRedisBackend.getc                 C   rÇ   r!   )r7   rF   )r0   Úkeysr5   r5   r6   rF   ]  rÉ   zRedisBackend.mgetc                 K   s>   t | jfi |¤Ž}| d¡}t|| j|i t| j|ƒfi |¤ŽS )NÚmax_retries)rµ   Úretry_policyrŽ   r	   r+   r   Úon_connection_error)r0   Úfunr1   ÚpolicyrÌ   rË   r5   r5   r6   r)   `  s   


þýzRedisBackend.ensurec                 C   s*   t |ƒ}t t ¡ ||pdt|dƒ¡ |S )NÚInfzin )Únextr=   ÚerrorÚE_LOSTr‘   r   )r0   rË   ÚexcÚ	intervalsÚretriesÚttsr5   r5   r6   rÍ   h  s   þz RedisBackend.on_connection_errorc                 K   s   | j | j||ffi |¤ŽS r!   )r)   Ú_set)r0   ri   rÄ   ÚstaterÌ   r5   r5   r6   r-   o  ó   zRedisBackend.setc                 C   sh   | j  ¡ �%}| jr| || j|¡ n| ||¡ | ||¡ | ¡  W d   ƒ d S 1 s-w   Y  d S r!   )r7   ÚpipelineÚexpiresÚsetexr-   ÚpublishÚexecute)r0   ri   rÄ   Úpiper5   r5   r6   rØ   r  s   
"úzRedisBackend._setc                    s    t t| ƒ |¡ | j |¡ d S r!   )r"   r   Úforgetr¢   rT   re   r3   r5   r6   rá   {  s   zRedisBackend.forgetc                 C   s   | j  |¡ d S r!   )r7   ÚdeleterÈ   r5   r5   r6   râ     ó   zRedisBackend.deletec                 C   rÇ   r!   )r7   ÚincrrÈ   r5   r5   r6   rä   ‚  rÉ   zRedisBackend.incrc                 C   s   | j  ||¡S r!   )r7   Úexpire)r0   ri   rÄ   r5   r5   r6   rå   …  s   zRedisBackend.expirec                 C   s   | j  |  |d¡d¡ d S )Nú.tr   )r7   rä   Úget_key_for_group)r0   Úgroup_idr[   r5   r5   r6   Úadd_to_chordˆ  rÚ   zRedisBackend.add_to_chordc           	      C   s>   ||ƒ\}}}}||v r|   |¡}||v rtd ||¡ƒ‚|S )NzDependency {0} raised {1!r})Úexception_to_pythonr   Úformat)	r0   ÚtupÚdecodeÚEXCEPTION_STATESÚPROPAGATE_STATESrÀ   ÚtidrÙ   Úretvalr5   r5   r6   Ú_unpack_chord_result‹  s   
z!RedisBackend._unpack_chord_resultc                 K   s   d S r!   r5   )r0   Úheader_resultÚbodyr2   r5   r5   r6   Úapply_chord•  s   zRedisBackend.apply_chordc              
      s˜  | j }|j|j}}|r|sd S | j}	|  |d¡}
|  |d¡}|  ||¡}|	 ¡ �7}| |
|  d|||g¡¡ 	|
¡ 
|¡}| jd urN| |
| j¡ || j¡}| ¡ d d… \}}}W d   ƒ n1 scw   Y  t|pldƒ}z—t|j|d�}|d | }||k�r| j| j‰ ‰|	 ¡ �}| |
d|¡ ¡ \}W d   ƒ n1 s¡w   Y  z5| ‡ ‡fdd	„|D ƒ¡ |	 ¡ �}| |
¡ |¡ ¡ \}}W d   ƒ n1 sÏw   Y  W W d S W W d S  t�y } zt d
|j|¡ |  |td |¡ƒ¡W  Y d }~W S d }~ww W d S  t�y& } zt d|j|¡ |  ||¡W  Y d }~S d }~w t�yK } zt d|j|¡ |  |td |¡ƒ¡W  Y d }~S d }~ww )Nz.jræ   r   é   r   )rŒ   Ú
chord_sizec                    s   g | ]}ˆ|ˆ ƒ‘qS r5   r5   )rA   rì   ©rí   Úunpackr5   r6   rC   Á  s    z5RedisBackend.on_chord_part_return.<locals>.<listcomp>z Chord callback for %r raised: %rzCallback error: {0!r}zChord %r raised: %rzJoin error: {0!r})rŒ   ÚidÚgroupr7   rç   Úencode_resultrÛ   ÚrpushÚencodeÚllenrŽ   rÜ   rå   rß   r‹   r   Úchordrí   rò   ÚlrangeÚdelayrâ   Ú	Exceptionr=   Ú	exceptionÚchord_error_from_stackr   rë   )r0   ÚrequestrÙ   r[   Ú	propagater2   rŒ   rð   Úgidr7   ÚjkeyÚtkeyrà   rÛ   rÀ   Ú
readycountÚ	totaldiffÚcallbackÚtotalÚreslrÔ   r5   rø   r6   Úon_chord_part_return�  s‚   
ý


þõ


þÿ
ý,ÿÿþ€ýó€þ€þz!RedisBackend.on_chord_part_returnc                 K   s   |   ¡ | jdi |¤Žd�S )N)r8   r5   )Ú_get_clientÚ	_get_pool©r0   Úparamsr5   r5   r6   Ú_create_clientØ  s   ÿzRedisBackend._create_clientc                 C   s   | j jS r!   )r�   ÚStrictRedisrO   r5   r5   r6   r  Ý  s   zRedisBackend._get_clientc                 K   s   | j di |¤ŽS rL   )ÚConnectionPoolr  r5   r5   r6   r  à  rã   zRedisBackend._get_poolc                 C   s   | j d u r
| jj| _ | j S r!   )r’   r�   r  rO   r5   r5   r6   r  ã  s   

zRedisBackend.ConnectionPoolc                 C   s   | j di | j¤ŽS rL   )r  r”   rO   r5   r5   r6   r7   é  s   zRedisBackend.clientr5   c                    s(   |si n|}t t| ƒ | jfd| ji¡S )NrÜ   )r"   r   Ú
__reduce__r�   rÜ   r/   r3   r5   r6   r  í  s   
ÿzRedisBackend.__reduce__g      @g      @c                 C   ó
   | j d S )Nr{   r°   rO   r5   r5   r6   r{   ó  ó   
zRedisBackend.hostc                 C   r  )Nr|   r°   rO   r5   r5   r6   r|   ÷  r  zRedisBackend.portc                 C   r  )Nr}   r°   rO   r5   r5   r6   r}   û  r  zRedisBackend.dbc                 C   r  )Nr~   r°   rO   r5   r5   r6   r~   ÿ  r  zRedisBackend.password)NNNNNNNr!   )r5   N),rl   rm   rn   Ú__doc__r    r�   r   Úsupports_autoexpireÚsupports_native_joinr#   r—   rÆ   rŽ   rF   r)   rÍ   r-   rØ   rá   râ   rä   rå   ré   r   rî   rï   rò   rõ   r  r  r  r  Úpropertyr  r
   r7   r  r   ÚPropertyr{   r|   r}   r~   ro   r5   r5   r3   r6   r   »   s\    þR7	
þ
	
ÿ;








r   c                       s@   e Zd ZdZeZ‡ fdd„Z‡ fdd„Zdd„ Zdd	„ Z‡  Z	S )
r   z!Redis sentinel task result store.c                    s0   | j d u rtt ¡ ƒ‚tt| ƒj|i |¤Ž d S r!   )r   r   ÚE_REDIS_SENTINEL_MISSINGr‘   r"   r   r#   r/   r3   r5   r6   r#   	  s   
zSentinelBackend.__init__c                    s’   |  d¡}t|g d�}|D ]}tt| ƒj||d�}|d  |¡ qdD ]}| |¡ q#dD ]}|d rF||d d v rF|d d  |¡||< q-|S )Nú;)Úhosts)r�   r¾   r"  )r{   r|   r}   r~   )r}   r~   r   )Úsplitrµ   r"   r   r—   Úappendr¶   rŽ   )r0   r�   r¾   Úchunksr”   Úchunkr`   Úparamr3   r5   r6   r—     s   

ÿ€z SentinelBackend._params_from_urlc                 K   sb   |  ¡ }| d¡}| jj di ¡}| dd¡}| di ¡}| jjdd„ |D ƒf||dœ|¤Ž}|S )	Nr"  Ú result_backend_transport_optionsÚmin_other_sentinelsr   Úsentinel_kwargsc                 S   s   g | ]
}|d  |d f‘qS )r{   r|   r5   )rA   Úcpr5   r5   r6   rC   ,  s    z:SentinelBackend._get_sentinel_instance.<locals>.<listcomp>)r)  r*  )Úcopyr¶   rŒ   r�   rŽ   r   ÚSentinel)r0   r  r”   r"  Úresult_backend_transport_optsr)  r*  Úsentinel_instancer5   r5   r6   Ú_get_sentinel_instance   s(   
ÿÿÿÿýüz&SentinelBackend._get_sentinel_instancec                 K   s@   | j di |¤Ž}| jj di ¡}| dd ¡}|j||  ¡ d�jS )Nr(  Úmaster_name)Úservice_nameÚredis_classr5   )r0  rŒ   r�   rŽ   Ú
master_forr  r8   )r0   r  r/  r.  r1  r5   r5   r6   r  3  s   ÿþýzSentinelBackend._get_pool)
rl   rm   rn   r  r   r#   r—   r0  r  ro   r5   r5   r3   r6   r     s    r   )Cr  Ú
__future__r   r   rb   Ú
contextlibr   Ú	functoolsr   r¤   r   r   r   Úkombu.utils.functionalr	   Úkombu.utils.objectsr
   Úkombu.utils.urlr   Úceleryr   Úcelery._stater   Úcelery.canvasr   Úcelery.exceptionsr   r   Úcelery.fiver   r   Úcelery.utilsr   Úcelery.utils.functionalr   Úcelery.utils.logr   Úcelery.utils.timer   Úasynchronousr   r   Úbaser   Úurllib.parser   ÚImportErrorÚurlparser�   Úredis.connectionÚkombu.transport.redisr   r   Ú__all__r�   r   r›   rœ   r¹   rš   rÓ   rN   rl   r=   r    r   r   r5   r5   r5   r6   Ú<module>   sj   þþÿa  K