diff --git a/doc/source/lib/galaxy.jobs.runners.lwr_client.rst b/doc/source/lib/galaxy.jobs.runners.lwr_client.rst deleted file mode 100644 index 1c0be4e8fac..00000000000 --- a/doc/source/lib/galaxy.jobs.runners.lwr_client.rst +++ /dev/null @@ -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: - - diff --git a/doc/source/lib/galaxy.jobs.runners.lwr_client.staging.rst b/doc/source/lib/galaxy.jobs.runners.lwr_client.staging.rst deleted file mode 100644 index 7ae6a269d3a..00000000000 --- a/doc/source/lib/galaxy.jobs.runners.lwr_client.staging.rst +++ /dev/null @@ -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: - - diff --git a/doc/source/lib/galaxy.jobs.runners.lwr_client.transport.rst b/doc/source/lib/galaxy.jobs.runners.lwr_client.transport.rst deleted file mode 100644 index f884599ba2d..00000000000 --- a/doc/source/lib/galaxy.jobs.runners.lwr_client.transport.rst +++ /dev/null @@ -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: - - diff --git a/doc/source/lib/galaxy.jobs.runners.rst b/doc/source/lib/galaxy.jobs.runners.rst index 30c2f598717..cdb20be1e44 100644 --- a/doc/source/lib/galaxy.jobs.runners.rst +++ b/doc/source/lib/galaxy.jobs.runners.rst @@ -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 ------------------------------ diff --git a/lib/galaxy/jobs/__init__.py b/lib/galaxy/jobs/__init__.py index 2ed8ee820e3..9429dec9ac3 100644 --- a/lib/galaxy/jobs/__init__.py +++ b/lib/galaxy/jobs/__init__.py @@ -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)) diff --git a/lib/galaxy/jobs/runners/__init__.py b/lib/galaxy/jobs/runners/__init__.py index dcc2a31d345..2836247bd3a 100644 --- a/lib/galaxy/jobs/runners/__init__.py +++ b/lib/galaxy/jobs/runners/__init__.py @@ -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 diff --git a/lib/galaxy/jobs/runners/lwr.py b/lib/galaxy/jobs/runners/lwr.py deleted file mode 100644 index d7df067e08f..00000000000 --- a/lib/galaxy/jobs/runners/lwr.py +++ /dev/null @@ -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 diff --git a/lib/galaxy/jobs/runners/lwr_client/__init__.py b/lib/galaxy/jobs/runners/lwr_client/__init__.py deleted file mode 100644 index 32cbdc082d3..00000000000 --- a/lib/galaxy/jobs/runners/lwr_client/__init__.py +++ /dev/null @@ -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 `_ -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:: - - = lwr://http://: - -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, -] diff --git a/lib/galaxy/jobs/runners/lwr_client/action_mapper.py b/lib/galaxy/jobs/runners/lwr_client/action_mapper.py deleted file mode 100644 index 13de35c865b..00000000000 --- a/lib/galaxy/jobs/runners/lwr_client/action_mapper.py +++ /dev/null @@ -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 -] diff --git a/lib/galaxy/jobs/runners/lwr_client/amqp_exchange.py b/lib/galaxy/jobs/runners/lwr_client/amqp_exchange.py deleted file mode 100644 index 96ba0db559d..00000000000 --- a/lib/galaxy/jobs/runners/lwr_client/amqp_exchange.py +++ /dev/null @@ -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 diff --git a/lib/galaxy/jobs/runners/lwr_client/amqp_exchange_factory.py b/lib/galaxy/jobs/runners/lwr_client/amqp_exchange_factory.py deleted file mode 100644 index c35a9648800..00000000000 --- a/lib/galaxy/jobs/runners/lwr_client/amqp_exchange_factory.py +++ /dev/null @@ -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 diff --git a/lib/galaxy/jobs/runners/lwr_client/client.py b/lib/galaxy/jobs/runners/lwr_client/client.py deleted file mode 100644 index 7656b05202c..00000000000 --- a/lib/galaxy/jobs/runners/lwr_client/client.py +++ /dev/null @@ -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 - ) diff --git a/lib/galaxy/jobs/runners/lwr_client/config_util.py b/lib/galaxy/jobs/runners/lwr_client/config_util.py deleted file mode 100644 index ab5453d0c8f..00000000000 --- a/lib/galaxy/jobs/runners/lwr_client/config_util.py +++ /dev/null @@ -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, -} diff --git a/lib/galaxy/jobs/runners/lwr_client/decorators.py b/lib/galaxy/jobs/runners/lwr_client/decorators.py deleted file mode 100644 index 94035f5ba1b..00000000000 --- a/lib/galaxy/jobs/runners/lwr_client/decorators.py +++ /dev/null @@ -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 diff --git a/lib/galaxy/jobs/runners/lwr_client/destination.py b/lib/galaxy/jobs/runners/lwr_client/destination.py deleted file mode 100644 index 53569f6f260..00000000000 --- a/lib/galaxy/jobs/runners/lwr_client/destination.py +++ /dev/null @@ -1,58 +0,0 @@ - -from re import match -from .util import filter_destination_params - -SUBMIT_PREFIX = "submit_" - - -def url_to_destination_params(url): - """Convert a legacy runner URL to a job destination - - >>> params_simple = url_to_destination_params("http://localhost:8913/") - >>> params_simple["url"] - 'http://localhost:8913/' - >>> params_simple["private_token"] is None - True - >>> advanced_url = "https://1234x@example.com:8914/managers/longqueue" - >>> params_advanced = url_to_destination_params(advanced_url) - >>> params_advanced["url"] - 'https://example.com:8914/managers/longqueue/' - >>> params_advanced["private_token"] - '1234x' - >>> runner_url = "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) diff --git a/lib/galaxy/jobs/runners/lwr_client/interface.py b/lib/galaxy/jobs/runners/lwr_client/interface.py deleted file mode 100644 index 95538d224da..00000000000 --- a/lib/galaxy/jobs/runners/lwr_client/interface.py +++ /dev/null @@ -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 diff --git a/lib/galaxy/jobs/runners/lwr_client/job_directory.py b/lib/galaxy/jobs/runners/lwr_client/job_directory.py deleted file mode 100644 index b07e1c539ee..00000000000 --- a/lib/galaxy/jobs/runners/lwr_client/job_directory.py +++ /dev/null @@ -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) diff --git a/lib/galaxy/jobs/runners/lwr_client/manager.py b/lib/galaxy/jobs/runners/lwr_client/manager.py deleted file mode 100644 index fb6d6c63a84..00000000000 --- a/lib/galaxy/jobs/runners/lwr_client/manager.py +++ /dev/null @@ -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] diff --git a/lib/galaxy/jobs/runners/lwr_client/object_client.py b/lib/galaxy/jobs/runners/lwr_client/object_client.py deleted file mode 100644 index 5ec9acdb4f5..00000000000 --- a/lib/galaxy/jobs/runners/lwr_client/object_client.py +++ /dev/null @@ -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) diff --git a/lib/galaxy/jobs/runners/lwr_client/path_mapper.py b/lib/galaxy/jobs/runners/lwr_client/path_mapper.py deleted file mode 100644 index 758c0901638..00000000000 --- a/lib/galaxy/jobs/runners/lwr_client/path_mapper.py +++ /dev/null @@ -1,98 +0,0 @@ -import os.path -from .action_mapper import FileActionMapper -from .action_mapper import path_type -from .util import PathHelper - -from galaxy.util import in_directory - - -class PathMapper(object): - """ Ties together a FileActionMapper and remote job configuration returned - by the 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] diff --git a/lib/galaxy/jobs/runners/lwr_client/setup_handler.py b/lib/galaxy/jobs/runners/lwr_client/setup_handler.py deleted file mode 100644 index 19906edf667..00000000000 --- a/lib/galaxy/jobs/runners/lwr_client/setup_handler.py +++ /dev/null @@ -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] diff --git a/lib/galaxy/jobs/runners/lwr_client/staging/__init__.py b/lib/galaxy/jobs/runners/lwr_client/staging/__init__.py deleted file mode 100644 index c2a51da1870..00000000000 --- a/lib/galaxy/jobs/runners/lwr_client/staging/__init__.py +++ /dev/null @@ -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)) diff --git a/lib/galaxy/jobs/runners/lwr_client/staging/down.py b/lib/galaxy/jobs/runners/lwr_client/staging/down.py deleted file mode 100644 index 08e9db410a2..00000000000 --- a/lib/galaxy/jobs/runners/lwr_client/staging/down.py +++ /dev/null @@ -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] diff --git a/lib/galaxy/jobs/runners/lwr_client/staging/up.py b/lib/galaxy/jobs/runners/lwr_client/staging/up.py deleted file mode 100644 index c8996fd8f4e..00000000000 --- a/lib/galaxy/jobs/runners/lwr_client/staging/up.py +++ /dev/null @@ -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] diff --git a/lib/galaxy/jobs/runners/lwr_client/transport/__init__.py b/lib/galaxy/jobs/runners/lwr_client/transport/__init__.py deleted file mode 100644 index f7119bbdc93..00000000000 --- a/lib/galaxy/jobs/runners/lwr_client/transport/__init__.py +++ /dev/null @@ -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] diff --git a/lib/galaxy/jobs/runners/lwr_client/transport/curl.py b/lib/galaxy/jobs/runners/lwr_client/transport/curl.py deleted file mode 100644 index 307af0b1fce..00000000000 --- a/lib/galaxy/jobs/runners/lwr_client/transport/curl.py +++ /dev/null @@ -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] diff --git a/lib/galaxy/jobs/runners/lwr_client/transport/standard.py b/lib/galaxy/jobs/runners/lwr_client/transport/standard.py deleted file mode 100644 index 1bd504ace40..00000000000 --- a/lib/galaxy/jobs/runners/lwr_client/transport/standard.py +++ /dev/null @@ -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() diff --git a/lib/galaxy/jobs/runners/lwr_client/util.py b/lib/galaxy/jobs/runners/lwr_client/util.py deleted file mode 100644 index 62f5aba3ff9..00000000000 --- a/lib/galaxy/jobs/runners/lwr_client/util.py +++ /dev/null @@ -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 diff --git a/lib/galaxy/tools/evaluation.py b/lib/galaxy/tools/evaluation.py index 4b00d004363..3a9e4460d77 100644 --- a/lib/galaxy/tools/evaluation.py +++ b/lib/galaxy/tools/evaluation.py @@ -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 diff --git a/lib/galaxy/tools/wrappers.py b/lib/galaxy/tools/wrappers.py index 1fc5b75b97a..43583842008 100644 --- a/lib/galaxy/tools/wrappers.py +++ b/lib/galaxy/tools/wrappers.py @@ -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