Merge pull request #7006 from nuwang/chained_dynamic_destinations

Added support for chaining dynamic destinations
This commit is contained in:
Marius van den Beek
2018-11-16 16:06:58 +01:00
committed by GitHub
13 changed files with 201 additions and 7 deletions
+24
View File
@@ -652,6 +652,30 @@
<param id="num_jobs">50</param>
<!-- <param id="job_states">queued</param> -->
</destination>
<destination id="burst_if_queued" runner="dynamic">
<!-- Dynamic destinations can be chained together to create more
complex rules. In this example, the built-in burst stock rule
determines whether to burst, and if so, directs to burst_if_size,
a user-defined dynamic destination. This destination in turn will
conditionally route it to a remote pulsar node if the input size
is below a certain threshold, or route to local if not. -->
<param id="type">burst</param>
<param id="from_destination_ids">local</param>
<param id="to_destination_id">burst_if_size</param>
<param id="num_jobs">2</param>
<param id="job_states">queued</param>
</destination>
<destination id="burst_if_size" runner="dynamic">
<param id="type">python</param>
<param id="function">to_destination_if_size</param>
<!-- Also demonstrates a destination level override of the
rules_module. This rules_module will take precedence over the
plugin level rules module when resolving the dynamic function -->
<param id="rules_module">galaxycloudrunner.rules</param>
<param id="max_size">1g</param>
<param id="to_destination_id">pulsar</param>
<param id="fallback_destination_id">local</param>
</destination>
<destination id="docker_dispatch" runner="dynamic">
<!-- Follow dynamic destination type will send all tool's that
support docker to static destination defined by
+22 -7
View File
@@ -205,24 +205,39 @@ class JobRunnerMapper(object):
job_destination = self.job_config.get_destination(job_destination_rep)
return job_destination
def __cache_job_destination(self, params, raw_job_destination=None):
def __determine_job_destination(self, params, raw_job_destination=None):
if raw_job_destination is None:
raw_job_destination = self.job_wrapper.tool.get_job_destination(params)
if raw_job_destination.runner == DYNAMIC_RUNNER_NAME:
job_destination = self.__handle_dynamic_job_destination(raw_job_destination)
log.debug("(%s) Mapped job to destination id: %s", self.job_wrapper.job_id, job_destination.id)
# Recursively handle chained dynamic destinations
if job_destination.runner == DYNAMIC_RUNNER_NAME:
return self.__determine_job_destination(params, raw_job_destination=job_destination)
else:
job_destination = raw_job_destination
log.debug("(%s) Mapped job to destination id: %s", self.job_wrapper.job_id, job_destination.id)
self.cached_job_destination = job_destination
log.debug("(%s) Mapped job to destination id: %s", self.job_wrapper.job_id, job_destination.id)
return job_destination
def __cache_job_destination(self, params, raw_job_destination=None):
self.cached_job_destination = self.__determine_job_destination(
params, raw_job_destination=raw_job_destination)
return self.cached_job_destination
def get_job_destination(self, params):
"""
Cache the job_destination to avoid recalculation.
cached_job_destination is a public property that is sometimes
externally set to short-circuit the mapper, such as during resubmits.
get_job_destination will respect that and not run the mapper if so.
"""
if not hasattr(self, 'cached_job_destination'):
self.__cache_job_destination(params)
return self.__cache_job_destination(params)
return self.cached_job_destination
def cache_job_destination(self, raw_job_destination):
self.__cache_job_destination(None, raw_job_destination=raw_job_destination)
return self.cached_job_destination
"""
Force update of cached_job_destination to mapper determined job
destination, overwriting any externally set cached_job_destination
"""
return self.__cache_job_destination(
None, raw_job_destination=raw_job_destination)
@@ -0,0 +1,28 @@
<?xml version="1.0"?>
<job_conf>
<plugins>
<plugin id="local" type="runner" load="galaxy.jobs.runners.local:LocalJobRunner" workers="2"/>
<plugin id="dynamic" type="runner">
<param id="rules_module">integration.chained_dyndest_rules.module1</param>
</plugin>
</plugins>
<handlers>
<handler id="main"/>
</handlers>
<destinations default="dyn_dest1">
<destination id="dyn_dest1" runner="dynamic">
<param id="type">python</param>
<param id="function">dyndest_chain_1</param>
</destination>
<destination id="dyn_dest2" runner="dynamic">
<param id="type">python</param>
<param id="function">dyndest_chain_2</param>
<!-- test overriding rules_module at destination level -->
<param id="rules_module">integration.chained_dyndest_rules.module2</param>
<param id="tmp_dir_prefix">from1</param>
</destination>
</destinations>
</job_conf>
@@ -0,0 +1,15 @@
def dyndest_chain_1():
# Check whether chaining dynamic job destinations work
return "dyn_dest2"
def dyndest_chain_2():
# Return an invalid destination as this module's function
# should never be called
return "invalid_destination"
def dyndest_chain_3():
# Return an invalid destination as this module's function
# should never be called
return "invalid_destination"
@@ -0,0 +1,23 @@
from galaxy.jobs import JobDestination
def dyndest_chain_1():
# Return an invalid destination as this module's function
# should never be called
return "invalid_destination"
def dyndest_chain_2(tmp_dir_prefix):
# Chain to yet a third
return JobDestination(
runner="dynamic",
params={'type': 'python',
'function': 'dyndest_chain_3',
'rules_module': 'integration.chained_dyndest_rules.module3',
'tmp_dir_prefix_two': '%sand2' % tmp_dir_prefix})
def dyndest_chain_3():
# Return an invalid destination as this module's function
# should never be called
return "invalid_destination"
@@ -0,0 +1,19 @@
from galaxy.jobs import JobDestination
def dyndest_chain_1():
# Return an invalid destination as this module's function
# should never be called
return "invalid_destination"
def dyndest_chain_2():
# Return an invalid destination as this module's function
# should never be called
return "invalid_destination"
def dyndest_chain_3(tmp_dir_prefix_two):
tmp_dir = '$(mktemp %sand3XXXXXXXXXXXX)' % tmp_dir_prefix_two
return JobDestination(runner="local",
params={'tmp_dir': tmp_dir})
@@ -0,0 +1,32 @@
"""Integration tests for chained dynamic job destinations."""
import os
import tempfile
from base.populators import (
skip_without_tool,
)
from .test_job_environments import BaseJobEnvironmentIntegrationTestCase
SCRIPT_DIRECTORY = os.path.abspath(os.path.dirname(__file__))
CHAINED_DYNDESTS_JOB_CONFIG = os.path.join(SCRIPT_DIRECTORY, "chained_dyndest_job_conf.xml")
class ChainedDynamicDestinationIntegrationTestCase(BaseJobEnvironmentIntegrationTestCase):
@classmethod
def handle_galaxy_config_kwds(cls, config):
cls.jobs_directory = tempfile.mkdtemp()
config["jobs_directory"] = cls.jobs_directory
config["job_config_file"] = CHAINED_DYNDESTS_JOB_CONFIG
@skip_without_tool("job_environment_default")
def test_default_environment_1801(self):
job_env = self._run_and_get_environment_properties()
# Since dynamic destinations compute final tmp_dir parameter to be
# $(mktemp from1and2and3XXXXXXXXXXXX), tmpdir should start
# with from1and2and3.
basename = os.path.basename(job_env.tmp)
assert basename.startswith("from1and2and3"), job_env.tmp
+23
View File
@@ -36,6 +36,12 @@ def test_dynamic_mapping():
assert mapper.job_config.rule_response == "local_runner"
def test_chained_dynamic_mapping():
mapper = __mapper(__dynamic_destination(dict(function="dynamic_chain_1")))
assert mapper.get_job_destination({}) is DYNAMICALLY_GENERATED_DESTINATION
assert mapper.job_config.rule_response == "final_destination"
def test_dynamic_mapping_priorities():
mapper = __mapper(__dynamic_destination(dict(function="tophat")))
assert mapper.get_job_destination({}) is DYNAMICALLY_GENERATED_DESTINATION
@@ -97,6 +103,23 @@ def test_dynamic_mapping_rule_module_override():
assert mapper.job_config.rule_response == "new_rules_package"
def test_dynamic_mapping_externally_set_job_destination():
mapper = __mapper(__dynamic_destination(dict(function="upload")))
# Initially, the mapper should not have a cached destination
assert not hasattr(mapper, 'cached_job_destination')
# Overwrite with an externally set job destination
manually_set_destination = JobDestination(runner="dynamic")
mapper.cached_job_destination = manually_set_destination
destination = mapper.get_job_destination({})
assert destination == manually_set_destination
assert mapper.cached_job_destination == manually_set_destination
# Force overwrite with mapper determined destination
mapper.cache_job_destination(None)
assert mapper.cached_job_destination is not None
assert mapper.cached_job_destination != manually_set_destination
assert mapper.job_config.rule_response == "local_runner"
def __assert_mapper_errors_with_message(mapper, message):
exception = None
try:
+15
View File
@@ -1,3 +1,4 @@
from galaxy.jobs import JobDestination
def upload():
@@ -15,6 +16,20 @@ def tool1():
return 'tool1_dest_id'
def dynamic_chain_1():
# Check whether chaining dynamic job destinations work
return JobDestination(runner="dynamic",
params={'type': 'python',
'function': 'dynamic_chain_2',
'test_param': 'my_test_param'})
def dynamic_chain_2(test_param):
# Check whether chaining dynamic job destinations work
assert test_param == "my_test_param"
return "final_destination"
def check_rule_params(
job_id,
tool,