reverted to 7f9397eb30; needs some cleanups but then we are back online again.... ;)

This commit is contained in:
Johannes Findeisen 2013-12-11 05:16:11 +01:00
commit 34e73a56f2
5 changed files with 52 additions and 62 deletions

View file

@ -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

View file

@ -20,7 +20,7 @@ along with this program. If not, see <http://www.gnu.org/licenses/>.
"""
__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()

View file

@ -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 }
]
}
},

View file

@ -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)

View file

@ -21,8 +21,6 @@ along with this program. If not, see <http://www.gnu.org/licenses/>.
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))
msg, task = self.taskInfos[0]
del self.taskInfos[0]
if task:
logger.debug("Starting Task Execution...")
task.execute(msg)
self.taskInfos.task_done()
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()