sphinx-compatible docstrings (7 modules remaining...)
This commit is contained in:
+185
-41
@@ -1,5 +1,13 @@
|
||||
#!/usr/bin/env python
|
||||
# -*- coding: utf-8 -*-
|
||||
"""
|
||||
| This file is part of the web2py Web Framework
|
||||
| Copyrighted by Massimo Di Pierro <mdipierro@cs.depaul.edu>
|
||||
| License: LGPLv3 (http://www.gnu.org/licenses/lgpl.html)
|
||||
|
||||
Background processes made simple
|
||||
---------------------------------
|
||||
"""
|
||||
|
||||
USAGE = """
|
||||
## Example
|
||||
@@ -118,6 +126,9 @@ CALLABLETYPES = (types.LambdaType, types.FunctionType,
|
||||
|
||||
|
||||
class Task(object):
|
||||
"""Defines a "task" object that gets passed from the main thread to the
|
||||
executor's one
|
||||
"""
|
||||
def __init__(self, app, function, timeout, args='[]', vars='{}', **kwargs):
|
||||
logger.debug(' new task allocated: %s.%s', app, function)
|
||||
self.app = app
|
||||
@@ -132,6 +143,9 @@ class Task(object):
|
||||
|
||||
|
||||
class TaskReport(object):
|
||||
"""Defines a "task report" object that gets passed from the executor's
|
||||
thread to the main one
|
||||
"""
|
||||
def __init__(self, status, result=None, output=None, tb=None):
|
||||
logger.debug(' new task report: %s', status)
|
||||
if tb:
|
||||
@@ -184,7 +198,7 @@ def _decode_dict(dct):
|
||||
|
||||
|
||||
def executor(queue, task, out):
|
||||
""" the background process """
|
||||
"""The function used to execute tasks in the background process"""
|
||||
logger.debug(' task started')
|
||||
|
||||
class LogOutput(object):
|
||||
@@ -249,20 +263,28 @@ def executor(queue, task, out):
|
||||
|
||||
|
||||
class MetaScheduler(threading.Thread):
|
||||
"""Base class documenting scheduler's base methods"""
|
||||
|
||||
def __init__(self):
|
||||
threading.Thread.__init__(self)
|
||||
self.process = None # the background process
|
||||
self.have_heartbeat = True # set to False to kill
|
||||
self.empty_runs = 0
|
||||
|
||||
|
||||
def async(self, task):
|
||||
"""
|
||||
starts the background process and returns:
|
||||
('ok',result,output)
|
||||
('error',exception,None)
|
||||
('timeout',None,None)
|
||||
('terminated',None,None)
|
||||
"""Starts the background process
|
||||
|
||||
Args:
|
||||
task : a `Task` object
|
||||
|
||||
Returns:
|
||||
tuple: containing::
|
||||
|
||||
('ok',result,output)
|
||||
('error',exception,None)
|
||||
('timeout',None,None)
|
||||
('terminated',None,None)
|
||||
|
||||
"""
|
||||
db = self.db
|
||||
sr = db.scheduler_run
|
||||
@@ -332,22 +354,27 @@ class MetaScheduler(threading.Thread):
|
||||
return tr
|
||||
|
||||
def die(self):
|
||||
"""Forces termination of the worker process along with any running
|
||||
task"""
|
||||
logger.info('die!')
|
||||
self.have_heartbeat = False
|
||||
self.terminate_process()
|
||||
|
||||
def give_up(self):
|
||||
"""Waits for any running task to be executed, then exits the worker
|
||||
process"""
|
||||
logger.info('Giving up as soon as possible!')
|
||||
self.have_heartbeat = False
|
||||
|
||||
def terminate_process(self):
|
||||
"""Terminates any running tasks (internal use only)"""
|
||||
try:
|
||||
self.process.terminate()
|
||||
except:
|
||||
pass # no process to terminate
|
||||
|
||||
def run(self):
|
||||
""" the thread that sends heartbeat """
|
||||
"""This is executed by the main thread to send heartbeats"""
|
||||
counter = 0
|
||||
while self.have_heartbeat:
|
||||
self.send_heartbeat(counter)
|
||||
@@ -361,6 +388,7 @@ class MetaScheduler(threading.Thread):
|
||||
time.sleep(1)
|
||||
|
||||
def pop_task(self):
|
||||
"""Fetches a task ready to be executed"""
|
||||
return Task(
|
||||
app=None,
|
||||
function='demo_function',
|
||||
@@ -369,6 +397,7 @@ class MetaScheduler(threading.Thread):
|
||||
vars='{}')
|
||||
|
||||
def report_task(self, task, task_report):
|
||||
"""Creates a task report"""
|
||||
print 'reporting task'
|
||||
pass
|
||||
|
||||
@@ -376,6 +405,8 @@ class MetaScheduler(threading.Thread):
|
||||
pass
|
||||
|
||||
def loop(self):
|
||||
"""Main loop, fetching tasks and starting executor's background
|
||||
processes"""
|
||||
try:
|
||||
self.start_heartbeats()
|
||||
while True and self.have_heartbeat:
|
||||
@@ -406,7 +437,8 @@ WORKER_STATUS = (ACTIVE, PICK, DISABLED, TERMINATE, KILL, STOP_TASK)
|
||||
|
||||
class TYPE(object):
|
||||
"""
|
||||
validator that check whether field is valid json and validate its type
|
||||
Validator that checks whether field is valid json and validates its type.
|
||||
Used for `args` and `vars` of the scheduler_task table
|
||||
"""
|
||||
|
||||
def __init__(self, myclass=list, parse=False):
|
||||
@@ -430,6 +462,32 @@ class TYPE(object):
|
||||
|
||||
|
||||
class Scheduler(MetaScheduler):
|
||||
"""Scheduler object
|
||||
|
||||
Args:
|
||||
db: DAL connection where Scheduler will create its tables
|
||||
tasks(dict): either a dict containing name-->func or None.
|
||||
If None, functions will be searched in the environment
|
||||
migrate(bool): turn migration on/off for the Scheduler's tables
|
||||
worker_name(str): force worker_name to identify each process.
|
||||
Leave it to None to autoassign a name (hostname#pid)
|
||||
group_names(list): process tasks belonging to this group
|
||||
heartbeat(int): how many seconds the worker sleeps between one execution
|
||||
and the following one. Indirectly sets how many seconds will pass
|
||||
between checks for new tasks
|
||||
max_empty_runs(int): how many loops are allowed to pass without
|
||||
processing any tasks before exiting the process. 0 to keep always
|
||||
the process alive
|
||||
discard_results(bool): Scheduler stores executions's details into the
|
||||
scheduler_run table. By default, only if there is a result the
|
||||
details are kept. Turning this to True means discarding results
|
||||
even for tasks that return something
|
||||
utc_time(bool): do all datetime calculations assuming UTC as the
|
||||
timezone. Remember to pass `start_time` and `stop_time` to tasks
|
||||
accordingly
|
||||
|
||||
"""
|
||||
|
||||
def __init__(self, db, tasks=None, migrate=True,
|
||||
worker_name=None, group_names=['main'], heartbeat=HEARTBEAT,
|
||||
max_empty_runs=0, discard_results=False, utc_time=False):
|
||||
@@ -467,9 +525,11 @@ class Scheduler(MetaScheduler):
|
||||
return True
|
||||
|
||||
def now(self):
|
||||
"""Shortcut that fetches current time based on UTC preferences"""
|
||||
return self.utc_time and datetime.datetime.utcnow() or datetime.datetime.now()
|
||||
|
||||
def set_requirements(self, scheduler_task):
|
||||
"""Called to set defaults for lazy_tables connections"""
|
||||
from gluon import current
|
||||
if hasattr(current, 'request'):
|
||||
scheduler_task.application_name.default = '%s/%s' % (
|
||||
@@ -477,6 +537,7 @@ class Scheduler(MetaScheduler):
|
||||
)
|
||||
|
||||
def define_tables(self, db, migrate):
|
||||
"""Defines Scheduler tables structure"""
|
||||
from gluon.dal import DEFAULT
|
||||
logger.debug('defining tables (migrate=%s)', migrate)
|
||||
now = self.now
|
||||
@@ -550,6 +611,23 @@ class Scheduler(MetaScheduler):
|
||||
db.commit()
|
||||
|
||||
def loop(self, worker_name=None):
|
||||
"""Main loop
|
||||
|
||||
This works basically as a neverending loop that:
|
||||
|
||||
- checks if the worker is ready to process tasks (is not DISABLED)
|
||||
- pops a task from the queue
|
||||
- if there is a task:
|
||||
|
||||
- spawns the executor background process
|
||||
- waits for the process to be finished
|
||||
- sleeps `heartbeat` seconds
|
||||
- if there is not a task:
|
||||
|
||||
- checks for max_empty_runs
|
||||
- sleeps `heartbeat` seconds
|
||||
|
||||
"""
|
||||
signal.signal(signal.SIGTERM, lambda signum, stack_frame: sys.exit(1))
|
||||
try:
|
||||
self.start_heartbeats()
|
||||
@@ -581,6 +659,10 @@ class Scheduler(MetaScheduler):
|
||||
self.die()
|
||||
|
||||
def wrapped_assign_tasks(self, db):
|
||||
"""Commodity function to call `assign_tasks` and trap exceptions
|
||||
If an exception is raised, assume it happened because of database
|
||||
contention and retries `assign_task` after 0.5 seconds
|
||||
"""
|
||||
logger.debug('Assigning tasks...')
|
||||
db.commit() #db.commit() only for Mysql
|
||||
x = 0
|
||||
@@ -597,6 +679,10 @@ class Scheduler(MetaScheduler):
|
||||
time.sleep(0.5)
|
||||
|
||||
def wrapped_pop_task(self):
|
||||
"""Commodity function to call `pop_task` and trap exceptions
|
||||
If an exception is raised, assume it happened because of database
|
||||
contention and retries `pop_task` after 0.5 seconds
|
||||
"""
|
||||
db = self.db
|
||||
db.commit() #another nifty db.commit() only for Mysql
|
||||
x = 0
|
||||
@@ -612,6 +698,7 @@ class Scheduler(MetaScheduler):
|
||||
time.sleep(0.5)
|
||||
|
||||
def pop_task(self, db):
|
||||
"""Grabs a task ready to be executed from the queue"""
|
||||
now = self.now()
|
||||
st = self.db.scheduler_task
|
||||
if self.is_a_ticker and self.do_assign_tasks:
|
||||
@@ -685,6 +772,8 @@ class Scheduler(MetaScheduler):
|
||||
uuid=task.uuid)
|
||||
|
||||
def report_task(self, task, task_report):
|
||||
"""Takes care of storing the result according to preferences
|
||||
and deals with logic for repeating tasks"""
|
||||
db = self.db
|
||||
now = self.now()
|
||||
while True:
|
||||
@@ -744,12 +833,25 @@ class Scheduler(MetaScheduler):
|
||||
time.sleep(0.5)
|
||||
|
||||
def adj_hibernation(self):
|
||||
"""Used to increase the "sleep" interval for DISABLED workers"""
|
||||
if self.worker_status[0] == DISABLED:
|
||||
wk_st = self.worker_status[1]
|
||||
hibernation = wk_st + 1 if wk_st < MAXHIBERNATION else MAXHIBERNATION
|
||||
self.worker_status[1] = hibernation
|
||||
|
||||
def send_heartbeat(self, counter):
|
||||
"""This function is vital for proper coordination among available
|
||||
workers.
|
||||
It:
|
||||
|
||||
- sends the heartbeat
|
||||
- elects a ticker among available workers (the only process that
|
||||
effectively dispatch tasks to workers)
|
||||
- deals with worker's statuses
|
||||
- does "housecleaning" for dead workers
|
||||
- triggers tasks assignment to workers
|
||||
|
||||
"""
|
||||
if not self.db_thread:
|
||||
logger.debug('thread building own DAL object')
|
||||
self.db_thread = DAL(
|
||||
@@ -828,6 +930,10 @@ class Scheduler(MetaScheduler):
|
||||
self.sleep()
|
||||
|
||||
def being_a_ticker(self):
|
||||
"""Elects a TICKER process that assigns tasks to available workers.
|
||||
Does its best to elect a worker that is not busy processing other tasks
|
||||
to allow a proper distribution of tasks among all active workers ASAP
|
||||
"""
|
||||
db = self.db_thread
|
||||
sw = db.scheduler_worker
|
||||
all_active = db(
|
||||
@@ -857,6 +963,11 @@ class Scheduler(MetaScheduler):
|
||||
return False
|
||||
|
||||
def assign_tasks(self, db):
|
||||
"""Assigns task to workers, that can then pop them from the queue
|
||||
|
||||
Deals with group_name(s) logic, in order to assign linearly tasks
|
||||
to available workers for those groups
|
||||
"""
|
||||
sw, st = db.scheduler_worker, db.scheduler_task
|
||||
now = self.now()
|
||||
all_workers = db(sw.status == ACTIVE).select()
|
||||
@@ -934,10 +1045,13 @@ class Scheduler(MetaScheduler):
|
||||
logger.info('TICKER: tasks are %s', x)
|
||||
|
||||
def sleep(self):
|
||||
"""Calculates the number of seconds to sleep according to worker's
|
||||
status and `heartbeat` parameter"""
|
||||
time.sleep(self.heartbeat * self.worker_status[1])
|
||||
# should only sleep until next available task
|
||||
|
||||
def set_worker_status(self, group_names=None, action=ACTIVE):
|
||||
"""Internal function to set worker's status"""
|
||||
if not group_names:
|
||||
group_names = self.group_names
|
||||
elif isinstance(group_names, str):
|
||||
@@ -948,32 +1062,47 @@ class Scheduler(MetaScheduler):
|
||||
).update(status=action)
|
||||
|
||||
def disable(self, group_names=None):
|
||||
"""Sets DISABLED on the workers processing `group_names` tasks.
|
||||
A DISABLED worker will be kept alive but it won't be able to process
|
||||
any waiting tasks, essentially putting it to sleep.
|
||||
By default, all group_names of Scheduler's instantation are selected"""
|
||||
self.set_worker_status(group_names=group_names,action=DISABLED)
|
||||
|
||||
def resume(self, group_names=None):
|
||||
"""Wakes a worker up (it will be able to process queued tasks)"""
|
||||
self.set_worker_status(group_names=group_names,action=ACTIVE)
|
||||
|
||||
def terminate(self, group_names=None):
|
||||
"""Sets TERMINATE as worker status. The worker will wait for any
|
||||
currently running tasks to be executed and then it will exit gracefully
|
||||
"""
|
||||
self.set_worker_status(group_names=group_names,action=TERMINATE)
|
||||
|
||||
def kill(self, group_names=None):
|
||||
"""Sets KILL as worker status. The worker will be killed even if it's
|
||||
processing a task."""
|
||||
self.set_worker_status(group_names=group_names,action=KILL)
|
||||
|
||||
def queue_task(self, function, pargs=[], pvars={}, **kwargs):
|
||||
"""
|
||||
Queue tasks. This takes care of handling the validation of all
|
||||
values.
|
||||
:param function: the function (anything callable with a __name__)
|
||||
:param pargs: "raw" args to be passed to the function. Automatically
|
||||
jsonified.
|
||||
:param pvars: "raw" kwargs to be passed to the function. Automatically
|
||||
jsonified
|
||||
:param kwargs: all the scheduler_task columns. args and vars here should be
|
||||
in json format already, they will override pargs and pvars
|
||||
parameters
|
||||
|
||||
returns a dict just as a normal validate_and_insert, plus a uuid key holding
|
||||
the uuid of the queued task. If validation is not passed, both id and uuid
|
||||
will be None, and you'll get an "error" dict holding the errors found.
|
||||
Args:
|
||||
function: the function (anything callable with a __name__)
|
||||
pargs: "raw" args to be passed to the function. Automatically
|
||||
jsonified.
|
||||
pvars: "raw" kwargs to be passed to the function. Automatically
|
||||
jsonified
|
||||
kwargs: all the parameters available (basically, every
|
||||
`scheduler_task` column). If args and vars are here, they should
|
||||
be jsonified already, and they will override pargs and pvars
|
||||
|
||||
Returns:
|
||||
a dict just as a normal validate_and_insert(), plus a uuid key
|
||||
holding the uuid of the queued task. If validation is not passed
|
||||
( i.e. some parameters are invalid) both id and uuid will be None,
|
||||
and you'll get an "error" dict holding the errors found.
|
||||
"""
|
||||
if hasattr(function, '__name__'):
|
||||
function = function.__name__
|
||||
@@ -999,17 +1128,23 @@ class Scheduler(MetaScheduler):
|
||||
|
||||
def task_status(self, ref, output=False):
|
||||
"""
|
||||
Shortcut for task status retrieval
|
||||
Retrieves task status and optionally the result of the task
|
||||
|
||||
:param ref: can be
|
||||
- integer --> lookup will be done by scheduler_task.id
|
||||
- string --> lookup will be done by scheduler_task.uuid
|
||||
- query --> lookup as you wish (as in db.scheduler_task.task_name == 'test1')
|
||||
:param output: fetch also the scheduler_run record
|
||||
Args:
|
||||
ref: can be
|
||||
|
||||
Returns a single Row object, for the last queued task
|
||||
If output == True, returns also the last scheduler_run record
|
||||
scheduler_run record is fetched by a left join, so it can
|
||||
- an integer : lookup will be done by scheduler_task.id
|
||||
- a string : lookup will be done by scheduler_task.uuid
|
||||
- a `Query` : lookup as you wish, e.g. ::
|
||||
|
||||
db.scheduler_task.task_name == 'test1'
|
||||
|
||||
output(bool): if `True`, fetch also the scheduler_run record
|
||||
|
||||
Returns:
|
||||
a single Row object, for the last queued task.
|
||||
If output == True, returns also the last scheduler_run record.
|
||||
The scheduler_run record is fetched by a left join, so it can
|
||||
have all fields == None
|
||||
|
||||
"""
|
||||
@@ -1044,21 +1179,27 @@ class Scheduler(MetaScheduler):
|
||||
return row
|
||||
|
||||
def stop_task(self, ref):
|
||||
"""
|
||||
Experimental!!!
|
||||
Shortcut for task termination.
|
||||
If the task is RUNNING it will terminate it --> execution will be set as FAILED
|
||||
If the task is QUEUED, its stop_time will be set as to "now",
|
||||
the enabled flag will be set to False, status to STOPPED
|
||||
"""Shortcut for task termination.
|
||||
|
||||
If the task is RUNNING it will terminate it, meaning that status
|
||||
will be set as FAILED.
|
||||
|
||||
If the task is QUEUED, its stop_time will be set as to "now",
|
||||
the enabled flag will be set to False, and the status to STOPPED
|
||||
|
||||
Args:
|
||||
ref: can be
|
||||
|
||||
- an integer : lookup will be done by scheduler_task.id
|
||||
- a string : lookup will be done by scheduler_task.uuid
|
||||
|
||||
:param ref: can be
|
||||
- integer --> lookup will be done by scheduler_task.id
|
||||
- string --> lookup will be done by scheduler_task.uuid
|
||||
Returns:
|
||||
- 1 if task was stopped (meaning an update has been done)
|
||||
- None if task was not found, or if task was not RUNNING or QUEUED
|
||||
|
||||
Note:
|
||||
Experimental
|
||||
"""
|
||||
from gluon.dal import Query
|
||||
st, sw = self.db.scheduler_task, self.db.scheduler_worker
|
||||
if isinstance(ref, int):
|
||||
q = st.id == ref
|
||||
@@ -1080,7 +1221,10 @@ class Scheduler(MetaScheduler):
|
||||
|
||||
def main():
|
||||
"""
|
||||
allows to run worker without python web2py.py .... by simply python this.py
|
||||
allows to run worker without python web2py.py .... by simply::
|
||||
|
||||
python gluon/scheduler.py
|
||||
|
||||
"""
|
||||
parser = optparse.OptionParser()
|
||||
parser.add_option(
|
||||
|
||||
Reference in New Issue
Block a user