mirror of
https://github.com/galaxyproject/galaxy.git
synced 2026-09-24 16:30:27 +08:00
WIP job handlers as uWSGI mules refactor
This commit is contained in:
+4
-1
@@ -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):
|
||||
|
||||
@@ -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).
|
||||
|
||||
@@ -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()
|
||||
|
||||
@@ -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
|
||||
|
||||
+145
-125
@@ -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):
|
||||
|
||||
@@ -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)
|
||||
@@ -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')
|
||||
Reference in New Issue
Block a user