From e951455be37a787a26678e66e607004a9b8ce954 Mon Sep 17 00:00:00 2001 From: Nate Coraor Date: Tue, 8 Sep 2015 12:08:40 -0400 Subject: [PATCH 1/8] With Pulsar, ensure AMQP consumers receive published messages. --- config/job_conf.xml.sample_advanced | 13 +++ lib/galaxy/jobs/runners/pulsar.py | 11 +- lib/pulsar/client/action_mapper.py | 22 +++- lib/pulsar/client/amqp_exchange.py | 103 +++++++++++++++++- lib/pulsar/client/amqp_exchange_factory.py | 17 ++- lib/pulsar/client/client.py | 20 ++-- lib/pulsar/client/config_util.py | 5 +- lib/pulsar/client/decorators.py | 5 +- lib/pulsar/client/interface.py | 7 +- lib/pulsar/client/manager.py | 31 ++++++ lib/pulsar/client/staging/__init__.py | 36 +++++-- lib/pulsar/client/staging/down.py | 6 +- lib/pulsar/client/staging/up.py | 26 +++-- lib/pulsar/client/transport/curl.py | 34 ++++-- lib/pulsar/client/util.py | 116 +++++++++++++++++++-- 15 files changed, 395 insertions(+), 57 deletions(-) diff --git a/config/job_conf.xml.sample_advanced b/config/job_conf.xml.sample_advanced index 444b7d2787e..17a9bf97f6e 100644 --- a/config/job_conf.xml.sample_advanced +++ b/config/job_conf.xml.sample_advanced @@ -54,6 +54,19 @@ galaxy_infrastructure_url is set in galaxy.ini. --> http://localhost:8080 + + + + + + diff --git a/lib/galaxy/jobs/runners/pulsar.py b/lib/galaxy/jobs/runners/pulsar.py index 8067b3528cb..eaf24a0903a 100644 --- a/lib/galaxy/jobs/runners/pulsar.py +++ b/lib/galaxy/jobs/runners/pulsar.py @@ -61,6 +61,14 @@ PULSAR_PARAM_SPECS = dict( map=specs.to_str_or_none, default=None, ), + persistence_directory=dict( + map=specs.to_str_or_none, + default=None, + ), + amqp_acknowledge=dict( + map=specs.to_bool_or_none, + default=None + ), amqp_consumer_timeout=dict( map=lambda val: None if val == "None" else float(val), default=None, @@ -145,7 +153,7 @@ class PulsarJobRunner( AsynchronousJobRunner ): def __init_client_manager( self ): client_manager_kwargs = {} - for kwd in 'manager', 'cache', 'transport': + for kwd in 'manager', 'cache', 'transport', 'persistence_directory': client_manager_kwargs[ kwd ] = self.runner_params[ kwd ] for kwd in self.runner_params.keys(): if kwd.startswith( 'amqp_' ): @@ -611,6 +619,7 @@ class PulsarMQJobRunner( PulsarJobRunner ): # This is a message queue driven runner, don't monitor # just setup required callback. self.client_manager.ensure_has_status_update_callback(self.__async_update) + self.client_manager.ensure_has_ack_consumers() def __async_update( self, full_status ): job_id = None diff --git a/lib/pulsar/client/action_mapper.py b/lib/pulsar/client/action_mapper.py index e81b4cce9fb..7d19eb55054 100644 --- a/lib/pulsar/client/action_mapper.py +++ b/lib/pulsar/client/action_mapper.py @@ -83,7 +83,9 @@ class FileActionMapper(object): ... f.close() ... mock_client = Bunch(default_file_action=default_action, action_config_path=f.name, files_endpoint=None) ... mapper = FileActionMapper(mock_client) - ... mapper = FileActionMapper(config=mapper.to_dict()) # Serialize and deserialize it to make sure still works + ... as_dict = config=mapper.to_dict() + ... # print(as_dict["paths"]) + ... mapper = FileActionMapper(config=as_dict) # Serialize and deserialize it to make sure still works ... unlink(f.name) ... return mapper >>> mapper = mapper_for(default_action='none', config_contents=json_string) @@ -171,7 +173,7 @@ class FileActionMapper(object): ssh_user=self.ssh_user, ssh_port=self.ssh_port, ssh_host=self.ssh_host, - paths=map(lambda m: m.to_dict(), self.mappers) + paths=list(map(lambda m: m.to_dict(), self.mappers)) ) def __client_to_config(self, client): @@ -261,7 +263,7 @@ class BaseAction(object): if self.staging_needed: # To ensure uniqueness, prepend unique prefix to each name prefix = unique_path_prefix(self.path) - for path, name in unstructured_map.iteritems(): + for path, name in unstructured_map.items(): unstructured_map[path] = join(prefix, name) else: path_rewrites = {} @@ -280,6 +282,16 @@ class BaseAction(object): def staging_action_local(self): return self.staging == STAGING_ACTION_LOCAL + def to_dict(self): + return dict(action_type=self.action_type) + + def __str__(self): + as_dict = self.to_dict() + attribute_str = "" + for key, value in as_dict.items(): + attribute_str += "%s=%s" % (key, value) + return "FileAction[%s]" % attribute_str + class NoneAction(BaseAction): """ This action indicates the corresponding path does not require any @@ -446,7 +458,7 @@ class PubkeyAuthenticatedTransferAction(BaseAction): def __serialize_ssh_key(self): f = tempfile.NamedTemporaryFile(delete=False) if self.ssh_key is not None: - f.write(self.ssh_key) + f.write(self.ssh_key.encode("utf-8")) else: raise Exception("SSH_KEY not available") return f.name @@ -644,7 +656,7 @@ MAPPER_CLASS_DICT = dict(map(lambda c: (c.match_type, c), MAPPER_CLASSES)) def mappers_from_dicts(mapper_def_list): - return map(lambda m: _mappper_from_dict(m), mapper_def_list) + return list(map(lambda m: _mappper_from_dict(m), mapper_def_list)) def _mappper_from_dict(mapper_dict): diff --git a/lib/pulsar/client/amqp_exchange.py b/lib/pulsar/client/amqp_exchange.py index ae22c82f5af..65b1e8010b0 100644 --- a/lib/pulsar/client/amqp_exchange.py +++ b/lib/pulsar/client/amqp_exchange.py @@ -3,7 +3,7 @@ import uuid import socket import logging import threading -from time import sleep +from time import sleep, time try: import kombu @@ -26,6 +26,14 @@ DEFAULT_RECONNECT_CONSUMER_WAIT = 1 DEFAULT_HEARTBEAT_WAIT = 1 DEFAULT_HEARTBEAT_JOIN_TIMEOUT = 10 +ACK_QUEUE_SUFFIX = "_ack" +ACK_UUID_KEY = 'acknowledge_uuid' +ACK_QUEUE_KEY = 'acknowledge_queue' +ACK_UUID_RESPONSE_KEY = 'acknowledge_uuid_response' +ACK_FORCE_NOACK_KEY = 'force_noack' +DEFAULT_ACK_MANAGER_SLEEP = 15 +DEFAULT_REPUBLISH_TIME = 30 + class PulsarExchange(object): """ Utility for publishing and consuming structured Pulsar queues using kombu. @@ -47,6 +55,9 @@ class PulsarExchange(object): connect_ssl=None, timeout=DEFAULT_TIMEOUT, publish_kwds={}, + publish_uuid_store=None, + consume_uuid_store=None, + republish_time=DEFAULT_REPUBLISH_TIME, ): """ """ @@ -57,6 +68,7 @@ class PulsarExchange(object): self.__connect_ssl = connect_ssl self.__exchange = kombu.Exchange(DEFAULT_EXCHANGE_NAME, DEFAULT_EXCHANGE_TYPE) self.__timeout = timeout + self.__republish_time = republish_time # Be sure to log message publishing failures. if publish_kwds.get("retry", False): if "retry_policy" not in publish_kwds: @@ -64,20 +76,36 @@ class PulsarExchange(object): if "errback" not in publish_kwds["retry_policy"]: publish_kwds["retry_policy"]["errback"] = self.__publish_errback self.__publish_kwds = publish_kwds + self.publish_uuid_store = publish_uuid_store + self.consume_uuid_store = consume_uuid_store + self.publish_ack_lock = threading.Lock() @property def url(self): return self.__url + @property + def acks_enabled(self): + return self.publish_uuid_store is not None + def consume(self, queue_name, callback, check=True, connection_kwargs={}): queue = self.__queue(queue_name) log.debug("Consuming queue '%s'", queue) + callbacks = [self.__ack_callback] + if callback is not None: + callbacks.append(callback) while check: heartbeat_thread = None try: with self.connection(self.__url, heartbeat=DEFAULT_HEARTBEAT, **connection_kwargs) as connection: - with kombu.Consumer(connection, queues=[queue], callbacks=[callback], accept=['json']): + with kombu.Consumer(connection, queues=[queue], callbacks=callbacks, accept=['json']): heartbeat_thread = self.__start_heartbeat(queue_name, connection) + # Ack manager should sleep before checking for + # repbulishes, but if that changes, need to drain the + # queue once before the ack manager starts doing its + # thing + if self.acks_enabled and queue_name.endswith(ACK_QUEUE_SUFFIX): + ack_manager_thread = self.__start_ack_manager(queue_name) while check and connection.connected: try: connection.drain_events(timeout=self.__timeout) @@ -90,6 +118,41 @@ class PulsarExchange(object): raise log.info("Done consuming queue %s" % queue_name) + def __ack_callback(self, body, message): + if ACK_UUID_KEY in body: + # The consumer of a normal queue has received a message requiring + # acknowledgement + ack_uuid = body[ACK_UUID_KEY] + ack_queue = body[ACK_QUEUE_KEY] + response = {ACK_UUID_RESPONSE_KEY: ack_uuid} + log.debug('Acknowledging UUID %s on queue %s', ack_uuid, ack_queue) + self.publish(ack_queue, response) + if self.consume_uuid_store is None: + log.warning('Received an ack request (UUID: %s, response queue: ' + '%s) but ack UUID persistence is not enabled, check your ' + 'config', ack_uuid, ack_queue) + elif ack_uuid not in self.consume_uuid_store: + # This message has not been seen before, store the uuid so it + # is not operated on more than once + self.consume_uuid_store[ack_uuid] = time() + else: + # This message has been seen before, prevent downstream + # callbacks from processing normally by acknowledging it here, + # still send the ack reply + log.warning('Message with UUID %s on queue %s has already ' + 'been performed, skipping callback', ack_uuid, ack_queue) + message.ack() + elif ACK_UUID_RESPONSE_KEY in body: + # The consumer of an ack queue has received an ack, remove it from the store + ack_uuid = body[ACK_UUID_RESPONSE_KEY] + log.debug('Got acknowledgement for UUID %s, will remove from store', ack_uuid) + try: + with self.publish_ack_lock: + del self.publish_uuid_store[ack_uuid] + except KeyError: + log.warning('Cannot remove UUID %s from store, already removed', ack_uuid) + message.ack() + def __handle_io_error(self, exc, heartbeat_thread): # In testing, errno is None log.warning('Got %s, will retry: %s', exc.__class__.__name__, exc) @@ -120,6 +183,16 @@ class PulsarExchange(object): key = self.__queue_name(name) publish_log_prefix = self.__publish_log_prefex(transaction_uuid) log.debug("%sBegin publishing to key %s", publish_log_prefix, key) + if (self.acks_enabled and not name.endswith(ACK_QUEUE_SUFFIX) + and ACK_FORCE_NOACK_KEY not in payload): + # Publishing a message on a normal queue and it's not a republish + # (or explicitly forced do-not-ack), so add ack keys + ack_uuid = str(transaction_uuid) + ack_queue = name + ACK_QUEUE_SUFFIX + payload[ACK_UUID_KEY] = ack_uuid + payload[ACK_QUEUE_KEY] = ack_queue + self.publish_uuid_store[ack_uuid] = payload + log.debug('Requesting acknowledgement of UUID %s on queue %s', ack_uuid, ack_queue) with self.connection(self.__url) as connection: with pools.producers[connection].acquire() as producer: log.debug("%sHave producer for publishing to key %s", publish_log_prefix, key) @@ -134,6 +207,25 @@ class PulsarExchange(object): ) log.debug("%sPublished to key %s", publish_log_prefix, key) + def ack_manager(self, queue_name): + log.debug('Acknowledgement manager thread alive') + resubmit_queue = queue_name[:-len(ACK_QUEUE_SUFFIX)] + try: + while True: + sleep(DEFAULT_ACK_MANAGER_SLEEP) + with self.publish_ack_lock: + for unack_uuid in self.publish_uuid_store.keys(): + if self.publish_uuid_store.get_time(unack_uuid) < time() - self.__republish_time: + log.debug('UUID %s has not been acknowledged, republishing original message', unack_uuid) + payload = self.publish_uuid_store[unack_uuid] + payload[ACK_FORCE_NOACK_KEY] = True + self.publish(resubmit_queue, payload) + self.publish_uuid_store.set_time(unack_uuid) + except: + log.exception("Problem with acknowledgement manager, leaving ack_manager method in problematic state!") + raise + log.debug('Acknowledgedment manager thread exiting') + def __prepare_publish_kwds(self, publish_log_prefix): if "retry_policy" in self.__publish_kwds: publish_kwds = copy.deepcopy(self.__publish_kwds) @@ -182,3 +274,10 @@ class PulsarExchange(object): thread = threading.Thread(name=thread_name, target=self.heartbeat, args=(connection,)) thread.start() return thread + + def __start_ack_manager(self, queue_name): + thread_name = "acknowledgement-manager-%s" % (self.__queue_name(queue_name)) + thread = threading.Thread(name=thread_name, target=self.ack_manager, args=(queue_name,)) + thread.daemon = True + thread.start() + return thread diff --git a/lib/pulsar/client/amqp_exchange_factory.py b/lib/pulsar/client/amqp_exchange_factory.py index df26f8f5853..4fb67e9b25f 100644 --- a/lib/pulsar/client/amqp_exchange_factory.py +++ b/lib/pulsar/client/amqp_exchange_factory.py @@ -1,5 +1,5 @@ from .amqp_exchange import PulsarExchange -from .util import filter_destination_params +from .util import filter_destination_params, MessageQueueUUIDStore def get_exchange(url, manager_name, params): @@ -9,6 +9,8 @@ def get_exchange(url, manager_name, params): connect_ssl=connect_ssl, publish_kwds=parse_amqp_publish_kwds(params) ) + if params.get('amqp_acknowledge', True): + exchange_kwds.update(parse_ack_kwds(params, manager_name)) timeout = params.get('amqp_consumer_timeout', False) if timeout is not False: exchange_kwds['timeout'] = timeout @@ -39,3 +41,16 @@ def parse_amqp_publish_kwds(params): if retry_policy_params: all_publish_params["retry_policy"] = retry_policy_params return all_publish_params + + +def parse_ack_kwds(params, manager_name): + ack_params = {} + persistence_directory = params.get('persistence_directory', None) + if persistence_directory: + subdirs = ['amqp_ack-%s' % manager_name] + ack_params['publish_uuid_store'] = MessageQueueUUIDStore(persistence_directory, subdirs=subdirs + ['publish']) + ack_params['consume_uuid_store'] = MessageQueueUUIDStore(persistence_directory, subdirs=subdirs + ['consume']) + republish_time = params.get('amqp_ack_republish_time', None) + if republish_time: + ack_params['republish_time'] = int(republish_time) + return ack_params diff --git a/lib/pulsar/client/client.py b/lib/pulsar/client/client.py index 27bb22ab37b..d7f3439bac9 100644 --- a/lib/pulsar/client/client.py +++ b/lib/pulsar/client/client.py @@ -1,12 +1,14 @@ import os -from json import dumps -from json import loads + +from six import string_types from .destination import submit_params from .setup_handler import build as build_setup_handler from .job_directory import RemoteJobDirectory from .decorators import parseJson from .decorators import retry +from .util import json_dumps +from .util import json_loads from .util import copy from .util import ensure_directory from .util import to_base64_json @@ -109,19 +111,19 @@ class JobClient(BaseJobClient): launch_params = dict(command_line=command_line, job_id=self.job_id) submit_params_dict = submit_params(self.destination_params) if submit_params_dict: - launch_params['params'] = dumps(submit_params_dict) + launch_params['params'] = json_dumps(submit_params_dict) if dependencies_description: - launch_params['dependencies_description'] = dumps(dependencies_description.to_dict()) + launch_params['dependencies_description'] = json_dumps(dependencies_description.to_dict()) if env: - launch_params['env'] = dumps(env) + launch_params['env'] = json_dumps(env) if remote_staging: - launch_params['remote_staging'] = dumps(remote_staging) + launch_params['remote_staging'] = json_dumps(remote_staging) if job_config and self.setup_handler.local: # Setup not yet called, job properties were inferred from # destination arguments. Hence, must have Pulsar setup job # before queueing. setup_params = _setup_params_from_job_config(job_config) - launch_params["setup_params"] = dumps(setup_params) + launch_params["setup_params"] = json_dumps(setup_params) return self._raw_execute("submit", launch_params) def full_status(self): @@ -174,10 +176,12 @@ class JobClient(BaseJobClient): # action type == 'message' should either copy or transfer # depending on default not just fallback to transfer. if action_type in ['transfer', 'message']: + if isinstance(contents, string_types): + contents = contents.encode("utf-8") return self._upload_file(args, contents, input_path) elif action_type == 'copy': path_response = self._raw_execute('path', args) - pulsar_path = loads(path_response)['path'] + pulsar_path = json_loads(path_response)['path'] copy(path, pulsar_path) return {'path': pulsar_path} diff --git a/lib/pulsar/client/config_util.py b/lib/pulsar/client/config_util.py index ab5453d0c8f..21103bff746 100644 --- a/lib/pulsar/client/config_util.py +++ b/lib/pulsar/client/config_util.py @@ -1,12 +1,14 @@ """ Generic interface for reading YAML/INI/JSON config files into nested dictionaries. """ +import codecs try: from galaxy import eggs eggs.require('PyYAML') except Exception: # If not in Galaxy, ignore this. pass + try: import yaml except ImportError: @@ -70,8 +72,9 @@ def __read_ini(path): def __read_json(path): + reader = codecs.getreader("utf-8") with open(path, "rb") as f: - return json.load(f) + return json.load(reader(f)) EXT_READERS = { CONFIG_TYPE_JSON: __read_json, diff --git a/lib/pulsar/client/decorators.py b/lib/pulsar/client/decorators.py index 94035f5ba1b..659b5c633d4 100644 --- a/lib/pulsar/client/decorators.py +++ b/lib/pulsar/client/decorators.py @@ -1,5 +1,6 @@ import time -import json + +from .util import json_loads MAX_RETRY_COUNT = 5 RETRY_SLEEP_TIME = 0.1 @@ -10,7 +11,7 @@ class parseJson(object): def __call__(self, func): def replacement(*args, **kwargs): response = func(*args, **kwargs) - return json.loads(response) + return json_loads(response) return replacement diff --git a/lib/pulsar/client/interface.py b/lib/pulsar/client/interface.py index af13a02e4f9..9d625f72c60 100644 --- a/lib/pulsar/client/interface.py +++ b/lib/pulsar/client/interface.py @@ -2,10 +2,7 @@ from abc import ABCMeta from abc import abstractmethod from string import Template -try: - from StringIO import StringIO as BytesIO -except ImportError: - from io import BytesIO +from six import BytesIO try: from six import text_type except ImportError: @@ -145,7 +142,7 @@ class LocalPulsarInterface(PulsarInterface): def __build_body(self, data, input_path): if data is not None: - return BytesIO(data.encode('utf-8')) + return BytesIO(data) elif input_path is not None: return open(input_path, 'rb') else: diff --git a/lib/pulsar/client/manager.py b/lib/pulsar/client/manager.py index f403fd17250..984a4f3d001 100644 --- a/lib/pulsar/client/manager.py +++ b/lib/pulsar/client/manager.py @@ -88,9 +88,13 @@ class MessageQueueClientManager(object): self.status_cache = {} self.callback_lock = threading.Lock() self.callback_thread = None + self.ack_consumer_threads = {} self.active = True def callback_wrapper(self, callback, body, message): + if message.acknowledged: + log.info("Message is already acknowledged (by an upstream callback?), Pulsar client will not handle this message") + return if not self.active: log.debug("Obtained update message for inactive client manager, attempting requeue.") try: @@ -136,14 +140,41 @@ class MessageQueueClientManager(object): thread.start() self.callback_thread = thread + def ack_consumer(self, queue_name): + try: + self.exchange.consume(queue_name + '_ack', None, check=self) + except Exception: + log.exception("Exception while handling %s acknowledgement messages, this shouldn't really happen. Handler should be restarted.", queue_name) + finally: + log.debug("Leaving Pulsar client %s acknowledgement thread, no additional acknowledgements will be processed.", queue_name) + + def ensure_has_ack_consumers(self): + with self.callback_lock: + for name in ('setup', 'kill'): + if name in self.ack_consumer_threads: + return + + run = functools.partial(self.ack_consumer, name) + thread = threading.Thread( + name="pulsar_client_%s_%s_ack" % (self.manager_name, name), + target=run + ) + thread.daemon = False # Lets not interrupt processing of this. + thread.start() + self.ack_consumer_threads[name] = thread + def shutdown(self, ensure_cleanup=False): self.active = False if ensure_cleanup: self.callback_thread.join() + for v in self.ack_consumer_threads.values(): + v.join() def __nonzero__(self): return self.active + __bool__ = __nonzero__ # Both needed Py2 v 3 + def get_client(self, destination_params, job_id, **kwargs): if job_id is None: raise Exception("Cannot generate Pulsar client for empty job_id.") diff --git a/lib/pulsar/client/staging/__init__.py b/lib/pulsar/client/staging/__init__.py index 7b0a47c3c9b..c121bf28be7 100644 --- a/lib/pulsar/client/staging/__init__.py +++ b/lib/pulsar/client/staging/__init__.py @@ -1,3 +1,4 @@ +import re from os.path import basename from os.path import join from os.path import dirname @@ -6,6 +7,7 @@ from os import sep from ..util import PathHelper COMMAND_VERSION_FILENAME = "COMMAND_VERSION" +DEFAULT_DYNAMIC_COLLECTION_PATTERN = [r"primary_.*|galaxy.json|metadata_.*|dataset_\d+\.dat|__instrument_.*|dataset_\d+_files.+"] class ClientJobDescription(object): @@ -50,12 +52,12 @@ class ClientJobDescription(object): def __init__( self, - tool, command_line, - config_files, - input_files, - client_outputs, - working_directory, + tool=None, + config_files=[], + input_files=[], + client_outputs=None, + working_directory=None, # More sensible default? dependencies_description=None, env=[], arbitrary_files=None, @@ -65,7 +67,7 @@ class ClientJobDescription(object): self.command_line = command_line self.config_files = config_files self.input_files = input_files - self.client_outputs = client_outputs + self.client_outputs = client_outputs or ClientOutputs() self.working_directory = working_directory self.dependencies_description = dependencies_description self.env = env @@ -95,18 +97,28 @@ class ClientOutputs(object): runner client. """ - def __init__(self, working_directory, output_files, work_dir_outputs=None, version_file=None): + def __init__( + self, + working_directory=None, + output_files=[], + work_dir_outputs=None, + version_file=None, + dynamic_outputs=None + ): self.working_directory = working_directory - self.work_dir_outputs = work_dir_outputs - self.output_files = output_files + self.work_dir_outputs = work_dir_outputs or [] + self.output_files = output_files or [] self.version_file = version_file + self.dynamic_outputs = dynamic_outputs or DEFAULT_DYNAMIC_COLLECTION_PATTERN + self.__dynamic_patterns = list(map(re.compile, self.dynamic_outputs)) def to_dict(self): return dict( working_directory=self.working_directory, work_dir_outputs=self.work_dir_outputs, output_files=self.output_files, - version_file=self.version_file + version_file=self.version_file, + dynamic_outputs=self.dynamic_outputs, ) @staticmethod @@ -116,8 +128,12 @@ class ClientOutputs(object): work_dir_outputs=config_dict.get('work_dir_outputs'), output_files=config_dict.get('output_files'), version_file=config_dict.get('version_file'), + dynamic_outputs=config_dict.get('dynamic_outputs'), ) + def dynamic_match(self, filename): + return any(map(lambda pattern: pattern.match(filename), self.__dynamic_patterns)) + class PulsarOutputs(object): """ Abstraction describing the output files PRODUCED by the remote Pulsar diff --git a/lib/pulsar/client/staging/down.py b/lib/pulsar/client/staging/down.py index cfb901d5eb0..37f2e4965c3 100644 --- a/lib/pulsar/client/staging/down.py +++ b/lib/pulsar/client/staging/down.py @@ -94,7 +94,7 @@ class ResultsCollector(object): if output_generated: self._attempt_collect_output('output', output_file) - for galaxy_path, pulsar in self.pulsar_outputs.output_extras(output_file).iteritems(): + for galaxy_path, pulsar in self.pulsar_outputs.output_extras(output_file).items(): self._attempt_collect_output('output', path=galaxy_path, name=pulsar) # else not output generated, do not attempt download. @@ -110,7 +110,8 @@ class ResultsCollector(object): for name in self.working_directory_contents: if name in self.downloaded_working_directory_files: continue - if COPY_FROM_WORKING_DIRECTORY_PATTERN.match(name): + if self.client_outputs.dynamic_match(name): + log.debug("collecting dynamic output %s" % name) output_file = join(working_directory, self.pulsar_outputs.path_helper.local_name(name)) if self._attempt_collect_output(output_type='output_workdir', path=output_file, name=name): self.downloaded_working_directory_files.append(name) @@ -128,6 +129,7 @@ class ResultsCollector(object): return collected def _collect_output(self, output_type, action, name): + log.info("collecting output %s with action %s" % (name, action)) return self.output_collector.collect_output(self, output_type, action, name) diff --git a/lib/pulsar/client/staging/up.py b/lib/pulsar/client/staging/up.py index 6fd01136aca..46ebe319290 100644 --- a/lib/pulsar/client/staging/up.py +++ b/lib/pulsar/client/staging/up.py @@ -65,9 +65,14 @@ class FileStager(object): self.config_files = client_job_description.config_files self.input_files = client_job_description.input_files self.output_files = client_job_description.output_files - self.tool_id = client_job_description.tool.id - self.tool_version = client_job_description.tool.version - self.tool_dir = abspath(client_job_description.tool.tool_dir) + if client_job_description.tool is not None: + self.tool_id = client_job_description.tool.id + self.tool_version = client_job_description.tool.version + self.tool_dir = abspath(client_job_description.tool.tool_dir) + else: + self.tool_id = None + self.tool_version = None + self.tool_dir = None self.working_directory = client_job_description.working_directory self.version_file = client_job_description.version_file self.arbitrary_files = client_job_description.arbitrary_files @@ -142,7 +147,7 @@ class FileStager(object): for path in paths: if path not in referenced_arbitrary_path_mappers: referenced_arbitrary_path_mappers[path] = mapper - for path, mapper in referenced_arbitrary_path_mappers.iteritems(): + for path, mapper in referenced_arbitrary_path_mappers.items(): action = self.action_mapper.action(path, path_type.UNSTRUCTURED, mapper) unstructured_map = action.unstructured_map(self.path_helper) self.arbitrary_files.update(unstructured_map) @@ -152,7 +157,7 @@ class FileStager(object): self.transfer_tracker.handle_transfer(referenced_tool_file, path_type.TOOL) def __upload_arbitrary_files(self): - for path, name in self.arbitrary_files.iteritems(): + for path, name in self.arbitrary_files.items(): self.transfer_tracker.handle_transfer(path, path_type.UNSTRUCTURED, name=name) def __upload_input_files(self): @@ -180,11 +185,17 @@ class FileStager(object): def __upload_working_directory_files(self): # Task manager stages files into working directory, these need to be # uploaded if present. - working_directory_files = listdir(self.working_directory) if exists(self.working_directory) else [] + working_directory_files = self.__working_directory_files() for working_directory_file in working_directory_files: path = join(self.working_directory, working_directory_file) self.transfer_tracker.handle_transfer(path, path_type.WORKDIR) + def __working_directory_files(self): + if self.working_directory and exists(self.working_directory): + return listdir(self.working_directory) + else: + return [] + def __initialize_version_file_rename(self): version_file = self.version_file if version_file: @@ -298,6 +309,9 @@ class JobInputs(object): Full path to directory to search. """ + if directory is None: + return [] + pattern = r"(%s%s\S+)" % (directory, sep) return self.find_pattern_references(pattern) diff --git a/lib/pulsar/client/transport/curl.py b/lib/pulsar/client/transport/curl.py index 983b9786f0b..d2f0f50cf52 100644 --- a/lib/pulsar/client/transport/curl.py +++ b/lib/pulsar/client/transport/curl.py @@ -1,7 +1,13 @@ +import logging + try: - from cStringIO import StringIO + from galaxy import eggs + eggs.require("six") except ImportError: - from io import StringIO + pass + +from six import string_types +from six import BytesIO try: from pycurl import Curl, HTTP_CODE @@ -19,6 +25,8 @@ NO_SUCH_FILE_MESSAGE = "Attempt to post file %s to URL %s, but file does not exi POST_FAILED_MESSAGE = "Failed to post_file properly for url %s, remote server returned status code of %s." GET_FAILED_MESSAGE = "Failed to get_file properly for url %s, remote server returned status code of %s." +log = logging.getLogger(__name__) + class PycurlTransport(object): @@ -36,7 +44,7 @@ class PycurlTransport(object): c.setopt(c.INFILESIZE, filesize) if data: c.setopt(c.POST, 1) - if type(data).__name__ == 'unicode': + if isinstance(data, string_types): data = data.encode('UTF-8') c.setopt(c.POSTFIELDS, data) c.perform() @@ -62,21 +70,31 @@ def post_file(url, path): def get_file(url, path): - buf = _open_output(path) + if path and os.path.exists(path): + buf = _open_output(path, 'ab') + size = os.path.getsize(path) + success_codes = (200, 206) + else: + buf = _open_output(path) + size = 0 + success_codes = (200,) try: c = _new_curl_object_for_url(url) c.setopt(c.WRITEFUNCTION, buf.write) + if size > 0: + log.info('transfer of %s will resume at %s bytes', url, size) + c.setopt(c.RESUME_FROM, size) c.perform() - status_code = c.getinfo(HTTP_CODE) - if int(status_code) != 200: + status_code = int(c.getinfo(HTTP_CODE)) + if status_code not in success_codes: message = GET_FAILED_MESSAGE % (url, status_code) raise Exception(message) finally: buf.close() -def _open_output(output_path): - return open(output_path, 'wb') if output_path else StringIO() +def _open_output(output_path, mode='wb'): + return open(output_path, mode) if output_path else BytesIO() def _new_curl_object_for_url(url): diff --git a/lib/pulsar/client/util.py b/lib/pulsar/client/util.py index 04908db6256..ae382c178f4 100644 --- a/lib/pulsar/client/util.py +++ b/lib/pulsar/client/util.py @@ -1,19 +1,46 @@ +from functools import wraps from threading import Lock, Event from weakref import WeakValueDictionary from os import walk from os import curdir +from os import listdir +from os import makedirs +from os import unlink from os.path import relpath from os.path import join +from os.path import abspath +from os.path import exists +from errno import ENOENT, EEXIST import os.path import hashlib import shutil import json -import base64 +import sys + +from six import binary_type + +# Variant of base64 compat layer inspired by BSD code from Bcfg2 +# https://github.com/Bcfg2/bcfg2/blob/maint/src/lib/Bcfg2/Compat.py +if sys.version_info >= (3, 0): + from base64 import b64encode as _b64encode, b64decode as _b64decode + + @wraps(_b64encode) + def b64encode(val, **kwargs): + try: + return _b64encode(val, **kwargs) + except TypeError: + return _b64encode(val.encode('UTF-8'), **kwargs).decode('UTF-8') + + @wraps(_b64decode) + def b64decode(val, **kwargs): + return _b64decode(val.encode('UTF-8'), **kwargs).decode('UTF-8') +else: + from base64 import b64encode, b64decode def unique_path_prefix(path): m = hashlib.md5() - m.update(path) + m.update(path.encode('utf-8')) return m.hexdigest() @@ -75,15 +102,17 @@ def filter_destination_params(destination_params, prefix): def to_base64_json(data): """ - >>> x = from_base64_json(to_base64_json(dict(a=5))) - >>> x["a"] + >>> enc = to_base64_json(dict(a=5)) + >>> dec = from_base64_json(enc) + >>> dec["a"] 5 """ - return base64.b64encode(json.dumps(data)) + dumped = json_dumps(data) + return b64encode(dumped) def from_base64_json(data): - return json.loads(base64.b64decode(data)) + return json.loads(b64decode(data)) class PathHelper(object): @@ -175,3 +204,78 @@ class EventHolder(object): def fail(self): self.failed = True + + +def json_loads(obj): + if isinstance(obj, binary_type): + obj = obj.decode("utf-8") + return json.loads(obj) + + +def json_dumps(obj): + if isinstance(obj, binary_type): + obj = obj.decode("utf-8") + return json.dumps(obj, cls=ClientJsonEncoder) + + +class ClientJsonEncoder(json.JSONEncoder): + + def default(self, obj): + if isinstance(obj, binary_type): + return obj.decode("utf-8") + return json.JSONEncoder.default(self, obj) + + +class MessageQueueUUIDStore(object): + """Persistent dict-like object for persisting message queue UUIDs that are + awaiting acknowledgement or that have been operated on. + """ + + def __init__(self, persistence_directory, subdirs=None): + if subdirs is None: + subdirs = ['acknowledge_uuids'] + self.__store = abspath(join(persistence_directory, *subdirs)) + try: + makedirs(self.__store) + except (OSError, IOError) as exc: + if exc.errno != EEXIST: + raise + + def __path(self, item): + return join(self.__store, item) + + def __contains__(self, item): + return exists(self.__path(item)) + + def __setitem__(self, key, value): + open(self.__path(key), 'w').write(json.dumps(value)) + + def __getitem__(self, key): + return json.loads(open(self.__path(key)).read()) + + def __delitem__(self, key): + try: + unlink(self.__path(key)) + except (OSError, IOError) as exc: + if exc.errno == ENOENT: + raise KeyError(key) + raise + + def keys(self): + return iter(os.listdir(self.__store)) + + def get_time(self, key): + try: + return os.stat(self.__path(key)).st_mtime + except (OSError, IOError) as exc: + if exc.errno == ENOENT: + raise KeyError(key) + raise + + def set_time(self, key): + try: + os.utime(self.__path(key), None) + except (OSError, IOError) as exc: + if exc.errno == ENOENT: + raise KeyError(key) + raise From 5bae90156fa7998b575981ac452db21f34683085 Mon Sep 17 00:00:00 2001 From: Nate Coraor Date: Thu, 10 Sep 2015 13:55:41 -0400 Subject: [PATCH 2/8] Fix a couple flake8 errors including a supposed PEP-8 violation that is actually a PEP-8-suggested acceptable option of how to handle the offending situation. =P --- lib/pulsar/client/amqp_exchange.py | 4 ++-- lib/pulsar/client/util.py | 2 +- 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/lib/pulsar/client/amqp_exchange.py b/lib/pulsar/client/amqp_exchange.py index 65b1e8010b0..bde9b36e381 100644 --- a/lib/pulsar/client/amqp_exchange.py +++ b/lib/pulsar/client/amqp_exchange.py @@ -105,7 +105,7 @@ class PulsarExchange(object): # queue once before the ack manager starts doing its # thing if self.acks_enabled and queue_name.endswith(ACK_QUEUE_SUFFIX): - ack_manager_thread = self.__start_ack_manager(queue_name) + self.__start_ack_manager(queue_name) while check and connection.connected: try: connection.drain_events(timeout=self.__timeout) @@ -184,7 +184,7 @@ class PulsarExchange(object): publish_log_prefix = self.__publish_log_prefex(transaction_uuid) log.debug("%sBegin publishing to key %s", publish_log_prefix, key) if (self.acks_enabled and not name.endswith(ACK_QUEUE_SUFFIX) - and ACK_FORCE_NOACK_KEY not in payload): + and ACK_FORCE_NOACK_KEY not in payload): # Publishing a message on a normal queue and it's not a republish # (or explicitly forced do-not-ack), so add ack keys ack_uuid = str(transaction_uuid) diff --git a/lib/pulsar/client/util.py b/lib/pulsar/client/util.py index ae382c178f4..35696befdc1 100644 --- a/lib/pulsar/client/util.py +++ b/lib/pulsar/client/util.py @@ -262,7 +262,7 @@ class MessageQueueUUIDStore(object): raise def keys(self): - return iter(os.listdir(self.__store)) + return iter(listdir(self.__store)) def get_time(self, key): try: From 8778da43cb641a925c421b188aabe053be56e313 Mon Sep 17 00:00:00 2001 From: Nate Coraor Date: Thu, 10 Sep 2015 15:13:40 -0400 Subject: [PATCH 3/8] Disable amqp_acknowledge by default, additional fixes from Pulsar's pickier flake8 checks. --- config/job_conf.xml.sample_advanced | 2 +- lib/pulsar/client/amqp_exchange.py | 18 +++++++++--------- lib/pulsar/client/amqp_exchange_factory.py | 2 +- lib/pulsar/client/manager.py | 18 +++++++++++++----- 4 files changed, 24 insertions(+), 16 deletions(-) diff --git a/config/job_conf.xml.sample_advanced b/config/job_conf.xml.sample_advanced index 17a9bf97f6e..08df8099af7 100644 --- a/config/job_conf.xml.sample_advanced +++ b/config/job_conf.xml.sample_advanced @@ -58,7 +58,7 @@ the AMQP server, so Galaxy/Pulsar can request that the consumer acknowledge messages and will resend them if acknowledgement is not received after a configurable timeout. --> - +