mirror of
https://github.com/galaxyproject/galaxy.git
synced 2026-09-24 16:30:27 +08:00
Merge branch 'dev' into fix_status_matching
This commit is contained in:
@@ -858,6 +858,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.
|
||||
|
||||
+47
-8
@@ -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
|
||||
|
||||
@@ -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
|
||||
@@ -372,6 +373,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))
|
||||
|
||||
+19
-25
@@ -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)
|
||||
|
||||
@@ -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):
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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)
|
||||
@@ -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()
|
||||
|
||||
Reference in New Issue
Block a user