From c604d4da4d85cf7fef8200de63bdd6f81560e047 Mon Sep 17 00:00:00 2001 From: "Rafael.Timmerberg" Date: Wed, 11 Dec 2013 00:53:58 +0100 Subject: [PATCH] changed from event to Queue.Queue, still not working... --- bin/linspector | 1 - linspector/tasks/task.py | 24 +++++++++++++++++++----- 2 files changed, 19 insertions(+), 6 deletions(-) diff --git a/bin/linspector b/bin/linspector index 028be95..7cccc89 100755 --- a/bin/linspector +++ b/bin/linspector @@ -115,7 +115,6 @@ def main(): exit() task = TaskExecutor.Instance() - logger.debug("blah") #logger.debug("task executor instance :%s", str(task)) job_count = 0 diff --git a/linspector/tasks/task.py b/linspector/tasks/task.py index e147138..66dbc51 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 Queue import Queue from linspector.utils.singleton import Singleton KEY_TYPE = "type" @@ -79,7 +80,7 @@ class TaskExecutor(object): def __init__(self): logger.debug("init taskExecutor") self.event = Event() - self.taskInfos = [] + self.taskInfos = Queue() task_thread = Thread(target=self._run_worker_thread) self._instantEnd = False self._running = True @@ -87,7 +88,9 @@ class TaskExecutor(object): task_thread.start() 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() @@ -104,6 +107,14 @@ class TaskExecutor(object): except Exception, e: logger.error("Error " + str(e)) + ''' + try: + logger.debug("try to get queued task") + msg, task = self.taskInfos.get() + logger.debug("got task %s for msg: %s", str(task), str(msg)) + task.execute(msg) + except Exception, e: + logger.error("Error " + str(e)) logger.debug("shutting down TaskExecutor!") def is_instant_end(self): @@ -122,7 +133,10 @@ class TaskExecutor(object): self.event.set() def schedule_task(self, msg, task): - logger.debug("appending task '%s' for msg: %s", str(task), str(msg)) - self.taskInfos.append((msg, task)) - logger.debug("%s task stored", str(len(self.taskInfos))) - self.event.set() + try: + logger.debug("appending task '%s' for msg: %s", str(task), str(msg)) + self.taskInfos.put((msg, task)) + except Exception, e: + logger.debug("queue is probably full: %s", str(e)) + #logger.debug("%s task Queued", str(len(self.taskInfos))) + #self.event.set()