o
    wvXj|_  ã                   @   s(  d Z ddlmZmZ ddlZddlZddlZddlZddlZddl	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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"m#Z# ddl$m%Z%m&Z&m'Z'm(Z(m)Z) ddl*m+Z+m,Z, ddl-m.Z.m/Z/ ddl0m1Z1m2Z2 ddl3m4Z4m5Z5 dZ6eddƒZ7e1e8ƒZ9e9j:e9j;e9j<e9j=f\Z:Z;Z<Z=dZ>G dd„ de?ƒZ@G dd„ deAƒZBee'G dd„ deAƒƒƒZCG dd„ deAƒZDG d d!„ d!eDƒZEG d"d#„ d#eAƒZFG d$d%„ d%eƒZGzeƒ  W n eH�y   dZIY n	w G d&d'„ d'eƒZId*d(d)„ZJdS )+zThe periodic task scheduler.é    )Úabsolute_importÚunicode_literalsN)Útimegm)Ú
namedtuple)Útotal_ordering)ÚEventÚThread)Úensure_multiprocessing)Úreset_signals)ÚProcess)Úmaybe_evaluateÚreprcall)Úcached_propertyé   )Ú__version__Ú	platformsÚsignals)ÚitemsÚ	monotonicÚpython_2_unicode_compatibleÚreraiseÚvalues)ÚcrontabÚmaybe_schedule)Úload_extension_class_namesÚsymbol_by_name)Ú
get_loggerÚiter_open_logger_fds)Úhumanize_secondsÚmaybe_make_aware)ÚSchedulingErrorÚScheduleEntryÚ	SchedulerÚPersistentSchedulerÚServiceÚEmbeddedServiceÚevent_t)ÚtimeÚpriorityÚentryi,  c                   @   s   e Zd ZdZdS )r    z*An error occurred while scheduling a task.N)Ú__name__Ú
__module__Ú__qualname__Ú__doc__© r.   r.   úH/var/www/html/myproject/venv/lib/python3.10/site-packages/celery/beat.pyr    .   s    r    c                   @   s(   e Zd ZdZdd„ Zdd„ Zdd„ ZdS )	ÚBeatLazyFuncap  An lazy function declared in 'beat_schedule' and called before sending to worker.

    Example:

        beat_schedule = {
            'test-every-5-minutes': {
                'task': 'test',
                'schedule': 300,
                'kwargs': {
                    "current": BeatCallBack(datetime.datetime.now)
                }
            }
        }

    c                 O   s   || _ ||dœ| _d S )N)ÚargsÚkwargs©Ú_funcÚ_func_params)ÚselfÚfuncr1   r2   r.   r.   r/   Ú__init__C   s   þzBeatLazyFunc.__init__c                 C   ó   |   ¡ S ©N)Údelay©r6   r.   r.   r/   Ú__call__J   ó   zBeatLazyFunc.__call__c                 C   s   | j | jd i | jd ¤ŽS )Nr1   r2   r3   r<   r.   r.   r/   r;   M   s   zBeatLazyFunc.delayN)r*   r+   r,   r-   r8   r=   r;   r.   r.   r.   r/   r0   2   s
    r0   c                   @   s¢   e Zd ZdZdZdZdZdZdZdZ	dZ
			ddd„Zdd	„ ZeZdd
d„Ze Z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dd„ ZdS ) r!   aâ  An entry in the scheduler.

    Arguments:
        name (str): see :attr:`name`.
        schedule (~celery.schedules.schedule): see :attr:`schedule`.
        args (Tuple): see :attr:`args`.
        kwargs (Dict): see :attr:`kwargs`.
        options (Dict): see :attr:`options`.
        last_run_at (~datetime.datetime): see :attr:`last_run_at`.
        total_run_count (int): see :attr:`total_run_count`.
        relative (bool): Is the time relative to when the server starts?
    Nr   r.   Fc                 C   sb   |
