|
Package Dropbox ::
Package web2py ::
Package gluon ::
Module scheduler
|
|
1
2
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
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
124 self.vars = vars
125 self.__dict__.update(kwargs)
126
128 return '<Task: %s>' % self.function
129
130
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
144 return '<TaskReport: %s>' % self.status
145
146
148 """ test function """
149 for i in range(argv[0]):
150 print 'click', i
151 time.sleep(1)
152 return 'done'
153
154
155
156
157
158
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
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
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
212
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
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
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
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
398 """
399 validator that check whether field is valid json and validate its type
400 """
401
402 - def __init__(self, myclass=list, parse=False):
405
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
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
437
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
452 return self.utc_time and datetime.datetime.utcnow() or datetime.datetime.now()
453
460
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
571
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
587 now = self.now()
588 st = self.db.scheduler_task
589 if self.is_a_ticker and self.do_assign_tasks:
590
591
592 self.wrapped_assign_tasks(db)
593 return None
594
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
602 db.commit()
603 logger.debug(' work to do %s', task.id)
604 else:
605 if self.greedy and self.is_a_ticker:
606
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
617 run_again = True
618 else:
619
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,
641 vars=task.vars,
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
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
660
661
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
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
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
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
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]
735 else:
736 if mybackedstatus.status == DISABLED:
737
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
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
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
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
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
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
817 sw, st = db.scheduler_worker, db.scheduler_task
818 now = self.now()
819 all_workers = db(sw.status == ACTIVE).select()
820
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
832
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
846
847
848
849
850
851
852
853
854
855
856
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
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
885 if x > 0:
886 self.empty_runs = 0
887
888
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
894 time.sleep(self.heartbeat * self.worker_status[1])
895
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
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
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