diff --git a/lib/galaxy/jobs/handler.py b/lib/galaxy/jobs/handler.py index 8306bd2ede1..86cdd868870 100644 --- a/lib/galaxy/jobs/handler.py +++ b/lib/galaxy/jobs/handler.py @@ -51,6 +51,7 @@ class JobHandler(object): def start(self): self.job_queue.start() + self.job_stop_queue.start() def shutdown(self): self.job_queue.shutdown() @@ -91,6 +92,7 @@ class JobHandlerQueue(Monitors, object): """ Starts the JobHandler's thread after checking for any unhandled jobs. """ + log.debug('Handler queue starting for jobs assigned to handler: %s', self.app.config.server_name) # Recover jobs at startup self.__check_jobs_at_startup() # Start the queue @@ -721,7 +723,12 @@ class JobHandlerStopQueue(Monitors): self.waiting = [] name = "JobHandlerStopQueue.monitor_thread" - self._init_monitor_thread(name, start=True, config=app.config) + self._init_monitor_thread(name, config=app.config) + log.info("job handler stop queue started") + + def start(self): + # Start the queue + self.monitor_thread.start() log.info("job handler stop queue started") def monitor(self): diff --git a/lib/galaxy/jobs/manager.py b/lib/galaxy/jobs/manager.py index 6805ca4ba1f..17625cc432b 100644 --- a/lib/galaxy/jobs/manager.py +++ b/lib/galaxy/jobs/manager.py @@ -25,7 +25,7 @@ class JobManager(object): self.app = app self.job_lock = False if self.app.is_job_handler(): - log.debug("Starting job handler") + log.debug("Initializing job handler") self.job_handler = handler.JobHandler(app) self.job_stop_queue = self.job_handler.job_stop_queue elif app.application_stack.has_pool(app.application_stack.pools.JOB_HANDLERS): diff --git a/lib/galaxy/workflow/scheduling_manager.py b/lib/galaxy/workflow/scheduling_manager.py index 882be090489..9aae8b86dcd 100644 --- a/lib/galaxy/workflow/scheduling_manager.py +++ b/lib/galaxy/workflow/scheduling_manager.py @@ -216,6 +216,7 @@ class WorkflowSchedulingManager(object, ConfiguresHandlers): def __start_request_monitor(self): self.request_monitor = WorkflowRequestMonitor(self.app, self) + self.app.application_stack.register_postfork_function(self.request_monitor.start) class WorkflowRequestMonitor(Monitors, object): @@ -223,7 +224,7 @@ class WorkflowRequestMonitor(Monitors, object): def __init__(self, app, workflow_scheduling_manager): self.app = app self.workflow_scheduling_manager = workflow_scheduling_manager - self._init_monitor_thread(name="WorkflowRequestMonitor.monitor_thread", target=self.__monitor, start=True, config=app.config) + self._init_monitor_thread(name="WorkflowRequestMonitor.monitor_thread", target=self.__monitor, config=app.config) def __monitor(self): to_monitor = self.workflow_scheduling_manager.active_workflow_schedulers @@ -276,5 +277,8 @@ class WorkflowRequestMonitor(Monitors, object): handler=handler, ) + def start(self): + self.monitor_thread.start() + def shutdown(self): self.shutdown_monitor()