WIP: Job handlers as UWSGI mules, attempt at unprogrammed mules.

This commit is contained in:
Nate Coraor
2017-08-21 21:45:28 -04:00
parent c7663a4e2e
commit a06d6cf616
13 changed files with 151 additions and 68 deletions
+5 -3
View File
@@ -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(
+5 -1
View File
@@ -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
+13 -3
View File
@@ -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()
+1 -1
View File
@@ -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
+2 -2
View File
@@ -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()
+1 -1
View File
@@ -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 )
+1 -1
View File
@@ -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
+1 -1
View File
@@ -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):
+3 -5
View File
@@ -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
+62 -17
View File
@@ -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)
+2 -2
View File
@@ -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)
+53 -30
View File
@@ -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()
+2 -1
View File
@@ -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()