o
    wvXjB  ã                   @   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mZ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! ddl"m#Z# dZ$dZ%ee&ƒZ'edg d¢ƒZ(dd„ Z)dd„ Z*G dd„ deƒZ+dd„ Z,dd„ Z-e-ƒ dd „ ƒZ.e-d!d"d#efgd$�dšd&d'„ƒZ/d(d)„ Z0e-d*d+d,�d-d.„ ƒZ1e!j2j3fd/d0„Z4e!j5j6e!j7j6fd1d2„Z8e,d3d+d,�d›d4d5„ƒZ9e,d3d6efgd7d8�d9d:„ ƒZ:e,d;efd<efgd=d>�d?d<„ ƒZ;e,d;efd@e<fdAe<fgdBd>�dœdCdD„ƒZ=e-ƒ dEdF„ ƒZ>e,ƒ d�dGdH„ƒZ?e,ƒ dIdJ„ ƒZ@e,ƒ dKdL„ ƒZAe,ƒ dMdN„ ƒZBe-d%dO�d�dPdQ„ƒZCe-dRdS�dTdU„ ƒZDe-ƒ dVdW„ ƒZEe-dXdY�dZd[„ ƒZFd\d]„ ZGe-d^dY�d_d`„ ƒZHe-dadY�dbdc„ ƒZIe-dddY�dedf„ ƒZJe-dgdhdidj�dždkdl„ƒZKe-dmdnefdoeLfdpeLfgdqdr�dŸdvdw„ƒZMe-ƒ dxdy„ ƒZNe-dzeLfgd{d>�d d|d}„ƒZOe,d~eLfgdd>�d¡d€d�„ƒZPe,d~eLfgdd>�d¡d‚dƒ„ƒZQe,ƒ d¢d„d…„ƒZRe,d†eLfd‡eLfgdˆd>�d£d‰dŠ„ƒZSe,ƒ d¤dŒd�„ƒZTe,dŽefd�efd�efd‘efgd’d>�		dœd“d”„ƒZUe,dŽefgd•d>�d–d—„ ƒZVe-ƒ d˜d™„ ƒZWdS )¥z.Worker remote control command implementations.é    )Úabsolute_importÚunicode_literalsN)Ú
namedtuple)ÚTERM_SIGNAME)Ú	safe_repr)ÚWorkerShutdown)ÚUserDictÚitemsÚstring_tÚtext_t)Úsignals)Ú
maybe_list)Ú
get_logger)ÚjsonifyÚ	strtobool)Úrateé   ©Ústate)ÚRequest)ÚPanel)ÚexchangeÚrouting_keyÚ
rate_limitÚcontroller_info_t)ÚaliasÚtypeÚvisibleÚdefault_timeoutÚhelpÚ	signatureÚargsÚvariadicc                 C   ó   d| iS )NÚok© ©Úvaluer%   r%   úR/var/www/html/myproject/venv/lib/python3.10/site-packages/celery/worker/control.pyr$   "   ó   r$   c                 C   r#   )NÚerrorr%   r&   r%   r%   r(   Únok&   r)   r+   c                   @   s8   e Zd ZdZi Zi Zedd„ ƒZe			d
dd	„ƒZdS )r   z+Global registry of remote control commands.c                 O   s(   |r| j di |¤Ž|Ž S | j di |¤ŽS )Nr%   )Ú	_register)Úclsr!   Úkwargsr%   r%   r(   Úregister0   s   zPanel.registerNÚcontrolTç      ð?c
              
      s"   ‡ ‡‡‡‡‡‡‡‡‡	f
dd„}
|
S )Nc              	      s^   ˆp| j }ˆp| jpd ¡  d¡d }| ˆj|< tˆ ˆˆ	ˆ|ˆˆˆƒˆj|< ˆ r-| ˆjˆ < | S )NÚ Ú
r   )Ú__name__Ú__doc__ÚstripÚsplitÚdatar   Úmeta)ÚfunÚcontrol_nameÚ_help©
r   r!   r-   r   r   Únamer    r   r"   r   r%   r(   Ú_inner;   s   


