mirror of
https://github.com/galaxyproject/galaxy.git
synced 2026-09-24 16:30:27 +08:00
Move logic for creating LwrExchange out of MessageQueueClientManager...
So can reuse server side when connecting server components to AMQP.
This commit is contained in:
@@ -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_")
|
||||
@@ -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
|
||||
|
||||
@@ -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])
|
||||
|
||||
Reference in New Issue
Block a user