Refactor scheduler
This commit is contained in:
@@ -2,46 +2,31 @@ from __future__ import absolute_import
|
||||
|
||||
import logging
|
||||
import atexit
|
||||
import sys
|
||||
|
||||
from .runtime import scheduler
|
||||
|
||||
from django.db.models.signals import post_syncdb
|
||||
from django.dispatch import receiver
|
||||
|
||||
from south.signals import pre_migrate
|
||||
|
||||
from signaler.signals import pre_collectstatic
|
||||
from project_tools.api import register_tool
|
||||
from navigation.api import bind_links
|
||||
|
||||
from .links import scheduler_tool_link, scheduler_list, job_list
|
||||
from .literals import SHUTDOWN_COMMANDS
|
||||
from .api import LocalScheduler
|
||||
|
||||
from .links import job_list
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
# TODO: shutdown scheduler on pre_syncdb to avoid accessing non existing models
|
||||
|
||||
@receiver(post_syncdb, dispatch_uid='scheduler_shutdown_post_syncdb')
|
||||
def scheduler_shutdown_post_syncdb(sender, **kwargs):
|
||||
logger.debug('Scheduler shut down on post syncdb signal')
|
||||
scheduler.shutdown()
|
||||
|
||||
|
||||
@receiver(pre_collectstatic, dispatch_uid='sheduler_shutdown_pre_collectstatic')
|
||||
def sheduler_shutdown_pre_collectstatic(sender, **kwargs):
|
||||
logger.debug('Scheduler shut down on collectstatic signal')
|
||||
scheduler.shutdown()
|
||||
|
||||
|
||||
@receiver(pre_migrate, dispatch_uid='sheduler_shutdown_pre_migrate')
|
||||
def sheduler_shutdown_pre_migrate(sender, **kwargs):
|
||||
logger.debug('Scheduler shut down on pre_migrate signal')
|
||||
scheduler.shutdown()
|
||||
|
||||
|
||||
def schedule_shutdown_on_exit():
|
||||
logger.debug('Scheduler shut down on exit')
|
||||
scheduler.shutdown()
|
||||
logger.debug('Schedulers shut down on exit')
|
||||
LocalScheduler.shutdown_all()
|
||||
|
||||
|
||||
register_tool(job_list)
|
||||
if any([command in sys.argv for command in SHUTDOWN_COMMANDS]):
|
||||
logger.debug('Schedulers shut down on SHUTDOWN_COMMAND')
|
||||
# Shutdown any scheduler already running
|
||||
LocalScheduler.shutdown_all()
|
||||
# Prevent any new scheduler afterwards to start
|
||||
LocalScheduler.lockdown()
|
||||
|
||||
register_tool(scheduler_tool_link)
|
||||
atexit.register(schedule_shutdown_on_exit)
|
||||
bind_links([LocalScheduler, 'scheduler_list', 'job_list'], scheduler_list, menu_name='secondary_menu')
|
||||
bind_links([LocalScheduler], job_list)
|
||||
|
||||
@@ -1,30 +1,147 @@
|
||||
from __future__ import absolute_import
|
||||
|
||||
from .runtime import scheduler
|
||||
from .exceptions import AlreadyScheduled
|
||||
import logging
|
||||
|
||||
registered_jobs = {}
|
||||
from apscheduler.scheduler import Scheduler as OriginalScheduler
|
||||
|
||||
from django.utils.translation import ugettext_lazy as _
|
||||
|
||||
from .exceptions import AlreadyScheduled, UnknownJobClass
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
def register_interval_job(name, title, func, weeks=0, days=0, hours=0, minutes=0,
|
||||
seconds=0, start_date=None, args=None,
|
||||
kwargs=None, job_name=None, **options):
|
||||
class SchedulerJobBase(object):
|
||||
job_type = u''
|
||||
|
||||
if name in registered_jobs:
|
||||
raise AlreadyScheduled
|
||||
def __init__(self, name, label, function, *args, **kwargs):
|
||||
self.scheduler = None
|
||||
self.name = name
|
||||
self.label = label
|
||||
self.function = function
|
||||
self.args = args
|
||||
self.kwargs = kwargs
|
||||
|
||||
job = scheduler.add_interval_job(func=func, weeks=weeks, days=days,
|
||||
hours=hours, minutes=minutes, seconds=seconds,
|
||||
start_date=start_date, args=args, kwargs=kwargs, **options)
|
||||
def stop(self):
|
||||
self.scheduler.stop_job(self)
|
||||
|
||||
registered_jobs[name] = {'title': title, 'job': job}
|
||||
@property
|
||||
def running(self):
|
||||
if self.scheduler:
|
||||
return self.scheduler.running
|
||||
else:
|
||||
return False
|
||||
|
||||
@property
|
||||
def start_date(self):
|
||||
return self._job.trigger.start_date
|
||||
|
||||
|
||||
def remove_job(name):
|
||||
if name in registered_jobs:
|
||||
scheduler.unschedule_job(registered_jobs[name]['job'])
|
||||
registered_jobs.pop(name)
|
||||
class IntervalJob(SchedulerJobBase):
|
||||
job_type = _(u'Interval job')
|
||||
|
||||
def start(self, scheduler):
|
||||
scheduler.add_job(self)
|
||||
|
||||
|
||||
def get_job_list():
|
||||
return registered_jobs.values()
|
||||
class DateJob(SchedulerJobBase):
|
||||
job_type = _(u'Date job')
|
||||
|
||||
def start(self, scheduler):
|
||||
scheduler.add_job(self)
|
||||
|
||||
|
||||
class LocalScheduler(object):
|
||||
scheduler_registry = {}
|
||||
lockdown = False
|
||||
|
||||
@classmethod
|
||||
def get(cls, name):
|
||||
return cls.scheduler_registry[name]
|
||||
|
||||
@classmethod
|
||||
def get_all(cls):
|
||||
return cls.scheduler_registry.values()
|
||||
|
||||
@classmethod
|
||||
def shutdown_all(cls):
|
||||
for scheduler in cls.scheduler_registry.values():
|
||||
scheduler.stop()
|
||||
|
||||
@classmethod
|
||||
def lockdown(cls):
|
||||
cls.lockdown = True
|
||||
|
||||
def __init__(self, name, label=None):
|
||||
self.scheduled_jobs = {}
|
||||
self._scheduler = None
|
||||
self.name = name
|
||||
self.label = label
|
||||
self.__class__.scheduler_registry[self.name] = self
|
||||
|
||||
def start(self):
|
||||
logger.debug('starting scheduler: %s' % self.name)
|
||||
if not self.__class__.lockdown:
|
||||
self._scheduler = OriginalScheduler()
|
||||
for job in self.scheduled_jobs.values():
|
||||
self._schedule_job(job)
|
||||
|
||||
self._scheduler.start()
|
||||
else:
|
||||
logger.debug('lockdown in effect')
|
||||
|
||||
def stop(self):
|
||||
if self._scheduler:
|
||||
self._scheduler.shutdown()
|
||||
del self._scheduler
|
||||
self._scheduler = None
|
||||
|
||||
@property
|
||||
def running(self):
|
||||
if self._scheduler:
|
||||
return self._scheduler.running
|
||||
else:
|
||||
return False
|
||||
|
||||
def clear(self):
|
||||
for job in self.scheduled_jobs.values():
|
||||
self.stop_job(job)
|
||||
|
||||
def stop_job(self, job):
|
||||
self._scheduler.unschedule_job(job._job)
|
||||
del(self.scheduled_jobs[job.name])
|
||||
job.scheduler = None
|
||||
|
||||
def _schedule_job(self, job):
|
||||
if isinstance(job, IntervalJob):
|
||||
job._job = self._scheduler.add_interval_job(job.function, *job.args, **job.kwargs)
|
||||
elif isinstance(job, DateJob):
|
||||
job._job = self._scheduler.add_date_job(job.function, *job.args, **job.kwargs)
|
||||
else:
|
||||
raise UnknownJobClass
|
||||
|
||||
def add_job(self, job):
|
||||
if job.scheduler or job.name in self.scheduled_jobs.keys():
|
||||
raise AlreadyScheduled
|
||||
|
||||
if self._scheduler:
|
||||
self._scheduler_job(job)
|
||||
|
||||
job.scheduler = self
|
||||
self.scheduled_jobs[job.name] = job
|
||||
|
||||
def add_interval_job(self, name, label, function, *args, **kwargs):
|
||||
job = IntervalJob(name=name, label=label, function=function, *args, **kwargs)
|
||||
self.add_job(job)
|
||||
return job
|
||||
|
||||
def add_date_job(self, name, label, function, *args, **kwargs):
|
||||
job = DateJob(name=name, label=label, function=function, *args, **kwargs)
|
||||
self.add_job(job)
|
||||
return job
|
||||
|
||||
def get_job_list(self):
|
||||
return self.scheduled_jobs.values()
|
||||
|
||||
def __unicode__(self):
|
||||
return unicode(self.label or self.name)
|
||||
|
||||
@@ -1,2 +1,14 @@
|
||||
class AlreadyScheduled(Exception):
|
||||
"""
|
||||
Raised when trying to schedule a Job instance of anything after it was
|
||||
already scheduled in any other scheduler
|
||||
"""
|
||||
pass
|
||||
|
||||
|
||||
class UnknownJobClass(Exception):
|
||||
"""
|
||||
Raised when trying to schedule a Job that is not of a a type:
|
||||
IntervalJob or DateJob
|
||||
"""
|
||||
pass
|
||||
|
||||
@@ -4,6 +4,8 @@ from django.utils.translation import ugettext_lazy as _
|
||||
|
||||
from navigation.api import Link
|
||||
|
||||
from .permissions import PERMISSION_VIEW_JOB_LIST
|
||||
from .permissions import PERMISSION_VIEW_JOB_LIST, PERMISSION_VIEW_SCHEDULER_LIST
|
||||
|
||||
job_list = Link(text=_(u'interval job list'), view='job_list', icon='time.png', permissions=[PERMISSION_VIEW_JOB_LIST])
|
||||
scheduler_tool_link = Link(text=_(u'local schedulers'), view='scheduler_list', icon='time.png', permissions=[PERMISSION_VIEW_SCHEDULER_LIST])
|
||||
scheduler_list = Link(text=_(u'scheduler list'), view='scheduler_list', sprite='time', permissions=[PERMISSION_VIEW_SCHEDULER_LIST])
|
||||
job_list = Link(text=_(u'interval job list'), view='job_list', args='object.name', sprite='timeline_marker', permissions=[PERMISSION_VIEW_JOB_LIST])
|
||||
|
||||
1
apps/scheduler/literals.py
Normal file
1
apps/scheduler/literals.py
Normal file
@@ -0,0 +1 @@
|
||||
SHUTDOWN_COMMANDS = ['syncdb', 'migrate', 'schemamigration', 'datamigration', 'collectstatic', 'shell', 'shell_plus']
|
||||
@@ -5,4 +5,5 @@ from django.utils.translation import ugettext_lazy as _
|
||||
from permissions.models import PermissionNamespace, Permission
|
||||
|
||||
namespace = PermissionNamespace('scheduler', _(u'Scheduler'))
|
||||
PERMISSION_VIEW_JOB_LIST = Permission.objects.register(namespace, 'jobs_list', _(u'View the interval job list'))
|
||||
PERMISSION_VIEW_SCHEDULER_LIST = Permission.objects.register(namespace, 'schedulers_list', _(u'View the local scheduler list'))
|
||||
PERMISSION_VIEW_JOB_LIST = Permission.objects.register(namespace, 'jobs_list', _(u'View the local scheduler job list'))
|
||||
|
||||
@@ -1,4 +0,0 @@
|
||||
from apscheduler.scheduler import Scheduler
|
||||
|
||||
scheduler = Scheduler()
|
||||
scheduler.start()
|
||||
@@ -1,5 +1,6 @@
|
||||
from django.conf.urls.defaults import patterns, url
|
||||
|
||||
urlpatterns = patterns('scheduler.views',
|
||||
url(r'^list/$', 'job_list', (), 'job_list'),
|
||||
url(r'^scheduler/list/$', 'scheduler_list', (), 'scheduler_list'),
|
||||
url(r'^scheduler/(?P<scheduler_name>\w+)/job/list/$', 'job_list', (), 'job_list'),
|
||||
)
|
||||
|
||||
@@ -3,32 +3,67 @@ from __future__ import absolute_import
|
||||
from django.shortcuts import render_to_response
|
||||
from django.template import RequestContext
|
||||
from django.utils.translation import ugettext_lazy as _
|
||||
from django.http import Http404
|
||||
|
||||
from permissions.models import Permission
|
||||
from common.utils import encapsulate
|
||||
|
||||
from .permissions import PERMISSION_VIEW_JOB_LIST
|
||||
from .api import get_job_list
|
||||
from .permissions import PERMISSION_VIEW_SCHEDULER_LIST, PERMISSION_VIEW_JOB_LIST
|
||||
from .api import LocalScheduler
|
||||
|
||||
|
||||
def job_list(request):
|
||||
Permission.objects.check_permissions(request.user, [PERMISSION_VIEW_JOB_LIST])
|
||||
def scheduler_list(request):
|
||||
Permission.objects.check_permissions(request.user, [PERMISSION_VIEW_SCHEDULER_LIST])
|
||||
|
||||
context = {
|
||||
'object_list': get_job_list(),
|
||||
'title': _(u'interval jobs'),
|
||||
'object_list': LocalScheduler.get_all(),
|
||||
'title': _(u'local schedulers'),
|
||||
'extra_columns': [
|
||||
{
|
||||
'name': _(u'name'),
|
||||
'attribute': 'name'
|
||||
},
|
||||
{
|
||||
'name': _(u'label'),
|
||||
'attribute': encapsulate(lambda job: job['title'])
|
||||
'attribute': 'label'
|
||||
},
|
||||
{
|
||||
'name': _(u'running'),
|
||||
'attribute': 'running'
|
||||
},
|
||||
],
|
||||
'hide_object': True,
|
||||
}
|
||||
|
||||
return render_to_response('generic_list.html', context,
|
||||
context_instance=RequestContext(request))
|
||||
|
||||
|
||||
def job_list(request, scheduler_name):
|
||||
Permission.objects.check_permissions(request.user, [PERMISSION_VIEW_JOB_LIST])
|
||||
try:
|
||||
scheduler = LocalScheduler.get(scheduler_name)
|
||||
except:
|
||||
raise Http404
|
||||
|
||||
context = {
|
||||
'object_list': scheduler.get_job_list(),
|
||||
'title': _(u'local jobs in scheduler: %s') % scheduler,
|
||||
'extra_columns': [
|
||||
{
|
||||
'name': _(u'name'),
|
||||
'attribute': 'name'
|
||||
},
|
||||
{
|
||||
'name': _(u'label'),
|
||||
'attribute': 'label'
|
||||
},
|
||||
{
|
||||
'name': _(u'start date time'),
|
||||
'attribute': encapsulate(lambda job: job['job'].trigger.start_date)
|
||||
'attribute': 'start_date'
|
||||
},
|
||||
{
|
||||
'name': _(u'interval'),
|
||||
'attribute': encapsulate(lambda job: job['job'].trigger.interval)
|
||||
'name': _(u'type'),
|
||||
'attribute': 'job_type'
|
||||
},
|
||||
],
|
||||
'hide_object': True,
|
||||
|
||||
Reference in New Issue
Block a user