| _ || _|| _|| _|r|ni | _|r|ni | _t||	| j d�| _|p(|  ¡ | _	|p-d| _
d S )N)Úappr   )r?   ÚnameÚtaskr1   r2   Úoptionsr   ÚscheduleÚdefault_nowÚlast_run_atÚtotal_run_count)r6   r@   rA   rE   rF   rC   r1   r2   rB   Úrelativer?   r.   r.   r/   r8   v   s   zScheduleEntry.__init__c                 C   s   | j r| j  ¡ S | j ¡ S r:   )rC   Únowr?   r<   r.   r.   r/   rD   ƒ   s   zScheduleEntry.default_nowc                 C   s(   | j di t| |p|  ¡ | jd d�¤ŽS )z8Return new instance, with date and count fields updated.r   )rE   rF   Nr.   )Ú	__class__ÚdictrD   rF   )r6   rE   r.   r.   r/   Ú_next_instance‡   s
   


ýzScheduleEntry._next_instancec              	   C   s*   | j | j| j| j| j| j| j| j| jffS r:   )	rI   r@   rA   rE   rF   rC   r1   r2   rB   r<   r.   r.   r/   Ú
__reduce__�   s   þzScheduleEntry.__reduce__c                 C   s&   | j  |j|j|j|j|jdœ¡ dS )zžUpdate values from another entry.

        Will only update "editable" fields:
            ``task``, ``schedule``, ``args``, ``kwargs``, ``options``.
        )rA   rC   r1   r2   rB   N)Ú__dict__ÚupdaterA   rC   r1   r2   rB   ©r6   Úotherr.   r.   r/   rN   –   s
   ýzScheduleEntry.updatec                 C   s   | j  | j¡S )z-See :meth:`~celery.schedule.schedule.is_due`.)rC   Úis_duerE   r<   r.   r.   r/   rQ   ¢   s   zScheduleEntry.is_duec                 C   s   t tt| ƒƒƒS r:   )Úiterr   Úvarsr<   r.   r.   r/   Ú__iter__¦   s   zScheduleEntry.__iter__c                 C   s,   dj | t| j| jp
d| jpi ƒt| ƒjd�S )Nz%<{name}: {0.name} {call} {0.schedule}r.   )Úcallr@   )Úformatr   rA   r1   r2   Útyper*   r<   r.   r.   r/   Ú__repr__©   s
   ýzScheduleEntry.__repr__c                 C   s   t |tƒrt| ƒt|ƒk S tS r:   )Ú
isinstancer!   ÚidÚNotImplementedrO   r.   r.   r/   Ú__lt__°   s   
zScheduleEntry.__lt__c                 C   s(   dD ]}t | |ƒt ||ƒkr dS qdS )N)rA   r1   r2   rB   rC   FT)Úgetattr)r6   rP   Úattrr.   r.   r/   Úeditable_fields_equal»   s
   ÿz#ScheduleEntry.editable_fields_equalc                 C   s
   |   |¡S )z™Test schedule entries equality.

        Will only compare "editable" fields:
        ``task``, ``schedule``, ``args``, ``kwargs``, ``options``.
        )r_   rO   r.   r.   r/   Ú__eq__Á   ó   
zScheduleEntry.__eq__c                 C   s
   | |k S )z›Test schedule entries inequality.

        Will only compare "editable" fields:
        ``task``, ``schedule``, ``args``, ``kwargs``, ``options``.
        r.   rO   r.   r.   r/   Ú__ne__É   ra   zScheduleEntry.__ne__)
