Remove submodule, just put Dependencies in ./libs
This commit is contained in:
@@ -0,0 +1,176 @@
|
||||
"""
|
||||
This module contains the expressions applicable for CronTrigger's fields.
|
||||
"""
|
||||
from calendar import monthrange
|
||||
import re
|
||||
|
||||
from apscheduler.util import asint
|
||||
|
||||
__all__ = ('AllExpression', 'RangeExpression', 'WeekdayRangeExpression',
|
||||
'WeekdayPositionExpression')
|
||||
|
||||
WEEKDAYS = ['mon', 'tue', 'wed', 'thu', 'fri', 'sat', 'sun']
|
||||
|
||||
|
||||
class AllExpression(object):
|
||||
value_re = re.compile(r'\*(?:/(?P<step>\d+))?$')
|
||||
|
||||
def __init__(self, step=None):
|
||||
self.step = asint(step)
|
||||
if self.step == 0:
|
||||
raise ValueError('Increment must be higher than 0')
|
||||
|
||||
def get_next_value(self, date, field):
|
||||
start = field.get_value(date)
|
||||
minval = field.get_min(date)
|
||||
maxval = field.get_max(date)
|
||||
start = max(start, minval)
|
||||
|
||||
if not self.step:
|
||||
next = start
|
||||
else:
|
||||
distance_to_next = (self.step - (start - minval)) % self.step
|
||||
next = start + distance_to_next
|
||||
|
||||
if next <= maxval:
|
||||
return next
|
||||
|
||||
def __str__(self):
|
||||
if self.step:
|
||||
return '*/%d' % self.step
|
||||
return '*'
|
||||
|
||||
def __repr__(self):
|
||||
return "%s(%s)" % (self.__class__.__name__, self.step)
|
||||
|
||||
|
||||
class RangeExpression(AllExpression):
|
||||
value_re = re.compile(
|
||||
r'(?P<first>\d+)(?:-(?P<last>\d+))?(?:/(?P<step>\d+))?$')
|
||||
|
||||
def __init__(self, first, last=None, step=None):
|
||||
AllExpression.__init__(self, step)
|
||||
first = asint(first)
|
||||
last = asint(last)
|
||||
if last is None and step is None:
|
||||
last = first
|
||||
if last is not None and first > last:
|
||||
raise ValueError('The minimum value in a range must not be '
|
||||
'higher than the maximum')
|
||||
self.first = first
|
||||
self.last = last
|
||||
|
||||
def get_next_value(self, date, field):
|
||||
start = field.get_value(date)
|
||||
minval = field.get_min(date)
|
||||
maxval = field.get_max(date)
|
||||
|
||||
# Apply range limits
|
||||
minval = max(minval, self.first)
|
||||
if self.last is not None:
|
||||
maxval = min(maxval, self.last)
|
||||
start = max(start, minval)
|
||||
|
||||
if not self.step:
|
||||
next = start
|
||||
else:
|
||||
distance_to_next = (self.step - (start - minval)) % self.step
|
||||
next = start + distance_to_next
|
||||
|
||||
if next <= maxval:
|
||||
return next
|
||||
|
||||
def __str__(self):
|
||||
if self.last != self.first and self.last is not None:
|
||||
range = '%d-%d' % (self.first, self.last)
|
||||
else:
|
||||
range = str(self.first)
|
||||
|
||||
if self.step:
|
||||
return '%s/%d' % (range, self.step)
|
||||
return range
|
||||
|
||||
def __repr__(self):
|
||||
args = [str(self.first)]
|
||||
if self.last != self.first and self.last is not None or self.step:
|
||||
args.append(str(self.last))
|
||||
if self.step:
|
||||
args.append(str(self.step))
|
||||
return "%s(%s)" % (self.__class__.__name__, ', '.join(args))
|
||||
|
||||
|
||||
class WeekdayRangeExpression(RangeExpression):
|
||||
value_re = re.compile(r'(?P<first>[a-z]+)(?:-(?P<last>[a-z]+))?',
|
||||
re.IGNORECASE)
|
||||
|
||||
def __init__(self, first, last=None):
|
||||
try:
|
||||
first_num = WEEKDAYS.index(first.lower())
|
||||
except ValueError:
|
||||
raise ValueError('Invalid weekday name "%s"' % first)
|
||||
|
||||
if last:
|
||||
try:
|
||||
last_num = WEEKDAYS.index(last.lower())
|
||||
except ValueError:
|
||||
raise ValueError('Invalid weekday name "%s"' % last)
|
||||
else:
|
||||
last_num = None
|
||||
|
||||
RangeExpression.__init__(self, first_num, last_num)
|
||||
|
||||
def __str__(self):
|
||||
if self.last != self.first and self.last is not None:
|
||||
return '%s-%s' % (WEEKDAYS[self.first], WEEKDAYS[self.last])
|
||||
return WEEKDAYS[self.first]
|
||||
|
||||
def __repr__(self):
|
||||
args = ["'%s'" % WEEKDAYS[self.first]]
|
||||
if self.last != self.first and self.last is not None:
|
||||
args.append("'%s'" % WEEKDAYS[self.last])
|
||||
return "%s(%s)" % (self.__class__.__name__, ', '.join(args))
|
||||
|
||||
|
||||
class WeekdayPositionExpression(AllExpression):
|
||||
options = ['1st', '2nd', '3rd', '4th', '5th', 'last']
|
||||
value_re = re.compile(r'(?P<option_name>%s) +(?P<weekday_name>(?:\d+|\w+))'
|
||||
% '|'.join(options), re.IGNORECASE)
|
||||
|
||||
def __init__(self, option_name, weekday_name):
|
||||
try:
|
||||
self.option_num = self.options.index(option_name.lower())
|
||||
except ValueError:
|
||||
raise ValueError('Invalid weekday position "%s"' % option_name)
|
||||
|
||||
try:
|
||||
self.weekday = WEEKDAYS.index(weekday_name.lower())
|
||||
except ValueError:
|
||||
raise ValueError('Invalid weekday name "%s"' % weekday_name)
|
||||
|
||||
def get_next_value(self, date, field):
|
||||
# Figure out the weekday of the month's first day and the number
|
||||
# of days in that month
|
||||
first_day_wday, last_day = monthrange(date.year, date.month)
|
||||
|
||||
# Calculate which day of the month is the first of the target weekdays
|
||||
first_hit_day = self.weekday - first_day_wday + 1
|
||||
if first_hit_day <= 0:
|
||||
first_hit_day += 7
|
||||
|
||||
# Calculate what day of the month the target weekday would be
|
||||
if self.option_num < 5:
|
||||
target_day = first_hit_day + self.option_num * 7
|
||||
else:
|
||||
target_day = first_hit_day + ((last_day - first_hit_day) / 7) * 7
|
||||
|
||||
if target_day <= last_day and target_day >= date.day:
|
||||
return target_day
|
||||
|
||||
def __str__(self):
|
||||
return '%s %s' % (self.options[self.option_num],
|
||||
WEEKDAYS[self.weekday])
|
||||
|
||||
def __repr__(self):
|
||||
return "%s('%s', '%s')" % (self.__class__.__name__,
|
||||
self.options[self.option_num],
|
||||
WEEKDAYS[self.weekday])
|
||||
@@ -0,0 +1,92 @@
|
||||
"""
|
||||
Fields represent :class:`~apscheduler.triggers.CronTrigger` options which map
|
||||
to :class:`~datetime.datetime` fields.
|
||||
"""
|
||||
from calendar import monthrange
|
||||
|
||||
from apscheduler.expressions import *
|
||||
|
||||
__all__ = ('BaseField', 'WeekField', 'DayOfMonthField', 'DayOfWeekField')
|
||||
|
||||
MIN_VALUES = {'year': 1970, 'month': 1, 'day': 1, 'week': 1,
|
||||
'day_of_week': 0, 'hour': 0, 'minute': 0, 'second': 0}
|
||||
MAX_VALUES = {'year': 2 ** 63, 'month': 12, 'day:': 31, 'week': 53,
|
||||
'day_of_week': 6, 'hour': 23, 'minute': 59, 'second': 59}
|
||||
|
||||
class BaseField(object):
|
||||
REAL = True
|
||||
COMPILERS = [AllExpression, RangeExpression]
|
||||
|
||||
def __init__(self, name, exprs):
|
||||
self.name = name
|
||||
self.compile_expressions(exprs)
|
||||
|
||||
def get_min(self, dateval):
|
||||
return MIN_VALUES[self.name]
|
||||
|
||||
def get_max(self, dateval):
|
||||
return MAX_VALUES[self.name]
|
||||
|
||||
def get_value(self, dateval):
|
||||
return getattr(dateval, self.name)
|
||||
|
||||
def get_next_value(self, dateval):
|
||||
smallest = None
|
||||
for expr in self.expressions:
|
||||
value = expr.get_next_value(dateval, self)
|
||||
if smallest is None or (value is not None and value < smallest):
|
||||
smallest = value
|
||||
|
||||
return smallest
|
||||
|
||||
def compile_expressions(self, exprs):
|
||||
self.expressions = []
|
||||
|
||||
# Split a comma-separated expression list, if any
|
||||
exprs = str(exprs).strip()
|
||||
if ',' in exprs:
|
||||
for expr in exprs.split(','):
|
||||
self.compile_expression(expr)
|
||||
else:
|
||||
self.compile_expression(exprs)
|
||||
|
||||
def compile_expression(self, expr):
|
||||
for compiler in self.COMPILERS:
|
||||
match = compiler.value_re.match(expr)
|
||||
if match:
|
||||
compiled_expr = compiler(**match.groupdict())
|
||||
self.expressions.append(compiled_expr)
|
||||
return
|
||||
|
||||
raise ValueError('Unrecognized expression "%s" for field "%s"' %
|
||||
(expr, self.name))
|
||||
|
||||
def __str__(self):
|
||||
expr_strings = (str(e) for e in self.expressions)
|
||||
return ','.join(expr_strings)
|
||||
|
||||
def __repr__(self):
|
||||
return "%s('%s', '%s')" % (self.__class__.__name__, self.name,
|
||||
str(self))
|
||||
|
||||
|
||||
class WeekField(BaseField):
|
||||
REAL = False
|
||||
|
||||
def get_value(self, dateval):
|
||||
return dateval.isocalendar()[1]
|
||||
|
||||
|
||||
class DayOfMonthField(BaseField):
|
||||
COMPILERS = BaseField.COMPILERS + [WeekdayPositionExpression]
|
||||
|
||||
def get_max(self, dateval):
|
||||
return monthrange(dateval.year, dateval.month)[1]
|
||||
|
||||
|
||||
class DayOfWeekField(BaseField):
|
||||
REAL = False
|
||||
COMPILERS = BaseField.COMPILERS + [WeekdayRangeExpression]
|
||||
|
||||
def get_value(self, dateval):
|
||||
return dateval.weekday()
|
||||
@@ -0,0 +1,407 @@
|
||||
"""
|
||||
This module is the main part of the library, and is the only module that
|
||||
regular users should be concerned with.
|
||||
"""
|
||||
from threading import Thread, Event, Lock
|
||||
from datetime import datetime, timedelta
|
||||
from logging import getLogger
|
||||
import os
|
||||
|
||||
from apscheduler.util import time_difference, asbool
|
||||
from apscheduler.triggers import DateTrigger, IntervalTrigger, CronTrigger
|
||||
|
||||
|
||||
logger = getLogger(__name__)
|
||||
|
||||
|
||||
class Job(object):
|
||||
"""
|
||||
Represents a task scheduled in the scheduler.
|
||||
"""
|
||||
|
||||
def __init__(self, trigger, func, args, kwargs):
|
||||
self.thread = None
|
||||
self.trigger = trigger
|
||||
self.func = func
|
||||
self.args = args
|
||||
self.kwargs = kwargs
|
||||
if hasattr(func, '__name__'):
|
||||
self.name = func.__name__
|
||||
else:
|
||||
self.name = str(func)
|
||||
|
||||
def run(self):
|
||||
"""
|
||||
Starts the execution of this job in a separate thread.
|
||||
"""
|
||||
if (self.thread and self.thread.isAlive()):
|
||||
logger.info('Skipping run of job %s (previously triggered '
|
||||
'instance is still running)', self)
|
||||
else:
|
||||
self.thread = Thread(target=self.run_in_thread)
|
||||
self.thread.setDaemon(False)
|
||||
self.thread.start()
|
||||
|
||||
def run_in_thread(self):
|
||||
"""
|
||||
Runs the associated callable.
|
||||
This method is executed in a dedicated thread.
|
||||
"""
|
||||
try:
|
||||
self.func(*self.args, **self.kwargs)
|
||||
except:
|
||||
logger.exception('Error executing job "%s"', self)
|
||||
raise
|
||||
|
||||
def __str__(self):
|
||||
return '%s: %s' % (self.name, repr(self.trigger))
|
||||
|
||||
def __repr__(self):
|
||||
return '%s(%s, %s)' % (self.__class__.__name__, self.name,
|
||||
repr(self.trigger))
|
||||
|
||||
|
||||
class SchedulerShutdownError(Exception):
|
||||
"""
|
||||
Thrown when attempting to use the scheduler after
|
||||
it's been shut down.
|
||||
"""
|
||||
|
||||
def __init__(self):
|
||||
Exception.__init__(self, 'Scheduler has already been shut down')
|
||||
|
||||
|
||||
class SchedulerAlreadyRunningError(Exception):
|
||||
"""
|
||||
Thrown when attempting to start the scheduler, but it's already running.
|
||||
"""
|
||||
|
||||
def __init__(self):
|
||||
Exception.__init__(self, 'Scheduler is already running')
|
||||
|
||||
|
||||
class Scheduler(object):
|
||||
"""
|
||||
This class is responsible for scheduling jobs and triggering
|
||||
their execution.
|
||||
"""
|
||||
|
||||
stopped = False
|
||||
thread = None
|
||||
misfire_grace_time = 1
|
||||
daemonic = True
|
||||
|
||||
def __init__(self, **config):
|
||||
self.jobs = []
|
||||
self.jobs_lock = Lock()
|
||||
self.wakeup = Event()
|
||||
self.configure(config)
|
||||
|
||||
def configure(self, config):
|
||||
"""
|
||||
Updates the configuration with the given options.
|
||||
"""
|
||||
for key, val in config.items():
|
||||
if key.startswith('apscheduler.'):
|
||||
key = key[12:]
|
||||
if key == 'misfire_grace_time':
|
||||
self.misfire_grace_time = int(val)
|
||||
elif key == 'daemonic':
|
||||
self.daemonic = asbool(val)
|
||||
|
||||
def start(self):
|
||||
"""
|
||||
Starts the scheduler in a new thread.
|
||||
"""
|
||||
if self.thread and self.thread.isAlive():
|
||||
raise SchedulerAlreadyRunningError
|
||||
|
||||
self.stopped = False
|
||||
self.thread = Thread(target=self.run, name='APScheduler')
|
||||
self.thread.setDaemon(self.daemonic)
|
||||
self.thread.start()
|
||||
logger.info('Scheduler started')
|
||||
|
||||
def shutdown(self, timeout=0):
|
||||
"""
|
||||
Shuts down the scheduler and terminates the thread.
|
||||
Does not terminate any currently running jobs.
|
||||
|
||||
:param timeout: time (in seconds) to wait for the scheduler thread to
|
||||
terminate, 0 to wait forever, None to skip waiting
|
||||
"""
|
||||
if self.stopped or not self.thread.isAlive():
|
||||
return
|
||||
|
||||
logger.info('Scheduler shutting down')
|
||||
self.stopped = True
|
||||
self.wakeup.set()
|
||||
if timeout is not None:
|
||||
self.thread.join(timeout)
|
||||
self.jobs = []
|
||||
|
||||
def cron_schedule(self, year='*', month='*', day='*', week='*',
|
||||
day_of_week='*', hour='*', minute='*', second='*',
|
||||
args=None, kwargs=None):
|
||||
"""
|
||||
Decorator that causes its host function to be scheduled
|
||||
according to the given parameters.
|
||||
This decorator does not wrap its host function.
|
||||
The scheduled function will be called without any arguments.
|
||||
See :meth:`add_cron_job` for more information.
|
||||
"""
|
||||
def inner(func):
|
||||
self.add_cron_job(func, year, month, day, week, day_of_week, hour,
|
||||
minute, second, args, kwargs)
|
||||
return func
|
||||
return inner
|
||||
|
||||
def interval_schedule(self, weeks=0, days=0, hours=0, minutes=0, seconds=0,
|
||||
start_date=None, repeat=0, args=None, kwargs=None):
|
||||
"""
|
||||
Decorator that causes its host function to be scheduled
|
||||
for execution on specified intervals.
|
||||
This decorator does not wrap its host function.
|
||||
The scheduled function will be called without any arguments.
|
||||
Note that the default repeat value is 0, which means to repeat forever.
|
||||
See :meth:`add_delayed_job` for more information.
|
||||
"""
|
||||
def inner(func):
|
||||
self.add_interval_job(func, weeks, days, hours, minutes, seconds,
|
||||
start_date, repeat, args, kwargs)
|
||||
return func
|
||||
return inner
|
||||
|
||||
def _add_job(self, trigger, func, args, kwargs):
|
||||
"""
|
||||
Adds a Job to the job list and notifies the scheduler thread.
|
||||
|
||||
:param trigger: trigger for the given callable
|
||||
:param args: list of positional arguments to call func with
|
||||
:param kwargs: dict of keyword arguments to call func with
|
||||
:return: the scheduled job
|
||||
:rtype: Job
|
||||
"""
|
||||
if self.stopped:
|
||||
raise SchedulerShutdownError
|
||||
if not hasattr(func, '__call__'):
|
||||
raise TypeError('func must be callable')
|
||||
|
||||
if args is None:
|
||||
args = []
|
||||
if kwargs is None:
|
||||
kwargs = {}
|
||||
|
||||
job = Job(trigger, func, args, kwargs)
|
||||
self.jobs_lock.acquire()
|
||||
try:
|
||||
self.jobs.append(job)
|
||||
finally:
|
||||
self.jobs_lock.release()
|
||||
logger.info('Added job "%s"', job)
|
||||
|
||||
# Notify the scheduler about the new job
|
||||
self.wakeup.set()
|
||||
|
||||
return job
|
||||
|
||||
def add_date_job(self, func, date, args=None, kwargs=None):
|
||||
"""
|
||||
Adds a job to be completed on a specific date and time.
|
||||
|
||||
:param func: callable to run
|
||||
:param args: positional arguments to call func with
|
||||
:param kwargs: keyword arguments to call func with
|
||||
"""
|
||||
trigger = DateTrigger(date)
|
||||
return self._add_job(trigger, func, args, kwargs)
|
||||
|
||||
def add_interval_job(self, func, weeks=0, days=0, hours=0, minutes=0,
|
||||
seconds=0, start_date=None, repeat=0, args=None,
|
||||
kwargs=None):
|
||||
"""
|
||||
Adds a job to be completed on specified intervals.
|
||||
|
||||
:param func: callable to run
|
||||
:param weeks: number of weeks to wait
|
||||
:param days: number of days to wait
|
||||
:param hours: number of hours to wait
|
||||
:param minutes: number of minutes to wait
|
||||
:param seconds: number of seconds to wait
|
||||
:param start_date: when to first execute the job and start the
|
||||
counter (default is after the given interval)
|
||||
:param repeat: number of times the job will be run (0 = repeat
|
||||
indefinitely)
|
||||
:param args: list of positional arguments to call func with
|
||||
:param kwargs: dict of keyword arguments to call func with
|
||||
"""
|
||||
interval = timedelta(weeks=weeks, days=days, hours=hours,
|
||||
minutes=minutes, seconds=seconds)
|
||||
trigger = IntervalTrigger(interval, repeat, start_date)
|
||||
return self._add_job(trigger, func, args, kwargs)
|
||||
|
||||
def add_cron_job(self, func, year='*', month='*', day='*', week='*',
|
||||
day_of_week='*', hour='*', minute='*', second='*',
|
||||
args=None, kwargs=None):
|
||||
"""
|
||||
Adds a job to be completed on times that match the given expressions.
|
||||
|
||||
:param func: callable to run
|
||||
:param year: year to run on
|
||||
:param month: month to run on (0 = January)
|
||||
:param day: day of month to run on
|
||||
:param week: week of the year to run on
|
||||
:param day_of_week: weekday to run on (0 = Monday)
|
||||
:param hour: hour to run on
|
||||
:param second: second to run on
|
||||
:param args: list of positional arguments to call func with
|
||||
:param kwargs: dict of keyword arguments to call func with
|
||||
:return: the scheduled job
|
||||
:rtype: Job
|
||||
"""
|
||||
trigger = CronTrigger(year=year, month=month, day=day, week=week,
|
||||
day_of_week=day_of_week, hour=hour,
|
||||
minute=minute, second=second)
|
||||
return self._add_job(trigger, func, args, kwargs)
|
||||
|
||||
def is_job_active(self, job):
|
||||
"""
|
||||
Determines if the given job is still on the job list.
|
||||
|
||||
:return: True if the job is still active, False if not
|
||||
"""
|
||||
self.jobs_lock.acquire()
|
||||
try:
|
||||
return job in self.jobs
|
||||
finally:
|
||||
self.jobs_lock.release()
|
||||
|
||||
def unschedule_job(self, job):
|
||||
"""
|
||||
Removes a job, preventing it from being fired any more.
|
||||
"""
|
||||
self.jobs_lock.acquire()
|
||||
try:
|
||||
self.jobs.remove(job)
|
||||
finally:
|
||||
self.jobs_lock.release()
|
||||
logger.info('Removed job "%s"', job)
|
||||
self.wakeup.set()
|
||||
|
||||
def unschedule_func(self, func):
|
||||
"""
|
||||
Removes all jobs that would execute the given function.
|
||||
"""
|
||||
self.jobs_lock.acquire()
|
||||
try:
|
||||
remove_list = [job for job in self.jobs if job.func == func]
|
||||
for job in remove_list:
|
||||
self.jobs.remove(job)
|
||||
logger.info('Removed job "%s"', job)
|
||||
finally:
|
||||
self.jobs_lock.release()
|
||||
|
||||
# Have the scheduler calculate a new wakeup time
|
||||
self.wakeup.set()
|
||||
|
||||
def dump_jobs(self):
|
||||
"""
|
||||
Gives a textual listing of all jobs currently scheduled on this
|
||||
scheduler.
|
||||
|
||||
:rtype: str
|
||||
"""
|
||||
job_strs = []
|
||||
now = datetime.now()
|
||||
self.jobs_lock.acquire()
|
||||
try:
|
||||
for job in self.jobs:
|
||||
next_fire_time = job.trigger.get_next_fire_time(now)
|
||||
job_str = '%s (next fire time: %s)' % (str(job),
|
||||
next_fire_time)
|
||||
job_strs.append(job_str)
|
||||
finally:
|
||||
self.jobs_lock.release()
|
||||
|
||||
if job_strs:
|
||||
return os.linesep.join(job_strs)
|
||||
return 'No jobs currently scheduled.'
|
||||
|
||||
def _get_next_wakeup_time(self, now):
|
||||
"""
|
||||
Determines the time of the next job execution, and removes finished
|
||||
jobs.
|
||||
|
||||
:param now: the result of datetime.now(), generated elsewhere for
|
||||
consistency.
|
||||
"""
|
||||
next_wakeup = None
|
||||
finished_jobs = []
|
||||
|
||||
self.jobs_lock.acquire()
|
||||
try:
|
||||
for job in self.jobs:
|
||||
next_run = job.trigger.get_next_fire_time(now)
|
||||
if next_run is None:
|
||||
finished_jobs.append(job)
|
||||
elif next_run and (next_wakeup is None or \
|
||||
next_run < next_wakeup):
|
||||
next_wakeup = next_run
|
||||
|
||||
# Clear out any finished jobs
|
||||
for job in finished_jobs:
|
||||
self.jobs.remove(job)
|
||||
logger.info('Removed finished job "%s"', job)
|
||||
finally:
|
||||
self.jobs_lock.release()
|
||||
|
||||
return next_wakeup
|
||||
|
||||
def _get_current_jobs(self):
|
||||
"""
|
||||
Determines which jobs should be executed right now.
|
||||
"""
|
||||
current_jobs = []
|
||||
now = datetime.now()
|
||||
start = now - timedelta(seconds=self.misfire_grace_time)
|
||||
|
||||
self.jobs_lock.acquire()
|
||||
try:
|
||||
for job in self.jobs:
|
||||
next_run = job.trigger.get_next_fire_time(start)
|
||||
if next_run:
|
||||
time_diff = time_difference(now, next_run)
|
||||
if next_run < now and time_diff <= self.misfire_grace_time:
|
||||
current_jobs.append(job)
|
||||
finally:
|
||||
self.jobs_lock.release()
|
||||
|
||||
return current_jobs
|
||||
|
||||
def run(self):
|
||||
"""
|
||||
Runs the main loop of the scheduler.
|
||||
"""
|
||||
self.wakeup.clear()
|
||||
while not self.stopped:
|
||||
# Execute any jobs scheduled to be run right now
|
||||
for job in self._get_current_jobs():
|
||||
logger.debug('Executing job "%s"', job)
|
||||
job.run()
|
||||
|
||||
# Figure out when the next job should be run, and
|
||||
# adjust the wait time accordingly
|
||||
now = datetime.now()
|
||||
next_wakeup_time = self._get_next_wakeup_time(now)
|
||||
|
||||
# Sleep until the next job is scheduled to be run,
|
||||
# or a new job is added, or the scheduler is stopped
|
||||
if next_wakeup_time is not None:
|
||||
wait_seconds = time_difference(next_wakeup_time, now)
|
||||
logger.debug('Next wakeup is due at %s (in %f seconds)',
|
||||
next_wakeup_time, wait_seconds)
|
||||
self.wakeup.wait(wait_seconds)
|
||||
else:
|
||||
logger.debug('No jobs; waiting until a job is added')
|
||||
self.wakeup.wait()
|
||||
self.wakeup.clear()
|
||||
@@ -0,0 +1,171 @@
|
||||
"""
|
||||
Triggers determine the times when a job should be executed.
|
||||
"""
|
||||
from datetime import datetime, timedelta
|
||||
from math import ceil
|
||||
|
||||
from apscheduler.fields import *
|
||||
from apscheduler.util import *
|
||||
|
||||
__all__ = ('CronTrigger', 'DateTrigger', 'IntervalTrigger')
|
||||
|
||||
|
||||
class CronTrigger(object):
|
||||
FIELD_NAMES = ('year', 'month', 'day', 'week', 'day_of_week', 'hour',
|
||||
'minute', 'second')
|
||||
FIELDS_MAP = {'year': BaseField,
|
||||
'month': BaseField,
|
||||
'week': WeekField,
|
||||
'day': DayOfMonthField,
|
||||
'day_of_week': DayOfWeekField,
|
||||
'hour': BaseField,
|
||||
'minute': BaseField,
|
||||
'second': BaseField}
|
||||
|
||||
def __init__(self, **values):
|
||||
self.fields = []
|
||||
for field_name in self.FIELD_NAMES:
|
||||
exprs = values.get(field_name) or '*'
|
||||
field_class = self.FIELDS_MAP[field_name]
|
||||
field = field_class(field_name, exprs)
|
||||
self.fields.append(field)
|
||||
|
||||
def _increment_field_value(self, dateval, fieldnum):
|
||||
"""
|
||||
Increments the designated field and resets all less significant fields
|
||||
to their minimum values.
|
||||
|
||||
:type dateval: datetime
|
||||
:type fieldnum: int
|
||||
:type amount: int
|
||||
:rtype: tuple
|
||||
:return: a tuple containing the new date, and the number of the field
|
||||
that was actually incremented
|
||||
"""
|
||||
i = 0
|
||||
values = {}
|
||||
while i < len(self.fields):
|
||||
field = self.fields[i]
|
||||
if not field.REAL:
|
||||
if i == fieldnum:
|
||||
fieldnum -= 1
|
||||
i -= 1
|
||||
else:
|
||||
i += 1
|
||||
continue
|
||||
|
||||
if i < fieldnum:
|
||||
values[field.name] = field.get_value(dateval)
|
||||
i += 1
|
||||
elif i > fieldnum:
|
||||
values[field.name] = field.get_min(dateval)
|
||||
i += 1
|
||||
else:
|
||||
value = field.get_value(dateval)
|
||||
maxval = field.get_max(dateval)
|
||||
if value == maxval:
|
||||
fieldnum -= 1
|
||||
i -= 1
|
||||
else:
|
||||
values[field.name] = value + 1
|
||||
i += 1
|
||||
|
||||
return datetime(**values), fieldnum
|
||||
|
||||
def _set_field_value(self, dateval, fieldnum, new_value):
|
||||
values = {}
|
||||
for i, field in enumerate(self.fields):
|
||||
if field.REAL:
|
||||
if i < fieldnum:
|
||||
values[field.name] = field.get_value(dateval)
|
||||
elif i > fieldnum:
|
||||
values[field.name] = field.get_min(dateval)
|
||||
else:
|
||||
values[field.name] = new_value
|
||||
|
||||
return datetime(**values)
|
||||
|
||||
def get_next_fire_time(self, start_date):
|
||||
next_date = datetime_ceil(start_date)
|
||||
fieldnum = 0
|
||||
while 0 <= fieldnum < len(self.fields):
|
||||
field = self.fields[fieldnum]
|
||||
curr_value = field.get_value(next_date)
|
||||
next_value = field.get_next_value(next_date)
|
||||
|
||||
if next_value is None:
|
||||
# No valid value was found
|
||||
next_date, fieldnum = self._increment_field_value(next_date,
|
||||
fieldnum - 1)
|
||||
elif next_value > curr_value:
|
||||
# A valid, but higher than the starting value, was found
|
||||
if field.REAL:
|
||||
next_date = self._set_field_value(next_date, fieldnum,
|
||||
next_value)
|
||||
fieldnum += 1
|
||||
else:
|
||||
next_date, fieldnum = self._increment_field_value(next_date,
|
||||
fieldnum)
|
||||
else:
|
||||
# A valid value was found, no changes necessary
|
||||
fieldnum += 1
|
||||
|
||||
if fieldnum >= 0:
|
||||
return next_date
|
||||
|
||||
def __repr__(self):
|
||||
field_reprs = ("%s='%s'" % (f.name, str(f)) for f in self.fields
|
||||
if str(f) != '*')
|
||||
return '%s(%s)' % (self.__class__.__name__, ', '.join(field_reprs))
|
||||
|
||||
|
||||
class DateTrigger(object):
|
||||
def __init__(self, run_date):
|
||||
self.run_date = convert_to_datetime(run_date)
|
||||
|
||||
def get_next_fire_time(self, start_date):
|
||||
if self.run_date >= start_date:
|
||||
return self.run_date
|
||||
|
||||
def __repr__(self):
|
||||
return '%s(%s)' % (self.__class__.__name__, repr(self.run_date))
|
||||
|
||||
|
||||
class IntervalTrigger(object):
|
||||
def __init__(self, interval, repeat, start_date=None):
|
||||
if not isinstance(interval, timedelta):
|
||||
raise TypeError('interval must be a timedelta')
|
||||
if repeat < 0:
|
||||
raise ValueError('Illegal value for repeat; expected >= 0, '
|
||||
'received %s' % repeat)
|
||||
|
||||
self.interval = interval
|
||||
self.interval_length = timedelta_seconds(self.interval)
|
||||
if self.interval_length == 0:
|
||||
self.interval = timedelta(seconds=1)
|
||||
self.interval_length = 1
|
||||
self.repeat = repeat
|
||||
if start_date is None:
|
||||
self.first_fire_date = datetime.now() + self.interval
|
||||
else:
|
||||
self.first_fire_date = convert_to_datetime(start_date)
|
||||
self.first_fire_date -= timedelta(microseconds=\
|
||||
self.first_fire_date.microsecond)
|
||||
if repeat > 0:
|
||||
self.last_fire_date = self.first_fire_date + interval * (repeat - 1)
|
||||
else:
|
||||
self.last_fire_date = None
|
||||
|
||||
def get_next_fire_time(self, start_date):
|
||||
if start_date < self.first_fire_date:
|
||||
return self.first_fire_date
|
||||
if self.last_fire_date and start_date > self.last_fire_date:
|
||||
return None
|
||||
timediff_seconds = timedelta_seconds(start_date - self.first_fire_date)
|
||||
next_interval_num = int(ceil(timediff_seconds / self.interval_length))
|
||||
return self.first_fire_date + self.interval * next_interval_num
|
||||
|
||||
def __repr__(self):
|
||||
return "%s(interval=%s, repeat=%d, start_date=%s)" % (
|
||||
self.__class__.__name__, repr(self.interval), self.repeat,
|
||||
repr(self.first_fire_date))
|
||||
@@ -0,0 +1,91 @@
|
||||
"""
|
||||
This module contains several handy functions primarily meant for internal use.
|
||||
"""
|
||||
|
||||
from datetime import date, datetime, timedelta
|
||||
from time import mktime
|
||||
|
||||
__all__ = ('asint', 'asbool', 'convert_to_datetime', 'timedelta_seconds',
|
||||
'time_difference', 'datetime_ceil')
|
||||
|
||||
|
||||
def asint(text):
|
||||
"""
|
||||
Safely converts a string to an integer, returning None if the string
|
||||
is None.
|
||||
|
||||
:type text: str
|
||||
:rtype: int
|
||||
"""
|
||||
if text is not None:
|
||||
return int(text)
|
||||
|
||||
|
||||
def asbool(obj):
|
||||
"""
|
||||
Interprets an object as a boolean value.
|
||||
|
||||
:rtype: bool
|
||||
"""
|
||||
if isinstance(obj, str):
|
||||
obj = obj.strip().lower()
|
||||
if obj in ('true', 'yes', 'on', 'y', 't', '1'):
|
||||
return True
|
||||
if obj in ('false', 'no', 'off', 'n', 'f', '0'):
|
||||
return False
|
||||
raise ValueError('Unable to interpret value "%s" as boolean' % obj)
|
||||
return bool(obj)
|
||||
|
||||
|
||||
def convert_to_datetime(dateval):
|
||||
"""
|
||||
Converts a date object to a datetime object.
|
||||
If an actual datetime object is passed, it is returned unmodified.
|
||||
|
||||
:type dateval: date
|
||||
:rtype: datetime
|
||||
"""
|
||||
if isinstance(dateval, datetime):
|
||||
return dateval
|
||||
elif isinstance(dateval, date):
|
||||
return datetime.fromordinal(dateval.toordinal())
|
||||
raise TypeError('Expected date, got %s instead' % type(dateval))
|
||||
|
||||
|
||||
def timedelta_seconds(delta):
|
||||
"""
|
||||
Converts the given timedelta to seconds.
|
||||
|
||||
:type delta: timedelta
|
||||
:rtype: float
|
||||
"""
|
||||
return delta.days * 24 * 60 * 60 + delta.seconds + \
|
||||
delta.microseconds / 1000000.0
|
||||
|
||||
|
||||
def time_difference(date1, date2):
|
||||
"""
|
||||
Returns the time difference in seconds between the given two
|
||||
datetime objects. The difference is calculated as: date1 - date2.
|
||||
|
||||
:param date1: the later datetime
|
||||
:type date1: datetime
|
||||
:param date2: the earlier datetime
|
||||
:type date2: datetime
|
||||
:rtype: float
|
||||
"""
|
||||
later = mktime(date1.timetuple())
|
||||
earlier = mktime(date2.timetuple())
|
||||
return int(later - earlier)
|
||||
|
||||
|
||||
def datetime_ceil(dateval):
|
||||
"""
|
||||
Rounds the given datetime object upwards.
|
||||
|
||||
:type dateval: datetime
|
||||
"""
|
||||
if dateval.microsecond > 0:
|
||||
return dateval + timedelta(seconds=1,
|
||||
microseconds=-dateval.microsecond)
|
||||
return dateval
|
||||
Reference in New Issue
Block a user