diff --git a/job_conf.xml.sample_advanced b/job_conf.xml.sample_advanced
index 0672a52c719..ad14177b546 100644
--- a/job_conf.xml.sample_advanced
+++ b/job_conf.xml.sample_advanced
@@ -19,27 +19,30 @@
/sge/lib/libdrmaa.so
-
-
-
+
+
+
+
+
+
-
-
-
- amqp://guest:guest@localhost:5672//
-
+
+
+ amqp://guest:guest@localhost:5672//
+
http://localhost:8080
-
+
-
-
-
+
+
+
@@ -213,87 +224,71 @@
foo
-
- https://windowshost.examle.com:8913/
-
+
+ https://examle.com:8913/
+
123456789changeme
-
-
+
-
-
-
+
-
-
-
-
-
-
+
-
+
+
+
- /path/to/remote/lwr/lwr_staging/
-
- remote_transfer
SecureShell
diff --git a/lib/galaxy/jobs/runners/pulsar.py b/lib/galaxy/jobs/runners/pulsar.py
new file mode 100644
index 00000000000..dc4b6e5f1bb
--- /dev/null
+++ b/lib/galaxy/jobs/runners/pulsar.py
@@ -0,0 +1,707 @@
+from __future__ import absolute_import # Need to import pulsar_client absolutely.
+
+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 pulsar.client import build_client_manager
+from pulsar.client import url_to_destination_params
+from pulsar.client import finish_job as pulsar_finish_job
+from pulsar.client import submit_job as pulsar_submit_job
+from pulsar.client import ClientJobDescription
+from pulsar.client import PulsarOutputs
+from pulsar.client import ClientOutputs
+from pulsar.client import PathMapper
+
+log = logging.getLogger( __name__ )
+
+__all__ = [ 'PulsarLegacyJobRunner', 'PulsarRESTJobRunner', 'PulsarMQJobRunner' ]
+
+NO_REMOTE_GALAXY_FOR_METADATA_MESSAGE = "Pulsar misconfiguration - Pulsar client configured to set metadata remotely, but remote Pulsar isn't properly configured with a galaxy_home directory."
+NO_REMOTE_DATATYPES_CONFIG = "Pulsar client is configured to use remote datatypes configuration when setting metadata externally, but Pulsar 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"
+
+PULSAR_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,
+ ),
+ amqp_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,
+ ),
+)
+
+
+PARAMETER_SPECIFICATION_REQUIRED = object()
+PARAMETER_SPECIFICATION_IGNORED = object()
+
+
+class PulsarJobRunner( AsynchronousJobRunner ):
+ """
+ Pulsar Job Runner
+ """
+ runner_name = "PulsarJobRunner"
+
+ def __init__( self, app, nworkers, **kwds ):
+ """Start the job runner """
+ super( PulsarJobRunner, self ).__init__( app, nworkers, runner_param_specs=PULSAR_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()
+ self._monitor()
+
+ def _monitor( self ):
+ # Extension point allow MQ variant to setup callback instead
+ self._init_monitor_thread()
+
+ def __init_client_manager( self ):
+ client_manager_kwargs = {}
+ for kwd in '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="pulsar", 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_status(job_state, status)
+ return job_state
+
+ def _update_job_state_for_status(self, job_state, pulsar_status):
+ if pulsar_status == "complete":
+ self.mark_as_finished(job_state)
+ return None
+ if pulsar_status == "failed":
+ self.fail_job(job_state)
+ return None
+ if pulsar_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 queue_job(self, job_wrapper):
+ job_destination = job_wrapper.job_destination
+ self._populate_parameter_defaults( 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 = PulsarJobRunner.__dependencies_description( client, job_wrapper )
+ rewrite_paths = not PulsarJobRunner.__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 = pulsar_submit_job(client, client_job_description, remote_job_config)
+ log.info("Pulsar 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
+
+ pulsar_job_state = AsynchronousJobState()
+ pulsar_job_state.job_wrapper = job_wrapper
+ pulsar_job_state.job_id = job_id
+ pulsar_job_state.old_state = True
+ pulsar_job_state.running = False
+ pulsar_job_state.job_destination = job_destination
+ self.monitor_job(pulsar_job_state)
+
+ def __prepare_job(self, job_wrapper, job_destination):
+ """ Build command-line and Pulsar 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 = PulsarJobRunner.__rewrite_parameters( client )
+ prepare_kwds = {}
+ if rewrite_parameters:
+ compute_environment = PulsarComputeEnvironment( 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 = PulsarJobRunner.__remote_metadata( client )
+ dependency_resolution = PulsarJobRunner.__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 Pulsar, always worked for Pulsar 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 _populate_parameter_defaults( self, job_destination ):
+ updated = False
+ params = job_destination.params
+ for key, value in self.destination_defaults.iteritems():
+ if key in params:
+ if value is PARAMETER_SPECIFICATION_IGNORED:
+ log.warn( "Pulsar runner in selected configuration ignores parameter %s" % key )
+ continue
+ #if self.runner_params.get( key, None ):
+ # # Let plugin define defaults for some parameters -
+ # # for instance that way jobs_directory can be
+ # # configured next to AMQP url (where it belongs).
+ # params[ key ] = self.runner_params[ key ]
+ # continue
+
+ if not value:
+ continue
+
+ if value is PARAMETER_SPECIFICATION_REQUIRED:
+ raise Exception( "Pulsar destination does not define required parameter %s" % key )
+ elif value is not PARAMETER_SPECIFICATION_IGNORED:
+ params[ key ] = value
+ updated = True
+ return updated
+
+ 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)
+ pulsar_outputs = PulsarOutputs.from_status_response(run_results)
+ # Use Pulsar 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,
+ pulsar_outputs=pulsar_outputs )
+ failed = pulsar_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 PulsarJobRunner.__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, 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, 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
+ pulsar_url = job.job_runner_name
+ job_id = job.job_runner_external_id
+ log.debug("Attempt remote Pulsar kill of job with url %s and id %s" % (pulsar_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( "(Pulsar/%s) is still in running state, adding to the Pulsar 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( PulsarJobRunner, 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( pulsar_client, job_wrapper ):
+ dependency_resolution = PulsarJobRunner.__dependency_resolution( pulsar_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( pulsar_client ):
+ dependency_resolution = pulsar_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( pulsar_client ):
+ remote_metadata = string_as_bool_or_none( pulsar_client.destination_params.get( "remote_metadata", False ) )
+ return remote_metadata
+
+ @staticmethod
+ def __use_remote_datatypes_conf( pulsar_client ):
+ """ When setting remote metadata, use integrated datatypes from this
+ Galaxy instance or use the datatypes config configured via the remote
+ Pulsar.
+
+ 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( pulsar_client.destination_params.get( "use_remote_datatypes", False ) )
+ return use_remote_datatypes
+
+ @staticmethod
+ def __rewrite_parameters( pulsar_client ):
+ return string_as_bool_or_none( pulsar_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
+ # Pulsar disables from_work_dir copying as part of the job command
+ # line we need to take the list of output locations on the Pulsar
+ # server (produced by self.get_output_files(job_wrapper)) and for
+ # each work_dir output substitute the effective path on the Pulsar
+ # 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 pulsar_workdir_path, real_path in work_dir_outputs:
+ if real_path == output.real_path:
+ output.false_path = pulsar_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, 'universe_wsgi.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 PulsarJobRunner.__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 PulsarLegacyJobRunner( PulsarJobRunner ):
+ destination_defaults = dict(
+ rewrite_parameters="false",
+ dependency_resolution="local",
+ )
+
+
+class PulsarMQJobRunner( PulsarJobRunner ):
+ destination_defaults = dict(
+ default_file_action="remote_transfer",
+ rewrite_parameters="true",
+ dependency_resolution="remote",
+ jobs_directory=PARAMETER_SPECIFICATION_REQUIRED,
+ url=PARAMETER_SPECIFICATION_IGNORED,
+ private_token=PARAMETER_SPECIFICATION_IGNORED
+ )
+
+ def _monitor( self ):
+ # 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)
+
+ def __async_update( self, full_status ):
+ 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_status(job_state, full_status[ "status" ] )
+ except Exception:
+ log.exception( "Failed to update Pulsar job status for job_id %s" % job_id )
+ raise
+ # Nothing else to do? - Attempt to fail the job?
+
+
+class PulsarRESTJobRunner( PulsarJobRunner ):
+ destination_defaults = dict(
+ default_file_action="transfer",
+ rewrite_parameters="true",
+ dependency_resolution="remote",
+ url=PARAMETER_SPECIFICATION_REQUIRED,
+ )
+
+
+class PulsarComputeEnvironment( ComputeEnvironment ):
+
+ def __init__( self, pulsar_client, job_wrapper, remote_job_config ):
+ self.pulsar_client = pulsar_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(pulsar_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/util/__init__.py b/lib/galaxy/jobs/runners/util/__init__.py
index 79ec002b4b9..aff63231255 100644
--- a/lib/galaxy/jobs/runners/util/__init__.py
+++ b/lib/galaxy/jobs/runners/util/__init__.py
@@ -1,7 +1,7 @@
"""
This module and its submodules contains utilities for running external
processes and interfacing with job managers. This module should contain
-functionality shared between Galaxy and the LWR.
+functionality shared between Galaxy and the Pulsar.
"""
from galaxy.util.bunch import Bunch
diff --git a/lib/galaxy/jobs/runners/util/cli/factory.py b/lib/galaxy/jobs/runners/util/cli/factory.py
index 517389cd16c..6908e4f7ef7 100644
--- a/lib/galaxy/jobs/runners/util/cli/factory.py
+++ b/lib/galaxy/jobs/runners/util/cli/factory.py
@@ -5,7 +5,7 @@ try:
)
code_dir = 'lib'
except ImportError:
- from lwr.managers.util.cli import (
+ from pulsar.managers.util.cli import (
CliInterface,
split_params
)
diff --git a/lib/galaxy/jobs/runners/util/cli/job/slurm.py b/lib/galaxy/jobs/runners/util/cli/job/slurm.py
index 95350d73e1b..c3fe43793ef 100644
--- a/lib/galaxy/jobs/runners/util/cli/job/slurm.py
+++ b/lib/galaxy/jobs/runners/util/cli/job/slurm.py
@@ -5,7 +5,7 @@ try:
from galaxy.model import Job
job_states = Job.states
except ImportError:
- # Not in Galaxy, map Galaxy job states to LWR ones.
+ # Not in Galaxy, map Galaxy job states to Pulsar ones.
from galaxy.util import enum
job_states = enum(RUNNING='running', OK='complete', QUEUED='queued')
diff --git a/lib/galaxy/jobs/runners/util/cli/job/torque.py b/lib/galaxy/jobs/runners/util/cli/job/torque.py
index 9159827b670..6cf7bcae3f5 100644
--- a/lib/galaxy/jobs/runners/util/cli/job/torque.py
+++ b/lib/galaxy/jobs/runners/util/cli/job/torque.py
@@ -7,7 +7,7 @@ try:
from galaxy.model import Job
job_states = Job.states
except ImportError:
- # Not in Galaxy, map Galaxy job states to LWR ones.
+ # Not in Galaxy, map Galaxy job states to Pulsar ones.
from galaxy.util import enum
job_states = enum(RUNNING='running', OK='complete', QUEUED='queued')
diff --git a/lib/galaxy/objectstore/__init__.py b/lib/galaxy/objectstore/__init__.py
index 41698401d11..11fe27b3bc9 100644
--- a/lib/galaxy/objectstore/__init__.py
+++ b/lib/galaxy/objectstore/__init__.py
@@ -623,9 +623,9 @@ def build_object_store_from_config(config, fsmon=False, config_xml=None):
elif store == 'irods':
from .rods import IRODSObjectStore
return IRODSObjectStore(config=config, config_xml=config_xml)
- elif store == 'lwr':
- from .lwr import LwrObjectStore
- return LwrObjectStore(config=config, config_xml=config_xml)
+ elif store == 'pulsar':
+ from .pulsar import PulsarObjectStore
+ return PulsarObjectStore(config=config, config_xml=config_xml)
else:
log.error("Unrecognized object store definition: {0}".format(store))
diff --git a/lib/galaxy/objectstore/lwr.py b/lib/galaxy/objectstore/pulsar.py
similarity index 50%
rename from lib/galaxy/objectstore/lwr.py
rename to lib/galaxy/objectstore/pulsar.py
index d051deeb52d..08f51ea14dd 100644
--- a/lib/galaxy/objectstore/lwr.py
+++ b/lib/galaxy/objectstore/pulsar.py
@@ -1,20 +1,17 @@
-from __future__ import absolute_import # Need to import lwr_client absolutely.
+from __future__ import absolute_import # Need to import pulsar_client absolutely.
from ..objectstore import ObjectStore
-try:
- from galaxy.jobs.runners.lwr_client.manager import ObjectStoreClientManager
-except ImportError:
- from lwr.lwr_client.manager import ObjectStoreClientManager
+from pulsar.client.manager import ObjectStoreClientManager
-class LwrObjectStore(ObjectStore):
+class PulsarObjectStore(ObjectStore):
"""
- Object store implementation that delegates to a remote LWR server.
+ Object store implementation that delegates to a remote Pulsar server.
This may be more aspirational than practical for now, it would be good to
Galaxy to a point that a handler thread could be setup that doesn't attempt
to access the disk files returned by a (this) object store - just passing
- them along to the LWR unmodified. That modification - along with this
- implementation and LWR job destinations would then allow Galaxy to fully
+ them along to the Pulsar unmodified. That modification - along with this
+ implementation and Pulsar job destinations would then allow Galaxy to fully
manage jobs on remote servers with completely different mount points.
This implementation should be considered beta and may be dropped from
@@ -22,38 +19,38 @@ class LwrObjectStore(ObjectStore):
"""
def __init__(self, config, config_xml):
- self.lwr_client = self.__build_lwr_client(config_xml)
+ self.pulsar_client = self.__build_pulsar_client(config_xml)
def exists(self, obj, **kwds):
- return self.lwr_client.exists(**self.__build_kwds(obj, **kwds))
+ return self.pulsar_client.exists(**self.__build_kwds(obj, **kwds))
def file_ready(self, obj, **kwds):
- return self.lwr_client.file_ready(**self.__build_kwds(obj, **kwds))
+ return self.pulsar_client.file_ready(**self.__build_kwds(obj, **kwds))
def create(self, obj, **kwds):
- return self.lwr_client.create(**self.__build_kwds(obj, **kwds))
+ return self.pulsar_client.create(**self.__build_kwds(obj, **kwds))
def empty(self, obj, **kwds):
- return self.lwr_client.empty(**self.__build_kwds(obj, **kwds))
+ return self.pulsar_client.empty(**self.__build_kwds(obj, **kwds))
def size(self, obj, **kwds):
- return self.lwr_client.size(**self.__build_kwds(obj, **kwds))
+ return self.pulsar_client.size(**self.__build_kwds(obj, **kwds))
def delete(self, obj, **kwds):
- return self.lwr_client.delete(**self.__build_kwds(obj, **kwds))
+ return self.pulsar_client.delete(**self.__build_kwds(obj, **kwds))
# TODO: Optimize get_data.
def get_data(self, obj, **kwds):
- return self.lwr_client.get_data(**self.__build_kwds(obj, **kwds))
+ return self.pulsar_client.get_data(**self.__build_kwds(obj, **kwds))
def get_filename(self, obj, **kwds):
- return self.lwr_client.get_filename(**self.__build_kwds(obj, **kwds))
+ return self.pulsar_client.get_filename(**self.__build_kwds(obj, **kwds))
def update_from_file(self, obj, **kwds):
- return self.lwr_client.update_from_file(**self.__build_kwds(obj, **kwds))
+ return self.pulsar_client.update_from_file(**self.__build_kwds(obj, **kwds))
def get_store_usage_percent(self):
- return self.lwr_client.get_store_usage_percent()
+ return self.pulsar_client.get_store_usage_percent()
def get_object_url(self, obj, extra_dir=None, extra_dir_at_root=False, alt_name=None):
return None
@@ -63,14 +60,14 @@ class LwrObjectStore(ObjectStore):
return kwds
pass
- def __build_lwr_client(self, config_xml):
+ def __build_pulsar_client(self, config_xml):
url = config_xml.get("url")
private_token = config_xml.get("private_token", None)
transport = config_xml.get("transport", None)
manager_options = dict(transport=transport)
client_options = dict(url=url, private_token=private_token)
- lwr_client = ObjectStoreClientManager(**manager_options).get_client(client_options)
- return lwr_client
+ pulsar_client = ObjectStoreClientManager(**manager_options).get_client(client_options)
+ return pulsar_client
def shutdown(self):
pass
diff --git a/lib/galaxy/tools/deps/dependencies.py b/lib/galaxy/tools/deps/dependencies.py
index 060b89cbbf7..6dd372a76fa 100644
--- a/lib/galaxy/tools/deps/dependencies.py
+++ b/lib/galaxy/tools/deps/dependencies.py
@@ -8,7 +8,7 @@ class DependenciesDescription(object):
related context required to resolve dependencies via the
ToolShedPackageDependencyResolver.
- This is meant to enable remote resolution of dependencies, by the LWR or
+ This is meant to enable remote resolution of dependencies, by the Pulsar or
other potential remote execution mechanisms.
"""
@@ -39,7 +39,7 @@ class DependenciesDescription(object):
@staticmethod
def _toolshed_install_dependency_from_dict(as_dict):
- # Rather than requiring full models in LWR, just use simple objects
+ # Rather than requiring full models in Pulsar, just use simple objects
# containing only properties and associations used to resolve
# dependencies for tool execution.
repository_object = bunch.Bunch(
diff --git a/lib/pulsar/__init__.py b/lib/pulsar/__init__.py
new file mode 100644
index 00000000000..e69de29bb2d
diff --git a/lib/pulsar/client/__init__.py b/lib/pulsar/client/__init__.py
new file mode 100644
index 00000000000..46acfe63445
--- /dev/null
+++ b/lib/pulsar/client/__init__.py
@@ -0,0 +1,62 @@
+"""
+pulsar client
+======
+
+This module contains logic for interfacing with an external Pulsar server.
+
+------------------
+Configuring Galaxy
+------------------
+
+Galaxy job runners are configured in Galaxy's ``job_conf.xml`` file. See ``job_conf.xml.sample_advanced``
+in your Galaxy code base or on
+`Bitbucket `_
+for information on how to configure Galaxy to interact with the Pulsar.
+
+Galaxy also supports an older, less rich configuration of job runners directly
+in its main ``universe_wsgi.ini`` file. The following section describes how to
+configure Galaxy to communicate with the Pulsar in this legacy mode.
+
+Legacy
+------
+
+A Galaxy tool can be configured to be executed remotely via Pulsar by
+adding a line to the ``universe_wsgi.ini`` file under the
+``galaxy:tool_runners`` section with the format::
+
+ = pulsar://http://:
+
+As an example, if a host named remotehost is running the Pulsar server
+application on port ``8913``, then the tool with id ``test_tool`` can
+be configured to run remotely on remotehost by adding the following
+line to ``universe.ini``::
+
+ test_tool = pulsar://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 PulsarOutputs
+from .staging import ClientOutputs
+from .client import OutputNotFoundException
+from .manager import build_client_manager
+from .destination import url_to_destination_params
+from .path_mapper import PathMapper
+
+__all__ = [
+ build_client_manager,
+ OutputNotFoundException,
+ url_to_destination_params,
+ finish_job,
+ submit_job,
+ ClientJobDescription,
+ PulsarOutputs,
+ ClientOutputs,
+ PathMapper,
+]
diff --git a/lib/pulsar/client/action_mapper.py b/lib/pulsar/client/action_mapper.py
new file mode 100644
index 00000000000..893ac1a4409
--- /dev/null
+++ b/lib/pulsar/client/action_mapper.py
@@ -0,0 +1,567 @@
+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 Pulsar client (i.e. Galaxy server) and remote Pulsar server with the same
+ paths. """
+ action_type = "none"
+ staging = STAGING_ACTION_NONE
+
+ def to_dict(self):
+ return dict(path=self.path, action_type=self.action_type)
+
+ @classmethod
+ def from_dict(cls, action_dict):
+ return NoneAction(path=action_dict["path"])
+
+ def path_rewrite(self, path_helper, path=None):
+ return None
+
+
+class RewriteAction(BaseAction):
+ """ This actin indicates the Pulsar server should simply rewrite the path
+ to the specified file.
+ """
+ action_spec = dict(
+ source_directory=REQUIRED_ACTION_KWD,
+ destination_directory=REQUIRED_ACTION_KWD
+ )
+ action_type = "rewrite"
+ staging = STAGING_ACTION_NONE
+
+ def __init__(self, path, file_lister=None, source_directory=None, destination_directory=None):
+ self.path = path
+ self.file_lister = file_lister or DEFAULT_FILE_LISTER
+ self.source_directory = source_directory
+ self.destination_directory = destination_directory
+
+ def to_dict(self):
+ return dict(
+ path=self.path,
+ action_type=self.action_type,
+ source_directory=self.source_directory,
+ destination_directory=self.destination_directory,
+ )
+
+ @classmethod
+ def from_dict(cls, action_dict):
+ return RewriteAction(
+ path=action_dict["path"],
+ source_directory=action_dict["source_directory"],
+ destination_directory=action_dict["destination_directory"],
+ )
+
+ def path_rewrite(self, path_helper, path=None):
+ if not path:
+ path = self.path
+ new_path = path_helper.from_posix_with_new_base(self.path, self.source_directory, self.destination_directory)
+ return None if new_path == self.path else new_path
+
+
+class TransferAction(BaseAction):
+ """ This actions indicates that the Pulsar client should initiate an HTTP
+ transfer of the corresponding path to the remote Pulsar server before
+ launching the job. """
+ action_type = "transfer"
+ staging = STAGING_ACTION_LOCAL
+
+
+class CopyAction(BaseAction):
+ """ This action indicates that the Pulsar client should execute a file system
+ copy of the corresponding path to the Pulsar staging directory prior to
+ launching the corresponding job. """
+ action_type = "copy"
+ staging = STAGING_ACTION_LOCAL
+
+
+class RemoteCopyAction(BaseAction):
+ """ This action indicates the Pulsar server should copy the file before
+ execution via direct file system copy. This is like a CopyAction, but
+ it indicates the action should occur on the Pulsar server instead of on
+ the client.
+ """
+ action_type = "remote_copy"
+ staging = STAGING_ACTION_REMOTE
+
+ def to_dict(self):
+ return dict(path=self.path, action_type=self.action_type)
+
+ @classmethod
+ def from_dict(cls, action_dict):
+ return RemoteCopyAction(path=action_dict["path"])
+
+ def write_to_path(self, path):
+ galaxy.util.copy_to_path(open(self.path, "rb"), path)
+
+ def write_from_path(self, pulsar_path):
+ destination = self.path
+ parent_directory = dirname(destination)
+ if not exists(parent_directory):
+ makedirs(parent_directory)
+ with open(pulsar_path, "rb") as f:
+ galaxy.util.copy_to_path(f, destination)
+
+
+class RemoteTransferAction(BaseAction):
+ """ This action indicates the Pulsar server should copy the file before
+ execution via direct file system copy. This is like a CopyAction, but
+ it indicates the action should occur on the Pulsar server instead of on
+ the client.
+ """
+ action_type = "remote_transfer"
+ staging = STAGING_ACTION_REMOTE
+
+ def __init__(self, path, file_lister=None, url=None):
+ super(RemoteTransferAction, self).__init__(path, file_lister=file_lister)
+ self.url = url
+
+ def to_dict(self):
+ return dict(path=self.path, action_type=self.action_type, url=self.url)
+
+ @classmethod
+ def from_dict(cls, action_dict):
+ return RemoteTransferAction(path=action_dict["path"], url=action_dict["url"])
+
+ def write_to_path(self, path):
+ get_file(self.url, path)
+
+ def write_from_path(self, pulsar_path):
+ post_file(self.url, pulsar_path)
+
+
+class 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/pulsar/client/amqp_exchange.py b/lib/pulsar/client/amqp_exchange.py
new file mode 100644
index 00000000000..fd12ace19c6
--- /dev/null
+++ b/lib/pulsar/client/amqp_exchange.py
@@ -0,0 +1,139 @@
+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 = "pulsar"
+DEFAULT_EXCHANGE_TYPE = "direct"
+# Set timeout to periodically give up looking and check if polling should end.
+DEFAULT_TIMEOUT = 0.2
+DEFAULT_HEARTBEAT = 580
+
+DEFAULT_RECONNECT_CONSUMER_WAIT = 1
+DEFAULT_HEARTBEAT_WAIT = 1
+
+
+class PulsarExchange(object):
+ """ Utility for publishing and consuming structured Pulsar queues using kombu.
+ This is shared between the server and client - an exchange should be setup
+ for each manager (or in the case of the client, each manager one wished to
+ communicate with.)
+
+ Each Pulsar manager is defined solely by name in the scheme, so only one Pulsar
+ should target each AMQP endpoint or care should be taken that unique
+ manager names are used across Pulsar servers targetting same AMQP endpoint -
+ and in particular only one such Pulsar should define an default manager with
+ name _default_.
+ """
+
+ def __init__(
+ self,
+ url,
+ manager_name,
+ connect_ssl=None,
+ timeout=DEFAULT_TIMEOUT,
+ publish_kwds={},
+ ):
+ """
+ """
+ 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), 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 = "pulsar_"
+ else:
+ key_prefix = "pulsar_%s_" % self.__manager_name
+ return key_prefix
+
+ def __start_heartbeat(self, queue_name, connection):
+ thread_name = "consume-heartbeat-%s" % (self.__queue_name(queue_name))
+ thread = threading.Thread(name=thread_name, target=self.heartbeat, args=(connection,))
+ thread.start()
+ return thread
diff --git a/lib/pulsar/client/amqp_exchange_factory.py b/lib/pulsar/client/amqp_exchange_factory.py
new file mode 100644
index 00000000000..df26f8f5853
--- /dev/null
+++ b/lib/pulsar/client/amqp_exchange_factory.py
@@ -0,0 +1,41 @@
+from .amqp_exchange import PulsarExchange
+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 = PulsarExchange(url, **exchange_kwds)
+ return exchange
+
+
+def parse_amqp_connect_ssl_params(params):
+ ssl_params = filter_destination_params(params, "amqp_connect_ssl_")
+ if not ssl_params:
+ return
+
+ ssl = __import__('ssl')
+ if 'cert_reqs' in ssl_params:
+ value = ssl_params['cert_reqs']
+ ssl_params['cert_reqs'] = getattr(ssl, value.upper())
+ return ssl_params
+
+
+def parse_amqp_publish_kwds(params):
+ all_publish_params = filter_destination_params(params, "amqp_publish_")
+ retry_policy_params = {}
+ for key in all_publish_params.keys():
+ if key.startswith("retry_"):
+ value = all_publish_params[key]
+ retry_policy_params[key[len("retry_"):]] = value
+ del all_publish_params[key]
+ if retry_policy_params:
+ all_publish_params["retry_policy"] = retry_policy_params
+ return all_publish_params
diff --git a/lib/pulsar/client/client.py b/lib/pulsar/client/client.py
new file mode 100644
index 00000000000..0cdd55c916e
--- /dev/null
+++ b/lib/pulsar/client/client.py
@@ -0,0 +1,381 @@
+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
+from .action_mapper import path_type
+
+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 Pulsar server to run this job.
+ """
+ setup_args = {"job_id": self.job_id}
+ if tool_id:
+ setup_args["tool_id"] = tool_id
+ if tool_version:
+ setup_args["tool_version"] = tool_version
+ return self.setup_handler.setup(**setup_args)
+
+ @property
+ def prefer_local_staging(self):
+ # If doing a job directory is defined, calculate paths here and stage
+ # remotely.
+ return self.job_directory is None
+
+
+class JobClient(BaseJobClient):
+ """
+ Objects of this client class perform low-level communication with a remote Pulsar server.
+
+ **Parameters**
+
+ destination_params : dict or str
+ connection parameters, either url with dict containing url (and optionally `private_token`).
+ job_id : str
+ Galaxy job/task id.
+ """
+
+ def __init__(self, destination_params, job_id, job_manager_interface):
+ super(JobClient, self).__init__(destination_params, job_id)
+ self.job_manager_interface = job_manager_interface
+
+ def launch(self, command_line, dependencies_description=None, env=[], remote_staging=[], job_config=None):
+ """
+ Queue up the execution of the supplied `command_line` on the remote
+ server. Called launch for historical reasons, should be renamed to
+ enqueue or something like that.
+
+ **Parameters**
+
+ command_line : str
+ Command to execute.
+ """
+ launch_params = dict(command_line=command_line, job_id=self.job_id)
+ submit_params_dict = submit_params(self.destination_params)
+ if submit_params_dict:
+ launch_params['params'] = 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 Pulsar setup job
+ # before queueing.
+ setup_params = _setup_params_from_job_config(job_config)
+ launch_params["setup_params"] = dumps(setup_params)
+ return self._raw_execute("submit", launch_params)
+
+ def full_status(self):
+ """ Return a dictionary summarizing final state of job.
+ """
+ return self.raw_check_complete()
+
+ def kill(self):
+ """
+ Cancel remote job, either removing from the queue or killing it.
+ """
+ return self._raw_execute("cancel", {"job_id": self.job_id})
+
+ @retry()
+ @parseJson()
+ def raw_check_complete(self):
+ """
+ Get check_complete response from the remote server.
+ """
+ check_complete_response = self._raw_execute("status", {"job_id": self.job_id})
+ return check_complete_response
+
+ def get_status(self):
+ check_complete_response = self.raw_check_complete()
+ # Older Pulsar instances won't set status so use 'complete', at some
+ # point drop backward compatibility.
+ status = check_complete_response.get("status", None)
+ return status
+
+ def clean(self):
+ """
+ Cleanup the remote job.
+ """
+ self._raw_execute("clean", {"job_id": self.job_id})
+
+ @parseJson()
+ def remote_setup(self, **setup_args):
+ """
+ Setup remote Pulsar server to run this job.
+ """
+ return self._raw_execute("setup", setup_args)
+
+ def put_file(self, path, input_type, name=None, contents=None, action_type='transfer'):
+ if not name:
+ name = os.path.basename(path)
+ args = {"job_id": self.job_id, "name": name, "type": input_type}
+ input_path = path
+ if contents:
+ input_path = None
+ if action_type == 'transfer':
+ return self._upload_file(args, contents, input_path)
+ elif action_type == 'copy':
+ pulsar_path = self._raw_execute('path', args)
+ copy(path, pulsar_path)
+ return {'path': pulsar_path}
+
+ def fetch_output(self, path, name, working_directory, action_type, output_type):
+ """
+ Fetch (transfer, copy, etc...) an output from the remote Pulsar server.
+
+ **Parameters**
+
+ path : str
+ Local path of the dataset.
+ name : str
+ Remote name of file (i.e. path relative to remote staging output
+ or working directory).
+ working_directory : str
+ Local working_directory for the job.
+ action_type : str
+ Where to find file on Pulsar (output_workdir or output). legacy is also
+ an option in this case Pulsar is asked for location - this will only be
+ used if targetting an older Pulsar server that didn't return statuses
+ allowing this to be inferred.
+ """
+ if output_type == 'output_workdir':
+ self._fetch_work_dir_output(name, working_directory, path, action_type=action_type)
+ elif output_type == 'output':
+ self._fetch_output(path=path, name=name, action_type=action_type)
+ else:
+ raise Exception("Unknown output_type %s" % output_type)
+
+ def _raw_execute(self, command, args={}, data=None, input_path=None, output_path=None):
+ return self.job_manager_interface.execute(command, args, data, input_path, output_path)
+
+ def _fetch_output(self, path, name=None, check_exists_remotely=False, action_type='transfer'):
+ if not name:
+ # Extra files will send in the path.
+ name = os.path.basename(path)
+
+ self.__populate_output_path(name, path, action_type)
+
+ def _fetch_work_dir_output(self, name, working_directory, output_path, action_type='transfer'):
+ ensure_directory(output_path)
+ if action_type == 'transfer':
+ self.__raw_download_output(name, self.job_id, path_type.OUTPUT_WORKDIR, output_path)
+ else: # Even if action is none - Pulsar has a different work_dir so this needs to be copied.
+ pulsar_path = self._output_path(name, self.job_id, path_type.OUTPUT_WORKDIR)['path']
+ copy(pulsar_path, output_path)
+
+ def __populate_output_path(self, name, output_path, action_type):
+ ensure_directory(output_path)
+ if action_type == 'transfer':
+ self.__raw_download_output(name, self.job_id, path_type.OUTPUT, output_path)
+ elif action_type == 'copy':
+ pulsar_path = self._output_path(name, self.job_id, path_type.OUTPUT)['path']
+ copy(pulsar_path, output_path)
+
+ @parseJson()
+ def _upload_file(self, args, contents, input_path):
+ return self._raw_execute("upload_file", args, contents, input_path)
+
+ @parseJson()
+ def _output_path(self, name, job_id, output_type):
+ return self._raw_execute("path",
+ {"name": name,
+ "job_id": self.job_id,
+ "type": output_type})
+
+ @retry()
+ def __raw_download_output(self, name, job_id, output_type, output_path):
+ output_params = {
+ "name": name,
+ "job_id": self.job_id,
+ "type": output_type
+ }
+ self._raw_execute("download_output", output_params, output_path=output_path)
+
+
+class BaseMessageJobClient(BaseJobClient):
+
+ def __init__(self, destination_params, job_id, client_manager):
+ super(BaseMessageJobClient, self).__init__(destination_params, job_id)
+ if not self.job_directory:
+ error_message = "Message-queue based Pulsar client requires destination define a remote job_directory to stage files into."
+ raise Exception(error_message)
+ self.client_manager = client_manager
+
+ def clean(self):
+ del self.client_manager.status_cache[self.job_id]
+
+ def full_status(self):
+ full_status = self.client_manager.status_cache.get(self.job_id, None)
+ if full_status is None:
+ raise Exception("full_status() called before a final status was properly cached with cilent manager.")
+ return full_status
+
+ def _build_setup_message(self, command_line, dependencies_description, env, remote_staging, job_config):
+ """
+ """
+ launch_params = dict(command_line=command_line, job_id=self.job_id)
+ submit_params_dict = submit_params(self.destination_params)
+ if submit_params_dict:
+ launch_params['submit_params'] = submit_params_dict
+ if dependencies_description:
+ launch_params['dependencies_description'] = dependencies_description.to_dict()
+ if env:
+ launch_params['env'] = env
+ if remote_staging:
+ launch_params['remote_staging'] = remote_staging
+ if job_config and self.setup_handler.local:
+ # Setup not yet called, job properties were inferred from
+ # destination arguments. Hence, must have Pulsar setup job
+ # before queueing.
+ setup_params = _setup_params_from_job_config(job_config)
+ launch_params["setup_params"] = setup_params
+ return launch_params
+
+
+class MessageJobClient(BaseMessageJobClient):
+
+ def launch(self, command_line, dependencies_description=None, env=[], remote_staging=[], job_config=None):
+ """
+ """
+ launch_params = self._build_setup_message(
+ command_line,
+ dependencies_description=dependencies_description,
+ env=env,
+ remote_staging=remote_staging,
+ job_config=job_config
+ )
+ response = self.client_manager.exchange.publish("setup", launch_params)
+ log.info("Job published to setup message queue.")
+ return response
+
+ def kill(self):
+ self.client_manager.exchange.publish("kill", dict(job_id=self.job_id))
+
+
+class MessageCLIJobClient(BaseMessageJobClient):
+
+ def __init__(self, destination_params, job_id, client_manager, shell):
+ super(MessageCLIJobClient, self).__init__(destination_params, job_id, client_manager)
+ self.remote_pulsar_path = destination_params["remote_pulsar_path"]
+ self.shell = shell
+
+ def launch(self, command_line, dependencies_description=None, env=[], remote_staging=[], job_config=None):
+ """
+ """
+ launch_params = self._build_setup_message(
+ command_line,
+ dependencies_description=dependencies_description,
+ env=env,
+ remote_staging=remote_staging,
+ job_config=job_config
+ )
+ base64_message = to_base64_json(launch_params)
+ submit_command = os.path.join(self.remote_pulsar_path, "scripts", "submit.bash")
+ # TODO: Allow configuration of manager, app, and ini path...
+ self.shell.execute("nohup %s --base64 %s &" % (submit_command, base64_message))
+
+ def kill(self):
+ # TODO
+ pass
+
+
+class InputCachingJobClient(JobClient):
+ """
+ Beta client that cache's staged files to prevent duplication.
+ """
+
+ def __init__(self, destination_params, job_id, job_manager_interface, client_cacher):
+ super(InputCachingJobClient, self).__init__(destination_params, job_id, job_manager_interface)
+ self.client_cacher = client_cacher
+
+ @parseJson()
+ def _upload_file(self, args, contents, input_path):
+ action = "upload_file"
+ if contents:
+ input_path = None
+ return self._raw_execute(action, args, contents, input_path)
+ else:
+ event_holder = self.client_cacher.acquire_event(input_path)
+ cache_required = self.cache_required(input_path)
+ if cache_required:
+ self.client_cacher.queue_transfer(self, input_path)
+ while not event_holder.failed:
+ available = self.file_available(input_path)
+ if available['ready']:
+ token = available['token']
+ args["cache_token"] = token
+ return self._raw_execute(action, args)
+ event_holder.event.wait(30)
+ if event_holder.failed:
+ raise Exception("Failed to transfer file %s" % input_path)
+
+ @parseJson()
+ def cache_required(self, path):
+ return self._raw_execute("cache_required", {"path": path})
+
+ @parseJson()
+ def cache_insert(self, path):
+ return self._raw_execute("cache_insert", {"path": path}, None, path)
+
+ @parseJson()
+ def file_available(self, path):
+ return self._raw_execute("file_available", {"path": path})
+
+
+def _setup_params_from_job_config(job_config):
+ job_id = job_config.get("job_id", None)
+ tool_id = job_config.get("tool_id", None)
+ tool_version = job_config.get("tool_version", None)
+ return dict(
+ job_id=job_id,
+ tool_id=tool_id,
+ tool_version=tool_version
+ )
diff --git a/lib/pulsar/client/config_util.py b/lib/pulsar/client/config_util.py
new file mode 100644
index 00000000000..ab5453d0c8f
--- /dev/null
+++ b/lib/pulsar/client/config_util.py
@@ -0,0 +1,80 @@
+""" 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/pulsar/client/decorators.py b/lib/pulsar/client/decorators.py
new file mode 100644
index 00000000000..94035f5ba1b
--- /dev/null
+++ b/lib/pulsar/client/decorators.py
@@ -0,0 +1,35 @@
+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/pulsar/client/destination.py b/lib/pulsar/client/destination.py
new file mode 100644
index 00000000000..cee7a40f630
--- /dev/null
+++ b/lib/pulsar/client/destination.py
@@ -0,0 +1,58 @@
+
+from re import match
+from .util import filter_destination_params
+
+SUBMIT_PREFIX = "submit_"
+
+
+def url_to_destination_params(url):
+ """Convert a legacy runner URL to a job destination
+
+ >>> params_simple = url_to_destination_params("http://localhost:8913/")
+ >>> params_simple["url"]
+ 'http://localhost:8913/'
+ >>> params_simple["private_token"] is None
+ True
+ >>> advanced_url = "https://1234x@example.com:8914/managers/longqueue"
+ >>> params_advanced = url_to_destination_params(advanced_url)
+ >>> params_advanced["url"]
+ 'https://example.com:8914/managers/longqueue/'
+ >>> params_advanced["private_token"]
+ '1234x'
+ >>> runner_url = "pulsar://http://localhost:8913/"
+ >>> runner_params = url_to_destination_params(runner_url)
+ >>> runner_params['url']
+ 'http://localhost:8913/'
+ """
+
+ if url.startswith("pulsar://"):
+ url = url[len("pulsar://"):]
+
+ if not url.endswith("/"):
+ url += "/"
+
+ # Check for private token embedded in the URL. A URL of the form
+ # https://moo@cow:8913 will try to contact https://cow:8913
+ # with a private key of moo
+ private_token_format = "https?://(.*)@.*/?"
+ private_token_match = match(private_token_format, url)
+ private_token = None
+ if private_token_match:
+ private_token = private_token_match.group(1)
+ url = url.replace("%s@" % private_token, '', 1)
+
+ destination_args = {"url": url,
+ "private_token": private_token}
+
+ return destination_args
+
+
+def submit_params(destination_params):
+ """
+
+ >>> destination_params = {"private_token": "12345", "submit_native_specification": "-q batch"}
+ >>> result = submit_params(destination_params)
+ >>> result
+ {'native_specification': '-q batch'}
+ """
+ return filter_destination_params(destination_params, SUBMIT_PREFIX)
diff --git a/lib/pulsar/client/interface.py b/lib/pulsar/client/interface.py
new file mode 100644
index 00000000000..af13a02e4f9
--- /dev/null
+++ b/lib/pulsar/client/interface.py
@@ -0,0 +1,152 @@
+from abc import ABCMeta
+from abc import abstractmethod
+from string import Template
+
+try:
+ from StringIO import StringIO as BytesIO
+except ImportError:
+ from io import BytesIO
+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 PulsarInterface(object):
+ """
+ Abstract base class describes how synchronous client communicates with
+ (potentially remote) Pulsar procedures. Obvious implementation is HTTP based
+ but Pulsar objects wrapped in routes can also be directly communicated with
+ if in memory.
+ """
+ __metaclass__ = ABCMeta
+
+ @abstractmethod
+ def execute(self, command, args={}, data=None, input_path=None, output_path=None):
+ """
+ Execute the correspond command against configured Pulsar job manager. Arguments are
+ method parameters and data or input_path describe essentially POST bodies. If command
+ results in a file, resulting path should be specified as output_path.
+ """
+
+
+COMMAND_TO_PATH = {
+ "path": Template("jobs/${job_id}/files/path"),
+ "upload_file": Template("jobs/${job_id}/files"),
+ "download_output": Template("jobs/${job_id}/files"),
+
+ "setup": Template("jobs"),
+ "clean": Template("jobs/${job_id}"),
+ "status": Template("jobs/${job_id}/status"),
+ "cancel": Template("jobs/${job_id}/cancel"),
+ "submit": Template("jobs/${job_id}/submit"),
+
+ "file_available": Template("cache/status"),
+ "cache_required": Template("cache"),
+ "cache_insert": Template("cache"),
+
+ "object_store_exists": Template("objects/${object_id}/exists"),
+ "object_store_file_ready": Template("objects/${object_id}/file_ready"),
+ "object_store_update_from_file": Template("objects/${object_id}"),
+ "object_store_create": Template("objects/${object_id}"),
+ "object_store_empty": Template("objects/${object_id}/empty"),
+ "object_store_size": Template("objects/${object_id}/size"),
+ "object_store_delete": Template("objects/${object_id}"),
+ "object_store_get_data": Template("objects/${object_id}"),
+ "object_store_get_filename": Template("objects/${object_id}/filename"),
+ "object_store_get_store_usage_percent": Template("object_store_usage_percent")
+}
+
+COMMAND_TO_METHOD = {
+ "upload_file": "POST",
+ "download_output": "GET",
+
+ "setup": "POST",
+ "submit": "POST",
+ "clean": "DELETE",
+ "cancel": "PUT",
+
+ "object_store_update_from_file": "PUT",
+ "object_store_create": "POST",
+ "object_store_delete": "DELETE",
+
+ "file_available": "GET",
+ "cache_required": "PUT",
+ "cache_insert": "POST",
+}
+
+
+class HttpPulsarInterface(PulsarInterface):
+
+ def __init__(self, destination_params, transport):
+ self.transport = transport
+ remote_host = destination_params.get("url")
+ assert remote_host is not None, "Failed to determine url for Pulsar client."
+ if not remote_host.endswith("/"):
+ remote_host = "%s/" % remote_host
+ if not remote_host.startswith("http"):
+ remote_host = "http://%s" % remote_host
+ self.remote_host = remote_host
+ self.private_token = destination_params.get("private_token", None)
+
+ def execute(self, command, args={}, data=None, input_path=None, output_path=None):
+ url = self.__build_url(command, args)
+ method = COMMAND_TO_METHOD.get(command, None) # Default to GET is no data, POST otherwise
+ response = self.transport.execute(url, method=method, data=data, input_path=input_path, output_path=output_path)
+ return response
+
+ def __build_url(self, command, args):
+ path = COMMAND_TO_PATH.get(command, Template(command)).safe_substitute(args)
+ if self.private_token:
+ args["private_token"] = self.private_token
+ arg_bytes = dict([(k, text_type(args[k]).encode('utf-8')) for k in args])
+ data = urlencode(arg_bytes)
+ url = self.remote_host + path + "?" + data
+ return url
+
+
+class LocalPulsarInterface(PulsarInterface):
+
+ def __init__(self, destination_params, job_manager=None, file_cache=None, object_store=None):
+ self.job_manager = job_manager
+ self.file_cache = file_cache
+ self.object_store = object_store
+
+ def __app_args(self):
+ # Arguments that would be specified from PulsarApp if running
+ # in web server.
+ return {
+ 'manager': self.job_manager,
+ 'file_cache': self.file_cache,
+ 'object_store': self.object_store,
+ 'ip': None
+ }
+
+ def execute(self, command, args={}, data=None, input_path=None, output_path=None):
+ # If data set, should be unicode (on Python 2) or str (on Python 3).
+ from pulsar.web import routes
+ from pulsar.web.framework import build_func_args
+ controller = getattr(routes, command)
+ action = controller.func
+ body_args = dict(body=self.__build_body(data, input_path))
+ args = build_func_args(action, args.copy(), self.__app_args(), body_args)
+ result = action(**args)
+ if controller.response_type != 'file':
+ return controller.body(result)
+ else:
+ # TODO: Add to Galaxy.
+ from galaxy.util import copy_to_path
+ with open(result, 'rb') as result_file:
+ copy_to_path(result_file, output_path)
+
+ def __build_body(self, data, input_path):
+ if data is not None:
+ return BytesIO(data.encode('utf-8'))
+ elif input_path is not None:
+ return open(input_path, 'rb')
+ else:
+ return None
diff --git a/lib/pulsar/client/job_directory.py b/lib/pulsar/client/job_directory.py
new file mode 100644
index 00000000000..c4026aead9a
--- /dev/null
+++ b/lib/pulsar/client/job_directory.py
@@ -0,0 +1,135 @@
+"""
+"""
+import os.path
+from collections import deque
+import posixpath
+
+from .util import PathHelper
+from galaxy.util import in_directory
+
+from logging import getLogger
+log = getLogger(__name__)
+
+
+TYPES_TO_METHOD = dict(
+ input="inputs_directory",
+ unstructured="unstructured_files_directory",
+ config="configs_directory",
+ tool="tool_files_directory",
+ workdir="working_directory",
+ output="outputs_directory",
+ output_workdir="working_directory",
+)
+
+
+class RemoteJobDirectory(object):
+ """ Representation of a (potentially) remote Pulsar-style staging directory.
+ """
+
+ def __init__(self, remote_staging_directory, remote_id, remote_sep):
+ self.path_helper = PathHelper(remote_sep)
+ self.job_directory = self.path_helper.remote_join(
+ remote_staging_directory,
+ remote_id
+ )
+
+ def working_directory(self):
+ return self._sub_dir('working')
+
+ def inputs_directory(self):
+ return self._sub_dir('inputs')
+
+ def outputs_directory(self):
+ return self._sub_dir('outputs')
+
+ def configs_directory(self):
+ return self._sub_dir('configs')
+
+ def tool_files_directory(self):
+ return self._sub_dir('tool_files')
+
+ def unstructured_files_directory(self):
+ return self._sub_dir('unstructured')
+
+ @property
+ def path(self):
+ return self.job_directory
+
+ @property
+ def separator(self):
+ return self.path_helper.separator
+
+ def calculate_path(self, remote_relative_path, input_type):
+ """ Only for used by Pulsar client, should override for managers to
+ enforce security and make the directory if needed.
+ """
+ directory, allow_nested_files = self._directory_for_file_type(input_type)
+ return self.path_helper.remote_join(directory, remote_relative_path)
+
+ def _directory_for_file_type(self, file_type):
+ allow_nested_files = False
+ # work_dir and input_extra are types used by legacy clients...
+ # Obviously this client won't be legacy because this is in the
+ # client module, but this code is reused on server which may
+ # serve legacy clients.
+ allow_nested_files = file_type in ['input', 'unstructured', 'output', 'output_workdir']
+ directory_function = getattr(self, TYPES_TO_METHOD.get(file_type, None), None)
+ if not directory_function:
+ raise Exception("Unknown file_type specified %s" % file_type)
+ return directory_function(), allow_nested_files
+
+ def _sub_dir(self, name):
+ return self.path_helper.remote_join(self.job_directory, name)
+
+
+def get_mapped_file(directory, remote_path, allow_nested_files=False, local_path_module=os.path, mkdir=True):
+ """
+
+ >>> import ntpath
+ >>> get_mapped_file(r'C:\\pulsar\\staging\\101', 'dataset_1_files/moo/cow', allow_nested_files=True, local_path_module=ntpath, mkdir=False)
+ 'C:\\\\pulsar\\\\staging\\\\101\\\\dataset_1_files\\\\moo\\\\cow'
+ >>> get_mapped_file(r'C:\\pulsar\\staging\\101', 'dataset_1_files/moo/cow', allow_nested_files=False, local_path_module=ntpath)
+ 'C:\\\\pulsar\\\\staging\\\\101\\\\cow'
+ >>> get_mapped_file(r'C:\\pulsar\\staging\\101', '../cow', allow_nested_files=True, local_path_module=ntpath, mkdir=False)
+ Traceback (most recent call last):
+ Exception: Attempt to read or write file outside an authorized directory.
+ """
+ if not allow_nested_files:
+ name = local_path_module.basename(remote_path)
+ path = local_path_module.join(directory, name)
+ else:
+ local_rel_path = __posix_to_local_path(remote_path, local_path_module=local_path_module)
+ local_path = local_path_module.join(directory, local_rel_path)
+ verify_is_in_directory(local_path, directory, local_path_module=local_path_module)
+ local_directory = local_path_module.dirname(local_path)
+ if mkdir and not local_path_module.exists(local_directory):
+ os.makedirs(local_directory)
+ path = local_path
+ return path
+
+
+def __posix_to_local_path(path, local_path_module=os.path):
+ """
+ Converts a posix path (coming from Galaxy), to a local path (be it posix or Windows).
+
+ >>> import ntpath
+ >>> __posix_to_local_path('dataset_1_files/moo/cow', local_path_module=ntpath)
+ 'dataset_1_files\\\\moo\\\\cow'
+ >>> import posixpath
+ >>> __posix_to_local_path('dataset_1_files/moo/cow', local_path_module=posixpath)
+ 'dataset_1_files/moo/cow'
+ """
+ partial_path = deque()
+ while True:
+ if not path or path == '/':
+ break
+ (path, base) = posixpath.split(path)
+ partial_path.appendleft(base)
+ return local_path_module.join(*partial_path)
+
+
+def verify_is_in_directory(path, directory, local_path_module=os.path):
+ if not in_directory(path, directory, local_path_module):
+ msg = "Attempt to read or write file outside an authorized directory."
+ log.warn("%s Attempted path: %s, valid directory: %s" % (msg, path, directory))
+ raise Exception(msg)
diff --git a/lib/pulsar/client/manager.py b/lib/pulsar/client/manager.py
new file mode 100644
index 00000000000..066d9a126c1
--- /dev/null
+++ b/lib/pulsar/client/manager.py
@@ -0,0 +1,223 @@
+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 HttpPulsarInterface
+from .interface import LocalPulsarInterface
+from .object_client import ObjectStoreClient
+from .transport import get_transport
+from .util import TransferEventManager
+from .destination import url_to_destination_params
+from .amqp_exchange_factory import get_exchange
+
+
+from logging import getLogger
+log = getLogger(__name__)
+
+DEFAULT_TRANSFER_THREADS = 2
+
+
+def build_client_manager(**kwargs):
+ if 'job_manager' in kwargs:
+ return ClientManager(**kwargs) # TODO: Consider more separation here.
+ elif kwargs.get('amqp_url', None):
+ return MessageQueueClientManager(**kwargs)
+ else:
+ return ClientManager(**kwargs)
+
+
+class ClientManager(object):
+ """
+ Factory to create Pulsar clients, used to manage potential shared
+ state between multiple client connections.
+ """
+ def __init__(self, **kwds):
+ if 'job_manager' in kwds:
+ self.job_manager_interface_class = LocalPulsarInterface
+ self.job_manager_interface_args = dict(job_manager=kwds['job_manager'], file_cache=kwds['file_cache'])
+ else:
+ self.job_manager_interface_class = HttpPulsarInterface
+ transport_type = kwds.get('transport', None)
+ transport = get_transport(transport_type)
+ self.job_manager_interface_args = dict(transport=transport)
+ cache = kwds.get('cache', None)
+ if cache is None:
+ cache = _environ_default_int('PULSAR_CACHE_TRANSFERS')
+ if cache:
+ log.info("Setting Pulsar client class to caching variant.")
+ self.client_cacher = ClientCacher(**kwds)
+ self.client_class = InputCachingJobClient
+ self.extra_client_kwds = {"client_cacher": self.client_cacher}
+ else:
+ log.info("Setting Pulsar client class to standard, non-caching variant.")
+ self.client_class = JobClient
+ self.extra_client_kwds = {}
+
+ def get_client(self, destination_params, job_id, **kwargs):
+ destination_params = _parse_destination_params(destination_params)
+ destination_params.update(**kwargs)
+ job_manager_interface_class = self.job_manager_interface_class
+ job_manager_interface_args = dict(destination_params=destination_params, **self.job_manager_interface_args)
+ job_manager_interface = job_manager_interface_class(**job_manager_interface_args)
+ return self.client_class(destination_params, job_id, job_manager_interface, **self.extra_client_kwds)
+
+ def shutdown(self):
+ pass
+
+
+try:
+ from galaxy.jobs.runners.util.cli import factory as cli_factory
+except ImportError:
+ from pulsar.managers.util.cli import factory as cli_factory
+
+
+class MessageQueueClientManager(object):
+
+ def __init__(self, **kwds):
+ self.url = kwds.get('amqp_url')
+ self.manager_name = kwds.get("manager", None) or "_default_"
+ self.exchange = get_exchange(self.url, self.manager_name, kwds)
+ self.status_cache = {}
+ self.callback_lock = threading.Lock()
+ self.callback_thread = None
+ self.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 Pulsar.")
+ callback(body)
+ except Exception:
+ log.exception("Failure processing job status update message.")
+ except BaseException as e:
+ log.exception("Failure processing job status update message - BaseException type %s" % type(e))
+ finally:
+ message.ack()
+
+ def run():
+ self.exchange.consume("status_update", callback_wrapper, check=self)
+ log.debug("Leaving Pulsar client status update thread, no additional Pulsar updates will be processed.")
+
+ thread = threading.Thread(
+ name="pulsar_client_%s_status_update_callback" % self.manager_name,
+ target=run
+ )
+ thread.daemon = False # Lets not interrupt processing of this.
+ thread.start()
+ self.callback_thread = thread
+
+ def 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 Pulsar client for empty job_id.")
+ destination_params = _parse_destination_params(destination_params)
+ destination_params.update(**kwargs)
+ if 'shell_plugin' in destination_params:
+ shell = cli_factory.get_shell(destination_params)
+ return MessageCLIJobClient(destination_params, job_id, self, shell)
+ else:
+ return MessageJobClient(destination_params, job_id, self)
+
+
+class ObjectStoreClientManager(object):
+
+ def __init__(self, **kwds):
+ if 'object_store' in kwds:
+ self.interface_class = LocalPulsarInterface
+ self.interface_args = dict(object_store=kwds['object_store'])
+ else:
+ self.interface_class = HttpPulsarInterface
+ transport_type = kwds.get('transport', None)
+ transport = get_transport(transport_type)
+ self.interface_args = dict(transport=transport)
+ self.extra_client_kwds = {}
+
+ def get_client(self, client_params):
+ interface_class = self.interface_class
+ interface_args = dict(destination_params=client_params, **self.interface_args)
+ interface = interface_class(**interface_args)
+ return ObjectStoreClient(interface)
+
+
+class ClientCacher(object):
+
+ def __init__(self, **kwds):
+ self.event_manager = TransferEventManager()
+ default_transfer_threads = _environ_default_int('PULSAR_CACHE_THREADS', DEFAULT_TRANSFER_THREADS)
+ num_transfer_threads = int(kwds.get('transfer_threads', default_transfer_threads))
+ self.__init_transfer_threads(num_transfer_threads)
+
+ def queue_transfer(self, client, path):
+ self.transfer_queue.put((client, path))
+
+ def acquire_event(self, input_path):
+ return self.event_manager.acquire_event(input_path)
+
+ def _transfer_worker(self):
+ while True:
+ transfer_info = self.transfer_queue.get()
+ try:
+ self.__perform_transfer(transfer_info)
+ except BaseException as e:
+ log.warn("Transfer failed.")
+ log.exception(e)
+ pass
+ self.transfer_queue.task_done()
+
+ def __perform_transfer(self, transfer_info):
+ (client, path) = transfer_info
+ event_holder = self.event_manager.acquire_event(path, force_clear=True)
+ failed = True
+ try:
+ client.cache_insert(path)
+ failed = False
+ finally:
+ event_holder.failed = failed
+ event_holder.release()
+
+ def __init_transfer_threads(self, num_transfer_threads):
+ self.num_transfer_threads = num_transfer_threads
+ self.transfer_queue = Queue()
+ for i in range(num_transfer_threads):
+ t = threading.Thread(target=self._transfer_worker)
+ t.daemon = True
+ t.start()
+
+
+def _parse_destination_params(destination_params):
+ 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, HttpPulsarInterface]
diff --git a/lib/pulsar/client/object_client.py b/lib/pulsar/client/object_client.py
new file mode 100644
index 00000000000..a9d795f1b6c
--- /dev/null
+++ b/lib/pulsar/client/object_client.py
@@ -0,0 +1,53 @@
+from .decorators import parseJson
+
+
+class ObjectStoreClient(object):
+
+ def __init__(self, pulsar_interface):
+ self.pulsar_interface = pulsar_interface
+
+ @parseJson()
+ def exists(self, **kwds):
+ return self._raw_execute("object_store_exists", args=self.__data(**kwds))
+
+ @parseJson()
+ def file_ready(self, **kwds):
+ return self._raw_execute("object_store_file_ready", args=self.__data(**kwds))
+
+ @parseJson()
+ def create(self, **kwds):
+ return self._raw_execute("object_store_create", args=self.__data(**kwds))
+
+ @parseJson()
+ def empty(self, **kwds):
+ return self._raw_execute("object_store_empty", args=self.__data(**kwds))
+
+ @parseJson()
+ def size(self, **kwds):
+ return self._raw_execute("object_store_size", args=self.__data(**kwds))
+
+ @parseJson()
+ def delete(self, **kwds):
+ return self._raw_execute("object_store_delete", args=self.__data(**kwds))
+
+ @parseJson()
+ def get_data(self, **kwds):
+ return self._raw_execute("object_store_get_data", args=self.__data(**kwds))
+
+ @parseJson()
+ def get_filename(self, **kwds):
+ return self._raw_execute("object_store_get_filename", args=self.__data(**kwds))
+
+ @parseJson()
+ def update_from_file(self, **kwds):
+ return self._raw_execute("object_store_update_from_file", args=self.__data(**kwds))
+
+ @parseJson()
+ def get_store_usage_percent(self):
+ return self._raw_execute("object_store_get_store_usage_percent", args={})
+
+ def __data(self, **kwds):
+ return kwds
+
+ def _raw_execute(self, command, args={}):
+ return self.pulsar_interface.execute(command, args, data=None, input_path=None, output_path=None)
diff --git a/lib/pulsar/client/path_mapper.py b/lib/pulsar/client/path_mapper.py
new file mode 100644
index 00000000000..02089117b5a
--- /dev/null
+++ b/lib/pulsar/client/path_mapper.py
@@ -0,0 +1,98 @@
+import os.path
+from .action_mapper import FileActionMapper
+from .action_mapper import path_type
+from .util import PathHelper
+
+from galaxy.util import in_directory
+
+
+class PathMapper(object):
+ """ Ties together a FileActionMapper and remote job configuration returned
+ by the Pulsar setup method to pre-determine the location of files for staging
+ on the remote Pulsar server.
+
+ This is not useful when rewrite_paths (as has traditionally been done with
+ the Pulsar) because when doing that the Pulsar determines the paths as files are
+ uploaded. When rewrite_paths is disabled however, the destination of files
+ needs to be determined prior to transfer so an object of this class can be
+ used.
+ """
+
+ def __init__(
+ self,
+ client,
+ remote_job_config,
+ local_working_directory,
+ action_mapper=None,
+ ):
+ self.local_working_directory = local_working_directory
+ if not action_mapper:
+ action_mapper = FileActionMapper(client)
+ self.action_mapper = action_mapper
+ self.input_directory = remote_job_config["inputs_directory"]
+ self.output_directory = remote_job_config["outputs_directory"]
+ self.working_directory = remote_job_config["working_directory"]
+ self.unstructured_files_directory = remote_job_config["unstructured_files_directory"]
+ self.config_directory = remote_job_config["configs_directory"]
+ separator = remote_job_config["system_properties"]["separator"]
+ self.path_helper = PathHelper(separator)
+
+ def remote_output_path_rewrite(self, local_path):
+ output_type = path_type.OUTPUT
+ if in_directory(local_path, self.local_working_directory):
+ output_type = path_type.OUTPUT_WORKDIR
+ remote_path = self.__remote_path_rewrite(local_path, output_type)
+ return remote_path
+
+ def remote_input_path_rewrite(self, local_path):
+ remote_path = self.__remote_path_rewrite(local_path, path_type.INPUT)
+ return remote_path
+
+ def remote_version_path_rewrite(self, local_path):
+ remote_path = self.__remote_path_rewrite(local_path, path_type.OUTPUT, name="COMMAND_VERSION")
+ return remote_path
+
+ def check_for_arbitrary_rewrite(self, local_path):
+ path = str(local_path) # Use false_path if needed.
+ action = self.action_mapper.action(path, path_type.UNSTRUCTURED)
+ if not action.staging_needed:
+ return action.path_rewrite(self.path_helper), []
+ unique_names = action.unstructured_map()
+ name = unique_names[path]
+ remote_path = self.path_helper.remote_join(self.unstructured_files_directory, name)
+ return remote_path, unique_names
+
+ def __remote_path_rewrite(self, dataset_path, dataset_path_type, name=None):
+ """ Return remote path of this file (if staging is required) else None.
+ """
+ path = str(dataset_path) # Use false_path if needed.
+ action = self.action_mapper.action(path, dataset_path_type)
+ if action.staging_needed:
+ if name is None:
+ name = os.path.basename(path)
+ remote_directory = self.__remote_directory(dataset_path_type)
+ remote_path_rewrite = self.path_helper.remote_join(remote_directory, name)
+ else:
+ # Actions which don't require staging MUST define a path_rewrite
+ # method.
+ remote_path_rewrite = action.path_rewrite(self.path_helper)
+
+ return remote_path_rewrite
+
+ def __action(self, dataset_path, dataset_path_type):
+ path = str(dataset_path) # Use false_path if needed.
+ action = self.action_mapper.action(path, dataset_path_type)
+ return action
+
+ def __remote_directory(self, dataset_path_type):
+ if dataset_path_type in [path_type.OUTPUT]:
+ return self.output_directory
+ elif dataset_path_type in [path_type.WORKDIR, path_type.OUTPUT_WORKDIR]:
+ return self.working_directory
+ elif dataset_path_type in [path_type.INPUT]:
+ return self.input_directory
+ else:
+ message = "PathMapper cannot handle path type %s" % dataset_path_type
+ raise Exception(message)
+
+__all__ = [PathMapper]
diff --git a/lib/pulsar/client/setup_handler.py b/lib/pulsar/client/setup_handler.py
new file mode 100644
index 00000000000..b61931eb7d8
--- /dev/null
+++ b/lib/pulsar/client/setup_handler.py
@@ -0,0 +1,103 @@
+import os
+from .util import filter_destination_params
+
+REMOTE_SYSTEM_PROPERTY_PREFIX = "remote_property_"
+
+
+def build(client, destination_args):
+ """ Build a SetupHandler object for client from destination parameters.
+ """
+ # Have defined a remote job directory, lets do the setup locally.
+ if client.job_directory:
+ handler = LocalSetupHandler(client, destination_args)
+ else:
+ handler = RemoteSetupHandler(client)
+ return handler
+
+
+class LocalSetupHandler(object):
+ """ Parse destination params to infer job setup parameters (input/output
+ directories, etc...). Default is to get this configuration data from the
+ remote Pulsar server.
+
+ Downside of this approach is that it requires more and more dependent
+ configuraiton of Galaxy. Upside is that it is asynchronous and thus makes
+ message queue driven configurations possible.
+
+ Remote system properties (such as galaxy_home) can be specified in
+ destination args by prefixing property with remote_property_ (e.g.
+ remote_property_galaxy_home).
+ """
+
+ def __init__(self, client, destination_args):
+ self.client = client
+ system_properties = self.__build_system_properties(destination_args)
+ system_properties["separator"] = client.job_directory.separator
+ self.system_properties = system_properties
+ self.jobs_directory = destination_args["jobs_directory"]
+
+ def setup(self, job_id, tool_id=None, tool_version=None):
+ return build_job_config(
+ job_id=job_id,
+ job_directory=self.client.job_directory,
+ system_properties=self.system_properties,
+ tool_id=tool_id,
+ tool_version=tool_version,
+ )
+
+ @property
+ def local(self):
+ """
+ """
+ return True
+
+ def __build_system_properties(self, destination_params):
+ return filter_destination_params(destination_params, REMOTE_SYSTEM_PROPERTY_PREFIX)
+
+
+class RemoteSetupHandler(object):
+ """ Default behavior. Fetch setup information from remote Pulsar server.
+ """
+ def __init__(self, client):
+ self.client = client
+
+ def setup(self, **setup_args):
+ return self.client.remote_setup(**setup_args)
+
+ @property
+ def local(self):
+ """
+ """
+ return False
+
+
+def build_job_config(job_id, job_directory, system_properties={}, tool_id=None, tool_version=None):
+ """
+ """
+ inputs_directory = job_directory.inputs_directory()
+ working_directory = job_directory.working_directory()
+ outputs_directory = job_directory.outputs_directory()
+ configs_directory = job_directory.configs_directory()
+ tools_directory = job_directory.tool_files_directory()
+ unstructured_files_directory = job_directory.unstructured_files_directory()
+ sep = system_properties.get("sep", os.sep)
+ job_config = {
+ "working_directory": working_directory,
+ "outputs_directory": outputs_directory,
+ "configs_directory": configs_directory,
+ "tools_directory": tools_directory,
+ "inputs_directory": inputs_directory,
+ "unstructured_files_directory": unstructured_files_directory,
+ # Poorly named legacy attribute. Drop at some point.
+ "path_separator": sep,
+ "job_id": job_id,
+ "system_properties": system_properties,
+ }
+ if tool_id:
+ job_config["tool_id"] = tool_id
+ if tool_version:
+ job_config["tool_version"] = tool_version
+ return job_config
+
+
+__all__ = [build_job_config, build]
diff --git a/lib/pulsar/client/staging/__init__.py b/lib/pulsar/client/staging/__init__.py
new file mode 100644
index 00000000000..7b0a47c3c9b
--- /dev/null
+++ b/lib/pulsar/client/staging/__init__.py
@@ -0,0 +1,161 @@
+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 PulsarOutputs(object):
+ """ Abstraction describing the output files PRODUCED by the remote Pulsar
+ server. """
+
+ def __init__(self, working_directory_contents, output_directory_contents, remote_separator=sep):
+ self.working_directory_contents = working_directory_contents
+ self.output_directory_contents = output_directory_contents
+ self.path_helper = PathHelper(remote_separator)
+
+ @staticmethod
+ def from_status_response(complete_response):
+ # Default to None instead of [] to distinguish between empty contents and it not set
+ # by the Pulsar - older Pulsar instances will not set these in complete response.
+ working_directory_contents = complete_response.get("working_directory_contents")
+ output_directory_contents = complete_response.get("outputs_directory_contents")
+ # Older (pre-2014) Pulsar servers will not include separator in response,
+ # so this should only be used when reasoning about outputs in
+ # subdirectories (which was not previously supported prior to that).
+ remote_separator = complete_response.get("system_properties", {}).get("separator", sep)
+ return PulsarOutputs(
+ working_directory_contents,
+ output_directory_contents,
+ remote_separator
+ )
+
+ def has_output_file(self, output_file):
+ return basename(output_file) in self.output_directory_contents
+
+ def output_extras(self, output_file):
+ """
+ Returns dict mapping local path to remote name.
+ """
+ output_directory = dirname(output_file)
+
+ def local_path(name):
+ return join(output_directory, self.path_helper.local_name(name))
+
+ files_directory = "%s_files%s" % (basename(output_file)[0:-len(".dat")], self.path_helper.separator)
+ names = filter(lambda o: o.startswith(files_directory), self.output_directory_contents)
+ return dict(map(lambda name: (local_path(name), name), names))
diff --git a/lib/pulsar/client/staging/down.py b/lib/pulsar/client/staging/down.py
new file mode 100644
index 00000000000..e3adf9d7fa3
--- /dev/null
+++ b/lib/pulsar/client/staging/down.py
@@ -0,0 +1,155 @@
+from os.path import join
+from os.path import relpath
+from re import compile
+from contextlib import contextmanager
+
+from ..staging import COMMAND_VERSION_FILENAME
+from ..action_mapper import FileActionMapper
+
+
+from logging import getLogger
+log = getLogger(__name__)
+
+# All output files marked with from_work_dir attributes will copied or downloaded
+# this pattern picks up attiditional files to copy back - such as those
+# associated with multiple outputs and metadata configuration. Set to .* to just
+# copy everything
+COPY_FROM_WORKING_DIRECTORY_PATTERN = compile(r"primary_.*|galaxy.json|metadata_.*|dataset_\d+\.dat|__instrument_.*|dataset_\d+_files.+")
+
+
+def finish_job(client, cleanup_job, job_completed_normally, client_outputs, pulsar_outputs):
+ """ Responsible for downloading results from remote server and cleaning up
+ Pulsar staging directory (if needed.)
+ """
+ collection_failure_exceptions = []
+ if job_completed_normally:
+ output_collector = ClientOutputCollector(client)
+ action_mapper = FileActionMapper(client)
+ results_stager = ResultsCollector(output_collector, action_mapper, client_outputs, pulsar_outputs)
+ collection_failure_exceptions = results_stager.collect()
+ __clean(collection_failure_exceptions, cleanup_job, client)
+ return collection_failure_exceptions
+
+
+class ClientOutputCollector(object):
+
+ def __init__(self, client):
+ self.client = client
+
+ def collect_output(self, results_collector, output_type, action, name):
+ # This output should have been handled by the Pulsar.
+ if not action.staging_action_local:
+ return False
+
+ working_directory = results_collector.client_outputs.working_directory
+ self.client.fetch_output(
+ path=action.path,
+ name=name,
+ working_directory=working_directory,
+ output_type=output_type,
+ action_type=action.action_type
+ )
+ return True
+
+
+class ResultsCollector(object):
+
+ def __init__(self, output_collector, action_mapper, client_outputs, pulsar_outputs):
+ self.output_collector = output_collector
+ self.action_mapper = action_mapper
+ self.client_outputs = client_outputs
+ self.pulsar_outputs = pulsar_outputs
+ self.downloaded_working_directory_files = []
+ self.exception_tracker = DownloadExceptionTracker()
+ self.output_files = client_outputs.output_files
+ self.working_directory_contents = pulsar_outputs.working_directory_contents or []
+
+ def collect(self):
+ self.__collect_working_directory_outputs()
+ self.__collect_outputs()
+ self.__collect_version_file()
+ self.__collect_other_working_directory_files()
+ return self.exception_tracker.collection_failure_exceptions
+
+ def __collect_working_directory_outputs(self):
+ working_directory = self.client_outputs.working_directory
+ # Fetch explicit working directory outputs.
+ for source_file, output_file in self.client_outputs.work_dir_outputs:
+ name = relpath(source_file, working_directory)
+ pulsar = self.pulsar_outputs.path_helper.remote_name(name)
+ if self._attempt_collect_output('output_workdir', path=output_file, name=pulsar):
+ self.downloaded_working_directory_files.append(pulsar)
+ # Remove from full output_files list so don't try to download directly.
+ try:
+ self.output_files.remove(output_file)
+ except ValueError:
+ raise Exception("Failed to remove %s from %s" % (output_file, self.output_files))
+
+ def __collect_outputs(self):
+ # Legacy Pulsar not returning list of files, iterate over the list of
+ # expected outputs for tool.
+ for output_file in self.output_files:
+ # Fetch output directly...
+ output_generated = self.pulsar_outputs.has_output_file(output_file)
+ if output_generated:
+ self._attempt_collect_output('output', output_file)
+
+ for galaxy_path, pulsar in self.pulsar_outputs.output_extras(output_file).iteritems():
+ self._attempt_collect_output('output', path=galaxy_path, name=pulsar)
+ # else not output generated, do not attempt download.
+
+ def __collect_version_file(self):
+ version_file = self.client_outputs.version_file
+ pulsar_output_directory_contents = self.pulsar_outputs.output_directory_contents
+ if version_file and COMMAND_VERSION_FILENAME in pulsar_output_directory_contents:
+ self._attempt_collect_output('output', version_file, name=COMMAND_VERSION_FILENAME)
+
+ def __collect_other_working_directory_files(self):
+ working_directory = self.client_outputs.working_directory
+ # Fetch remaining working directory outputs of interest.
+ for name in self.working_directory_contents:
+ if name in self.downloaded_working_directory_files:
+ continue
+ if COPY_FROM_WORKING_DIRECTORY_PATTERN.match(name):
+ output_file = join(working_directory, self.pulsar_outputs.path_helper.local_name(name))
+ if self._attempt_collect_output(output_type='output_workdir', path=output_file, name=name):
+ self.downloaded_working_directory_files.append(name)
+
+ def _attempt_collect_output(self, output_type, path, name=None):
+ # path is final path on galaxy server (client)
+ # name is the 'name' of the file on the Pulsar server (possible a relative)
+ # path.
+ collected = False
+ with self.exception_tracker():
+ action = self.action_mapper.action(path, output_type)
+ if self._collect_output(output_type, action, name):
+ collected = True
+
+ return collected
+
+ def _collect_output(self, output_type, action, name):
+ return self.output_collector.collect_output(self, output_type, action, name)
+
+
+class DownloadExceptionTracker(object):
+
+ def __init__(self):
+ self.collection_failure_exceptions = []
+
+ @contextmanager
+ def __call__(self):
+ try:
+ yield
+ except Exception as e:
+ self.collection_failure_exceptions.append(e)
+
+
+def __clean(collection_failure_exceptions, cleanup_job, client):
+ failed = (len(collection_failure_exceptions) > 0)
+ if (not failed and cleanup_job != "never") or cleanup_job == "always":
+ try:
+ client.clean()
+ except Exception:
+ log.warn("Failed to cleanup remote Pulsar job")
+
+__all__ = [finish_job]
diff --git a/lib/pulsar/client/staging/up.py b/lib/pulsar/client/staging/up.py
new file mode 100644
index 00000000000..9712e69eb87
--- /dev/null
+++ b/lib/pulsar/client/staging/up.py
@@ -0,0 +1,420 @@
+from os.path import abspath, basename, join, exists
+from os.path import dirname
+from os.path import relpath
+from os import listdir, sep
+from re import findall
+from io import open
+
+from ..staging import COMMAND_VERSION_FILENAME
+from ..action_mapper import FileActionMapper
+from ..action_mapper import path_type
+from ..action_mapper import MessageAction
+from ..util import PathHelper
+from ..util import directory_files
+
+from logging import getLogger
+log = getLogger(__name__)
+
+
+def submit_job(client, client_job_description, job_config=None):
+ """
+ """
+ file_stager = FileStager(client, client_job_description, job_config)
+ rebuilt_command_line = file_stager.get_command_line()
+ job_id = file_stager.job_id
+ launch_kwds = dict(
+ command_line=rebuilt_command_line,
+ dependencies_description=client_job_description.dependencies_description,
+ env=client_job_description.env,
+ )
+ if file_stager.job_config:
+ launch_kwds["job_config"] = file_stager.job_config
+ remote_staging = {}
+ remote_staging_actions = file_stager.transfer_tracker.remote_staging_actions
+ if remote_staging_actions:
+ remote_staging["setup"] = remote_staging_actions
+ # Somehow make the following optional.
+ remote_staging["action_mapper"] = file_stager.action_mapper.to_dict()
+ remote_staging["client_outputs"] = client_job_description.client_outputs.to_dict()
+
+ if remote_staging:
+ launch_kwds["remote_staging"] = remote_staging
+
+ client.launch(**launch_kwds)
+ return job_id
+
+
+class FileStager(object):
+ """
+ Objects of the FileStager class interact with an Pulsar client object to
+ stage the files required to run jobs on a remote Pulsar server.
+
+ **Parameters**
+
+ client : JobClient
+ Pulsar client object.
+ client_job_description : client_job_description
+ Description of client view of job to stage and execute remotely.
+ """
+
+ def __init__(self, client, client_job_description, job_config):
+ """
+ """
+ self.client = client
+ self.command_line = client_job_description.command_line
+ self.config_files = client_job_description.config_files
+ self.input_files = client_job_description.input_files
+ self.output_files = client_job_description.output_files
+ 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 Pulsar server.
+ self.job_inputs = JobInputs(self.command_line, self.config_files)
+
+ self.action_mapper = FileActionMapper(client)
+
+ self.__handle_setup(job_config)
+
+ self.transfer_tracker = TransferTracker(client, self.path_helper, self.action_mapper, self.job_inputs, rewrite_paths=self.rewrite_paths)
+
+ self.__initialize_referenced_tool_files()
+ if self.rewrite_paths:
+ self.__initialize_referenced_arbitrary_files()
+
+ self.__upload_tool_files()
+ self.__upload_input_files()
+ self.__upload_working_directory_files()
+ self.__upload_arbitrary_files()
+
+ if self.rewrite_paths:
+ self.__initialize_output_file_renames()
+ self.__initialize_task_output_file_renames()
+ self.__initialize_config_file_renames()
+ self.__initialize_version_file_rename()
+
+ self.__handle_rewrites()
+
+ self.__upload_rewritten_config_files()
+
+ def __handle_setup(self, job_config):
+ if not job_config:
+ job_config = self.client.setup(self.tool_id, self.tool_version)
+
+ self.new_working_directory = job_config['working_directory']
+ self.new_outputs_directory = job_config['outputs_directory']
+ # Default configs_directory to match remote working_directory to mimic
+ # behavior of older Pulsar servers.
+ self.new_configs_directory = job_config.get('configs_directory', self.new_working_directory)
+ self.remote_separator = self.__parse_remote_separator(job_config)
+ self.path_helper = PathHelper(self.remote_separator)
+ # If remote Pulsar server assigned job id, use that otherwise
+ # just use local job_id assigned.
+ galaxy_job_id = self.client.job_id
+ self.job_id = job_config.get('job_id', galaxy_job_id)
+ if self.job_id != galaxy_job_id:
+ # Remote Pulsar server assigned an id different than the
+ # Galaxy job id, update client to reflect this.
+ self.client.job_id = self.job_id
+ self.job_config = job_config
+
+ def __parse_remote_separator(self, job_config):
+ separator = job_config.get("system_properties", {}).get("separator", None)
+ if not separator: # Legacy Pulsar
+ separator = job_config["path_separator"] # Poorly named
+ return separator
+
+ def __initialize_referenced_tool_files(self):
+ self.referenced_tool_files = self.job_inputs.find_referenced_subfiles(self.tool_dir)
+
+ def __initialize_referenced_arbitrary_files(self):
+ referenced_arbitrary_path_mappers = dict()
+ for mapper in self.action_mapper.unstructured_mappers():
+ mapper_pattern = mapper.to_pattern()
+ # TODO: Make more sophisticated, allow parent directories,
+ # grabbing sibbling files based on patterns, etc...
+ paths = self.job_inputs.find_pattern_references(mapper_pattern)
+ for path in paths:
+ if path not in referenced_arbitrary_path_mappers:
+ referenced_arbitrary_path_mappers[path] = mapper
+ for path, mapper in referenced_arbitrary_path_mappers.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 = "Pulsar: __upload_input_file called on empty or missing dataset." + \
+ " So such file: [%s]" % input_file
+ log.debug(message)
+
+ def __upload_input_extra_files(self, input_file):
+ files_path = "%s_files" % input_file[0:-len(".dat")]
+ if exists(files_path) and self.__stage_input(files_path):
+ for extra_file_name in directory_files(files_path):
+ extra_file_path = join(files_path, extra_file_name)
+ remote_name = self.path_helper.remote_name(relpath(extra_file_path, dirname(files_path)))
+ self.transfer_tracker.handle_transfer(extra_file_path, path_type.INPUT, name=remote_name)
+
+ def __upload_working_directory_files(self):
+ # Task manager stages files into working directory, these need to be
+ # uploaded if present.
+ working_directory_files = 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/pulsar/client/transport/__init__.py b/lib/pulsar/client/transport/__init__.py
new file mode 100644
index 00000000000..1f1925e3c1c
--- /dev/null
+++ b/lib/pulsar/client/transport/__init__.py
@@ -0,0 +1,31 @@
+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('PULSAR_CURL_TRANSPORT', "0")
+ # If PULSAR_CURL_TRANSPORT is unset or set to 0, use default,
+ # else use curl.
+ if use_curl.isdigit() and not int(use_curl):
+ transport_type = 'urllib'
+ else:
+ transport_type = 'curl'
+ return transport_type
+
+# 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/pulsar/client/transport/curl.py b/lib/pulsar/client/transport/curl.py
new file mode 100644
index 00000000000..9594c719541
--- /dev/null
+++ b/lib/pulsar/client/transport/curl.py
@@ -0,0 +1,71 @@
+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 Pulsar client but pycurl is unavailable."
+
+
+class PycurlTransport(object):
+
+ def execute(self, url, method=None, data=None, input_path=None, output_path=None):
+ buf = _open_output(output_path)
+ try:
+ c = _new_curl_object()
+ c.setopt(c.URL, url.encode('ascii'))
+ c.setopt(c.WRITEFUNCTION, buf.write)
+ if method:
+ c.setopt(c.CUSTOMREQUEST, method)
+ if input_path:
+ c.setopt(c.UPLOAD, 1)
+ c.setopt(c.READFUNCTION, open(input_path, 'rb').read)
+ filesize = 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/pulsar/client/transport/standard.py b/lib/pulsar/client/transport/standard.py
new file mode 100644
index 00000000000..999ecce5af9
--- /dev/null
+++ b/lib/pulsar/client/transport/standard.py
@@ -0,0 +1,51 @@
+"""
+Pulsar HTTP Client layer based on Python Standard Library (urllib2)
+"""
+from __future__ import with_statement
+from os.path import getsize
+import mmap
+try:
+ from urllib2 import urlopen
+except ImportError:
+ from urllib.request import urlopen
+try:
+ from urllib2 import Request
+except ImportError:
+ from urllib.request import Request
+
+
+class Urllib2Transport(object):
+
+ def _url_open(self, request, data):
+ return urlopen(request, data)
+
+ def execute(self, url, method=None, data=None, input_path=None, output_path=None):
+ request = self.__request(url, data, method)
+ input = None
+ try:
+ if input_path:
+ 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()
+
+ def __request(self, url, data, method):
+ request = Request(url=url, data=data)
+ if method:
+ request.get_method = lambda: method
+ return request
diff --git a/lib/pulsar/client/util.py b/lib/pulsar/client/util.py
new file mode 100644
index 00000000000..04908db6256
--- /dev/null
+++ b/lib/pulsar/client/util.py
@@ -0,0 +1,177 @@
+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 Pulsar designed to be used with
+ >>> # posix client.
+ >>> posix_path_helper = PathHelper("/", local_path_module=posixpath)
+ >>> windows_slash = "\\\\"
+ >>> len(windows_slash)
+ 1
+ >>> nt_path_helper = PathHelper(windows_slash, local_path_module=posixpath)
+ >>> posix_path_helper.remote_name("moo/cow")
+ 'moo/cow'
+ >>> nt_path_helper.remote_name("moo/cow")
+ 'moo\\\\cow'
+ >>> posix_path_helper.local_name("moo/cow")
+ 'moo/cow'
+ >>> nt_path_helper.local_name("moo\\\\cow")
+ 'moo/cow'
+ >>> posix_path_helper.from_posix_with_new_base("/galaxy/data/bowtie/hg19.fa", "/galaxy/data/", "/work/galaxy/data")
+ '/work/galaxy/data/bowtie/hg19.fa'
+ >>> posix_path_helper.from_posix_with_new_base("/galaxy/data/bowtie/hg19.fa", "/galaxy/data", "/work/galaxy/data")
+ '/work/galaxy/data/bowtie/hg19.fa'
+ >>> posix_path_helper.from_posix_with_new_base("/galaxy/data/bowtie/hg19.fa", "/galaxy/data", "/work/galaxy/data/")
+ '/work/galaxy/data/bowtie/hg19.fa'
+ '''
+
+ def __init__(self, separator, local_path_module=os.path):
+ self.separator = separator
+ self.local_join = local_path_module.join
+ self.local_sep = local_path_module.sep
+
+ def remote_name(self, local_name):
+ return self.remote_join(*local_name.split(self.local_sep))
+
+ def local_name(self, remote_name):
+ return self.local_join(*remote_name.split(self.separator))
+
+ def remote_join(self, *args):
+ return self.separator.join(args)
+
+ def from_posix_with_new_base(self, posix_path, old_base, new_base):
+ # TODO: Test with new_base as a windows path against nt_path_helper.
+ if old_base.endswith("/"):
+ old_base = old_base[:-1]
+ if not posix_path.startswith(old_base):
+ message_template = "Cannot compute new path for file %s, does not start with %s."
+ message = message_template % (posix_path, old_base)
+ raise Exception(message)
+ stripped_path = posix_path[len(old_base):]
+ while stripped_path.startswith("/"):
+ stripped_path = stripped_path[1:]
+ path_parts = stripped_path.split(self.separator)
+ if new_base.endswith(self.separator):
+ new_base = new_base[:-len(self.separator)]
+ return self.remote_join(new_base, *path_parts)
+
+
+class TransferEventManager(object):
+
+ def __init__(self):
+ self.events = WeakValueDictionary(dict())
+ self.events_lock = Lock()
+
+ def acquire_event(self, path, force_clear=False):
+ with self.events_lock:
+ if path in self.events:
+ event_holder = self.events[path]
+ else:
+ event_holder = EventHolder(Event(), path, self)
+ self.events[path] = event_holder
+ if force_clear:
+ event_holder.event.clear()
+ return event_holder
+
+
+class EventHolder(object):
+
+ def __init__(self, event, path, condition_manager):
+ self.event = event
+ self.path = path
+ self.condition_manager = condition_manager
+ self.failed = False
+
+ def release(self):
+ self.event.set()
+
+ def fail(self):
+ self.failed = True