o
    ›¨Êh¬M  ã                   @   s  d Z ddlZddl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 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'„ Z*e'd(d)d*�d+d,„ ƒZ+ej,j-fd-d.„Z.ej/j0ej1j0fd/d0„Z2e&d1d)d*�dŸd2d3„ƒZ3e&d4d5d*�dŸd6d7„ƒZ4dŸd8d9„Z5e&d1d:e6fgd;d<�d=d>„ ƒZ7e&d?e6fd@e6fgdAdB�dCd@„ ƒZ8e&d?e6fdDe9fdEe9fgdFdB�d dGdH„ƒZ:e'ƒ dIdJ„ ƒZ;e&ƒ d¡dKdL„ƒZ<e&ƒ dMdN„ ƒZ=e&ƒ dOdP„ ƒZ>e&ƒ dQdR„ ƒZ?e'd#dS�d¡dTdU„ƒZ@e'dVdW�dXdY„ ƒZAe'ƒ dZd[„ ƒZBe'd\d]�d^d_„ ƒZCd`da„ ZDe'dbd]�dcdd„ ƒZEe'ded]�dždfdg„ƒZFe'dhd]�didj„ ƒZGe'dkdldmdn�d¢dodp„ƒZHe'dqdre6fdseIfdteIfgdudv�d£dzd{„ƒZJe'ƒ d|d}„ ƒZKe'd~eIfgddB�d¤d€d�„ƒZLe&d‚eIfgdƒdB�d¥d„d…„ƒZMe&d‚eIfgdƒdB�d¥d†d‡„ƒZNe&ƒ d¦dˆd‰„ƒZOe&dŠeIfd‹eIfgdŒdB�d§d�dŽ„ƒZPe&ƒ d¨d�d‘„ƒZQe&d’e6fd“e6fd”e6fd•e6fgd–dB�		d d—d˜„ƒZRe&d’e6fgd™dB�dšd›„ ƒZSe'ƒ dœd�„ ƒZTdS )©z.Worker remote control command implementations.é    N)ÚUserDictÚdefaultdictÚ
namedtuple)ÚTERM_SIGNAME)Ú	safe_repr)ÚWorkerShutdown)Ú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!   úG/var/www/html/env/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$   Úregister,   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$   Ú_inner7   s   


