This commit is contained in:
Aysam Guerler
2015-02-06 11:06:29 -05:00
2 changed files with 36 additions and 8 deletions
+13 -4
View File
@@ -27,8 +27,12 @@ def all_control_queues_for_declare(config):
Refactor later to actually persist this somewhere instead of building it repeatedly.
"""
return [Queue('control.%s' % q, galaxy_exchange, routing_key='control') for
q in config.server_names]
possible_uwsgi_queues = []
if config.is_uwsgi:
import uwsgi
possible_uwsgi_queues = [Queue("control.%s.%s" % (config.server_name, wkr['id']), galaxy_exchange, routing_key='control') for wkr in uwsgi.workers()]
return possible_uwsgi_queues + [Queue('c ntrol.%s' % q, galaxy_exchange, routing_key='control') for
q in config.server_names]
def control_queue_from_config(config):
@@ -36,5 +40,10 @@ def control_queue_from_config(config):
Returns a Queue instance with the correct name and routing key for this
galaxy process's config
"""
return Queue("control.%s" % config.server_name, galaxy_exchange,
routing_key='control')
if config.is_uwsgi:
import uwsgi
return Queue("control.%s.%s" % (config.server_name, uwsgi.worker_id()), galaxy_exchange,
routing_key='control')
else:
return Queue("control.%s" % config.server_name, galaxy_exchange,
routing_key='control')
+23 -4
View File
@@ -124,6 +124,9 @@ def app_factory( global_conf, **kwargs ):
def postfork_setup():
from galaxy.app import app
from galaxy.queue_worker import GalaxyQueueWorker
if app.config.is_uwsgi:
import uwsgi
app.config.server_name += ".%s" % uwsgi.worker_id()
app.control_worker = GalaxyQueueWorker(app, galaxy.queues.control_queue_from_config(app.config),
galaxy.queue_worker.control_message_to_task)
app.control_worker.daemon = True
@@ -248,12 +251,28 @@ def populate_api_routes( webapp, app ):
webapp.add_route( '/visualization/show/:visualization_name', controller='visualization', action='render', visualization_name=None )
# Deprecated in favor of POST /api/workflows with 'workflow' in payload.
webapp.mapper.connect( 'import_workflow_deprecated', '/api/workflows/upload', controller='workflows', action='import_new_workflow_deprecated', conditions=dict( method=['POST'] ) )
webapp.mapper.connect( 'workflow_dict', '/api/workflows/{workflow_id}/download', controller='workflows', action='workflow_dict', conditions=dict( method=['GET'] ) )
webapp.mapper.connect( 'import_workflow_deprecated',
'/api/workflows/upload',
controller='workflows',
action='import_new_workflow_deprecated',
conditions=dict( method=['POST'] ) )
webapp.mapper.connect( 'workflow_dict',
'/api/workflows/{workflow_id}/download',
controller='workflows',
action='workflow_dict',
conditions=dict( method=['GET'] ) )
# Preserve the following download route for now for dependent applications -- deprecate at some point
webapp.mapper.connect( 'workflow_dict', '/api/workflows/download/{workflow_id}', controller='workflows', action='workflow_dict', conditions=dict( method=['GET'] ) )
webapp.mapper.connect( 'workflow_dict',
'/api/workflows/download/{workflow_id}',
controller='workflows',
action='workflow_dict',
conditions=dict( method=['GET'] ) )
# Deprecated in favor of POST /api/workflows with shared_workflow_id in payload.
webapp.mapper.connect( 'import_shared_workflow_deprecated', '/api/workflows/import', controller='workflows', action='import_shared_workflow_deprecated', conditions=dict( method=['POST'] ) )
webapp.mapper.connect( 'import_shared_workflow_deprecated',
'/api/workflows/import',
controller='workflows',
action='import_shared_workflow_deprecated',
conditions=dict( method=['POST'] ) )
# API refers to usages and invocations - these mean the same thing but the
# usage routes should be considered deprecated.