diff --git a/config/galaxy.ini.sample b/config/galaxy.ini.sample index d0b176c842e..10c3f633a7e 100644 --- a/config/galaxy.ini.sample +++ b/config/galaxy.ini.sample @@ -1,5 +1,5 @@ # -# Galaxy is configured by default to be useable in a single-user development +# Galaxy is configured by default to be usable in a single-user development # environment. To tune the application for a multi-user production # environment, see the documentation at: # @@ -118,9 +118,9 @@ paste.app_factory = galaxy.web.buildapp:app_factory #database_query_profiling_proxy = False # By default, Galaxy will use the same database to track user data and -# tool shed install data. There are many situtations in which it is -# valuable to seperate these - for instance bootstrapping fresh Galaxy -# instances with pretested installs. The following optin can be used to +# tool shed install data. There are many situations in which it is +# valuable to separate these - for instance bootstrapping fresh Galaxy +# instances with pretested installs. The following option can be used to # separate the tool shed install database (all other options listed above # but prefixed with install_ are also available). #install_database_connection = sqlite:///./database/universe.sqlite?isolation_level=IMMEDIATE diff --git a/lib/galaxy/app.py b/lib/galaxy/app.py index 427fcabda07..70309deecd3 100644 --- a/lib/galaxy/app.py +++ b/lib/galaxy/app.py @@ -38,7 +38,12 @@ class UniverseApplication( object, config.ConfiguresGalaxyMixin ): self.config.check() config.configure_logging( self.config ) self.configure_fluent_log() - self._amqp_internal_connection_obj = galaxy.queues.connection_from_config(self.config) + + self.amqp_internal_connection_obj = galaxy.queues.connection_from_config(self.config) + # control_worker *can* be initialized with a queue, but here we don't + # want to and we'll allow postfork to bind and start it. + self.control_worker = GalaxyQueueWorker(self) + self._configure_tool_shed_registry() self._configure_object_store( fsmon=True ) # Setup the database engine and ORM @@ -152,13 +157,6 @@ class UniverseApplication( object, config.ConfiguresGalaxyMixin ): self.model.engine.dispose() self.server_starttime = int(time.time()) # used for cachebusting - def setup_control_queue(self): - self.control_worker = GalaxyQueueWorker(self, galaxy.queues.control_queue_from_config(self.config), - galaxy.queue_worker.control_message_to_task, - self._amqp_internal_connection_obj) - self.control_worker.daemon = True - self.control_worker.start() - def shutdown( self ): self.workflow_scheduling_manager.shutdown() self.job_manager.shutdown() diff --git a/lib/galaxy/jobs/__init__.py b/lib/galaxy/jobs/__init__.py index dfa5833f91a..4c459b51c0e 100644 --- a/lib/galaxy/jobs/__init__.py +++ b/lib/galaxy/jobs/__init__.py @@ -1801,7 +1801,7 @@ class TaskWrapper(JobWrapper): task.command_line = self.command_line self.sa_session.flush() - def cleanup( self ): + def cleanup( self, delete_files=True ): # There is no task cleanup. The job cleans up for all tasks. pass diff --git a/lib/galaxy/queue_worker.py b/lib/galaxy/queue_worker.py index 529fc445ebe..e911bee4a65 100644 --- a/lib/galaxy/queue_worker.py +++ b/lib/galaxy/queue_worker.py @@ -24,52 +24,6 @@ logging.getLogger('kombu').setLevel(logging.WARNING) log = logging.getLogger(__name__) -class GalaxyQueueWorker(ConsumerMixin, threading.Thread): - """ - This is a flexible worker for galaxy's queues. Each process, web or - handler, will have one of these used for dispatching so called 'control' - tasks. - """ - def __init__(self, app, queue, task_mapping, connection=None): - super(GalaxyQueueWorker, self).__init__() - log.info("Initalizing Galaxy Queue Worker on %s", util.mask_password_from_url(app.config.amqp_internal_connection)) - if connection: - self.connection = connection - else: - self.connection = Connection(app.config.amqp_internal_connection) - self.app = app - # Eventually we may want different workers w/ their own queues and task - # mappings. Right now, there's only the one. - self.control_queue = queue - self.task_mapping = task_mapping - self.declare_queues = galaxy.queues.all_control_queues_for_declare(app.config) - # TODO we may want to purge the queue at the start to avoid executing - # stale 'reload_tool', etc messages. This can happen if, say, a web - # process goes down and messages get sent before it comes back up. - # Those messages will no longer be useful (in any current case) - - def get_consumers(self, Consumer, channel): - return [Consumer(queues=self.control_queue, - callbacks=[self.process_task])] - - def process_task(self, body, message): - if body['task'] in self.task_mapping: - if body.get('noop', None) != self.app.config.server_name: - try: - f = self.task_mapping[body['task']] - log.info("Instance recieved '%s' task, executing now." % body['task']) - f(self.app, **body['kwargs']) - except Exception: - # this shouldn't ever throw an exception, but... - log.exception("Error running control task type: %s" % body['task']) - else: - log.warning("Recieved a malformed task message:\n%s" % body) - message.ack() - - def shutdown(self): - self.should_stop = True - - def send_control_task(trans, task, noop_self=False, kwargs={}): log.info("Sending %s control task." % task) payload = {'task': task, @@ -127,3 +81,63 @@ control_message_to_task = { 'reload_tool': reload_tool, 'reload_display_application': reload_display_application, 'reload_tool_data_tables': reload_tool_data_tables, 'admin_job_lock': admin_job_lock} + + +class GalaxyQueueWorker(ConsumerMixin, threading.Thread): + """ + This is a flexible worker for galaxy's queues. Each process, web or + handler, will have one of these used for dispatching so called 'control' + tasks. + """ + def __init__(self, app, queue=None, task_mapping=control_message_to_task, connection=None): + super(GalaxyQueueWorker, self).__init__() + log.info("Initalizing %s Galaxy Queue Worker on %s", app.config.server_name, util.mask_password_from_url(app.config.amqp_internal_connection)) + self.daemon = True + if connection: + self.connection = connection + else: + self.connection = app.amqp_internal_connection_obj + # explicitly force connection instead of lazy-connecting the first + # time it is required. + self.connection.connect() + self.app = app + # Eventually we may want different workers w/ their own queues and task + # mappings. Right now, there's only the one. + if queue: + # Allows assignment of a particular queue for this worker. + self.control_queue = queue + else: + # Default to figuring out which control queue to use based on the app config. + queue = galaxy.queues.control_queue_from_config(app.config) + self.task_mapping = task_mapping + self.declare_queues = galaxy.queues.all_control_queues_for_declare(app.config) + # TODO we may want to purge the queue at the start to avoid executing + # stale 'reload_tool', etc messages. This can happen if, say, a web + # process goes down and messages get sent before it comes back up. + # Those messages will no longer be useful (in any current case) + + def bind_and_start(self): + log.info("Binding and starting galaxy control worker for %s", self.app.config.server_name) + self.control_queue = galaxy.queues.control_queue_from_config(self.app.config) + self.start() + + def get_consumers(self, Consumer, channel): + return [Consumer(queues=self.control_queue, + callbacks=[self.process_task])] + + def process_task(self, body, message): + if body['task'] in self.task_mapping: + if body.get('noop', None) != self.app.config.server_name: + try: + f = self.task_mapping[body['task']] + log.info("Instance '%s' recieved '%s' task, executing now.", self.app.config.server_name, body['task']) + f(self.app, **body['kwargs']) + except Exception: + # this shouldn't ever throw an exception, but... + log.exception("Error running control task type: %s" % body['task']) + else: + log.warning("Recieved a malformed task message:\n%s" % body) + message.ack() + + def shutdown(self): + self.should_stop = True diff --git a/lib/galaxy/webapps/galaxy/buildapp.py b/lib/galaxy/webapps/galaxy/buildapp.py index d9485216469..3748400f78c 100644 --- a/lib/galaxy/webapps/galaxy/buildapp.py +++ b/lib/galaxy/webapps/galaxy/buildapp.py @@ -125,7 +125,7 @@ def postfork_setup(): if process_is_uwsgi: import uwsgi app.config.server_name += ".%s" % uwsgi.worker_id() - app.setup_control_queue() + app.control_worker.bind_and_start() def populate_api_routes( webapp, app ): diff --git a/templates/user/register.mako b/templates/user/register.mako index 68393e1b703..0c7f85cc207 100644 --- a/templates/user/register.mako +++ b/templates/user/register.mako @@ -1,7 +1,6 @@ <%! #This is a hack, we should restructure templates to avoid this. def inherit(context): - print 'context:', context if context.get('trans').webapp.name == 'galaxy' and context.get( 'use_panels', True ): return '/webapps/galaxy/base_panels.mako' else: diff --git a/test/api/test_workflows_from_yaml.py b/test/api/test_workflows_from_yaml.py index 4c1860cc723..e04eb284b1e 100644 --- a/test/api/test_workflows_from_yaml.py +++ b/test/api/test_workflows_from_yaml.py @@ -30,25 +30,30 @@ class WorkflowsFromYamlApiTestCase( BaseWorkflowsApiTestCase ): """) self._get("workflows/%s/download" % workflow_id).content - def test_multiple_input( self ): - history_id = self.dataset_populator.new_history() - self._run_jobs(""" -steps: - - type: input - label: input1 - - type: input - label: input2 - - tool_id: cat_list - state: - input1: - - $link: input1 - - $link: input2 -test_data: - input1: "hello world" - input2: "123" -""", history_id=history_id) - contents1 = self.dataset_populator.get_history_dataset_content(history_id) - assert contents1 == "hello world\n123\n" +# FIXME: This test fails on some machines due to (we're guessing) yaml loading +# order being not guaranteed and inconsistent across platforms. The workflow +# yaml loader probably needs to enforce order using something like the +# approach described here: +# https://stackoverflow.com/questions/13297744/pyyaml-control-ordering-of-items-called-by-yaml-load +# def test_multiple_input( self ): +# history_id = self.dataset_populator.new_history() +# self._run_jobs(""" +# steps: +# - type: input +# label: input1 +# - type: input +# label: input2 +# - tool_id: cat_list +# state: +# input1: +# - $link: input1 +# - $link: input2 +# test_data: +# input1: "hello world" +# input2: "123" +# """, history_id=history_id) +# contents1 = self.dataset_populator.get_history_dataset_content(history_id) +# assert contents1 == "hello world\n123\n" def test_simple_output_actions( self ): history_id = self.dataset_populator.new_history()