From 859109c37879814ad6ef44233097d8d027322de6 Mon Sep 17 00:00:00 2001 From: Roberto Rosario Date: Wed, 1 Aug 2012 01:41:37 -0400 Subject: [PATCH] Refactor scheduler --- apps/scheduler/__init__.py | 51 ++++-------- apps/scheduler/api.py | 153 ++++++++++++++++++++++++++++++---- apps/scheduler/exceptions.py | 12 +++ apps/scheduler/links.py | 6 +- apps/scheduler/literals.py | 1 + apps/scheduler/permissions.py | 3 +- apps/scheduler/runtime.py | 4 - apps/scheduler/urls.py | 3 +- apps/scheduler/views.py | 57 ++++++++++--- 9 files changed, 220 insertions(+), 70 deletions(-) create mode 100644 apps/scheduler/literals.py delete mode 100644 apps/scheduler/runtime.py diff --git a/apps/scheduler/__init__.py b/apps/scheduler/__init__.py index e9e693b024..25fcbf6190 100644 --- a/apps/scheduler/__init__.py +++ b/apps/scheduler/__init__.py @@ -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) diff --git a/apps/scheduler/api.py b/apps/scheduler/api.py index 6ce39bc3c0..aeebd6fbec 100644 --- a/apps/scheduler/api.py +++ b/apps/scheduler/api.py @@ -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) diff --git a/apps/scheduler/exceptions.py b/apps/scheduler/exceptions.py index f30d9fe815..6d6515a76c 100644 --- a/apps/scheduler/exceptions.py +++ b/apps/scheduler/exceptions.py @@ -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 diff --git a/apps/scheduler/links.py b/apps/scheduler/links.py index 7808f3331c..9ad1a1ff10 100644 --- a/apps/scheduler/links.py +++ b/apps/scheduler/links.py @@ -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]) diff --git a/apps/scheduler/literals.py b/apps/scheduler/literals.py new file mode 100644 index 0000000000..b56b5148d7 --- /dev/null +++ b/apps/scheduler/literals.py @@ -0,0 +1 @@ +SHUTDOWN_COMMANDS = ['syncdb', 'migrate', 'schemamigration', 'datamigration', 'collectstatic', 'shell', 'shell_plus'] diff --git a/apps/scheduler/permissions.py b/apps/scheduler/permissions.py index 203f675ff4..2ef2343811 100644 --- a/apps/scheduler/permissions.py +++ b/apps/scheduler/permissions.py @@ -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')) diff --git a/apps/scheduler/runtime.py b/apps/scheduler/runtime.py deleted file mode 100644 index a9440e946b..0000000000 --- a/apps/scheduler/runtime.py +++ /dev/null @@ -1,4 +0,0 @@ -from apscheduler.scheduler import Scheduler - -scheduler = Scheduler() -scheduler.start() diff --git a/apps/scheduler/urls.py b/apps/scheduler/urls.py index fde9602994..3630a34bbc 100644 --- a/apps/scheduler/urls.py +++ b/apps/scheduler/urls.py @@ -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\w+)/job/list/$', 'job_list', (), 'job_list'), ) diff --git a/apps/scheduler/views.py b/apps/scheduler/views.py index 597632eca7..fb622ce818 100644 --- a/apps/scheduler/views.py +++ b/apps/scheduler/views.py @@ -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,