diff --git a/bin/linspector b/bin/linspector index c1788b4..d476d9b 100755 --- a/bin/linspector +++ b/bin/linspector @@ -35,6 +35,7 @@ from linspector.core.interface import LinspectorInterface from linspector.core.job import Job from linspector.core.scheduler import Scheduler from linspector.frontends.lish import LishFrontend +from linspector.tasks.task import TaskExecutor logger = logging.getLogger(__name__) @@ -120,6 +121,7 @@ def main(): scheduler = Scheduler({"apscheduler.threadpool.core_threads": args.corethreads, "apscheduler.threadpool.max_threads": args.threads}) scheduler.start() + TaskExecutor.Instance() start_date = datetime.datetime.now() time_delta = 0 @@ -156,6 +158,7 @@ def main(): logger.debug("shutting down scheduler") + shutdown_wait = True if "shutdown_wait" in core: shutdown_wait = core["shutdown_wait"] diff --git a/linspector/core/job.py b/linspector/core/job.py index 0ca4eb5..2dded36 100644 --- a/linspector/core/job.py +++ b/linspector/core/job.py @@ -23,6 +23,7 @@ along with this program. If not, see . from datetime import datetime from binascii import crc32 from logging import getLogger +from linspector.tasks.task import TaskExecutor logger = getLogger(__name__) @@ -112,9 +113,7 @@ class Job: def handle_alarm(self): for member in self.members: for task in member.get_tasks(): - task.execute(str(self.get_message())) - #print task - logger.info("DO TASK EXECUTION HERE! NOT IMPLEMENTED!") + TaskExecutor.Instance().schedule_task(self.get_message(), task) def handle_call(self): logger.debug("handle call") diff --git a/linspector/tasks/task.py b/linspector/tasks/task.py index b8989c4..768e1de 100644 --- a/linspector/tasks/task.py +++ b/linspector/tasks/task.py @@ -21,6 +21,7 @@ along with this program. If not, see . from logging import getLogger from threading import Event, Thread +from linspector.utils.singleton import Singleton KEY_TYPE = "type" KEY_ARGS = "args" @@ -66,14 +67,8 @@ class Task(object): raise e - +@Singleton class TaskExecutor(object): - _instance = None - def __new__(cls, *args, **kwargs): - if not cls._instance: - cls._instance = super(TaskExecutor, cls).__init__() - return cls._instance - def __init__(self): self.event = Event() self.taskInfos = []