diff --git a/config/galaxy.yml.sample b/config/galaxy.yml.sample index a8bd9c428c8..e9d91554771 100644 --- a/config/galaxy.yml.sample +++ b/config/galaxy.yml.sample @@ -1019,9 +1019,14 @@ galaxy: # 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 + # them. Set to 0 to disable this and terminate without waiting. Among + # others, these threads include the job handler workers, which are + # responsible for preparing/submitting and collecting/finishing jobs, + # and which can cause job errors if not shut down cleanly. If using + # supervisord, consider also increasing the value of `stopwaitsecs`. + # If using job handler mules, consider also setting the `mule-reload- + # mercy` uWSGI option. See the Galaxy Admin Documentation for more. + #monitor_thread_join_timeout: 30 # Write thread status periodically to 'heartbeat.log', (careful, uses # disk space rapidly!). Useful to determine why your processes may be diff --git a/doc/source/admin/galaxy_options.rst b/doc/source/admin/galaxy_options.rst index acc8b0b51e4..e68276436ad 100644 --- a/doc/source/admin/galaxy_options.rst +++ b/doc/source/admin/galaxy_options.rst @@ -2097,9 +2097,15 @@ :Description: 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. -:Default: ``5`` + them. Set to 0 to disable this and terminate without waiting. + Among others, these threads include the job handler workers, which + are responsible for preparing/submitting and collecting/finishing + jobs, and which can cause job errors if not shut down cleanly. If + using supervisord, consider also increasing the value of + `stopwaitsecs`. If using job handler mules, consider also setting + the `mule-reload-mercy` uWSGI option. See the Galaxy Admin + Documentation for more. +:Default: ``30`` :Type: int diff --git a/doc/source/admin/scaling.md b/doc/source/admin/scaling.md index 3a16d5122fd..5af4d2ae955 100644 --- a/doc/source/admin/scaling.md +++ b/doc/source/admin/scaling.md @@ -407,6 +407,22 @@ option, `enable-threads` is set implicitly. This option enables the Python GIL a by Galaxy itself for various non-web tasks), which Galaxy uses extensively. Setting it explicitly, however, is harmless and can prevent strange difficult-to-debug situations if `threads` is accidentally unset. +**Worker/Mule shutdown/reload mercy** + +By default, uWSGI will wait up to 60 seconds for web workers and mules to terminate. This is generally safe for +servicing web requests, but some parts of Galaxy's job preparation/submission and collection/finishing operations can +take quite a bit of time to complete and are not entirely reentrant: job errors or state inconsistencies can occur if +interrupted (although every effort has been made to minimize such possibilities). By default, Galaxy will wait up to 30 +seconds for the threads allocated for these operations to terminate after instructing them to shut down. You can change +this behavior by increasing the value of `monitor_thread_join_timeout` in the `galaxy` section of `galaxy.yml`. The +maximum amount of time that Galaxy will take to shut down job runner workers is `monitor_thread_join_timeout * +runner_plugin_count` since each plugin is shut down sequentially (`runner_plugin_count` is the number of ``s in +your `job_conf.xml`). + +Thus you should set the appropriate uWSGI `*-restart-mercy` option to a value higher than the maximum job runner worker +shutdown time. If using **uWSGI all-in-one**, set `worker-reload-mercy`, and if using **uWSGI + Mule job handling**, set +`mule-reload-mercy` (both in the `uwsgi` section of `galaxy.yml`). + **Signals** The signal handling options (`die-on-term` and `hook-master-start` with `unix_signal` values) are not required but, if @@ -494,6 +510,7 @@ umask = 022 autostart = true autorestart = true startsecs = 10 +stopwaitsecs = 65 user = galaxy numprocs = 1 stopsignal = INT @@ -506,6 +523,9 @@ must stay up for this long before we consider it OK. If the process crashes soon made to your local installation) supervisord will try again a couple of times to restart the process before giving up and marking it as failed. This is one of the many ways supervisord is much friendly for managing these sorts of tasks. +The value of `stopwaitsecs` should be at least as large as the smallest value of uWSGI's `reload-mercy`, +`worker-reload-mercy`, and `mule-reload-mercy` options, all of which default to `60`. + If using the **uWSGI + Webless** scenario, you'll need to addtionally define job handlers to start. There's no simple way to activate a virtualenv when using supervisor, but you can simulate the effects by setting `$PATH` and `$VIRTUAL_ENV`: @@ -520,6 +540,7 @@ umask = 022 autostart = true autorestart = true startsecs = 15 +stopwaitsecs = 35 user = galaxy environment = VIRTUAL_ENV="/srv/galaxy/venv",PATH="/srv/galaxy/venv/bin:%(ENV_PATH)s" ``` @@ -529,6 +550,10 @@ substitution in the `command` and `process_name` fields. We've set `numprocs = 3 processes. Supervisord will loop over `0..numprocs` and launch `handler0`, `handler1`, and `handler2` processes automatically for us, templating out the command string so each handler receives a different log file and name. +The value of `stopwaitsecs` should be at least as large as `monitor_thread_join_timeout * runner_plugin_count`, which +is `30` in the default configuration (`monitor_thread_join_timeout` is a Galaxy configuration option and +`runner_plugin_count` is the number of ``s in your `job_conf.xml`). + Lastly, collect the tasks defined above into a single group. If you are not using webless handlers this is as simple as: ```ini diff --git a/lib/galaxy/config.py b/lib/galaxy/config.py index 263e02f752f..cf23dca1854 100644 --- a/lib/galaxy/config.py +++ b/lib/galaxy/config.py @@ -425,7 +425,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.monitor_thread_join_timeout = int(kwargs.get("monitor_thread_join_timeout", 30)) 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 e4a999b9709..1ef5e5cfc1f 100644 --- a/lib/galaxy/jobs/handler.py +++ b/lib/galaxy/jobs/handler.py @@ -965,8 +965,12 @@ class DefaultJobDispatcher(object): job_wrapper.fail(DEFAULT_JOB_PUT_FAILURE_MESSAGE) def shutdown(self): - for runner in self.job_runners.values(): + failures = [] + for name, runner in self.job_runners.items(): try: runner.shutdown() except Exception: - raise Exception("Failed to shutdown runner %s" % runner) + failures.append(name) + log.exception("Failed to shutdown runner %s", name) + if failures: + raise Exception("Failed to shutdown runners: %s" % ', '.join(failures)) diff --git a/lib/galaxy/jobs/runners/__init__.py b/lib/galaxy/jobs/runners/__init__.py index 0ba8ec84114..3c36b7ae18f 100644 --- a/lib/galaxy/jobs/runners/__init__.py +++ b/lib/galaxy/jobs/runners/__init__.py @@ -6,8 +6,10 @@ import logging import os import string import subprocess +import sys import threading import time +import traceback from six.moves.queue import ( Empty, @@ -85,10 +87,21 @@ class BaseJobRunner(object): log.debug('Starting %s %s workers' % (self.nworkers, self.runner_name)) for i in range(self.nworkers): worker = threading.Thread(name="%s.work_thread-%d" % (self.runner_name, i), target=self.run_next) - worker.setDaemon(True) + worker.daemon = True self.app.application_stack.register_postfork_function(worker.start) self.work_threads.append(worker) + def _alive_worker_threads(self, cycle=False): + # yield endlessly as long as there are alive threads if cycle is True + alive = True + while alive: + alive = False + for thread in self.work_threads: + if thread.is_alive(): + if cycle: + alive = True + yield thread + def run_next(self): """Run the next item in the work queue (a job waiting to run) """ @@ -129,22 +142,36 @@ class BaseJobRunner(object): def shutdown(self): """Attempts to gracefully shut down the worker threads """ - log.info("%s: Sending stop signal to %s worker threads" % (self.runner_name, len(self.work_threads))) + log.info("%s: Sending stop signal to %s job worker threads", self.runner_name, len(self.work_threads)) 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: + log.info("Waiting up to %d seconds for job worker threads to shutdown...", join_timeout) + start = time.time() + # NOTE: threads that have already joined by now are not going to be logged + for thread in self._alive_worker_threads(cycle=True): + if time.time() > (start + join_timeout): + break try: - thread.join(join_timeout) - except Exception as e: - exception = e - log.exception("Faild to shutdown worker thread") + thread.join(2) + except Exception: + log.exception("Caught exception attempting to shutdown job worker thread %s:", thread.name) + if not thread.is_alive(): + log.debug("Job worker thread terminated: %s", thread.name) + else: + log.info("All job worker threads shutdown cleanly") + return - if exception: - raise exception + for thread in self._alive_worker_threads(): + try: + frame = sys._current_frames()[thread.ident] + except KeyError: + # thread is now stopped + continue + log.warning("Timed out waiting for job worker thread %s to terminate, shutdown will be unclean! Thread " + "stack is:\n%s", thread.name, ''.join(traceback.format_stack(frame))) # Most runners should override the legacy URL handler methods and destination param method def url_to_destination(self, url): diff --git a/lib/galaxy/webapps/galaxy/config_schema.yml b/lib/galaxy/webapps/galaxy/config_schema.yml index 9bbdadfa459..10bb176bc33 100644 --- a/lib/galaxy/webapps/galaxy/config_schema.yml +++ b/lib/galaxy/webapps/galaxy/config_schema.yml @@ -1563,12 +1563,17 @@ mapping: monitor_thread_join_timeout: type: int - default: 5 + default: 30 required: false desc: | 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. + terminate without waiting. Among others, these threads include the job handler + workers, which are responsible for preparing/submitting and collecting/finishing + jobs, and which can cause job errors if not shut down cleanly. If using + supervisord, consider also increasing the value of `stopwaitsecs`. If using job + handler mules, consider also setting the `mule-reload-mercy` uWSGI option. See + the Galaxy Admin Documentation for more. use_heartbeat: type: bool diff --git a/test/base/driver_util.py b/test/base/driver_util.py index cae48fdc5c9..e2a26883388 100644 --- a/test/base/driver_util.py +++ b/test/base/driver_util.py @@ -228,6 +228,7 @@ def setup_galaxy_config( user_library_import_dir=user_library_import_dir, webhooks_dir=TEST_WEBHOOKS_DIR, logging=LOGGING_CONFIG_DEFAULT, + monitor_thread_join_timeout=5, ) config.update(database_conf(tmpdir, prefer_template_database=prefer_template_database)) config.update(install_database_conf(tmpdir, default_merged=default_install_db_merged))