WIP: attempted to support workers shared between pools, but locking

won't allow for this, since you can't wait on messages from a specific
farm.
This commit is contained in:
Nate Coraor
2017-08-21 21:48:01 -04:00
parent ece3f242fb
commit 27a0b1ba40
9 changed files with 230 additions and 58 deletions
-2
View File
@@ -540,8 +540,6 @@ class Configuration(object):
self.job_manager = kwargs.get('job_manager', self.server_name).strip()
self.job_handlers = [x.strip() for x in kwargs.get('job_handlers', self.server_name).split(',')]
self.default_job_handlers = [x.strip() for x in kwargs.get('default_job_handlers', ','.join(self.job_handlers)).split(',')]
self.job_handler_count = int( kwargs.get( 'job_handler_count', 1 ) )
self.job_handler_pool_name = kwargs.get('job_handler_pool_name', 'job-handlers').strip()
# Store per-tool runner configs
self.tool_handlers = self.__read_tool_job_config(global_conf_parser, 'galaxy:tool_handlers', 'name')
self.tool_runners = self.__read_tool_job_config(global_conf_parser, 'galaxy:tool_runners', 'url')
+1
View File
@@ -170,6 +170,7 @@ class JobConfiguration(object, ConfiguresHandlers):
"""This is a stop-gap solution until the job conf loading is rewritten
for YAML job conf support.
"""
# FIXME: this will not support the case where uWSGI is being used w/o handler mules
if self.app.application_stack.handle_jobs:
self.handlers[self.app.config.server_name] = (self.app.config.server_name,)
self.default_handler_id = self.app.config.server_name
+2 -1
View File
@@ -99,7 +99,8 @@ class MessageJobQueue(object):
def put(self, job_id, tool_id):
msg = JobHandlerMessage(task='setup', job_id=job_id)
self.app.application_stack.send_message(self.app.config.job_handler_pool_name, msg)
# TODO: send to a specific pool
self.app.application_stack.send_message(self.app.application_stack.purposes.JOB_HANDLER, msg)
def put_stop(self, *args):
pass
+39
View File
@@ -0,0 +1,39 @@
"""Generic config parsing into dictionary.
"""
import errno
import logging
try:
import yaml
except ImportError:
yaml = None
log = logging.getLogger(__name__)
def parse_config(config_file, root, default=None):
if default:
conf = default.copy()
else:
conf = {}
try:
load_func = __load_func(config_file)
with open(config_file) as fh:
c = load_func(fh)
if not c:
c = {}
conf.update(c.get(root, {}))
except (OSError, IOError) as exc:
if exc.errno == errno.ENOENT:
log.warning("config file '%s' does not exist, running with default config", config_file)
else:
raise
return conf
def __load_func(path):
if path.endswith('yaml') or path.endswith('.yml'):
if not yaml:
raise RuntimeError("The 'yaml' module could not be imported, please install PyYAML to read '%s'" % path)
return yaml.load
+5 -3
View File
@@ -101,10 +101,12 @@ class ConfiguresHandlers:
:return: bool
"""
stack_handles_jobs = self.app.application_stack.process_in_pool(self.app.config.job_handler_pool_name)
if stack_handles_jobs is not None:
if self.app.application_stack.has_purpose(self.app.application_stack.purposes.JOB_HANDLER):
# Handlers started as uWSGI mules do not require configuration in the job conf
return stack_handles_jobs
return True
#if stack_handles_jobs is not None:
# # Handlers started as uWSGI mules do not require configuration in the job conf
# return stack_handles_jobs
for collection in self.handlers.values():
if server_name in collection:
return True
+158 -35
View File
@@ -6,6 +6,7 @@ import inspect
import logging
import os
import signal
import socket
import threading
# The uwsgi module is automatically injected by the parent uwsgi
@@ -16,22 +17,25 @@ try:
except ImportError:
uwsgi = None
try:
from uwsgidecorators import postfork as uwsgi_postfork
except:
uwsgi_postfork = lambda x: x # noqa: E731
if uwsgi is not None and hasattr(uwsgi, 'numproc'):
print("WARNING: This is a uwsgi process but the uwsgidecorators library"
" is unavailable. This is likely due to using an external (not"
" in Galaxy's virtualenv) uwsgi and you may experience errors. "
"HINT:\n {venv}/bin/pip install uwsgidecorators".format(
venv=os.environ.get('VIRTUAL_ENV', '/path/to/venv')))
#try:
# from uwsgidecorators import postfork as uwsgi_postfork
#except:
# uwsgi_postfork = lambda x: x # noqa: E731
# if uwsgi is not None and hasattr(uwsgi, 'numproc'):
# print("WARNING: This is a uwsgi process but the uwsgidecorators library"
# " is unavailable. This is likely due to using an external (not"
# " in Galaxy's virtualenv) uwsgi and you may experience errors. "
# "HINT:\n {venv}/bin/pip install uwsgidecorators".format(
# venv=os.environ.get('VIRTUAL_ENV', '/path/to/venv')))
from six import string_types
from .message import ApplicationStackMessage, JobHandlerMessage, decode
from .transport import ApplicationStackTransport, UWSGIFarmMessageTransport
from galaxy.util.bunch import Bunch
from galaxy.util.configdict import parse_config
log = logging.getLogger(__name__)
@@ -94,30 +98,105 @@ class ApplicationStack(object):
transport_class = ApplicationStackTransport
log_filter_class = ApplicationStackLogFilter
purposes = Bunch(
WEB_WORKER = 'web-worker',
JOB_HANDLER = 'job-handler',
)
default_conf = {
'processes': 1,
'threads': 4,
'server_name': '{server_name}',
'workers': [],
'pools': [],
}
web_worker_pool = {
'name': 'web-workers',
'purpose': 'web-worker',
}
@classmethod
def register_postfork_function(cls, f, *args, **kwargs):
f(*args, **kwargs)
def __init__(self, app=None):
self.app = app
# FIXME: hardcoded path
self._conf = parse_config('config/stack_conf.yml', 'stack', default=self.default_conf)
self.server_name_template = '{server_name}'
self.pools = []
self.pool_names = []
self._purposes = []
def _postfork_init(self):
self.pools = []
self.pool_names = []
self._purposes = []
for pool in self._conf.get('pools', []) + [ApplicationStack.web_worker_pool]:
if self.in_pool(pool['name']):
self.pools.append(pool)
self.pool_names.append(pool['name'])
self._purposes.append(pool['purpose'])
for pool in self.pools:
if pool.get('server_name', None):
self.server_name_template = pool['server_name']
break
else:
if self._conf.get('server_name', None):
# default if the process is not in a pool and there is an unpooled server name set in the conf
self.server_name_template = self._conf['server_name']
if self.pools:
# default if the process is in a pool
self.server_name_template = '{server_name}.{pool_name}.{process_num}'
def start(self):
pass
self._postfork_init()
@property
def pool_name(self):
# TODO: in the future, allow for mapping job handlers in the job conf by something other than server_name, such
# as pool name and process number. but for now, the job-handlers pool should be the first one specified under a
# worker set's 'pools' list, if more than one pool is specified for a given worker set
try:
return self.pool_names[0]
except:
return None
def in_pool(self, pool_name):
return None
def has_purpose(self, purpose):
return purpose in self._purposes
# FIXME: used?
def workers(self):
return []
def allowed_middleware(self, middleware):
if hasattr(middleware, '__name__'):
middleware = middleware.__name__
return middleware not in self.prohibited_middleware
def set_postfork_server_name(self, app):
pass
def set_server_name(self, app, caller_tmpl_dict):
tmpl_dict = {
'server_name': app.config.server_name,
'pool_name': self.pool_name,
'process_num': 1,
'id': 1,
'hostname': socket.gethostname().split('.', 1)[0],
'fqdn': socket.getfqdn(),
}
tmpl_dict.update(caller_tmpl_dict)
app.config.server_name = self.server_name_template.format(**tmpl_dict)
log.debug('######## server_name is: %s', app.config.server_name)
def process_in_pool(self, pool_name):
return None
def set_postfork_server_name(self, app):
self.set_server_name(app, tmpl_dict, {})
def register_message_handler(self, func, name=None):
pass
@@ -128,6 +207,13 @@ class ApplicationStack(object):
def send_message(self, dest, msg=None, target=None, params=None, **kwargs):
pass
def send_purpose_message(self, purpose, **kwargs):
for pool in self.pools:
if purpose == pool['purpose']:
return self.send_message(pool['name'], **kwargs)
else:
raise RuntimeError('No pools defined for purpose: %s', purpose)
def shutdown(self):
pass
@@ -136,12 +222,13 @@ class MessageApplicationStack(ApplicationStack):
def __init__(self, app=None):
super(MessageApplicationStack, self).__init__(app=app)
self.dispatcher = ApplicationStackMessageDispatcher()
self.transport = self.transport_class(app, dispatcher=self.dispatcher)
self.transport = self.transport_class(app, stack=self, dispatcher=self.dispatcher)
#if app:
# log.debug('######## registering self.start')
# self.register_postfork_function(self.start)
def start(self):
super(MessageApplicationStack, self).start()
self.transport.start()
def register_message_handler(self, func, name=None):
@@ -182,8 +269,21 @@ class UWSGIApplicationStack(MessageApplicationStack):
# FIXME: this is copied into UWSGIFarmMessageTransport
shutdown_msg = '__SHUTDOWN__'
postfork_functions = []
# TODO: used?
handler_prefix = 'mule-handler-'
default_conf = {
'processes': 1,
'threads': 4,
'server_name': '{server_name}.{process_num}',
'workers': [{
'load': 'lib/galaxy/main.py',
'processes': 1,
'pools': ['job-handlers']
}],
'pools': [{
'name': 'job-handlers',
'purpose': 'job-handler',
}]
}
@classmethod
def register_postfork_function(cls, f, *args, **kwargs):
@@ -222,6 +322,16 @@ class UWSGIApplicationStack(MessageApplicationStack):
# uWSGI master will restart us
self.app.exit = True
def in_pool(self, pool_name):
if uwsgi.worker_id() > 0 and pool_name == ApplicationStack.web_worker_pool['name']:
return True
else:
return self._in_farm(pool_name)
#return self.transport._farm_name == pool_name
def workers(self):
return uwsgi.workers()
# FIXME: these are copied into UWSGIFarmMessageTransport
@property
def _farms(self):
@@ -242,6 +352,9 @@ class UWSGIApplicationStack(MessageApplicationStack):
self._mules_list = [mules] if isinstance(mules, string_types) or mules is True else mules
return self._mules_list
def _in_farm(self, farm_name):
return uwsgi.mule_id() in self._farms.get(farm_name, [])
def start(self):
self._is_mule = uwsgi.mule_id() > 0
if self._is_mule:
@@ -249,17 +362,20 @@ class UWSGIApplicationStack(MessageApplicationStack):
super(UWSGIApplicationStack, self).start()
def set_postfork_server_name(self, app):
tmpl_dict = {
'id': 'worker%s' % uwsgi.worker_id() if uwsgi.mule_id() == 0 else 'mule%s' % uwsgi.mule_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()
tmpl_dict['process_num'] = uwsgi.worker_id()
elif self.pool_name in self._farms:
tmpl_dict['process_num'] = self._farms[self.pool_name].index(uwsgi.mule_id()) + 1
log.debug('######## self._farms[self.pool_name]: %s', self._farms[self.pool_name])
self.set_server_name(app, tmpl_dict)
# FIXME: used?
def workers(self):
return uwsgi.workers()
def process_in_pool(self, pool_name):
return self.transport._farm_name == pool_name
#@property
#def pool_name(self):
# # could get this from self._conf, or could get it from uwsgi.opts
# return self.transport._farm_name
def shutdown(self):
log.debug('######## STACK SHUTDOWN CALLED')
@@ -309,17 +425,24 @@ def register_postfork_function(f, *args, **kwargs):
application_stack_class().register_postfork_function(f, *args, **kwargs)
def process_in_pool(pool_name):
return application_stack_instance().process_in_pool(pool_name)
def _mules():
mules = uwsgi.opt.get('mule', [])
return [mules] if isinstance(mules, string_types) or mules is True else mules
@uwsgi_postfork
#@uwsgi_postfork
def _do_uwsgi_postfork():
import os
# FIXME: _mules duplicated again
for i, mule in enumerate(_mules()):
if mule is not True and i + 1 == uwsgi.mule_id():
# mules will inherit the postfork function list and call them immediately upon fork, but programmed mules
# should not do that (they will call the postfork functions in-place as they start up after exec())
UWSGIApplicationStack.postfork_functions = []
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 = []
log.debug('######## postfork called, pid %s mule %s NOW functions are: %s' % (os.getpid(), uwsgi.mule_id(), UWSGIApplicationStack.postfork_functions))
for f, args, kwargs in [t for t in UWSGIApplicationStack.postfork_functions]:
f(*args, **kwargs)
if uwsgi:
uwsgi.post_fork_hook = _do_uwsgi_postfork
+19 -16
View File
@@ -27,10 +27,11 @@ class ApplicationStackTransport(object):
self.dispatcher_thread = threading.Thread(name=self.__class__.__name__ + ".dispatcher_thread", target=self._dispatch_messages)
#self.dispatcher_thread.daemon = True
def __init__(self, app, dispatcher=None):
def __init__(self, app, stack, dispatcher=None):
""" Pre-fork initialization.
"""
self.app = app
self.stack = stack
self.can_run = False
self.running = False
self.dispatcher = dispatcher
@@ -95,8 +96,8 @@ class UWSGIFarmMessageTransport(ApplicationStackTransport):
# uwsgi.set_option('locks', len(self.lock_map))
# log.debug('Created %s uWSGI locks' % len(self.lock_map))
def __init__(self, app, dispatcher=None):
super(UWSGIFarmMessageTransport, self).__init__(app, dispatcher=dispatcher)
def __init__(self, app, stack, dispatcher=None):
super(UWSGIFarmMessageTransport, self).__init__(app, stack, dispatcher=dispatcher)
self._farms_dict = None
self._mules_list = None
self._is_mule = False
@@ -135,6 +136,7 @@ class UWSGIFarmMessageTransport(ApplicationStackTransport):
uwsgi.unlock(self._locks.index(name_or_id))
def _farm_recv_msg_lock_num(self):
# need one lock per farm... except mules can have multiple farms... ack
return self._locks.index('RECV_MSG_FARM_' + self._farm_name)
def _dispatch_messages(self):
@@ -167,20 +169,19 @@ class UWSGIFarmMessageTransport(ApplicationStackTransport):
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'
raise RuntimeError('Mule %s is not in a farm! Set `farm = <pool_name>:%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 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()
log.info('######## Mule transport started, worker id: %s, mule id: %s, farm names: %s, server name: %s', uwsgi.worker_id(), uwsgi.mule_id(), self.stack.farm_names, self.app.config.server_name)
#self._send_all_messages()
@property
def _farm_name(self):
for name, mules in self._farms.items():
if uwsgi.mule_id() in mules:
return name
return None
#@property
#def _farm_name(self):
# for name, mules in self._farms.items():
# if uwsgi.mule_id() in mules:
# 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
@@ -196,6 +197,8 @@ class UWSGIFarmMessageTransport(ApplicationStackTransport):
def send_message(self, msg, dest):
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()
#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()
log.debug('######## Sending message to farm %s: %s', dest, msg)
uwsgi.farm_msg(dest, msg)
+2 -1
View File
@@ -163,6 +163,7 @@ def find_ini(supplied_ini, galaxy_root):
return guess
# FIXME: probably don't need this now that server names are fixed up in app. Do we really need fine grained control over mule names? Maybe for multihost mules? Do we support multihost mules? Could set the mule server name template in galaxy.ini...
def set_server_name(server_name):
default = None if not uwsgi else MULE_SERVER_NAME_DEFAULT
if not server_name:
@@ -241,7 +242,7 @@ def main():
args = arg_parser.parse_args()
if args.log_file:
os.environ["GALAXY_CONFIG_LOG_DESTINATION"] = os.path.abspath(args.log_file)
set_server_name(args.server_name)
#set_server_name(args.server_name)
pid_file = args.pid_file
log.setLevel(logging.DEBUG)
+4
View File
@@ -12,6 +12,7 @@ from galaxy.config import (
parse_dependency_options,
)
from galaxy.script import main_factory
from galaxy.util.configdict import parse_config
DESCRIPTION = "Script to determine uWSGI command line arguments."
@@ -19,6 +20,9 @@ COMMAND_TEMPLATE = '{virtualenv}--ini-paste {galaxy_ini} --paste-logger --die-on
def _get_uwsgi_args(args, kwargs):
# FIXME: hardcoded
stack_conf = parse_config('config/stack_conf.yml', 'stack', {})
# FIXME: these belong in stack conf
handlerct = int(kwargs.get('job_handler_count', 1))
pool_name = kwargs.get('job_handler_pool_name', 'job-handlers')
virtualenv = os.environ.get('VIRTUAL_ENV', None)