Eliminate LWR from Galaxy.

This commit is contained in:
John Chilton
2015-10-06 12:33:36 +01:00
parent 3687d45a03
commit c3ef5467c7
30 changed files with 6 additions and 4057 deletions
@@ -1,132 +0,0 @@
galaxy.jobs.runners.lwr_client package
======================================
.. automodule:: galaxy.jobs.runners.lwr_client
:members:
:undoc-members:
:show-inheritance:
Subpackages
-----------
.. toctree::
galaxy.jobs.runners.lwr_client.staging
galaxy.jobs.runners.lwr_client.transport
Submodules
----------
galaxy.jobs.runners.lwr_client.action_mapper module
---------------------------------------------------
.. automodule:: galaxy.jobs.runners.lwr_client.action_mapper
:members:
:undoc-members:
:show-inheritance:
galaxy.jobs.runners.lwr_client.amqp_exchange module
---------------------------------------------------
.. automodule:: galaxy.jobs.runners.lwr_client.amqp_exchange
:members:
:undoc-members:
:show-inheritance:
galaxy.jobs.runners.lwr_client.amqp_exchange_factory module
-----------------------------------------------------------
.. automodule:: galaxy.jobs.runners.lwr_client.amqp_exchange_factory
:members:
:undoc-members:
:show-inheritance:
galaxy.jobs.runners.lwr_client.client module
--------------------------------------------
.. automodule:: galaxy.jobs.runners.lwr_client.client
:members:
:undoc-members:
:show-inheritance:
galaxy.jobs.runners.lwr_client.config_util module
-------------------------------------------------
.. automodule:: galaxy.jobs.runners.lwr_client.config_util
:members:
:undoc-members:
:show-inheritance:
galaxy.jobs.runners.lwr_client.decorators module
------------------------------------------------
.. automodule:: galaxy.jobs.runners.lwr_client.decorators
:members:
:undoc-members:
:show-inheritance:
galaxy.jobs.runners.lwr_client.destination module
-------------------------------------------------
.. automodule:: galaxy.jobs.runners.lwr_client.destination
:members:
:undoc-members:
:show-inheritance:
galaxy.jobs.runners.lwr_client.interface module
-----------------------------------------------
.. automodule:: galaxy.jobs.runners.lwr_client.interface
:members:
:undoc-members:
:show-inheritance:
galaxy.jobs.runners.lwr_client.job_directory module
---------------------------------------------------
.. automodule:: galaxy.jobs.runners.lwr_client.job_directory
:members:
:undoc-members:
:show-inheritance:
galaxy.jobs.runners.lwr_client.manager module
---------------------------------------------
.. automodule:: galaxy.jobs.runners.lwr_client.manager
:members:
:undoc-members:
:show-inheritance:
galaxy.jobs.runners.lwr_client.object_client module
---------------------------------------------------
.. automodule:: galaxy.jobs.runners.lwr_client.object_client
:members:
:undoc-members:
:show-inheritance:
galaxy.jobs.runners.lwr_client.path_mapper module
-------------------------------------------------
.. automodule:: galaxy.jobs.runners.lwr_client.path_mapper
:members:
:undoc-members:
:show-inheritance:
galaxy.jobs.runners.lwr_client.setup_handler module
---------------------------------------------------
.. automodule:: galaxy.jobs.runners.lwr_client.setup_handler
:members:
:undoc-members:
:show-inheritance:
galaxy.jobs.runners.lwr_client.util module
------------------------------------------
.. automodule:: galaxy.jobs.runners.lwr_client.util
:members:
:undoc-members:
:show-inheritance:
@@ -1,28 +0,0 @@
galaxy.jobs.runners.lwr_client.staging package
==============================================
.. automodule:: galaxy.jobs.runners.lwr_client.staging
:members:
:undoc-members:
:show-inheritance:
Submodules
----------
galaxy.jobs.runners.lwr_client.staging.down module
--------------------------------------------------
.. automodule:: galaxy.jobs.runners.lwr_client.staging.down
:members:
:undoc-members:
:show-inheritance:
galaxy.jobs.runners.lwr_client.staging.up module
------------------------------------------------
.. automodule:: galaxy.jobs.runners.lwr_client.staging.up
:members:
:undoc-members:
:show-inheritance:
@@ -1,28 +0,0 @@
galaxy.jobs.runners.lwr_client.transport package
================================================
.. automodule:: galaxy.jobs.runners.lwr_client.transport
:members:
:undoc-members:
:show-inheritance:
Submodules
----------
galaxy.jobs.runners.lwr_client.transport.curl module
----------------------------------------------------
.. automodule:: galaxy.jobs.runners.lwr_client.transport.curl
:members:
:undoc-members:
:show-inheritance:
galaxy.jobs.runners.lwr_client.transport.standard module
--------------------------------------------------------
.. automodule:: galaxy.jobs.runners.lwr_client.transport.standard
:members:
:undoc-members:
:show-inheritance:
-8
View File
@@ -11,7 +11,6 @@ Subpackages
.. toctree::
galaxy.jobs.runners.lwr_client
galaxy.jobs.runners.state_handlers
galaxy.jobs.runners.util
@@ -50,13 +49,6 @@ galaxy.jobs.runners.local module
:undoc-members:
:show-inheritance:
galaxy.jobs.runners.lwr module
------------------------------
.. automodule:: galaxy.jobs.runners.lwr
:members:
:undoc-members:
:show-inheritance:
galaxy.jobs.runners.pbs module
------------------------------
+2 -3
View File
@@ -289,9 +289,8 @@ class JobConfiguration( object ):
"""
log.debug('Loading job configuration from %s' % self.app.config.config_file)
# Always load local and lwr
self.runner_plugins = [dict(id='local', load='local', workers=self.app.config.local_job_queue_workers),
dict(id='lwr', load='lwr', workers=self.app.config.cluster_job_queue_workers)]
# Always load local
self.runner_plugins = [dict(id='local', load='local', workers=self.app.config.local_job_queue_workers)]
# Load tasks if configured
if self.app.config.use_tasked_jobs:
self.runner_plugins.append(dict(id='tasks', load='tasks', workers=self.app.config.local_task_queue_workers))
+2 -4
View File
@@ -188,8 +188,6 @@ class BaseJobRunner( object ):
raise NotImplementedError()
def build_command_line( self, job_wrapper, include_metadata=False, include_work_dir_outputs=True ):
# TODO: Eliminate extra kwds no longer used (since LWR skips
# abstraction and calls build_command directly).
container = self._find_container( job_wrapper )
return build_command(
self,
@@ -248,8 +246,8 @@ class BaseJobRunner( object ):
def _handle_metadata_externally( self, job_wrapper, resolve_requirements=False ):
"""
Set metadata externally. Used by the local and lwr job runners where this
shouldn't be attached to command-line to execute.
Set metadata externally. Used by the Pulsar job runner where this
shouldn't be attached to command line to execute.
"""
# run the metadata setting script here
# this is terminate-able when output dataset/job is deleted
-648
View File
@@ -1,648 +0,0 @@
import logging
from galaxy import model
from galaxy.jobs.runners import AsynchronousJobState, AsynchronousJobRunner
from galaxy.jobs import ComputeEnvironment
from galaxy.jobs import JobDestination
from galaxy.jobs.command_factory import build_command
from galaxy.tools.deps import dependencies
from galaxy.util import string_as_bool_or_none
from galaxy.util.bunch import Bunch
from galaxy.util import specs
import errno
from time import sleep
import os
from .lwr_client import build_client_manager
from .lwr_client import url_to_destination_params
from .lwr_client import finish_job as lwr_finish_job
from .lwr_client import submit_job as lwr_submit_job
from .lwr_client import ClientJobDescription
from .lwr_client import LwrOutputs
from .lwr_client import ClientOutputs
from .lwr_client import PathMapper
log = logging.getLogger( __name__ )
__all__ = [ 'LwrJobRunner' ]
NO_REMOTE_GALAXY_FOR_METADATA_MESSAGE = "LWR misconfiguration - LWR client configured to set metadata remotely, but remote LWR isn't properly configured with a galaxy_home directory."
NO_REMOTE_DATATYPES_CONFIG = "LWR client is configured to use remote datatypes configuration when setting metadata externally, but LWR is not configured with this information. Defaulting to datatypes_conf.xml."
GENERIC_REMOTE_ERROR = "Failed to communicate with remote job server."
# Is there a good way to infer some default for this? Can only use
# url_for from web threads. https://gist.github.com/jmchilton/9098762
DEFAULT_GALAXY_URL = "http://localhost:8080"
LWR_PARAM_SPECS = dict(
transport=dict(
map=specs.to_str_or_none,
valid=specs.is_in("urllib", "curl", None),
default=None
),
cache=dict(
map=specs.to_bool_or_none,
default=None,
),
url=dict(
map=specs.to_str_or_none,
default=None,
),
galaxy_url=dict(
map=specs.to_str_or_none,
default=DEFAULT_GALAXY_URL,
),
manager=dict(
map=specs.to_str_or_none,
default=None,
),
amqp_consumer_timeout=dict(
map=lambda val: None if val == "None" else float(val),
default=None,
),
amqp_connect_ssl_ca_certs=dict(
map=specs.to_str_or_none,
default=None,
),
amqp_connect_ssl_keyfile=dict(
map=specs.to_str_or_none,
default=None,
),
amqp_connect_ssl_certfile=dict(
map=specs.to_str_or_none,
default=None,
),
amqp_connect_ssl_cert_reqs=dict(
map=specs.to_str_or_none,
default=None,
),
# http://kombu.readthedocs.org/en/latest/reference/kombu.html#kombu.Producer.publish
amqp_publish_retry=dict(
map=specs.to_bool,
default=False,
),
amqp_publish_priority=dict(
map=int,
valid=lambda x: 0 <= x and x <= 9,
default=0,
),
# http://kombu.readthedocs.org/en/latest/reference/kombu.html#kombu.Exchange.delivery_mode
amqp_publish_delivery_mode=dict(
map=str,
valid=specs.is_in("transient", "persistent"),
default="persistent",
),
amqp_publish_retry_max_retries=dict(
map=int,
default=None,
),
amqp_publish_retry_interval_start=dict(
map=int,
default=None,
),
amqp_publish_retry_interval_step=dict(
map=int,
default=None,
),
amqp_publish_retry_interval_max=dict(
map=int,
default=None,
),
)
class LwrJobRunner( AsynchronousJobRunner ):
"""
LWR Job Runner
"""
runner_name = "LWRRunner"
def __init__( self, app, nworkers, **kwds ):
"""Start the job runner """
super( LwrJobRunner, self ).__init__( app, nworkers, runner_param_specs=LWR_PARAM_SPECS, **kwds )
self._init_worker_threads()
galaxy_url = self.runner_params.galaxy_url
if galaxy_url:
galaxy_url = galaxy_url.rstrip("/")
self.galaxy_url = galaxy_url
self.__init_client_manager()
if self.runner_params.url:
# This is a message queue driven runner, don't monitor
# just setup required callback.
self.client_manager.ensure_has_status_update_callback(self.__async_update)
else:
self._init_monitor_thread()
def __init_client_manager( self ):
client_manager_kwargs = {}
for kwd in 'url', 'manager', 'cache', 'transport':
client_manager_kwargs[ kwd ] = self.runner_params[ kwd ]
for kwd in self.runner_params.keys():
if kwd.startswith( 'amqp_' ):
client_manager_kwargs[ kwd ] = self.runner_params[ kwd ]
self.client_manager = build_client_manager(**client_manager_kwargs)
def url_to_destination( self, url ):
"""Convert a legacy URL to a job destination"""
return JobDestination( runner="lwr", params=url_to_destination_params( url ) )
def check_watched_item(self, job_state):
try:
client = self.get_client_from_state(job_state)
status = client.get_status()
except Exception:
# An orphaned job was put into the queue at app startup, so remote server went down
# either way we are done I guess.
self.mark_as_finished(job_state)
return None
job_state = self.__update_job_state_for_lwr_status(job_state, status)
return job_state
def __update_job_state_for_lwr_status(self, job_state, lwr_status):
if lwr_status == "complete":
self.mark_as_finished(job_state)
return None
if lwr_status == "failed":
self.fail_job(job_state)
return None
if lwr_status == "running" and not job_state.running:
job_state.running = True
job_state.job_wrapper.change_state( model.Job.states.RUNNING )
return job_state
def __async_update( self, full_status ):
while not hasattr( self.app, 'job_manager' ):
# The status update thread can start consuming before app is done initializing
log.debug( 'Received a status update message before app is initialized, waiting 5 seconds' )
sleep( 5 )
job_id = None
try:
job_id = full_status[ "job_id" ]
job, job_wrapper = self.app.job_manager.job_handler.job_queue.job_pair_for_id( job_id )
job_state = self.__job_state( job, job_wrapper )
self.__update_job_state_for_lwr_status(job_state, full_status["status"])
except Exception:
log.exception( "Failed to update LWR job status for job_id %s" % job_id )
raise
# Nothing else to do? - Attempt to fail the job?
def queue_job(self, job_wrapper):
job_destination = job_wrapper.job_destination
command_line, client, remote_job_config, compute_environment = self.__prepare_job( job_wrapper, job_destination )
if not command_line:
return
try:
dependencies_description = LwrJobRunner.__dependencies_description( client, job_wrapper )
rewrite_paths = not LwrJobRunner.__rewrite_parameters( client )
unstructured_path_rewrites = {}
if compute_environment:
unstructured_path_rewrites = compute_environment.unstructured_path_rewrites
client_job_description = ClientJobDescription(
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,
tool=job_wrapper.tool,
config_files=job_wrapper.extra_filenames,
dependencies_description=dependencies_description,
env=client.env,
rewrite_paths=rewrite_paths,
arbitrary_files=unstructured_path_rewrites,
)
job_id = lwr_submit_job(client, client_job_description, remote_job_config)
log.info("lwr job submitted with job_id %s" % job_id)
job_wrapper.set_job_destination( job_destination, job_id )
job_wrapper.change_state( model.Job.states.QUEUED )
except Exception:
job_wrapper.fail( "failure running job", exception=True )
log.exception("failure running job %d" % job_wrapper.job_id)
return
lwr_job_state = AsynchronousJobState()
lwr_job_state.job_wrapper = job_wrapper
lwr_job_state.job_id = job_id
lwr_job_state.old_state = True
lwr_job_state.running = False
lwr_job_state.job_destination = job_destination
self.monitor_job(lwr_job_state)
def __prepare_job(self, job_wrapper, job_destination):
""" Build command-line and LWR client for this job. """
command_line = None
client = None
remote_job_config = None
compute_environment = None
try:
client = self.get_client_from_wrapper(job_wrapper)
tool = job_wrapper.tool
remote_job_config = client.setup(tool.id, tool.version)
rewrite_parameters = LwrJobRunner.__rewrite_parameters( client )
prepare_kwds = {}
if rewrite_parameters:
compute_environment = LwrComputeEnvironment( client, job_wrapper, remote_job_config )
prepare_kwds[ 'compute_environment' ] = compute_environment
job_wrapper.prepare( **prepare_kwds )
self.__prepare_input_files_locally(job_wrapper)
remote_metadata = LwrJobRunner.__remote_metadata( client )
dependency_resolution = LwrJobRunner.__dependency_resolution( client )
metadata_kwds = self.__build_metadata_configuration(client, job_wrapper, remote_metadata, remote_job_config)
remote_command_params = dict(
working_directory=remote_job_config['working_directory'],
metadata_kwds=metadata_kwds,
dependency_resolution=dependency_resolution,
)
remote_working_directory = remote_job_config['working_directory']
# TODO: Following defs work for LWR, always worked for LWR but should be
# calculated at some other level.
remote_job_directory = os.path.abspath(os.path.join(remote_working_directory, os.path.pardir))
remote_tool_directory = os.path.abspath(os.path.join(remote_job_directory, "tool_files"))
container = self._find_container(
job_wrapper,
compute_working_directory=remote_working_directory,
compute_tool_directory=remote_tool_directory,
compute_job_directory=remote_job_directory,
)
command_line = build_command(
self,
job_wrapper=job_wrapper,
container=container,
include_metadata=remote_metadata,
include_work_dir_outputs=False,
remote_command_params=remote_command_params,
)
except Exception:
job_wrapper.fail( "failure preparing job", exception=True )
log.exception("failure running job %d" % job_wrapper.job_id)
# If we were able to get a command line, run the job
if not command_line:
job_wrapper.finish( '', '' )
return command_line, client, remote_job_config, compute_environment
def __prepare_input_files_locally(self, job_wrapper):
"""Run task splitting commands locally."""
prepare_input_files_cmds = getattr(job_wrapper, 'prepare_input_files_cmds', None)
if prepare_input_files_cmds is not None:
for cmd in prepare_input_files_cmds: # run the commands to stage the input files
if 0 != os.system(cmd):
raise Exception('Error running file staging command: %s' % cmd)
job_wrapper.prepare_input_files_cmds = None # prevent them from being used in-line
def get_output_files(self, job_wrapper):
output_paths = job_wrapper.get_output_fnames()
return [ str( o ) for o in output_paths ] # Force job_path from DatasetPath objects.
def get_input_files(self, job_wrapper):
input_paths = job_wrapper.get_input_paths()
return [ str( i ) for i in input_paths ] # Force job_path from DatasetPath objects.
def get_client_from_wrapper(self, job_wrapper):
job_id = job_wrapper.job_id
if hasattr(job_wrapper, 'task_id'):
job_id = "%s_%s" % (job_id, job_wrapper.task_id)
params = job_wrapper.job_destination.params.copy()
for key, value in params.iteritems():
if value:
params[key] = model.User.expand_user_properties( job_wrapper.get_job().user, value )
env = getattr( job_wrapper.job_destination, "env", [] )
return self.get_client( params, job_id, env )
def get_client_from_state(self, job_state):
job_destination_params = job_state.job_destination.params
job_id = job_state.job_id
return self.get_client( job_destination_params, job_id )
def get_client( self, job_destination_params, job_id, env=[] ):
# Cannot use url_for outside of web thread.
# files_endpoint = url_for( controller="job_files", job_id=encoded_job_id )
encoded_job_id = self.app.security.encode_id(job_id)
job_key = self.app.security.encode_id( job_id, kind="jobs_files" )
files_endpoint = "%s/api/jobs/%s/files?job_key=%s" % (
self.galaxy_url,
encoded_job_id,
job_key
)
get_client_kwds = dict(
job_id=str( job_id ),
files_endpoint=files_endpoint,
env=env
)
return self.client_manager.get_client( job_destination_params, **get_client_kwds )
def finish_job( self, job_state ):
stderr = stdout = ''
job_wrapper = job_state.job_wrapper
try:
client = self.get_client_from_state(job_state)
run_results = client.full_status()
remote_working_directory = run_results.get("working_directory", None)
stdout = run_results.get('stdout', '')
stderr = run_results.get('stderr', '')
exit_code = run_results.get('returncode', None)
lwr_outputs = LwrOutputs.from_status_response(run_results)
# Use LWR client code to transfer/copy files back
# and cleanup job if needed.
completed_normally = \
job_wrapper.get_state() not in [ model.Job.states.ERROR, model.Job.states.DELETED ]
cleanup_job = self.app.config.cleanup_job
client_outputs = self.__client_outputs(client, job_wrapper)
finish_args = dict( client=client,
job_completed_normally=completed_normally,
cleanup_job=cleanup_job,
client_outputs=client_outputs,
lwr_outputs=lwr_outputs )
failed = lwr_finish_job( **finish_args )
if failed:
job_wrapper.fail("Failed to find or download one or more job outputs from remote server.", exception=True)
except Exception:
message = GENERIC_REMOTE_ERROR
job_wrapper.fail( message, exception=True )
log.exception("failure finishing job %d" % job_wrapper.job_id)
return
if not LwrJobRunner.__remote_metadata( client ):
self._handle_metadata_externally( job_wrapper, resolve_requirements=True )
# Finish the job
try:
job_wrapper.finish(
stdout,
stderr,
exit_code,
remote_working_directory=remote_working_directory
)
except Exception:
log.exception("Job wrapper finish method failed")
job_wrapper.fail("Unable to finish job", exception=True)
def fail_job( self, job_state ):
"""
Seperated out so we can use the worker threads for it.
"""
self.stop_job( self.sa_session.query( self.app.model.Job ).get( job_state.job_wrapper.job_id ) )
job_state.job_wrapper.fail( getattr( job_state, "fail_message", GENERIC_REMOTE_ERROR ) )
def check_pid( self, pid ):
try:
os.kill( pid, 0 )
return True
except OSError as e:
if e.errno == errno.ESRCH:
log.debug( "check_pid(): PID %d is dead" % pid )
else:
log.warning( "check_pid(): Got errno %s when attempting to check PID %d: %s" % ( errno.errorcode[e.errno], pid, e.strerror ) )
return False
def stop_job( self, job ):
# if our local job has JobExternalOutputMetadata associated, then our primary job has to have already finished
job_ext_output_metadata = job.get_external_output_metadata()
if job_ext_output_metadata:
pid = job_ext_output_metadata[0].job_runner_external_pid # every JobExternalOutputMetadata has a pid set, we just need to take from one of them
if pid in [ None, '' ]:
log.warning( "stop_job(): %s: no PID in database for job, unable to stop" % job.id )
return
pid = int( pid )
if not self.check_pid( pid ):
log.warning( "stop_job(): %s: PID %d was already dead or can't be signaled" % ( job.id, pid ) )
return
for sig in [ 15, 9 ]:
try:
os.killpg( pid, sig )
except OSError as e:
log.warning( "stop_job(): %s: Got errno %s when attempting to signal %d to PID %d: %s" % ( job.id, errno.errorcode[e.errno], sig, pid, e.strerror ) )
return # give up
sleep( 2 )
if not self.check_pid( pid ):
log.debug( "stop_job(): %s: PID %d successfully killed with signal %d" % ( job.id, pid, sig ) )
return
else:
log.warning( "stop_job(): %s: PID %d refuses to die after signaling TERM/KILL" % ( job.id, pid ) )
else:
# Remote kill
lwr_url = job.job_runner_name
job_id = job.job_runner_external_id
log.debug("Attempt remote lwr kill of job with url %s and id %s" % (lwr_url, job_id))
client = self.get_client(job.destination_params, job_id)
client.kill()
def recover( self, job, job_wrapper ):
"""Recovers jobs stuck in the queued/running state when Galaxy started"""
job_state = self.__job_state( job, job_wrapper )
job_wrapper.command_line = job.get_command_line()
state = job.get_state()
if state in [model.Job.states.RUNNING, model.Job.states.QUEUED]:
log.debug( "(LWR/%s) is still in running state, adding to the LWR queue" % ( job.get_id()) )
job_state.old_state = True
job_state.running = state == model.Job.states.RUNNING
self.monitor_queue.put( job_state )
def shutdown( self ):
super( LwrJobRunner, self ).shutdown()
self.client_manager.shutdown()
def __job_state( self, job, job_wrapper ):
job_state = AsynchronousJobState()
# TODO: Determine why this is set when using normal message queue updates
# but not CLI submitted MQ updates...
raw_job_id = job.get_job_runner_external_id() or job_wrapper.job_id
job_state.job_id = str( raw_job_id )
job_state.runner_url = job_wrapper.get_job_runner_url()
job_state.job_destination = job_wrapper.job_destination
job_state.job_wrapper = job_wrapper
return job_state
def __client_outputs( self, client, job_wrapper ):
work_dir_outputs = self.get_work_dir_outputs( job_wrapper )
output_files = self.get_output_files( job_wrapper )
client_outputs = ClientOutputs(
working_directory=job_wrapper.working_directory,
work_dir_outputs=work_dir_outputs,
output_files=output_files,
version_file=job_wrapper.get_version_string_path(),
)
return client_outputs
@staticmethod
def __dependencies_description( lwr_client, job_wrapper ):
dependency_resolution = LwrJobRunner.__dependency_resolution( lwr_client )
remote_dependency_resolution = dependency_resolution == "remote"
if not remote_dependency_resolution:
return None
requirements = job_wrapper.tool.requirements or []
installed_tool_dependencies = job_wrapper.tool.installed_tool_dependencies or []
return dependencies.DependenciesDescription(
requirements=requirements,
installed_tool_dependencies=installed_tool_dependencies,
)
@staticmethod
def __dependency_resolution( lwr_client ):
dependency_resolution = lwr_client.destination_params.get( "dependency_resolution", "local" )
if dependency_resolution not in ["none", "local", "remote"]:
raise Exception("Unknown dependency_resolution value encountered %s" % dependency_resolution)
return dependency_resolution
@staticmethod
def __remote_metadata( lwr_client ):
remote_metadata = string_as_bool_or_none( lwr_client.destination_params.get( "remote_metadata", False ) )
return remote_metadata
@staticmethod
def __use_remote_datatypes_conf( lwr_client ):
""" When setting remote metadata, use integrated datatypes from this
Galaxy instance or use the datatypes config configured via the remote
LWR.
Both options are broken in different ways for same reason - datatypes
may not match. One can push the local datatypes config to the remote
server - but there is no guarentee these datatypes will be defined
there. Alternatively, one can use the remote datatype config - but
there is no guarentee that it will contain all the datatypes available
to this Galaxy.
"""
use_remote_datatypes = string_as_bool_or_none( lwr_client.destination_params.get( "use_remote_datatypes", False ) )
return use_remote_datatypes
@staticmethod
def __rewrite_parameters( lwr_client ):
return string_as_bool_or_none( lwr_client.destination_params.get( "rewrite_parameters", False ) ) or False
def __build_metadata_configuration(self, client, job_wrapper, remote_metadata, remote_job_config):
metadata_kwds = {}
if remote_metadata:
remote_system_properties = remote_job_config.get("system_properties", {})
remote_galaxy_home = remote_system_properties.get("galaxy_home", None)
if not remote_galaxy_home:
raise Exception(NO_REMOTE_GALAXY_FOR_METADATA_MESSAGE)
metadata_kwds['exec_dir'] = remote_galaxy_home
outputs_directory = remote_job_config['outputs_directory']
configs_directory = remote_job_config['configs_directory']
working_directory = remote_job_config['working_directory']
# For metadata calculation, we need to build a list of of output
# file objects with real path indicating location on Galaxy server
# and false path indicating location on compute server. Since the
# LWR disables from_work_dir copying as part of the job command
# line we need to take the list of output locations on the LWR
# server (produced by self.get_output_files(job_wrapper)) and for
# each work_dir output substitute the effective path on the LWR
# server relative to the remote working directory as the
# false_path to send the metadata command generation module.
work_dir_outputs = self.get_work_dir_outputs(job_wrapper, job_working_directory=working_directory)
outputs = [Bunch(false_path=os.path.join(outputs_directory, os.path.basename(path)), real_path=path) for path in self.get_output_files(job_wrapper)]
for output in outputs:
for lwr_workdir_path, real_path in work_dir_outputs:
if real_path == output.real_path:
output.false_path = lwr_workdir_path
metadata_kwds['output_fnames'] = outputs
metadata_kwds['compute_tmp_dir'] = working_directory
metadata_kwds['config_root'] = remote_galaxy_home
default_config_file = os.path.join(remote_galaxy_home, 'galaxy.ini')
metadata_kwds['config_file'] = remote_system_properties.get('galaxy_config_file', default_config_file)
metadata_kwds['dataset_files_path'] = remote_system_properties.get('galaxy_dataset_files_path', None)
if LwrJobRunner.__use_remote_datatypes_conf( client ):
remote_datatypes_config = remote_system_properties.get('galaxy_datatypes_config_file', None)
if not remote_datatypes_config:
log.warn(NO_REMOTE_DATATYPES_CONFIG)
remote_datatypes_config = os.path.join(remote_galaxy_home, 'datatypes_conf.xml')
metadata_kwds['datatypes_config'] = remote_datatypes_config
else:
integrates_datatypes_config = self.app.datatypes_registry.integrated_datatypes_configs
# Ensure this file gets pushed out to the remote config dir.
job_wrapper.extra_filenames.append(integrates_datatypes_config)
metadata_kwds['datatypes_config'] = os.path.join(configs_directory, os.path.basename(integrates_datatypes_config))
return metadata_kwds
class LwrComputeEnvironment( ComputeEnvironment ):
def __init__( self, lwr_client, job_wrapper, remote_job_config ):
self.lwr_client = lwr_client
self.job_wrapper = job_wrapper
self.local_path_config = job_wrapper.default_compute_environment()
self.unstructured_path_rewrites = {}
# job_wrapper.prepare is going to expunge the job backing the following
# computations, so precalculate these paths.
self._wrapper_input_paths = self.local_path_config.input_paths()
self._wrapper_output_paths = self.local_path_config.output_paths()
self.path_mapper = PathMapper(lwr_client, remote_job_config, self.local_path_config.working_directory())
self._config_directory = remote_job_config[ "configs_directory" ]
self._working_directory = remote_job_config[ "working_directory" ]
self._sep = remote_job_config[ "system_properties" ][ "separator" ]
self._tool_dir = remote_job_config[ "tools_directory" ]
version_path = self.local_path_config.version_path()
new_version_path = self.path_mapper.remote_version_path_rewrite(version_path)
if new_version_path:
version_path = new_version_path
self._version_path = version_path
def output_paths( self ):
local_output_paths = self._wrapper_output_paths
results = []
for local_output_path in local_output_paths:
wrapper_path = str( local_output_path )
remote_path = self.path_mapper.remote_output_path_rewrite( wrapper_path )
results.append( self._dataset_path( local_output_path, remote_path ) )
return results
def input_paths( self ):
local_input_paths = self._wrapper_input_paths
results = []
for local_input_path in local_input_paths:
wrapper_path = str( local_input_path )
# This will over-copy in some cases. For instance in the case of task
# splitting, this input will be copied even though only the work dir
# input will actually be used.
remote_path = self.path_mapper.remote_input_path_rewrite( wrapper_path )
results.append( self._dataset_path( local_input_path, remote_path ) )
return results
def _dataset_path( self, local_dataset_path, remote_path ):
remote_extra_files_path = None
if remote_path:
remote_extra_files_path = "%s_files" % remote_path[ 0:-len( ".dat" ) ]
return local_dataset_path.with_path_for_job( remote_path, remote_extra_files_path )
def working_directory( self ):
return self._working_directory
def config_directory( self ):
return self._config_directory
def new_file_path( self ):
return self.working_directory() # Problems with doing this?
def sep( self ):
return self._sep
def version_path( self ):
return self._version_path
def rewriter( self, parameter_value ):
unstructured_path_rewrites = self.unstructured_path_rewrites
if parameter_value in unstructured_path_rewrites:
# Path previously mapped, use previous mapping.
return unstructured_path_rewrites[ parameter_value ]
if parameter_value in unstructured_path_rewrites.itervalues():
# Path is a rewritten remote path (this might never occur,
# consider dropping check...)
return parameter_value
rewrite, new_unstructured_path_rewrites = self.path_mapper.check_for_arbitrary_rewrite( parameter_value )
if rewrite:
unstructured_path_rewrites.update(new_unstructured_path_rewrites)
return rewrite
else:
# Did need to rewrite, use original path or value.
return parameter_value
def unstructured_path_rewriter( self ):
return self.rewriter
@@ -1,62 +0,0 @@
"""
lwr_client
==========
This module contains logic for interfacing with an external LWR server.
------------------
Configuring Galaxy
------------------
Galaxy job runners are configured in Galaxy's ``job_conf.xml`` file. See ``job_conf.xml.sample_advanced``
in your Galaxy code base or on
`Bitbucket <https://bitbucket.org/galaxy/galaxy-dist/src/tip/job_conf.xml.sample_advanced?at=default>`_
for information on how to configure Galaxy to interact with the LWR.
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 LWR in this legacy mode.
Legacy
------
A Galaxy tool can be configured to be executed remotely via LWR by
adding a line to the ``galaxy.ini`` file under the ``galaxy:tool_runners``
section with the format::
<tool_id> = lwr://http://<lwr_host>:<lwr_port>
As an example, if a host named remotehost is running the LWR 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 = lwr://http://remotehost:8913
Remember this must be added after the ``[galaxy:tool_runners]`` header
in the ``universe.ini`` file.
"""
from .staging.down import finish_job
from .staging.up import submit_job
from .staging import ClientJobDescription
from .staging import LwrOutputs
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,
LwrOutputs,
ClientOutputs,
PathMapper,
]
@@ -1,567 +0,0 @@
from json import load
from os import makedirs
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
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]
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)
... mapper = FileActionMapper(config=mapper.to_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.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,
paths=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
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 action.action_type == "remote_transfer":
url_base = self.files_endpoint
if not url_base:
raise Exception("Attempted to use remote_transfer action with defining a files_endpoint")
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
REQUIRED_ACTION_KWD = object()
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.iteritems():
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
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 LWR client (i.e. Galaxy server) and remote LWR 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 LWR 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 LWR client should initiate an HTTP
transfer of the corresponding path to the remote LWR server before
launching the job. """
action_type = "transfer"
staging = STAGING_ACTION_LOCAL
class CopyAction(BaseAction):
""" This action indicates that the LWR client should execute a file system
copy of the corresponding path to the LWR staging directory prior to
launching the corresponding job. """
action_type = "copy"
staging = STAGING_ACTION_LOCAL
class RemoteCopyAction(BaseAction):
""" This action indicates the LWR 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 LWR 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, lwr_path):
destination = self.path
parent_directory = dirname(destination)
if not exists(parent_directory):
makedirs(parent_directory)
with open(lwr_path, "rb") as f:
galaxy.util.copy_to_path(f, destination)
class RemoteTransferAction(BaseAction):
""" This action indicates the LWR 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 LWR server instead of on
the client.
"""
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, lwr_path):
post_file(self.url, lwr_path)
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]
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)
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 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,
]
actions = dict([(clazz.action_type, clazz) for clazz in ACTION_CLASSES])
__all__ = [
FileActionMapper,
path_type,
from_dict,
MessageAction,
RemoteTransferAction, # For testing
]
@@ -1,139 +0,0 @@
try:
import kombu
from kombu import pools
except ImportError:
kombu = None
import socket
import logging
import threading
from time import sleep
log = logging.getLogger(__name__)
KOMBU_UNAVAILABLE = "Attempting to bind to AMQP message queue, but kombu dependency unavailable"
DEFAULT_EXCHANGE_NAME = "lwr"
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
class LwrExchange(object):
""" Utility for publishing and consuming structured LWR 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 LWR manager is defined solely by name in the scheme, so only one LWR
should target each AMQP endpoint or care should be taken that unique
manager names are used across LWR servers targetting same AMQP endpoint -
and in particular only one such LWR should define an default manager with
name _default_.
"""
def __init__(
self,
url,
manager_name,
connect_ssl=None,
timeout=DEFAULT_TIMEOUT,
publish_kwds={},
):
"""
"""
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
# Be sure to log message publishing failures.
if publish_kwds.get("retry", False):
if "retry_policy" not in publish_kwds:
publish_kwds["retry_policy"] = {}
if "errback" not in publish_kwds["retry_policy"]:
publish_kwds["retry_policy"]["errback"] = self.__publish_errback
self.__publish_kwds = publish_kwds
@property
def url(self):
return self.__url
def consume(self, queue_name, callback, check=True, connection_kwargs={}):
queue = self.__queue(queue_name)
log.debug("Consuming queue '%s'", queue)
while check:
heartbeat_thread = None
try:
with self.connection(self.__url, heartbeat=DEFAULT_HEARTBEAT, **connection_kwargs) as connection:
with kombu.Consumer(connection, queues=[queue], callbacks=[callback], accept=['json']):
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:
# In testing, errno is None
log.warning('Got %s, will retry: %s', exc.__class__.__name__, exc)
if heartbeat_thread:
heartbeat_thread.join()
sleep(DEFAULT_RECONNECT_CONSUMER_WAIT)
def heartbeat(self, connection):
log.debug('AMQP heartbeat thread alive')
while connection.connected:
connection.heartbeat_check()
sleep(DEFAULT_HEARTBEAT_WAIT)
log.debug('AMQP heartbeat thread exiting')
def publish(self, name, payload):
with self.connection(self.__url) as connection:
with pools.producers[connection].acquire() as producer:
key = self.__queue_name(name)
producer.publish(
payload,
serializer='json',
exchange=self.__exchange,
declare=[self.__exchange],
routing_key=key,
**self.__publish_kwds
)
def __publish_errback(self, exc, interval):
log.error("Connection error while publishing: %r", exc, exc_info=1)
log.info("Retrying in %s seconds", interval)
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 = "lwr_"
else:
key_prefix = "lwr_%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
@@ -1,41 +0,0 @@
from .amqp_exchange import LwrExchange
from .util import filter_destination_params
def get_exchange(url, manager_name, params):
connect_ssl = parse_amqp_connect_ssl_params(params)
exchange_kwds = dict(
manager_name=manager_name,
connect_ssl=connect_ssl,
publish_kwds=parse_amqp_publish_kwds(params)
)
timeout = params.get('amqp_consumer_timeout', False)
if timeout is not False:
exchange_kwds['timeout'] = timeout
exchange = LwrExchange(url, **exchange_kwds)
return exchange
def parse_amqp_connect_ssl_params(params):
ssl_params = filter_destination_params(params, "amqp_connect_ssl_")
if not ssl_params:
return
ssl = __import__('ssl')
if 'cert_reqs' in ssl_params:
value = ssl_params['cert_reqs']
ssl_params['cert_reqs'] = getattr(ssl, value.upper())
return ssl_params
def parse_amqp_publish_kwds(params):
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
@@ -1,428 +0,0 @@
import os
from json import dumps
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 copy
from .util import ensure_directory
from .util import to_base64_json
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):
self.destination_params = destination_params
self.job_id = job_id
if "jobs_directory" in (destination_params or {}):
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
self.env = destination_params.get("env", [])
self.files_endpoint = destination_params.get("files_endpoint", None)
self.job_directory = job_directory
self.default_file_action = self.destination_params.get("default_file_action", "transfer")
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 LWR 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 LWR 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'] = dumps(submit_params_dict)
if dependencies_description:
launch_params['dependencies_description'] = dumps(dependencies_description.to_dict())
if env:
launch_params['env'] = dumps(env)
if remote_staging:
launch_params['remote_staging'] = 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 LWR setup job
# before queueing.
setup_params = _setup_params_from_job_config(job_config)
launch_params["setup_params"] = dumps(setup_params)
return self._raw_execute("launch", 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("kill", {"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("check_complete", {"job_id": self.job_id})
return check_complete_response
def get_status(self):
check_complete_response = self.raw_check_complete()
# Older LWR instances won't set status so use 'complete', at some
# point drop backward compatibility.
status = check_complete_response.get("status", None)
if status in ["status", None]:
# LEGACY: Bug in certains older LWR instances returned literal
# "status".
complete = check_complete_response["complete"] == "true"
old_status = "complete" if complete else "running"
status = old_status
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 LWR 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, "input_type": input_type}
input_path = path
if contents:
input_path = None
if action_type == 'transfer':
return self._upload_file(args, contents, input_path)
elif action_type == 'copy':
lwr_path = self._raw_execute('input_path', args)
copy(path, lwr_path)
return {'path': lwr_path}
def fetch_output(self, path, name, working_directory, action_type, output_type):
"""
Fetch (transfer, copy, etc...) an output from the remote LWR 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 LWR (output_workdir or output). legacy is also
an option in this case LWR is asked for location - this will only be
used if targetting an older LWR server that didn't return statuses
allowing this to be inferred.
"""
if output_type == 'legacy':
self._fetch_output_legacy(path, working_directory, action_type=action_type)
elif 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)
# Deprecated
def _fetch_output_legacy(self, path, working_directory, action_type='transfer'):
# Needs to determine if output is task/working directory or standard.
name = os.path.basename(path)
output_type = self._get_output_type(name)
if output_type == "none":
# Just make sure the file was created.
if not os.path.exists(path):
raise OutputNotFoundException(path)
return
elif output_type in ["task"]:
path = os.path.join(working_directory, name)
self.__populate_output_path(name, path, output_type, action_type)
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)
output_type = "direct" # Task/from_work_dir outputs now handled with fetch_work_dir_output
self.__populate_output_path(name, path, output_type, 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, "work_dir", output_path)
else: # Even if action is none - LWR has a different work_dir so this needs to be copied.
lwr_path = self._output_path(name, self.job_id, 'work_dir')['path']
copy(lwr_path, output_path)
def __populate_output_path(self, name, output_path, output_type, action_type):
ensure_directory(output_path)
if action_type == 'transfer':
self.__raw_download_output(name, self.job_id, output_type, output_path)
elif action_type == 'copy':
lwr_path = self._output_path(name, self.job_id, output_type)['path']
copy(lwr_path, output_path)
@parseJson()
def _upload_file(self, args, contents, input_path):
return self._raw_execute(self._upload_file_action(args), args, contents, input_path)
def _upload_file_action(self, args):
# Hack for backward compatibility, instead of using new upload_file
# path. Use old paths.
input_type = args['input_type']
action = {
# For backward compatibility just target upload_input_extra for all
# inputs, it allows nested inputs. Want to do away with distinction
# inputs and extra inputs.
'input': 'upload_extra_input',
'config': 'upload_config_file',
'workdir': 'upload_working_directory_file',
'tool': 'upload_tool_file',
'unstructured': 'upload_unstructured_file',
}[input_type]
del args['input_type']
return action
@parseJson()
def _get_output_type(self, name):
return self._raw_execute("get_output_type", {"name": name,
"job_id": self.job_id})
@parseJson()
def _output_path(self, name, job_id, output_type):
return self._raw_execute("output_path",
{"name": name,
"job_id": self.job_id,
"output_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,
"output_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 LWR 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
if job_config and self.setup_handler.local:
# Setup not yet called, job properties were inferred from
# destination arguments. Hence, must have LWR 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_lwr_path = destination_params["remote_lwr_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_lwr_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 = self._upload_file_action(args)
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
)
@@ -1,80 +0,0 @@
""" Generic interface for reading YAML/INI/JSON config files into nested dictionaries.
"""
try:
from galaxy import eggs
eggs.require('PyYAML')
except Exception:
# If not in Galaxy, ignore this.
pass
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):
with open(path, "rb") as f:
return json.load(f)
EXT_READERS = {
CONFIG_TYPE_JSON: __read_json,
CONFIG_TYPE_YAML: __read_yaml,
CONFIG_TYPE_INI: __read_ini,
}
@@ -1,35 +0,0 @@
import time
import json
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
@@ -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 = "lwr://http://localhost:8913/"
>>> runner_params = url_to_destination_params(runner_url)
>>> runner_params['url']
'http://localhost:8913/'
"""
if url.startswith("lwr://"):
url = url[len("lwr://"):]
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)
@@ -1,102 +0,0 @@
from abc import ABCMeta
from abc import abstractmethod
try:
from StringIO import StringIO as BytesIO
except ImportError:
from io import BytesIO
try:
from six import text_type
except ImportError:
from galaxy.util import unicodify as text_type
try:
from urllib import urlencode
except ImportError:
from urllib.parse import urlencode
class LwrInteface(object):
"""
Abstract base class describes how synchronous client communicates with
(potentially remote) LWR procedures. Obvious implementation is HTTP based
but LWR 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 LWR 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.
"""
class HttpLwrInterface(LwrInteface):
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 LWR 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_key = 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)
response = self.transport.execute(url, data=data, input_path=input_path, output_path=output_path)
return response
def __build_url(self, command, args):
if self.private_key:
args["private_key"] = self.private_key
arg_bytes = dict([(k, text_type(args[k]).encode('utf-8')) for k in args])
data = urlencode(arg_bytes)
url = self.remote_host + command + "?" + data
return url
class LocalLwrInterface(LwrInteface):
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 LwrApp 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 lwr.web import routes
from lwr.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.encode('utf-8'))
elif input_path is not None:
return open(input_path, 'rb')
else:
return None
@@ -1,137 +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",
input_extra="inputs_directory",
unstructured="unstructured_files_directory",
config="configs_directory",
tool="tool_files_directory",
work_dir="working_directory",
workdir="working_directory",
output="outputs_directory",
output_workdir="working_directory",
)
class RemoteJobDirectory(object):
""" Representation of a (potentially) remote LWR-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 LWR 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', 'input_extra', '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:\\lwr\\staging\\101', 'dataset_1_files/moo/cow', allow_nested_files=True, local_path_module=ntpath, mkdir=False)
'C:\\\\lwr\\\\staging\\\\101\\\\dataset_1_files\\\\moo\\\\cow'
>>> get_mapped_file(r'C:\\lwr\\staging\\101', 'dataset_1_files/moo/cow', allow_nested_files=False, local_path_module=ntpath)
'C:\\\\lwr\\\\staging\\\\101\\\\cow'
>>> get_mapped_file(r'C:\\lwr\\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)
@@ -1,223 +0,0 @@
import threading
try:
from Queue import Queue
except ImportError:
from queue import Queue
from os import getenv
from .client import JobClient
from .client import InputCachingJobClient
from .client import MessageJobClient
from .client import MessageCLIJobClient
from .interface import HttpLwrInterface
from .interface import LocalLwrInterface
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('url', None):
return MessageQueueClientManager(**kwargs)
else:
return ClientManager(**kwargs)
class ClientManager(object):
"""
Factory to create LWR 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 = LocalLwrInterface
self.job_manager_interface_args = dict(job_manager=kwds['job_manager'], file_cache=kwds['file_cache'])
else:
self.job_manager_interface_class = HttpLwrInterface
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('LWR_CACHE_TRANSFERS')
if cache:
log.info("Setting LWR 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 LWR 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):
pass
try:
from galaxy.jobs.runners.util.cli import factory as cli_factory
except ImportError:
from lwr.managers.util.cli import factory as cli_factory
class MessageQueueClientManager(object):
def __init__(self, **kwds):
self.url = kwds.get('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.active = True
def ensure_has_status_update_callback(self, callback):
with self.callback_lock:
if self.callback_thread is not None:
return
def callback_wrapper(body, message):
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 LWR.")
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 run():
self.exchange.consume("status_update", callback_wrapper, check=self)
log.debug("Leaving LWR client status update thread, no additional LWR updates will be processed.")
thread = threading.Thread(
name="lwr_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 shutdown(self):
self.active = False
def __nonzero__(self):
return self.active
def get_client(self, destination_params, job_id, **kwargs):
if job_id is None:
raise Exception("Cannot generate LWR 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 = LocalLwrInterface
self.interface_args = dict(object_store=kwds['object_store'])
else:
self.interface_class = HttpLwrInterface
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('LWR_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):
try:
unicode_type = unicode
except NameError:
unicode_type = str
if isinstance(destination_params, str) or isinstance(destination_params, unicode_type):
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, HttpLwrInterface]
@@ -1,53 +0,0 @@
from .decorators import parseJson
class ObjectStoreClient(object):
def __init__(self, lwr_interface):
self.lwr_interface = lwr_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.lwr_interface.execute(command, args, data=None, input_path=None, output_path=None)
@@ -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 LWR setup method to pre-determine the location of files for staging
on the remote LWR server.
This is not useful when rewrite_paths (as has traditionally been done with
the LWR) because when doing that the LWR 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]
@@ -1,103 +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 LWR 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 LWR 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 = {
"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]
@@ -1,174 +0,0 @@
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"
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,
tool,
command_line,
config_files,
input_files,
client_outputs,
working_directory,
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
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, output_files, work_dir_outputs=None, version_file=None):
self.working_directory = working_directory
self.work_dir_outputs = work_dir_outputs
self.output_files = output_files
self.version_file = version_file
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
)
@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'),
)
class LwrOutputs(object):
""" Abstraction describing the output files PRODUCED by the remote LWR
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 LWR - older LWR instances will not set these in complete response.
working_directory_contents = complete_response.get("working_directory_contents", None)
output_directory_contents = complete_response.get("outputs_directory_contents", None)
# Older (pre-2014) LWR 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 LwrOutputs(
working_directory_contents,
output_directory_contents,
remote_separator
)
def has_output_file(self, output_file):
if self.output_directory_contents is None:
# Legacy LWR doesn't report this, return None indicating unsure if
# output was generated.
return None
else:
return basename(output_file) in self.output_directory_contents
def has_output_directory_listing(self):
return self.output_directory_contents is not None
def output_extras(self, output_file):
"""
Returns dict mapping local path to remote name.
"""
if not self.has_output_directory_listing():
# Fetching $output.extra_files_path is not supported with legacy
# LWR (pre-2014) severs.
return {}
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))
@@ -1,162 +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, lwr_outputs):
""" Responsible for downloading results from remote server and cleaning up
LWR 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, lwr_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 LWR.
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, lwr_outputs):
self.output_collector = output_collector
self.action_mapper = action_mapper
self.client_outputs = client_outputs
self.lwr_outputs = lwr_outputs
self.downloaded_working_directory_files = []
self.exception_tracker = DownloadExceptionTracker()
self.output_files = client_outputs.output_files
self.working_directory_contents = lwr_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)
lwr_name = self.lwr_outputs.path_helper.remote_name(name)
if self._attempt_collect_output('output_workdir', path=output_file, name=lwr_name):
self.downloaded_working_directory_files.append(lwr_name)
# 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 LWR 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.lwr_outputs.has_output_file(output_file)
if output_generated is None:
self._attempt_collect_output('legacy', output_file)
elif output_generated:
self._attempt_collect_output('output', output_file)
for galaxy_path, lwr_name in self.lwr_outputs.output_extras(output_file).iteritems():
self._attempt_collect_output('output', path=galaxy_path, name=lwr_name)
# else not output generated, do not attempt download.
def __collect_version_file(self):
version_file = self.client_outputs.version_file
# output_directory_contents may be none for legacy LWR servers.
lwr_output_directory_contents = (self.lwr_outputs.output_directory_contents or [])
if version_file and COMMAND_VERSION_FILENAME in lwr_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 COPY_FROM_WORKING_DIRECTORY_PATTERN.match(name):
output_file = join(working_directory, self.lwr_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 LWR server (possible a relative)
# path.
collected = False
with self.exception_tracker():
# output_action_type cannot be 'legacy' but output_type may be
# eventually drop support for legacy mode (where type wasn't known)
# ahead of time.
output_action_type = 'output_workdir' if output_type == 'output_workdir' else 'output'
action = self.action_mapper.action(path, output_action_type)
if self._collect_output(output_type, action, name):
collected = True
return collected
def _collect_output(self, output_type, action, name):
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 LWR job")
__all__ = [finish_job]
@@ -1,420 +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 LWR client object to
stage the files required to run jobs on a remote LWR server.
**Parameters**
client : JobClient
LWR 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
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)
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 LWR 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 LWR 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 LWR 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 LWR 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 LWR
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.iteritems():
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.iteritems():
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 = "LWR: __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 = listdir(self.working_directory) if exists(self.working_directory) else []
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 __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.
"""
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)
get_path = lambda: 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)
get_path = lambda: 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]
@@ -1,31 +0,0 @@
from .standard import Urllib2Transport
from .curl import PycurlTransport
import os
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('LWR_CURL_TRANSPORT', "0")
# If LWR_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
# TODO: Provide urllib implementation if these unavailable,
# also explore a requests+poster option.
from .curl import get_file
from .curl import post_file
__all__ = [get_transport, get_file, post_file]
@@ -1,69 +0,0 @@
try:
from cStringIO import StringIO
except ImportError:
from io import StringIO
try:
from pycurl import Curl
except ImportError:
pass
from os.path import getsize
PYCURL_UNAVAILABLE_MESSAGE = \
"You are attempting to use the Pycurl version of the LWR client but pycurl is unavailable."
class PycurlTransport(object):
def execute(self, url, data=None, input_path=None, output_path=None):
buf = _open_output(output_path)
try:
c = _new_curl_object()
c.setopt(c.URL, url.encode('ascii'))
c.setopt(c.WRITEFUNCTION, buf.write)
if input_path:
c.setopt(c.UPLOAD, 1)
c.setopt(c.READFUNCTION, open(input_path, 'rb').read)
filesize = getsize(input_path)
c.setopt(c.INFILESIZE, filesize)
if data:
c.setopt(c.POST, 1)
if type(data).__name__ == 'unicode':
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):
c = _new_curl_object()
c.setopt(c.URL, url.encode('ascii'))
c.setopt(c.HTTPPOST, [("file", (c.FORM_FILE, path.encode('ascii')))])
c.perform()
def get_file(url, path):
buf = _open_output(path)
try:
c = _new_curl_object()
c.setopt(c.URL, url.encode('ascii'))
c.setopt(c.WRITEFUNCTION, buf.write)
c.perform()
finally:
buf.close()
def _open_output(output_path):
return open(output_path, 'wb') if output_path else StringIO()
def _new_curl_object():
try:
return Curl()
except NameError:
raise ImportError(PYCURL_UNAVAILABLE_MESSAGE)
___all__ = [PycurlTransport, post_file, get_file]
@@ -1,45 +0,0 @@
"""
LWR 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, data=None, input_path=None, output_path=None):
request = Request(url=url, data=data)
input = None
try:
if input_path:
if getsize(input_path):
input = open(input_path, 'rb')
data = mmap.mmap(input.fileno(), 0, access=mmap.ACCESS_READ)
else:
data = b""
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()
-177
View File
@@ -1,177 +0,0 @@
from threading import Lock, Event
from weakref import WeakValueDictionary
from os import walk
from os import curdir
from os.path import relpath
from os.path import join
import os.path
import hashlib
import shutil
import json
import base64
def unique_path_prefix(path):
m = hashlib.md5()
m.update(path)
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):
"""
>>> x = from_base64_json(to_base64_json(dict(a=5)))
>>> x["a"]
5
"""
return base64.b64encode(json.dumps(data))
def from_base64_json(data):
return json.loads(base64.b64decode(data))
class PathHelper(object):
'''
>>> import posixpath
>>> # Forcing local path to posixpath because LWR 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
+1 -1
View File
@@ -316,7 +316,7 @@ class ToolEvaluator( object ):
try:
open( dataset_path.false_path, 'w' ).close()
except EnvironmentError:
pass # May well not exist - e.g. LWR.
pass # May well not exist - e.g. Pulsar.
else:
param_dict[name] = DatasetFilenameWrapper( hda )
# Provide access to a path to store additional files
+1 -1
View File
@@ -9,7 +9,7 @@ log = getLogger( __name__ )
# Fields in .log files corresponding to paths, must have one of the following
# field names and all such fields are assumed to be paths. This is to allow
# remote ComputeEnvironments (such as one used by LWR) determine what values to
# remote ComputeEnvironments (such as one used by Pulsar) determine what values to
# rewrite or transfer...
PATH_ATTRIBUTES = [ "path" ]
# ... by default though - don't rewrite anything (if no ComputeEnviornment