From 36e0e01b6b9da78ae6f37f556385befba513bc14 Mon Sep 17 00:00:00 2001 From: Nate Coraor Date: Mon, 13 Mar 2017 17:15:11 -0400 Subject: [PATCH] Run job handlers as uWSGI mules --- lib/galaxy/config.py | 24 ++++++++++++++++-------- lib/galaxy/jobs/__init__.py | 26 ++++++++++++++++++-------- lib/galaxy/jobs/handler.py | 23 +++++++++++++++++++++++ lib/galaxy/model/__init__.py | 9 ++++++++- lib/galaxy/util/handlers.py | 10 ++++++++++ 5 files changed, 75 insertions(+), 17 deletions(-) diff --git a/lib/galaxy/config.py b/lib/galaxy/config.py index 4408dabef24..146cd033518 100644 --- a/lib/galaxy/config.py +++ b/lib/galaxy/config.py @@ -20,6 +20,10 @@ from datetime import timedelta from six import string_types from six.moves import configparser +try: + import uwsgi +except: + uwsgi = None from galaxy.containers import parse_containers_config from galaxy.exceptions import ConfigurationError @@ -499,14 +503,17 @@ class Configuration(object): if self.heartbeat_log is None: self.heartbeat_log = 'heartbeat_{server_name}.log' # Determine which 'server:' this is - self.server_name = 'main' - for arg in sys.argv: - # Crummy, but PasteScript does not give you a way to determine this - if arg.lower().startswith('--server-name='): - self.server_name = arg.split('=', 1)[-1] - # Allow explicit override of server name in confg params - if "server_name" in kwargs: - self.server_name = kwargs.get("server_name") + if uwsgi and uwsgi.mule_id() > 0: + self.server_name = 'mule%d' % uwsgi.mule_id() + else: + self.server_name = 'main' + for arg in sys.argv: + # Crummy, but PasteScript does not give you a way to determine this + if arg.lower().startswith('--server-name='): + self.server_name = arg.split('=', 1)[-1] + # Allow explicit override of server name in confg params + if "server_name" in kwargs: + self.server_name = kwargs.get("server_name") # Store all configured server names self.server_names = [] for section in global_conf_parser.sections(): @@ -539,6 +546,7 @@ 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 ) ) # 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') diff --git a/lib/galaxy/jobs/__init__.py b/lib/galaxy/jobs/__init__.py index d0d8914d3d4..ce7ce6444c2 100644 --- a/lib/galaxy/jobs/__init__.py +++ b/lib/galaxy/jobs/__init__.py @@ -19,6 +19,10 @@ 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 @@ -166,6 +170,14 @@ class JobConfiguration(object, ConfiguresHandlers): except Exception as e: raise config_exception(e, job_config_file) + def __config_default_handlers(self): + """This is a stop-gap solution until the job conf loading is rewritten + for YAML job conf support. + """ + if not uwsgi: + self.handlers[self.app.config.server_name] = (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). @@ -205,9 +217,11 @@ class JobConfiguration(object, ConfiguresHandlers): # Must define at least one handler to have a default. if not self.handlers: - raise ValueError("Job configuration file defines no valid handler elements.") - # Determine the default handler(s) - self.default_handler_id = self._get_default(self.app.config, handlers_conf, list(self.handlers.keys())) + self.__config_default_handlers() + log.info("No handlers defined or empty handlers group, will use default handlers: %s", self.handlers) + else: + # Determine the default handler(s) + self.default_handler_id = self.__get_default(handlers, list(self.handlers.keys())) # Parse destinations destinations = root.find('destinations') @@ -341,11 +355,7 @@ class JobConfiguration(object, ConfiguresHandlers): self.runner_plugins.append(dict(id=runner, load=runner, workers=self.app.config.cluster_job_queue_workers)) # Set the handlers - for id in self.app.config.job_handlers: - self.handlers[id] = (id,) - - self.handlers['default_job_handlers'] = self.app.config.default_job_handlers - self.default_handler_id = 'default_job_handlers' + self.__config_default_handlers() # Set tool handler configs for id, tool_handlers in self.app.config.tool_handlers.items(): diff --git a/lib/galaxy/jobs/handler.py b/lib/galaxy/jobs/handler.py index 399987f4ba2..b3702116504 100644 --- a/lib/galaxy/jobs/handler.py +++ b/lib/galaxy/jobs/handler.py @@ -3,6 +3,7 @@ Galaxy job handler, prepares, runs, tracks, and finishes Galaxy jobs """ import datetime +import json import os import time import logging @@ -10,6 +11,10 @@ import threading from Queue import Queue, Empty from sqlalchemy.sql.expression import and_, or_, select, func, true, null +try: + import uwsgi +except ImportError: + uwsgi = None from galaxy import model from galaxy.util.sleeper import Sleeper @@ -76,6 +81,9 @@ class JobHandlerQueue(object): self.running = True self.monitor_thread = threading.Thread(name="JobHandlerQueue.monitor_thread", target=self.__monitor) self.monitor_thread.setDaemon(True) + if uwsgi: + self.mule_thread = threading.Thread( name="JobHandlerQueue.mule_thread", target=self.__mule ) + self.mule_thread.setDaemon( True ) def start(self): """ @@ -85,6 +93,8 @@ class JobHandlerQueue(object): self.__check_jobs_at_startup() # Start the queue self.monitor_thread.start() + if uwsgi: + self.mule_thread.start() log.info("job handler queue started") def job_wrapper(self, job, use_persisted_destination=False): @@ -171,6 +181,19 @@ class JobHandlerQueue(object): job_wrapper.job_runner_mapper.cached_job_destination = job_destination return job_wrapper + def __mule( self ): + while self.running: + try: + msg = uwsgi.mule_get_msg() + job_dict = json.loads(msg) + job = self.sa_session.query( model.Job ).get( job_dict['job_id'] ) + job.handler = self.app.config.server_name + self.sa_session.add( job ) + self.sa_session.flush() + except: + log.exception( "Exception in mule message handling" ) + time.sleep( 1 ) + def __monitor(self): """ Continually iterate the waiting jobs, checking is each is ready to diff --git a/lib/galaxy/model/__init__.py b/lib/galaxy/model/__init__.py index 9b06e090d68..ebd7f384a40 100644 --- a/lib/galaxy/model/__init__.py +++ b/lib/galaxy/model/__init__.py @@ -23,6 +23,10 @@ from sqlalchemy import (and_, func, join, not_, or_, select, true, type_coerce, types) from sqlalchemy.ext import hybrid from sqlalchemy.orm import aliased, joinedload, object_session +try: + import uwsgi +except ImportError: + uwsgi = None import galaxy.model.metadata import galaxy.model.orm.now @@ -650,7 +654,10 @@ class Job(object, JobLike, Dictifiable): self.imported = imported def set_handler(self, handler): - self.handler = handler + if uwsgi and handler is None: + uwsgi.mule_msg(json.dumps({'job_id': self.id, 'state': Job.states.NEW})) + else: + self.handler = handler def set_params(self, params): self.params = params diff --git a/lib/galaxy/util/handlers.py b/lib/galaxy/util/handlers.py index e494f9df5ed..9dee5d23761 100644 --- a/lib/galaxy/util/handlers.py +++ b/lib/galaxy/util/handlers.py @@ -8,6 +8,12 @@ import logging import os import random +try: + import uwsgi +except ImportError: + uwsgi = None + + log = logging.getLogger(__name__) @@ -100,6 +106,8 @@ class ConfiguresHandlers: :return: bool """ + if uwsgi and self.app.config.server_name.startswith('mule'): + return True for collection in self.handlers.values(): if server_name in collection: return True @@ -127,6 +135,8 @@ class ConfiguresHandlers: :returns: str -- A valid job handler ID. """ + if id_or_tag is None and self.default_handler_id is None: + return None if id_or_tag is None: id_or_tag = self.default_handler_id return self._get_single_item(self.handlers[id_or_tag], index=index)