NNNNNr.   NNFNr:   )r*   r+   r,   r-   r@   rC   r1   r2   rB   rE   rF   r8   rD   Ú_default_nowrK   Ú__next__ÚnextrL   rN   rQ   rT   rX   r\   r_   r`   rb   r.   r.   r.   r/   r!   Q   s4    
þ
r!   c                   @   sD  e Zd ZdZeZdZeZdZ	dZ
dZdZeZ		d>dd„Zdd	„ Zd?d
d„Zd@dd„Zdd„ Zefdd„Zeejfdd„Zeeejejfdd„Zdd„ Zdd„ Zdd„ ZdAd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(d0d1„ Z)d2d3„ Z*d4d5„ Z+d6d7„ Z,e-e+e,ƒZe.d8d9„ ƒZ/e.d:d;„ ƒZ0e-d<d=„ ƒZ1dS )Br"   aË  Scheduler for periodic tasks.

    The :program:`celery beat` program may instantiate this class
    multiple times for introspection purposes, but then with the
    ``lazy`` argument set.  It's important for subclasses to
    be idempotent when this argument is set.

    Arguments:
        schedule (~celery.schedules.schedule): see :attr:`schedule`.
        max_interval (int): see :attr:`max_interval`.
        lazy (bool): Don't set up the schedule.
    Né´   r   Fc                 K   st   || _ t|d u r
i n|ƒ| _|p|jjp| j| _|p|jj| _d | _d | _	|d u r-|jj
n|| _|s8|  ¡  d S d S r:   )r?   r   ÚdataÚconfÚbeat_max_loop_intervalÚmax_intervalÚamqpÚProducerÚ_heapÚold_schedulersÚbeat_sync_everyÚsync_every_tasksÚsetup_schedule)r6   r?   rC   rj   rl   Úlazyrp   r2   r.   r.   r/   r8   ó   s    ÿþþÿzScheduler.__init__c                 C   sJ   i }| j jjr| j jjsd|vrdtdddƒddidœ|d< |  |¡ d S )Nzcelery.backend_cleanupÚ0Ú4Ú*ÚexpiresiÀ¨  )rA   rC   rB   )r?   rh   Úresult_expiresÚbackendÚsupports_autoexpirer   Úupdate_from_dict)r6   rg   Úentriesr.   r.   r/   Úinstall_default_entries  s   
ÿ

ýz!Scheduler.install_default_entriesc              
   C   st   t d|j|jƒ z
| j||dd�}W n ty/ } ztd|t ¡ dd� W Y d }~d S d }~ww td|j|j	ƒ d S )Nz#Scheduler: Sending due task %s (%s)F)ÚproducerÚadvancezMessage Error: %s
%sT©Úexc_infoz%s sent. id->%s)
Úinfor@   rA   Úapply_asyncÚ	ExceptionÚerrorÚ	tracebackÚformat_stackÚdebugrZ   )r6   r)   r}   ÚresultÚexcr.   r.   r/   Úapply_entry  s   
ÿ€ÿzScheduler.apply_entryç{®Gáz„¿c                 C   s   |r
|dkr
|| S |S )Nr   r.   )r6   ÚnÚdriftr.   r.   r/   Úadjust  s   zScheduler.adjustc                 C   s   |  ¡ S r:   )rQ   )r6   r)   r.   r.   r/   rQ     r>   zScheduler.is_duec                 C   s4   | j }t| ¡ ƒ}|| ¡ ƒ|jd  ||ƒpd S )z9Return a utc timestamp, make sure heapq in currect order.g    €„.Ar   )rŽ   r   rD   ÚutctimetupleÚmicrosecond)r6   r)   Únext_time_to_runÚmktimerŽ   Úas_nowr.   r.   r/   Ú_when   s   
ÿ
þzScheduler._whenc                 C   s\   d}g | _ t| jƒD ]}| ¡ \}}| j  ||  ||rdn|¡p!d||ƒ¡ q
|| j ƒ dS )z:Populate the heap with the data contained in the schedule.é   r   N)rm   r   rC   rQ   Úappendr”   )r6   r&   Úheapifyr(   r)   rQ   Únext_call_delayr.   r.   r/   Úpopulate_heap*  s   
þûzScheduler.populate_heapc                 C   sâ   | j }| j}| jdu s|  | j| j¡st | j¡| _|  ¡  | j}|s%|S |d }|d }	|  |	¡\}
}|
rh||ƒ}||u r\|  	|	¡}| j
|	| jd� ||||  ||¡|d |ƒƒ dS |||ƒ ||d |ƒS |||ƒpn||ƒS )z­Run a tick - one iteration of the scheduler.

        Executes one due task per call.

        Returns:
            float: preferred delay in seconds for next call.
        Nr   é   )r}   r   )rŽ   rj   rm   Úschedules_equalrn   rC   Úcopyr™   rQ   ÚreserverŠ   r}   r”   )r6   r&   ÚminÚheappopÚheappushrŽ   rj   ÚHÚeventr)   rQ   r‘   ÚverifyÚ
next_entryr.   r.   r/   Útick:  s2   	
ÿ
ÿ
zScheduler.tickc                 C   s€   ||  u rd u rdS  |d u s|d u rdS t | ¡ ƒt | ¡ ƒkr$dS | ¡ D ]\}}| |¡}|s6 dS ||kr= dS q(dS )NTF)ÚsetÚkeysr   Úget)r6   Úold_schedulesÚnew_schedulesr@   Ú	old_entryÚ	new_entryr.   r.   r/   r›   `  s   ÿ
ÿzScheduler.schedules_equalc                 C   s,   | j  ptƒ | j  | jkp| jo| j| jkS r:   )Ú
_last_syncr   Ú
sync_everyrp   Ú_tasks_since_syncr<   r.   r.   r/   Úshould_synco  s   ÿ
üzScheduler.should_syncc                 C   s   t |ƒ }| j|j< |S r:   )re   rC   r@   )r6   r)   r¬   r.   r.   r/   r�   w  s   zScheduler.reserveTc           	   
   K   sb  |r|   |¡n|}| jj |j¡}zŽzVdd„ |jpg D ƒ}dd„ |j ¡ D ƒ}|rH|j||fd|i|j	¤ŽW W |  j
