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