þ
zPanel._register.<locals>._innerr!   )r)   r:   r   r   r   r   r   r   r   r   r;   r!   r9   r$   r(   2   s   
zPanel._register)	NNr,   Tr-   NNNN)	r0   Ú
__module__Ú__qualname__r1   r4   r5   Úclassmethodr+   r(   r!   r!   r!   r$   r   &   s    
þr   c                  K   ó   t jdddi| ¤ŽS )Nr   r,   r!   ©r   r+   ©r*   r!   r!   r$   Úcontrol_commandD   ó   rB   c                  K   r?   )Nr   Úinspectr!   r@   rA   r!   r!   r$   Úinspect_commandH   rC   rE   c                 C   s   t | j ¡ ƒS )z6Information about Celery installation for bug reports.)r    ÚappÚ	bugreportr   r!   r!   r$   ÚreportN   ó   rH   Ú	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.)rK   )Ú	keyfilterÚunknown_type_filter)r   rF   ÚconfÚtableÚ_wanted_config_keyr   )r   rK   r*   r!   r!   r$   rN   T   s   þrN   c                 C   s   t | tƒo
|  d¡ S )NÚ__)Ú
isinstanceÚstrÚ
startswith)Úkeyr!   r!   r$   rP   `   s   rP   Ú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>l   s    ÿÿzquery_task.<locals>.<dictcomp>)Ú_find_requests_by_idr	   )r   rV   r*   r!   r!   r$   Ú
query_taskf   s   
þr^   c              	   c   s0   � | D ]}z||ƒV  W q t y   Y qw d S ©N)ÚKeyError)rV   Úget_requestÚtask_idr!   r!   r$   r]   r   s   €ÿýr]   c                 C   s   || ƒrdS || ƒrdS dS )NÚactiveÚreservedÚreadyr!   )ÚrequestÚ	is_activeÚis_reservedr!   r!   r$   rX   {   s
   rX   rb   c                 K   sR   t t|ƒpg ƒd}}t| |||fi |¤Ž}t|tƒr!d|v r!|S td|› 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``).
    Nr    ztasks z flagged as revoked)Úsetr	   Ú_revokerR   Údictr    )r   rb   Ú	terminateÚsignalr*   Útask_idsr!   r!   r$   Úrevoke…   s
   ro   Úheadersz/[key1=value1 [key2=value2 [... [keyN=valueN]]]]c                 K   s,  t  |pt¡}t|tƒrdd„ |D ƒ}| ¡ D ]\}}ttj 	|¡p#g ƒtt|ƒƒ }|tj|< q|s;t
d|› d�ƒS ttjƒ}	ttƒ}
|	D ]=}t|dƒrƒ|jrƒ| ¡ D ].\}}||jv r‚t|ƒ}t|j| ƒ}t|ƒt|ƒ@ }|r‚|
|  |¡ |j| jj|d� qTqF|
sŽt
d|› d�ƒS t
d|
› d�ƒS )	aâ  Revoke task by header (or list of headers).

    Keyword Arguments:
        headers(dictionary): Dictionary that contains stamping scheme name as keys and stamps as values.
                             If headers is a list, it will be converted to a dictionary.
        terminate (bool): Also terminate the process if the task is active.
        signal (str): Name of signal to use for terminate (e.g., ``KILL``).
    Sample headers input:
        {'mtask_id': [id1, id2, id3]}
    c                 S   s&   i | ]}|  d ¡d |  d ¡d “qS )ú=r   r   )r3   )rZ   Úhr!   r!   r$   r\   ±   s   & z-revoke_by_stamped_headers.<locals>.<dictcomp>zheaders z' flagged as revoked, but not terminatedÚstamps©rm   z were not terminatedz revoked)Ú_signalsÚsignumr   rR   ÚlistÚitemsr	   Úworker_stateÚrevoked_stampsÚgetr    Úactive_requestsr   ri   Úhasattrrs   Úupdaterl   ÚconsumerÚpool)r   rp   rl   rm   r*   rv   Úheaderrs   Úupdated_stampsr|   Ú#terminated_scheme_to_stamps_mappingr[   Úexpected_header_keyÚexpected_header_valueÚactual_headerÚmatching_stamps_for_requestr!   r!   r$   Úrevoke_by_stamped_headers›   s0   
 

€rˆ   c           
      K   s¼   t |ƒ}tƒ }tj |¡ |rQt |pt¡}t|ƒD ]&}|j	|vr@| 
|j	¡ t d|j	|¡ |j| jj|d� t |ƒ|kr@ nq|sGtdƒS td d |¡¡ƒS d |¡}	t d|	¡ |S )NzTerminating %s (%s)rt   zterminate: tasks unknownzterminate: {}z, zTasks flagged as revoked: %s)Úlenri   ry   Úrevokedr~   ru   rv   r   r]   rW   ÚaddÚloggerrY   rl   r   r€   r    ÚformatÚjoin)
r   rn   rl   rm   r*   ÚsizeÚ
terminatedrv   rf   Úidstrr!   r!   r$   rj   Õ   s&   
€
rj   rm   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)rl   rm   )ro   )r   rm   rb   r*   r!   r!   r$   rl   í   s   rl   Ú	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.app.task.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: 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'   rF   Útasksr   r`   rŒ   r&   r   Úreset_rate_limitsrY   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)
rF   r—   r`   rŒ   r&   r'   Úsoft_time_limitÚ
time_limitrY   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)rF   rŸ   r#   ©r   r*   r!   r!   r$   rŸ   =  rI   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)r   ÚgossipÚelection)r   rW   ÚtopicÚactionr*   r!   r!   r$   r¢   C  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)r   Úevent_dispatcherÚgroupsr‹   rŒ   rY   r    ©r   Ú
dispatcherr!   r!   r$   Úenable_eventsP  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)r   r¥   r¦   ÚdiscardrŒ   rY   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¬   )rŒ   Údebugr   r¥   Úsendry   ÚSOFTWARE_INFOr§   r!   r!   r$   Ú	heartbeatf  s   
r²   )r   c                 K   sJ   || j kr#t d|¡ |rtj |¡ tj ¡  tjj| jj	 
¡ dœS dS )zRequest mingle sync-data.zsync with %s)rŠ   rŸ   N)ÚhostnamerŒ   rY   ry   rŠ   r~   ÚpurgeÚ_datarF   rŸ   Úforward)r   Ú	from_noderŠ   r*   r!   r!   r$   Úhellop  s   


þúr¸   gš™™™™™É?)r   c                 K   s   t dƒS )zPing worker(s).Úpong)r    r    r!   r!   r$   Úping‚  s   rº   c                 K   s   | j j ¡ S )z&Request worker statistics/information.)r   Ú
controllerÚstatsr    r!   r!   r$   r¼   ˆ  s   r¼   Údump_schedule)r   c                 K   s   t t| jjƒƒS )z0List of currently scheduled ETA/countdown tasks.)rw   Ú_iter_schedule_requestsr   Útimerr    r!   r!   r$   Ú	scheduledŽ  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Úpriorityrf   )ÚscheduleÚqueueÚentryr   Ú
IndexErrorÚ	TypeErrorrR   r   rÁ   Ú	isoformatrÂ   rY   )r¿   ÚwaitingÚarg0r!   r!   r$   r¾   ”  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   s   g | ]}|  ¡ ‘qS r!   ©rY   ©rZ   rf   r!   r!   r$   Ú
<listcomp>¬  s    zreserved.<locals>.<listcomp>)Útsetry   Úreserved_requestsr|   )r   r*   Úreserved_tasksr!   r!   r$   rd   £  s   

ÿÿrd   Údump_activec                    s   ‡ fdd„|   tj¡D ƒS )z'List of tasks currently being executed.c                    s   g | ]}|j ˆ d �‘qS )©ÚsaferÌ   rÍ   rÓ   r!   r$   rÎ   ²  s    ÿzactive.<locals>.<listcomp>)rÏ   ry   r|   )r   rÔ   r*   r!   rÓ   r$   rc   ¯  s   

ÿrc   Údump_revokedc                 K   s
   t tjƒS )zList of revoked task-ids.)rw   ry   rŠ   r    r!   r!   r$   rŠ   ¶  s   
rŠ   Ú
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)rT   ©rZ   rž   r!   r!   r$   Ú	<genexpr>Ì  s   € 
ÿ
ÿzregistered.<locals>.<genexpr>c                    sB   ‡ fdd„ˆD ƒ}|rdd„ |  ¡ D ƒ}d ˆ jd |¡¡S ˆ jS )Nc                    s.   i | ]}t ˆ |d ƒd ur|tt ˆ |d ƒƒ“qS r_   )ÚgetattrrS   )rZ   Úfield©rž   r!   r$   r\   Ð  s
    ÿz5registered.<locals>._extract_info.<locals>.<dictcomp>c                 S   s   g | ]}d   |¡‘qS )rq   )rŽ   )rZ   Úfr!   r!   r$   rÎ   Õ  s    z5registered.<locals>._extract_info.<locals>.<listcomp>z{} [{}]ú )rx   r�   r:   rŽ   )rž   ÚfieldsrY   )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>)rF   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 r_   r!   )Úv©Úobjectsr!   r$   Ú<lambda>õ  s    zobjgraph.<locals>.<lambda>)rç   Ú	highlightÚfilenameró   )	ÚobjgraphÚImportErrorrŒ   rY   Ú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   s4   | j jjr	tdƒS | j j |¡ | j  |¡ tdƒS )z!Grow pool by n processes/threads.zJpool_grow is not supported with autoscale. Adjust autoscale range instead.zpool will grow)r   r»   Ú
autoscalerr'   r€   ÚgrowÚ_update_prefetch_countr    ©r   r  r*   r!   r!   r$   Ú	pool_grow  s
   
r  c                 K   s6   | j jjr	tdƒS | j j |¡ | j  | ¡ tdƒS )z#Shrink pool by n processes/threads.zLpool_shrink is not supported with autoscale. Adjust autoscale range instead.zpool will shrink)r   r»   r  r'   r€   Ú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)rF   rN   Úworker_pool_restartsr   r»   Úreloadr    r–   )r   Úmodulesr  r  r*   r!   r!   r$   Úpool_restart,  s   
r  ÚmaxÚminz[max [min]]c                 C   s:   | j jj}|r| ||¡\}}td|› d|› �ƒS tdƒ‚)zModify autoscale settings.zautoscale now max=z min=zAutoscale not enabled)r   r»   r  r~   r    r–   )r   r  r  r  Úmax_Úmin_r!   r!   r$   Ú	autoscale6  s
   
r  úGot shutdown from remotec                 K   s   t  |¡ t|ƒ‚)zShutdown worker(s).)rŒ   Úwarningr   )r   Úmsgr*   r!   r!   r$   ÚshutdownC  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 )r   Ú	call_soonÚadd_task_queuer    )r   rÄ   r   r  r   Úoptionsr!   r!   r$   Úadd_consumerL  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 )r   r  Úcancel_task_queuer    )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)rk   Úas_dict)rZ   rÄ   r!   r!   r$   rÎ   n  s    ÿz!active_queues.<locals>.<listcomp>)r   Útask_consumerÚqueuesr   r!   r!   r$   Úactive_queuesj  s
   ÿr*  )F)FN)NNNr_   )NF)rè   ré   r   )ré   )r   )NFN)NN)r  )Ur1   r  rö   Úcollectionsr   r   r   Úbilliard.commonr   Úkombu.utils.encodingr   Úcelery.exceptionsr   Úcelery.platformsr   ru   Úcelery.utils.functionalr	   Úcelery.utils.logr
   Úcelery.utils.serializationr   r   Úcelery.utils.timer   r.   r   ry   rf   r   Ú__all__râ   r0   rŒ   r   r    r'   r   rB   rE   rH   rN   rP   r^   ÚrequestsÚ__getitem__r]   r|   Ú__contains__rÐ   rX   ro   rˆ   rj   rS   rl   r   Úfloatr�   rŸ   r¢   r©   r«   r²   r¸   rº   r¼   rÀ   r¾   rd   rc   rŠ   rå   Úintrô   rþ   r  r  r  r  r  r  r"  r%  r*  r!   r!   r!   r$   Ú<module>   s,   
ýþ
	
ÿ

þ
þþ
6ý
þ
$þ





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