d7  _
|  ¡ rG|  ¡  S S | j|j||fd|i|j	¤ŽW W |  j
d7  _
|  ¡ rh|  ¡  S S  ty‹ } ztttdj||d�ƒt ¡ d	 ƒ W Y d }~nd }~ww W |  j
d7  _
|  ¡ rž|  ¡  d S d S |  j
d7  _
|  ¡ r°|  ¡  w w )
Nc                 S   s    g | ]}t |tƒr|ƒ n|‘qS r.   ©rY   r0   )Ú.0Úvr.   r.   r/   Ú
<listcomp>ƒ  s     z)Scheduler.apply_async.<locals>.<listcomp>c                 S   s&   i | ]\}}|t |tƒr|ƒ n|“qS r.   r±   )r²   Úkr³   r.   r.   r/   Ú
<dictcomp>„  s   & z)Scheduler.apply_async.<locals>.<dictcomp>r}   r   z-Couldn't apply scheduled task {0.name}: {exc})r‰   rš   )r�   r?   Útasksr¨   rA   r1   r2   r   r‚   rB   r¯   r°   Ú_do_syncÚ	send_taskrƒ   r   r    rV   Úsysr€   )	r6   r)   r}   r~   r2   rA   Ú
entry_argsÚentry_kwargsr‰   r.   r.   r/   r‚   {  sV   ÿþ
ÿ÷ÿþ
ÿúÿÿ
þ€ÿÿÿ
ÿzScheduler.apply_asyncc                 O   s   | j j|i |¤ŽS r:   )r?   r¹   ©r6   r1   r2   r.   r.   r/   r¹   –  ó   zScheduler.send_taskc                 C   s    |   | j¡ |  | jjj¡ d S r:   )r|   rg   Úmerge_inplacer?   rh   Úbeat_scheduler<   r.   r.   r/   rq   ™  s   zScheduler.setup_schedulec                 C   s6   zt dƒ |  ¡  W tƒ | _d| _d S tƒ | _d| _w )Nzbeat: Synchronizing schedule...r   )r‡   Úsyncr   r­   r¯   r<   r.   r.   r/   r¸   �  s   

ÿzScheduler._do_syncc                 C   s   d S r:   r.   r<   r.   r.   r/   rÁ   ¥  s   zScheduler.syncc                 C   s   |   ¡  d S r:   )rÁ   r<   r.   r.   r/   Úclose¨  s   zScheduler.closec                 K   s&   | j dd| ji|¤Ž}|| j|j< |S )Nr?   r.   )ÚEntryr?   rC   r@   )r6   r2   r)   r.   r.   r/   Úadd«  s   zScheduler.addc                 C   s4   t || jƒr| j|_|S | jdi t||| jd�¤ŽS ©N)r@   r?   r.   )rY   rÃ   r?   rJ   )r6   r@   r)   r.   r.   r/   Ú_maybe_entry°  s   zScheduler._maybe_entryc                    s"   ˆ j  ‡ fdd„t|ƒD ƒ¡ d S )Nc                    s   i | ]\}}|ˆ   ||¡“qS r.   )rÆ   )r²   r@   r)   r<   r.   r/   r¶   ·  s    ÿÿz.Scheduler.update_from_dict.<locals>.<dictcomp>)rC   rN   r   )r6   Údict_r.   r<   r/   rz   ¶  s   þzScheduler.update_from_dictc              	   C   s‚   | j }t|ƒt|ƒ}}||A D ]}| |d ¡ q|D ]#}| jdi t|| || jd�¤Ž}| |¡r:||  |¡ q|||< qd S rÅ   )rC   r¦   ÚpoprÃ   rJ   r?   r¨   rN   )r6   ÚbrC   ÚAÚBÚkeyr)   r.   r.   r/   r¿   ¼  s    

