diff --git a/lib/galaxy/app.py b/lib/galaxy/app.py index f1674282ca4..65e26c77357 100644 --- a/lib/galaxy/app.py +++ b/lib/galaxy/app.py @@ -56,6 +56,8 @@ class UniverseApplication(object, config.ConfiguresGalaxyMixin): self.startup_timer = ExecutionTimer() self.new_installation = False self.application_stack = application_stack_instance(app=self) + # A lot of postfork initialization depends on the server name, ensure it is set immediately after forking before other postfork functions + self.application_stack.register_postfork_function(self.application_stack.set_postfork_server_name, self) self.application_stack.register_postfork_function(self.application_stack.start) # Read config file and check for errors self.config = config.Configuration(**kwargs) @@ -189,10 +191,10 @@ class UniverseApplication(object, config.ConfiguresGalaxyMixin): # Start the job manager from galaxy.jobs import manager self.job_manager = manager.JobManager(self) - self.job_manager.start() + self.application_stack.register_postfork_function(self.job_manager.start) # FIXME: These are exposed directly for backward compatibility - self.job_queue = self.job_manager.job_queue - self.job_stop_queue = self.job_manager.job_stop_queue + #self.job_queue = self.job_manager.job_queue + #self.job_stop_queue = self.job_manager.job_stop_queue self.proxy_manager = ProxyManager(self.config) # Initialize the external service types self.external_service_types = external_service_types.ExternalServiceTypesCollection( diff --git a/lib/galaxy/config.py b/lib/galaxy/config.py index 0f6b56d754a..e7de8590d56 100644 --- a/lib/galaxy/config.py +++ b/lib/galaxy/config.py @@ -28,7 +28,7 @@ from galaxy.util import listify from galaxy.util import string_as_bool from galaxy.util.dbkeys import GenomeBuilds from galaxy.web.formatting import expand_pretty_datetime_format -from galaxy.web.stack import register_postfork_function +from galaxy.web.stack import application_stack_log_filter, register_postfork_function from .version import VERSION_MAJOR log = logging.getLogger(__name__) @@ -907,7 +907,11 @@ def configure_logging(config): formatter = logging.Formatter(format) # Hook everything up handler.setFormatter(formatter) + handler.addFilter(application_stack_log_filter()()) root.addHandler(handler) + else: + for h in root.handlers: + h.addFilter(application_stack_log_filter()()) # If sentry is configured, also log to it if getattr(config, "sentry_dsn", None): from raven.handlers.logging import SentryHandler diff --git a/lib/galaxy/jobs/manager.py b/lib/galaxy/jobs/manager.py index d16d2005973..43a50d27ad4 100644 --- a/lib/galaxy/jobs/manager.py +++ b/lib/galaxy/jobs/manager.py @@ -23,12 +23,22 @@ class JobManager(object): def __init__(self, app): self.app = app + if app.application_stack.setup_jobs_with_msg: + # defer setup to postfork + log.debug('######### registering manager init function') + self.app.application_stack.register_postfork_function(self.init) + else: + self.init() + + def init(self): + log.debug("Initializing job manager interface") if self.app.is_job_handler(): log.debug("Starting job handler") - self.job_handler = handler.JobHandler(app) + self.job_handler = handler.JobHandler(self.app) self.job_stop_queue = self.job_handler.job_stop_queue - elif app.application_stack.setup_jobs_with_msg: - self.job_handler = MessageJobHandler( app ) + elif self.app.application_stack.setup_jobs_with_msg: + # not a handler, but notification is via the application stack + self.job_handler = MessageJobHandler( self.app ) self.job_stop_queue = NoopQueue() else: self.job_handler = NoopHandler() diff --git a/lib/galaxy/tools/actions/__init__.py b/lib/galaxy/tools/actions/__init__.py index 13a5cc561ae..c6762780b53 100644 --- a/lib/galaxy/tools/actions/__init__.py +++ b/lib/galaxy/tools/actions/__init__.py @@ -510,7 +510,7 @@ class DefaultToolAction(object): trans.response.send_redirect(url_for(controller='tool_runner', action='redirect', redirect_url=redirect_url)) else: # Put the job in the queue if tracking in memory - app.job_queue.put(job.id, job.tool_id) + app.job_manager.job_queue.put(job.id, job.tool_id) trans.log_event("Added job to the job queue, id: %s" % str(job.id), tool_id=job.tool_id) return job, out_data diff --git a/lib/galaxy/tools/actions/history_imp_exp.py b/lib/galaxy/tools/actions/history_imp_exp.py index a67b1b509ee..222f0b14e28 100644 --- a/lib/galaxy/tools/actions/history_imp_exp.py +++ b/lib/galaxy/tools/actions/history_imp_exp.py @@ -59,7 +59,7 @@ class ImportHistoryToolAction(ToolAction): trans.sa_session.flush() # Queue the job for execution - trans.app.job_queue.put(job.id, tool.id) + trans.app.job_manager.job_queue.put(job.id, tool.id) trans.log_event("Added import history job to the job queue, id: %s" % str(job.id), tool_id=job.tool_id) return job, odict() @@ -138,7 +138,7 @@ class ExportHistoryToolAction(ToolAction): trans.sa_session.flush() # Queue the job for execution - trans.app.job_queue.put(job.id, tool.id) + trans.app.job_manager.job_queue.put(job.id, tool.id) trans.log_event("Added export history job to the job queue, id: %s" % str(job.id), tool_id=job.tool_id) return job, odict() diff --git a/lib/galaxy/tools/actions/metadata.py b/lib/galaxy/tools/actions/metadata.py index a86dee634da..6ccd9f874ef 100644 --- a/lib/galaxy/tools/actions/metadata.py +++ b/lib/galaxy/tools/actions/metadata.py @@ -101,7 +101,7 @@ class SetMetadataToolAction(ToolAction): sa_session.flush() # Queue the job for execution - app.job_queue.put(job.id, tool.id) + app.job_manager.job_queue.put(job.id, tool.id) # FIXME: need to add event logging to app and log events there rather than trans. # trans.log_event( "Added set external metadata job to the job queue, id: %s" % str(job.id), tool_id=job.tool_id ) diff --git a/lib/galaxy/tools/actions/model_operations.py b/lib/galaxy/tools/actions/model_operations.py index d87c44e33a6..4edcd6ce141 100644 --- a/lib/galaxy/tools/actions/model_operations.py +++ b/lib/galaxy/tools/actions/model_operations.py @@ -61,7 +61,7 @@ class ModelOperationToolAction(DefaultToolAction): trans.sa_session.flush() # ensure job.id are available # Queue the job for execution - # trans.app.job_queue.put( job.id, tool.id ) + # trans.app.job_manager.job_queue.put( job.id, tool.id ) # trans.log_event( "Added database job action to the job queue, id: %s" % str(job.id), tool_id=job.tool_id ) log.info("Calling produce_outputs, tool is %s" % tool) return job, out_data diff --git a/lib/galaxy/tools/actions/upload_common.py b/lib/galaxy/tools/actions/upload_common.py index 16f1a8982f6..2b88516caf5 100644 --- a/lib/galaxy/tools/actions/upload_common.py +++ b/lib/galaxy/tools/actions/upload_common.py @@ -418,7 +418,7 @@ def create_job(trans, params, tool, json_file_path, data_list, folder=None, hist trans.sa_session.flush() # Queue the job for execution - trans.app.job_queue.put(job.id, job.tool_id) + trans.app.job_manager.job_queue.put(job.id, job.tool_id) trans.log_event("Added job to the job queue, id: %s" % str(job.id), tool_id=job.tool_id) output = odict() for i, v in enumerate(data_list): diff --git a/lib/galaxy/util/handlers.py b/lib/galaxy/util/handlers.py index b583c2e6e6f..5656ea17d2d 100644 --- a/lib/galaxy/util/handlers.py +++ b/lib/galaxy/util/handlers.py @@ -8,8 +8,6 @@ import logging import os import random -from galaxy.web.stack import process_in_pool - log = logging.getLogger(__name__) @@ -103,10 +101,10 @@ class ConfiguresHandlers: :return: bool """ - stack_handler = process_in_pool(self.app.config.job_handler_pool_name) - if stack_handler is not None: + stack_handles_jobs = self.app.application_stack.process_in_pool(self.app.config.job_handler_pool_name) + if stack_handles_jobs is not None: # Handlers started as uWSGI mules do not require configuration in the job conf - return stack_handler + return stack_handles_jobs for collection in self.handlers.values(): if server_name in collection: return True diff --git a/lib/galaxy/web/stack/__init__.py b/lib/galaxy/web/stack/__init__.py index e920bcd8509..9dc8f18deda 100644 --- a/lib/galaxy/web/stack/__init__.py +++ b/lib/galaxy/web/stack/__init__.py @@ -70,6 +70,20 @@ class ApplicationStackMessageDispatcher(object): self.__funcs[msg.target](msg) +class ApplicationStackLogFilter(logging.Filter): + def filter(self, record): + record.worker_id = None + record.mule_id = None + return True + + +class UWSGILogFilter(logging.Filter): + def filter(self, record): + record.worker_id = uwsgi.worker_id() + record.mule_id = uwsgi.mule_id() + return True + + class ApplicationStack(object): name = None prohibited_middleware = frozenset() @@ -78,6 +92,7 @@ class ApplicationStack(object): handle_jobs = False # used by galaxy.jobs to determine whether this process handles jobs transport_class = ApplicationStackTransport + log_filter_class = ApplicationStackLogFilter @classmethod def register_postfork_function(cls, f, *args, **kwargs): @@ -130,7 +145,6 @@ class MessageApplicationStack(ApplicationStack): self.transport.start() def register_message_handler(self, func, name=None): - log.debug('######## %s register_message_handler called', uwsgi.mule_id()) self.dispatcher.register_func(func, name) self.transport.start_if_needed() @@ -164,6 +178,7 @@ class UWSGIApplicationStack(MessageApplicationStack): setup_jobs_with_msg = True transport_class = UWSGIFarmMessageTransport + log_filter_class = UWSGILogFilter # FIXME: this is copied into UWSGIFarmMessageTransport shutdown_msg = '__SHUTDOWN__' postfork_functions = [] @@ -172,27 +187,40 @@ class UWSGIApplicationStack(MessageApplicationStack): @classmethod def register_postfork_function(cls, f, *args, **kwargs): - cls.postfork_functions.append((f, args, kwargs)) + if uwsgi.mule_id() == 0: + cls.postfork_functions.append((f, args, kwargs)) + else: + # mules are forked from the master and run the master's postfork functions immediately before the forked + # process is replaced. that is prevented in the _do_uwsgi_postfork function, and because mules are + # standalone non-forking processes, they should run postfork functions immediately + f(*args, **kwargs) def __init__(self, app=None): super(UWSGIApplicationStack, self).__init__(app=app) + self._farms_dict = None + self._mules_list = None + self._is_mule = None def __register_signal_handlers(self): for name in ('TERM', 'INT', 'HUP'): + log.debug('######## registered signal handler for SIG%s', name) sig = getattr(signal, 'SIG%s' % name) signal.signal(sig, self._handle_signal) def _handle_signal(self, signum, frame): - if signum in (signal.SIGTERM, signal.SIGINT): - log.info('######## Mule %s received SIGINT/SIGTERM, shutting down', uwsgi.mule_id()) + if signum == signal.SIGTERM: + log.info('######## Mule %s received SIGTERM, shutting down gracefully', uwsgi.mule_id()) self.shutdown() - # This terminates the application loop in the handler entrypoint - self.app.exit = True + elif signum == signal.SIGINT: + log.info('######## Mule %s received SIGINT, shutting down immediately', uwsgi.mule_id()) + self.shutdown() + ## This terminates the application loop in the handler entrypoint + #self.app.exit = True elif signum == signal.SIGHUP: log.debug('######## Mule %s received SIGHUP, restarting', uwsgi.mule_id()) - self.shutdown() - # uWSGI master will restart us - self.app.exit = True + #self.shutdown() + ## uWSGI master will restart us + #self.app.exit = True # FIXME: these are copied into UWSGIFarmMessageTransport @property @@ -211,15 +239,20 @@ class UWSGIApplicationStack(MessageApplicationStack): if self._mules_list is None: self._mules_list = [] mules = uwsgi.opt.get('mule', []) - self._mules_list = [mules] if isinstance(mules, string_types) else mules + self._mules_list = [mules] if isinstance(mules, string_types) or mules is True else mules return self._mules_list def start(self): - self.__register_signal_handlers() + self._is_mule = uwsgi.mule_id() > 0 + if self._is_mule: + self.__register_signal_handlers() super(UWSGIApplicationStack, self).start() def set_postfork_server_name(self, app): - app.config.server_name += ".%d" % uwsgi.worker_id() + if uwsgi.mule_id() == 0: + app.config.server_name += ".worker%d" % uwsgi.worker_id() + else: + app.config.server_name += ".mule%d" % uwsgi.mule_id() # FIXME: used? def workers(self): @@ -229,11 +262,14 @@ class UWSGIApplicationStack(MessageApplicationStack): return self.transport._farm_name == pool_name def shutdown(self): + log.debug('######## STACK SHUTDOWN CALLED') super(UWSGIApplicationStack, self).shutdown() - for farm in self._farms: - for mule in self._mules: - # This will possibly generate more than we need, but that's ok - self.transport.send_message(self.shutdown_msg, farm) + # FIXME: blech + #if not self._is_mule: + # for farm in self._farms: + # for mule in self._mules: + # # This will possibly generate more than we need, but that's ok + # self.transport.send_message(self.shutdown_msg, farm) class PasteApplicationStack(ApplicationStack): @@ -265,6 +301,10 @@ def application_stack_instance(app=None): return stack_class(app=app) +def application_stack_log_filter(): + return application_stack_class().log_filter_class + + def register_postfork_function(f, *args, **kwargs): application_stack_class().register_postfork_function(f, *args, **kwargs) @@ -274,6 +314,11 @@ def process_in_pool(pool_name): @uwsgi_postfork -def _do_postfork(): +def _do_uwsgi_postfork(): + import os + log.debug('######## postfork called, pid %s mule %s functions are: %s' % (os.getpid(), uwsgi.mule_id(), UWSGIApplicationStack.postfork_functions)) + #if uwsgi.mule_id() > 0: + # # mules will inherit the postfork function list and call them immediately upon fork, but should not do that + # UWSGIApplicationStack.postfork_functions = [] for f, args, kwargs in [t for t in UWSGIApplicationStack.postfork_functions]: f(*args, **kwargs) diff --git a/lib/galaxy/web/stack/message.py b/lib/galaxy/web/stack/message.py index d9932011ac2..ba5a8516f59 100644 --- a/lib/galaxy/web/stack/message.py +++ b/lib/galaxy/web/stack/message.py @@ -98,6 +98,6 @@ class JobHandlerMessage(ApplicationStackMessage): def decode(msg_str): - d = json.loads(s) + d = json.loads(msg_str) cls = d.pop('__classname__') - return globals()[cls](d) + return globals()[cls](**d) diff --git a/lib/galaxy/web/stack/transport.py b/lib/galaxy/web/stack/transport.py index 429d0e4138f..b93edfdb8af 100644 --- a/lib/galaxy/web/stack/transport.py +++ b/lib/galaxy/web/stack/transport.py @@ -4,6 +4,10 @@ from __future__ import absolute_import import logging import threading +try: + from queue import Empty, Queue +except ImportError: + from Queue import Empty, Queue try: import uwsgi @@ -17,9 +21,11 @@ log = logging.getLogger(__name__) class ApplicationStackTransport(object): + shutdown_msg = '__SHUTDOWN__' + def __init_dispatcher_thread(self): self.dispatcher_thread = threading.Thread(name=self.__class__.__name__ + ".dispatcher_thread", target=self._dispatch_messages) - self.dispatcher_thread.daemon = True + #self.dispatcher_thread.daemon = True def __init__(self, app, dispatcher=None): """ Pre-fork initialization. @@ -36,19 +42,10 @@ class ApplicationStackTransport(object): def start_if_needed(self): # Don't unnecessarily start a thread that we don't need. - log.debug('######## start_if_needed called') - log.debug('######## %s' % self.can_run) - log.debug('######## %s' % self.running) - log.debug('######## %s' % self.dispatcher_thread.is_alive()) - log.debug('######## %s' % self.dispatcher) - log.debug('######## %s' % self.dispatcher.handler_count) - import traceback - traceback.print_stack() - # FIXME: can_run is False here in the mule, but start() was called way back when the app was loaded and it was True then, what's going on here? if self.can_run and not self.running and not self.dispatcher_thread.is_alive() and self.dispatcher and self.dispatcher.handler_count: self.running = True self.dispatcher_thread.start() - log.debug('######## Web stack IPC message dispatcher thread started') + log.debug('######## Web stack IPC message dispatcher thread started in mule %s', uwsgi.mule_id()) def stop_if_unneeded(self): if self.can_run and self.running and self.dispatcher_thread.is_alive() and self.dispatcher and not self.dispatcher.handler_count: @@ -67,12 +64,20 @@ class ApplicationStackTransport(object): pass def shutdown(self): + log.debug('######## TRANSPORT SHUTDOWN CALLED') self.running = False - self.dispatcher_thread.join() + if self.dispatcher_thread.is_alive(): + # FIXME + self.send_message(self.shutdown_msg, 'job-handlers') + self.dispatcher_thread.join() + log.debug('######## Joined dispatcher thread') class UWSGIFarmMessageTransport(ApplicationStackTransport): - shutdown_msg = '__SHUTDOWN__' + """ Communication via uWSGI Mule Farm messages. Communication is unidirectional (workers -> mules). + """ + # FIXME: do you need a clear "this is a producer, this is a consumer flag? + #shutdown_msg = '__SHUTDOWN__' # Define any static lock names here, additional locks will be appended for each configured farm's message handler _locks = [] @@ -94,6 +99,8 @@ class UWSGIFarmMessageTransport(ApplicationStackTransport): super(UWSGIFarmMessageTransport, self).__init__(app, dispatcher=dispatcher) self._farms_dict = None self._mules_list = None + self._is_mule = False + self._msg_queue = Queue() self.__initialize_locks() @property @@ -112,7 +119,7 @@ class UWSGIFarmMessageTransport(ApplicationStackTransport): if self._mules_list is None: self._mules_list = [] mules = uwsgi.opt.get('mule', []) - self._mules_list = [mules] if isinstance(mules, string_types) else mules + self._mules_list = [mules] if isinstance(mules, string_types) or mules is True else mules return self._mules_list def __lock(self, name_or_id): @@ -139,31 +146,34 @@ class UWSGIFarmMessageTransport(ApplicationStackTransport): try: log.debug('######## Mule %s acquired message receive lock, waiting for new message', uwsgi.mule_id()) msg = uwsgi.farm_get_msg() - if msg == self.shutdown_msg: - # all you need to do is pass here, self.running should already be set False by the signal handler calling the shutdown method defined in the superclass - log.debug('######## SHUTTING DOWN %s', uwsgi.mule_id()) log.debug('######## Mule %s received message: %s', uwsgi.mule_id(), msg) + if msg == self.shutdown_msg or msg is None: + # all you need to do is pass here, self.running should already be set False by the signal handler calling the shutdown method defined in the superclass + self.running = False + log.debug('######## SHUTTING DOWN %s', uwsgi.mule_id()) except: log.exception( "Exception in mule message handling" ) finally: self.__unlock(lock) - if msg: + log.debug('######## Mule %s released lock', uwsgi.mule_id()) + if msg != self.shutdown_msg and msg is not None: self.dispatcher.dispatch(msg) log.info('######## Mule %s message handler shutting down', uwsgi.mule_id()) def start(self): """ Post-fork initialization. """ - if uwsgi.mule_id() == 0: - # this is the main process - return - if not uwsgi.in_farm(): - raise RuntimeError('Mule %s is not in a farm! Set `farm = %s:%s` in uWSGI configuration' - % (uwsgi.mule_id(), - self.app.config.job_handler_pool_name, - ','.join(map(str, range(1, len(filter(lambda x: x.endswith('galaxy/main.py'), self._mules)) + 1))))) + # TODO: what happens if workers > 1?? + self._is_mule = uwsgi.mule_id() > 0 + if self._is_mule: + if not uwsgi.in_farm(): + raise RuntimeError('Mule %s is not in a farm! Set `farm = %s:%s` in uWSGI configuration' + % (uwsgi.mule_id(), + self.app.config.job_handler_pool_name, + ','.join(map(str, range(1, len(filter(lambda x: x.endswith('galaxy/main.py'), self._mules)) + 1))))) super(UWSGIFarmMessageTransport, self).start() - log.info('######## Mule started, mule id: %s, farm name: %s, server name: %s', uwsgi.mule_id(), self._farm_name, self.app.config.server_name) + log.info('######## Mule transport started, worker id: %s, mule id: %s, farm name: %s, server name: %s', uwsgi.worker_id(), uwsgi.mule_id(), self._farm_name, self.app.config.server_name) + self._send_all_messages() @property def _farm_name(self): @@ -172,7 +182,20 @@ class UWSGIFarmMessageTransport(ApplicationStackTransport): return name return None + def _send_all_messages(self): + # the sender doesn't have a running thread, all we are concerned with here is whether or not we've forked yet + if self.can_run and not self._is_mule: + while True: + try: + msg, dest = self._msg_queue.get_nowait() + except Empty: + break + log.debug('######## Sending message in mule %s to farm %s: %s', uwsgi.mule_id(), dest, msg) + uwsgi.farm_msg(dest, msg) + log.debug('######## Message sent') + + def send_message(self, msg, dest): - log.debug('######## Sending message to farm %s: %s', dest, msg) - uwsgi.farm_msg(dest, msg) - log.debug('######## Message sent') + log.debug('######## Queing message in mule %s to farm %s: %s', uwsgi.mule_id(), dest, msg) + self._msg_queue.put((msg, dest)) + self._send_all_messages() diff --git a/lib/galaxy/webapps/galaxy/buildapp.py b/lib/galaxy/webapps/galaxy/buildapp.py index d009bb6a9f9..9faa25c93ac 100644 --- a/lib/galaxy/webapps/galaxy/buildapp.py +++ b/lib/galaxy/webapps/galaxy/buildapp.py @@ -181,7 +181,8 @@ def uwsgi_app_factory(): def postfork_setup(): from galaxy.app import app - app.application_stack.set_postfork_server_name(app) + # FIXME + #app.application_stack.set_postfork_server_name(app) app.control_worker.bind_and_start()