diff --git a/lib/galaxy/managers/workflows.py b/lib/galaxy/managers/workflows.py index 508b4e8d38f..65831f50f50 100644 --- a/lib/galaxy/managers/workflows.py +++ b/lib/galaxy/managers/workflows.py @@ -8,6 +8,7 @@ from collections import namedtuple from gxformat2 import ( from_galaxy_native, ImporterGalaxyInterface, + ImportOptions, python_to_workflow, ) from six import string_types @@ -250,13 +251,15 @@ class WorkflowContentsManager(UsesAnnotations): so workflows can be extracted. """ workflow_class = as_dict.get("class", None) - if workflow_class == "GalaxyWorkflow": + if workflow_class == "GalaxyWorkflow" or "$graph" in as_dict: if not self.app.config.enable_beta_workflow_format: raise exceptions.ConfigDoesNotAllowException("Format2 workflows not enabled.") # Format 2 Galaxy workflow. galaxy_interface = Format2ConverterGalaxyInterface() - as_dict = python_to_workflow(as_dict, galaxy_interface, workflow_directory=None) + import_options = ImportOptions() + import_options.deduplicate_subworkflows = True + as_dict = python_to_workflow(as_dict, galaxy_interface, workflow_directory=None, import_options=import_options) return as_dict def build_workflow_from_dict( @@ -370,12 +373,21 @@ class WorkflowContentsManager(UsesAnnotations): # but do need to use to make connections steps_by_external_id = {} + # Preload dependent workflows with locally defined content_ids. + subworkflows = data.get("subworkflows") + subworkflow_id_map = None + if subworkflows: + subworkflow_id_map = {} + for key, subworkflow_dict in subworkflows.items(): + subworkflow = self.__build_embedded_subworkflow(trans, subworkflow_dict, **kwds) + subworkflow_id_map[key] = subworkflow + # Keep track of tools required by the workflow that are not available in # the local Galaxy instance. Each tuple in the list of missing_tool_tups # will be ( tool_id, tool_name, tool_version ). missing_tool_tups = [] for step_dict in self.__walk_step_dicts(data): - self.__load_subworkflows(trans, step_dict) + self.__load_subworkflows(trans, step_dict, subworkflow_id_map, **kwds) for step_dict in self.__walk_step_dicts(data): module, step = self.__module_from_dict(trans, steps, steps_by_external_id, step_dict, **kwds) @@ -932,11 +944,11 @@ class WorkflowContentsManager(UsesAnnotations): yield step_dict - def __load_subworkflows(self, trans, step_dict): + def __load_subworkflows(self, trans, step_dict, subworkflow_id_map, **kwds): step_type = step_dict.get("type", None) if step_type == "subworkflow": subworkflow = self.__load_subworkflow_from_step_dict( - trans, step_dict + trans, step_dict, subworkflow_id_map, **kwds ) step_dict["subworkflow"] = subworkflow @@ -990,7 +1002,7 @@ class WorkflowContentsManager(UsesAnnotations): trans.sa_session.add(m) return module, step - def __load_subworkflow_from_step_dict(self, trans, step_dict): + def __load_subworkflow_from_step_dict(self, trans, step_dict, subworkflow_id_map, **kwds): embedded_subworkflow = step_dict.get("subworkflow", None) subworkflow_id = step_dict.get("content_id", None) if embedded_subworkflow and subworkflow_id: @@ -1000,11 +1012,10 @@ class WorkflowContentsManager(UsesAnnotations): raise Exception("Subworkflow step must define either subworkflow or content_id.") if embedded_subworkflow: - subworkflow = self.build_workflow_from_dict( - trans, - embedded_subworkflow, - create_stored_workflow=False, - ).workflow + subworkflow = self.__build_embedded_subworkflow(trans, embedded_subworkflow, **kwds) + elif subworkflow_id_map is not None: + # Interpret content_id as a workflow local thing. + subworkflow = subworkflow_id_map[subworkflow_id[1:]] else: workflow_manager = WorkflowsManager(self.app) subworkflow = workflow_manager.get_owned_workflow( @@ -1013,6 +1024,12 @@ class WorkflowContentsManager(UsesAnnotations): return subworkflow + def __build_embedded_subworkflow(self, trans, data, **kwds): + subworkflow = self.build_workflow_from_dict( + trans, data, create_stored_workflow=False, fill_defaults=kwds.get("fill_defaults", False) + ).workflow + return subworkflow + def __connect_workflow_steps(self, steps, steps_by_external_id): """ Second pass to deal with connections between steps. diff --git a/test/api/test_workflows_from_yaml.py b/test/api/test_workflows_from_yaml.py index 78d84d553f1..ff67c99b3bc 100644 --- a/test/api/test_workflows_from_yaml.py +++ b/test/api/test_workflows_from_yaml.py @@ -177,6 +177,47 @@ steps: # content = self.dataset_populator.get_history_dataset_content( history_id ) # self.assertEqual("chr5\t131424298\t131424460\tCCDS4149.1_cds_0_0_chr5_131424299_f\t0\t+\n", content) + def test_subworkflow_duplicate(self): + duplicate_subworkflow_invocate_wf = """ +format-version: "v2.0" +$graph: +- id: nested + class: GalaxyWorkflow + inputs: + inner_input: data + outputs: + inner_output: + outputSource: inner_cat/out_file1 + steps: + inner_cat: + tool_id: cat + in: + input1: inner_input + queries_0|input2: inner_input + +- id: main + class: GalaxyWorkflow + inputs: + outer_input: data + steps: + - tool_id: cat + label: outer_cat + in: + input1: outer_input + - run: '#nested' + label: nested_workflow_1 + in: + inner_input: outer_cat/out_file1 + - run: '#nested' + label: nested_workflow_2 + in: + inner_input: nested_workflow_1/inner_output +""" + history_id = self.dataset_populator.new_history() + self._run_jobs(duplicate_subworkflow_invocate_wf, test_data={"outer_input": "hello world"}, history_id=history_id, client_convert=False) + content = self.dataset_populator.get_history_dataset_content(history_id) + assert content == "hello world\nhello world\nhello world\nhello world\n" + def test_pause(self): workflow_id = self._upload_yaml_workflow(""" class: GalaxyWorkflow diff --git a/test/base/populators.py b/test/base/populators.py index f438a4c478e..deda3c35dab 100644 --- a/test/base/populators.py +++ b/test/base/populators.py @@ -689,7 +689,7 @@ class WorkflowPopulator(BaseWorkflowPopulator, ImporterGalaxyInterface): } data.update(**kwds) upload_response = self._post("workflows", data=data) - assert upload_response.status_code == 200, upload_response + assert upload_response.status_code == 200, upload_response.content return upload_response.json()