From b8c97f108681e0e0963efc95a16af13be84642f9 Mon Sep 17 00:00:00 2001 From: Nuwan Goonasekera Date: Mon, 12 Nov 2018 23:08:16 +0530 Subject: [PATCH 1/8] Added support for chaining dynamic destinations --- lib/galaxy/jobs/mapper.py | 21 ++++++++++++++------- test/unit/jobs/test_mapper.py | 6 ++++++ test/unit/jobs/test_rules/10_site.py | 15 +++++++++++++++ 3 files changed, 35 insertions(+), 7 deletions(-) diff --git a/lib/galaxy/jobs/mapper.py b/lib/galaxy/jobs/mapper.py index 84e0d068258..4351e4d5d64 100644 --- a/lib/galaxy/jobs/mapper.py +++ b/lib/galaxy/jobs/mapper.py @@ -205,24 +205,31 @@ 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) + # Recursively handle chained dynamic destinations + job_destination = 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 + return job_destination + + def __cache_job_destination(self, params, raw_job_destination=None): + if not hasattr(self, 'cached_job_destination'): + 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. """ - if not hasattr(self, 'cached_job_destination'): - self.__cache_job_destination(params) - return self.cached_job_destination + return self.__cache_job_destination(params) def cache_job_destination(self, raw_job_destination): - self.__cache_job_destination(None, raw_job_destination=raw_job_destination) - return self.cached_job_destination + return self.__cache_job_destination( + None, raw_job_destination=raw_job_destination) diff --git a/test/unit/jobs/test_mapper.py b/test/unit/jobs/test_mapper.py index 30e02a77dc4..cac783a475d 100644 --- a/test/unit/jobs/test_mapper.py +++ b/test/unit/jobs/test_mapper.py @@ -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 diff --git a/test/unit/jobs/test_rules/10_site.py b/test/unit/jobs/test_rules/10_site.py index a3578dac359..0712df11d87 100644 --- a/test/unit/jobs/test_rules/10_site.py +++ b/test/unit/jobs/test_rules/10_site.py @@ -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, From f46ab5a5dc21c1e8dd115ec3d9a37e0e3ee26a56 Mon Sep 17 00:00:00 2001 From: Nuwan Goonasekera Date: Tue, 13 Nov 2018 23:56:39 +0530 Subject: [PATCH 2/8] Fix debug logging --- lib/galaxy/jobs/mapper.py | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/lib/galaxy/jobs/mapper.py b/lib/galaxy/jobs/mapper.py index 4351e4d5d64..c9ab9a93b53 100644 --- a/lib/galaxy/jobs/mapper.py +++ b/lib/galaxy/jobs/mapper.py @@ -210,12 +210,13 @@ class JobRunnerMapper(object): 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 - job_destination = self.__determine_job_destination( - params, raw_job_destination=job_destination) + 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) + 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): From f6c0e6ff9fdbf66a2af6d6538ad816274ecb93b3 Mon Sep 17 00:00:00 2001 From: Nuwan Goonasekera Date: Thu, 15 Nov 2018 13:08:04 +0530 Subject: [PATCH 3/8] Fixed issue with job_destination caching and added documentation --- lib/galaxy/jobs/mapper.py | 15 +++++++++++---- 1 file changed, 11 insertions(+), 4 deletions(-) diff --git a/lib/galaxy/jobs/mapper.py b/lib/galaxy/jobs/mapper.py index c9ab9a93b53..bb0236dad14 100644 --- a/lib/galaxy/jobs/mapper.py +++ b/lib/galaxy/jobs/mapper.py @@ -220,17 +220,24 @@ class JobRunnerMapper(object): return job_destination def __cache_job_destination(self, params, raw_job_destination=None): - if not hasattr(self, 'cached_job_destination'): - self.cached_job_destination = self.__determine_job_destination( + 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. """ - return self.__cache_job_destination(params) + if not hasattr(self, 'cached_job_destination'): + return self.__cache_job_destination(params) + return self.cached_job_destination def cache_job_destination(self, raw_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) From c114ba4f43cb4e73fbff470f1283aed103e91fc2 Mon Sep 17 00:00:00 2001 From: Nuwan Goonasekera Date: Thu, 15 Nov 2018 13:27:49 +0530 Subject: [PATCH 4/8] Added test for cached_job_destination to prevent regression --- lib/galaxy/jobs/mapper.py | 2 +- test/unit/jobs/test_mapper.py | 17 +++++++++++++++++ 2 files changed, 18 insertions(+), 1 deletion(-) diff --git a/lib/galaxy/jobs/mapper.py b/lib/galaxy/jobs/mapper.py index bb0236dad14..563309cd483 100644 --- a/lib/galaxy/jobs/mapper.py +++ b/lib/galaxy/jobs/mapper.py @@ -221,7 +221,7 @@ class JobRunnerMapper(object): 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) + params, raw_job_destination=raw_job_destination) return self.cached_job_destination def get_job_destination(self, params): diff --git a/test/unit/jobs/test_mapper.py b/test/unit/jobs/test_mapper.py index cac783a475d..9f0e70d6bbc 100644 --- a/test/unit/jobs/test_mapper.py +++ b/test/unit/jobs/test_mapper.py @@ -103,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: From 766039adbd17490bf6e758311329b03917fefa18 Mon Sep 17 00:00:00 2001 From: Nuwan Goonasekera Date: Thu, 15 Nov 2018 18:51:06 +0530 Subject: [PATCH 5/8] Added usage example of dynamic rule chaining --- config/job_conf.xml.sample_advanced | 24 ++++++++++++++++++++++++ 1 file changed, 24 insertions(+) diff --git a/config/job_conf.xml.sample_advanced b/config/job_conf.xml.sample_advanced index 921cc063aec..361d3693526 100644 --- a/config/job_conf.xml.sample_advanced +++ b/config/job_conf.xml.sample_advanced @@ -652,6 +652,30 @@ 50 + + + burst + local + burst_if_size + 2 + queued + + + python + to_destination_if_size + + galaxycloudrunner.rules + 1g + pulsar + local + + integration.chained_dyndest_rules.module2 from1 diff --git a/test/integration/chained_dyndest_rules/module1/__init__.py b/test/integration/chained_dyndest_rules/module1/__init__.py new file mode 100644 index 00000000000..e69de29bb2d diff --git a/test/integration/chained_dyndest_rules/module1/rules.py b/test/integration/chained_dyndest_rules/module1/rules.py new file mode 100644 index 00000000000..8a70a8069e5 --- /dev/null +++ b/test/integration/chained_dyndest_rules/module1/rules.py @@ -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" diff --git a/test/integration/chained_dyndest_rules/module2/__init__.py b/test/integration/chained_dyndest_rules/module2/__init__.py new file mode 100644 index 00000000000..e69de29bb2d diff --git a/test/integration/chained_dyndest_rules/module2/rules.py b/test/integration/chained_dyndest_rules/module2/rules.py new file mode 100644 index 00000000000..afe0d6dcc0c --- /dev/null +++ b/test/integration/chained_dyndest_rules/module2/rules.py @@ -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" diff --git a/test/integration/chained_dyndest_rules/module3/__init__.py b/test/integration/chained_dyndest_rules/module3/__init__.py new file mode 100644 index 00000000000..e69de29bb2d diff --git a/test/integration/chained_dyndest_rules/module3/rules.py b/test/integration/chained_dyndest_rules/module3/rules.py new file mode 100644 index 00000000000..72b3f9b68f5 --- /dev/null +++ b/test/integration/chained_dyndest_rules/module3/rules.py @@ -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}) diff --git a/test/integration/chained_dyndest_rules/rules.py b/test/integration/chained_dyndest_rules/rules.py deleted file mode 100644 index c877864ca6a..00000000000 --- a/test/integration/chained_dyndest_rules/rules.py +++ /dev/null @@ -1,21 +0,0 @@ -from galaxy.jobs import JobDestination - - -def dyndest_chain_1(): - # Check whether chaining dynamic job destinations work - return "dyn_dest2" - - -def dyndest_chain_2(tmp_dir_prefix): - # Chain to yet a third - return JobDestination( - runner="dynamic", - params={'type': 'python', - 'function': 'dyndest_chain_3', - 'tmp_dir_prefix': '%sand2' % tmp_dir_prefix}) - - -def dyndest_chain_3(tmp_dir_prefix): - tmp_dir = '$(mktemp %sand3XXXXXXXXXXXX)' % tmp_dir_prefix - return JobDestination(runner="local", - params={'tmp_dir': tmp_dir}) From ae7869170762b08425c264b0d9777df0efd977ea Mon Sep 17 00:00:00 2001 From: Marius van den Beek Date: Fri, 16 Nov 2018 17:12:45 +0530 Subject: [PATCH 8/8] Update test/integration/test_chained_dynamic_destinations.py Co-Authored-By: nuwang --- test/integration/test_chained_dynamic_destinations.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/test/integration/test_chained_dynamic_destinations.py b/test/integration/test_chained_dynamic_destinations.py index e018c629d22..2e3b8c8ef23 100644 --- a/test/integration/test_chained_dynamic_destinations.py +++ b/test/integration/test_chained_dynamic_destinations.py @@ -1,4 +1,4 @@ -"""Integration tests for the Pulsar embedded runner.""" +"""Integration tests for chained dynamic job destinations.""" import os import tempfile