Merge branch 'release_15.07'

This commit is contained in:
Dannon Baker
2015-09-23 11:38:31 -04:00
22 changed files with 436 additions and 83 deletions
@@ -118,7 +118,7 @@ return Backbone.View.extend({
.on('mousedown', function(e) { e.preventDefault(); });
// handle text editing event
it.find('#text-content').on('keyup', function(e) {
it.find('#text-content').on('change input', function(e) {
self.model.set('url_paste', $(e.target).val());
self.model.set('file_size', $(e.target).val().length);
});
+13
View File
@@ -54,6 +54,19 @@
galaxy_infrastructure_url is set in galaxy.ini.
-->
<param id="galaxy_url">http://localhost:8080</param>
<!-- AMQP does not guarantee that a published message is received by
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. -->
<!-- <param id="amqp_acknowledge">False</param> -->
<!-- Galaxy reuses Pulsar's persistence_directory parameter (via the
Pulsar client lib) to store a record of received
acknowledgements, and to keep track of messages which have not
been acknowledged. -->
<!-- <param id="persistence_directory">/path/to/dir</param> -->
<!-- Number of seconds to wait for an acknowledgement before
republishing a message. -->
<!-- <param id="amqp_republish_time">30</param> -->
<!-- Pulsar job manager to communicate with (see Pulsar
docs for information on job managers). -->
<!-- <param id="manager">_default_</param> -->
+4
View File
@@ -1036,6 +1036,10 @@ class JobWrapper( object ):
def change_state( self, state, info=False ):
job = self.get_job()
self.sa_session.refresh( job )
if job.state in model.Job.terminal_states:
log.warning( "(%s) Ignoring state change from '%s' to '%s' for job "
"that is already terminal", job.id, job.state, state )
return
for dataset_assoc in job.output_datasets + job.output_library_datasets:
dataset = dataset_assoc.dataset
self.sa_session.refresh( dataset )
+10 -1
View File
@@ -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
+3
View File
@@ -317,6 +317,9 @@ class Job( object, HasJobMetrics, Dictifiable ):
PAUSED = 'paused',
DELETED = 'deleted',
DELETED_NEW = 'deleted_new' )
terminal_states = [ states.OK,
states.ERROR,
states.DELETED ]
# Please include an accessor (get/set pair) for any new columns/members.
def __init__( self ):
self.session_id = None
@@ -222,6 +222,7 @@ class AdminController( BaseUIController, Admin ):
@web.expose
@web.require_admin
def edit_category( self, trans, **kwd ):
'''Handle requests to edit TS category name or description'''
message = escape( kwd.get( 'message', '' ) )
status = kwd.get( 'status', 'done' )
id = kwd.get( 'id', None )
@@ -248,19 +249,19 @@ class AdminController( BaseUIController, Admin ):
else:
category.name = new_name
flush_needed = True
if original_category_description != new_description:
category.description = new_description
if not flush_needed:
flush_needed = True
if flush_needed:
trans.sa_session.add( category )
trans.sa_session.flush()
if original_category_name != new_name:
# Update the Tool Shed's repository registry.
trans.app.repository_registry.edit_category_entry( original_category_name, new_name )
message = "The information has been saved for category '%s'" % escape( category.name )
status = 'done'
return trans.response.send_redirect( web.url_for( controller='admin',
if original_category_description != new_description:
category.description = new_description
if not flush_needed:
flush_needed = True
if flush_needed:
trans.sa_session.add( category )
trans.sa_session.flush()
if original_category_name != new_name:
# Update the Tool Shed's repository registry.
trans.app.repository_registry.edit_category_entry( original_category_name, new_name )
message = "The information has been saved for category '%s'" % escape( category.name )
status = 'done'
return trans.response.send_redirect( web.url_for( controller='admin',
action='manage_categories',
message=message,
status=status ) )
+17 -5
View File
@@ -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):
+107 -9
View File
@@ -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,27 +68,46 @@ 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:
publish_kwds["retry_policy"] = {}
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()
@staticmethod
def __publish_errback(exc, interval, publish_log_prefix=""):
log.error("%sConnection error while publishing: %r", publish_log_prefix, exc, exc_info=1)
log.info("%sRetrying in %s seconds", publish_log_prefix, interval)
@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
self.__start_ack_manager(queue_name)
while check and connection.connected:
try:
connection.drain_events(timeout=self.__timeout)
@@ -90,6 +120,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 +185,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,21 +209,36 @@ 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)
def errback(exc, interval):
return self.__publish_errback(exc, interval, publish_log_prefix)
return PulsarExchange.__publish_errback(exc, interval, publish_log_prefix)
publish_kwds["retry_policy"]["errback"] = errback
else:
publish_kwds = self.__publish_kwds
return publish_kwds
def __publish_errback(self, exc, interval, publish_log_prefix=""):
log.error("%sConnection error while publishing: %r", publish_log_prefix, exc, exc_info=1)
log.info("%sRetrying in %s seconds", publish_log_prefix, interval)
def __publish_log_prefex(self, transaction_uuid=None):
prefix = ""
if transaction_uuid:
@@ -182,3 +272,11 @@ 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):
if self.acks_enabled and queue_name.endswith(ACK_QUEUE_SUFFIX):
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
+16 -1
View File
@@ -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', False):
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
+12 -8
View File
@@ -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}
+4 -1
View File
@@ -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,
+3 -2
View File
@@ -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
+2 -5
View File
@@ -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:
+41 -2
View File
@@ -88,9 +88,14 @@ 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:
@@ -117,9 +122,12 @@ class MessageQueueClientManager(object):
try:
self.exchange.consume("status_update", callback_wrapper, check=self)
except Exception:
log.exception("Exception while handling status update messages, this shouldn't really happen. Handler should be restarted.")
log.exception("Exception while handling status update messages, "
"this shouldn't really happen. Handler should be "
"restarted.")
finally:
log.debug("Leaving Pulsar client status update thread, no additional Pulsar updates will be processed.")
log.debug("Leaving Pulsar client status update thread, no "
"additional Pulsar updates will be processed.")
def ensure_has_status_update_callback(self, callback):
with self.callback_lock:
@@ -136,14 +144,45 @@ 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.")
+26 -10
View File
@@ -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
+4 -2
View File
@@ -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)
+20 -6
View File
@@ -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)
+26 -8
View File
@@ -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):
+110 -6
View File
@@ -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(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
File diff suppressed because one or more lines are too long
File diff suppressed because one or more lines are too long
+1 -1
View File
@@ -81,7 +81,7 @@ def inherit(context):
<form name="login" id="login" action="${form_action}" method="post" >
<div class="form-row">
<label>Username / Email Address:</label>
<input type="text" name="login" value="${login | h}" size="40"/>
<input type="text" name="login" value="${login or ''| h}" size="40"/>
<input type="hidden" name="redirect" value="${redirect | h}" size="40"/>
</div>
<div class="form-row">