þ
zPanel._register.<locals>._innerr%   )r-   r>   r   r   r   r   r   r    r!   r"   r?   r%   r=   r(   r,   6   s   
zPanel._register)	NNr0   Tr1   NNNN)	r4   Ú
__module__Ú__qualname__r5   r8   r9   Úclassmethodr/   r,   r%   r%   r%   r(   r   *   s    
þr   c                  K   ó   t jdddi| ¤ŽS )Nr   r0   r%   ©r   r/   ©r.   r%   r%   r(   Úcontrol_commandH   ó   rF   c                  K   rC   )Nr   Úinspectr%   rD   rE   r%   r%   r(   Úinspect_commandL   rG   rI   c                 C   s   t | j ¡ ƒS )z6Information about Celery installation for bug reports.)r$   ÚappÚ	bugreportr   r%   r%   r(   ÚreportR   ó   rL   Ú	dump_confz[include_defaults=False]Úwith_defaults)r   r    r!   Fc                 K   s   t | jjj|d�ttd�S )zList configuration.)rO   )Ú	keyfilterÚunknown_type_filter)r   rJ   ÚconfÚtableÚ_wanted_config_keyr   )r   rO   r.   r%   r%   r(   rR   X   s   þrR   c                 C   s   t | tƒo
|  d¡ S )NÚ__)Ú
isinstancer
   Ú
startswith)Úkeyr%   r%   r(   rT   d   s   rT   Úidsz[id1 [id2 [... [idN]]]])r"   r    c                 K   s   dd„ t t|ƒƒD ƒS )z!Query for task information by id.c                 S   s    i | ]}|j t|ƒ| ¡ f“qS r%   )ÚidÚ_state_of_taskÚinfo)Ú.0Úreqr%   r%   r(   Ú
<dictcomp>p   s    ÿÿzquery_task.<locals>.<dictcomp>)Ú_find_requests_by_idr   )r   rY   r.   r%   r%   r(   Ú
query_taskj   s   
þra   c              	   c   s0   � | D ]}z||ƒV  W q t y   Y qw d S ©N)ÚKeyError)rY   Úget_requestÚtask_idr%   r%   r(   r`   v   s   €ÿýr`   c                 C   s   || ƒrdS || ƒrdS dS )NÚactiveÚreservedÚreadyr%   )ÚrequestÚ	is_activeÚis_reservedr%   r%   r(   r[      s
   r[   re   c                 K   sÜ   t t|ƒpg ƒd}}t|ƒ}t ƒ }tj |¡ |r\t |pt¡}t	|ƒD ]&}	|	j
|vrK| |	j
¡ t d|	j
|¡ |	j| jj|d� t|ƒ|krK nq%|sRtdƒS td d |¡¡ƒS d |¡}
t d|
¡ td |
¡ƒS )	zÝRevoke task by task id (or list of ids).

    Keyword Arguments:
        terminate (bool): Also terminate the process if the task is active.
        signal (str): Name of signal to use for terminate (e.g., ``KILL``).
    NzTerminating %s (%s))Úsignalzterminate: tasks unknownzterminate: {0}z, zTasks flagged as revoked: %sztasks {0} flagged as revoked)Úsetr   ÚlenÚworker_stateÚrevokedÚupdateÚ_signalsÚsignumr   r`   rZ   ÚaddÚloggerr\   Ú	terminateÚconsumerÚpoolr$   ÚformatÚjoin)r   re   rv   rl   r.   Útask_idsÚsizeÚ
terminatedrs   ri   Úidstrr%   r%   r(   Úrevoke‰   s(   
€
r   rl   z <signal> [id1 [id2 [... [idN]]]])r"   r!   r    c                 K   s   t | |d|d�S )z+Terminate task by task id (or list of ids).T)rv   rl   )r   )r   rl   re   r.   r%   r%   r(   rv   °   s   rv   Ú	task_namer   z0<task_name> <rate_limit (e.g., 5/s | 5/m | 5/h)>)r!   r    c              
   K   s¶   zt |ƒ W n ty } ztd |¡ƒW  Y d}~S d}~ww z	|| jj| _W n ty>   tj	d|dd� tdƒ Y S w | j
 ¡  |sPt d|¡ tdƒS t d	||¡ td
ƒS )zýTell worker(s) to modify the rate limit for a task by type.

    See Also:
        :attr:`celery.task.base.Task.rate_limit`.

    Arguments:
        task_name (str): Type of task to set rate limit for.
        rate_limit (int, str): New rate limit.
    z Invalid rate limit string: {0!r}Nz&Rate limit attempt for unknown task %sT©Úexc_infoúunknown taskz)Rate limits disabled for tasks of type %sz rate limit disabled successfullyz(New rate limit for tasks of type %s: %s.znew rate limit set successfully)r   Ú
ValueErrorr+   ry   rJ   Útasksr   rc   ru   r*   rw   Úreset_rate_limitsr\   r$   )r   r€   r   r.   Úexcr%   r%   r(   r   º   s,   €ÿÿý
ÿÚsoftÚhardz#<task_name> <soft_secs> [hard_secs]c                 K   s`   z| j j| }W n ty   tjd|dd� tdƒ Y S w ||_||_t d|||¡ t	dƒS )zÍTell worker(s) to modify the time limit for task by type.

    Arguments:
        task_name (str): Name of task to change.
        hard (float): Hard time limit.
        soft (float): Soft time limit.
    z-Change time limit attempt for unknown task %sTr�   rƒ   z5New time limits for tasks of type %s: soft=%s hard=%sztime limits set successfully)