ûzScheduler.merge_inplacec                 C   s   dd„ }| j  || jjj¡S )Nc                 S   s   t d| |ƒ d S )Nz9beat: Connection error: %s. Trying again in %s seconds...)r„   )r‰   Úintervalr.   r.   r/   Ú_error_handlerÏ  s   ÿz3Scheduler._ensure_connected.<locals>._error_handler)Ú
connectionÚensure_connectionr?   rh   Úbroker_connection_max_retries)r6   rÎ   r.   r.   r/   Ú_ensure_connectedÌ  s   
ÿzScheduler._ensure_connectedc                 C   s   | j S r:   ©rg   r<   r.   r.   r/   Úget_schedule×  s   zScheduler.get_schedulec                 C   s
   || _ d S r:   rÓ   ©r6   rC   r.   r.   r/   Úset_scheduleÚ  s   
zScheduler.set_schedulec                 C   s
   | j  ¡ S r:   )r?   Úconnection_for_writer<   r.   r.   r/   rÏ   Þ  s   
zScheduler.connectionc                 C   s   | j |  ¡ dd�S )NF)Úauto_declare)rl   rÒ   r<   r.   r.   r/   r}   â  s   zScheduler.producerc                 C   s   dS )NÚ r.   r<   r.   r.   r/   r�   æ  s   zScheduler.info)NNNFNr:   )r‹   )NT)2r*   r+   r,   r-   r!   rÃ   rC   ÚDEFAULT_MAX_INTERVALrj   r®   rp   r­   r¯   Úloggerr8   r|   rŠ   rŽ   rQ   r   r”   r&   Úheapqr—   r™   rž   rŸ   r    r¥   r›   r°   r�   r‚   r¹   rq   r¸   rÁ   rÂ   rÄ   rÆ   rz   r¿   rÒ   rÔ   rÖ   Úpropertyr   rÏ   r}   r�   r.   r.   r.   r/   r"   Ò   sZ    
ÿ




ÿ&



r"   c                   @   s‚   e Zd ZdZeZdZ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eeeƒZdd„ Zdd„ Zedd„ ƒZdS )r#   z+Scheduler backed by :mod:`shelve` database.)rÙ   z.dbz.datz.bakz.dirNc                 O   s(   |  d¡| _tj| g|¢R i |¤Ž d S )NÚschedule_filename)r¨   rÞ   r"   r8   r½   r.   r.   r/   r8   ó  s   zPersistentScheduler.__init__c              	   C   sL   | j D ] }t tj¡� t | j| ¡ W d   ƒ n1 sw   Y  qd S r:   )Úknown_suffixesr   Úignore_errnoÚerrnoÚENOENTÚosÚremoverÞ   )r6   Úsuffixr.   r.   r/   Ú
_remove_db÷  s   
ÿ€ÿzPersistentScheduler._remove_dbc                 C   s   | j j| jdd�S )NT)Ú	writeback)ÚpersistenceÚopenrÞ   r<   r.   r.   r/   Ú_open_scheduleü  r¾   z"PersistentScheduler._open_schedulec                 C   s"   t d| j|dd� |  ¡  |  ¡ S )Nz'Removing corrupted schedule file %r: %rTr   )r„   rÞ   ræ   rê   )r6   r‰   r.   r.   r/   Ú _destroy_open_corrupted_scheduleÿ  s
   ÿz4PersistentScheduler._destroy_open_corrupted_schedulec              
   C   sb  z|   ¡ | _| j ¡  W n ty$ } z|  |¡| _W Y d }~nd }~ww |  ¡  | jjj}| j 	t
