Merge pull request #2052 from jmchilton/pulsar-deps

Pulsar-As-A-Dependency
This commit is contained in:
Nate Coraor
2016-04-01 17:14:34 -04:00
35 changed files with 10 additions and 4177 deletions
+1 -1
View File
@@ -5,4 +5,4 @@ set -e
flake8 --exclude `paste -sd, .ci/flake8_blacklist.txt` .
# Apply stricter rules for the directories shared with Pulsar
flake8 --ignore= --max-line-length=150 lib/galaxy/jobs/runners/util/ lib/pulsar/
flake8 --ignore= --max-line-length=150 lib/galaxy/jobs/runners/util/
-1
View File
@@ -1,6 +1,5 @@
lib/galaxy/util/
lib/galaxy/jobs/runners/util/
lib/pulsar/
lib/galaxy/tools/parser/
lib/galaxy/tools/lint.py
lib/galaxy/tools/lint_util.py
-1
View File
@@ -10,5 +10,4 @@ lib
log_tempfile
mimeparse
psyco_full
pulsar
tool_shed
-132
View File
@@ -1,132 +0,0 @@
pulsar.client package
=====================
.. automodule:: pulsar.client
:members:
:undoc-members:
:show-inheritance:
Subpackages
-----------
.. toctree::
pulsar.client.staging
pulsar.client.transport
Submodules
----------
pulsar.client.action_mapper module
----------------------------------
.. automodule:: pulsar.client.action_mapper
:members:
:undoc-members:
:show-inheritance:
pulsar.client.amqp_exchange module
----------------------------------
.. automodule:: pulsar.client.amqp_exchange
:members:
:undoc-members:
:show-inheritance:
pulsar.client.amqp_exchange_factory module
------------------------------------------
.. automodule:: pulsar.client.amqp_exchange_factory
:members:
:undoc-members:
:show-inheritance:
pulsar.client.client module
---------------------------
.. automodule:: pulsar.client.client
:members:
:undoc-members:
:show-inheritance:
pulsar.client.config_util module
--------------------------------
.. automodule:: pulsar.client.config_util
:members:
:undoc-members:
:show-inheritance:
pulsar.client.decorators module
-------------------------------
.. automodule:: pulsar.client.decorators
:members:
:undoc-members:
:show-inheritance:
pulsar.client.destination module
--------------------------------
.. automodule:: pulsar.client.destination
:members:
:undoc-members:
:show-inheritance:
pulsar.client.interface module
------------------------------
.. automodule:: pulsar.client.interface
:members:
:undoc-members:
:show-inheritance:
pulsar.client.job_directory module
----------------------------------
.. automodule:: pulsar.client.job_directory
:members:
:undoc-members:
:show-inheritance:
pulsar.client.manager module
----------------------------
.. automodule:: pulsar.client.manager
:members:
:undoc-members:
:show-inheritance:
pulsar.client.object_client module
----------------------------------
.. automodule:: pulsar.client.object_client
:members:
:undoc-members:
:show-inheritance:
pulsar.client.path_mapper module
--------------------------------
.. automodule:: pulsar.client.path_mapper
:members:
:undoc-members:
:show-inheritance:
pulsar.client.setup_handler module
----------------------------------
.. automodule:: pulsar.client.setup_handler
:members:
:undoc-members:
:show-inheritance:
pulsar.client.util module
-------------------------
.. automodule:: pulsar.client.util
:members:
:undoc-members:
:show-inheritance:
-28
View File
@@ -1,28 +0,0 @@
pulsar.client.staging package
=============================
.. automodule:: pulsar.client.staging
:members:
:undoc-members:
:show-inheritance:
Submodules
----------
pulsar.client.staging.down module
---------------------------------
.. automodule:: pulsar.client.staging.down
:members:
:undoc-members:
:show-inheritance:
pulsar.client.staging.up module
-------------------------------
.. automodule:: pulsar.client.staging.up
:members:
:undoc-members:
:show-inheritance:
@@ -1,52 +0,0 @@
pulsar.client.transport package
===============================
.. automodule:: pulsar.client.transport
:members:
:undoc-members:
:show-inheritance:
Submodules
----------
pulsar.client.transport.curl module
-----------------------------------
.. automodule:: pulsar.client.transport.curl
:members:
:undoc-members:
:show-inheritance:
pulsar.client.transport.poster module
-------------------------------------
.. automodule:: pulsar.client.transport.poster
:members:
:undoc-members:
:show-inheritance:
pulsar.client.transport.requests module
---------------------------------------
.. automodule:: pulsar.client.transport.requests
:members:
:undoc-members:
:show-inheritance:
pulsar.client.transport.ssh module
----------------------------------
.. automodule:: pulsar.client.transport.ssh
:members:
:undoc-members:
:show-inheritance:
pulsar.client.transport.standard module
---------------------------------------
.. automodule:: pulsar.client.transport.standard
:members:
:undoc-members:
:show-inheritance:
-15
View File
@@ -1,15 +0,0 @@
pulsar package
==============
.. automodule:: pulsar
:members:
:undoc-members:
:show-inheritance:
Subpackages
-----------
.. toctree::
pulsar.client
@@ -41,13 +41,16 @@ kombu==3.0.30
amqp==1.4.8
anyjson==0.3.3
# Pulsar requirements
psutil==4.1.0
pulsar-galaxy-lib==0.7.0.dev1
# sqlalchemy-migrate and dependencies
sqlalchemy-migrate==0.10.0
decorator==4.0.2
Tempita==0.5.3dev
sqlparse==0.1.16
pbr==1.8.0
# six is also a Pulsar client dep
six==1.9.0
Parsley==1.3
nose==1.3.7
+4
View File
@@ -39,6 +39,10 @@ requests
# kombu and dependencies
kombu
# Pulsar requirements
psutil
pulsar-galaxy-lib==0.7.0.dev1
# sqlalchemy-migrate and dependencies
sqlalchemy-migrate
decorator
+1 -1
View File
@@ -217,7 +217,7 @@ class PulsarJobRunner( AsynchronousJobRunner ):
command_line=command_line,
input_files=self.get_input_files(job_wrapper),
client_outputs=self.__client_outputs(client, job_wrapper),
working_directory=job_wrapper.working_directory,
working_directory=job_wrapper.tool_working_directory,
tool=job_wrapper.tool,
config_files=job_wrapper.extra_filenames,
dependencies_description=dependencies_description,
View File
-62
View File
@@ -1,62 +0,0 @@
"""
pulsar client
=================
This module contains logic for interfacing with an external Pulsar server.
------------------
Configuring Galaxy
------------------
Galaxy job runners are configured in Galaxy's ``job_conf.xml`` file. See ``job_conf.xml.sample_advanced``
in your Galaxy code base or on
`Bitbucket <https://bitbucket.org/galaxy/galaxy-dist/src/tip/config/job_conf.xml.sample_advanced?at=default>`_
for information on how to configure Galaxy to interact with the Pulsar.
Galaxy also supports an older, less rich configuration of job runners directly
in its main ``galaxy.ini`` file. The following section describes how to
configure Galaxy to communicate with the Pulsar in this legacy mode.
Legacy
------
A Galaxy tool can be configured to be executed remotely via Pulsar by
adding a line to the ``galaxy.ini`` file under the ``galaxy:tool_runners``
section with the format::
<tool_id> = pulsar://http://<pulsar_host>:<pulsar_port>
As an example, if a host named remotehost is running the Pulsar server
application on port ``8913``, then the tool with id ``test_tool`` can
be configured to run remotely on remotehost by adding the following
line to ``galaxy.ini``::
test_tool = pulsar://http://remotehost:8913
Remember this must be added after the ``[galaxy:tool_runners]`` header
in the ``galaxy.ini`` file.
"""
from .staging.down import finish_job
from .staging.up import submit_job
from .staging import ClientJobDescription
from .staging import PulsarOutputs
from .staging import ClientOutputs
from .client import OutputNotFoundException
from .manager import build_client_manager
from .destination import url_to_destination_params
from .path_mapper import PathMapper
__all__ = [
'build_client_manager',
'OutputNotFoundException',
'url_to_destination_params',
'finish_job',
'submit_job',
'ClientJobDescription',
'PulsarOutputs',
'ClientOutputs',
'PathMapper',
]
-708
View File
@@ -1,708 +0,0 @@
from contextlib import contextmanager
from json import load
from os import makedirs
from os import unlink
from os.path import exists
from os.path import abspath
from os.path import dirname
from os.path import join
from os.path import basename
from os.path import sep
import fnmatch
from re import compile
from re import escape
import galaxy.util
from galaxy.util.bunch import Bunch
from .config_util import read_file
from .util import directory_files
from .util import unique_path_prefix
from .transport import get_file
from .transport import post_file
from .transport import rsync_get_file, scp_get_file
from .transport import rsync_post_file, scp_post_file
import tempfile
DEFAULT_MAPPED_ACTION = 'transfer' # Not really clear to me what this should be, exception?
DEFAULT_PATH_MAPPER_TYPE = 'prefix'
STAGING_ACTION_REMOTE = "remote"
STAGING_ACTION_LOCAL = "local"
STAGING_ACTION_NONE = None
STAGING_ACTION_DEFAULT = "default"
# Poor man's enum.
path_type = Bunch(
# Galaxy input datasets and extra files.
INPUT="input",
# Galaxy config and param files.
CONFIG="config",
# Files from tool's tool_dir (for now just wrapper if available).
TOOL="tool",
# Input work dir files - e.g. metadata files, task-split input files, etc..
WORKDIR="workdir",
# Galaxy output datasets in their final home.
OUTPUT="output",
# Galaxy from_work_dir output paths and other files (e.g. galaxy.json)
OUTPUT_WORKDIR="output_workdir",
# Other fixed tool parameter paths (likely coming from tool data, but not
# nessecarily). Not sure this is the best name...
UNSTRUCTURED="unstructured",
)
ACTION_DEFAULT_PATH_TYPES = [
path_type.INPUT,
path_type.CONFIG,
path_type.TOOL,
path_type.WORKDIR,
path_type.OUTPUT,
path_type.OUTPUT_WORKDIR,
]
ALL_PATH_TYPES = ACTION_DEFAULT_PATH_TYPES + [path_type.UNSTRUCTURED]
MISSING_FILES_ENDPOINT_ERROR = "Attempted to use remote_transfer action without defining a files_endpoint."
MISSING_SSH_KEY_ERROR = "Attempt to use file transfer action requiring an SSH key without specifying a ssh_key."
class FileActionMapper(object):
"""
Objects of this class define how paths are mapped to actions.
>>> json_string = r'''{"paths": [ \
{"path": "/opt/galaxy", "action": "none"}, \
{"path": "/galaxy/data", "action": "transfer"}, \
{"path": "/cool/bamfiles/**/*.bam", "action": "copy", "match_type": "glob"}, \
{"path": ".*/dataset_\\\\d+.dat", "action": "copy", "match_type": "regex"} \
]}'''
>>> from tempfile import NamedTemporaryFile
>>> from os import unlink
>>> def mapper_for(default_action, config_contents):
... f = NamedTemporaryFile(delete=False)
... f.write(config_contents.encode('UTF-8'))
... f.close()
... mock_client = Bunch(default_file_action=default_action, action_config_path=f.name, files_endpoint=None)
... mapper = FileActionMapper(mock_client)
... 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)
>>> # Test first config line above, implicit path prefix mapper
>>> action = mapper.action('/opt/galaxy/tools/filters/catWrapper.py', 'input')
>>> action.action_type == u'none'
True
>>> action.staging_needed
False
>>> # Test another (2nd) mapper, this one with a different action
>>> action = mapper.action('/galaxy/data/files/000/dataset_1.dat', 'input')
>>> action.action_type == u'transfer'
True
>>> action.staging_needed
True
>>> # Always at least copy work_dir outputs.
>>> action = mapper.action('/opt/galaxy/database/working_directory/45.sh', 'workdir')
>>> action.action_type == u'copy'
True
>>> action.staging_needed
True
>>> # Test glob mapper (matching test)
>>> mapper.action('/cool/bamfiles/projectABC/study1/patient3.bam', 'input').action_type == u'copy'
True
>>> # Test glob mapper (non-matching test)
>>> mapper.action('/cool/bamfiles/projectABC/study1/patient3.bam.bai', 'input').action_type == u'none'
True
>>> # Regex mapper test.
>>> mapper.action('/old/galaxy/data/dataset_10245.dat', 'input').action_type == u'copy'
True
>>> # Doesn't map unstructured paths by default
>>> mapper.action('/old/galaxy/data/dataset_10245.dat', 'unstructured').action_type == u'none'
True
>>> input_only_mapper = mapper_for(default_action="none", config_contents=r'''{"paths": [ \
{"path": "/", "action": "transfer", "path_types": "input"} \
] }''')
>>> input_only_mapper.action('/dataset_1.dat', 'input').action_type == u'transfer'
True
>>> input_only_mapper.action('/dataset_1.dat', 'output').action_type == u'none'
True
>>> unstructured_mapper = mapper_for(default_action="none", config_contents=r'''{"paths": [ \
{"path": "/", "action": "transfer", "path_types": "*any*"} \
] }''')
>>> unstructured_mapper.action('/old/galaxy/data/dataset_10245.dat', 'unstructured').action_type == u'transfer'
True
"""
def __init__(self, client=None, config=None):
if config is None and client is None:
message = "FileActionMapper must be constructed from either a client or a config dictionary."
raise Exception(message)
if config is None:
config = self.__client_to_config(client)
self.default_action = config.get("default_action", "transfer")
self.ssh_key = config.get("ssh_key", None)
self.ssh_user = config.get("ssh_user", None)
self.ssh_host = config.get("ssh_host", None)
self.ssh_port = config.get("ssh_port", None)
self.mappers = mappers_from_dicts(config.get("paths", []))
self.files_endpoint = config.get("files_endpoint", None)
def action(self, path, type, mapper=None):
mapper = self.__find_mapper(path, type, mapper)
action_class = self.__action_class(path, type, mapper)
file_lister = DEFAULT_FILE_LISTER
action_kwds = {}
if mapper:
file_lister = mapper.file_lister
action_kwds = mapper.action_kwds
action = action_class(path, file_lister=file_lister, **action_kwds)
self.__process_action(action, type)
return action
def unstructured_mappers(self):
""" Return mappers that will map 'unstructured' files (i.e. go beyond
mapping inputs, outputs, and config files).
"""
return filter(lambda m: path_type.UNSTRUCTURED in m.path_types, self.mappers)
def to_dict(self):
return dict(
default_action=self.default_action,
files_endpoint=self.files_endpoint,
ssh_key=self.ssh_key,
ssh_user=self.ssh_user,
ssh_port=self.ssh_port,
ssh_host=self.ssh_host,
paths=list(map(lambda m: m.to_dict(), self.mappers))
)
def __client_to_config(self, client):
action_config_path = client.action_config_path
if action_config_path:
config = read_file(action_config_path)
else:
config = dict()
config["default_action"] = client.default_file_action
config["files_endpoint"] = client.files_endpoint
for attr in ['ssh_key', 'ssh_user', 'ssh_port', 'ssh_host']:
if hasattr(client, attr):
config[attr] = getattr(client, attr)
return config
def __load_action_config(self, path):
config = load(open(path, 'rb'))
self.mappers = mappers_from_dicts(config.get('paths', []))
def __find_mapper(self, path, type, mapper=None):
if not mapper:
normalized_path = abspath(path)
for query_mapper in self.mappers:
if query_mapper.matches(normalized_path, type):
mapper = query_mapper
break
return mapper
def __action_class(self, path, type, mapper):
action_type = self.default_action if type in ACTION_DEFAULT_PATH_TYPES else "none"
if mapper:
action_type = mapper.action_type
if type in ["workdir", "output_workdir"] and action_type == "none":
# We are changing the working_directory relative to what
# Galaxy would use, these need to be copied over.
action_type = "copy"
action_class = actions.get(action_type, None)
if action_class is None:
message_template = "Unknown action_type encountered %s while trying to map path %s"
message_args = (action_type, path)
raise Exception(message_template % message_args)
return action_class
def __process_action(self, action, file_type):
""" Extension point to populate extra action information after an
action has been created.
"""
if getattr(action, "inject_url", False):
self.__inject_url(action, file_type)
if getattr(action, "inject_ssh_properties", False):
self.__inject_ssh_properties(action)
def __inject_url(self, action, file_type):
url_base = self.files_endpoint
if not url_base:
raise Exception(MISSING_FILES_ENDPOINT_ERROR)
if "?" not in url_base:
url_base = "%s?" % url_base
# TODO: URL encode path.
url = "%s&path=%s&file_type=%s" % (url_base, action.path, file_type)
action.url = url
def __inject_ssh_properties(self, action):
for attr in ["ssh_key", "ssh_host", "ssh_port", "ssh_user"]:
action_attr = getattr(action, attr)
if action_attr == UNSET_ACTION_KWD:
client_default_attr = getattr(self, attr, None)
setattr(action, attr, client_default_attr)
if action.ssh_key is None:
raise Exception(MISSING_SSH_KEY_ERROR)
REQUIRED_ACTION_KWD = object()
UNSET_ACTION_KWD = "__UNSET__"
class BaseAction(object):
action_spec = {}
def __init__(self, path, file_lister=None):
self.path = path
self.file_lister = file_lister or DEFAULT_FILE_LISTER
def unstructured_map(self, path_helper):
unstructured_map = self.file_lister.unstructured_map(self.path)
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.items():
unstructured_map[path] = join(prefix, name)
else:
path_rewrites = {}
for path in unstructured_map:
rewrite = self.path_rewrite(path_helper, path)
if rewrite:
path_rewrites[path] = rewrite
unstructured_map = path_rewrites
return unstructured_map
@property
def staging_needed(self):
return self.staging != STAGING_ACTION_NONE
@property
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
additional action. This should indicate paths that are available both on
the Pulsar client (i.e. Galaxy server) and remote Pulsar server with the same
paths. """
action_type = "none"
staging = STAGING_ACTION_NONE
def to_dict(self):
return dict(path=self.path, action_type=self.action_type)
@classmethod
def from_dict(cls, action_dict):
return NoneAction(path=action_dict["path"])
def path_rewrite(self, path_helper, path=None):
return None
class RewriteAction(BaseAction):
""" This actin indicates the Pulsar server should simply rewrite the path
to the specified file.
"""
action_spec = dict(
source_directory=REQUIRED_ACTION_KWD,
destination_directory=REQUIRED_ACTION_KWD
)
action_type = "rewrite"
staging = STAGING_ACTION_NONE
def __init__(self, path, file_lister=None, source_directory=None, destination_directory=None):
self.path = path
self.file_lister = file_lister or DEFAULT_FILE_LISTER
self.source_directory = source_directory
self.destination_directory = destination_directory
def to_dict(self):
return dict(
path=self.path,
action_type=self.action_type,
source_directory=self.source_directory,
destination_directory=self.destination_directory,
)
@classmethod
def from_dict(cls, action_dict):
return RewriteAction(
path=action_dict["path"],
source_directory=action_dict["source_directory"],
destination_directory=action_dict["destination_directory"],
)
def path_rewrite(self, path_helper, path=None):
if not path:
path = self.path
new_path = path_helper.from_posix_with_new_base(self.path, self.source_directory, self.destination_directory)
return None if new_path == self.path else new_path
class TransferAction(BaseAction):
""" This actions indicates that the Pulsar client should initiate an HTTP
transfer of the corresponding path to the remote Pulsar server before
launching the job. """
action_type = "transfer"
staging = STAGING_ACTION_LOCAL
class CopyAction(BaseAction):
""" This action indicates that the Pulsar client should execute a file system
copy of the corresponding path to the Pulsar staging directory prior to
launching the corresponding job. """
action_type = "copy"
staging = STAGING_ACTION_LOCAL
class RemoteCopyAction(BaseAction):
""" This action indicates the Pulsar server should copy the file before
execution via direct file system copy. This is like a CopyAction, but
it indicates the action should occur on the Pulsar server instead of on
the client.
"""
action_type = "remote_copy"
staging = STAGING_ACTION_REMOTE
def to_dict(self):
return dict(path=self.path, action_type=self.action_type)
@classmethod
def from_dict(cls, action_dict):
return RemoteCopyAction(path=action_dict["path"])
def write_to_path(self, path):
galaxy.util.copy_to_path(open(self.path, "rb"), path)
def write_from_path(self, pulsar_path):
destination = self.path
parent_directory = dirname(destination)
if not exists(parent_directory):
makedirs(parent_directory)
with open(pulsar_path, "rb") as f:
galaxy.util.copy_to_path(f, destination)
class RemoteTransferAction(BaseAction):
""" This action indicates the Pulsar server should copy the file before
execution via direct file system copy. This is like a CopyAction, but
it indicates the action should occur on the Pulsar server instead of on
the client.
"""
inject_url = True
action_type = "remote_transfer"
staging = STAGING_ACTION_REMOTE
def __init__(self, path, file_lister=None, url=None):
super(RemoteTransferAction, self).__init__(path, file_lister=file_lister)
self.url = url
def to_dict(self):
return dict(path=self.path, action_type=self.action_type, url=self.url)
@classmethod
def from_dict(cls, action_dict):
return RemoteTransferAction(path=action_dict["path"], url=action_dict["url"])
def write_to_path(self, path):
get_file(self.url, path)
def write_from_path(self, pulsar_path):
post_file(self.url, pulsar_path)
class PubkeyAuthenticatedTransferAction(BaseAction):
"""Base class for file transfers requiring an SSH public/private key
"""
inject_ssh_properties = True
action_spec = dict(
ssh_key=UNSET_ACTION_KWD,
ssh_user=UNSET_ACTION_KWD,
ssh_host=UNSET_ACTION_KWD,
ssh_port=UNSET_ACTION_KWD,
)
staging = STAGING_ACTION_REMOTE
def __init__(self, path, file_lister=None, ssh_user=UNSET_ACTION_KWD,
ssh_host=UNSET_ACTION_KWD, ssh_port=UNSET_ACTION_KWD, ssh_key=UNSET_ACTION_KWD):
super(PubkeyAuthenticatedTransferAction, self).__init__(path, file_lister=file_lister)
self.ssh_user = ssh_user
self.ssh_host = ssh_host
self.ssh_port = ssh_port
self.ssh_key = ssh_key
def to_dict(self):
return dict(path=self.path, action_type=self.action_type,
ssh_user=self.ssh_user, ssh_host=self.ssh_host,
ssh_port=self.ssh_port)
@contextmanager
def _serialized_key(self):
key_file = self.__serialize_ssh_key()
yield key_file
self.__cleanup_ssh_key(key_file)
def __serialize_ssh_key(self):
f = tempfile.NamedTemporaryFile(delete=False)
if self.ssh_key is not None:
f.write(self.ssh_key.encode("utf-8"))
else:
raise Exception("SSH_KEY not available")
return f.name
def __cleanup_ssh_key(self, keyfile):
if exists(keyfile):
unlink(keyfile)
class RsyncTransferAction(PubkeyAuthenticatedTransferAction):
action_type = "remote_rsync_transfer"
@classmethod
def from_dict(cls, action_dict):
return RsyncTransferAction(path=action_dict["path"],
ssh_user=action_dict["ssh_user"],
ssh_host=action_dict["ssh_host"],
ssh_port=action_dict["ssh_port"],
ssh_key=action_dict["ssh_key"])
def write_to_path(self, path):
with self._serialized_key() as key_file:
rsync_get_file(self.path, path, self.ssh_user, self.ssh_host,
self.ssh_port, key_file)
def write_from_path(self, pulsar_path):
with self._serialized_key() as key_file:
rsync_post_file(pulsar_path, self.path, self.ssh_user,
self.ssh_host, self.ssh_port, key_file)
class ScpTransferAction(PubkeyAuthenticatedTransferAction):
action_type = "remote_scp_transfer"
@classmethod
def from_dict(cls, action_dict):
return ScpTransferAction(path=action_dict["path"],
ssh_user=action_dict["ssh_user"],
ssh_host=action_dict["ssh_host"],
ssh_port=action_dict["ssh_port"],
ssh_key=action_dict["ssh_key"])
def write_to_path(self, path):
with self._serialized_key() as key_file:
scp_get_file(self.path, path, self.ssh_user, self.ssh_host,
self.ssh_port, key_file)
def write_from_path(self, pulsar_path):
with self._serialized_key() as key_file:
scp_post_file(pulsar_path, self.path, self.ssh_user, self.ssh_host,
self.ssh_port, key_file)
class MessageAction(object):
""" Sort of pseudo action describing "files" store in memory and
transferred via message (HTTP, Python-call, MQ, etc...)
"""
action_type = "message"
staging = STAGING_ACTION_DEFAULT
def __init__(self, contents, client=None):
self.contents = contents
self.client = client
@property
def staging_needed(self):
return True
@property
def staging_action_local(self):
# Ekkk, cannot be called if created through from_dict.
# Shouldn't be a problem the way it is used - but is an
# object design problem.
return self.client.prefer_local_staging
def to_dict(self):
return dict(contents=self.contents, action_type=MessageAction.action_type)
@classmethod
def from_dict(cls, action_dict):
return MessageAction(contents=action_dict["contents"])
def write_to_path(self, path):
open(path, "w").write(self.contents)
DICTIFIABLE_ACTION_CLASSES = [RemoteCopyAction, RemoteTransferAction, MessageAction, RsyncTransferAction, ScpTransferAction]
def from_dict(action_dict):
action_type = action_dict.get("action_type", None)
target_class = None
for action_class in DICTIFIABLE_ACTION_CLASSES:
if action_type == action_class.action_type:
target_class = action_class
if not target_class:
message = "Failed to recover action from dictionary - invalid action type specified %s." % action_type
raise Exception(message)
return target_class.from_dict(action_dict)
class BasePathMapper(object):
def __init__(self, config):
action_type = config.get('action', DEFAULT_MAPPED_ACTION)
action_class = actions.get(action_type, None)
action_kwds = action_class.action_spec.copy()
for key, value in action_kwds.items():
if key in config:
action_kwds[key] = config[key]
elif value is REQUIRED_ACTION_KWD:
message_template = "action_type %s requires key word argument %s"
message = message_template % (action_type, key)
raise Exception(message)
else:
action_kwds[key] = value
self.action_type = action_type
self.action_kwds = action_kwds
path_types_str = config.get('path_types', "*defaults*")
path_types_str = path_types_str.replace("*defaults*", ",".join(ACTION_DEFAULT_PATH_TYPES))
path_types_str = path_types_str.replace("*any*", ",".join(ALL_PATH_TYPES))
self.path_types = path_types_str.split(",")
self.file_lister = FileLister(config)
def matches(self, path, path_type):
path_type_matches = path_type in self.path_types
return path_type_matches and self._path_matches(path)
def _extend_base_dict(self, **kwds):
base_dict = dict(
action=self.action_type,
path_types=",".join(self.path_types),
match_type=self.match_type
)
base_dict.update(self.file_lister.to_dict())
base_dict.update(self.action_kwds)
base_dict.update(**kwds)
return base_dict
class PrefixPathMapper(BasePathMapper):
match_type = 'prefix'
def __init__(self, config):
super(PrefixPathMapper, self).__init__(config)
self.prefix_path = abspath(config['path'])
def _path_matches(self, path):
return path.startswith(self.prefix_path)
def to_pattern(self):
pattern_str = "(%s%s[^\s,\"\']+)" % (escape(self.prefix_path), escape(sep))
return compile(pattern_str)
def to_dict(self):
return self._extend_base_dict(path=self.prefix_path)
class GlobPathMapper(BasePathMapper):
match_type = 'glob'
def __init__(self, config):
super(GlobPathMapper, self).__init__(config)
self.glob_path = config['path']
def _path_matches(self, path):
return fnmatch.fnmatch(path, self.glob_path)
def to_pattern(self):
return compile(fnmatch.translate(self.glob_path))
def to_dict(self):
return self._extend_base_dict(path=self.glob_path)
class RegexPathMapper(BasePathMapper):
match_type = 'regex'
def __init__(self, config):
super(RegexPathMapper, self).__init__(config)
self.pattern_raw = config['path']
self.pattern = compile(self.pattern_raw)
def _path_matches(self, path):
return self.pattern.match(path) is not None
def to_pattern(self):
return self.pattern
def to_dict(self):
return self._extend_base_dict(path=self.pattern_raw)
MAPPER_CLASSES = [PrefixPathMapper, GlobPathMapper, RegexPathMapper]
MAPPER_CLASS_DICT = dict(map(lambda c: (c.match_type, c), MAPPER_CLASSES))
def mappers_from_dicts(mapper_def_list):
return list(map(lambda m: _mappper_from_dict(m), mapper_def_list))
def _mappper_from_dict(mapper_dict):
map_type = mapper_dict.get('match_type', DEFAULT_PATH_MAPPER_TYPE)
return MAPPER_CLASS_DICT[map_type](mapper_dict)
class FileLister(object):
def __init__(self, config):
self.depth = int(config.get("depth", "0"))
def to_dict(self):
return dict(
depth=self.depth
)
def unstructured_map(self, path):
depth = self.depth
if self.depth == 0:
return {path: basename(path)}
else:
while depth > 0:
path = dirname(path)
depth -= 1
return dict([(join(path, f), f) for f in directory_files(path)])
DEFAULT_FILE_LISTER = FileLister(dict(depth=0))
ACTION_CLASSES = [
NoneAction,
RewriteAction,
TransferAction,
CopyAction,
RemoteCopyAction,
RemoteTransferAction,
RsyncTransferAction,
ScpTransferAction,
]
actions = dict([(clazz.action_type, clazz) for clazz in ACTION_CLASSES])
__all__ = [
'FileActionMapper',
'path_type',
'from_dict',
'MessageAction',
'RemoteTransferAction', # For testing
]
-286
View File
@@ -1,286 +0,0 @@
import copy
import uuid
import socket
import logging
import threading
from time import sleep, time
try:
import kombu
from kombu import pools
except ImportError:
kombu = None
log = logging.getLogger(__name__)
KOMBU_UNAVAILABLE = "Attempting to bind to AMQP message queue, but kombu dependency unavailable"
DEFAULT_EXCHANGE_NAME = "pulsar"
DEFAULT_EXCHANGE_TYPE = "direct"
# Set timeout to periodically give up looking and check if polling should end.
DEFAULT_TIMEOUT = 0.2
DEFAULT_HEARTBEAT = 580
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_SUBMIT_QUEUE_KEY = 'acknowledge_submit_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.
This is shared between the server and client - an exchange should be setup
for each manager (or in the case of the client, each manager one wished to
communicate with.)
Each Pulsar manager is defined solely by name in the scheme, so only one Pulsar
should target each AMQP endpoint or care should be taken that unique
manager names are used across Pulsar servers targetting same AMQP endpoint -
and in particular only one such Pulsar should define an default manager with
name _default_.
"""
def __init__(
self,
url,
manager_name,
connect_ssl=None,
timeout=DEFAULT_TIMEOUT,
publish_kwds={},
publish_uuid_store=None,
consume_uuid_store=None,
republish_time=DEFAULT_REPUBLISH_TIME,
):
"""
"""
if not kombu:
raise Exception(KOMBU_UNAVAILABLE)
self.__url = url
self.__manager_name = manager_name
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"] = {}
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()
# 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.ack_manager_thread = self.__start_ack_manager()
@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=callbacks, accept=['json']):
heartbeat_thread = self.__start_heartbeat(queue_name, connection)
while check and connection.connected:
try:
connection.drain_events(timeout=self.__timeout)
except socket.timeout:
pass
except (IOError, socket.error) as exc:
self.__handle_io_error(exc, heartbeat_thread)
except BaseException:
log.exception("Problem consuming queue, consumer quitting in problematic fashion!")
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)
try:
if heartbeat_thread:
heartbeat_thread.join(DEFAULT_HEARTBEAT_JOIN_TIMEOUT)
except Exception:
log.exception("Failed to join heartbeat thread, this is bad?")
try:
sleep(DEFAULT_RECONNECT_CONSUMER_WAIT)
except Exception:
log.exception("Interrupted sleep while waiting to reconnect to message queue, may restart unless problems encountered.")
def heartbeat(self, connection):
log.debug('AMQP heartbeat thread alive')
try:
while connection.connected:
connection.heartbeat_check()
sleep(DEFAULT_HEARTBEAT_WAIT)
except BaseException:
log.exception("Problem with heartbeat, leaving heartbeat method in problematic state!")
raise
log.debug('AMQP heartbeat thread exiting')
def publish(self, name, payload):
# Consider optionally disabling if throughput becomes main concern.
transaction_uuid = uuid.uuid1()
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
payload[ACK_SUBMIT_QUEUE_KEY] = name
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)
publish_kwds = self.__prepare_publish_kwds(publish_log_prefix)
producer.publish(
payload,
serializer='json',
exchange=self.__exchange,
declare=[self.__exchange],
routing_key=key,
**publish_kwds
)
log.debug("%sPublished to key %s", publish_log_prefix, key)
def ack_manager(self):
log.debug('Acknowledgement manager thread alive')
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:
payload = self.publish_uuid_store[unack_uuid]
payload[ACK_FORCE_NOACK_KEY] = True
resubmit_queue = payload[ACK_SUBMIT_QUEUE_KEY]
log.debug('UUID %s has not been acknowledged, '
'republishing original message on queue %s',
unack_uuid, resubmit_queue)
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 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_log_prefex(self, transaction_uuid=None):
prefix = ""
if transaction_uuid:
prefix = "[publish:%s] " % str(transaction_uuid)
return prefix
def connection(self, connection_string, **kwargs):
if "ssl" not in kwargs:
kwargs["ssl"] = self.__connect_ssl
return kombu.Connection(connection_string, **kwargs)
def __queue(self, name):
queue_name = self.__queue_name(name)
queue = kombu.Queue(queue_name, self.__exchange, routing_key=queue_name)
return queue
def __queue_name(self, name):
key_prefix = self.__key_prefix()
queue_name = '%s_%s' % (key_prefix, name)
return queue_name
def __key_prefix(self):
if self.__manager_name == "_default_":
key_prefix = "pulsar_"
else:
key_prefix = "pulsar_%s_" % self.__manager_name
return key_prefix
def __start_heartbeat(self, queue_name, connection):
thread_name = "consume-heartbeat-%s" % (self.__queue_name(queue_name))
thread = threading.Thread(name=thread_name, target=self.heartbeat, args=(connection,))
thread.start()
return thread
def __start_ack_manager(self):
if self.acks_enabled:
thread_name = "acknowledgement-manager"
thread = threading.Thread(name=thread_name, target=self.ack_manager)
thread.daemon = True
thread.start()
return thread
@@ -1,56 +0,0 @@
from .amqp_exchange import PulsarExchange
from .util import filter_destination_params, MessageQueueUUIDStore
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)
)
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
exchange = PulsarExchange(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):
all_publish_params = filter_destination_params(params, "amqp_publish_")
retry_policy_params = {}
for key in all_publish_params.keys():
if key.startswith("retry_"):
value = all_publish_params[key]
retry_policy_params[key[len("retry_"):]] = value
del all_publish_params[key]
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
-400
View File
@@ -1,400 +0,0 @@
import os
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
from .action_mapper import (
path_type,
actions,
)
import logging
log = logging.getLogger(__name__)
CACHE_WAIT_SECONDS = 3
class OutputNotFoundException(Exception):
def __init__(self, path):
self.path = path
def __str__(self):
return "No remote output found for path %s" % self.path
class BaseJobClient(object):
def __init__(self, destination_params, job_id):
destination_params = destination_params or {}
self.destination_params = destination_params
self.job_id = job_id
if "jobs_directory" in destination_params:
staging_directory = destination_params["jobs_directory"]
sep = destination_params.get("remote_sep", os.sep)
job_directory = RemoteJobDirectory(
remote_staging_directory=staging_directory,
remote_id=job_id,
remote_sep=sep,
)
else:
job_directory = None
for attr in ["ssh_key", "ssh_user", "ssh_host", "ssh_port"]:
setattr(self, attr, destination_params.get(attr, None))
self.env = destination_params.get("env", [])
self.files_endpoint = destination_params.get("files_endpoint", None)
self.job_directory = job_directory
default_file_action = self.destination_params.get("default_file_action", "transfer")
if default_file_action not in actions:
raise Exception("Unknown Pulsar default file action type %s" % default_file_action)
self.default_file_action = default_file_action
self.action_config_path = self.destination_params.get("file_action_config", None)
self.setup_handler = build_setup_handler(self, destination_params)
def setup(self, tool_id=None, tool_version=None):
"""
Setup remote Pulsar server to run this job.
"""
setup_args = {"job_id": self.job_id}
if tool_id:
setup_args["tool_id"] = tool_id
if tool_version:
setup_args["tool_version"] = tool_version
return self.setup_handler.setup(**setup_args)
@property
def prefer_local_staging(self):
# If doing a job directory is defined, calculate paths here and stage
# remotely.
return self.job_directory is None
class JobClient(BaseJobClient):
"""
Objects of this client class perform low-level communication with a remote Pulsar server.
**Parameters**
destination_params : dict or str
connection parameters, either url with dict containing url (and optionally `private_token`).
job_id : str
Galaxy job/task id.
"""
def __init__(self, destination_params, job_id, job_manager_interface):
super(JobClient, self).__init__(destination_params, job_id)
self.job_manager_interface = job_manager_interface
def launch(self, command_line, dependencies_description=None, env=[], remote_staging=[], job_config=None):
"""
Queue up the execution of the supplied `command_line` on the remote
server. Called launch for historical reasons, should be renamed to
enqueue or something like that.
**Parameters**
command_line : str
Command to execute.
"""
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'] = json_dumps(submit_params_dict)
if dependencies_description:
launch_params['dependencies_description'] = json_dumps(dependencies_description.to_dict())
if env:
launch_params['env'] = json_dumps(env)
if 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"] = json_dumps(setup_params)
return self._raw_execute("submit", launch_params)
def full_status(self):
""" Return a dictionary summarizing final state of job.
"""
return self.raw_check_complete()
def kill(self):
"""
Cancel remote job, either removing from the queue or killing it.
"""
return self._raw_execute("cancel", {"job_id": self.job_id})
@retry()
@parseJson()
def raw_check_complete(self):
"""
Get check_complete response from the remote server.
"""
check_complete_response = self._raw_execute("status", {"job_id": self.job_id})
return check_complete_response
def get_status(self):
check_complete_response = self.raw_check_complete()
# Older Pulsar instances won't set status so use 'complete', at some
# point drop backward compatibility.
status = check_complete_response.get("status", None)
return status
def clean(self):
"""
Cleanup the remote job.
"""
self._raw_execute("clean", {"job_id": self.job_id})
@parseJson()
def remote_setup(self, **setup_args):
"""
Setup remote Pulsar server to run this job.
"""
return self._raw_execute("setup", setup_args)
def put_file(self, path, input_type, name=None, contents=None, action_type='transfer'):
if not name:
name = os.path.basename(path)
args = {"job_id": self.job_id, "name": name, "type": input_type}
input_path = path
if contents:
input_path = None
# 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 = json_loads(path_response)['path']
copy(path, pulsar_path)
return {'path': pulsar_path}
def fetch_output(self, path, name, working_directory, action_type, output_type):
"""
Fetch (transfer, copy, etc...) an output from the remote Pulsar server.
**Parameters**
path : str
Local path of the dataset.
name : str
Remote name of file (i.e. path relative to remote staging output
or working directory).
working_directory : str
Local working_directory for the job.
action_type : str
Where to find file on Pulsar (output_workdir or output). legacy is also
an option in this case Pulsar is asked for location - this will only be
used if targetting an older Pulsar server that didn't return statuses
allowing this to be inferred.
"""
if output_type == 'output_workdir':
self._fetch_work_dir_output(name, working_directory, path, action_type=action_type)
elif output_type == 'output':
self._fetch_output(path=path, name=name, action_type=action_type)
else:
raise Exception("Unknown output_type %s" % output_type)
def _raw_execute(self, command, args={}, data=None, input_path=None, output_path=None):
return self.job_manager_interface.execute(command, args, data, input_path, output_path)
def _fetch_output(self, path, name=None, check_exists_remotely=False, action_type='transfer'):
if not name:
# Extra files will send in the path.
name = os.path.basename(path)
self.__populate_output_path(name, path, action_type)
def _fetch_work_dir_output(self, name, working_directory, output_path, action_type='transfer'):
ensure_directory(output_path)
if action_type == 'transfer':
self.__raw_download_output(name, self.job_id, path_type.OUTPUT_WORKDIR, output_path)
else: # Even if action is none - Pulsar has a different work_dir so this needs to be copied.
pulsar_path = self._output_path(name, self.job_id, path_type.OUTPUT_WORKDIR)['path']
copy(pulsar_path, output_path)
def __populate_output_path(self, name, output_path, action_type):
ensure_directory(output_path)
if action_type == 'transfer':
self.__raw_download_output(name, self.job_id, path_type.OUTPUT, output_path)
elif action_type == 'copy':
pulsar_path = self._output_path(name, self.job_id, path_type.OUTPUT)['path']
copy(pulsar_path, output_path)
@parseJson()
def _upload_file(self, args, contents, input_path):
return self._raw_execute("upload_file", args, contents, input_path)
@parseJson()
def _output_path(self, name, job_id, output_type):
return self._raw_execute("path",
{"name": name,
"job_id": self.job_id,
"type": output_type})
@retry()
def __raw_download_output(self, name, job_id, output_type, output_path):
output_params = {
"name": name,
"job_id": self.job_id,
"type": output_type
}
self._raw_execute("download_output", output_params, output_path=output_path)
class BaseMessageJobClient(BaseJobClient):
def __init__(self, destination_params, job_id, client_manager):
super(BaseMessageJobClient, self).__init__(destination_params, job_id)
if not self.job_directory:
error_message = "Message-queue based Pulsar client requires destination define a remote job_directory to stage files into."
raise Exception(error_message)
self.client_manager = client_manager
def clean(self):
del self.client_manager.status_cache[self.job_id]
def full_status(self):
full_status = self.client_manager.status_cache.get(self.job_id, None)
if full_status is None:
raise Exception("full_status() called before a final status was properly cached with cilent manager.")
return full_status
def _build_setup_message(self, command_line, dependencies_description, env, remote_staging, job_config):
"""
"""
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['submit_params'] = submit_params_dict
if dependencies_description:
launch_params['dependencies_description'] = dependencies_description.to_dict()
if env:
launch_params['env'] = env
if remote_staging:
launch_params['remote_staging'] = remote_staging
launch_params['remote_staging']['ssh_key'] = self.ssh_key
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"] = setup_params
return launch_params
class MessageJobClient(BaseMessageJobClient):
def launch(self, command_line, dependencies_description=None, env=[], remote_staging=[], job_config=None):
"""
"""
launch_params = self._build_setup_message(
command_line,
dependencies_description=dependencies_description,
env=env,
remote_staging=remote_staging,
job_config=job_config
)
response = self.client_manager.exchange.publish("setup", launch_params)
log.info("Job published to setup message queue.")
return response
def kill(self):
self.client_manager.exchange.publish("kill", dict(job_id=self.job_id))
class MessageCLIJobClient(BaseMessageJobClient):
def __init__(self, destination_params, job_id, client_manager, shell):
super(MessageCLIJobClient, self).__init__(destination_params, job_id, client_manager)
self.remote_pulsar_path = destination_params["remote_pulsar_path"]
self.shell = shell
def launch(self, command_line, dependencies_description=None, env=[], remote_staging=[], job_config=None):
"""
"""
launch_params = self._build_setup_message(
command_line,
dependencies_description=dependencies_description,
env=env,
remote_staging=remote_staging,
job_config=job_config
)
base64_message = to_base64_json(launch_params)
submit_command = os.path.join(self.remote_pulsar_path, "scripts", "submit.bash")
# TODO: Allow configuration of manager, app, and ini path...
self.shell.execute("nohup %s --base64 %s &" % (submit_command, base64_message))
def kill(self):
# TODO
pass
class InputCachingJobClient(JobClient):
"""
Beta client that cache's staged files to prevent duplication.
"""
def __init__(self, destination_params, job_id, job_manager_interface, client_cacher):
super(InputCachingJobClient, self).__init__(destination_params, job_id, job_manager_interface)
self.client_cacher = client_cacher
@parseJson()
def _upload_file(self, args, contents, input_path):
action = "upload_file"
if contents:
input_path = None
return self._raw_execute(action, args, contents, input_path)
else:
event_holder = self.client_cacher.acquire_event(input_path)
cache_required = self.cache_required(input_path)
if cache_required:
self.client_cacher.queue_transfer(self, input_path)
while not event_holder.failed:
available = self.file_available(input_path)
if available['ready']:
token = available['token']
args["cache_token"] = token
return self._raw_execute(action, args)
event_holder.event.wait(30)
if event_holder.failed:
raise Exception("Failed to transfer file %s" % input_path)
@parseJson()
def cache_required(self, path):
return self._raw_execute("cache_required", {"path": path})
@parseJson()
def cache_insert(self, path):
return self._raw_execute("cache_insert", {"path": path}, None, path)
@parseJson()
def file_available(self, path):
return self._raw_execute("file_available", {"path": path})
def _setup_params_from_job_config(job_config):
job_id = job_config.get("job_id", None)
tool_id = job_config.get("tool_id", None)
tool_version = job_config.get("tool_version", None)
return dict(
job_id=job_id,
tool_id=tool_id,
tool_version=tool_version
)
-77
View File
@@ -1,77 +0,0 @@
""" Generic interface for reading YAML/INI/JSON config files into nested dictionaries.
"""
import codecs
try:
import yaml
except ImportError:
yaml = None
try:
from ConfigParser import ConfigParser
except ImportError:
from configparser import ConfigParser
import json
CONFIG_TYPE_JSON = "json"
CONFIG_TYPE_YAML = "yaml"
CONFIG_TYPE_INI = "ini"
DEFAULT_CONFIG_TYPE = CONFIG_TYPE_YAML
JSON_EXTS = [".json"]
YAML_EXTS = [".yaml", ".yml"]
INI_EXTS = [".ini"]
EXT_MAP = {
CONFIG_TYPE_JSON: JSON_EXTS,
CONFIG_TYPE_YAML: YAML_EXTS,
CONFIG_TYPE_INI: INI_EXTS,
}
def read_file(path, type=None, default_type=DEFAULT_CONFIG_TYPE):
if path is None:
raise ValueError("Undefined path supplied.")
config_type = __find_type(path, type, default_type)
return EXT_READERS[config_type](path)
def __find_type(path, explicit_type, default_type):
if explicit_type:
return explicit_type
for config_type, config_exts in EXT_MAP.items():
for ext in config_exts:
if path.endswith(ext):
return config_type
return default_type
def __read_yaml(path):
if yaml is None:
raise ImportError("Attempting to read YAML configuration file - but PyYAML dependency unavailable.")
with open(path, "rb") as f:
return yaml.load(f)
def __read_ini(path):
config = ConfigParser()
config.read(path)
return config._sections
def __read_json(path):
reader = codecs.getreader("utf-8")
with open(path, "rb") as f:
return json.load(reader(f))
EXT_READERS = {
CONFIG_TYPE_JSON: __read_json,
CONFIG_TYPE_YAML: __read_yaml,
CONFIG_TYPE_INI: __read_ini,
}
-36
View File
@@ -1,36 +0,0 @@
import time
from .util import json_loads
MAX_RETRY_COUNT = 5
RETRY_SLEEP_TIME = 0.1
class parseJson(object):
def __call__(self, func):
def replacement(*args, **kwargs):
response = func(*args, **kwargs)
return json_loads(response)
return replacement
class retry(object):
def __call__(self, func):
def replacement(*args, **kwargs):
max_count = MAX_RETRY_COUNT
count = 0
while True:
count += 1
try:
return func(*args, **kwargs)
except:
if count >= max_count:
raise
else:
time.sleep(RETRY_SLEEP_TIME)
continue
return replacement
-58
View File
@@ -1,58 +0,0 @@
from re import match
from .util import filter_destination_params
SUBMIT_PREFIX = "submit_"
def url_to_destination_params(url):
"""Convert a legacy runner URL to a job destination
>>> params_simple = url_to_destination_params("http://localhost:8913/")
>>> params_simple["url"]
'http://localhost:8913/'
>>> params_simple["private_token"] is None
True
>>> advanced_url = "https://1234x@example.com:8914/managers/longqueue"
>>> params_advanced = url_to_destination_params(advanced_url)
>>> params_advanced["url"]
'https://example.com:8914/managers/longqueue/'
>>> params_advanced["private_token"]
'1234x'
>>> runner_url = "pulsar://http://localhost:8913/"
>>> runner_params = url_to_destination_params(runner_url)
>>> runner_params['url']
'http://localhost:8913/'
"""
if url.startswith("pulsar://"):
url = url[len("pulsar://"):]
if not url.endswith("/"):
url += "/"
# Check for private token embedded in the URL. A URL of the form
# https://moo@cow:8913 will try to contact https://cow:8913
# with a private key of moo
private_token_format = "https?://(.*)@.*/?"
private_token_match = match(private_token_format, url)
private_token = None
if private_token_match:
private_token = private_token_match.group(1)
url = url.replace("%s@" % private_token, '', 1)
destination_args = {"url": url,
"private_token": private_token}
return destination_args
def submit_params(destination_params):
"""
>>> destination_params = {"private_token": "12345", "submit_native_specification": "-q batch"}
>>> result = submit_params(destination_params)
>>> result
{'native_specification': '-q batch'}
"""
return filter_destination_params(destination_params, SUBMIT_PREFIX)
-147
View File
@@ -1,147 +0,0 @@
from abc import ABCMeta
from abc import abstractmethod
from string import Template
from six import BytesIO
from six import text_type
try:
from urllib import urlencode
except ImportError:
from urllib.parse import urlencode
class PulsarInterface(object):
"""
Abstract base class describes how synchronous client communicates with
(potentially remote) Pulsar procedures. Obvious implementation is HTTP based
but Pulsar objects wrapped in routes can also be directly communicated with
if in memory.
"""
__metaclass__ = ABCMeta
@abstractmethod
def execute(self, command, args={}, data=None, input_path=None, output_path=None):
"""
Execute the correspond command against configured Pulsar job manager. Arguments are
method parameters and data or input_path describe essentially POST bodies. If command
results in a file, resulting path should be specified as output_path.
"""
COMMAND_TO_PATH = {
"path": Template("jobs/${job_id}/files/path"),
"upload_file": Template("jobs/${job_id}/files"),
"download_output": Template("jobs/${job_id}/files"),
"setup": Template("jobs"),
"clean": Template("jobs/${job_id}"),
"status": Template("jobs/${job_id}/status"),
"cancel": Template("jobs/${job_id}/cancel"),
"submit": Template("jobs/${job_id}/submit"),
"file_available": Template("cache/status"),
"cache_required": Template("cache"),
"cache_insert": Template("cache"),
"object_store_exists": Template("objects/${object_id}/exists"),
"object_store_file_ready": Template("objects/${object_id}/file_ready"),
"object_store_update_from_file": Template("objects/${object_id}"),
"object_store_create": Template("objects/${object_id}"),
"object_store_empty": Template("objects/${object_id}/empty"),
"object_store_size": Template("objects/${object_id}/size"),
"object_store_delete": Template("objects/${object_id}"),
"object_store_get_data": Template("objects/${object_id}"),
"object_store_get_filename": Template("objects/${object_id}/filename"),
"object_store_get_store_usage_percent": Template("object_store_usage_percent")
}
COMMAND_TO_METHOD = {
"upload_file": "POST",
"download_output": "GET",
"setup": "POST",
"submit": "POST",
"clean": "DELETE",
"cancel": "PUT",
"object_store_update_from_file": "PUT",
"object_store_create": "POST",
"object_store_delete": "DELETE",
"file_available": "GET",
"cache_required": "PUT",
"cache_insert": "POST",
}
class HttpPulsarInterface(PulsarInterface):
def __init__(self, destination_params, transport):
self.transport = transport
remote_host = destination_params.get("url")
assert remote_host is not None, "Failed to determine url for Pulsar client."
if not remote_host.endswith("/"):
remote_host = "%s/" % remote_host
if not remote_host.startswith("http"):
remote_host = "http://%s" % remote_host
self.remote_host = remote_host
self.private_token = destination_params.get("private_token", None)
def execute(self, command, args={}, data=None, input_path=None, output_path=None):
url = self.__build_url(command, args)
method = COMMAND_TO_METHOD.get(command, None) # Default to GET is no data, POST otherwise
response = self.transport.execute(url, method=method, data=data, input_path=input_path, output_path=output_path)
return response
def __build_url(self, command, args):
path = COMMAND_TO_PATH.get(command, Template(command)).safe_substitute(args)
if self.private_token:
args["private_token"] = self.private_token
arg_bytes = dict([(k, text_type(args[k]).encode('utf-8')) for k in args])
data = urlencode(arg_bytes)
url = self.remote_host + path + "?" + data
return url
class LocalPulsarInterface(PulsarInterface):
def __init__(self, destination_params, job_manager=None, file_cache=None, object_store=None):
self.job_manager = job_manager
self.file_cache = file_cache
self.object_store = object_store
def __app_args(self):
# Arguments that would be specified from PulsarApp if running
# in web server.
return {
'manager': self.job_manager,
'file_cache': self.file_cache,
'object_store': self.object_store,
'ip': None
}
def execute(self, command, args={}, data=None, input_path=None, output_path=None):
# If data set, should be unicode (on Python 2) or str (on Python 3).
from pulsar.web import routes
from pulsar.web.framework import build_func_args
controller = getattr(routes, command)
action = controller.func
body_args = dict(body=self.__build_body(data, input_path))
args = build_func_args(action, args.copy(), self.__app_args(), body_args)
result = action(**args)
if controller.response_type != 'file':
return controller.body(result)
else:
# TODO: Add to Galaxy.
from galaxy.util import copy_to_path
with open(result, 'rb') as result_file:
copy_to_path(result_file, output_path)
def __build_body(self, data, input_path):
if data is not None:
return BytesIO(data)
elif input_path is not None:
return open(input_path, 'rb')
else:
return None
-135
View File
@@ -1,135 +0,0 @@
"""
"""
import os.path
from collections import deque
import posixpath
from .util import PathHelper
from galaxy.util import in_directory
from logging import getLogger
log = getLogger(__name__)
TYPES_TO_METHOD = dict(
input="inputs_directory",
unstructured="unstructured_files_directory",
config="configs_directory",
tool="tool_files_directory",
workdir="working_directory",
output="outputs_directory",
output_workdir="working_directory",
)
class RemoteJobDirectory(object):
""" Representation of a (potentially) remote Pulsar-style staging directory.
"""
def __init__(self, remote_staging_directory, remote_id, remote_sep):
self.path_helper = PathHelper(remote_sep)
self.job_directory = self.path_helper.remote_join(
remote_staging_directory,
remote_id
)
def working_directory(self):
return self._sub_dir('working')
def inputs_directory(self):
return self._sub_dir('inputs')
def outputs_directory(self):
return self._sub_dir('outputs')
def configs_directory(self):
return self._sub_dir('configs')
def tool_files_directory(self):
return self._sub_dir('tool_files')
def unstructured_files_directory(self):
return self._sub_dir('unstructured')
@property
def path(self):
return self.job_directory
@property
def separator(self):
return self.path_helper.separator
def calculate_path(self, remote_relative_path, input_type):
""" Only for used by Pulsar client, should override for managers to
enforce security and make the directory if needed.
"""
directory, allow_nested_files = self._directory_for_file_type(input_type)
return self.path_helper.remote_join(directory, remote_relative_path)
def _directory_for_file_type(self, file_type):
allow_nested_files = False
# work_dir and input_extra are types used by legacy clients...
# Obviously this client won't be legacy because this is in the
# client module, but this code is reused on server which may
# serve legacy clients.
allow_nested_files = file_type in ['input', 'unstructured', 'output', 'output_workdir']
directory_function = getattr(self, TYPES_TO_METHOD.get(file_type, None), None)
if not directory_function:
raise Exception("Unknown file_type specified %s" % file_type)
return directory_function(), allow_nested_files
def _sub_dir(self, name):
return self.path_helper.remote_join(self.job_directory, name)
def get_mapped_file(directory, remote_path, allow_nested_files=False, local_path_module=os.path, mkdir=True):
"""
>>> import ntpath
>>> get_mapped_file(r'C:\\pulsar\\staging\\101', 'dataset_1_files/moo/cow', allow_nested_files=True, local_path_module=ntpath, mkdir=False)
'C:\\\\pulsar\\\\staging\\\\101\\\\dataset_1_files\\\\moo\\\\cow'
>>> get_mapped_file(r'C:\\pulsar\\staging\\101', 'dataset_1_files/moo/cow', allow_nested_files=False, local_path_module=ntpath)
'C:\\\\pulsar\\\\staging\\\\101\\\\cow'
>>> get_mapped_file(r'C:\\pulsar\\staging\\101', '../cow', allow_nested_files=True, local_path_module=ntpath, mkdir=False)
Traceback (most recent call last):
Exception: Attempt to read or write file outside an authorized directory.
"""
if not allow_nested_files:
name = local_path_module.basename(remote_path)
path = local_path_module.join(directory, name)
else:
local_rel_path = __posix_to_local_path(remote_path, local_path_module=local_path_module)
local_path = local_path_module.join(directory, local_rel_path)
verify_is_in_directory(local_path, directory, local_path_module=local_path_module)
local_directory = local_path_module.dirname(local_path)
if mkdir and not local_path_module.exists(local_directory):
os.makedirs(local_directory)
path = local_path
return path
def __posix_to_local_path(path, local_path_module=os.path):
"""
Converts a posix path (coming from Galaxy), to a local path (be it posix or Windows).
>>> import ntpath
>>> __posix_to_local_path('dataset_1_files/moo/cow', local_path_module=ntpath)
'dataset_1_files\\\\moo\\\\cow'
>>> import posixpath
>>> __posix_to_local_path('dataset_1_files/moo/cow', local_path_module=posixpath)
'dataset_1_files/moo/cow'
"""
partial_path = deque()
while True:
if not path or path == '/':
break
(path, base) = posixpath.split(path)
partial_path.appendleft(base)
return local_path_module.join(*partial_path)
def verify_is_in_directory(path, directory, local_path_module=os.path):
if not in_directory(path, directory, local_path_module):
msg = "Attempt to read or write file outside an authorized directory."
log.warn("%s Attempted path: %s, valid directory: %s" % (msg, path, directory))
raise Exception(msg)
-281
View File
@@ -1,281 +0,0 @@
import threading
import functools
try:
from Queue import Queue
except ImportError:
from queue import Queue
from os import getenv
from six import string_types
from .client import JobClient
from .client import InputCachingJobClient
from .client import MessageJobClient
from .client import MessageCLIJobClient
from .interface import HttpPulsarInterface
from .interface import LocalPulsarInterface
from .object_client import ObjectStoreClient
from .transport import get_transport
from .util import TransferEventManager
from .destination import url_to_destination_params
from .amqp_exchange_factory import get_exchange
from logging import getLogger
log = getLogger(__name__)
DEFAULT_TRANSFER_THREADS = 2
def build_client_manager(**kwargs):
if 'job_manager' in kwargs:
return ClientManager(**kwargs) # TODO: Consider more separation here.
elif kwargs.get('amqp_url', None):
return MessageQueueClientManager(**kwargs)
else:
return ClientManager(**kwargs)
class ClientManager(object):
"""
Factory to create Pulsar clients, used to manage potential shared
state between multiple client connections.
"""
def __init__(self, **kwds):
if 'job_manager' in kwds:
self.job_manager_interface_class = LocalPulsarInterface
self.job_manager_interface_args = dict(job_manager=kwds['job_manager'], file_cache=kwds['file_cache'])
else:
self.job_manager_interface_class = HttpPulsarInterface
transport_type = kwds.get('transport', None)
transport = get_transport(transport_type)
self.job_manager_interface_args = dict(transport=transport)
cache = kwds.get('cache', None)
if cache is None:
cache = _environ_default_int('PULSAR_CACHE_TRANSFERS')
if cache:
log.info("Setting Pulsar client class to caching variant.")
self.client_cacher = ClientCacher(**kwds)
self.client_class = InputCachingJobClient
self.extra_client_kwds = {"client_cacher": self.client_cacher}
else:
log.info("Setting Pulsar client class to standard, non-caching variant.")
self.client_class = JobClient
self.extra_client_kwds = {}
def get_client(self, destination_params, job_id, **kwargs):
destination_params = _parse_destination_params(destination_params)
destination_params.update(**kwargs)
job_manager_interface_class = self.job_manager_interface_class
job_manager_interface_args = dict(destination_params=destination_params, **self.job_manager_interface_args)
job_manager_interface = job_manager_interface_class(**job_manager_interface_args)
return self.client_class(destination_params, job_id, job_manager_interface, **self.extra_client_kwds)
def shutdown(self, ensure_cleanup=False):
pass
try:
from galaxy.jobs.runners.util.cli import factory as cli_factory
except ImportError:
from pulsar.managers.util.cli import factory as cli_factory
class MessageQueueClientManager(object):
def __init__(self, **kwds):
self.url = kwds.get('amqp_url')
self.manager_name = kwds.get("manager", None) or "_default_"
self.exchange = get_exchange(self.url, self.manager_name, kwds)
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:
message.requeue()
log.debug("Requeue succeeded, will likely be handled next time consumer is enabled.")
except Exception:
log.debug("Requeue failed, message may be lost?")
return
try:
if "job_id" in body:
job_id = body["job_id"]
self.status_cache[job_id] = body
log.debug("Handling asynchronous status update from remote Pulsar.")
callback(body)
except Exception:
log.exception("Failure processing job status update message.")
except BaseException as e:
log.exception("Failure processing job status update message - BaseException type %s" % type(e))
finally:
message.ack()
def callback_consumer(self, callback_wrapper):
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.")
finally:
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:
if self.callback_thread is not None:
return
callback_wrapper = functools.partial(self.callback_wrapper, callback)
run = functools.partial(self.callback_consumer, callback_wrapper)
thread = threading.Thread(
name="pulsar_client_%s_status_update_callback" % self.manager_name,
target=run
)
thread.daemon = False # Lets not interrupt processing of this.
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.")
destination_params = _parse_destination_params(destination_params)
destination_params.update(**kwargs)
if 'shell_plugin' in destination_params:
shell = cli_factory.get_shell(destination_params)
return MessageCLIJobClient(destination_params, job_id, self, shell)
else:
return MessageJobClient(destination_params, job_id, self)
class ObjectStoreClientManager(object):
def __init__(self, **kwds):
if 'object_store' in kwds:
self.interface_class = LocalPulsarInterface
self.interface_args = dict(object_store=kwds['object_store'])
else:
self.interface_class = HttpPulsarInterface
transport_type = kwds.get('transport', None)
transport = get_transport(transport_type)
self.interface_args = dict(transport=transport)
self.extra_client_kwds = {}
def get_client(self, client_params):
interface_class = self.interface_class
interface_args = dict(destination_params=client_params, **self.interface_args)
interface = interface_class(**interface_args)
return ObjectStoreClient(interface)
class ClientCacher(object):
def __init__(self, **kwds):
self.event_manager = TransferEventManager()
default_transfer_threads = _environ_default_int('PULSAR_CACHE_THREADS', DEFAULT_TRANSFER_THREADS)
num_transfer_threads = int(kwds.get('transfer_threads', default_transfer_threads))
self.__init_transfer_threads(num_transfer_threads)
def queue_transfer(self, client, path):
self.transfer_queue.put((client, path))
def acquire_event(self, input_path):
return self.event_manager.acquire_event(input_path)
def _transfer_worker(self):
while True:
transfer_info = self.transfer_queue.get()
try:
self.__perform_transfer(transfer_info)
except BaseException as e:
log.warn("Transfer failed.")
log.exception(e)
pass
self.transfer_queue.task_done()
def __perform_transfer(self, transfer_info):
(client, path) = transfer_info
event_holder = self.event_manager.acquire_event(path, force_clear=True)
failed = True
try:
client.cache_insert(path)
failed = False
finally:
event_holder.failed = failed
event_holder.release()
def __init_transfer_threads(self, num_transfer_threads):
self.num_transfer_threads = num_transfer_threads
self.transfer_queue = Queue()
for i in range(num_transfer_threads):
t = threading.Thread(target=self._transfer_worker)
t.daemon = True
t.start()
def _parse_destination_params(destination_params):
if isinstance(destination_params, string_types):
destination_params = url_to_destination_params(destination_params)
return destination_params
def _environ_default_int(variable, default="0"):
val = getenv(variable, default)
int_val = int(default)
if str(val).isdigit():
int_val = int(val)
return int_val
__all__ = [
'ClientManager',
'ObjectStoreClientManager',
'HttpPulsarInterface'
]
-53
View File
@@ -1,53 +0,0 @@
from .decorators import parseJson
class ObjectStoreClient(object):
def __init__(self, pulsar_interface):
self.pulsar_interface = pulsar_interface
@parseJson()
def exists(self, **kwds):
return self._raw_execute("object_store_exists", args=self.__data(**kwds))
@parseJson()
def file_ready(self, **kwds):
return self._raw_execute("object_store_file_ready", args=self.__data(**kwds))
@parseJson()
def create(self, **kwds):
return self._raw_execute("object_store_create", args=self.__data(**kwds))
@parseJson()
def empty(self, **kwds):
return self._raw_execute("object_store_empty", args=self.__data(**kwds))
@parseJson()
def size(self, **kwds):
return self._raw_execute("object_store_size", args=self.__data(**kwds))
@parseJson()
def delete(self, **kwds):
return self._raw_execute("object_store_delete", args=self.__data(**kwds))
@parseJson()
def get_data(self, **kwds):
return self._raw_execute("object_store_get_data", args=self.__data(**kwds))
@parseJson()
def get_filename(self, **kwds):
return self._raw_execute("object_store_get_filename", args=self.__data(**kwds))
@parseJson()
def update_from_file(self, **kwds):
return self._raw_execute("object_store_update_from_file", args=self.__data(**kwds))
@parseJson()
def get_store_usage_percent(self):
return self._raw_execute("object_store_get_store_usage_percent", args={})
def __data(self, **kwds):
return kwds
def _raw_execute(self, command, args={}):
return self.pulsar_interface.execute(command, args, data=None, input_path=None, output_path=None)
-98
View File
@@ -1,98 +0,0 @@
import os.path
from .action_mapper import FileActionMapper
from .action_mapper import path_type
from .util import PathHelper
from galaxy.util import in_directory
class PathMapper(object):
""" Ties together a FileActionMapper and remote job configuration returned
by the Pulsar setup method to pre-determine the location of files for staging
on the remote Pulsar server.
This is not useful when rewrite_paths (as has traditionally been done with
the Pulsar) because when doing that the Pulsar determines the paths as files are
uploaded. When rewrite_paths is disabled however, the destination of files
needs to be determined prior to transfer so an object of this class can be
used.
"""
def __init__(
self,
client,
remote_job_config,
local_working_directory,
action_mapper=None,
):
self.local_working_directory = local_working_directory
if not action_mapper:
action_mapper = FileActionMapper(client)
self.action_mapper = action_mapper
self.input_directory = remote_job_config["inputs_directory"]
self.output_directory = remote_job_config["outputs_directory"]
self.working_directory = remote_job_config["working_directory"]
self.unstructured_files_directory = remote_job_config["unstructured_files_directory"]
self.config_directory = remote_job_config["configs_directory"]
separator = remote_job_config["system_properties"]["separator"]
self.path_helper = PathHelper(separator)
def remote_output_path_rewrite(self, local_path):
output_type = path_type.OUTPUT
if in_directory(local_path, self.local_working_directory):
output_type = path_type.OUTPUT_WORKDIR
remote_path = self.__remote_path_rewrite(local_path, output_type)
return remote_path
def remote_input_path_rewrite(self, local_path):
remote_path = self.__remote_path_rewrite(local_path, path_type.INPUT)
return remote_path
def remote_version_path_rewrite(self, local_path):
remote_path = self.__remote_path_rewrite(local_path, path_type.OUTPUT, name="COMMAND_VERSION")
return remote_path
def check_for_arbitrary_rewrite(self, local_path):
path = str(local_path) # Use false_path if needed.
action = self.action_mapper.action(path, path_type.UNSTRUCTURED)
if not action.staging_needed:
return action.path_rewrite(self.path_helper), []
unique_names = action.unstructured_map()
name = unique_names[path]
remote_path = self.path_helper.remote_join(self.unstructured_files_directory, name)
return remote_path, unique_names
def __remote_path_rewrite(self, dataset_path, dataset_path_type, name=None):
""" Return remote path of this file (if staging is required) else None.
"""
path = str(dataset_path) # Use false_path if needed.
action = self.action_mapper.action(path, dataset_path_type)
if action.staging_needed:
if name is None:
name = os.path.basename(path)
remote_directory = self.__remote_directory(dataset_path_type)
remote_path_rewrite = self.path_helper.remote_join(remote_directory, name)
else:
# Actions which don't require staging MUST define a path_rewrite
# method.
remote_path_rewrite = action.path_rewrite(self.path_helper)
return remote_path_rewrite
def __action(self, dataset_path, dataset_path_type):
path = str(dataset_path) # Use false_path if needed.
action = self.action_mapper.action(path, dataset_path_type)
return action
def __remote_directory(self, dataset_path_type):
if dataset_path_type in [path_type.OUTPUT]:
return self.output_directory
elif dataset_path_type in [path_type.WORKDIR, path_type.OUTPUT_WORKDIR]:
return self.working_directory
elif dataset_path_type in [path_type.INPUT]:
return self.input_directory
else:
message = "PathMapper cannot handle path type %s" % dataset_path_type
raise Exception(message)
__all__ = ['PathMapper']
-104
View File
@@ -1,104 +0,0 @@
import os
from .util import filter_destination_params
REMOTE_SYSTEM_PROPERTY_PREFIX = "remote_property_"
def build(client, destination_args):
""" Build a SetupHandler object for client from destination parameters.
"""
# Have defined a remote job directory, lets do the setup locally.
if client.job_directory:
handler = LocalSetupHandler(client, destination_args)
else:
handler = RemoteSetupHandler(client)
return handler
class LocalSetupHandler(object):
""" Parse destination params to infer job setup parameters (input/output
directories, etc...). Default is to get this configuration data from the
remote Pulsar server.
Downside of this approach is that it requires more and more dependent
configuraiton of Galaxy. Upside is that it is asynchronous and thus makes
message queue driven configurations possible.
Remote system properties (such as galaxy_home) can be specified in
destination args by prefixing property with remote_property_ (e.g.
remote_property_galaxy_home).
"""
def __init__(self, client, destination_args):
self.client = client
system_properties = self.__build_system_properties(destination_args)
system_properties["separator"] = client.job_directory.separator
self.system_properties = system_properties
self.jobs_directory = destination_args["jobs_directory"]
def setup(self, job_id, tool_id=None, tool_version=None):
return build_job_config(
job_id=job_id,
job_directory=self.client.job_directory,
system_properties=self.system_properties,
tool_id=tool_id,
tool_version=tool_version,
)
@property
def local(self):
"""
"""
return True
def __build_system_properties(self, destination_params):
return filter_destination_params(destination_params, REMOTE_SYSTEM_PROPERTY_PREFIX)
class RemoteSetupHandler(object):
""" Default behavior. Fetch setup information from remote Pulsar server.
"""
def __init__(self, client):
self.client = client
def setup(self, **setup_args):
return self.client.remote_setup(**setup_args)
@property
def local(self):
"""
"""
return False
def build_job_config(job_id, job_directory, system_properties={}, tool_id=None, tool_version=None):
"""
"""
inputs_directory = job_directory.inputs_directory()
working_directory = job_directory.working_directory()
outputs_directory = job_directory.outputs_directory()
configs_directory = job_directory.configs_directory()
tools_directory = job_directory.tool_files_directory()
unstructured_files_directory = job_directory.unstructured_files_directory()
sep = system_properties.get("sep", os.sep)
job_config = {
"job_directory": job_directory.path,
"working_directory": working_directory,
"outputs_directory": outputs_directory,
"configs_directory": configs_directory,
"tools_directory": tools_directory,
"inputs_directory": inputs_directory,
"unstructured_files_directory": unstructured_files_directory,
# Poorly named legacy attribute. Drop at some point.
"path_separator": sep,
"job_id": job_id,
"system_properties": system_properties,
}
if tool_id:
job_config["tool_id"] = tool_id
if tool_version:
job_config["tool_version"] = tool_version
return job_config
__all__ = ['build_job_config', 'build']
-177
View File
@@ -1,177 +0,0 @@
import re
from os.path import basename
from os.path import join
from os.path import dirname
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):
""" A description of how client views job - command_line, inputs, etc..
**Parameters**
command_line : str
The local command line to execute, this will be rewritten for
the remote server.
config_files : list
List of Galaxy 'configfile's produced for this job. These will
be rewritten and sent to remote server.
input_files : list
List of input files used by job. These will be transferred and
references rewritten.
client_outputs : ClientOutputs
Description of outputs produced by job (at least output files along
with optional version string and working directory outputs.
tool_dir : str
Directory containing tool to execute (if a wrapper is used, it will
be transferred to remote server).
working_directory : str
Local path created by Galaxy for running this job.
dependencies_description : list
galaxy.tools.deps.dependencies.DependencyDescription object describing
tool dependency context for remote depenency resolution.
env: list
List of dict object describing environment variables to populate.
version_file : str
Path to version file expected on the client server
arbitrary_files : dict()
Additional non-input, non-tool, non-config, non-working directory files
to transfer before staging job. This is most likely data indices but
can be anything. For now these are copied into staging working
directory but this will be reworked to find a better, more robust
location.
rewrite_paths : boolean
Indicates whether paths should be rewritten in job inputs (command_line
and config files) while staging files).
"""
def __init__(
self,
command_line,
tool=None,
config_files=[],
input_files=[],
client_outputs=None,
working_directory=None, # More sensible default?
dependencies_description=None,
env=[],
arbitrary_files=None,
rewrite_paths=True,
):
self.tool = tool
self.command_line = command_line
self.config_files = config_files
self.input_files = input_files
self.client_outputs = client_outputs or ClientOutputs()
self.working_directory = working_directory
self.dependencies_description = dependencies_description
self.env = env
self.rewrite_paths = rewrite_paths
self.arbitrary_files = arbitrary_files or {}
@property
def output_files(self):
return self.client_outputs.output_files
@property
def version_file(self):
return self.client_outputs.version_file
@property
def tool_dependencies(self):
if not self.remote_dependency_resolution:
return None
return dict(
requirements=(self.tool.requirements or []),
installed_tool_dependencies=(self.tool.installed_tool_dependencies or [])
)
class ClientOutputs(object):
""" Abstraction describing the output datasets EXPECTED by the Galaxy job
runner client.
"""
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 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,
dynamic_outputs=self.dynamic_outputs,
)
@staticmethod
def from_dict(config_dict):
return ClientOutputs(
working_directory=config_dict.get('working_directory'),
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
server. """
def __init__(self, working_directory_contents, output_directory_contents, remote_separator=sep):
self.working_directory_contents = working_directory_contents
self.output_directory_contents = output_directory_contents
self.path_helper = PathHelper(remote_separator)
@staticmethod
def from_status_response(complete_response):
# Default to None instead of [] to distinguish between empty contents and it not set
# by the Pulsar - older Pulsar instances will not set these in complete response.
working_directory_contents = complete_response.get("working_directory_contents")
output_directory_contents = complete_response.get("outputs_directory_contents")
# Older (pre-2014) Pulsar servers will not include separator in response,
# so this should only be used when reasoning about outputs in
# subdirectories (which was not previously supported prior to that).
remote_separator = complete_response.get("system_properties", {}).get("separator", sep)
return PulsarOutputs(
working_directory_contents,
output_directory_contents,
remote_separator
)
def has_output_file(self, output_file):
return basename(output_file) in self.output_directory_contents
def output_extras(self, output_file):
"""
Returns dict mapping local path to remote name.
"""
output_directory = dirname(output_file)
def local_path(name):
return join(output_directory, self.path_helper.local_name(name))
files_directory = "%s_files%s" % (basename(output_file)[0:-len(".dat")], self.path_helper.separator)
names = filter(lambda o: o.startswith(files_directory), self.output_directory_contents)
return dict(map(lambda name: (local_path(name), name), names))
-157
View File
@@ -1,157 +0,0 @@
from os.path import join
from os.path import relpath
from re import compile
from contextlib import contextmanager
from ..staging import COMMAND_VERSION_FILENAME
from ..action_mapper import FileActionMapper
from logging import getLogger
log = getLogger(__name__)
# All output files marked with from_work_dir attributes will copied or downloaded
# this pattern picks up attiditional files to copy back - such as those
# associated with multiple outputs and metadata configuration. Set to .* to just
# copy everything
COPY_FROM_WORKING_DIRECTORY_PATTERN = compile(r"primary_.*|galaxy.json|metadata_.*|dataset_\d+\.dat|__instrument_.*|dataset_\d+_files.+")
def finish_job(client, cleanup_job, job_completed_normally, client_outputs, pulsar_outputs):
""" Responsible for downloading results from remote server and cleaning up
Pulsar staging directory (if needed.)
"""
collection_failure_exceptions = []
if job_completed_normally:
output_collector = ClientOutputCollector(client)
action_mapper = FileActionMapper(client)
results_stager = ResultsCollector(output_collector, action_mapper, client_outputs, pulsar_outputs)
collection_failure_exceptions = results_stager.collect()
__clean(collection_failure_exceptions, cleanup_job, client)
return collection_failure_exceptions
class ClientOutputCollector(object):
def __init__(self, client):
self.client = client
def collect_output(self, results_collector, output_type, action, name):
# This output should have been handled by the Pulsar.
if not action.staging_action_local:
return False
working_directory = results_collector.client_outputs.working_directory
self.client.fetch_output(
path=action.path,
name=name,
working_directory=working_directory,
output_type=output_type,
action_type=action.action_type
)
return True
class ResultsCollector(object):
def __init__(self, output_collector, action_mapper, client_outputs, pulsar_outputs):
self.output_collector = output_collector
self.action_mapper = action_mapper
self.client_outputs = client_outputs
self.pulsar_outputs = pulsar_outputs
self.downloaded_working_directory_files = []
self.exception_tracker = DownloadExceptionTracker()
self.output_files = client_outputs.output_files
self.working_directory_contents = pulsar_outputs.working_directory_contents or []
def collect(self):
self.__collect_working_directory_outputs()
self.__collect_outputs()
self.__collect_version_file()
self.__collect_other_working_directory_files()
return self.exception_tracker.collection_failure_exceptions
def __collect_working_directory_outputs(self):
working_directory = self.client_outputs.working_directory
# Fetch explicit working directory outputs.
for source_file, output_file in self.client_outputs.work_dir_outputs:
name = relpath(source_file, working_directory)
pulsar = self.pulsar_outputs.path_helper.remote_name(name)
if self._attempt_collect_output('output_workdir', path=output_file, name=pulsar):
self.downloaded_working_directory_files.append(pulsar)
# Remove from full output_files list so don't try to download directly.
try:
self.output_files.remove(output_file)
except ValueError:
raise Exception("Failed to remove %s from %s" % (output_file, self.output_files))
def __collect_outputs(self):
# Legacy Pulsar not returning list of files, iterate over the list of
# expected outputs for tool.
for output_file in self.output_files:
# Fetch output directly...
output_generated = self.pulsar_outputs.has_output_file(output_file)
if output_generated:
self._attempt_collect_output('output', output_file)
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.
def __collect_version_file(self):
version_file = self.client_outputs.version_file
pulsar_output_directory_contents = self.pulsar_outputs.output_directory_contents
if version_file and COMMAND_VERSION_FILENAME in pulsar_output_directory_contents:
self._attempt_collect_output('output', version_file, name=COMMAND_VERSION_FILENAME)
def __collect_other_working_directory_files(self):
working_directory = self.client_outputs.working_directory
# Fetch remaining working directory outputs of interest.
for name in self.working_directory_contents:
if name in self.downloaded_working_directory_files:
continue
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)
def _attempt_collect_output(self, output_type, path, name=None):
# path is final path on galaxy server (client)
# name is the 'name' of the file on the Pulsar server (possible a relative)
# path.
collected = False
with self.exception_tracker():
action = self.action_mapper.action(path, output_type)
if self._collect_output(output_type, action, name):
collected = True
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)
class DownloadExceptionTracker(object):
def __init__(self):
self.collection_failure_exceptions = []
@contextmanager
def __call__(self):
try:
yield
except Exception as e:
self.collection_failure_exceptions.append(e)
def __clean(collection_failure_exceptions, cleanup_job, client):
failed = (len(collection_failure_exceptions) > 0)
if (not failed and cleanup_job != "never") or cleanup_job == "always":
try:
client.clean()
except Exception:
log.warn("Failed to cleanup remote Pulsar job")
__all__ = ['finish_job']
-438
View File
@@ -1,438 +0,0 @@
from os.path import abspath, basename, join, exists
from os.path import dirname
from os.path import relpath
from os import listdir, sep
from re import findall
from io import open
from ..staging import COMMAND_VERSION_FILENAME
from ..action_mapper import FileActionMapper
from ..action_mapper import path_type
from ..action_mapper import MessageAction
from ..util import PathHelper
from ..util import directory_files
from logging import getLogger
log = getLogger(__name__)
def submit_job(client, client_job_description, job_config=None):
"""
"""
file_stager = FileStager(client, client_job_description, job_config)
rebuilt_command_line = file_stager.get_command_line()
job_id = file_stager.job_id
launch_kwds = dict(
command_line=rebuilt_command_line,
dependencies_description=client_job_description.dependencies_description,
env=client_job_description.env,
)
if file_stager.job_config:
launch_kwds["job_config"] = file_stager.job_config
remote_staging = {}
remote_staging_actions = file_stager.transfer_tracker.remote_staging_actions
if remote_staging_actions:
remote_staging["setup"] = remote_staging_actions
# Somehow make the following optional.
remote_staging["action_mapper"] = file_stager.action_mapper.to_dict()
remote_staging["client_outputs"] = client_job_description.client_outputs.to_dict()
if remote_staging:
launch_kwds["remote_staging"] = remote_staging
client.launch(**launch_kwds)
return job_id
class FileStager(object):
"""
Objects of the FileStager class interact with an Pulsar client object to
stage the files required to run jobs on a remote Pulsar server.
**Parameters**
client : JobClient
Pulsar client object.
client_job_description : client_job_description
Description of client view of job to stage and execute remotely.
"""
def __init__(self, client, client_job_description, job_config):
"""
"""
self.client = client
self.command_line = client_job_description.command_line
self.config_files = client_job_description.config_files
self.input_files = client_job_description.input_files
self.output_files = client_job_description.output_files
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
self.rewrite_paths = client_job_description.rewrite_paths
# Setup job inputs, these will need to be rewritten before
# shipping off to remote Pulsar server.
self.job_inputs = JobInputs(self.command_line, self.config_files)
self.action_mapper = FileActionMapper(client)
self.__handle_setup(job_config)
self.transfer_tracker = TransferTracker(client, self.path_helper, self.action_mapper, self.job_inputs, rewrite_paths=self.rewrite_paths)
self.__initialize_referenced_tool_files()
if self.rewrite_paths:
self.__initialize_referenced_arbitrary_files()
self.__upload_tool_files()
self.__upload_input_files()
self.__upload_working_directory_files()
self.__upload_arbitrary_files()
if self.rewrite_paths:
self.__initialize_output_file_renames()
self.__initialize_task_output_file_renames()
self.__initialize_config_file_renames()
self.__initialize_version_file_rename()
self.__handle_rewrites()
self.__upload_rewritten_config_files()
def __handle_setup(self, job_config):
if not job_config:
job_config = self.client.setup(self.tool_id, self.tool_version)
self.new_working_directory = job_config['working_directory']
self.new_outputs_directory = job_config['outputs_directory']
# Default configs_directory to match remote working_directory to mimic
# behavior of older Pulsar servers.
self.new_configs_directory = job_config.get('configs_directory', self.new_working_directory)
self.remote_separator = self.__parse_remote_separator(job_config)
self.path_helper = PathHelper(self.remote_separator)
# If remote Pulsar server assigned job id, use that otherwise
# just use local job_id assigned.
galaxy_job_id = self.client.job_id
self.job_id = job_config.get('job_id', galaxy_job_id)
if self.job_id != galaxy_job_id:
# Remote Pulsar server assigned an id different than the
# Galaxy job id, update client to reflect this.
self.client.job_id = self.job_id
self.job_config = job_config
def __parse_remote_separator(self, job_config):
separator = job_config.get("system_properties", {}).get("separator", None)
if not separator: # Legacy Pulsar
separator = job_config["path_separator"] # Poorly named
return separator
def __initialize_referenced_tool_files(self):
self.referenced_tool_files = self.job_inputs.find_referenced_subfiles(self.tool_dir)
def __initialize_referenced_arbitrary_files(self):
referenced_arbitrary_path_mappers = dict()
for mapper in self.action_mapper.unstructured_mappers():
mapper_pattern = mapper.to_pattern()
# TODO: Make more sophisticated, allow parent directories,
# grabbing sibbling files based on patterns, etc...
paths = self.job_inputs.find_pattern_references(mapper_pattern)
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.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)
def __upload_tool_files(self):
for referenced_tool_file in self.referenced_tool_files:
self.transfer_tracker.handle_transfer(referenced_tool_file, path_type.TOOL)
def __upload_arbitrary_files(self):
for path, name in self.arbitrary_files.items():
self.transfer_tracker.handle_transfer(path, path_type.UNSTRUCTURED, name=name)
def __upload_input_files(self):
for input_file in self.input_files:
self.__upload_input_file(input_file)
self.__upload_input_extra_files(input_file)
def __upload_input_file(self, input_file):
if self.__stage_input(input_file):
if exists(input_file):
self.transfer_tracker.handle_transfer(input_file, path_type.INPUT)
else:
message = "Pulsar: __upload_input_file called on empty or missing dataset." + \
" So such file: [%s]" % input_file
log.debug(message)
def __upload_input_extra_files(self, input_file):
files_path = "%s_files" % input_file[0:-len(".dat")]
if exists(files_path) and self.__stage_input(files_path):
for extra_file_name in directory_files(files_path):
extra_file_path = join(files_path, extra_file_name)
remote_name = self.path_helper.remote_name(relpath(extra_file_path, dirname(files_path)))
self.transfer_tracker.handle_transfer(extra_file_path, path_type.INPUT, name=remote_name)
def __upload_working_directory_files(self):
# Task manager stages files into working directory, these need to be
# uploaded if present.
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:
remote_path = self.path_helper.remote_join(self.new_outputs_directory, COMMAND_VERSION_FILENAME)
self.transfer_tracker.register_rewrite(version_file, remote_path, path_type.OUTPUT)
def __initialize_output_file_renames(self):
for output_file in self.output_files:
remote_path = self.path_helper.remote_join(self.new_outputs_directory, basename(output_file))
self.transfer_tracker.register_rewrite(output_file, remote_path, path_type.OUTPUT)
def __initialize_task_output_file_renames(self):
for output_file in self.output_files:
name = basename(output_file)
task_file = join(self.working_directory, name)
remote_path = self.path_helper.remote_join(self.new_working_directory, name)
self.transfer_tracker.register_rewrite(task_file, remote_path, path_type.OUTPUT_WORKDIR)
def __initialize_config_file_renames(self):
for config_file in self.config_files:
remote_path = self.path_helper.remote_join(self.new_configs_directory, basename(config_file))
self.transfer_tracker.register_rewrite(config_file, remote_path, path_type.CONFIG)
def __handle_rewrites(self):
"""
For each file that has been transferred and renamed, updated
command_line and configfiles to reflect that rewrite.
"""
self.transfer_tracker.rewrite_input_paths()
def __upload_rewritten_config_files(self):
for config_file, new_config_contents in self.job_inputs.config_files.items():
self.transfer_tracker.handle_transfer(config_file, type=path_type.CONFIG, contents=new_config_contents)
def get_command_line(self):
"""
Returns the rewritten version of the command line to execute suitable
for remote host.
"""
return self.job_inputs.command_line
def __stage_input(self, file_path):
# If we have disabled path rewriting, just assume everything needs to be transferred,
# else check to ensure the file is referenced before transferring it.
return (not self.rewrite_paths) or self.job_inputs.path_referenced(file_path)
class JobInputs(object):
"""
Abstractions over dynamic inputs created for a given job (namely the command to
execute and created configfiles).
**Parameters**
command_line : str
Local command to execute for this job. (To be rewritten.)
config_files : str
Config files created for this job. (To be rewritten.)
>>> import tempfile
>>> tf = tempfile.NamedTemporaryFile()
>>> def setup_inputs(tf):
... open(tf.name, "w").write(u"world /path/to/input the rest")
... inputs = JobInputs(u"hello /path/to/input", [tf.name])
... return inputs
>>> inputs = setup_inputs(tf)
>>> inputs.rewrite_paths(u"/path/to/input", u'C:\\input')
>>> inputs.command_line == u'hello C:\\\\input'
True
>>> inputs.config_files[tf.name] == u'world C:\\\\input the rest'
True
>>> tf.close()
>>> tf = tempfile.NamedTemporaryFile()
>>> inputs = setup_inputs(tf)
>>> inputs.find_referenced_subfiles('/path/to') == [u'/path/to/input']
True
>>> inputs.path_referenced('/path/to')
True
>>> inputs.path_referenced(u'/path/to')
True
>>> inputs.path_referenced('/path/to/input')
True
>>> inputs.path_referenced('/path/to/notinput')
False
>>> tf.close()
"""
def __init__(self, command_line, config_files):
self.command_line = command_line
self.config_files = {}
for config_file in config_files or []:
config_contents = _read(config_file)
self.config_files[config_file] = config_contents
def find_pattern_references(self, pattern):
referenced_files = set()
for input_contents in self.__items():
referenced_files.update(findall(pattern, input_contents))
return list(referenced_files)
def find_referenced_subfiles(self, directory):
"""
Return list of files below specified `directory` in job inputs. Could
use more sophisticated logic (match quotes to handle spaces, handle
subdirectories, etc...).
**Parameters**
directory : str
Full path to directory to search.
"""
if directory is None:
return []
pattern = r"(%s%s\S+)" % (directory, sep)
return self.find_pattern_references(pattern)
def path_referenced(self, path):
pattern = r"%s" % path
found = False
for input_contents in self.__items():
if findall(pattern, input_contents):
found = True
break
return found
def rewrite_paths(self, local_path, remote_path):
"""
Rewrite references to `local_path` with `remote_path` in job inputs.
"""
self.__rewrite_command_line(local_path, remote_path)
self.__rewrite_config_files(local_path, remote_path)
def __rewrite_command_line(self, local_path, remote_path):
self.command_line = self.command_line.replace(local_path, remote_path)
def __rewrite_config_files(self, local_path, remote_path):
for config_file, contents in self.config_files.items():
self.config_files[config_file] = contents.replace(local_path, remote_path)
def __items(self):
items = [self.command_line]
items.extend(self.config_files.values())
return items
class TransferTracker(object):
def __init__(self, client, path_helper, action_mapper, job_inputs, rewrite_paths):
self.client = client
self.path_helper = path_helper
self.action_mapper = action_mapper
self.job_inputs = job_inputs
self.rewrite_paths = rewrite_paths
self.file_renames = {}
self.remote_staging_actions = []
def handle_transfer(self, path, type, name=None, contents=None):
action = self.__action_for_transfer(path, type, contents)
if action.staging_needed:
local_action = action.staging_action_local
if local_action:
response = self.client.put_file(path, type, name=name, contents=contents, action_type=action.action_type)
def get_path():
return response['path']
else:
job_directory = self.client.job_directory
assert job_directory, "job directory required for action %s" % action
if not name:
name = basename(path)
self.__add_remote_staging_input(action, name, type)
def get_path():
return job_directory.calculate_path(name, type)
register = self.rewrite_paths or type == 'tool' # Even if inputs not rewritten, tool must be.
if register:
self.register_rewrite(path, get_path(), type, force=True)
elif self.rewrite_paths:
path_rewrite = action.path_rewrite(self.path_helper)
if path_rewrite:
self.register_rewrite(path, path_rewrite, type, force=True)
# else: # No action for this file
def __add_remote_staging_input(self, action, name, type):
input_dict = dict(
name=name,
type=type,
action=action.to_dict(),
)
self.remote_staging_actions.append(input_dict)
def __action_for_transfer(self, path, type, contents):
if contents:
# If contents loaded in memory, no need to write out file and copy,
# just transfer.
action = MessageAction(contents=contents, client=self.client)
else:
if not exists(path):
message = "handle_tranfer called on non-existent file - [%s]" % path
log.warn(message)
raise Exception(message)
action = self.__action(path, type)
return action
def register_rewrite(self, local_path, remote_path, type, force=False):
action = self.__action(local_path, type)
if action.staging_needed or force:
self.file_renames[local_path] = remote_path
def rewrite_input_paths(self):
"""
For each file that has been transferred and renamed, updated
command_line and configfiles to reflect that rewrite.
"""
for local_path, remote_path in self.file_renames.items():
self.job_inputs.rewrite_paths(local_path, remote_path)
def __action(self, path, type):
return self.action_mapper.action(path, type)
def _read(path):
"""
Utility method to quickly read small files (config files and tool
wrappers) into memory as bytes.
"""
input = open(path, "r", encoding="utf-8")
try:
return input.read()
finally:
input.close()
__all__ = ['submit_job']
-49
View File
@@ -1,49 +0,0 @@
from .standard import Urllib2Transport
from .curl import PycurlTransport
import os
from .ssh import rsync_get_file, scp_get_file
from .ssh import rsync_post_file, scp_post_file
from .curl import curl_available
from .requests import requests_multipart_post_available
if curl_available:
from .curl import get_file
from .curl import post_file
elif requests_multipart_post_available:
from .requests import get_file
from .requests import post_file
else:
from .poster import get_file
from .poster import post_file
def get_transport(transport_type=None, os_module=os):
transport_type = _get_transport_type(transport_type, os_module)
if transport_type == 'urllib':
transport = Urllib2Transport()
else:
transport = PycurlTransport()
return transport
def _get_transport_type(transport_type, os_module):
if not transport_type:
use_curl = os_module.getenv('PULSAR_CURL_TRANSPORT', "0")
# If PULSAR_CURL_TRANSPORT is unset or set to 0, use default,
# else use curl.
if use_curl.isdigit() and not int(use_curl):
transport_type = 'urllib'
else:
transport_type = 'curl'
return transport_type
__all__ = [
'get_transport',
'get_file',
'post_file',
'rsync_get_file',
'rsync_post_file',
'scp_get_file',
'scp_post_file'
]
-110
View File
@@ -1,110 +0,0 @@
import logging
from six import string_types
from six import BytesIO
try:
from pycurl import Curl, HTTP_CODE
curl_available = True
except ImportError:
curl_available = False
import os.path
PYCURL_UNAVAILABLE_MESSAGE = \
"You are attempting to use the Pycurl version of the Pulsar client but pycurl is unavailable."
NO_SUCH_FILE_MESSAGE = "Attempt to post file %s to URL %s, but file does not exist."
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):
def execute(self, url, method=None, data=None, input_path=None, output_path=None):
buf = _open_output(output_path)
try:
c = _new_curl_object_for_url(url)
c.setopt(c.WRITEFUNCTION, buf.write)
if method:
c.setopt(c.CUSTOMREQUEST, method)
if input_path:
c.setopt(c.UPLOAD, 1)
c.setopt(c.READFUNCTION, open(input_path, 'rb').read)
filesize = os.path.getsize(input_path)
c.setopt(c.INFILESIZE, filesize)
if data:
c.setopt(c.POST, 1)
if isinstance(data, string_types):
data = data.encode('UTF-8')
c.setopt(c.POSTFIELDS, data)
c.perform()
if not output_path:
return buf.getvalue()
finally:
buf.close()
def post_file(url, path):
if not os.path.exists(path):
# pycurl doesn't always produce a great exception for this,
# wrap it in a better one.
message = NO_SUCH_FILE_MESSAGE % (path, url)
raise Exception(message)
c = _new_curl_object_for_url(url)
c.setopt(c.HTTPPOST, [("file", (c.FORM_FILE, path.encode('ascii')))])
c.perform()
status_code = c.getinfo(HTTP_CODE)
if int(status_code) != 200:
message = POST_FAILED_MESSAGE % (url, status_code)
raise Exception(message)
def get_file(url, 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 = 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, mode='wb'):
return open(output_path, mode) if output_path else BytesIO()
def _new_curl_object_for_url(url):
c = _new_curl_object()
c.setopt(c.URL, url.encode('ascii'))
return c
def _new_curl_object():
try:
return Curl()
except NameError:
raise ImportError(PYCURL_UNAVAILABLE_MESSAGE)
__all__ = [
'PycurlTransport',
'post_file',
'get_file'
]
-53
View File
@@ -1,53 +0,0 @@
from __future__ import absolute_import
import logging
try:
from urllib2 import urlopen
except ImportError:
from urllib.request import urlopen
try:
from urllib2 import Request
except ImportError:
from urllib.request import Request
try:
import poster
except ImportError:
poster = None
POSTER_UNAVAILABLE_MESSAGE = "Pulsar configured to use poster module - but it is unavailable. Please install poster."
log = logging.getLogger(__name__)
if poster is not None:
poster.streaminghttp.register_openers()
def post_file(url, path):
__ensure_poster()
try:
datagen, headers = poster.encode.multipart_encode({"file": open(path, "rb")})
request = Request(url, datagen, headers)
return urlopen(request).read()
except:
log.exception("problem")
raise
def get_file(url, path):
__ensure_poster()
request = Request(url=url)
response = urlopen(request)
with open(path, 'wb') as output:
while True:
buffer = response.read(1024)
if not buffer:
break
output.write(buffer)
def __ensure_poster():
if poster is None:
raise ImportError(POSTER_UNAVAILABLE_MESSAGE)
-47
View File
@@ -1,47 +0,0 @@
from __future__ import absolute_import
import logging
try:
import requests
except ImportError:
requests = None
try:
import requests_toolbelt
requests_multipart_post_available = True
except ImportError:
requests_multipart_post_available = False
requests_toolbelt = None
REQUESTS_UNAVAILABLE_MESSAGE = "Pulsar configured to use requests module - but it is unavailable. Please install requests."
REQUESTS_TOOLBELT_UNAVAILABLE_MESSAGE = "Pulsar configured to use requests_toolbelt module - but it is unavailable. Please install requests_toolbelt."
log = logging.getLogger(__name__)
def post_file(url, path):
if requests_toolbelt is None:
raise ImportError(REQUESTS_TOOLBELT_UNAVAILABLE_MESSAGE)
__ensure_requests()
m = requests_toolbelt.MultipartEncoder(
fields={'file': ('filename', open(path, 'rb'))}
)
requests.post(url, data=m, headers={'Content-Type': m.content_type})
def get_file(url, path):
__ensure_requests()
r = requests.get(url, stream=True)
with open(path, 'wb') as f:
for chunk in r.iter_content(chunk_size=1024):
if chunk: # filter out keep-alive new chunks
f.write(chunk)
f.flush()
def __ensure_requests():
if requests is None:
raise ImportError(REQUESTS_UNAVAILABLE_MESSAGE)
-79
View File
@@ -1,79 +0,0 @@
import subprocess
import os
SSH_OPTIONS = ['-o', 'StrictHostKeyChecking=no', '-o', 'PreferredAuthentications=publickey', '-o', 'PubkeyAuthentication=yes']
def rsync_get_file(uri_from, uri_to, user, host, port, key):
cmd = [
'rsync',
'-e',
'ssh -i %s -p %s %s' % (key, port, ' '.join(SSH_OPTIONS)),
'%s@%s:%s' % (user, host, uri_from),
uri_to,
]
_call(cmd)
def rsync_post_file(uri_from, uri_to, user, host, port, key):
_ensure_dir(uri_to, key, port, user, host)
cmd = [
'rsync',
'-e',
'ssh -i %s -p %s %s' % (key, port, ' '.join(SSH_OPTIONS)),
uri_from,
'%s@%s:%s' % (user, host, uri_to),
]
_call(cmd)
def scp_get_file(uri_from, uri_to, user, host, port, key):
cmd = [
'scp',
'-P', str(port),
'-i', key
] + SSH_OPTIONS + [
'%s@%s:%s' % (user, host, uri_from),
uri_to,
]
_call(cmd)
def scp_post_file(uri_from, uri_to, user, host, port, key):
_ensure_dir(uri_to, key, port, user, host)
cmd = [
'scp',
'-P', str(port),
'-i', key,
] + SSH_OPTIONS + [
uri_from,
'%s@%s:%s' % (user, host, uri_to),
]
_call(cmd)
def _ensure_dir(uri_to, key, port, user, host):
directory = os.path.dirname(uri_to)
cmd = [
'ssh',
'-i', key,
'-p', str(port),
] + SSH_OPTIONS + [
'%s@%s' % (user, host),
'mkdir', '-p', directory,
]
_call(cmd)
def _call(cmd):
exit_code = subprocess.check_call(cmd)
if exit_code != 0:
raise Exception("%s exited with code %s" % (cmd[0], exit_code))
___all__ = [
'rsync_post_file',
'rsync_get_file',
'scp_post_file',
'scp_get_file'
]
-53
View File
@@ -1,53 +0,0 @@
"""
Pulsar HTTP Client layer based on Python Standard Library (urllib2)
"""
from __future__ import with_statement
from os.path import getsize
import mmap
try:
from urllib2 import urlopen
except ImportError:
from urllib.request import urlopen
try:
from urllib2 import Request
except ImportError:
from urllib.request import Request
class Urllib2Transport(object):
def _url_open(self, request, data):
return urlopen(request, data)
def execute(self, url, method=None, data=None, input_path=None, output_path=None):
request = self.__request(url, data, method)
input = None
try:
if input_path:
size = getsize(input_path)
if size:
input = open(input_path, 'rb')
data = mmap.mmap(input.fileno(), 0, access=mmap.ACCESS_READ)
else:
data = b""
request.add_header('Content-Length', str(size))
response = self._url_open(request, data)
finally:
if input:
input.close()
if output_path:
with open(output_path, 'wb') as output:
while True:
buffer = response.read(1024)
if not buffer:
break
output.write(buffer)
return response
else:
return response.read()
def __request(self, url, data, method):
request = Request(url=url, data=data)
if method:
request.get_method = lambda: method
return request
-281
View File
@@ -1,281 +0,0 @@
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 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.encode('utf-8'))
return m.hexdigest()
def copy(source, destination):
""" Copy file from source to destination if needed (skip if source
is destination).
"""
source = os.path.abspath(source)
destination = os.path.abspath(destination)
if source != destination:
shutil.copyfile(source, destination)
def ensure_directory(file_path):
directory = os.path.dirname(file_path)
if not os.path.exists(directory):
os.makedirs(directory)
def directory_files(directory):
"""
>>> from tempfile import mkdtemp
>>> from shutil import rmtree
>>> from os.path import join
>>> from os import makedirs
>>> tempdir = mkdtemp()
>>> with open(join(tempdir, "moo"), "w") as f: pass
>>> directory_files(tempdir)
['moo']
>>> subdir = join(tempdir, "cow", "sub1")
>>> makedirs(subdir)
>>> with open(join(subdir, "subfile1"), "w") as f: pass
>>> with open(join(subdir, "subfile2"), "w") as f: pass
>>> sorted(directory_files(tempdir))
['cow/sub1/subfile1', 'cow/sub1/subfile2', 'moo']
>>> rmtree(tempdir)
"""
contents = []
for path, _, files in walk(directory):
relative_path = relpath(path, directory)
for name in files:
# Return file1.txt, dataset_1_files/image.png, etc... don't
# include . in path.
if relative_path != curdir:
contents.append(join(relative_path, name))
else:
contents.append(name)
return contents
def filter_destination_params(destination_params, prefix):
destination_params = destination_params or {}
return dict([(key[len(prefix):], destination_params[key])
for key in destination_params
if key.startswith(prefix)])
def to_base64_json(data):
"""
>>> enc = to_base64_json(dict(a=5))
>>> dec = from_base64_json(enc)
>>> dec["a"]
5
"""
dumped = json_dumps(data)
return b64encode(dumped)
def from_base64_json(data):
return json.loads(b64decode(data))
class PathHelper(object):
'''
>>> import posixpath
>>> # Forcing local path to posixpath because Pulsar designed to be used with
>>> # posix client.
>>> posix_path_helper = PathHelper("/", local_path_module=posixpath)
>>> windows_slash = "\\\\"
>>> len(windows_slash)
1
>>> nt_path_helper = PathHelper(windows_slash, local_path_module=posixpath)
>>> posix_path_helper.remote_name("moo/cow")
'moo/cow'
>>> nt_path_helper.remote_name("moo/cow")
'moo\\\\cow'
>>> posix_path_helper.local_name("moo/cow")
'moo/cow'
>>> nt_path_helper.local_name("moo\\\\cow")
'moo/cow'
>>> posix_path_helper.from_posix_with_new_base("/galaxy/data/bowtie/hg19.fa", "/galaxy/data/", "/work/galaxy/data")
'/work/galaxy/data/bowtie/hg19.fa'
>>> posix_path_helper.from_posix_with_new_base("/galaxy/data/bowtie/hg19.fa", "/galaxy/data", "/work/galaxy/data")
'/work/galaxy/data/bowtie/hg19.fa'
>>> posix_path_helper.from_posix_with_new_base("/galaxy/data/bowtie/hg19.fa", "/galaxy/data", "/work/galaxy/data/")
'/work/galaxy/data/bowtie/hg19.fa'
'''
def __init__(self, separator, local_path_module=os.path):
self.separator = separator
self.local_join = local_path_module.join
self.local_sep = local_path_module.sep
def remote_name(self, local_name):
return self.remote_join(*local_name.split(self.local_sep))
def local_name(self, remote_name):
return self.local_join(*remote_name.split(self.separator))
def remote_join(self, *args):
return self.separator.join(args)
def from_posix_with_new_base(self, posix_path, old_base, new_base):
# TODO: Test with new_base as a windows path against nt_path_helper.
if old_base.endswith("/"):
old_base = old_base[:-1]
if not posix_path.startswith(old_base):
message_template = "Cannot compute new path for file %s, does not start with %s."
message = message_template % (posix_path, old_base)
raise Exception(message)
stripped_path = posix_path[len(old_base):]
while stripped_path.startswith("/"):
stripped_path = stripped_path[1:]
path_parts = stripped_path.split(self.separator)
if new_base.endswith(self.separator):
new_base = new_base[:-len(self.separator)]
return self.remote_join(new_base, *path_parts)
class TransferEventManager(object):
def __init__(self):
self.events = WeakValueDictionary(dict())
self.events_lock = Lock()
def acquire_event(self, path, force_clear=False):
with self.events_lock:
if path in self.events:
event_holder = self.events[path]
else:
event_holder = EventHolder(Event(), path, self)
self.events[path] = event_holder
if force_clear:
event_holder.event.clear()
return event_holder
class EventHolder(object):
def __init__(self, event, path, condition_manager):
self.event = event
self.path = path
self.condition_manager = condition_manager
self.failed = False
def release(self):
self.event.set()
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