From fbec530eb4d3b3fad1b40fba294b937605251cd0 Mon Sep 17 00:00:00 2001 From: Dannon Baker Date: Fri, 1 May 2015 10:35:58 -0400 Subject: [PATCH 1/7] Remove unused imports during queue refactoring. --- lib/galaxy/webapps/galaxy/buildapp.py | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/lib/galaxy/webapps/galaxy/buildapp.py b/lib/galaxy/webapps/galaxy/buildapp.py index f6148079644..bd4a166edb8 100644 --- a/lib/galaxy/webapps/galaxy/buildapp.py +++ b/lib/galaxy/webapps/galaxy/buildapp.py @@ -3,8 +3,6 @@ Provides factory methods to assemble the Galaxy web application """ import atexit -import os -import os.path import sys from paste import httpexceptions @@ -535,6 +533,7 @@ def populate_api_routes( webapp, app ): webapp.mapper.connect( "create", "/api/metrics", controller="metrics", action="create", conditions=dict( method=["POST"] ) ) + def _add_item_tags_controller( webapp, name_prefix, path_prefix, **kwd ): # Not just using map.resources because actions should be based on name not id controller = "%stags" % name_prefix From 81f864626a96cb50b3b943e7441a00dd9668dbcd Mon Sep 17 00:00:00 2001 From: Dannon Baker Date: Fri, 1 May 2015 13:34:39 -0400 Subject: [PATCH 2/7] Split initialization w/ early forced connection of queue_worker. This should resolve mapper initialization issues. --- lib/galaxy/app.py | 15 ++-- lib/galaxy/queue_worker.py | 101 ++++++++++++++------------ lib/galaxy/webapps/galaxy/buildapp.py | 2 +- 3 files changed, 65 insertions(+), 53 deletions(-) diff --git a/lib/galaxy/app.py b/lib/galaxy/app.py index 1a75554b460..52676903882 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,11 +157,9 @@ 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 + def bind_and_start_control_queue(self): + log.info("Binding and starting galaxy control worker for %s", app.config.server_name) + self.control_worker.control_queue = galaxy.queues.control_queue_from_config(self.config) self.control_worker.start() def shutdown( self ): diff --git a/lib/galaxy/queue_worker.py b/lib/galaxy/queue_worker.py index 529fc445ebe..e97fce3f5f1 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,58 @@ 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 Galaxy Queue Worker on %s", 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 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 diff --git a/lib/galaxy/webapps/galaxy/buildapp.py b/lib/galaxy/webapps/galaxy/buildapp.py index bd4a166edb8..15b98c78f31 100644 --- a/lib/galaxy/webapps/galaxy/buildapp.py +++ b/lib/galaxy/webapps/galaxy/buildapp.py @@ -124,7 +124,7 @@ def postfork_setup(): if app.config.is_uwsgi: import uwsgi app.config.server_name += ".%s" % uwsgi.worker_id() - app.setup_control_queue() + app.bind_and_start_control_queue() def populate_api_routes( webapp, app ): From d96035deecd177bb5835cb1170557ddead5da407 Mon Sep 17 00:00:00 2001 From: Dannon Baker Date: Fri, 1 May 2015 13:55:15 -0400 Subject: [PATCH 3/7] Slight refactoring, improve debugging (w/ app identifiers) for amqp workers --- lib/galaxy/app.py | 5 ----- lib/galaxy/queue_worker.py | 9 +++++++-- lib/galaxy/webapps/galaxy/buildapp.py | 2 +- 3 files changed, 8 insertions(+), 8 deletions(-) diff --git a/lib/galaxy/app.py b/lib/galaxy/app.py index 52676903882..018a5583ba0 100644 --- a/lib/galaxy/app.py +++ b/lib/galaxy/app.py @@ -157,11 +157,6 @@ class UniverseApplication( object, config.ConfiguresGalaxyMixin ): self.model.engine.dispose() self.server_starttime = int(time.time()) # used for cachebusting - def bind_and_start_control_queue(self): - log.info("Binding and starting galaxy control worker for %s", app.config.server_name) - self.control_worker.control_queue = galaxy.queues.control_queue_from_config(self.config) - self.control_worker.start() - def shutdown( self ): self.workflow_scheduling_manager.shutdown() self.job_manager.shutdown() diff --git a/lib/galaxy/queue_worker.py b/lib/galaxy/queue_worker.py index e97fce3f5f1..e911bee4a65 100644 --- a/lib/galaxy/queue_worker.py +++ b/lib/galaxy/queue_worker.py @@ -91,7 +91,7 @@ class GalaxyQueueWorker(ConsumerMixin, threading.Thread): """ def __init__(self, app, queue=None, task_mapping=control_message_to_task, connection=None): super(GalaxyQueueWorker, self).__init__() - log.info("Initalizing Galaxy Queue Worker on %s", util.mask_password_from_url(app.config.amqp_internal_connection)) + 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 @@ -116,6 +116,11 @@ class GalaxyQueueWorker(ConsumerMixin, threading.Thread): # 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])] @@ -125,7 +130,7 @@ class GalaxyQueueWorker(ConsumerMixin, threading.Thread): 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']) + 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... diff --git a/lib/galaxy/webapps/galaxy/buildapp.py b/lib/galaxy/webapps/galaxy/buildapp.py index 15b98c78f31..06aaa4bc8d2 100644 --- a/lib/galaxy/webapps/galaxy/buildapp.py +++ b/lib/galaxy/webapps/galaxy/buildapp.py @@ -124,7 +124,7 @@ def postfork_setup(): if app.config.is_uwsgi: import uwsgi app.config.server_name += ".%s" % uwsgi.worker_id() - app.bind_and_start_control_queue() + app.control_worker.bind_and_start() def populate_api_routes( webapp, app ): From 06e85e5858133b35c8fd0e9908c42df4eb247b0c Mon Sep 17 00:00:00 2001 From: Daniel Blankenberg Date: Mon, 4 May 2015 12:15:56 -0400 Subject: [PATCH 4/7] Fix some minor typos in comment docs in config/galaxy.ini.sample. --- config/galaxy.ini.sample | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) 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 From 89b96984a2df9afed5e22cda3ee01c788c42f06d Mon Sep 17 00:00:00 2001 From: Dannon Baker Date: Mon, 4 May 2015 12:48:50 -0400 Subject: [PATCH 5/7] Remove test_workflow_from_yaml due to inconsistent yaml loading (makes it a bad test). Details in comment. --- test/api/test_workflows_from_yaml.py | 43 ++++++++++++++++------------ 1 file changed, 24 insertions(+), 19 deletions(-) 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() From e616a51d9e8f2162ea36414c4f0d812adc66e2a3 Mon Sep 17 00:00:00 2001 From: Dannon Baker Date: Mon, 4 May 2015 16:09:23 -0400 Subject: [PATCH 6/7] remove debugging statement left in registration code --- templates/user/register.mako | 1 - 1 file changed, 1 deletion(-) 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: From 0e49e3a00cd0081781694d4c4f962168e97e0b94 Mon Sep 17 00:00:00 2001 From: John Chilton Date: Fri, 1 May 2015 20:46:34 -0400 Subject: [PATCH 7/7] Fix scary traceback when using task splitting jobs fail. Regression probably introduced in af846d7b3923a5ef5c1ced7723ac6385e119a1df, must not be to harmful since noone has noticed it though. --- lib/galaxy/jobs/__init__.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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