diff --git a/ndscheduler/corescheduler/core/base.py b/ndscheduler/corescheduler/core/base.py
index 84efe290..ca1b7fe5 100644
--- a/ndscheduler/corescheduler/core/base.py
+++ b/ndscheduler/corescheduler/core/base.py
@@ -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
@@ -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):
@@ -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
@@ -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)
diff --git a/ndscheduler/corescheduler/datastore/base.py b/ndscheduler/corescheduler/datastore/base.py
index be10772c..3de9f052 100644
--- a/ndscheduler/corescheduler/datastore/base.py
+++ b/ndscheduler/corescheduler/datastore/base.py
@@ -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
@@ -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):
diff --git a/ndscheduler/corescheduler/scheduler_manager.py b/ndscheduler/corescheduler/scheduler_manager.py
index eb44247c..fb1f4696 100644
--- a/ndscheduler/corescheduler/scheduler_manager.py
+++ b/ndscheduler/corescheduler/scheduler_manager.py
@@ -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
+ 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.
@@ -165,5 +164,8 @@ def modify_job(self, job_id, **kwargs):
- 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.
"""
self.sched.modify_scheduler_job(job_id, **kwargs)
diff --git a/ndscheduler/corescheduler/scheduler_manager_test.py b/ndscheduler/corescheduler/scheduler_manager_test.py
index c8d1bbeb..05748e77 100644
--- a/ndscheduler/corescheduler/scheduler_manager_test.py
+++ b/ndscheduler/corescheduler/scheduler_manager_test.py
@@ -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)
@@ -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')
@@ -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)
@@ -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'])
diff --git a/ndscheduler/server/handlers/README.md b/ndscheduler/server/handlers/README.md
index 7b1649c8..9bedd60d 100644
--- a/ndscheduler/server/handlers/README.md
+++ b/ndscheduler/server/handlers/README.md
@@ -28,7 +28,8 @@
$ curl -X POST localhost:8888/api/v1/jobs \
--header "Content-Type:application/json" \
-d "{\"job_class_string\": \"simple_scheduler.jobs.sample_job.AwesomeJob\", \
- \"name\": \"My first job\", \"minute\": \"*/1\", \"pub_args\": [\"arg1\", 2]}"
+ \"name\": \"My first job\", \"pub_args\": [\"arg1\", 2], \
+ \"trigger\": \"cron\", \"trigger_params\": {\"minute\": \"*/1\"}}"
# Get all jobs
$ curl localhost:8888/api/v1/jobs
@@ -37,6 +38,50 @@
## REST APIs
### Jobs
+#### Difference between a Cron job and an Interval job
+
+Version 2 of the API supports next to the already existing Cron jobs also
+Interval jobs. Therefore the API of the jobs has to be changed. A job uses the
+string parameter ```trigger``` that indicates the job type. Possible values are
+ ```cron``` or ```interval```. The parameters for the chosen job type can be set
+by the ```trigger_params```.
+
+* **Cron Job**
+
+ A cron job is defined by the ```trigger``` = ```cron```. A cron
+ job needs the parameters ```month```, ```day_of_week```, ```day```,
+ ```hour```, ```minute``` parameters within the ```trigger_params```.
+
+ {
+ "job_id": "d8f376e858a411e4b6ae22000ac58d05",
+ "job_class_string": "simple_scheduler.jobs.clean_apns.CleanAPNsJob",
+ "name": "Clean APNs",
+ "pub_args": ["arg1": 1, "arg2": 2],
+ "trigger": "cron",
+ "trigger_params": {
+ "month": "*",
+ "day_of_week": "*",
+ "day": "*",
+ "hour": "*/1",
+ "minute": "*"}
+
+
+* **Interval Job**
+
+ An interval job is defined by the ```trigger``` = ```interval```. An interval
+ job needs the ```interval``` parameter within the ```trigger_params```. The
+ interval is defined in seconds.
+
+ {
+ "job_id": "d8f376e858a411e4b6ae22000ac58d05",
+ "job_class_string": "simple_scheduler.jobs.clean_apns.CleanAPNsJob",
+ "name": "Clean APNs",
+ "pub_args": ["arg1": 1, "arg2": 2],
+ "trigger": "interval",
+ "trigger_params": {
+ "interval": 3600}
+
+
#### Get all jobs
Returns json data for all jobs.
@@ -68,15 +113,17 @@
"job_class_string": "simple_scheduler.jobs.clean_apns.CleanAPNsJob",
"name": "Clean APNs",
"pub_args": ["arg1": 1, "arg2": 2],
- "month": "*",
- "day_of_week": "*",
- "day": "*",
- "hour": "*/1",
- "minute": "*"
+ "trigger": "cron",
+ "trigger_params": {
+ "month": "*",
+ "day_of_week": "*",
+ "day": "*",
+ "hour": "*/1",
+ "minute": "*"}
},
...]
}
-
+
* **Error Response:**
* **Code:** 400 Bad Request
@@ -126,7 +173,26 @@
* **Success Response:**
- * **Code:** 200 OK
+ * **Return for a Cron job**
+ **Code:** 200 OK
+ **Content:**
+
+ {
+ "job_id": "d8f376e858a411e4b6ae22000ac58d05",
+ "job_class_string": "simple_scheduler.jobs.clean_apns.CleanAPNsJob",
+ "name": "Clean APNs",
+ "pub_args": ["arg1": 1, "arg2": 2],
+ "trigger": "cron",
+ "trigger_params": {
+ "month": "*",
+ "day_of_week": "*",
+ "day": "*",
+ "hour": "*/1",
+ "minute": "*"}
+ }
+
+ * **Return for an Interval job**
+ **Code:** 200 OK
**Content:**
{
@@ -134,12 +200,11 @@
"job_class_string": "simple_scheduler.jobs.clean_apns.CleanAPNsJob",
"name": "Clean APNs",
"pub_args": ["arg1": 1, "arg2": 2],
- "month": "*",
- "day_of_week": "*",
- "day": "*",
- "hour": "*/1",
- "minute": "*"
+ "trigger": "interval",
+ "trigger_params": {
+ "interval": 3600}
}
+
* **Error Response:**
@@ -184,19 +249,33 @@
None
-* **Data Params**
+* **Data Params for a Cron job**
{
"job_class_string": "simple_scheduler.jobs.clean_apns.CleanAPNsJob",
"name": "Clean APNs",
"pub_args": ["arg1": 1, "arg2": 2],
- "month": "*",
- "day_of_week": "*",
- "day": "*",
- "hour": "*/1",
- "minute": "*"
+ "trigger": "cron",
+ "trigger_params": {
+ "month": "*",
+ "day_of_week": "*",
+ "day": "*",
+ "hour": "*/1",
+ "minute": "*"}
}
+* **Data Params for an Interval job**
+
+ {
+ "job_class_string": "simple_scheduler.jobs.clean_apns.CleanAPNsJob",
+ "name": "Clean APNs",
+ "pub_args": ["arg1": 1, "arg2": 2],
+ "trigger": "interval",
+ "trigger_params": {
+ "interval": 3600}
+ }
+
+
Required fields: `job_class_string` and `name`
* **Success Response:**
@@ -226,11 +305,9 @@ Required fields: `job_class_string` and `name`
"job_class_string": "simple_scheduler.jobs.clean_apns.CleanAPNsJob",
"name": "Clean APNs",
"pub_args": ["arg1": 1, "arg2": 2],
- "month": "*",
- "day_of_week": "*",
- "day": "*",
- "hour": "*/1",
- "minute": "*"
+ "trigger": "interval",
+ "trigger_params": {
+ "interval": 3600}
},
success : function(r) {
console.log(r);
@@ -310,17 +387,31 @@ Required fields: `job_class_string` and `name`
None
-* **Data Params**
+* **Data Params for a Cron job**
{
"job_class_string": "simple_scheduler.jobs.clean_apns.CleanAPNsJob",
"name": "Clean APNs",
"pub_args": ["arg1": 1, "arg2": 2],
- "month": "*",
- "day_of_week": "*",
- "day": "*",
- "hour": "*/1",
- "minute": "*"
+ "trigger": "cron",
+ "trigger_params": {
+ "month": "*",
+ "day_of_week": "*",
+ "day": "*",
+ "hour": "*/1",
+ "minute": "*"}
+ }
+
+
+* **Data Params for an Interval job**
+
+ {
+ "job_class_string": "simple_scheduler.jobs.clean_apns.CleanAPNsJob",
+ "name": "Clean APNs",
+ "pub_args": ["arg1": 1, "arg2": 2],
+ "trigger": "interval",
+ "trigger_params": {
+ "interval": 3600}
}
* **Success Response:**
@@ -353,6 +444,14 @@ Required fields: `job_class_string` and `name`
url: "/api/v1/jobs/d8f376e858a411e4b6ae22000ac58d05",
dataType: "json",
type : "PUT",
+ data: {
+ "job_class_string": "simple_scheduler.jobs.clean_apns.CleanAPNsJob",
+ "name": "Clean APNs",
+ "pub_args": ["arg1": 1, "arg2": 2],
+ "trigger": "interval",
+ "trigger_params": {
+ "interval": 3600}
+ },
success : function(r) {
console.log(r);
}
@@ -505,12 +604,14 @@ Required fields: `job_class_string` and `name`
execution_id: "7252d7a6a80f11e58bcc02ba903740c3",
hostname: "",
job: {
- day: "*",
- day_of_week: "*",
- hour: "*",
job_id: "bb0dec52797f11e4a14122000a150f89",
- minute: "*/5",
- month: "*",
+ trigger: "cron",
+ trigger_params: {
+ day: "*",
+ day_of_week: "*",
+ hour: "*",
+ minute: "*/5",
+ month: "*"},
name: "Poll sendgrid for bounces and spamreports",
pub_args: [],
job_class_string: "simple_scheduler.jobs.sample_job.ImportDataJob",
@@ -582,12 +683,14 @@ Required fields: `job_class_string` and `name`
execution_id: "7252d7a6a80f11e58bcc02ba903740c3",
hostname: "",
job: {
- day: "*",
- day_of_week: "*",
- hour: "*",
job_id: "bb0dec52797f11e4a14122000a150f89",
- minute: "*/5",
- month: "*",
+ trigger: "cron",
+ trigger_params: {
+ day: "*",
+ day_of_week: "*",
+ hour: "*",
+ minute: "*/5",
+ month: "*"},
name: "Poll sendgrid for bounces and spamreports",
pub_args: [],
job_class_string: "simple_scheduler.jobs.sample_job.ImportDataJob",
diff --git a/ndscheduler/server/handlers/audit_logs.py b/ndscheduler/server/handlers/audit_logs.py
index 9a2a2f30..d74cc2f5 100644
--- a/ndscheduler/server/handlers/audit_logs.py
+++ b/ndscheduler/server/handlers/audit_logs.py
@@ -47,6 +47,6 @@ def get_logs_yield(self):
def get(self):
"""Returns audit logs.
- Handles the endpoint GET /api/v1/logs.
+ Handles the endpoint GET /api/v2/logs.
"""
self.get_logs_yield()
diff --git a/ndscheduler/server/handlers/executions.py b/ndscheduler/server/handlers/executions.py
index 82a7b5a9..6b04b259 100644
--- a/ndscheduler/server/handlers/executions.py
+++ b/ndscheduler/server/handlers/executions.py
@@ -90,14 +90,14 @@ def get(self, execution_id=None):
"""Returns a execution or multiple executions.
Handles two endpoints:
- GET /api/v1/executions (when execution_id == None)
+ GET /api/v2/executions (when execution_id == None)
It takes two query string parameters:
- time_range_end - unix epoch timestamp. Default: now
- time_range_start - unix epoch timestamp. Default: 10 minutes ago.
These two parameters limit the executions to return:
time_range_start <= execution.scheduled_time <= time_range_end
- GET /api/v1/executions/{execution_id} (when execution_id != None)
+ GET /api/v2/executions/{execution_id} (when execution_id != None)
:param str execution_id: Execution id.
"""
@@ -163,7 +163,7 @@ def post(self, job_id):
"""Runs a job.
Handles an endpoint:
- POST /api/v1/executions
+ POST /api/v2/executions
Args:
job_id: String for job id.
@@ -175,6 +175,6 @@ def delete(self, job_id):
"""Stops a job execution.
Handles an endpoint:
- POST /api/v1/executions
+ POST /api/v2/executions
"""
raise tornado.web.HTTPError(501, 'Not implemented yet.')
diff --git a/ndscheduler/server/handlers/executions_test.py b/ndscheduler/server/handlers/executions_test.py
index 8d0871cb..a97e99b4 100644
--- a/ndscheduler/server/handlers/executions_test.py
+++ b/ndscheduler/server/handlers/executions_test.py
@@ -26,7 +26,7 @@ class ExecutionsTest(tornado.testing.AsyncHTTPTestCase):
def setUp(self, *args, **kwargs):
super(ExecutionsTest, self).setUp(*args, **kwargs)
self.server.start_scheduler()
- self.EXECUTIONS_URL = '/api/v1/executions'
+ self.EXECUTIONS_URL = '/api/v2/executions'
self.old_get_executions_yield = executions.Handler.get_executions_yield
self.old_get_execution_yield = executions.Handler.get_execution_yield
executions.Handler.get_executions_yield = mock_get_executions_yield
diff --git a/ndscheduler/server/handlers/jobs.py b/ndscheduler/server/handlers/jobs.py
index 2d0f3e04..29cf2dd9 100644
--- a/ndscheduler/server/handlers/jobs.py
+++ b/ndscheduler/server/handlers/jobs.py
@@ -6,6 +6,9 @@
import tornado.gen
import tornado.web
+import apscheduler.triggers.cron
+import apscheduler.triggers.interval
+
from ndscheduler.corescheduler import constants
from ndscheduler.corescheduler import utils
from ndscheduler.server.handlers import base
@@ -42,7 +45,19 @@ def _build_job_dict(self, job):
'job_class_string': utils.get_job_name(job),
'pub_args': utils.get_job_args(job)}
- return_dict.update(utils.get_cron_strings(job))
+ if isinstance(job.trigger, apscheduler.triggers.cron.CronTrigger):
+ return_dict["trigger"] = "cron"
+ return_dict["trigger_params"] = utils.get_cron_strings(job)
+ elif isinstance(job.trigger, apscheduler.triggers.interval.IntervalTrigger):
+ return_dict["trigger"] = "interval"
+ trigger_params = {
+ 'interval': job.trigger.interval.total_seconds()
+ }
+ return_dict["trigger_params"] = trigger_params
+ else:
+ return_dict["trigger"] = "unknown"
+ return_dict["trigger_params"] = {}
+
return return_dict
@tornado.concurrent.run_on_executor
@@ -106,8 +121,8 @@ def get(self, job_id=None):
"""Returns a job or multiple jobs.
Handles two endpoints:
- GET /api/v1/jobs (when job_id == None)
- GET /api/v1/jobs/{job_id} (when job_id != None)
+ GET /api/v2/jobs (when job_id == None)
+ GET /api/v2/jobs/{job_id} (when job_id != None)
:param str job_id: String for job id.
"""
@@ -123,7 +138,7 @@ def post(self):
add_job() is a non-blocking operation, but audit log is a blocking operation.
Handles an endpoint:
- POST /api/v1/jobs
+ POST /api/v2/jobs
"""
self._validate_post_data()
@@ -171,7 +186,7 @@ def delete(self, job_id):
"""Deletes a job.
Handles an endpoint:
- DELETE /api/v1/jobs/{job_id}
+ DELETE /api/v2/jobs/{job_id}
:param str job_id: Job id
"""
@@ -182,21 +197,58 @@ def delete(self, job_id):
self.set_status(200)
self.finish(response)
- def _generate_description_for_item(self, old_job, new_job, item):
+ def _generate_description_for_item(self, item_name, old_item, new_item):
"""Returns a diff for one field of a job.
-
- :param dict old_job: Dict for old job.
- :param dict new_job: Dict for new job after modification.
+ :param str item_name: Description of the item
+ :param old_item: Item of old job.
+ :param new_item: Item of modified job.
:return: String for description.
:rtype: str
"""
- if old_job[item] != new_job[item]:
+ if old_item != new_item:
return ('%s: %s =>'
- ' %s
') % (item, old_job[item], new_job[item])
+ ' %s
') % (item_name, old_item, new_item)
return ''
+ def _generate_trigger_description(self, trigger, trigger_params):
+ """Returns a trigger description for the job
+
+ :param str trigger: Trigger of the job
+ :param dict trigger_params: Parameters of the trigger
+
+ :return: Description of the trigger
+ :rtype: str
+ """
+
+ if trigger == 'cron':
+ descr = "%s: minute=\"%s\" hour=\"%s\" day=\"%s\" month=\"%s\" " \
+ "day_of_week=\"%s\"" % (trigger, trigger_params['minute'],
+ trigger_params['hour'], trigger_params['day'],
+ trigger_params['month'], trigger_params['day_of_week']
+ )
+ elif trigger == 'interval':
+ descr = "%s: interval=\"%d\"" % (trigger, trigger_params['interval'])
+ else:
+ descr = "unknown"
+
+ return descr
+
+ def _generate_pubargs_description(self, pub_args):
+ """Generates description text for pub_args.
+
+ :param pub_args: List or tuple of pup_args.
+
+ :return: String description of pup_args.
+ :rtype: str
+ """
+
+ if isinstance(pub_args, tuple):
+ return str(list(pub_args))
+ else:
+ return str(pub_args)
+
def _generate_description_for_modify(self, old_job, new_job):
"""Generates description text after modifying a job.
@@ -206,19 +258,24 @@ def _generate_description_for_modify(self, old_job, new_job):
:return: String for description.
:rtype: str
"""
- description = ''
- items = [
- 'name',
- 'job_class_string',
- 'pub_args',
- 'minute',
- 'hour',
- 'day',
- 'month',
- 'day_of_week'
- ]
- for item in items:
- description += self._generate_description_for_item(old_job, new_job, item)
+
+ description = self._generate_description_for_item('Name', old_job['name'], new_job['name'])
+ description += self._generate_description_for_item(
+ 'Job Class',
+ old_job['job_class_string'],
+ new_job['job_class_string']
+ )
+ description += self._generate_description_for_item(
+ 'Trigger',
+ self._generate_trigger_description(old_job['trigger'], old_job['trigger_params']),
+ self._generate_trigger_description(new_job['trigger'], new_job['trigger_params'])
+ )
+ description += self._generate_description_for_item(
+ 'Arguments',
+ self._generate_pubargs_description(old_job['pub_args']),
+ self._generate_pubargs_description(new_job['pub_args'])
+ )
+
return description
def _modify_job(self, job_id):
@@ -262,7 +319,7 @@ def put(self, job_id):
"""Modifies a job.
Handles an endpoint:
- PUT /api/v1/jobs/{job_id}
+ PUT /api/v2/jobs/{job_id}
:param str job_id: Job id.
"""
@@ -281,7 +338,7 @@ def patch(self, job_id):
pause_job() is a non-blocking operation, but audit log is a blocking operation.
Handles an endpoint:
- PATCH /api/v1/jobs/{job_id}
+ PATCH /api/v2/jobs/{job_id}
:param str job_id: Job id.
"""
@@ -308,7 +365,7 @@ def options(self, job_id):
resume_job() is a non-blocking operation, but audit log is a blocking operation.
Handles an endpoint:
- OPTIONS /api/v1/jobs/{job_id}
+ OPTIONS /api/v2/jobs/{job_id}
:param str job_id: Job id.
"""
@@ -336,18 +393,29 @@ def _validate_post_data(self):
:raises: HTTPError(400: Bad arguments).
"""
- all_required_fields = ['name', 'job_class_string']
+ all_required_fields = ['name', 'job_class_string', 'trigger', 'trigger_params']
for field in all_required_fields:
if field not in self.json_args:
raise tornado.web.HTTPError(400, reason='Require this parameter: %s' % field)
- at_least_one_required_fields = ['month', 'day', 'hour', 'minute', 'day_of_week']
- valid_cron_string = False
- for field in at_least_one_required_fields:
- if field in self.json_args:
- valid_cron_string = True
- break
-
- if not valid_cron_string:
- raise tornado.web.HTTPError(400, reason=('Require at least one of following parameters:'
- ' %s' % str(at_least_one_required_fields)))
+ # TODO better validating
+ if self.json_args['trigger'] == 'cron':
+ valid_cron_string = False
+ at_least_one_required_fields = ['month', 'day', 'hour', 'minute', 'day_of_week']
+ for field in at_least_one_required_fields:
+ if field in self.json_args['trigger_params']:
+ valid_cron_string = True
+ break
+
+ if not valid_cron_string:
+ raise tornado.web.HTTPError(
+ 400, reason=('Require at least one of following parameters in "trigger_params":'
+ ' %s' % str(at_least_one_required_fields))
+ )
+
+ if self.json_args['trigger'] == 'interval':
+ if 'interval' not in self.json_args['trigger_params']:
+ raise tornado.web.HTTPError(
+ 400,
+ reason=('Require the parameter "interval" in "trigger_params".')
+ )
diff --git a/ndscheduler/server/handlers/jobs_test.py b/ndscheduler/server/handlers/jobs_test.py
index 76d6b69c..1e9a5ea4 100644
--- a/ndscheduler/server/handlers/jobs_test.py
+++ b/ndscheduler/server/handlers/jobs_test.py
@@ -34,7 +34,7 @@ def setUp(self, *args, **kwargs):
super(JobsTest, self).setUp(*args, **kwargs)
self.server.start_scheduler()
- self.JOBS_URL = '/api/v1/jobs'
+ self.JOBS_URL = '/api/v2/jobs'
self.old_get_jobs_yield = jobs.Handler.get_jobs_yield
self.old_get_job_yield = jobs.Handler.get_job_yield
self.old_delete_job_yield = jobs.Handler.delete_job_yield
@@ -72,7 +72,10 @@ def test_add_job_success(self):
data = {
'job_class_string': 'hello.world',
'name': 'hello world job',
- 'minute': '*/5'}
+ 'trigger': 'cron',
+ 'trigger_params': {
+ 'minute': '*/5'}
+ }
response = self.fetch(self.JOBS_URL, method='POST', headers=headers,
body=json.dumps(data))
return_info = json.loads(response.body.decode())
@@ -92,7 +95,10 @@ def test_add_job_failed(self):
data = {
'job_class_string': 'hello.world',
- 'minute': '*/5'}
+ 'trigger': 'cron',
+ 'trigger_params': {
+ 'minute': '*/5'}
+ }
response = self.fetch(self.JOBS_URL, method='POST', headers=headers,
body=json.dumps(data))
self.assertEqual(response.code, 400)
@@ -102,7 +108,10 @@ def test_pause_resume_job(self):
data = {
'job_class_string': 'hello.world',
'name': 'hello world job',
- 'minute': '*/5'}
+ 'trigger': 'cron',
+ 'trigger_params': {
+ 'minute': '*/5'}
+ }
response = self.fetch(self.JOBS_URL, method='POST', headers=headers,
body=json.dumps(data))
return_info = json.loads(response.body.decode())
@@ -121,7 +130,10 @@ def test_get_jobs(self):
data = {
'job_class_string': 'hello.world',
'name': 'hello world job',
- 'minute': '*/5'}
+ 'trigger': 'cron',
+ 'trigger_params': {
+ 'minute': '*/5'}
+ }
response = self.fetch(self.JOBS_URL, method='POST', headers=headers,
body=json.dumps(data))
return_info = json.loads(response.body.decode())
@@ -134,14 +146,17 @@ def test_get_jobs(self):
job = return_info['jobs'][0]
self.assertEqual(job['job_class_string'], data['job_class_string'])
self.assertEqual(job['name'], data['name'])
- self.assertEqual(job['minute'], data['minute'])
+ self.assertEqual(job['trigger_params']['minute'], data['trigger_params']['minute'])
def test_delete_job(self):
headers = {'Content-Type': 'application/json; charset=UTF-8'}
data = {
'job_class_string': 'hello.world',
'name': 'hello world job',
- 'minute': '*/5'}
+ 'trigger': 'cron',
+ 'trigger_params': {
+ 'minute': '*/5'}
+ }
response = self.fetch(self.JOBS_URL, method='POST', headers=headers,
body=json.dumps(data))
return_info = json.loads(response.body.decode())
@@ -161,7 +176,10 @@ def test_modify_job(self):
data = {
'job_class_string': 'hello.world',
'name': 'hello world job',
- 'minute': '*/5'}
+ 'trigger': 'cron',
+ 'trigger_params': {
+ 'minute': '*/5'}
+ }
response = self.fetch(self.JOBS_URL, method='POST', headers=headers,
body=json.dumps(data))
return_info = json.loads(response.body.decode())
@@ -175,7 +193,10 @@ def test_modify_job(self):
data = {
'job_class_string': 'hello.world!!!!',
'name': 'hello world job~~~~',
- 'minute': '*/20'}
+ 'trigger': 'cron',
+ 'trigger_params': {
+ 'minute': '*/20'}
+ }
response = self.fetch(self.JOBS_URL + '/' + return_info['job_id'] + '?sync=true',
method='PUT', headers=headers, body=json.dumps(data))
self.assertEqual(response.code, 200)
diff --git a/ndscheduler/server/server.py b/ndscheduler/server/server.py
index 113b6430..29b8866c 100644
--- a/ndscheduler/server/server.py
+++ b/ndscheduler/server/server.py
@@ -23,7 +23,7 @@
class SchedulerServer:
- VERSION = 'v1'
+ VERSION = 'v2'
singleton = None
diff --git a/ndscheduler/static/css/main.css b/ndscheduler/static/css/main.css
index bd839586..07ca7e54 100644
--- a/ndscheduler/static/css/main.css
+++ b/ndscheduler/static/css/main.css
@@ -189,4 +189,9 @@ body {
width: 16px !important;
padding-right: 10px !important;
padding-left: 10px !important;
-}
\ No newline at end of file
+}
+
+
+.modal-tab-content {
+ margin-bottom: 10px;
+}
diff --git a/ndscheduler/static/js/config.js b/ndscheduler/static/js/config.js
index 55ebecef..a89aa4dd 100644
--- a/ndscheduler/static/js/config.js
+++ b/ndscheduler/static/js/config.js
@@ -7,7 +7,7 @@ define([], function() {
'use strict';
- var urlPrefix = '/api/v1';
+ var urlPrefix = '/api/v2';
return {
'jobs_url': urlPrefix + '/jobs',
diff --git a/ndscheduler/static/js/models/job.js b/ndscheduler/static/js/models/job.js
index 5f95a707..a6b2fcab 100644
--- a/ndscheduler/static/js/models/job.js
+++ b/ndscheduler/static/js/models/job.js
@@ -20,6 +20,35 @@ require.config({
}
});
+function secondsToStr( seconds_in ) {
+ let temp = seconds_in;
+ const years = Math.floor( temp / 31536000 ),
+ days = Math.floor( ( temp %= 31536000 ) / 86400 ),
+ hours = Math.floor( ( temp %= 86400 ) / 3600 ),
+ minutes = Math.floor( ( temp %= 3600 ) / 60 ),
+ seconds = temp % 60;
+
+ if ( days || hours || seconds || minutes ) {
+ return ( years ? years + "y " : "" ) +
+ ( days ? days + "d " : "" ) +
+ ( hours ? hours + "h " : "" ) +
+ ( minutes ? minutes + "m " : "" ) +
+ Number.parseFloat( seconds ).toFixed( 2 ) + "s";
+ }
+
+ return "< 1s";
+}
+
+function secondsToObj( seconds_in ) {
+ let temp = seconds_in;
+
+ return {'days': Math.floor( temp / 86400 ),
+ 'hours': Math.floor( ( temp %= 86400 ) / 3600 ),
+ 'minutes': Math.floor( ( temp %= 3600 ) / 60 ),
+ 'seconds': temp % 60}
+}
+
+
define(['backbone', 'vendor/moment-timezone-with-data'], function(backbone, moment) {
'use strict';
@@ -31,9 +60,19 @@ define(['backbone', 'vendor/moment-timezone-with-data'], function(backbone, mome
* @return {string} schedule string for this job.
*/
getScheduleString: function() {
- return 'minute: ' + this.get('minute') + ', hour: ' + this.get('hour') +
- ', day: ' + this.get('day') + ', month: ' + this.get('month') +
- ', day of week: ' + this.get('day_of_week');
+ var trig = this.get('trigger');
+ var params = this.get('trigger_params');
+
+ if (trig == 'cron')
+ return 'Cron: minute: ' + params.minute + ', hour: ' + params.hour +
+ ', day: ' + params.day + ', month: ' + params.month +
+ ', day of week: ' + params.day_of_week;
+ else if (trig == 'interval')
+ return 'Interval: ' + secondsToStr(params.interval);
+ else
+ return 'Unknown trigger type!';
+
+
},
/**
diff --git a/ndscheduler/static/js/templates/add-job.html b/ndscheduler/static/js/templates/add-job.html
index bf6dfc10..83450fc2 100644
--- a/ndscheduler/static/js/templates/add-job.html
+++ b/ndscheduler/static/js/templates/add-job.html
@@ -35,38 +35,80 @@