diff --git a/lib/galaxy/app.py b/lib/galaxy/app.py index 29af9854746..f1674282ca4 100644 --- a/lib/galaxy/app.py +++ b/lib/galaxy/app.py @@ -55,7 +55,8 @@ class UniverseApplication(object, config.ConfiguresGalaxyMixin): self.name = 'galaxy' self.startup_timer = ExecutionTimer() self.new_installation = False - self.application_stack = application_stack_instance() + self.application_stack = application_stack_instance(app=self) + self.application_stack.register_postfork_function(self.application_stack.start) # Read config file and check for errors self.config = config.Configuration(**kwargs) self.config.check() @@ -233,6 +234,8 @@ class UniverseApplication(object, config.ConfiguresGalaxyMixin): os.unlink(self.datatypes_registry.integrated_datatypes_configs) except: pass + self.application_stack.shutdown() + # This is used to signal the webless application loop to terminate self.exit = True def configure_fluent_log(self): diff --git a/lib/galaxy/jobs/__init__.py b/lib/galaxy/jobs/__init__.py index ce7ce6444c2..5988ade78e5 100644 --- a/lib/galaxy/jobs/__init__.py +++ b/lib/galaxy/jobs/__init__.py @@ -19,10 +19,6 @@ from tempfile import NamedTemporaryFile from xml.etree import ElementTree import six -try: - import uwsgi -except ImportError: - uwsgi = None import galaxy from galaxy import model, util @@ -174,9 +170,9 @@ class JobConfiguration(object, ConfiguresHandlers): """This is a stop-gap solution until the job conf loading is rewritten for YAML job conf support. """ - if not uwsgi: + 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, + self.default_handler_id = self.app.config.server_name def __parse_job_conf_xml(self, tree): """Loads the new-style job configuration from options in the job config file (by default, job_conf.xml). diff --git a/lib/galaxy/jobs/handler.py b/lib/galaxy/jobs/handler.py index 562492f26bc..38b39cb1252 100644 --- a/lib/galaxy/jobs/handler.py +++ b/lib/galaxy/jobs/handler.py @@ -3,7 +3,6 @@ Galaxy job handler, prepares, runs, tracks, and finishes Galaxy jobs """ import datetime -import json import os import time import logging @@ -16,6 +15,7 @@ from galaxy import model from galaxy.util.sleeper import Sleeper from galaxy.jobs import JobWrapper, TaskWrapper, JobDestination from galaxy.jobs.mapper import JobNotReadyException +from galaxy.web.stack.message import JobHandlerMessage log = logging.getLogger(__name__) @@ -86,7 +86,8 @@ class JobHandlerQueue(object): self.__check_jobs_at_startup() # Start the queue self.monitor_thread.start() - self.app.application_stack.initialize_msg_handler(self.app, self.handle_msg) + # The stack code is initialized in the application + self.app.application_stack.register_message_handler(self.handle_msg, name=JobHandlerMessage.target) log.info("job handler queue started") def job_wrapper(self, job, use_persisted_destination=False): @@ -157,8 +158,6 @@ class JobHandlerQueue(object): if self.sa_session.dirty: self.sa_session.flush() - # If we are using a message queue or application stack messaging for handler assignment, signal those as well. - def __recover_job_wrapper(self, job): # Already dispatched and running job_wrapper = self.job_wrapper(job) @@ -634,8 +633,8 @@ class JobHandlerQueue(object): return JOB_WAIT return JOB_READY - def _handle_setup_msg(self, job_dict): - job = self.sa_session.query(model.Job).get(job_dict['job_id']) + def _handle_setup_msg(self, job_id=None): + job = self.sa_session.query(model.Job).get(job_id) if job.handler is None: job.handler = self.app.config.server_name self.sa_session.add(job) @@ -647,12 +646,9 @@ class JobHandlerQueue(object): def handle_msg(self, msg): try: - job_dict = json.loads(msg) - assert 'msg_type' in job_dict, "Missing required 'msg_type' message parameter" - getattr(self, '_handle_%s_msg' % job_dict['msg_type'])(job_dict) + getattr(self, '_handle_%s_msg' % msg.task)(**msg.params) except: log.exception( "Exception in mule message handling" ) - time.sleep( 1 ) def put(self, job_id, tool_id): """Add a job to the queue (by job identifier)""" @@ -668,9 +664,10 @@ class JobHandlerQueue(object): else: log.info("sending stop signal to worker thread") self.running = False - self.app.application_stack.terminate_handler() if not self.app.config.track_jobs_in_database: self.queue.put(self.STOP_SIGNAL) + # A message could still be received while shutting down, should be ok since they will be picked up on next startup. + self.app.application_stack.deregister_message_handler(name=JobHandlerMessage.target) self.sleeper.wake() log.info("job handler queue stopped") self.dispatcher.shutdown() diff --git a/lib/galaxy/jobs/manager.py b/lib/galaxy/jobs/manager.py index 4e916f3214c..d16d2005973 100644 --- a/lib/galaxy/jobs/manager.py +++ b/lib/galaxy/jobs/manager.py @@ -2,13 +2,13 @@ Top-level Galaxy job manager, moves jobs to handler(s) """ -import json import logging from sqlalchemy.sql.expression import null from galaxy.jobs import handler, NoopQueue from galaxy.model import Job +from galaxy.web.stack.message import JobHandlerMessage log = logging.getLogger(__name__) @@ -89,13 +89,11 @@ class MessageJobQueue(object): self.app = app def put(self, job_id, tool_id): - # FIXME: uwsgi farm name hardcoded here - # TODO: probably need a single class that encodes and decodes messages - self.app.application_stack.send_msg(json.dumps({'msg_type': 'setup', 'job_id': job_id, 'state': Job.states.NEW}), self.app.config.job_handler_pool_name) + msg = JobHandlerMessage(task='setup', job_id=job_id) + self.app.application_stack.send_message(self.app.config.job_handler_pool_name, msg) def put_stop(self, *args): - return + pass def shutdown(self): - #self.app.application_stack.send_msg(json.dumps({'msg_type': 'shutdown'})) ... - return + pass diff --git a/lib/galaxy/web/stack/__init__.py b/lib/galaxy/web/stack/__init__.py index a86c52d9a04..e920bcd8509 100644 --- a/lib/galaxy/web/stack/__init__.py +++ b/lib/galaxy/web/stack/__init__.py @@ -1,6 +1,6 @@ """Web application stack operations """ -from __future__ import print_function +from __future__ import absolute_import, print_function import inspect import logging @@ -29,21 +29,67 @@ except: from six import string_types -from galaxy import model +from .message import ApplicationStackMessage, JobHandlerMessage, decode +from .transport import ApplicationStackTransport, UWSGIFarmMessageTransport log = logging.getLogger(__name__) +class ApplicationStackMessageDispatcher(object): + def __init__(self): + self.__funcs = {} + + def __func_name(self, func, name): + if not name: + name = func.__name__ + return name + + def register_func(self, func, name=None): + name = self.__func_name(func, name) + self.__funcs[name] = func + + def deregister_func(self, func=None, name=None): + name = self.__func_name(func, name) + del self.__func[name] + + @property + def handler_count(self): + return len(self.__funcs) + + def dispatch(self, msg_str): + msg = decode(msg_str) + try: + msg.validate() + except AssertionError as exc: + log.error('######## Invalid message received: %s, error: %s', msg_str, exc) + return + if msg.target not in self.__funcs: + log.error("######## Received message with target '%s' but no functions were registered with that name. Params were: %s", msg.target, msg.params) + else: + self.__funcs[msg.target](msg) + + class ApplicationStack(object): name = None prohibited_middleware = frozenset() - setup_jobs_with_msg = False + # TODO: this is a fairly clunky way of handling these cases + setup_jobs_with_msg = False # used in galaxy.jobs.manager to determine whether jobs should be sent via message + handle_jobs = False # used by galaxy.jobs to determine whether this process handles jobs + + transport_class = ApplicationStackTransport @classmethod def register_postfork_function(cls, f, *args, **kwargs): f(*args, **kwargs) + def __init__(self, app=None): + self.app = app + + def start(self): + pass + + # FIXME: used? def workers(self): return [] @@ -52,20 +98,61 @@ class ApplicationStack(object): middleware = middleware.__name__ return middleware not in self.prohibited_middleware - def process_in_pool(self, pool_name): - return None - - def initialize_msg_handler(self, app): - return None - - def send_msg(self, msg, dest): - pass - def set_postfork_server_name(self, app): pass + def process_in_pool(self, pool_name): + return None -class UWSGIApplicationStack(ApplicationStack): + def register_message_handler(self, func, name=None): + pass + + def deregister_message_handler(self, func=None, name=None): + pass + + def send_message(self, dest, msg=None, target=None, params=None, **kwargs): + pass + + def shutdown(self): + pass + + +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) + #if app: + # log.debug('######## registering self.start') + # self.register_postfork_function(self.start) + + def start(self): + 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() + + def deregister_message_handler(self, func=None, name=None): + self.dispatcher.deregister_func(func, name) + self.transport.shutdown_if_needed() + + def send_message(self, dest, msg=None, target=None, params=None, **kwargs): + assert msg is not None or target is not None, "Either 'msg' or 'target' parameters must be set" + if not msg: + msg = ApplicationStackMessage( + target=target, + params=params, + **kwargs + ) + self.transport.send_message(msg.encode(), dest) + + def shutdown(self): + self.transport.shutdown() + + +class UWSGIApplicationStack(MessageApplicationStack): """Interface to the uWSGI application stack. Supports running additional webless Galaxy workers as mules. Mules must be farmed to be communicable via uWSGI mule messaging, unfarmed mules are not supported. """ @@ -76,16 +163,38 @@ class UWSGIApplicationStack(ApplicationStack): ]) setup_jobs_with_msg = True + transport_class = UWSGIFarmMessageTransport + # FIXME: this is copied into UWSGIFarmMessageTransport + shutdown_msg = '__SHUTDOWN__' postfork_functions = [] + # TODO: used? handler_prefix = 'mule-handler-' - # Define any static lock names here, additional locks will be appended for each configured farm's message handler - _locks = [] - @classmethod def register_postfork_function(cls, f, *args, **kwargs): cls.postfork_functions.append((f, args, kwargs)) + def __init__(self, app=None): + super(UWSGIApplicationStack, self).__init__(app=app) + + def __register_signal_handlers(self): + for name in ('TERM', 'INT', 'HUP'): + 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()) + 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 + + # FIXME: these are copied into UWSGIFarmMessageTransport @property def _farms(self): if self._farms_dict is None: @@ -105,125 +214,36 @@ class UWSGIApplicationStack(ApplicationStack): self._mules_list = [mules] if isinstance(mules, string_types) else mules return self._mules_list - def __initialize_locks(self): - num = int(uwsgi.opt.get('locks', 0)) + 1 - farms = self._farms.keys() - need = len(farms) - #need = len(filter(lambda x: x.startswith('LOCK_'), dir(self))) - if num < need: - raise RuntimeError('Need %i uWSGI locks but only %i exist(s): Set `locks = %i` in uWSGI configuration' % (need, num, need - 1)) - sys.exit(1) - self._locks.extend(map(lambda x: 'RECV_MSG_FARM_' + x, farms)) - # This would be nice, but in my 2.0.15 uWSGI, the uwsgi module has no set_option function. And I don't know if it'd work. - #if len(self.lock_map) > 1: - # uwsgi.set_option('locks', len(self.lock_map)) - # log.debug('Created %s uWSGI locks' % len(self.lock_map)) - - def __register_signal_handlers(self): - for name in ('TERM', 'INT', 'HUP'): - sig = getattr(signal, 'SIG%s' % name) - signal.signal(sig, self._handle_signal) - - def __handle_msgs(self): - while self.running: - # We are going to do this a lot, so cache the lock number - lock = self._farm_recv_msg_lock_num() - try: - self._lock(lock) - log.debug('######## Mule %s acquired message receive lock, waiting for new message', uwsgi.mule_id()) - msg = uwsgi.farm_get_msg() - log.debug('######## Mule %s received message: %s', uwsgi.mule_id(), msg) - self._unlock(lock) - self.msg_handler(msg) - except: - log.exception( "Exception in mule message handling" ) - log.info('uWSGI Mule (id: %s) message handler shutting down', uwsgi.mule_id()) - - def __init__(self): - super(UWSGIApplicationStack, self).__init__() - self._farms_dict = None - self._mules_list = None - self.app = None - self.running = False - self.msg_handler = None - self.msg_thread = None - self.msg_thread = threading.Thread(name="UWSGIApplicationStack.mule_msg_thread", target=self.__handle_msgs) - self.msg_thread.daemon = True - self.__initialize_locks() - - def _lock(self, name_or_id): - try: - uwsgi.lock(name_or_id) - except TypeError: - uwsgi.lock(self._locks.index(name_or_id)) - - def _unlock(self, name_or_id): - try: - uwsgi.unlock(name_or_id) - except TypeError: - uwsgi.unlock(self._locks.index(name_or_id)) - - @property - def _farm_name(self): - for name, mules in self._farms.items(): - if uwsgi.mule_id() in mules: - return name - return None - - def _farm_recv_msg_lock_num(self): - return self._locks.index('RECV_MSG_FARM_' + self._farm_name) - - 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()) - # 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()) - # uWSGI master will restart us - self.app.exit = True - - def workers(self): - return uwsgi.workers() - - def process_in_pool(self, pool_name): - return self._farm_name == pool_name - - def initialize_msg_handler(self, app, handler_func): - """ Post-fork initialization. - """ - self.app = app - self.running = True - if not uwsgi.in_farm(): - raise RuntimeError('Mule %s is not in a farm! Set `farm = %s:%s` in uWSGI configuration' - % (uwsgi.mule_id(), - app.config.job_handler_pool_name, - ','.join(map(str, range(1, len(filter(lambda x: x.endswith('galaxy/web/stack/uwsgi_mule/handler.py'), self._mules)) + 1))))) + def start(self): self.__register_signal_handlers() - self.msg_handler = handler_func - self.msg_thread.start() - log.info('######## Job handler mule started, mule id: %s, farm name: %s, handler id: %s', uwsgi.mule_id(), self._farm_name, self.app.config.server_name) - - def send_msg(self, msg, dest): - # TODO: the handler farm name should be configurable - #log.debug('################## sending message from mule %s: %s', uwsgi.mule_id(), msg) - log.debug('######## Sending message to farm %s: %s', dest, msg) - uwsgi.farm_msg(dest, msg) - log.debug('######## Message sent') - - def terminate_handler(self): - self.running = False + super(UWSGIApplicationStack, self).start() def set_postfork_server_name(self, app): app.config.server_name += ".%d" % uwsgi.worker_id() + # FIXME: used? + def workers(self): + return uwsgi.workers() + + def process_in_pool(self, pool_name): + return self.transport._farm_name == pool_name + + def shutdown(self): + 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) + class PasteApplicationStack(ApplicationStack): name = 'Python Paste' + handle_jobs = True class WeblessApplicationStack(ApplicationStack): name = 'Webless' + handle_jobs = True def application_stack_class(): @@ -240,9 +260,9 @@ def application_stack_class(): return WeblessApplicationStack -def application_stack_instance(): +def application_stack_instance(app=None): stack_class = application_stack_class() - return stack_class() + return stack_class(app=app) def register_postfork_function(f, *args, **kwargs): diff --git a/lib/galaxy/web/stack/message.py b/lib/galaxy/web/stack/message.py new file mode 100644 index 00000000000..d9932011ac2 --- /dev/null +++ b/lib/galaxy/web/stack/message.py @@ -0,0 +1,103 @@ +"""Web Application Stack worker messaging +""" +from __future__ import absolute_import + +import json +import logging + + +log = logging.getLogger(__name__) + + +class ApplicationStackMessageEncoder(json.JSONEncoder): + def encode(self, o): + if isinstance(o, AttributeDict): + return o.serialize(self) + return json.JSONEncoder.encode(self, o) + + +class ApplicationStackMessage(dict): + def __init__(self, target=None, params=None, **kwargs): + """Any extra kwargs override values in params + """ + self['target'] = target + self['params'] = params or {} + for k, v in kwargs.items(): + self['params'][k] = v + + #@classmethod + #def from_string(cls, s) + # #kwargs['object_hook'] = ApplicationStackMessage + # d = json.loads(s) + # return cls(json.loads(s, *args, **kwargs)) + + def validate(self): + assert self['target'] is not None, "Missing 'target' parameter" + + def serialize(self, encoder): + #return json.JSONEncoder.encode(encoder, ApplicationStackMessage.clean(self)) + self.validate() + return json.JSONEncoder.encode(encoder, self) + + def dumps(self, *args, **kwargs): + #kwargs['cls'] = AttributeStackMessageEncoder + return json.dumps(self, *args, **kwargs) + + def dump(self, fp, *args, **kwargs): + #raise Exception("This isn't using the Encoder, what gives?") + #kwargs['cls'] = AttributeStackMessageEncoder + return json.dump(self, fp, *args, **kwargs) + + @property + def target(self): + return self['target'] + + @target.setter + def set_target(self, target): + self['target'] = target + + @property + def params(self): + return self['params'] + + @params.setter + def set_params(self, params): + self['params'] = params + + #property + #def class(self): + # return globals()[self['cls']] + + def encode(self): + self['__classname__'] = self.__class__.__name__ + return json.dumps(self) + + +# when we add additional messages this should become generalized and subclassed +class JobHandlerMessage(ApplicationStackMessage): + target = 'job_handler' + + def __init__(self, target=None, params=None, **kwargs): + super(JobHandlerMessage, self).__init__(params=params, **kwargs) + self['target'] = self.target + + def validate(self): + super(JobHandlerMessage, self).validate() + for param in ('task', 'job_id'): + assert param in self['params'], "Missing required parameter '%s'" % param + + @property + def task(self): + return self['params']['task'] + + @property + def params(self): + d = self['params'].copy() + del d['task'] + return d + + +def decode(msg_str): + d = json.loads(s) + cls = d.pop('__classname__') + return globals()[cls](d) diff --git a/lib/galaxy/web/stack/transport.py b/lib/galaxy/web/stack/transport.py new file mode 100644 index 00000000000..429d0e4138f --- /dev/null +++ b/lib/galaxy/web/stack/transport.py @@ -0,0 +1,178 @@ +"""Web application stack operations +""" +from __future__ import absolute_import + +import logging +import threading + +try: + import uwsgi +except ImportError: + uwsgi = None + +from six import string_types + + +log = logging.getLogger(__name__) + + +class ApplicationStackTransport(object): + 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 + + def __init__(self, app, dispatcher=None): + """ Pre-fork initialization. + """ + self.app = app + self.can_run = False + self.running = False + self.dispatcher = dispatcher + self.dispatcher_thread = None + self.__init_dispatcher_thread() + + def _dispatch_messages(self): + pass + + 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') + + 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: + self.running = False + self.dispatcher_thread.join() + self.__init_dispatcher_thread() + + def start(self): + """ Post-fork initialization. + """ + self.can_run = True + log.debug('######## start called') + self.start_if_needed() + + def send_message(self, msg, dest): + pass + + def shutdown(self): + self.running = False + self.dispatcher_thread.join() + + +class UWSGIFarmMessageTransport(ApplicationStackTransport): + shutdown_msg = '__SHUTDOWN__' + # Define any static lock names here, additional locks will be appended for each configured farm's message handler + _locks = [] + + def __initialize_locks(self): + num = int(uwsgi.opt.get('locks', 0)) + 1 + farms = self._farms.keys() + need = len(farms) + #need = len(filter(lambda x: x.startswith('LOCK_'), dir(self))) + if num < need: + raise RuntimeError('Need %i uWSGI locks but only %i exist(s): Set `locks = %i` in uWSGI configuration' % (need, num, need - 1)) + sys.exit(1) + self._locks.extend(map(lambda x: 'RECV_MSG_FARM_' + x, farms)) + # This would be nice, but in my 2.0.15 uWSGI, the uwsgi module has no set_option function. And I don't know if it'd work. + #if len(self.lock_map) > 1: + # 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) + self._farms_dict = None + self._mules_list = None + self.__initialize_locks() + + @property + def _farms(self): + if self._farms_dict is None: + self._farms_dict = {} + farms = uwsgi.opt.get('farm', []) + farms = [farms] if isinstance(farms, string_types) else farms + for farm in farms: + name, mules = farm.split(':', 1) + self._farms_dict[name] = [int(m) for m in mules.split(',')] + return self._farms_dict + + @property + def _mules(self): + 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 + return self._mules_list + + def __lock(self, name_or_id): + try: + uwsgi.lock(name_or_id) + except TypeError: + uwsgi.lock(self._locks.index(name_or_id)) + + def __unlock(self, name_or_id): + try: + uwsgi.unlock(name_or_id) + except TypeError: + uwsgi.unlock(self._locks.index(name_or_id)) + + def _farm_recv_msg_lock_num(self): + return self._locks.index('RECV_MSG_FARM_' + self._farm_name) + + def _dispatch_messages(self): + # We are going to do this a lot, so cache the lock number + lock = self._farm_recv_msg_lock_num() + while self.running: + msg = None + self.__lock(lock) + 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) + except: + log.exception( "Exception in mule message handling" ) + finally: + self.__unlock(lock) + if msg: + 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))))) + 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) + + @property + def _farm_name(self): + for name, mules in self._farms.items(): + if uwsgi.mule_id() in mules: + return name + return None + + 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')