dƒ¡}|d urI||krItd||ƒ | j ¡  | jjj}| j 	t
dƒ¡}|d urr||krrdddœ}td|| || ƒ | j ¡  | j t
dƒi ¡}|  | jjj¡ |  | j¡ | j t
d	ƒtt
dƒ|t
dƒ|i¡ |  ¡  td
d dd„ t|ƒD ƒ¡ ƒ d S )NÚtzz%Reset: Timezone changed from %r to %rÚutc_enabledÚenabledÚdisabled)TFz Reset: UTC changed from %s to %sr{   r   zCurrent schedule:
Ú
c                 s   s   � | ]}t |ƒV  qd S r:   )Úrepr)r²   r)   r.   r.   r/   Ú	<genexpr>(  s   € 
ÿz5PersistentScheduler.setup_schedule.<locals>.<genexpr>)rê   Ú_storer§   rƒ   rë   Ú_create_scheduler?   rh   Útimezoner¨   ÚstrÚwarningÚclearÚ
enable_utcÚ
setdefaultr¿   rÀ   r|   rC   rN   r   rÁ   r‡   Újoinr   )r6   r‰   rì   Ú	stored_tzÚutcÚ
stored_utcÚchoicesr{   r.   r.   r/   rq     sB   
€ÿ



ÿ
ýÿz"PersistentScheduler.setup_schedulec                 C   sî   dD ]r}z	| j tdƒ  W n. ty;   z	i | j tdƒ< W n ty6 } z|  |¡| _ W Y d }~Y qd }~ww Y  d S w tdƒ| j vrOtdƒ | j  ¡   d S tdƒ| j vrbtdƒ | j  ¡   d S tdƒ| j vrrtdƒ | j  ¡   d S d S )	N)r   rš   r{   r   z+DB Reset: Account for new __version__ fieldrì   z"DB Reset: Account for new tz fieldrí   z+DB Reset: Account for new utc_enabled field)ró   rö   ÚKeyErrorrë   r÷   rø   )r6   Ú_r‰   r.   r.   r/   rô   +  s6   €þÿï
ú
ý
ìz$PersistentScheduler._create_schedulec                 C   s   | j tdƒ S ©Nr{   ©ró   rö   r<   r.   r.   r/   rÔ   B  s   z PersistentScheduler.get_schedulec                 C   s   || j tdƒ< d S r  r  rÕ   r.   r.   r/   rÖ   E  r¾   z PersistentScheduler.set_schedulec                 C   s   | j d ur| j  ¡  d S d S r:   )ró   rÁ   r<   r.   r.   r/   rÁ   I  s   
ÿzPersistentScheduler.syncc                 C   s   |   ¡  | j ¡  d S r:   )rÁ   ró   rÂ   r<   r.   r.   r/   rÂ   M  s   zPersistentScheduler.closec                 C   s   dj | d�S )Nz$    . db -> {self.schedule_filename}r<   )rV   r<   r.   r.   r/   r�   Q  s   zPersistentScheduler.info)r*   r+   r,   r-   Úshelverè   rß   ró   r8   ræ   rê   rë   rq   rô   rÔ   rÖ   rÝ   rC   rÁ   rÂ   r�   r.   r.   r.   r/   r#   ë  s$    &
r#   c                   @   s`   e Zd ZdZeZ		ddd„Zdd„ Zddd	„Zd
d„ Z	ddd„Z
		ddd„Zedd„ ƒZdS )r$   zCelery periodic task service.Nc                 C   sB   || _ |p|jj| _|p| j| _|p|jj| _tƒ | _tƒ | _	d S r:   )
r?   rh   ri   rj   Úscheduler_clsÚbeat_schedule_filenamerÞ   r   Ú_is_shutdownÚ_is_stopped)r6   r?   rj   rÞ   r  r.   r.   r/   r8   [  s   ÿ
ÿzService.__init__c                 C   s   | j | j| j| j| jffS r:   )rI   rj   rÞ   r  r?   r<   r.   r.   r/   rL   g  s   ÿzService.__reduce__Fc              	   C   sì   t dƒ tdt| jjƒƒ tjj| d� |r"tjj| d� t	 
d¡ zNz/| j ¡ sQ| j ¡ }|rL|dkrLtdt|dd�ƒ t |¡ | j ¡ rL| j ¡  | j ¡ r)W n ttfyb   | j ¡  Y nw W |  ¡  d S W |  ¡  d S |  ¡  w )	Nzbeat: Starting...z#beat: Ticking with max interval->%s)Úsenderzcelery beatg        zbeat: Waking up %s.zin )Úprefix)r�   r‡   r   Ú	schedulerrj   r   Ú	beat_initÚsendÚbeat_embedded_initr   Úset_process_titler  Úis_setr¥   r'   Úsleepr°   r¸   ÚKeyboardInterruptÚ
SystemExitr¦   rÁ   )r6   Úembedded_processrÍ   r.   r.   r/   Ústartk  s6   
ÿ



ÿ



