Package Dropbox :: Package web2py :: Package gluon :: Module scheduler
[hide private]
[frames] | no frames]

Source Code for Module Dropbox.web2py.gluon.scheduler

   1  #!/usr/bin/env python 
   2  # -*- coding: utf-8 -*- 
   3   
   4  USAGE = """ 
   5  ## Example 
   6   
   7  For any existing app 
   8   
   9  Create File: app/models/scheduler.py ====== 
  10  from gluon.scheduler import Scheduler 
  11   
  12  def demo1(*args,**vars): 
  13      print 'you passed args=%s and vars=%s' % (args, vars) 
  14      return 'done!' 
  15   
  16  def demo2(): 
  17      1/0 
  18   
  19  scheduler = Scheduler(db,dict(demo1=demo1,demo2=demo2)) 
  20  ## run worker nodes with: 
  21   
  22     cd web2py 
  23     python web2py.py -K myapp 
  24  or 
  25     python gluon/scheduler.py -u sqlite://storage.sqlite \ 
  26                               -f applications/myapp/databases/ \ 
  27                               -t mytasks.py 
  28  (-h for info) 
  29  python scheduler.py -h 
  30   
  31  ## schedule jobs using 
  32  http://127.0.0.1:8000/myapp/appadmin/insert/db/scheduler_task 
  33   
  34  ## monitor scheduled jobs 
  35  http://127.0.0.1:8000/myapp/appadmin/select/db?query=db.scheduler_task.id>0 
  36   
  37  ## view completed jobs 
  38  http://127.0.0.1:8000/myapp/appadmin/select/db?query=db.scheduler_run.id>0 
  39   
  40  ## view workers 
  41  http://127.0.0.1:8000/myapp/appadmin/select/db?query=db.scheduler_worker.id>0 
  42   
  43  ## To install the scheduler as a permanent daemon on Linux (w/ Upstart), put 
  44  ## the following into /etc/init/web2py-scheduler.conf: 
  45  ## (This assumes your web2py instance is installed in <user>'s home directory, 
  46  ## running as <user>, with app <myapp>, on network interface eth0.) 
  47   
  48  description "web2py task scheduler" 
  49  start on (local-filesystems and net-device-up IFACE=eth0) 
  50  stop on shutdown 
  51  respawn limit 8 60 # Give up if restart occurs 8 times in 60 seconds. 
  52  exec sudo -u <user> python /home/<user>/web2py/web2py.py -K <myapp> 
  53  respawn 
  54   
  55  ## You can then start/stop/restart/check status of the daemon with: 
  56  sudo start web2py-scheduler 
  57  sudo stop web2py-scheduler 
  58  sudo restart web2py-scheduler 
  59  sudo status web2py-scheduler 
  60  """ 
  61   
  62  import os 
  63  import time 
  64  import multiprocessing 
  65  import sys 
  66  import threading 
  67  import traceback 
  68  import signal 
  69  import socket 
  70  import datetime 
  71  import logging 
  72  import optparse 
  73  import types 
  74  import Queue 
  75   
  76  if 'WEB2PY_PATH' in os.environ: 
  77      sys.path.append(os.environ['WEB2PY_PATH']) 
  78  else: 
  79      os.environ['WEB2PY_PATH'] = os.getcwd() 
  80   
  81  if not os.environ['WEB2PY_PATH'] in sys.path: 
  82      sys.path.append(os.environ['WEB2PY_PATH']) 
  83   
  84  try: 
  85      from gluon.contrib.simplejson import loads, dumps 
  86  except: 
  87      from simplejson import loads, dumps 
  88   
  89  logger = logging.getLogger('web2py.scheduler') 
  90   
  91  from gluon import DAL, Field, IS_NOT_EMPTY, IS_IN_SET, IS_NOT_IN_DB, IS_INT_IN_RANGE, IS_DATETIME 
  92  from gluon.utils import web2py_uuid 
  93   
  94   
  95  QUEUED = 'QUEUED' 
  96  ASSIGNED = 'ASSIGNED' 
  97  RUNNING = 'RUNNING' 
  98  COMPLETED = 'COMPLETED' 
  99  FAILED = 'FAILED' 
 100  TIMEOUT = 'TIMEOUT' 
 101  STOPPED = 'STOPPED' 
 102  ACTIVE = 'ACTIVE' 
 103  TERMINATE = 'TERMINATE' 
 104  DISABLED = 'DISABLED' 
 105  KILL = 'KILL' 
 106  EXPIRED = 'EXPIRED' 
 107  SECONDS = 1 
 108  HEARTBEAT = 3 * SECONDS 
 109  MAXHIBERNATION = 10 
 110  CLEAROUT = '!clear!' 
 111   
 112  CALLABLETYPES = (types.LambdaType, types.FunctionType, 
 113                   types.BuiltinFunctionType, 
 114                   types.MethodType, types.BuiltinMethodType) 
 115   
 116   
