From e767c4ed6be023fd7c4a5ae8d5203a28492a8a65 Mon Sep 17 00:00:00 2001 From: John Chilton Date: Mon, 2 Jun 2014 00:13:55 -0500 Subject: [PATCH] Move logic for creating LwrExchange out of MessageQueueClientManager... So can reuse server side when connecting server components to AMQP. --- .../lwr_client/amqp_exchange_factory.py | 32 +++++++++++++++++++ lib/galaxy/jobs/runners/lwr_client/manager.py | 15 ++------- lib/galaxy/jobs/runners/lwr_client/util.py | 16 ---------- 3 files changed, 34 insertions(+), 29 deletions(-) create mode 100644 lib/galaxy/jobs/runners/lwr_client/amqp_exchange_factory.py diff --git a/lib/galaxy/jobs/runners/lwr_client/amqp_exchange_factory.py b/lib/galaxy/jobs/runners/lwr_client/amqp_exchange_factory.py new file mode 100644 index 00000000000..71531c5b56a --- /dev/null +++ b/lib/galaxy/jobs/runners/lwr_client/amqp_exchange_factory.py @@ -0,0 +1,32 @@ +from .amqp_exchange import LwrExchange +from .util import filter_destination_params + + +def get_exchange(url, manager_name, params): + connect_ssl = parse_amqp_connect_ssl_params(params) + exchange_kwds = dict( + manager_name=manager_name, + connect_ssl=connect_ssl, + publish_kwds=parse_amqp_publish_kwds(params) + ) + timeout = params.get('amqp_consumer_timeout', False) + if timeout is not False: + exchange_kwds['timeout'] = timeout + exchange = LwrExchange(url, **exchange_kwds) + return exchange + + +def parse_amqp_connect_ssl_params(params): + ssl_params = filter_destination_params(params, "amqp_connect_ssl_") + if not ssl_params: + return + + ssl = __import__('ssl') + if 'cert_reqs' in ssl_params: + value = ssl_params['cert_reqs'] + ssl_params['cert_reqs'] = getattr(ssl, value.upper()) + return ssl_params + + +def parse_amqp_publish_kwds(params): + return filter_destination_params(params, "amqp_publish_") diff --git a/lib/galaxy/jobs/runners/lwr_client/manager.py b/lib/galaxy/jobs/runners/lwr_client/manager.py index b865949fa8a..ab83a49a1a1 100644 --- a/lib/galaxy/jobs/runners/lwr_client/manager.py +++ b/lib/galaxy/jobs/runners/lwr_client/manager.py @@ -13,10 +13,8 @@ from .interface import LocalLwrInterface from .object_client import ObjectStoreClient from .transport import get_transport from .util import TransferEventManager -from .util import parse_amqp_connect_ssl_params -from .util import parse_amqp_publish_kwds from .destination import url_to_destination_params -from .amqp_exchange import LwrExchange +from .amqp_exchange_factory import get_exchange from logging import getLogger @@ -78,16 +76,7 @@ class MessageQueueClientManager(object): def __init__(self, **kwds): self.url = kwds.get('url') self.manager_name = kwds.get("manager", None) or "_default_" - self.connect_ssl = parse_amqp_connect_ssl_params(kwds) - exchange_kwds = dict( - manager_name=self.manager_name, - connect_ssl=self.connect_ssl, - publish_kwds=parse_amqp_publish_kwds(kwds) - ) - timeout = kwds.get('amqp_consumer_timeout', False) - if timeout is not False: - exchange_kwds['timeout'] = timeout - self.exchange = LwrExchange(self.url, **exchange_kwds) + self.exchange = get_exchange(self.url, self.manager_name, kwds) self.status_cache = {} self.callback_lock = threading.Lock() self.callback_thread = None diff --git a/lib/galaxy/jobs/runners/lwr_client/util.py b/lib/galaxy/jobs/runners/lwr_client/util.py index 2c35f2ab8cc..ce1adfa2cf5 100644 --- a/lib/galaxy/jobs/runners/lwr_client/util.py +++ b/lib/galaxy/jobs/runners/lwr_client/util.py @@ -63,22 +63,6 @@ def directory_files(directory): return contents -def parse_amqp_connect_ssl_params(params): - ssl_params = filter_destination_params(params, "amqp_connect_ssl_") - if not ssl_params: - return - - ssl = __import__('ssl') - if 'cert_reqs' in ssl_params: - value = ssl_params['cert_reqs'] - ssl_params['cert_reqs'] = getattr(ssl, value.upper()) - return ssl_params - - -def parse_amqp_publish_kwds(params): - return filter_destination_params(params, "amqp_publish_") - - def filter_destination_params(destination_params, prefix): destination_params = destination_params or {} return dict([(key[len(prefix):], destination_params[key])