Merged in jmchilton/galaxy-central-fork-1 (pull request #422)

Implement Pulsar job runners.
This commit is contained in:
John Chilton
2014-06-23 21:27:51 -05:00
31 changed files with 3997 additions and 105 deletions
+68 -73
View File
@@ -19,27 +19,30 @@
<!-- Override the $DRMAA_LIBRARY_PATH environment variable -->
<param id="drmaa_library_path">/sge/lib/libdrmaa.so</param>
</plugin>
<plugin id="lwr" type="runner" load="galaxy.jobs.runners.lwr:LwrJobRunner">
<!-- More information on LWR can be found at https://lwr.readthedocs.org -->
<!-- Uncomment following line to use libcurl to perform HTTP calls (defaults to urllib) -->
<plugin id="cli" type="runner" load="galaxy.jobs.runners.cli:ShellJobRunner" />
<plugin id="condor" type="runner" load="galaxy.jobs.runners.condor:CondorJobRunner" />
<plugin id="slurm" type="runner" load="galaxy.jobs.runners.slurm:SlurmJobRunner" />
<!-- Pulsar runners (see more at https://pulsar.readthedocs.org) -->
<plugin id="pulsar_rest" type="runner" load="galaxy.jobs.runners.pulsar:PulsarRESTJobRunner">
<!-- Allow optimized HTTP calls with libcurl (defaults to urllib) -->
<!-- <param id="transport">curl</param> -->
<!-- *Experimental Caching*: Uncomment next parameters to enable
caching and specify the number of caching threads to enable on Galaxy
side. Likely will not work with newer features such as MQ support.
If this is enabled be sure to specify a `file_cache_dir` in the remote
LWR's main configuration file.
<!-- *Experimental Caching*: Next parameter enables caching.
Likely will not work with newer features such as MQ support.
If this is enabled be sure to specify a `file_cache_dir` in
the remote Pulsar's servers main configuration file.
-->
<!-- <param id="cache">True</param> -->
<!-- <param id="transfer_threads">2</param> -->
</plugin>
<plugin id="amqp_lwr" type="runner" load="galaxy.jobs.runners.lwr:LwrJobRunner">
<param id="url">amqp://guest:guest@localhost:5672//</param>
<!-- If using message queue driven LWR - the LWR will generally
initiate file transfers so a the URL of this Galaxy instance
must be configured. -->
<plugin id="pulsar_mq" type="runner" load="galaxy.jobs.runners.pulsar:PulsarMQJobRunner">
<!-- AMQP URL to connect to. -->
<param id="amqp_url">amqp://guest:guest@localhost:5672//</param>
<!-- URL remote Pulsar apps should transfer files to this Galaxy
instance to/from. -->
<param id="galaxy_url">http://localhost:8080</param>
<!-- If multiple managers configured on the LWR, specify which one
this plugin targets. -->
<!-- Pulsar job manager to communicate with (see Pulsar
docs for information on job managers). -->
<!-- <param id="manager">_default_</param> -->
<!-- The AMQP client can provide an SSL client certificate (e.g. for
validation), the following options configure that certificate
@@ -58,9 +61,17 @@
higher value (in seconds) (or `None` to use blocking connections). -->
<!-- <param id="amqp_consumer_timeout">None</param> -->
</plugin>
<plugin id="cli" type="runner" load="galaxy.jobs.runners.cli:ShellJobRunner" />
<plugin id="condor" type="runner" load="galaxy.jobs.runners.condor:CondorJobRunner" />
<plugin id="slurm" type="runner" load="galaxy.jobs.runners.slurm:SlurmJobRunner" />
<plugin id="pulsar_legacy" type="runner" load="galaxy.jobs.runners.pulsar:PulsarLegacyJobRunner">
<!-- Pulsar job runner with default parameters matching those
of old LWR job runner. If your Pulsar server is running on a
Windows machine for instance this runner should still be used.
These destinations still needs to target a Pulsar server,
older LWR plugins and destinations still work in Galaxy can
target LWR servers, but this support should be considered
deprecated and will disappear with a future release of Galaxy.
-->
</plugin>
</plugins>
<handlers default="handlers">
<!-- Additional job handlers - the id should match the name of a
@@ -125,8 +136,8 @@
$galaxy_root:ro,$tool_directory:ro,$working_directory:rw,$default_file_path:ro
If using the LWR, defaults will be even further restricted because the
LWR will (by default) stage all needed inputs into the job's job_directory
If using the Pulsar, defaults will be even further restricted because the
Pulsar will (by default) stage all needed inputs into the job's job_directory
(so there is not need to allow the docker container to read all the
files - let alone write over them). Defaults in this case becomes:
@@ -135,7 +146,7 @@
Python string.Template is used to expand volumes and values $defaults,
$galaxy_root, $default_file_path, $tool_directory, $working_directory,
are available to all jobs and $job_directory is also available for
LWR jobs.
Pulsar jobs.
-->
<!-- Control memory allocatable by docker container with following option:
-->
@@ -213,87 +224,71 @@
<!-- A destination that represents a method in the dynamic runner. -->
<param id="function">foo</param>
</destination>
<destination id="secure_lwr" runner="lwr">
<param id="url">https://windowshost.examle.com:8913/</param>
<!-- If set, private_token must match token remote LWR server configured with. -->
<destination id="secure_pulsar_rest_dest" runner="pulsar_rest">
<param id="url">https://examle.com:8913/</param>
<!-- If set, private_token must match token in remote Pulsar's
configuration. -->
<param id="private_token">123456789changeme</param>
<!-- Uncomment the following statement to disable file staging (e.g.
if there is a shared file system between Galaxy and the LWR
if there is a shared file system between Galaxy and the Pulsar
server). Alternatively action can be set to 'copy' - to replace
http transfers with file system copies, 'remote_transfer' to cause
the lwr to initiate HTTP transfers instead of Galaxy, or
'remote_copy' to cause lwr to initiate file system copies.
the Pulsar to initiate HTTP transfers instead of Galaxy, or
'remote_copy' to cause Pulsar to initiate file system copies.
If setting this to 'remote_transfer' be sure to specify a
'galaxy_url' attribute on the runner plugin above. -->
<!-- <param id="default_file_action">none</param> -->
<!-- The above option is just the default, the transfer behavior
none|copy|http can be configured on a per path basis via the
following file. See lib/galaxy/jobs/runners/lwr_client/action_mapper.py
for examples of how to configure this file. This is very beta
and nature of file will likely change.
following file. See Pulsar documentation for more details and
examples.
-->
<!-- <param id="file_action_config">file_actions.json</param> -->
<!-- Uncomment following option to disable Galaxy tool dependency
resolution and utilize remote LWR's configuraiton of tool
dependency resolution instead (same options as Galaxy for
dependency resolution are available in LWR). At a minimum
the remote LWR server should define a tool_dependencies_dir in
its `server.ini` configuration. The LWR will not attempt to
stage dependencies - so ensure the the required galaxy or tool
shed packages are available remotely (exact same tool shed
installed changesets are required).
<!-- <param id="file_action_config">file_actions.yaml</param> -->
<!-- The non-legacy Pulsar runners will attempt to resolve Galaxy
dependencies remotely - to enable this set a tool_dependency_dir
in Pulsar's configuration (can work with all the same dependency
resolutions mechanisms as Galaxy - tool Shed installs, Galaxy
packages, etc...). To disable this behavior, set the follow parameter
to none. To generate the dependency resolution command locally
set the following parameter local.
-->
<!-- <param id="dependency_resolution">remote</params> -->
<!-- Traditionally, the LWR allow Galaxy to generate a command line
as if it were going to run the command locally and then the
LWR client rewrites it after the fact using regular
expressions. Setting the following value to true causes the
LWR runner to insert itself into the command line generation
process and generate the correct command line from the get go.
This will likely be the default someday - but requires a newer
LWR version and is less well tested. -->
<!-- <param id="rewrite_parameters">true</params> -->
<!-- <param id="dependency_resolution">none</params> -->
<!-- Uncomment following option to enable setting metadata on remote
LWR server. The 'use_remote_datatypes' option is available for
Pulsar server. The 'use_remote_datatypes' option is available for
determining whether to use remotely configured datatypes or local
ones (both alternatives are a little brittle). -->
<!-- <param id="remote_metadata">true</param> -->
<!-- <param id="use_remote_datatypes">false</param> -->
<!-- <param id="remote_property_galaxy_home">/path/to/remote/galaxy-central</param> -->
<!-- If remote LWR server is configured to run jobs as the real user,
<!-- If remote Pulsar server is configured to run jobs as the real user,
uncomment the following line to pass the current Galaxy user
along. -->
<!-- <param id="submit_user">$__user_name__</param> -->
<!-- Various other submission parameters can be passed along to the LWR
whose use will depend on the remote LWR's configured job manager.
<!-- Various other submission parameters can be passed along to the Pulsar
whose use will depend on the remote Pulsar's configured job manager.
For instance:
-->
<!-- <param id="submit_native_specification">-P bignodes -R y -pe threads 8</param> -->
</destination>
<destination id="amqp_lwr_dest" runner="amqp_lwr" >
<!-- url and private_token are not valid when using MQ driven LWR. The plugin above
determines which queue/manager to target and the underlying MQ server should be
used to configure security.
<!-- <param id="submit_native_specification">-P bignodes -R y -pe threads 8</param> -->
<!-- Disable parameter rewriting and rewrite generated commands
instead. This may be required if remote host is Windows machine
but probably not otherwise.
-->
<!-- Traditionally, the LWR client sends request to LWR
server to populate various system properties. This
<!-- <param id="rewrite_parameters">false</params> -->
</destination>
<destination id="pulsar_mq_dest" runner="amqp_pulsar" >
<!-- The RESTful Pulsar client sends a request to Pulsar
to populate various system properties. This
extra step can be disabled and these calculated here
on client by uncommenting jobs_directory and
specifying any additional remote_property_ of
interest, this is not optional when using message
queues.
-->
<param id="jobs_directory">/path/to/remote/lwr/lwr_staging/</param>
<!-- Default the LWR send files to and pull files from Galaxy when
using message queues (in the more traditional mode Galaxy sends
files to and pull files from the LWR - this is obviously less
appropriate when using a message queue).
The default_file_action currently requires pycurl be available
to Galaxy (presumably in its virtualenv). Making this dependency
optional is an open task.
<param id="jobs_directory">/path/to/remote/pulsar/files/staging/</param>
<!-- Otherwise MQ and Legacy pulsar destinations can be supplied
all the same destination parameters as the RESTful client documented
above (though url and private_token are ignored when using a MQ).
-->
<param id="default_file_action">remote_transfer</param>
</destination>
<destination id="ssh_torque" runner="cli">
<param id="shell_plugin">SecureShell</param>
+707
View File
@@ -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
+1 -1
View File
@@ -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
+1 -1
View File
@@ -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
)
@@ -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')
@@ -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')
+3 -3
View File
@@ -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))
@@ -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
+2 -2
View File
@@ -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(
View File
+62
View File
@@ -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 <https://bitbucket.org/galaxy/galaxy-dist/src/tip/job_conf.xml.sample_advanced?at=default>`_
for information on how to configure Galaxy to interact with the 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::
<tool_id> = pulsar://http://<pulsar_host>:<pulsar_port>
As an example, if a host named remotehost is running the Pulsar server
application on port ``8913``, then the tool with id ``test_tool`` can
be configured to run remotely on remotehost by adding the following
line to ``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,
]
+567
View File
@@ -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
]
+139
View File
@@ -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
@@ -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
+381
View File
@@ -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
)
+80
View File
@@ -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,
}
+35
View File
@@ -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
+58
View File
@@ -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)
+152
View File
@@ -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
+135
View File
@@ -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)
+223
View File
@@ -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]
+53
View File
@@ -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)
+98
View File
@@ -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]
+103
View File
@@ -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]
+161
View File
@@ -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))
+155
View File
@@ -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]
+420
View File
@@ -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]
+31
View File
@@ -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]
+71
View File
@@ -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]
+51
View File
@@ -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
+177
View File
@@ -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