diff --git a/.ci/flake8_wrapper.sh b/.ci/flake8_wrapper.sh index d087ebb23d3..4924f630a8a 100755 --- a/.ci/flake8_wrapper.sh +++ b/.ci/flake8_wrapper.sh @@ -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/ diff --git a/.ci/py3_sources.txt b/.ci/py3_sources.txt index 75869f8cb23..75546357121 100644 --- a/.ci/py3_sources.txt +++ b/.ci/py3_sources.txt @@ -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 diff --git a/doc/source/lib/modules.rst b/doc/source/lib/modules.rst index 2f834b22134..a1f23167148 100644 --- a/doc/source/lib/modules.rst +++ b/doc/source/lib/modules.rst @@ -10,5 +10,4 @@ lib log_tempfile mimeparse psyco_full - pulsar tool_shed diff --git a/doc/source/lib/pulsar.client.rst b/doc/source/lib/pulsar.client.rst deleted file mode 100644 index 3a29f55ffae..00000000000 --- a/doc/source/lib/pulsar.client.rst +++ /dev/null @@ -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: - - diff --git a/doc/source/lib/pulsar.client.staging.rst b/doc/source/lib/pulsar.client.staging.rst deleted file mode 100644 index 4d33f6676b8..00000000000 --- a/doc/source/lib/pulsar.client.staging.rst +++ /dev/null @@ -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: - - diff --git a/doc/source/lib/pulsar.client.transport.rst b/doc/source/lib/pulsar.client.transport.rst deleted file mode 100644 index 8b68dbb15a5..00000000000 --- a/doc/source/lib/pulsar.client.transport.rst +++ /dev/null @@ -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: - - diff --git a/doc/source/lib/pulsar.rst b/doc/source/lib/pulsar.rst deleted file mode 100644 index c22c0a95bdb..00000000000 --- a/doc/source/lib/pulsar.rst +++ /dev/null @@ -1,15 +0,0 @@ -pulsar package -============== - -.. automodule:: pulsar - :members: - :undoc-members: - :show-inheritance: - -Subpackages ------------ - -.. toctree:: - - pulsar.client - diff --git a/lib/galaxy/dependencies/pinned-requirements.txt b/lib/galaxy/dependencies/pinned-requirements.txt index 5a5c312422a..2ee38fe4b4e 100644 --- a/lib/galaxy/dependencies/pinned-requirements.txt +++ b/lib/galaxy/dependencies/pinned-requirements.txt @@ -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 diff --git a/lib/galaxy/dependencies/requirements.txt b/lib/galaxy/dependencies/requirements.txt index 598dcb4b401..23ea3cd86ae 100644 --- a/lib/galaxy/dependencies/requirements.txt +++ b/lib/galaxy/dependencies/requirements.txt @@ -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 diff --git a/lib/galaxy/jobs/runners/pulsar.py b/lib/galaxy/jobs/runners/pulsar.py index b52a777c79f..1eb56e59ba1 100644 --- a/lib/galaxy/jobs/runners/pulsar.py +++ b/lib/galaxy/jobs/runners/pulsar.py @@ -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, diff --git a/lib/pulsar/__init__.py b/lib/pulsar/__init__.py deleted file mode 100644 index e69de29bb2d..00000000000 diff --git a/lib/pulsar/client/__init__.py b/lib/pulsar/client/__init__.py deleted file mode 100644 index 5ada39334bc..00000000000 --- a/lib/pulsar/client/__init__.py +++ /dev/null @@ -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 `_ -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:: - - = pulsar://http://: - -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', -] diff --git a/lib/pulsar/client/action_mapper.py b/lib/pulsar/client/action_mapper.py deleted file mode 100644 index 7d19eb55054..00000000000 --- a/lib/pulsar/client/action_mapper.py +++ /dev/null @@ -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 -] diff --git a/lib/pulsar/client/amqp_exchange.py b/lib/pulsar/client/amqp_exchange.py deleted file mode 100644 index 03ad85ffa95..00000000000 --- a/lib/pulsar/client/amqp_exchange.py +++ /dev/null @@ -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 diff --git a/lib/pulsar/client/amqp_exchange_factory.py b/lib/pulsar/client/amqp_exchange_factory.py deleted file mode 100644 index 72f917860c7..00000000000 --- a/lib/pulsar/client/amqp_exchange_factory.py +++ /dev/null @@ -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 diff --git a/lib/pulsar/client/client.py b/lib/pulsar/client/client.py deleted file mode 100644 index d7f3439bac9..00000000000 --- a/lib/pulsar/client/client.py +++ /dev/null @@ -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 - ) diff --git a/lib/pulsar/client/config_util.py b/lib/pulsar/client/config_util.py deleted file mode 100644 index a128ca1dda3..00000000000 --- a/lib/pulsar/client/config_util.py +++ /dev/null @@ -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, -} diff --git a/lib/pulsar/client/decorators.py b/lib/pulsar/client/decorators.py deleted file mode 100644 index 659b5c633d4..00000000000 --- a/lib/pulsar/client/decorators.py +++ /dev/null @@ -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 diff --git a/lib/pulsar/client/destination.py b/lib/pulsar/client/destination.py deleted file mode 100644 index cee7a40f630..00000000000 --- a/lib/pulsar/client/destination.py +++ /dev/null @@ -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) diff --git a/lib/pulsar/client/interface.py b/lib/pulsar/client/interface.py deleted file mode 100644 index c346c573890..00000000000 --- a/lib/pulsar/client/interface.py +++ /dev/null @@ -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 diff --git a/lib/pulsar/client/job_directory.py b/lib/pulsar/client/job_directory.py deleted file mode 100644 index c4026aead9a..00000000000 --- a/lib/pulsar/client/job_directory.py +++ /dev/null @@ -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) diff --git a/lib/pulsar/client/manager.py b/lib/pulsar/client/manager.py deleted file mode 100644 index d03b60a8411..00000000000 --- a/lib/pulsar/client/manager.py +++ /dev/null @@ -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' -] diff --git a/lib/pulsar/client/object_client.py b/lib/pulsar/client/object_client.py deleted file mode 100644 index a9d795f1b6c..00000000000 --- a/lib/pulsar/client/object_client.py +++ /dev/null @@ -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) diff --git a/lib/pulsar/client/path_mapper.py b/lib/pulsar/client/path_mapper.py deleted file mode 100644 index 97c1a770311..00000000000 --- a/lib/pulsar/client/path_mapper.py +++ /dev/null @@ -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'] diff --git a/lib/pulsar/client/setup_handler.py b/lib/pulsar/client/setup_handler.py deleted file mode 100644 index 2cc9d1afcc4..00000000000 --- a/lib/pulsar/client/setup_handler.py +++ /dev/null @@ -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'] diff --git a/lib/pulsar/client/staging/__init__.py b/lib/pulsar/client/staging/__init__.py deleted file mode 100644 index c121bf28be7..00000000000 --- a/lib/pulsar/client/staging/__init__.py +++ /dev/null @@ -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)) diff --git a/lib/pulsar/client/staging/down.py b/lib/pulsar/client/staging/down.py deleted file mode 100644 index 37f2e4965c3..00000000000 --- a/lib/pulsar/client/staging/down.py +++ /dev/null @@ -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'] diff --git a/lib/pulsar/client/staging/up.py b/lib/pulsar/client/staging/up.py deleted file mode 100644 index 46ebe319290..00000000000 --- a/lib/pulsar/client/staging/up.py +++ /dev/null @@ -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'] diff --git a/lib/pulsar/client/transport/__init__.py b/lib/pulsar/client/transport/__init__.py deleted file mode 100644 index 4ff9c7df25f..00000000000 --- a/lib/pulsar/client/transport/__init__.py +++ /dev/null @@ -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' -] diff --git a/lib/pulsar/client/transport/curl.py b/lib/pulsar/client/transport/curl.py deleted file mode 100644 index 633a57dbac6..00000000000 --- a/lib/pulsar/client/transport/curl.py +++ /dev/null @@ -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' -] diff --git a/lib/pulsar/client/transport/poster.py b/lib/pulsar/client/transport/poster.py deleted file mode 100644 index f8ea3f7d9d3..00000000000 --- a/lib/pulsar/client/transport/poster.py +++ /dev/null @@ -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) diff --git a/lib/pulsar/client/transport/requests.py b/lib/pulsar/client/transport/requests.py deleted file mode 100644 index 1e5c0c4ff97..00000000000 --- a/lib/pulsar/client/transport/requests.py +++ /dev/null @@ -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) diff --git a/lib/pulsar/client/transport/ssh.py b/lib/pulsar/client/transport/ssh.py deleted file mode 100644 index b0c7a3b13fd..00000000000 --- a/lib/pulsar/client/transport/ssh.py +++ /dev/null @@ -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' -] diff --git a/lib/pulsar/client/transport/standard.py b/lib/pulsar/client/transport/standard.py deleted file mode 100644 index 9a5c46372e4..00000000000 --- a/lib/pulsar/client/transport/standard.py +++ /dev/null @@ -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 diff --git a/lib/pulsar/client/util.py b/lib/pulsar/client/util.py deleted file mode 100644 index 35696befdc1..00000000000 --- a/lib/pulsar/client/util.py +++ /dev/null @@ -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