Skip to content
Closed
Show file tree
Hide file tree
Changes from 4 commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
67 changes: 38 additions & 29 deletions ndscheduler/corescheduler/core/base.py
Original file line number Diff line number Diff line change
Expand Up @@ -115,20 +115,17 @@ def run_scheduler_job(cls, job_class, job_id, execution_id, datastore, *args, **
description=job_class.get_failed_description(),
result=job_class.get_failed_result())

def add_scheduler_job(self, job_class_string, name, pub_args=None,
month=None, day_of_week=None, day=None,
hour=None, minute=None, **kwargs):
def add_scheduler_job(self, job_class_string, name, trigger, pub_args=None,
trigger_params=None, **kwargs):
"""Add a job. Job information will be persistent in postgres.
This is a NON-BLOCKING operation, as internally, apscheduler calls wakeup()
that is async.
:param str job_class_string: String for job class, e.g., myscheduler.jobs.a_job.NiceJob
:param str name: String for job name, e.g., Check Melissa job.
:param str trigger: String for job trigger passed to the scheduler.
:param list pub_args: List for arguments passed to publish method of a task.
:param str month: String for month cron string, e.g., */10
:param str day_of_week: String for day of week cron string, e.g., 1-6
:param str day: String for day cron string, e.g., */1
:param str hour: String for hour cron string, e.g., */2
:param str minute: String for minute cron string, e.g., */3
:param trigger_params: Dict of trigger parameters passed to the apscheduler
during adding jobs.
:param dict kwargs: Other keyword arguments passed to run_job function.
:return: String of job id, e.g., 6bca19736d374ef2b3df23eb278b512e
:rtype: str
Expand All @@ -144,9 +141,17 @@ def add_scheduler_job(self, job_class_string, name, pub_args=None,
datastore.db_config, datastore.table_names]
arguments.extend(pub_args)

self.add_job(self.run_job,
'cron', month=month, day=day, day_of_week=day_of_week, hour=hour,
minute=minute, args=arguments, kwargs=kwargs, name=name, id=job_id)
# Rename interval to seconds for calling apscheduler's reschedule API
if 'interval' in trigger_params:
trigger_params['seconds'] = int(trigger_params.pop('interval'))

self.add_job(func=self.run_job,
trigger=trigger,
args=arguments,
kwargs=kwargs,
name=name,
id=job_id,
**trigger_params)
return job_id

def modify_scheduler_job(self, job_id, **kwargs):
Expand All @@ -159,11 +164,10 @@ def modify_scheduler_job(self, job_id, **kwargs):
- job_class_string: String for job class string, e.g.,
myscheduler.jobs.a_job.NiceJob
- pub_args: List of arguments passed to the task.
- month: String for month cron string, e.g., */10
- day_of_week: String for day of week cron string, e.g., 1-6
- day: String for day cron string, e.g., */1
- hour: String for hour cron string, e.g., */2
- minute: String for minute cron string, e.g., */3
- trigger: String for job trigger passed to the scheduler.
- trigger_params: Dict of trigger parameters passed to the apscheduler during
adding jobs.

"""

# This is a BLOCKING operation
Expand All @@ -173,24 +177,29 @@ def modify_scheduler_job(self, job_id, **kwargs):
if 'job_class_string' in kwargs or 'pub_args' in kwargs:
args = list(job.args)
if 'job_class_string' in kwargs:
args[0] = kwargs['job_class_string']
args[0] = kwargs.pop('job_class_string')
# 'task_name' is not an argument for modify_job.
del kwargs['job_class_string']
if 'pub_args' in kwargs:
args = args[:constants.JOB_ARGS] + kwargs['pub_args']
del kwargs['pub_args']
args = args[:constants.JOB_ARGS] + kwargs.pop('pub_args')
kwargs['args'] = args

cron_keywords = ['month', 'day', 'hour', 'minute', 'day_of_week']
trigger_kwargs = {}
for cron_key in cron_keywords:
if cron_key in kwargs:
trigger_kwargs[cron_key] = kwargs[cron_key]
del kwargs[cron_key]
# Handle trigger rescheduling
# This is a NON-BLOCKING operation
if 'trigger' in kwargs and 'trigger_params' in kwargs:
# Handle trigger only if all parameter are available
trigger_params = kwargs.pop('trigger_params')

# Rename interval to seconds for calling apscheduler's reschedule API
if 'interval' in trigger_params:
trigger_params['seconds'] = int(trigger_params.pop('interval'))

job.reschedule(trigger=kwargs.pop('trigger'), **trigger_params)

if trigger_kwargs:
# This is a NON-BLOCKING operation
job.reschedule(trigger='cron', **trigger_kwargs)
else: # Remove if not all trigger parameters are available
if 'trigger' in kwargs:
kwargs.pop('trigger')
if 'trigger_params' in kwargs:
kwargs.pop('trigger_params')

# This is a NON-BLOCKING operation
job.modify(**kwargs)
Expand Down
18 changes: 17 additions & 1 deletion ndscheduler/corescheduler/datastore/base.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,9 @@

import dateutil.tz
import dateutil.parser
import apscheduler.triggers.cron
import apscheduler.triggers.interval

from apscheduler.jobstores import sqlalchemy as sched_sqlalchemy
from sqlalchemy import desc, select, MetaData

Expand Down Expand Up @@ -125,7 +128,20 @@ def _build_execution(self, row):
'name': job.name,
'task_name': utils.get_job_name(job),
'pub_args': utils.get_job_args(job)}
return_json['job'].update(utils.get_cron_strings(job))

