better IS_NOT_IN_DB, new scheduler API, after_connection dal callback, thanks Niphlod
This commit is contained in:
+48
-22
@@ -86,7 +86,9 @@ try:
|
||||
except:
|
||||
from simplejson import loads, dumps
|
||||
|
||||
logger = logging.getLogger('web2py.scheduler')
|
||||
IDENTIFIER = "%s#%s" % (socket.gethostname(),os.getpid())
|
||||
|
||||
logger = logging.getLogger('web2py.scheduler.%s' % IDENTIFIER)
|
||||
|
||||
from gluon import DAL, Field, IS_NOT_EMPTY, IS_IN_SET, IS_NOT_IN_DB, IS_INT_IN_RANGE, IS_DATETIME
|
||||
from gluon.utils import web2py_uuid
|
||||
@@ -103,6 +105,7 @@ ACTIVE = 'ACTIVE'
|
||||
TERMINATE = 'TERMINATE'
|
||||
DISABLED = 'DISABLED'
|
||||
KILL = 'KILL'
|
||||
PICK = 'PICK'
|
||||
EXPIRED = 'EXPIRED'
|
||||
SECONDS = 1
|
||||
HEARTBEAT = 3 * SECONDS
|
||||
@@ -391,7 +394,7 @@ class MetaScheduler(threading.Thread):
|
||||
|
||||
TASK_STATUS = (QUEUED, RUNNING, COMPLETED, FAILED, TIMEOUT, STOPPED, EXPIRED)
|
||||
RUN_STATUS = (RUNNING, COMPLETED, FAILED, TIMEOUT, STOPPED)
|
||||
WORKER_STATUS = (ACTIVE, DISABLED, TERMINATE, KILL)
|
||||
WORKER_STATUS = (ACTIVE, PICK, DISABLED, TERMINATE, KILL)
|
||||
|
||||
|
||||
class TYPE(object):
|
||||
@@ -431,8 +434,7 @@ class Scheduler(MetaScheduler):
|
||||
self.tasks = tasks
|
||||
self.group_names = group_names
|
||||
self.heartbeat = heartbeat
|
||||
self.worker_name = worker_name or socket.gethostname(
|
||||
) + '#' + str(os.getpid())
|
||||
self.worker_name = worker_name or IDENTIFIER
|
||||
#list containing status as recorded in the table plus a boost parameter
|
||||
#for hibernation (i.e. when someone stop the worker acting on the worker table)
|
||||
self.worker_status = [RUNNING, 1]
|
||||
@@ -565,7 +567,7 @@ class Scheduler(MetaScheduler):
|
||||
break
|
||||
except:
|
||||
db.rollback()
|
||||
logger.error('TICKER(%s): error assigning tasks', self.worker_name)
|
||||
logger.error('TICKER: error assigning tasks')
|
||||
x += 1
|
||||
time.sleep(0.5)
|
||||
|
||||
@@ -580,10 +582,10 @@ class Scheduler(MetaScheduler):
|
||||
break
|
||||
except:
|
||||
db.rollback()
|
||||
logger.error('%s: error popping tasks', self.worker_name)
|
||||
logger.error(' error popping tasks')
|
||||
x += 1
|
||||
time.sleep(0.5)
|
||||
|
||||
|
||||
def pop_task(self, db):
|
||||
now = self.now()
|
||||
st = self.db.scheduler_task
|
||||
@@ -605,7 +607,7 @@ class Scheduler(MetaScheduler):
|
||||
else:
|
||||
if self.greedy and self.is_a_ticker:
|
||||
#there are other tasks ready to be assigned
|
||||
logger.info('TICKER (%s): greedy loop', self.worker_name)
|
||||
logger.info('TICKER: greedy loop')
|
||||
self.wrapped_assign_tasks(db)
|
||||
else:
|
||||
logger.info('nothing to do')
|
||||
@@ -726,31 +728,32 @@ class Scheduler(MetaScheduler):
|
||||
sw, st = db.scheduler_worker, db.scheduler_task
|
||||
now = self.now()
|
||||
# record heartbeat
|
||||
mybackedstatus = db(
|
||||
sw.worker_name == self.worker_name).select().first()
|
||||
mybackedstatus = db(sw.worker_name == self.worker_name).select().first()
|
||||
if not mybackedstatus:
|
||||
sw.insert(status=ACTIVE, worker_name=self.worker_name,
|
||||
first_heartbeat=now, last_heartbeat=now,
|
||||
group_names=self.group_names)
|
||||
self.worker_status = [ACTIVE, 1] # activating the process
|
||||
mybackedstatus = ACTIVE
|
||||
else:
|
||||
if mybackedstatus.status == DISABLED:
|
||||
mybackedstatus = mybackedstatus.status
|
||||
if mybackedstatus == DISABLED:
|
||||
# keep sleeping
|
||||
self.worker_status[0] = DISABLED
|
||||
if self.worker_status[1] == MAXHIBERNATION:
|
||||
logger.debug('........recording heartbeat')
|
||||
db(sw.worker_name == self.worker_name).update(
|
||||
last_heartbeat=now)
|
||||
elif mybackedstatus.status == TERMINATE:
|
||||
elif mybackedstatus == TERMINATE:
|
||||
self.worker_status[0] = TERMINATE
|
||||
logger.debug("Waiting to terminate the current task")
|
||||
self.give_up()
|
||||
return
|
||||
elif mybackedstatus.status == KILL:
|
||||
elif mybackedstatus == KILL:
|
||||
self.worker_status[0] = KILL
|
||||
self.die()
|
||||
else:
|
||||
logger.debug('........recording heartbeat (%s)', self.worker_status[0])
|
||||
logger.debug('........recording heartbeat (%s) %s', self.worker_status[0], mybackedstatus)
|
||||
db(sw.worker_name == self.worker_name).update(
|
||||
last_heartbeat=now, status=ACTIVE)
|
||||
self.worker_status[1] = 1 # re-activating the process
|
||||
@@ -758,8 +761,7 @@ class Scheduler(MetaScheduler):
|
||||
self.worker_status[0] = ACTIVE
|
||||
|
||||
self.do_assign_tasks = False
|
||||
|
||||
if counter % 5 == 0:
|
||||
if counter % 5 == 0 or mybackedstatus == PICK:
|
||||
try:
|
||||
# delete inactive workers
|
||||
expiration = now - datetime.timedelta(seconds=self.heartbeat * 3)
|
||||
@@ -769,8 +771,7 @@ class Scheduler(MetaScheduler):
|
||||
' freeing workers that have not sent heartbeat')
|
||||
inactive_workers = db(
|
||||
((sw.last_heartbeat < expiration) & (sw.status == ACTIVE)) |
|
||||
((sw.last_heartbeat <
|
||||
departure) & (sw.status != ACTIVE))
|
||||
((sw.last_heartbeat < departure) & (sw.status != ACTIVE))
|
||||
)
|
||||
db(st.assigned_worker_name.belongs(
|
||||
inactive_workers._select(sw.worker_name)))(st.status == RUNNING)\
|
||||
@@ -800,7 +801,7 @@ class Scheduler(MetaScheduler):
|
||||
#only if this worker isn't busy, otherwise wait for a free one
|
||||
db(sw.worker_name == self.worker_name).update(is_ticker=True)
|
||||
db(sw.worker_name != self.worker_name).update(is_ticker=False)
|
||||
logger.info("TICKER(%s): I'm a ticker", self.worker_name)
|
||||
logger.info("TICKER: I'm a ticker")
|
||||
else:
|
||||
#giving up, only if I'm not alone
|
||||
if len(all_active) > 1:
|
||||
@@ -888,13 +889,35 @@ class Scheduler(MetaScheduler):
|
||||
#I'll be greedy only if tasks assigned are equal to the limit
|
||||
# (meaning there could be others ready to be assigned)
|
||||
self.greedy = x >= limit and True or False
|
||||
logger.info('TICKER(%s): workers are %s', self.worker_name, len(all_workers))
|
||||
logger.info('TICKER(%s): tasks are %s', self.worker_name, x)
|
||||
logger.info('TICKER: workers are %s', len(all_workers))
|
||||
logger.info('TICKER: tasks are %s', x)
|
||||
|
||||
def sleep(self):
|
||||
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):
|
||||
if not group_names:
|
||||
group_names = self.group_names
|
||||
elif isinstance(group_names, str):
|
||||
group_names = [group_names]
|
||||
for group in group_names:
|
||||
self.db(
|
||||
self.db.scheduler_worker.group_names.contains(group)
|
||||
).update(status=action)
|
||||
|
||||
def disable(self, group_names=None):
|
||||
self.set_worker_status(group_names=group_names,action=DISABLED)
|
||||
|
||||
def resume(self, group_names=None):
|
||||
self.set_worker_status(group_names=group_names,action=ACTIVE)
|
||||
|
||||
def terminate(self, group_names=None):
|
||||
self.set_worker_status(group_names=group_names,action=TERMINATE)
|
||||
|
||||
def kill(self, group_names=None):
|
||||
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
|
||||
@@ -917,6 +940,7 @@ class Scheduler(MetaScheduler):
|
||||
tvars = 'vars' in kwargs and kwargs.pop('vars') or dumps(pvars)
|
||||
tuuid = 'uuid' in kwargs and kwargs.pop('uuid') or web2py_uuid()
|
||||
tname = 'task_name' in kwargs and kwargs.pop('task_name') or function
|
||||
immediate = 'immediate' in kwargs and kwargs.pop('immediate') or None
|
||||
rtn = self.db.scheduler_task.validate_and_insert(
|
||||
function_name=function,
|
||||
task_name=tname,
|
||||
@@ -926,6 +950,8 @@ class Scheduler(MetaScheduler):
|
||||
**kwargs)
|
||||
if not rtn.errors:
|
||||
rtn.uuid = tuuid
|
||||
if immediate:
|
||||
self.db(self.db.scheduler_worker.is_ticker == True).update(status=PICK)
|
||||
else:
|
||||
rtn.uuid = None
|
||||
return rtn
|
||||
@@ -970,7 +996,7 @@ class Scheduler(MetaScheduler):
|
||||
left=left,
|
||||
limitby=(0, 1))
|
||||
).first()
|
||||
if output:
|
||||
if row and output:
|
||||
row.result = row.scheduler_run.run_result and \
|
||||
loads(row.scheduler_run.run_result,
|
||||
object_hook=_decode_dict) or None
|
||||
|
||||
Reference in New Issue
Block a user