diff --git a/CHANGELOG b/CHANGELOG
index 710e666..109afee 100644
--- a/CHANGELOG
+++ b/CHANGELOG
@@ -1,7 +1,6 @@
0.16 - Date ??? Next release is scheduled to christmas 2013
- Added a CHANGELOG file
- - Added multi core support by instantiating multiple scheduler instances
- Added "gui" command to lish providing an urwid (curses) based interface to
Linspector
- Added scheduling status to boot message
diff --git a/bin/linspector b/bin/linspector
index 9455ff1..d915f7a 100755
--- a/bin/linspector
+++ b/bin/linspector
@@ -20,7 +20,7 @@ along with this program. If not, see .
"""
-__version__ = "0.15.3/AMNESIA"
+__version__ = "0.15/AMNESIA"
import argparse
import datetime
@@ -29,10 +29,13 @@ import logging.handlers
import os
import os.path as path
import sys
+import multiprocessing
+from multiprocessing import Pool
from linspector.config.parser import FullConfigParser
from linspector.core.interface import LinspectorInterface
-from linspector.core.worker import LinspectorWorker
+from linspector.core.job import LinspectorJob
+from linspector.core.scheduler import LinspectorScheduler
from linspector.frontends.lish import LishFrontend
from linspector.tasks.task import TaskExecutor
@@ -60,14 +63,11 @@ def parse_args():
parser.add_argument("-c", "--logcount", default=5, type=int,
help="maximum number of logfiles in rotation (default: 5)")
- parser.add_argument("-i", "--instances", default=1, type=int,
- help="number of scheduler instances (default: 1)")
-
parser.add_argument("-m", "--logsize", default=10485760, type=int,
help="maximum logfile size in bytes (default: 10485760)")
- parser.add_argument("-t", "--threads", default=1000, type=int,
- help="maximum number of scheduler threads (default: 1000)")
+ parser.add_argument("-t", "--threads", default=3500, type=int,
+ help="maximum number of scheduler threads (default: 3500)")
parser.add_argument("-k", "--corethreads", default=0, type=int,
help="number of scheduler core threads (default: 0)")
@@ -91,6 +91,14 @@ def parse_args():
return parser.parse_args()
+pool = Pool(processes=16)
+print 'cpu_count() = %d\n' % multiprocessing.cpu_count()
+
+
+def handle_job(job):
+ result = pool.apply_async(job.handle_call())
+
+
def main():
global lin_conf
args = parse_args()
@@ -114,26 +122,21 @@ def main():
logger.error(msg)
exit()
- executor = TaskExecutor()
- #logger.debug("task executor instance :%s", str(task))
+ scheduler = LinspectorScheduler({"apscheduler.threadpool.core_threads": args.corethreads,
+ "apscheduler.threadpool.max_threads": args.threads})
+ scheduler.start()
+ TaskExecutor.Instance()
job_count = 0
for layout in lin_conf.get_enabled_layouts():
for hostgroup in layout.get_hostgroups():
job_count += (hostgroup.get_services().__len__() * hostgroup.get_hosts().__len__())
- workers = []
- for i in range(0, args.instances):
- process = LinspectorWorker(args.corethreads, args.threads)
- workers.append(process)
- process.daemon = True
-
start_date = datetime.datetime.now()
time_delta = 0
jobs = []
count = 0
percent = 0
- instance = 0
for layout in lin_conf.get_enabled_layouts():
for hostgroup in layout.get_hostgroups():
for service in hostgroup.get_services():
@@ -151,22 +154,19 @@ def main():
time_delta += float(args.delay)
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, executor=executor)
-
- instance += 1
- if instance == args.instances:
- instance = 0
+ job = LinspectorJob(service,
+ host,
+ hostgroup.get_members(),
+ core,
+ hostgroup)
+ scheduler_job = period.createJob(scheduler, job, handle_job, start_date=new_start_date)
+ if scheduler_job is not None:
+ job.set_job(scheduler_job)
+ jobs.append(job)
print("\nScheduled " + str(job_count) + " jobs")
- schedulers = []
- for worker in workers:
- schedulers.append(worker.get_scheduler())
- worker.start()
-
- interface = LinspectorInterface(jobs, schedulers, lin_conf, root_logger, __version__)
+ interface = LinspectorInterface(jobs, scheduler, lin_conf, root_logger, __version__)
if "jsonrpc_backend" in core and core["jsonrpc_backend"]:
from linspector.backends.jsonrpc import JsonrpcBackend
@@ -184,15 +184,13 @@ def main():
if "shutdown_wait" in core:
shutdown_wait = core["shutdown_wait"]
- for worker in workers:
- worker.shutdown(shutdown_wait)
- worker.terminate()
-
+ scheduler.shutdown(wait=shutdown_wait)
logging.shutdown()
if shutdown_wait:
TaskExecutor.Instance().stop()
else:
TaskExecutor.Instance().stop_immediately()
+
if __name__ == "__main__":
main()
\ No newline at end of file
diff --git a/examples/minimal.json b/examples/minimal.json
index fc4dbb8..78defdc 100644
--- a/examples/minimal.json
+++ b/examples/minimal.json
@@ -21,8 +21,7 @@
"members": ["root"],
"hosts": [ "a"],
"services":[
- { "class": "etc/dummy", "args": { "sleep": 1, "fail": true }, "periods": ["fast"], "threshold": 1 },
- { "class": "etc/dummy", "args": { "sleep": 1, "fail": false }, "periods": ["fast"], "threshold": 1 }
+ { "class": "etc/dummy", "args": { "sleep": 1, "fail": true }, "periods": ["fast"], "threshold": 1 }
]
}
},
@@ -47,4 +46,4 @@
}
}
}
-}
+}
\ No newline at end of file
diff --git a/linspector/core/job.py b/linspector/core/job.py
index 41fe9be..c0168b2 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, executor):
+ def __init__(self, service, host, members, core, hostgroup):
self.service = service
self.host = host
self.members = members
@@ -49,7 +49,6 @@ 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
@@ -116,11 +115,8 @@ class LinspectorJob:
for task in member.get_tasks():
if self.status.lower() in task.get_task_type().lower():
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):
logger.debug("handle call")
logger.debug(self.service)
diff --git a/linspector/tasks/task.py b/linspector/tasks/task.py
index 281487f..d5d1e60 100644
--- a/linspector/tasks/task.py
+++ b/linspector/tasks/task.py
@@ -21,8 +21,6 @@ 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
KEY_TYPE = "type"
@@ -76,11 +74,11 @@ class Task(object):
raise e
+@Singleton
class TaskExecutor(object):
def __init__(self):
- logger.debug("init taskExecutor")
self.event = Event()
- self.taskInfos = Queue()
+ self.taskInfos = []
task_thread = Thread(target=self._run_worker_thread)
self._instantEnd = False
self._running = True
@@ -88,17 +86,20 @@ 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:
+ self.event.clear()
+ self.event.wait()
+
try:
- 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()
+ msg, task = self.taskInfos[0]
+ del self.taskInfos[0]
+ if task:
+ logger.debug("Starting Task Execution...")
+ task.execute(msg)
+
except Exception, e:
- logger.error("Error: %s", e.message)
- logger.debug("shutting down TaskExecutor!")
+ logger.error("Error " + str(e))
def is_instant_end(self):
return self._instantEnd
@@ -108,16 +109,13 @@ 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))
-
+ self.taskInfos.append((msg, task))
+ self.event.set()