if isinstance(job.trigger, apscheduler.triggers.cron.CronTrigger):
return_json["trigger"] = "cron"
return_json["trigger_params"] = utils.get_cron_strings(job)
elif isinstance(job.trigger, apscheduler.triggers.interval.IntervalTrigger):
return_json["trigger"] = "interval"
trigger_params = {
'interval': job.trigger.interval.total_seconds()
}
return_json["trigger_params"] = trigger_params
else:
return_json["trigger"] = "unknown"
return_json["trigger_params"] = {}

return return_json

def get_time_isoformat_from_db(self, time_object):
Expand Down
17 changes: 8 additions & 9 deletions ndscheduler/corescheduler/scheduler_manager.py
Original file line number Diff line number Diff line change
Expand Up @@ -79,25 +79,24 @@ def stop(self):
#
# Manage jobs
#
def add_job(self, job_class_string, name, pub_args=None, month=None,
day_of_week=None, day=None, hour=None, minute=None, **kwargs):
def add_job(self, job_class_string, name, trigger, pub_args=None,
trigger_params=None, **kwargs):
"""Add a job. Job infomation will be persistent in the datastore.
This is a NON-BLOCKING operation, as internally, apscheduler calls wakeup()
that is async.
:param str job_class_string: String for job class, e.g., myscheduler.jobs.a_job.NiceJob
:param str name: String for job name, e.g., Check Melissa job.
:param str trigger: String for job trigger passed to the scheduler.
:param pub_args: List for arguments passed to publish method of a task.
:param str month: String for month cron string, e.g., */10
:param str day_of_week: String for day of week cron string, e.g., 1-6
:param str day: String for day cron string, e.g., */1
:param str hour: String for hour cron string, e.g., */2
:param str minute: String for minute cron string, e.g., */3
:param trigger_params: Dict of trigger parameters passed to the apscheduler
Comment thread
matrixx567 marked this conversation as resolved.
during adding jobs.
:param kwargs: Other keyword arguments passed to run_job function.
:return: String of job id, e.g., 6bca19736d374ef2b3df23eb278b512e
:rtype: str
"""
return self.sched.add_scheduler_job(job_class_string, name, pub_args, month, day_of_week,
day, hour, minute, **kwargs)
return self.sched.add_scheduler_job(job_class_string, name,
trigger, pub_args, trigger_params,
**kwargs)

def pause_job(self, job_id):
"""Pauses the schedule of a job.
Expand Down
53 changes: 31 additions & 22 deletions ndscheduler/corescheduler/scheduler_manager_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -27,15 +27,18 @@ def test_add_job_get_job(self):
task_name = 'hello.world'
name = 'it is hello world'
pub_args = ('1', '2', '3')
month = '*/1'
day_of_week = 'sat'
day = '*/2'
hour = '*/3'
minute = '*/4'
trigger_params = dict(
month='*/1',
day_of_week='sat',
day='*/2',
hour='*/3',
minute='*/4'
)
trigger = 'cron'

# non-blocking operation
job_id = self.scheduler.add_job(task_name, name, pub_args, month, day_of_week, day,
hour, minute, languages='en-us')
job_id = self.scheduler.add_job(task_name, name, trigger, pub_args,
trigger_params, languages='en-us')

self.assertTrue(len(job_id), 32)

Expand All @@ -50,17 +53,17 @@ def test_add_job_get_job(self):
# Year
self.assertEqual(str(job.trigger.fields[0]), '*')
# Month
self.assertEqual(str(job.trigger.fields[1]), month)
self.assertEqual(str(job.trigger.fields[1]), trigger_params['month'])
# day of month
self.assertEqual(str(job.trigger.fields[2]), day)
self.assertEqual(str(job.trigger.fields[2]), trigger_params['day'])
# week
self.assertEqual(str(job.trigger.fields[3]), '*')
# day of week
self.assertEqual(str(job.trigger.fields[4]), day_of_week)
self.assertEqual(str(job.trigger.fields[4]), trigger_params['day_of_week'])
# hour
self.assertEqual(str(job.trigger.fields[5]), hour)
self.assertEqual(str(job.trigger.fields[5]), trigger_params['hour'])
# minute
self.assertEqual(str(job.trigger.fields[6]), minute)
self.assertEqual(str(job.trigger.fields[6]), trigger_params['minute'])
# second
self.assertEqual(str(job.trigger.fields[7]), '0')

Expand All @@ -69,26 +72,32 @@ def test_add_job_modify_job(self):
job_class_string = 'hello.world2'
name = 'it is hello world 2'
pub_args = ('1', '2', '3')
month = '*/1'
day_of_week = 'sat'
day = '*/2'
hour = '*/3'
minute = '*/4'
trigger_params = dict(
month='*/1',
day_of_week='sat',
day='*/2',
hour='*/3',
minute='*/4'
)
trigger = 'cron'

# non-blocking operation
job_id = self.scheduler.add_job(job_class_string, name, pub_args, month, day_of_week,
day, hour, minute)
job_id = self.scheduler.add_job(job_class_string, name, trigger, pub_args, trigger_params)

self.assertTrue(len(job_id), 32)

job_class_string = 'hello.world1234'
args = ['5', '6', '7']
name = 'hello world 3'
month = '*/6'
trigger = 'cron'
trigger_params = dict(
month='*/6'
)

# non-blocking operation
self.scheduler.modify_job(job_id, name=name, job_class_string=job_class_string,
pub_args=args, month=month)
trigger=trigger,
pub_args=args, trigger_params=trigger_params)

# blocking operation
job = self.scheduler.get_job(job_id)
Expand All @@ -99,4 +108,4 @@ def test_add_job_modify_job(self):
None, None]
arguments += args
self.assertEqual(list(job.args), arguments)
self.assertEqual(str(job.trigger.fields[1]), month)
self.assertEqual(str(job.trigger.fields[1]), trigger_params['month'])
Loading