117 -class Task(object):
118 - def __init__(self, app, function, timeout, args='[]', vars='{}', **kwargs):
119 logger.debug(' new task allocated: %s.%s', app, function) 120 self.app = app 121 self.function = function 122 self.timeout = timeout 123 self.args = args # json 124 self.vars = vars # json 125 self.__dict__.update(kwargs)
126
127 - def __str__(self):
128 return '<Task: %s>' % self.function
129 130
131 -class TaskReport(object):
132 - def __init__(self, status, result=None, output=None, tb=None):
133 logger.debug(' new task report: %s', status) 134 if tb: 135 logger.debug(' traceback: %s', tb) 136 else: 137 logger.debug(' result: %s', result) 138 self.status = status 139 self.result = result 140 self.output = output 141 self.tb = tb
142
143 - def __str__(self):
144 return '<TaskReport: %s>' % self.status
145 146
147 -def demo_function(*argv, **kwargs):
148 """ test function """ 149 for i in range(argv[0]): 150 print 'click', i 151 time.sleep(1) 152 return 'done'
153 154 #the two functions below deal with simplejson decoding as unicode, esp for the dict decode 155 #and subsequent usage as function Keyword arguments unicode variable names won't work! 156 #borrowed from http://stackoverflow.com/questions/956867/how-to-get-string-objects-instead-unicode-ones-from-json-in-python 157 158
159 -def _decode_list(lst):
160 newlist = [] 161 for i in lst: 162 if isinstance(i, unicode): 163 i = i.encode('utf-8') 164 elif isinstance(i, list): 165 i = _decode_list(i) 166 newlist.append(i) 167 return newlist
168 169
170 -def _decode_dict(dct):
171 newdict = {} 172 for k, v in dct.iteritems(): 173 if isinstance(k, unicode): 174 k = k.encode('utf-8') 175 if isinstance(v, unicode): 176 v = v.encode('utf-8') 177 elif isinstance(v, list): 178 v = _decode_list(v) 179 newdict[k] = v 180 return newdict
181 182
183 -def executor(queue, task, out):
184 """ the background process """ 185 logger.debug(' task started') 186 187 class LogOutput(object): 188 """Facility to log output at intervals""" 189 def __init__(self, out_queue): 190 self.out_queue = out_queue 191 self.stdout = sys.stdout 192 sys.stdout = self
193 194 def __del__(self): 195 sys.stdout = self.stdout 196 197 def flush(self): 198 pass 199 200 def write(self, data): 201 self.out_queue.put(data) 202 203 stdout = LogOutput(out) 204 try: 205 if task.app: 206 os.chdir(os.environ['WEB2PY_PATH']) 207 from gluon.shell import env, parse_path_info 208 from gluon import current 209 level = logging.getLogger().getEffectiveLevel() 210 logging.getLogger().setLevel(logging.WARN) 211 # Get controller-specific subdirectory if task.app is of 212 # form 'app/controller' 213 (a, c, f) = parse_path_info(task.app) 214 _env = env(a=a, c=c, import_models=True) 215 logging.getLogger().setLevel(level) 216 f = task.function 217 functions = current._scheduler.tasks 218 if not functions: 219 #look into env 220 _function = _env.get(f) 221 else: 222 _function = functions.get(f) 223 if not isinstance(_function, CALLABLETYPES): 224 raise NameError( 225 "name '%s' not found in scheduler's environment" % f) 226 globals().update(_env) 227 args = loads(task.args) 228 vars = loads(task.vars, object_hook=_decode_dict) 229 result = dumps(_function(*args, **vars)) 230 else: 231 ### for testing purpose only 232 result = eval(task.function)( 233 *loads(task.args, object_hook=_decode_dict), 234 **loads(task.vars, object_hook=_decode_dict)) 235 queue.put(TaskReport(COMPLETED, result=result)) 236 except BaseException, e: 237 tb = traceback.format_exc() 238 queue.put(TaskReport(FAILED, tb=tb)) 239 del stdout 240 241
242 -class MetaScheduler(threading.Thread):
243 - def __init__(self):
244 threading.Thread.__init__(self) 245 self.process = None # the background process 246 self.have_heartbeat = True # set to False to kill 247 self.empty_runs = 0
248
249 - def async(self, task):
250 """ 251 starts the background process and returns: 252 ('ok',result,output) 253 ('error',exception,None) 254 ('timeout',None,None) 255 ('terminated',None,None) 256 """ 257 db = self.db 258 sr = db.scheduler_run 259 out = multiprocessing.Queue() 260 queue = multiprocessing.Queue(maxsize=1) 261 p = multiprocessing.Process(target=executor, args=(queue, task, out)) 262 self.process = p 263 logger.debug(' task starting') 264 p.start() 265 266 task_output = "" 267 tout = "" 268 269 try: 270 if task.sync_output > 0: 271 run_timeout = task.sync_output 272 else: 273 run_timeout = task.timeout 274 275 start = time.time() 276 277 while p.is_alive() and ( 278 not task.timeout or time.time() - start < task.timeout): 279 if tout: 280 try: 281 logger.debug(' partial output saved') 282 db(sr.id == task.run_id).update(run_output=task_output) 283 db.commit() 284 except: 285 pass 286 p.join(timeout=run_timeout) 287 tout = "" 288 while not out.empty(): 289 tout += out.get() 290 if tout: 291 logger.debug(' partial output: "%s"' % str(tout)) 292 if CLEAROUT in tout: 293 task_output = tout[ 294 tout.rfind(CLEAROUT) + len(CLEAROUT):] 295 else: 296 task_output += tout 297 except: 298 p.terminate() 299 p.join() 300 self.have_heartbeat = False 301 logger.debug(' task stopped by general exception') 302 tr = TaskReport(STOPPED) 303 else: 304 if p.is_alive(): 305 p.terminate() 306 logger.debug(' task timeout') 307 try: 308 # we try to get a traceback here 309 tr = queue.get(timeout=2) 310 tr.status = TIMEOUT 311 tr.output = task_output 312 except Queue.Empty: 313 tr = TaskReport(TIMEOUT) 314 elif queue.empty(): 315 self.have_heartbeat = False 316 logger.debug(' task stopped') 317 tr = TaskReport(STOPPED) 318 else: 319 logger.debug(' task completed or failed') 320 tr = queue.get() 321 tr.output = task_output 322 return tr
323
324 - def die(self):
325 logger.info('die!') 326 self.have_heartbeat = False 327 self.terminate_process()
328
329 - def give_up(self):
330 logger.info('Giving up as soon as possible!') 331 self.have_heartbeat = False
332
333 - def terminate_process(self):
334 try: 335 self.process.terminate() 336 except: 337 pass # no process to terminate
338
339 - def run(self):
340 """ the thread that sends heartbeat """ 341 counter = 0 342 while self.have_heartbeat: 343 self.send_heartbeat(counter) 344 counter += 1
345
346 - def start_heartbeats(self):
347 self.start()
348
349 - def send_heartbeat(self, counter):
350 print 'thum' 351 time.sleep(1)
352
353 - def pop_task(self):
354 return Task( 355 app=None, 356 function='demo_function', 357 timeout=7, 358 args='[2]', 359 vars='{}')
360
361 - def report_task(self, task, task_report):
362 print 'reporting task' 363 pass
364
365 - def sleep(self):
366 pass
367
368 - def loop(self):
369 try: 370 self.start_heartbeats() 371 while True and self.have_heartbeat: 372 logger.debug('looping...') 373 task = self.pop_task() 374 if task: 375 self.empty_runs = 0 376 self.report_task(task, self.async(task)) 377 else: 378 self.empty_runs += 1 379 logger.debug('sleeping...') 380 if self.max_empty_runs != 0: 381 logger.debug('empty runs %s/%s', 382 self.empty_runs, self.max_empty_runs) 383 if self.empty_runs >= self.max_empty_runs: 384 logger.info( 385 'empty runs limit reached, killing myself') 386 self.die() 387 self.sleep() 388 except KeyboardInterrupt: 389 self.die()
390 391 392 TASK_STATUS = (QUEUED, RUNNING, COMPLETED, FAILED, TIMEOUT, STOPPED, EXPIRED) 393 RUN_STATUS = (RUNNING, COMPLETED, FAILED, TIMEOUT, STOPPED) 394 WORKER_STATUS = (ACTIVE, DISABLED, TERMINATE, KILL) 395 396
397 -class TYPE(object):
398 """ 399 validator that check whether field is valid json and validate its type 400 """ 401
402 - def __init__(self, myclass=list, parse=False):
403 self.myclass = myclass 404 self.parse = parse
405
406 - def __call__(self, value):
407 from gluon import current 408 try: 409 obj = loads(value) 410 except: 411 return (value, current.T('invalid json')) 412 else: 413 if isinstance(obj, self.myclass): 414 if self.parse: 415 return (obj, None) 416 else: 417 return (value, None) 418 else: 419 return (value, current.T('Not of type: %s') % self.myclass)
420 421
422 -class Scheduler(MetaScheduler):
423 - def __init__(self, db, tasks=None, migrate=True, 424 worker_name=None, group_names=['main'], heartbeat=HEARTBEAT, 425 max_empty_runs=0, discard_results=False, utc_time=False):
426 427 MetaScheduler.__init__(self) 428 429 self.db = db 430 self.db_thread = None 431 self.tasks = tasks 432 self.group_names = group_names 433 self.heartbeat = heartbeat 434 self.worker_name = worker_name or socket.gethostname( 435 ) + '#' + str(os.getpid()) 436 #list containing status as recorded in the table plus a boost parameter 437 #for hibernation (i.e. when someone stop the worker acting on the worker table) 438 self.worker_status = [RUNNING, 1] 439 self.max_empty_runs = max_empty_runs 440 self.discard_results = discard_results 441 self.is_a_ticker = False 442 self.do_assign_tasks = False 443 self.greedy = False 444 self.utc_time = utc_time 445 446 from gluon import current 447 current._scheduler = self 448 449 self.define_tables(db, migrate=migrate)
450
451 - def now(self):
452 return self.utc_time and datetime.datetime.utcnow() or datetime.datetime.now()
453
454 - def set_requirements(self, scheduler_task):
455 from gluon import current 456 if hasattr(current, 'request'): 457 scheduler_task.application_name.default = '%s/%s' % ( 458 current.request.application, current.request.controller 459 )
460
461 - def define_tables(self, db, migrate):
462 from gluon.dal import DEFAULT 463 logger.debug('defining tables (migrate=%s)', migrate) 464 now = self.now 465 db.define_table( 466 'scheduler_task', 467 Field('application_name', requires=IS_NOT_EMPTY(), 468 default=None, writable=False), 469 Field('task_name', default=None), 470 Field('group_name', default='main'), 471 Field('status', requires=IS_IN_SET(TASK_STATUS), 472 default=QUEUED, writable=False), 473 Field('function_name', 474 requires=IS_IN_SET(sorted(self.tasks.keys())) 475 if self.tasks else DEFAULT), 476 Field('uuid', requires=IS_NOT_IN_DB(db, 'scheduler_task.uuid'), 477 unique=True, default=web2py_uuid), 478 Field('args', 'text', default='[]', requires=TYPE(list)), 479 Field('vars', 'text', default='{}', requires=TYPE(dict)), 480 Field('enabled', 'boolean', default=True), 481 Field('start_time', 'datetime', default=now, 482 requires=IS_DATETIME()), 483 Field('next_run_time', 'datetime', default=now), 484 Field('stop_time', 'datetime'), 485 Field('repeats', 'integer', default=1, comment="0=unlimited", 486 requires=IS_INT_IN_RANGE(0, None)), 487 Field('retry_failed', 'integer', default=0, comment="-1=unlimited", 488 requires=IS_INT_IN_RANGE(-1, None)), 489 Field('period', 'integer', default=60, comment='seconds', 490 requires=IS_INT_IN_RANGE(0, None)), 491 Field('timeout', 'integer', default=60, comment='seconds', 492 requires=IS_INT_IN_RANGE(0, None)), 493 Field('sync_output', 'integer', default=0, 494 comment="update output every n sec: 0=never", 495 requires=IS_INT_IN_RANGE(0, None)), 496 Field('times_run', 'integer', default=0, writable=False), 497 Field('times_failed', 'integer', default=0, writable=False), 498 Field('last_run_time', 'datetime', writable=False, readable=False), 499 Field('assigned_worker_name', default='', writable=False), 500 on_define=self.set_requirements, 501 migrate=migrate, format='%(task_name)s') 502 503 db.define_table( 504 'scheduler_run', 505 Field('scheduler_task', 'reference scheduler_task'), 506 Field('status', requires=IS_IN_SET(RUN_STATUS)), 507 Field('start_time', 'datetime'), 508 Field('stop_time', 'datetime'), 509 Field('run_output', 'text'), 510 Field('run_result', 'text'), 511 Field('traceback', 'text'), 512 Field('worker_name', default=self.worker_name), 513 migrate=migrate) 514 515 db.define_table( 516 'scheduler_worker', 517 Field('worker_name', unique=True), 518 Field('first_heartbeat', 'datetime'), 519 Field('last_heartbeat', 'datetime'), 520 Field('status', requires=IS_IN_SET(WORKER_STATUS)), 521 Field('is_ticker', 'boolean', default=False, writable=False), 522 Field('group_names', 'list:string', default=self.group_names), 523 migrate=migrate) 524 if migrate: 525 db.commit()
526
527 - def loop(self, worker_name=None):
528 signal.signal(signal.SIGTERM, lambda signum, stack_frame: sys.exit(1)) 529 try: 530 self.start_heartbeats() 531 while True and self.have_heartbeat: 532 if self.worker_status[0] == DISABLED: 533 logger.debug('Someone stopped me, sleeping until better times come (%s)', self.worker_status[1]) 534 self.sleep() 535 continue 536 logger.debug('looping...') 537 task = self.wrapped_pop_task() 538 if task: 539 self.empty_runs = 0 540 self.worker_status[0] = RUNNING 541 self.report_task(task, self.async(task)) 542 self.worker_status[0] = ACTIVE 543 else: 544 self.empty_runs += 1 545 logger.debug('sleeping...') 546 if self.max_empty_runs != 0: 547 logger.debug('empty runs %s/%s', 548 self.empty_runs, self.max_empty_runs) 549 if self.empty_runs >= self.max_empty_runs: 550 logger.info( 551 'empty runs limit reached, killing myself') 552 self.die() 553 self.sleep() 554 except (KeyboardInterrupt, SystemExit): 555 logger.info('catched') 556 self.die()
557
558 - def wrapped_assign_tasks(self, db):
559 db.commit() # ?don't know if it's useful, let's be completely sure 560 x = 0 561 while x < 10: 562 try: 563 self.assign_tasks(db) 564 db.commit() 565 break 566 except: 567 db.rollback() 568 logger.error('TICKER(%s): error assigning tasks', self.worker_name) 569 x += 1 570 time.sleep(0.5)
571
572 - def wrapped_pop_task(self):
573 db = self.db 574 x = 0 575 while x < 10: 576 try: 577 rtn = self.pop_task(db) 578 return rtn 579 break 580 except: 581 db.rollback() 582 logger.error('TICKER(%s): error popping tasks', self.worker_name) 583 x += 1 584 time.sleep(0.5)
585
586 - def pop_task(self, db):
587 now = self.now() 588 st = self.db.scheduler_task 589 if self.is_a_ticker and self.do_assign_tasks: 590 #I'm a ticker, and 5 loops passed without reassigning tasks, let's do 591 #that and loop again 592 self.wrapped_assign_tasks(db) 593 return None 594 #ready to process something 595 grabbed = db(st.assigned_worker_name == self.worker_name)( 596 st.status == ASSIGNED) 597 598 task = grabbed.select(limitby=(0, 1), orderby=st.next_run_time).first() 599 if task: 600 task.update_record(status=RUNNING, last_run_time=now) 601 #noone will touch my task! 602 db.commit() 603 logger.debug(' work to do %s', task.id) 604 else: 605 if self.greedy and self.is_a_ticker: 606 #there are other tasks ready to be assigned 607 logger.info('TICKER (%s): greedy loop', self.worker_name) 608 self.wrapped_assign_tasks(db) 609 else: 610 logger.info('nothing to do') 611 return None 612 next_run_time = task.last_run_time + datetime.timedelta( 613 seconds=task.period) 614 times_run = task.times_run + 1 615 if times_run < task.repeats or task.repeats == 0: 616 #need to run (repeating task) 617 run_again = True 618 else: 619 #no need to run again 620 run_again = False 621 run_id = 0 622 while True and not self.discard_results: 623 logger.debug(' new scheduler_run record') 624 try: 625 run_id = db.scheduler_run.insert( 626 scheduler_task=task.id, 627 status=RUNNING, 628 start_time=now, 629 worker_name=self.worker_name) 630 db.commit() 631 break 632 except: 633 time.sleep(0.5) 634 db.rollback() 635 logger.info('new task %(id)s "%(task_name)s" %(application_name)s.%(function_name)s' % task) 636 return Task( 637 app=task.application_name, 638 function=task.function_name, 639 timeout=task.timeout, 640 args=task.args, # in json 641 vars=task.vars, # in json 642 task_id=task.id, 643 run_id=run_id, 644 run_again=run_again, 645 next_run_time=next_run_time, 646 times_run=times_run, 647 stop_time=task.stop_time, 648 retry_failed=task.retry_failed, 649 times_failed=task.times_failed, 650 sync_output=task.sync_output)
651
652 - def report_task(self, task, task_report):
653 db = self.db 654 now = self.now() 655 while True: 656 try: 657 if not self.discard_results: 658 if task_report.result != 'null' or task_report.tb: 659 #result is 'null' as a string if task completed 660 #if it's stopped it's None as NoneType, so we record 661 #the STOPPED "run" anyway 662 logger.debug(' recording task report in db (%s)', 663 task_report.status) 664 db(db.scheduler_run.id == task.run_id).update( 665 status=task_report.status, 666 stop_time=now, 667 run_result=task_report.result, 668 run_output=task_report.output, 669 traceback=task_report.tb) 670 else: 671 logger.debug(' deleting task report in db because of no result') 672 db(db.scheduler_run.id == task.run_id).delete() 673 #if there is a stop_time and the following run would exceed it 674 is_expired = (task.stop_time 675 and task.next_run_time > task.stop_time 676 and True or False) 677 status = (task.run_again and is_expired and EXPIRED 678 or task.run_again and not is_expired 679 and QUEUED or COMPLETED) 680 if task_report.status == COMPLETED: 681 d = dict(status=status, 682 next_run_time=task.next_run_time, 683 times_run=task.times_run, 684 times_failed=0 685 ) 686 db(db.scheduler_task.id == task.task_id)( 687 db.scheduler_task.status == RUNNING).update(**d) 688 else: 689 st_mapping = {'FAILED': 'FAILED', 690 'TIMEOUT': 'TIMEOUT', 691 'STOPPED': 'QUEUED'}[task_report.status] 692 status = (task.retry_failed 693 and task.times_failed < task.retry_failed 694 and QUEUED or task.retry_failed == -1 695 and QUEUED or st_mapping) 696 db( 697 (db.scheduler_task.id == task.task_id) & 698 (db.scheduler_task.status == RUNNING) 699 ).update( 700 times_failed=db.scheduler_task.times_failed + 1, 701 next_run_time=task.next_run_time, 702 status=status 703 ) 704 db.commit() 705 logger.info('task completed (%s)', task_report.status) 706 break 707 except: 708 db.rollback() 709 time.sleep(0.5)
710
711 - def adj_hibernation(self):
712 if self.worker_status[0] == DISABLED: 713 wk_st = self.worker_status[1] 714 hibernation = wk_st + 1 if wk_st < MAXHIBERNATION else MAXHIBERNATION 715 self.worker_status[1] = hibernation
716
717 - def send_heartbeat(self, counter):
718 if not self.db_thread: 719 logger.debug('thread building own DAL object') 720 self.db_thread = DAL( 721 self.db._uri, folder=self.db._adapter.folder) 722 self.define_tables(self.db_thread, migrate=False) 723 try: 724 db = self.db_thread 725 sw, st = db.scheduler_worker, db.scheduler_task 726 now = self.now() 727 # record heartbeat 728 mybackedstatus = db( 729 sw.worker_name == self.worker_name).select().first() 730 if not mybackedstatus: 731 sw.insert(status=ACTIVE, worker_name=self.worker_name, 732 first_heartbeat=now, last_heartbeat=now, 733 group_names=self.group_names) 734 self.worker_status = [ACTIVE, 1] # activating the process 735 else: 736 if mybackedstatus.status == DISABLED: 737 # keep sleeping 738 self.worker_status[0] = DISABLED 739 if self.worker_status[1] == MAXHIBERNATION: 740 logger.debug('........recording heartbeat') 741 db(sw.worker_name == self.worker_name).update( 742 last_heartbeat=now) 743 elif mybackedstatus.status == TERMINATE: 744 self.worker_status[0] = TERMINATE 745 logger.debug("Waiting to terminate the current task") 746 self.give_up() 747 return 748 elif mybackedstatus.status == KILL: 749 self.worker_status[0] = KILL 750 self.die() 751 else: 752 logger.debug('........recording heartbeat (%s)', self.worker_status[0]) 753 db(sw.worker_name == self.worker_name).update( 754 last_heartbeat=now, status=ACTIVE) 755 self.worker_status[1] = 1 # re-activating the process 756 if self.worker_status[0] <> RUNNING: 757 self.worker_status[0] = ACTIVE 758 759 self.do_assign_tasks = False 760 761 if counter % 5 == 0: 762 try: 763 # delete inactive workers 764 expiration = now - datetime.timedelta(seconds=self.heartbeat * 3) 765 departure = now - datetime.timedelta( 766 seconds=self.heartbeat * 3 * MAXHIBERNATION) 767 logger.debug( 768 ' freeing workers that have not sent heartbeat') 769 inactive_workers = db( 770 ((sw.last_heartbeat < expiration) & (sw.status == ACTIVE)) | 771 ((sw.last_heartbeat < 772 departure) & (sw.status != ACTIVE)) 773 ) 774 db(st.assigned_worker_name.belongs( 775 inactive_workers._select(sw.worker_name)))(st.status == RUNNING)\ 776 .update(assigned_worker_name='', status=QUEUED) 777 inactive_workers.delete() 778 self.is_a_ticker = self.being_a_ticker() 779 if self.worker_status[0] == ACTIVE: 780 self.do_assign_tasks = True 781 except: 782 pass 783 db.commit() 784 except: 785 db.rollback() 786 self.adj_hibernation() 787 self.sleep()
788
789 - def being_a_ticker(self):
790 db = self.db_thread 791 sw = db.scheduler_worker 792 all_active = db( 793 (sw.worker_name != self.worker_name) & (sw.status == ACTIVE) 794 ).select() 795 ticker = all_active.find(lambda row: row.is_ticker is True).first() 796 not_busy = self.worker_status[0] == ACTIVE 797 if not ticker: 798 if not_busy: 799 #only if this worker isn't busy, otherwise wait for a free one 800 db(sw.worker_name == self.worker_name).update(is_ticker=True) 801 db(sw.worker_name != self.worker_name).update(is_ticker=False) 802 logger.info("TICKER(%s): I'm a ticker", self.worker_name) 803 else: 804 #giving up, only if I'm not alone 805 if len(all_active) > 1: 806 db(sw.worker_name == self.worker_name).update(is_ticker=False) 807 else: 808 not_busy = True 809 db.commit() 810 return not_busy 811 else: 812 logger.info( 813 "%s is a ticker, I'm a poor worker" % ticker.worker_name) 814 return False
815
816 - def assign_tasks(self, db):
817 sw, st = db.scheduler_worker, db.scheduler_task 818 now = self.now() 819 all_workers = db(sw.status == ACTIVE).select() 820 #build workers as dict of groups 821 wkgroups = {} 822 for w in all_workers: 823 group_names = w.group_names 824 for gname in group_names: 825 if gname not in wkgroups: 826 wkgroups[gname] = dict( 827 workers=[{'name': w.worker_name, 'c': 0}]) 828 else: 829 wkgroups[gname]['workers'].append( 830 {'name': w.worker_name, 'c': 0}) 831 #set queued tasks that expired between "runs" (i.e., you turned off 832 #the scheduler): then it wasn't expired, but now it is 833 db(st.status.belongs( 834 (QUEUED, ASSIGNED)))(st.stop_time < now).update(status=EXPIRED) 835 836 all_available = db( 837 (st.status.belongs((QUEUED, ASSIGNED))) & 838 ((st.times_run < st.repeats) | (st.repeats == 0)) & 839 (st.start_time <= now) & 840 ((st.stop_time == None) | (st.stop_time > now)) & 841 (st.next_run_time <= now) & 842 (st.enabled == True) 843 ) 844 limit = len(all_workers) * (50 / (len(wkgroups) or 1)) 845 #if there are a moltitude of tasks, let's figure out a maximum of tasks per worker. 846 #this can be adjusted with some added intelligence (like esteeming how many tasks will 847 #a worker complete before the ticker reassign them around, but the gain is quite small 848 #50 is quite a sweet spot also for fast tasks, with sane heartbeat values 849 #NB: ticker reassign tasks every 5 cycles, so if a worker completes his 50 tasks in less 850 #than heartbeat*5 seconds, it won't pick new tasks until heartbeat*5 seconds pass. 851 852 #If a worker is currently elaborating a long task, all other tasks assigned 853 #to him needs to be reassigned "freely" to other workers, that may be free. 854 #this shuffles up things a bit, in order to maintain the idea of a semi-linear scalability 855 856 #let's freeze it up 857 db.commit() 858 x = 0 859 for group in wkgroups.keys(): 860 tasks = all_available(st.group_name == group).select( 861 limitby=(0, limit), orderby = st.next_run_time) 862 #let's break up the queue evenly among workers 863 for task in tasks: 864 x += 1 865 gname = task.group_name 866 ws = wkgroups.get(gname) 867 if ws: 868 counter = 0 869 myw = 0 870 for i, w in enumerate(ws['workers']): 871 if w['c'] < counter: 872 myw = i 873 counter = w['c'] 874 d = dict( 875 status=ASSIGNED, 876 assigned_worker_name=wkgroups[gname]['workers'][myw]['name'] 877 ) 878 if not task.task_name: 879 d['task_name'] = task.function_name 880 task.update_record(**d) 881 wkgroups[gname]['workers'][myw]['c'] += 1 882 883 db.commit() 884 #I didn't report tasks but I'm working nonetheless!!!! 885 if x > 0: 886 self.empty_runs = 0 887 #I'll be greedy only if tasks assigned are equal to the limit 888 # (meaning there could be others ready to be assigned) 889 self.greedy = x >= limit and True or False 890 logger.info('TICKER(%s): workers are %s', self.worker_name, len(all_workers)) 891 logger.info('TICKER(%s): tasks are %s', self.worker_name, x)
892
893 - def sleep(self):
894 time.sleep(self.heartbeat * self.worker_status[1])
895 # should only sleep until next available task 896
897 - def queue_task(self, function, pargs=[], pvars={}, **kwargs):
898 """ 899 Queue tasks. This takes care of handling the validation of all 900 values. 901 :param function: the function (anything callable with a __name__) 902 :param pargs: "raw" args to be passed to the function. Automatically 903 jsonified. 904 :param pvars: "raw" kwargs to be passed to the function. Automatically 905 jsonified 906 :param kwargs: all the scheduler_task columns. args and vars here should be 907 in json format already, they will override pargs and pvars 908 909 returns a dict just as a normal validate_and_insert, plus a uuid key holding 910 the uuid of the queued task. If validation is not passed, both id and uuid 911 will be None, and you'll get an "error" dict holding the errors found. 912 """ 913 if hasattr(function, '__name__'): 914 function = function.__name__ 915 targs = 'args' in kwargs and kwargs.pop('args') or dumps(pargs) 916 tvars = 'vars' in kwargs and kwargs.pop('vars') or dumps(pvars) 917 tuuid = 'uuid' in kwargs and kwargs.pop('uuid') or web2py_uuid() 918 tname = 'task_name' in kwargs and kwargs.pop('task_name') or function 919 rtn = self.db.scheduler_task.validate_and_insert( 920 function_name=function, 921 task_name=tname, 922 args=targs, 923 vars=tvars, 924 uuid=tuuid, 925 **kwargs) 926 if not rtn.errors: 927 rtn.uuid = tuuid 928 else: 929 rtn.uuid = None 930 return rtn
931
932 - def task_status(self, ref, output=False):
933 """ 934 Shortcut for task status retrieval 935 936 :param ref: can be 937 - integer --> lookup will be done by scheduler_task.id 938 - string --> lookup will be done by scheduler_task.uuid 939 - query --> lookup as you wish (as in db.scheduler_task.task_name == 'test1') 940 :param output: fetch also the scheduler_run record 941 942 Returns a single Row object, for the last queued task 943 If output == True, returns also the last scheduler_run record 944 scheduler_run record is fetched by a left join, so it can 945 have all fields == None 946 947 """ 948 from gluon.dal import Query 949 sr, st = self.db.scheduler_run, self.db.scheduler_task 950 if isinstance(ref, int): 951 q = st.id == ref 952 elif isinstance(ref, str): 953 q = st.uuid == ref 954 elif isinstance(ref, Query): 955 q = ref 956 else: 957 raise SyntaxError( 958 "You can retrieve results only by id, uuid or Query") 959 fields = st.ALL 960 left = False 961 orderby = ~st.id 962 if output: 963 fields = st.ALL, sr.ALL 964 left = sr.on(sr.scheduler_task == st.id) 965 orderby = ~st.id | ~sr.id 966 row = self.db(q).select( 967 *fields, 968 **dict(orderby=orderby, 969 left=left, 970 limitby=(0, 1)) 971 ).first() 972 if output: 973 row.result = row.scheduler_run.run_result and \ 974 loads(row.scheduler_run.run_result, 975 object_hook=_decode_dict) or None 976 return row
977 978
979 -def main():
980 """ 981 allows to run worker without python web2py.py .... by simply python this.py 982 """ 983 parser = optparse.OptionParser() 984 parser.add_option( 985 "-w", "--worker_name", dest="worker_name", default=None, 986 help="start a worker with name") 987 parser.add_option( 988 "-b", "--heartbeat", dest="heartbeat", default=10, 989 type='int', help="heartbeat time in seconds (default 10)") 990 parser.add_option( 991 "-L", "--logger_level", dest="logger_level", 992 default=30, 993 type='int', 994 help="set debug output level (0-100, 0 means all, 100 means none;default is 30)") 995 parser.add_option("-E", "--empty-runs", 996 dest="max_empty_runs", 997 type='int', 998 default=0, 999 help="max loops with no grabbed tasks permitted (0 for never check)") 1000 parser.add_option( 1001 "-g", "--group_names", dest="group_names", 1002 default='main', 1003 help="comma separated list of groups to be picked by the worker") 1004 parser.add_option( 1005 "-f", "--db_folder", dest="db_folder", 1006 default='/Users/mdipierro/web2py/applications/scheduler/databases', 1007 help="location of the dal database folder") 1008 parser.add_option( 1009 "-u", "--db_uri", dest="db_uri", 1010 default='sqlite://storage.sqlite', 1011 help="database URI string (web2py DAL syntax)") 1012 parser.add_option( 1013 "-t", "--tasks", dest="tasks", default=None, 1014 help="file containing task files, must define" + 1015 "tasks = {'task_name':(lambda: 'output')} or similar set of tasks") 1016 parser.add_option( 1017 "-U", "--utc-time", dest="utc_time", default=False, 1018 help="work with UTC timestamps" 1019 ) 1020 (options, args) = parser.parse_args() 1021 if not options.tasks or not options.db_uri: 1022 print USAGE 1023 if options.tasks: 1024 path, filename = os.path.split(options.tasks) 1025 if filename.endswith('.py'): 1026 filename = filename[:-3] 1027 sys.path.append(path) 1028 print 'importing tasks...' 1029 tasks = __import__(filename, globals(), locals(), [], -1).tasks 1030 print 'tasks found: ' + ', '.join(tasks.keys()) 1031 else: 1032 tasks = {} 1033 group_names = [x.strip() for x in options.group_names.split(',')] 1034 1035 logging.getLogger().setLevel(options.logger_level) 1036 1037 print 'groups for this worker: ' + ', '.join(group_names) 1038 print 'connecting to database in folder: ' + options.db_folder or './' 1039 print 'using URI: ' + options.db_uri 1040 db = DAL(options.db_uri, folder=options.db_folder) 1041 print 'instantiating scheduler...' 1042 scheduler = Scheduler(db=db, 1043 worker_name=options.worker_name, 1044 tasks=tasks, 1045 migrate=True, 1046 group_names=group_names, 1047 heartbeat=options.heartbeat, 1048 max_empty_runs=options.max_empty_runs, 1049 utc_time=options.utc_time) 1050 signal.signal(signal.SIGTERM, lambda signum, stack_frame: sys.exit(1)) 1051 print 'starting main worker loop...' 1052 scheduler.loop()
1053 1054 if __name__ == '__main__': 1055 main() 1056