added TaskExecutor into the job to execute tasks in its own Thread, and stops blocking the job
This commit is contained in:
parent
8f9d9bc551
commit
1b089807bb
3 changed files with 7 additions and 10 deletions
|
|
@ -35,6 +35,7 @@ from linspector.core.interface import LinspectorInterface
|
||||||
from linspector.core.job import Job
|
from linspector.core.job import Job
|
||||||
from linspector.core.scheduler import Scheduler
|
from linspector.core.scheduler import Scheduler
|
||||||
from linspector.frontends.lish import LishFrontend
|
from linspector.frontends.lish import LishFrontend
|
||||||
|
from linspector.tasks.task import TaskExecutor
|
||||||
|
|
||||||
logger = logging.getLogger(__name__)
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
|
@ -120,6 +121,7 @@ def main():
|
||||||
scheduler = Scheduler({"apscheduler.threadpool.core_threads": args.corethreads,
|
scheduler = Scheduler({"apscheduler.threadpool.core_threads": args.corethreads,
|
||||||
"apscheduler.threadpool.max_threads": args.threads})
|
"apscheduler.threadpool.max_threads": args.threads})
|
||||||
scheduler.start()
|
scheduler.start()
|
||||||
|
TaskExecutor.Instance()
|
||||||
|
|
||||||
start_date = datetime.datetime.now()
|
start_date = datetime.datetime.now()
|
||||||
time_delta = 0
|
time_delta = 0
|
||||||
|
|
@ -156,6 +158,7 @@ def main():
|
||||||
|
|
||||||
logger.debug("shutting down scheduler")
|
logger.debug("shutting down scheduler")
|
||||||
|
|
||||||
|
|
||||||
shutdown_wait = True
|
shutdown_wait = True
|
||||||
if "shutdown_wait" in core:
|
if "shutdown_wait" in core:
|
||||||
shutdown_wait = core["shutdown_wait"]
|
shutdown_wait = core["shutdown_wait"]
|
||||||
|
|
|
||||||
|
|
@ -23,6 +23,7 @@ along with this program. If not, see <http://www.gnu.org/licenses/>.
|
||||||
from datetime import datetime
|
from datetime import datetime
|
||||||
from binascii import crc32
|
from binascii import crc32
|
||||||
from logging import getLogger
|
from logging import getLogger
|
||||||
|
from linspector.tasks.task import TaskExecutor
|
||||||
|
|
||||||
logger = getLogger(__name__)
|
logger = getLogger(__name__)
|
||||||
|
|
||||||
|
|
@ -112,9 +113,7 @@ class Job:
|
||||||
def handle_alarm(self):
|
def handle_alarm(self):
|
||||||
for member in self.members:
|
for member in self.members:
|
||||||
for task in member.get_tasks():
|
for task in member.get_tasks():
|
||||||
task.execute(str(self.get_message()))
|
TaskExecutor.Instance().schedule_task(self.get_message(), task)
|
||||||
#print task
|
|
||||||
logger.info("DO TASK EXECUTION HERE! NOT IMPLEMENTED!")
|
|
||||||
|
|
||||||
def handle_call(self):
|
def handle_call(self):
|
||||||
logger.debug("handle call")
|
logger.debug("handle call")
|
||||||
|
|
|
||||||
|
|
@ -21,6 +21,7 @@ along with this program. If not, see <http://www.gnu.org/licenses/>.
|
||||||
|
|
||||||
from logging import getLogger
|
from logging import getLogger
|
||||||
from threading import Event, Thread
|
from threading import Event, Thread
|
||||||
|
from linspector.utils.singleton import Singleton
|
||||||
|
|
||||||
KEY_TYPE = "type"
|
KEY_TYPE = "type"
|
||||||
KEY_ARGS = "args"
|
KEY_ARGS = "args"
|
||||||
|
|
@ -66,14 +67,8 @@ class Task(object):
|
||||||
raise e
|
raise e
|
||||||
|
|
||||||
|
|
||||||
|
@Singleton
|
||||||
class TaskExecutor(object):
|
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):
|
def __init__(self):
|
||||||
self.event = Event()
|
self.event = Event()
|
||||||
self.taskInfos = []
|
self.taskInfos = []
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue