From 3745c27c9ae589908d79f9b8816992fe081522d5 Mon Sep 17 00:00:00 2001 From: "Rafael.Timmerberg" Date: Wed, 11 Dec 2013 03:41:56 +0100 Subject: [PATCH] removed singleton and passed instance of TaskExecutor to workers --- bin/linspector | 4 ++-- linspector/core/job.py | 5 +++-- linspector/core/worker.py | 4 ++-- linspector/tasks/task.py | 31 ++++++------------------------- 4 files changed, 13 insertions(+), 31 deletions(-) diff --git a/bin/linspector b/bin/linspector index 7cccc89..9455ff1 100755 --- a/bin/linspector +++ b/bin/linspector @@ -114,7 +114,7 @@ def main(): logger.error(msg) exit() - task = TaskExecutor.Instance() + executor = TaskExecutor() #logger.debug("task executor instance :%s", str(task)) job_count = 0 @@ -153,7 +153,7 @@ def main(): new_start_date = start_date + datetime.timedelta(seconds=time_delta) worker = workers[instance] - worker.create_job(jobs, service, host, hostgroup, core, period, start_date=new_start_date) + worker.create_job(jobs, service, host, hostgroup, core, period, start_date=new_start_date, executor=executor) instance += 1 if instance == args.instances: diff --git a/linspector/core/job.py b/linspector/core/job.py index 259a5a8..41fe9be 100644 --- a/linspector/core/job.py +++ b/linspector/core/job.py @@ -39,7 +39,7 @@ logger = getLogger(__name__) class LinspectorJob: - def __init__(self, service, host, members, core, hostgroup): + def __init__(self, service, host, members, core, hostgroup, executor): self.service = service self.host = host self.members = members @@ -49,6 +49,7 @@ class LinspectorJob: self.enabled = True self.scheduler_job = None self.job_id = self.hex_string() + self.executor = executor """ NONE job was not executed OK when everything is fine @@ -117,7 +118,7 @@ class LinspectorJob: logger.debug("Executing Task of type: " + self.status) try: - TaskExecutor.Instance().schedule_task(job_information, task) + self.executor.schedule_task(job_information, task) except Exception, e: logger.error("Error while executing: " + str(e)) def handle_call(self): diff --git a/linspector/core/worker.py b/linspector/core/worker.py index ce0f1b2..48d6197 100644 --- a/linspector/core/worker.py +++ b/linspector/core/worker.py @@ -53,8 +53,8 @@ class LinspectorWorker(Process): def shutdown(self, wait=True): self.scheduler.shutdown(wait=wait) - def create_job(self, jobs, service, host, hostgroup, core, period, start_date=None): - job = LinspectorJob(service, host, hostgroup.get_members(), core, hostgroup) + def create_job(self, jobs, service, host, hostgroup, core, period, start_date=None, executor=None): + job = LinspectorJob(service, host, hostgroup.get_members(), core, hostgroup, executor) scheduler_job = period.createJob(self.scheduler, job, self.handle_job, start_date=start_date) if scheduler_job is not None: job.set_job(scheduler_job) diff --git a/linspector/tasks/task.py b/linspector/tasks/task.py index 66dbc51..281487f 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 time import sleep from Queue import Queue from linspector.utils.singleton import Singleton @@ -75,7 +76,6 @@ class Task(object): raise e -@Singleton class TaskExecutor(object): def __init__(self): logger.debug("init taskExecutor") @@ -90,31 +90,14 @@ class TaskExecutor(object): def _run_worker_thread(self): logger.debug("start running taskExcecutor worker Thread") while self.is_running() or not self.is_instant_end(): - ''' - if len(self.taskInfos) == 0: - logger.debug("waiting for jobs to add") - self.event.clear() - self.event.wait() - logger.debug("Task Executor waked up") - try: - logger.debug("trying to get task information") - msg, task = self.taskInfos[0] - del self.taskInfos[0] - if task: - logger.debug("Starting Task Execution of task: %s", str(task)) - task.execute(msg) - - except Exception, e: - logger.error("Error " + str(e)) - ''' - try: - logger.debug("try to get queued task") + logger.debug("try to get queued task from %s", self.taskInfos) msg, task = self.taskInfos.get() logger.debug("got task %s for msg: %s", str(task), str(msg)) task.execute(msg) + self.taskInfos.task_done() except Exception, e: - logger.error("Error " + str(e)) + logger.error("Error: %s", e.message) logger.debug("shutting down TaskExecutor!") def is_instant_end(self): @@ -125,18 +108,16 @@ class TaskExecutor(object): def stop(self): self._running = False - self.event.set() def stop_immediately(self): self._running = False self._instantEnd = True - self.event.set() def schedule_task(self, msg, task): try: logger.debug("appending task '%s' for msg: %s", str(task), str(msg)) self.taskInfos.put((msg, task)) + logger.debug("into queue: %s ", str(self.taskInfos)) except Exception, e: logger.debug("queue is probably full: %s", str(e)) - #logger.debug("%s task Queued", str(len(self.taskInfos))) - #self.event.set() +