mirror of
https://github.com/galaxyproject/galaxy.git
synced 2026-09-24 16:30:27 +08:00
Support for de-duplicated workflow import.
This commit is contained in:
@@ -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.
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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()
|
||||
|
||||
|
||||
|
||||
Reference in New Issue
Block a user