From 09c9f27f69ae72cac34b1286dfefce125ac404a6 Mon Sep 17 00:00:00 2001 From: John Chilton Date: Tue, 31 Oct 2017 09:42:54 -0400 Subject: [PATCH 1/2] More structured app shutdown. - Cleanup database engine when shutting down Galaxy app (this wasn't done previously and I think it can cause connection issues when testing across different databases). - Join monitor threads when shutting down (old behavior can be restored by setting wait time for monitor threads back to "0"). See note in galaxy.ini.sample. - Better error handling during shutdown, don't let a failure to cleanup one thing cause other things not to be cleaned up. - New monitor thread mixin based on jobs stuff - reuse for workflows so it properly wakes up workflow scheduling on shutdown. --- config/galaxy.ini.sample | 5 +++ lib/galaxy/app.py | 55 +++++++++++++++++++---- lib/galaxy/config.py | 1 + lib/galaxy/jobs/handler.py | 44 ++++++++---------- lib/galaxy/jobs/runners/__init__.py | 22 +++++++-- lib/galaxy/util/monitors.py | 44 ++++++++++++++++++ lib/galaxy/workflow/scheduling_manager.py | 32 ++++++------- 7 files changed, 151 insertions(+), 52 deletions(-) create mode 100644 lib/galaxy/util/monitors.py diff --git a/config/galaxy.ini.sample b/config/galaxy.ini.sample index da67d2e107a..759db03e98a 100644 --- a/config/galaxy.ini.sample +++ b/config/galaxy.ini.sample @@ -853,6 +853,11 @@ nglims_config_file = tool-data/nglims.yaml # public site. Enabled in the sample config for development. use_interactive = True +# When stopping Galaxy cleanly, how much time to give various monitoring/polling +# threads to finish before giving up on joining them. Set to 0 to disable this and +# restore the pre-18.01 default behavior. +#monitor_thread_join_timeout = 5 + # Write thread status periodically to 'heartbeat.log', (careful, uses disk # space rapidly!). Useful to determine why your processes may be consuming a # lot of CPU. diff --git a/lib/galaxy/app.py b/lib/galaxy/app.py index 3f51b4aff30..7e92bb829c8 100644 --- a/lib/galaxy/app.py +++ b/lib/galaxy/app.py @@ -212,19 +212,58 @@ class UniverseApplication(object, config.ConfiguresGalaxyMixin): log.info("Galaxy app startup finished %s" % self.startup_timer) def shutdown(self): - self.watchers.shutdown() - self.workflow_scheduling_manager.shutdown() - self.job_manager.shutdown() - self.object_store.shutdown() - if self.heartbeat: - self.heartbeat.shutdown() - self.update_repository_manager.shutdown() + exception = None try: - self.control_worker.shutdown() + self.watchers.shutdown() + except Exception as e: + exception = exception or e + log.exception("Failed to shutdown configuration watchers cleanly") + try: + self.workflow_scheduling_manager.shutdown() + except Exception as e: + exception = exception or e + log.exception("Failed to shutdown workflow scheduling manager cleanly") + try: + self.job_manager.shutdown() + except Exception as e: + exception = exception or e + log.exception("Failed to shutdown job manager cleanly") + try: + self.object_store.shutdown() + except Exception as e: + exception = exception or e + log.exception("Failed to shutdown object store cleanly") + try: + if self.heartbeat: + self.heartbeat.shutdown() + except Exception as e: + exception = exception or e + log.exception("Failed to shutdown heartbeat cleanly") + try: + self.update_repository_manager.shutdown() + except Exception as e: + exception = exception or e + log.exception("Failed to shutdown update repository manager cleanly") + + try: + try: + self.control_worker.shutdown() + except Exception as e: + exception = exception or e + log.exception("Failed to shutdown control worker cleanly") except AttributeError: # There is no control_worker pass + try: + self.model.engine.dispose() + except Exception as e: + exception = exception or e + log.exception("Failed to shutdown SA database engine cleanly") + + if exception: + raise exception + def configure_fluent_log(self): if self.config.fluent_log: from galaxy.util.logging.fluent_log import FluentTraceLogger diff --git a/lib/galaxy/config.py b/lib/galaxy/config.py index 9c24edb8445..8379b8f2ba1 100644 --- a/lib/galaxy/config.py +++ b/lib/galaxy/config.py @@ -371,6 +371,7 @@ class Configuration(object): self.use_heartbeat = string_as_bool(kwargs.get('use_heartbeat', 'False')) self.heartbeat_interval = int(kwargs.get('heartbeat_interval', 20)) self.heartbeat_log = kwargs.get('heartbeat_log', None) + self.monitor_thread_join_timeout = int(kwargs.get("monitor_thread_join_timeout", 5)) self.log_actions = string_as_bool(kwargs.get('log_actions', 'False')) self.log_events = string_as_bool(kwargs.get('log_events', 'False')) self.sanitize_all_html = string_as_bool(kwargs.get('sanitize_all_html', True)) diff --git a/lib/galaxy/jobs/handler.py b/lib/galaxy/jobs/handler.py index 2acb8b15f85..e478b2d79f2 100644 --- a/lib/galaxy/jobs/handler.py +++ b/lib/galaxy/jobs/handler.py @@ -4,7 +4,6 @@ Galaxy job handler, prepares, runs, tracks, and finishes Galaxy jobs import datetime import logging import os -import threading import time from Queue import ( Empty, @@ -27,7 +26,7 @@ from galaxy.jobs import ( TaskWrapper ) from galaxy.jobs.mapper import JobNotReadyException -from galaxy.util.sleeper import Sleeper +from galaxy.util.monitors import Monitors log = logging.getLogger(__name__) @@ -57,7 +56,7 @@ class JobHandler(object): self.job_stop_queue.shutdown() -class JobHandlerQueue(object): +class JobHandlerQueue(Monitors, object): """ Job Handler's Internal Queue, this is what actually implements waiting for jobs to be runnable and dispatching to a JobRunner. @@ -84,11 +83,8 @@ class JobHandlerQueue(object): self.waiting_jobs = [] # Contains wrappers of jobs that are limited or ready (so they aren't created unnecessarily/multiple times) self.job_wrappers = {} - # Helper for interruptable sleep - self.sleeper = Sleeper() - self.running = True - self.monitor_thread = threading.Thread(name="JobHandlerQueue.monitor_thread", target=self.__monitor) - self.monitor_thread.setDaemon(True) + name = "JobHandlerQueue.monitor_thread" + self._init_monitor_thread(name, target=self.__monitor, config=app.config) def start(self): """ @@ -203,7 +199,7 @@ class JobHandlerQueue(object): Continually iterate the waiting jobs, checking is each is ready to run and dispatching if so. """ - while self.running: + while self.monitor_running: try: # If jobs are locked, there's nothing to monitor and we skip # to the sleep. @@ -211,8 +207,7 @@ class JobHandlerQueue(object): self.__monitor_step() except Exception: log.exception("Exception in monitor_step") - # Sleep - self.sleeper.sleep(1) + self._monitor_sleep(1) def __monitor_step(self): """ @@ -670,15 +665,15 @@ class JobHandlerQueue(object): return else: log.info("sending stop signal to worker thread") - self.running = False + self.stop_monitoring() if not self.app.config.track_jobs_in_database: self.queue.put(self.STOP_SIGNAL) - self.sleeper.wake() + self.shutdown_monitor() log.info("job handler queue stopped") self.dispatcher.shutdown() -class JobHandlerStopQueue(object): +class JobHandlerStopQueue(Monitors): """ A queue for jobs which need to be terminated prematurely. """ @@ -699,12 +694,8 @@ class JobHandlerStopQueue(object): # Contains jobs that are waiting (only use from monitor thread) self.waiting = [] - # Helper for interruptable sleep - self.sleeper = Sleeper() - self.running = True - self.monitor_thread = threading.Thread(name="JobHandlerStopQueue.monitor_thread", target=self.monitor) - self.monitor_thread.setDaemon(True) - self.monitor_thread.start() + name = "JobHandlerStopQueue.monitor_thread" + self._init_monitor_thread(name, start=True, config=app.config) log.info("job handler stop queue started") def monitor(self): @@ -713,13 +704,13 @@ class JobHandlerStopQueue(object): """ # HACK: Delay until after forking, we need a way to do post fork notification!!! time.sleep(10) - while self.running: + while self.monitor_running: try: self.monitor_step() except Exception: log.exception("Exception in monitor_step") # Sleep - self.sleeper.sleep(1) + self._monitor_sleep(1) def monitor_step(self): """ @@ -778,10 +769,10 @@ class JobHandlerStopQueue(object): return else: log.info("sending stop signal to worker thread") - self.running = False + self.stop_monitoring() if not self.app.config.track_jobs_in_database: self.queue.put(self.STOP_SIGNAL) - self.sleeper.wake() + self.shutdown_monitor() log.info("job handler stop queue stopped") @@ -873,4 +864,7 @@ class DefaultJobDispatcher(object): def shutdown(self): for runner in self.job_runners.itervalues(): - runner.shutdown() + try: + runner.shutdown() + except Exception: + raise Exception("Failed to shutdown runner %s" % runner) diff --git a/lib/galaxy/jobs/runners/__init__.py b/lib/galaxy/jobs/runners/__init__.py index 675d7ba484b..b29e6e57f7d 100644 --- a/lib/galaxy/jobs/runners/__init__.py +++ b/lib/galaxy/jobs/runners/__init__.py @@ -29,6 +29,7 @@ from galaxy.util import ( shrink_stream_by_size ) from galaxy.util.bunch import Bunch +from galaxy.util.monitors import Monitors from .state_handler_factory import build_state_handlers @@ -135,6 +136,19 @@ class BaseJobRunner(object): for i in range(len(self.work_threads)): self.work_queue.put((STOP_SIGNAL, None)) + join_timeout = self.app.config.monitor_thread_join_timeout + if join_timeout > 0: + exception = None + for thread in self.work_threads: + try: + thread.join(join_timeout) + except Exception as e: + exception = e + log.exception("Faild to shutdown worker thread") + + if exception: + raise exception + # Most runners should override the legacy URL handler methods and destination param method def url_to_destination(self, url): """ @@ -506,7 +520,7 @@ class AsynchronousJobState(JobState): self.cleanup_file_attributes.append(attribute) -class AsynchronousJobRunner(BaseJobRunner): +class AsynchronousJobRunner(Monitors, BaseJobRunner): """Parent class for any job runner that runs jobs asynchronously (e.g. via a distributed resource manager). Provides general methods for having a thread to monitor the state of asynchronous jobs and submitting those jobs @@ -524,9 +538,8 @@ class AsynchronousJobRunner(BaseJobRunner): self.monitor_queue = Queue() def _init_monitor_thread(self): - self.monitor_thread = threading.Thread(name="%s.monitor_thread" % self.runner_name, target=self.monitor) - self.monitor_thread.setDaemon(True) - self.monitor_thread.start() + name = "%s.monitor_thread" % self.runner_name + super(AsynchronousJobRunner, self)._init_monitor_thread(name=name, target=self.monitor, start=True, config=self.app.config) def handle_stop(self): # DRMAA and SGE runners should override this and disconnect. @@ -565,6 +578,7 @@ class AsynchronousJobRunner(BaseJobRunner): log.info("%s: Sending stop signal to monitor thread" % self.runner_name) self.monitor_queue.put(STOP_SIGNAL) # Call the parent's shutdown method to stop workers + self.shutdown_monitor() super(AsynchronousJobRunner, self).shutdown() def check_watched_items(self): diff --git a/lib/galaxy/util/monitors.py b/lib/galaxy/util/monitors.py new file mode 100644 index 00000000000..63c2e1a43e9 --- /dev/null +++ b/lib/galaxy/util/monitors.py @@ -0,0 +1,44 @@ +from __future__ import absolute_import + +import logging +import threading + +from .sleeper import Sleeper + +log = logging.getLogger(__name__) + +DEFAULT_MONITOR_THREAD_JOIN_TIMEOUT = 5 + + +class Monitors: + + def _init_monitor_thread(self, name, target_name=None, target=None, start=False, config=None): + self.monitor_join_sleep = getattr(config, "monitor_thread_join_timeout", DEFAULT_MONITOR_THREAD_JOIN_TIMEOUT) + self.monitor_join = self.monitor_join_sleep > 0 + self.monitor_sleeper = Sleeper() + self.monitor_running = True + + if target is not None: + assert target_name is None + monitor_func = target + else: + target_name = target_name or "monitor" + monitor_func = getattr(self, target_name) + self.sleeper = Sleeper() + self.monitor_thread = threading.Thread(name=name, target=monitor_func) + self.monitor_thread.setDaemon(True) + if start: + self.monitor_thread.start() + + def stop_monitoring(self): + self.monitor_running = False + + def _monitor_sleep(self, sleep_amount): + self.sleeper.sleep(sleep_amount) + + def shutdown_monitor(self): + self.stop_monitoring() + self.sleeper.wake() + if self.monitor_join: + log.debug("Joining monitor thread") + self.monitor_thread.join(self.monitor_join_sleep) diff --git a/lib/galaxy/workflow/scheduling_manager.py b/lib/galaxy/workflow/scheduling_manager.py index 5b6d1ee52dd..a815f612323 100644 --- a/lib/galaxy/workflow/scheduling_manager.py +++ b/lib/galaxy/workflow/scheduling_manager.py @@ -1,13 +1,12 @@ import logging import os -import threading -import time from xml.etree import ElementTree import galaxy.workflow.schedulers from galaxy import model from galaxy.util import plugin_config from galaxy.util.handlers import ConfiguresHandlers +from galaxy.util.monitors import Monitors log = logging.getLogger(__name__) @@ -75,17 +74,23 @@ class WorkflowSchedulingManager(object, ConfiguresHandlers): return self.__has_handlers.get_handler(None, index=random_index) def shutdown(self): + exception = None for workflow_scheduler in self.workflow_schedulers.values(): try: workflow_scheduler.shutdown() - except Exception: + except Exception as e: + exception = exception or e log.exception(EXCEPTION_MESSAGE_SHUTDOWN) if self.request_monitor: try: self.request_monitor.shutdown() - except Exception: + except Exception as e: + exception = exception or e log.exception("Failed to shutdown workflow request monitor.") + if exception: + raise exception + def queue(self, workflow_invocation, request_params): workflow_invocation.state = model.WorkflowInvocation.states.NEW scheduler = request_params.get("scheduler", None) or self.default_scheduler_id @@ -170,32 +175,29 @@ class WorkflowSchedulingManager(object, ConfiguresHandlers): self.request_monitor = WorkflowRequestMonitor(self.app, self) -class WorkflowRequestMonitor(object): +class WorkflowRequestMonitor(Monitors, object): def __init__(self, app, workflow_scheduling_manager): self.app = app - self.active = True self.workflow_scheduling_manager = workflow_scheduling_manager - self.monitor_thread = threading.Thread(name="WorkflowRequestMonitor.monitor_thread", target=self.__monitor) - self.monitor_thread.setDaemon(True) - self.monitor_thread.start() + self._init_monitor_thread(name="WorkflowRequestMonitor.monitor_thread", target=self.__monitor, start=True, config=app.config) def __monitor(self): to_monitor = self.workflow_scheduling_manager.active_workflow_schedulers - while self.active: + while self.monitor_running: for workflow_scheduler_id, workflow_scheduler in to_monitor.items(): - if not self.active: + if not self.monitor_running: return self.__schedule(workflow_scheduler_id, workflow_scheduler) - # TODO: wake if stopped - time.sleep(1) + + self._monitor_sleep(1) def __schedule(self, workflow_scheduler_id, workflow_scheduler): invocation_ids = self.__active_invocation_ids(workflow_scheduler_id) for invocation_id in invocation_ids: self.__attempt_schedule(invocation_id, workflow_scheduler) - if not self.active: + if not self.monitor_running: return def __attempt_schedule(self, invocation_id, workflow_scheduler): @@ -232,4 +234,4 @@ class WorkflowRequestMonitor(object): ) def shutdown(self): - self.active = False + self.shutdown_monitor() From c3e89834087c6d6b0468d031356ac503525b17b0 Mon Sep 17 00:00:00 2001 From: John Chilton Date: Mon, 30 Oct 2017 23:58:50 -0400 Subject: [PATCH 2/2] Fix integration database encoding issue with templated database stuff. --- lib/galaxy/config.py | 1 + lib/galaxy/model/migrate/check.py | 13 ++++++++++++- 2 files changed, 13 insertions(+), 1 deletion(-) diff --git a/lib/galaxy/config.py b/lib/galaxy/config.py index 191492081b3..3c4b7c25913 100644 --- a/lib/galaxy/config.py +++ b/lib/galaxy/config.py @@ -133,6 +133,7 @@ class Configuration(object): self.database_create_tables = string_as_bool(kwargs.get("database_create_tables", "True")) self.database_query_profiling_proxy = string_as_bool(kwargs.get("database_query_profiling_proxy", "False")) self.database_template = kwargs.get("database_template", None) + self.database_encoding = kwargs.get("database_encoding", None) # Create new databases with this encoding. self.slow_query_log_threshold = float(kwargs.get("slow_query_log_threshold", 0)) # Don't set this to true for production databases, but probably should diff --git a/lib/galaxy/model/migrate/check.py b/lib/galaxy/model/migrate/check.py index c8bbf05c9dd..fb963393a71 100644 --- a/lib/galaxy/model/migrate/check.py +++ b/lib/galaxy/model/migrate/check.py @@ -34,7 +34,18 @@ def create_or_verify_database(url, galaxy_config_file, engine_options={}, app=No new_database = not database_exists(url) if new_database: template = app and getattr(app.config, "database_template", None) - create_database(url, template=template) + encoding = app and getattr(app.config, "database_encoding", None) + create_kwds = {} + + message = "Creating database for URI [%s]" % url + if template: + message += " from template [%s]" % template + create_kwds["template"] = template + if encoding: + message += " with encoding [%s]" % encoding + create_kwds["encoding"] = encoding + log.info(message) + create_database(url, **create_kwds) # Create engine and metadata engine = create_engine(url, **engine_options)