rJ   r…   rc   ru   r*   r+   Úsoft_time_limitÚ
time_limitr\   r$   )r   r€   r‰   rˆ   r.   Útaskr%   r%   r(   r‹   â   s   ÿýÿr‹   c                 K   s   d| j jjiS )z Get current logical clock value.Úclock)rJ   r�   r'   ©r   r.   r%   r%   r(   r�      rM   r�   c                 K   s"   | j jr| j j |||¡ dS dS )z¦Hold election.

    Arguments:
        id (str): Unique election id.
        topic (str): Election topic.
        action (str): Action to take for elected actor.
    N)rw   ÚgossipÚelection)r   rZ   ÚtopicÚactionr.   r%   r%   r(   r�     s   	ÿr�   c                 C   s>   | j j}|jrd|jvr|j d¡ t d¡ tdƒS tdƒS )z+Tell worker(s) to send task-related events.rŒ   z)Events of group {task} enabled by remote.ztask events enabledztask events already enabled)rw   Úevent_dispatcherÚgroupsrt   ru   r\   r$   ©r   Ú
dispatcherr%   r%   r(   Úenable_events  s   
r—   c                 C   s8   | j j}d|jv r|j d¡ t d¡ tdƒS tdƒS )z3Tell worker(s) to stop sending task-related events.rŒ   z*Events of group {task} disabled by remote.ztask events disabledztask events already disabled)rw   r“   r”   Údiscardru   r\   r$   r•   r%   r%   r(   Údisable_events  s   

r™   c                 C   s,   t  d¡ | jj}|jddditj¤Ž dS )z3Tell worker(s) to send event heartbeat immediately.zHeartbeat requested by remote.úworker-heartbeatÚfreqé   N)rš   )ru   Údebugrw   r“   Úsendro   ÚSOFTWARE_INFOr•   r%   r%   r(   Ú	heartbeat)  s   
r    )r   c                 K   s@   || j krt d|¡ |rtj |¡ tjj| jj 	¡ dœS dS )zRequest mingle sync-data.zsync with %s)rp   r�   N)
Úhostnameru   r\   ro   rp   rq   Ú_datarJ   r�   Úforward)r   Ú	from_noderp   r.   r%   r%   r(   Úhello3  s   

þür¥   gš™™™™™É?)r   c                 K   s   t dƒS )zPing worker(s).Úpong)r$   rŽ   r%   r%   r(   ÚpingC  s   r§   c                 K   s   | j j ¡ S )z&Request worker statistics/information.)rw   Ú
controllerÚstatsrŽ   r%   r%   r(   r©   I  s   r©   Údump_schedule)r   c                 K   s   t t| jjƒƒS )z0List of currently scheduled ETA/countdown tasks.)ÚlistÚ_iter_schedule_requestsrw   ÚtimerrŽ   r%   r%   r(   Ú	scheduledO  s   r®   c              
   c   sj   � | j jD ]-}z|jjd }W n ttfy   Y qw t|tƒr2|jr(|j 	¡ nd |j
| ¡ dœV  qd S )Nr   )ÚetaÚpriorityri   )ÚscheduleÚqueueÚentryr!   Ú
IndexErrorÚ	TypeErrorrV   r   r¯   Ú	isoformatr°   r\   )r­   ÚwaitingÚarg0r%   r%   r(   r¬   U  s   €ÿ
ý€ùr¬   Údump_reservedc                 K   s.   |   tj¡|   tj¡ }|sg S dd„ |D ƒS )zAList of currently reserved tasks, not including scheduled/active.c                 S   ó   g | ]}|  ¡ ‘qS r%   ©r\   ©r]   ri   r%   r%   r(   Ú
<listcomp>m  s    zreserved.<locals>.<listcomp>)Útsetro   Úreserved_requestsÚactive_requests)r   r.   Úreserved_tasksr%   r%   r(   rg   d  s   

ÿÿrg   Údump_activec                 K   s   dd„ |   tj¡D ƒS )z'List of tasks currently being executed.c                 S   rº   r%   r»   r¼   r%   r%   r(   r½   s  s    ÿzactive.<locals>.<listcomp>)r¾   ro   rÀ   rŽ   r%   r%   r(   rf   p  s   
ÿrf   Údump_revokedc                 K   s
   t tjƒS )zList of revoked task-ids.)r«   ro   rp   rŽ   r%   r%   r(   rp   w  s   
rp   Ú
dump_tasksÚtaskinfoitemsz[attr1 [attr2 [... [attrN]]]])r   r"   r    c                    sJ   | j j‰ˆpt‰|rˆndd„ ˆD ƒ}‡fdd„‰ ‡ ‡fdd„t|ƒD ƒS )zìList of registered tasks.

    Arguments:
        taskinfoitems (Sequence[str]): List of task attributes to include.
            Defaults to ``exchange,routing_key,rate_limit``.
        builtins (bool): Also include built-in tasks.
    c                 s   s   � | ]
}|  d ¡s|V  qdS )zcelery.N)rW   ©r]   rŒ   r%   r%   r(   Ú	<genexpr>�  s   € 
ÿ
ÿzregistered.<locals>.<genexpr>c                    sB   ‡ fdd„ˆD ƒ}|rdd„ t |ƒD ƒ}d ˆ jd |¡¡S ˆ jS )Nc                    s.   i | ]}t ˆ |d ƒd ur|tt ˆ |d ƒƒ“qS rb   )ÚgetattrÚstr)r]   Úfield©rŒ   r%   r(   r_   ‘  s
    ÿz5registered.<locals>._extract_info.<locals>.<dictcomp>c                 S   s   g | ]}d   |¡‘qS )ú=)rz   )r]   Úfr%   r%   r(   r½   –  s    z5registered.<locals>._extract_info.<locals>.<listcomp>z	{0} [{1}]ú )r	   ry   r>   rz   )rŒ   Úfieldsr\   )rÅ   rË   r(   Ú_extract_info�  s   
ÿz!registered.<locals>._extract_infoc                    s   g | ]}ˆ ˆ| ƒ‘qS r%   r%   rÆ   )rÐ   Úregr%   r(   r½   š  s    zregistered.<locals>.<listcomp>)rJ   r…   ÚDEFAULT_TASK_INFO_ITEMSÚsorted)r   rÅ   Úbuiltinsr.   r…   r%   )rÐ   rÑ   rÅ   r(   Ú
registered}  s   ÿ
rÕ   g      N@r   ÚnumÚ	max_depthz.[object_type=Request] [num=200 [max_depth=10]])r   r!   r    éÈ   é
   r   c                    sœ   zddl }W n ty   tdƒ‚w t d|¡ tjdddd��$}| |¡d|… ‰ |jˆ |‡ fd	d