ù€ÿ€þzService.startc                 C   ó   | j  ¡  | j ¡  d S r:   )r  rÂ   r  r¦   r<   r.   r.   r/   rÁ   ƒ  ó   
zService.syncc                 C   s*   t dƒ | j ¡  |o| j ¡  d S  d S )Nzbeat: Shutting down...)r�   r  r¦   r  Úwait)r6   r  r.   r.   r/   Ústop‡  s   
zService.stopúcelery.beat_schedulersc                 C   s4   | j }tt|ƒp	i ƒ}t| j|d�| j|| j|d�S )N)Úaliases)r?   rÞ   rj   rr   )rÞ   rJ   r   r   r  r?   rj   )r6   rr   Úextension_namespaceÚfilenamer  r.   r.   r/   Úget_schedulerŒ  s   
ÿüzService.get_schedulerc                 C   r9   r:   )r  r<   r.   r.   r/   r  ˜  s   zService.scheduler)NNN)F)Fr  )r*   r+   r,   r-   r#   r  r8   rL   r  rÁ   r  r  r   r  r.   r.   r.   r/   r$   V  s    
ÿ


ÿr$   c                       s0   e Zd ZdZ‡ fdd„Zdd„ Zdd„ Z‡  ZS )Ú	_Threadedz(Embedded task scheduler using threading.c                    s6   t t| ƒ ¡  || _t|fi |¤Ž| _d| _d| _d S )NTÚBeat)Úsuperr  r8   r?   r$   ÚserviceÚdaemonr@   ©r6   r?   r2   ©rI   r.   r/   r8      s
   
z_Threaded.__init__c                 C   r  r:   )r?   Úset_currentr"  r  r<   r.   r.   r/   Úrun§  r  z_Threaded.runc                 C   s   | j jdd� d S )NT)r  )r"  r  r<   r.   r.   r/   r  «  r¾   z_Threaded.stop)r*   r+   r,   r-   r8   r'  r  Ú__classcell__r.   r.   r%  r/   r  �  s
    r  c                       s,   e Zd Z‡ fdd„Zdd„ Zdd„ Z‡  ZS )Ú_Processc                    s0   t t| ƒ ¡  || _t|fi |¤Ž| _d| _d S )Nr   )r!  r)  r8   r?   r$   r"  r@   r$  r%  r.   r/   r8   ¶  s   
z_Process.__init__c                 C   sP   t dd� t tjtjtjgttƒ ƒ ¡ | j	 
¡  | j	 ¡  | jjdd� d S )NF)ÚfullT)r  )r
   r   Úclose_open_fdsrº   Ú	__stdin__Ú
__stdout__Ú
__stderr__Úlistr   r?   Úset_defaultr&  r"  r  r<   r.   r.   r/   r'  ¼  s   
ÿþ

z_Process.runc                 C   s   | j  ¡  |  ¡  d S r:   )r"  r  Ú	terminater<   r.   r.   r/   r  Å  s   
z_Process.stop)r*   r+   r,   r8   r'  r  r(  r.   r.   r%  r/   r)  ´  s    	r)  c                 K   s<   |  dd¡s
tdu rt| fddi|¤ŽS t| fd|i|¤ŽS )z»Return embedded clock service.

    Arguments:
        thread (bool): Run threaded instead of as a separate process.
            Uses :mod:`multiprocessing` by default, if available.
    ÚthreadFNrj   r   )rÈ   r)  r  )r?   rj   r2   r.   r.   r/   r%   Ê  s   r%   r:   )Kr-   Ú
__future__r   r   rœ   rá   rÜ   rã   r  rº   r'   r…   Úcalendarr   Úcollectionsr   Ú	functoolsr   Ú	threadingr   r   Úbilliardr	   Úbilliard.commonr
   Úbilliard.contextr   Úkombu.utils.functionalr   r   Úkombu.utils.objectsr   rÙ   r   r   r   Úfiver   r   r   r   r   Ú	schedulesr   r   Úutils.importsr   r   Ú	utils.logr   r   Ú
utils.timer   r   Ú__all__r&   r*   rÛ   r‡   r�   r„   r÷   rÚ   rƒ   r    Úobjectr0   r!   r"   r#   r$   r  ÚNotImplementedErrorr)  r%   r.   r.   r.   r/   Ú<module>   sd   
ÿ  kG
ÿ