„|jd� d|jiW  d  ƒ S 1 sGw   Y  dS )a  Create graph of uncollected objects (memory-leak debugging).

    Arguments:
        num (int): Max number of objects to graph.
        max_depth (int): Traverse at most n levels deep.
        type (str): Name of object to graph.  Default is ``"Request"``.
    r   NzRequires the objgraph libraryzDumping graph for type %rÚcobjgz.pngF)ÚprefixÚsuffixÚdeletec                    s   | ˆ v S rb   r%   )Úv©Úobjectsr%   r(   Ú<lambda>¶  s    zobjgraph.<locals>.<lambda>)r×   Ú	highlightÚfilenamerã   )	ÚobjgraphÚImportErrorru   r\   ÚtempfileÚNamedTemporaryFileÚby_typeÚshow_backrefsr>   )r   rÖ   r×   r   Ú	_objgraphÚfhr%   rß   r(   rä   Ÿ  s$   ÿÿý$ørä   c                 K   s   ddl m} |ƒ S )z Sample current RSS memory usage.r   )Ú
sample_mem)Úcelery.utils.debugrì   )r   r.   rì   r%   r%   r(   Ú	memsample¼  s   rî   Úsamplesz[n_samples=10]c                 K   s(   ddl m} t ¡ }|j|d� | ¡ S )z/Dump statistics of previous memsample requests.r   )r�   )Úfile)Úcelery.utilsr�   ÚioÚStringIOÚmemdumpÚgetvalue)r   rï   r.   r�   Úoutr%   r%   r(   rô   Ã  s   rô   Únz[N=1]c                 K   sD   | j jjr| j jj |¡ tdƒS | j j |¡ | j  |¡ tdƒS )z!Grow pool by n processes/threads.zpool will grow)rw   r¨   Ú
autoscalerÚforce_scale_uprx   ÚgrowÚ_update_prefetch_countr$   ©r   r÷   r.   r%   r%   r(   Ú	pool_growÑ  s   
þrý   c                 K   sF   | j jjr| j jj |¡ tdƒS | j j |¡ | j  | ¡ tdƒS )z#Shrink pool by n processes/threads.zpool will shrink)rw   r¨   rø   Úforce_scale_downrx   Úshrinkrû   r$   rü   r%   r%   r(   Úpool_shrinkß  s   
þr   c                 K   s.   | j jjr| jjj|||d� tdƒS tdƒ‚)zRestart execution pool.)Úreloaderzreload startedzPool restarts not enabled)rJ   rR   Úworker_pool_restartsrw   r¨   Úreloadr$   r„   )r   Úmodulesr  r  r.   r%   r%   r(   Úpool_restartí  s   
r  ÚmaxÚminz[max [min]]c                 C   s6   | j jj}|r| ||¡\}}td ||¡ƒS tdƒ‚)zModify autoscale settings.zautoscale now max={0} min={1}zAutoscale not enabled)rw   r¨   rø   rq   r$   ry   r„   )r   r  r  rø   Úmax_Úmin_r%   r%   r(   Ú	autoscale÷  s
   
r
  úGot shutdown from remotec                 K   s   t  |¡ t|ƒ‚)zShutdown worker(s).)ru   Úwarningr   )r   Úmsgr.   r%   r%   r(   Úshutdown  s   
r  r²   r   Úexchange_typer   z'<queue> [exchange [type [routing_key]]]c                 K   s2   | j j| j j|||pd|fi |¤Ž td |¡ƒS )z2Tell worker(s) to consume from task queue by name.Údirectzadd consumer {0})rw   Ú	call_soonÚadd_task_queuer$   ry   )r   r²   r   r  r   Úoptionsr%   r%   r(   Úadd_consumer  s   þþr  z<queue>c                 K   s    | j  | j j|¡ td |¡ƒS )z9Tell worker(s) to stop consuming from task queue by name.zno longer consuming from {0})rw   r  Úcancel_task_queuer$   ry   )r   r²   Ú_r%   r%   r(   Úcancel_consumer  s   ÿr  c                 C   s    | j jrdd„ | j jjD ƒS g S )z:List the task queues a worker is currently consuming from.c                 S   s   g | ]
}t |jd d�ƒ‘qS )T)Úrecurse)ÚdictÚas_dict)r]   r²   r%   r%   r(   r½   /  s    ÿz!active_queues.<locals>.<listcomp>)rw   Útask_consumerÚqueuesr   r%   r%   r(   Úactive_queues+  s
   ÿr  )F)FN)NNNrb   )NF)rØ   rÙ   r   )rÙ   )r   )NFN)NN)r  )Xr5   Ú
__future__r   r   rò   ræ   Úcollectionsr   Úbilliard.commonr   Úkombu.utils.encodingr   Úcelery.exceptionsr   Úcelery.fiver   r	   r
   r   Úcelery.platformsr   rr   Úcelery.utils.functionalr   Úcelery.utils.logr   Úcelery.utils.serializationr   r   Úcelery.utils.timer   r2   r   ro   ri   r   Ú__all__rÒ   r4   ru   r   r$   r+   r   rF   rI   rL   rR   rT   ra   ÚrequestsÚ__getitem__r`   rÀ   Ú__contains__r¿   r[   r   rv   r   Úfloatr‹   r�   r�   r—   r™   r    r¥   r§   r©   r®   r¬   rg   rf   rp   rÕ   Úinträ   rî   rô   rý   r   r  r
  r  r  r  r  r%   r%   r%   r(   Ú<module>   s$  
ýþ
	
ÿ

þ
þ#ý
þ
$þ





	





ýý
þ
þ
þ
	þ	üù	ÿ	þ
