mirror of
https://github.com/galaxyproject/galaxy.git
synced 2026-09-24 16:30:27 +08:00
Merge pull request #6776 from jmchilton/gxformat2_integration
Allow API import/export of format 2 workflows.
This commit is contained in:
@@ -382,6 +382,8 @@ class Configuration(object):
|
||||
# workflows built using these modules may not function in the
|
||||
# future.
|
||||
self.enable_beta_workflow_modules = string_as_bool(kwargs.get('enable_beta_workflow_modules', 'False'))
|
||||
# Enable use of gxformat2 workflows.
|
||||
self.enable_beta_workflow_format = string_as_bool(kwargs.get('enable_beta_workflow_format', 'False'))
|
||||
# These are not even beta - just experiments - don't use them unless
|
||||
# you want yours tools to be broken in the future.
|
||||
self.enable_beta_tool_formats = string_as_bool(kwargs.get('enable_beta_tool_formats', 'False'))
|
||||
|
||||
@@ -5,6 +5,12 @@ import logging
|
||||
import uuid
|
||||
from collections import namedtuple
|
||||
|
||||
from gxformat2 import (
|
||||
from_galaxy_native,
|
||||
ImporterGalaxyInterface,
|
||||
ImportOptions,
|
||||
python_to_workflow,
|
||||
)
|
||||
from six import string_types
|
||||
from sqlalchemy import and_
|
||||
from sqlalchemy.orm import joinedload, subqueryload
|
||||
@@ -235,6 +241,27 @@ class WorkflowContentsManager(UsesAnnotations):
|
||||
self.app = app
|
||||
self._resource_mapper_function = get_resource_mapper_function(app)
|
||||
|
||||
def normalize_workflow_format(self, as_dict):
|
||||
"""Process incoming workflow descriptions for consumption by other methods.
|
||||
|
||||
Currently this mostly means converting format 2 workflows into standard Galaxy
|
||||
workflow JSON for consumption for the rest of this module. In the future we will
|
||||
want to be a lot more percise about this - preserve the original description along
|
||||
side the data model and apply updates in a way that largely preserves YAML structure
|
||||
so workflows can be extracted.
|
||||
"""
|
||||
workflow_class = as_dict.get("class", None)
|
||||
if workflow_class == "GalaxyWorkflow" or "$graph" in as_dict or "yaml_content" 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()
|
||||
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(
|
||||
self,
|
||||
trans,
|
||||
@@ -346,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)
|
||||
@@ -380,6 +416,12 @@ class WorkflowContentsManager(UsesAnnotations):
|
||||
option describes the workflow in a context more tied to the current Galaxy instance and includes
|
||||
fields like 'url' and 'url' and actual unencoded step ids instead of 'order_index'.
|
||||
"""
|
||||
|
||||
def to_format_2(wf_dict, **kwds):
|
||||
if not trans.app.config.enable_beta_workflow_format:
|
||||
raise exceptions.ConfigDoesNotAllowException("Format2 workflows not enabled.")
|
||||
return from_galaxy_native(wf_dict, None, **kwds)
|
||||
|
||||
if version == '':
|
||||
version = None
|
||||
if version is not None:
|
||||
@@ -393,6 +435,12 @@ class WorkflowContentsManager(UsesAnnotations):
|
||||
wf_dict = self._workflow_to_dict_instance(stored, workflow=workflow, legacy=False)
|
||||
elif style == "run":
|
||||
wf_dict = self._workflow_to_dict_run(trans, stored, workflow=workflow)
|
||||
elif style == "format2":
|
||||
wf_dict = self._workflow_to_dict_export(trans, stored, workflow=workflow)
|
||||
wf_dict = to_format_2(wf_dict)
|
||||
elif style == "format2_wrapped_yaml":
|
||||
wf_dict = self._workflow_to_dict_export(trans, stored, workflow=workflow)
|
||||
wf_dict = to_format_2(wf_dict, json_wrapper=True)
|
||||
else:
|
||||
wf_dict = self._workflow_to_dict_export(trans, stored, workflow=workflow)
|
||||
if version:
|
||||
@@ -896,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
|
||||
|
||||
@@ -930,7 +978,8 @@ class WorkflowContentsManager(UsesAnnotations):
|
||||
|
||||
# Create the model class for the step
|
||||
steps.append(step)
|
||||
steps_by_external_id[step_dict['id']] = step
|
||||
external_id = step_dict["id"]
|
||||
steps_by_external_id[external_id] = step
|
||||
if 'workflow_outputs' in step_dict:
|
||||
workflow_outputs = step_dict['workflow_outputs']
|
||||
found_output_names = set([])
|
||||
@@ -954,7 +1003,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:
|
||||
@@ -964,11 +1013,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(
|
||||
@@ -977,6 +1025,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.
|
||||
|
||||
@@ -999,7 +1053,10 @@ class WorkflowContentsManager(UsesAnnotations):
|
||||
conn.input_step = step
|
||||
conn.input_name = input_name
|
||||
conn.output_name = conn_dict['output_name']
|
||||
conn.output_step = steps_by_external_id[conn_dict['id']]
|
||||
external_id = conn_dict['id']
|
||||
if external_id not in steps_by_external_id:
|
||||
raise KeyError("Failed to find external id %s in %s" % (external_id, steps_by_external_id.keys()))
|
||||
conn.output_step = steps_by_external_id[external_id]
|
||||
|
||||
input_subworkflow_step_index = conn_dict.get('input_subworkflow_step_id', None)
|
||||
if input_subworkflow_step_index is not None:
|
||||
@@ -1023,3 +1080,9 @@ class MissingToolsException(exceptions.MessageException):
|
||||
def __init__(self, workflow, errors):
|
||||
self.workflow = workflow
|
||||
self.errors = errors
|
||||
|
||||
|
||||
class Format2ConverterGalaxyInterface(ImporterGalaxyInterface):
|
||||
|
||||
def import_workflow(self, workflow, **kwds):
|
||||
raise NotImplementedError("Direct format 2 import of nested workflows is not yet implemented, use bioblend client.")
|
||||
|
||||
@@ -492,6 +492,7 @@ class WorkflowsAPIController(BaseAPIController, UsesStoredWorkflowMixin, UsesAnn
|
||||
stored_workflow = self.__get_stored_workflow(trans, id)
|
||||
workflow_dict = payload.get('workflow') or payload
|
||||
if workflow_dict:
|
||||
workflow_dict = self.__normalize_workflow(workflow_dict)
|
||||
new_workflow_name = workflow_dict.get('name') or workflow_dict.get('name')
|
||||
if new_workflow_name and new_workflow_name != stored_workflow.name:
|
||||
sanitized_name = sanitize_html(new_workflow_name)
|
||||
@@ -572,6 +573,7 @@ class WorkflowsAPIController(BaseAPIController, UsesStoredWorkflowMixin, UsesAnn
|
||||
raise exceptions.MessageException("The data content does not appear to be a valid workflow.")
|
||||
if not data:
|
||||
raise exceptions.MessageException("The data content is missing.")
|
||||
data = self.__normalize_workflow(data)
|
||||
workflow, missing_tool_tups = self._workflow_from_dict(trans, data, source=source)
|
||||
workflow = workflow.latest_workflow
|
||||
if workflow.has_errors:
|
||||
@@ -584,6 +586,7 @@ class WorkflowsAPIController(BaseAPIController, UsesStoredWorkflowMixin, UsesAnn
|
||||
|
||||
def __api_import_new_workflow(self, trans, payload, **kwd):
|
||||
data = payload['workflow']
|
||||
data = self.__normalize_workflow(data)
|
||||
import_tools = util.string_as_bool(payload.get("import_tools", False))
|
||||
if import_tools and not trans.user_is_admin:
|
||||
raise exceptions.AdminRequiredException()
|
||||
@@ -651,6 +654,9 @@ class WorkflowsAPIController(BaseAPIController, UsesStoredWorkflowMixin, UsesAnn
|
||||
'fill_defaults': fill_defaults,
|
||||
}
|
||||
|
||||
def __normalize_workflow(self, as_dict):
|
||||
return self.workflow_contents_manager.normalize_workflow_format(as_dict)
|
||||
|
||||
@expose_api
|
||||
def import_shared_workflow_deprecated(self, trans, payload, **kwd):
|
||||
"""
|
||||
|
||||
@@ -2107,6 +2107,13 @@ mapping:
|
||||
Enable beta workflow modules that should not yet be considered part of Galaxy's
|
||||
stable API.
|
||||
|
||||
enable_beta_workflow_format:
|
||||
type: bool
|
||||
default: false
|
||||
required: false
|
||||
desc: |
|
||||
Enable import and export of workflows as Galaxy Format 2 workflows.
|
||||
|
||||
force_beta_workflow_scheduled_min_steps:
|
||||
type: int
|
||||
default: 250
|
||||
|
||||
@@ -38,6 +38,7 @@ if [ ! -z "$GALAXY_RUN_WITH_TEST_TOOLS" ];
|
||||
then
|
||||
export GALAXY_CONFIG_OVERRIDE_TOOL_CONFIG_FILE="test/functional/tools/samples_tool_conf.xml"
|
||||
export GALAXY_CONFIG_ENABLE_BETA_WORKFLOW_MODULES="true"
|
||||
export GALAXY_CONFIG_ENABLE_BETA_WORKFLOW_FORMAT="true"
|
||||
export GALAXY_CONFIG_OVERRIDE_ENABLE_BETA_TOOL_FORMATS="true"
|
||||
export GALAXY_CONFIG_OVERRIDE_WEBHOOKS_DIR="test/functional/webhooks"
|
||||
fi
|
||||
|
||||
+276
-286
File diff suppressed because it is too large
Load Diff
@@ -22,13 +22,17 @@ class WorkflowsFromYamlApiTestCase(BaseWorkflowsApiTestCase):
|
||||
def setUp(self):
|
||||
super(WorkflowsFromYamlApiTestCase, self).setUp()
|
||||
|
||||
def _upload_and_download(self, yaml_workflow):
|
||||
workflow_id = self._upload_yaml_workflow(yaml_workflow)
|
||||
workflow = self._get("workflows/%s/download" % workflow_id).json()
|
||||
return workflow
|
||||
def _upload_and_download(self, yaml_workflow, **kwds):
|
||||
style = None
|
||||
if "style" in kwds:
|
||||
style = kwds.pop("style")
|
||||
workflow_id = self._upload_yaml_workflow(yaml_workflow, **kwds)
|
||||
return self.workflow_populator.download_workflow(workflow_id, style=style)
|
||||
|
||||
def test_simple_upload(self):
|
||||
workflow = self._upload_and_download(WORKFLOW_SIMPLE_CAT_AND_RANDOM_LINES)
|
||||
workflow = self._upload_and_download(WORKFLOW_SIMPLE_CAT_AND_RANDOM_LINES, client_convert=False)
|
||||
|
||||
assert workflow["annotation"].startswith("Simple workflow that ")
|
||||
|
||||
tool_count = {'random_lines1': 0, 'cat1': 0}
|
||||
input_found = False
|
||||
@@ -47,6 +51,10 @@ class WorkflowsFromYamlApiTestCase(BaseWorkflowsApiTestCase):
|
||||
assert tool_count['random_lines1'] == 1
|
||||
assert tool_count['cat1'] == 2
|
||||
|
||||
workflow_as_format2 = self._upload_and_download(WORKFLOW_SIMPLE_CAT_AND_RANDOM_LINES, client_convert=False, style="format2")
|
||||
assert workflow_as_format2["doc"].startswith("Simple workflow that")
|
||||
|
||||
|
||||
# FIXME: This test fails on some machines due to (we're guessing) yaml.safe_loading
|
||||
# order being not guaranteed and inconsistent across platforms. The workflow
|
||||
# yaml.safe_loader probably needs to enforce order using something like the
|
||||
@@ -86,15 +94,16 @@ input1: "hello world"
|
||||
|
||||
def test_inputs_to_steps(self):
|
||||
history_id = self.dataset_populator.new_history()
|
||||
self._run_jobs(WORKFLOW_SIMPLE_CAT_TWICE, test_data={"input1": "hello world"}, history_id=history_id)
|
||||
self._run_jobs(WORKFLOW_SIMPLE_CAT_TWICE, test_data={"input1": "hello world"}, history_id=history_id, round_trip_format_conversion=True)
|
||||
contents1 = self.dataset_populator.get_history_dataset_content(history_id)
|
||||
self.assertEqual(contents1.strip(), "hello world\nhello world")
|
||||
|
||||
def test_outputs(self):
|
||||
workflow_id = self._upload_yaml_workflow(WORKFLOW_WITH_OUTPUTS)
|
||||
workflow_id = self._upload_yaml_workflow(WORKFLOW_WITH_OUTPUTS, round_trip_format_conversion=True)
|
||||
workflow = self._get("workflows/%s/download" % workflow_id).json()
|
||||
self.assertEqual(workflow["steps"]["1"]["workflow_outputs"][0]["output_name"], "out_file1")
|
||||
self.assertEqual(workflow["steps"]["1"]["workflow_outputs"][0]["label"], "wf_output_1")
|
||||
workflow = self.workflow_populator.download_workflow(workflow_id, style="format2")
|
||||
|
||||
def test_runtime_inputs(self):
|
||||
workflow = self._upload_and_download(WORKFLOW_RUNTIME_PARAMETER_SIMPLE)
|
||||
@@ -116,11 +125,12 @@ class: GalaxyWorkflow
|
||||
inputs:
|
||||
outer_input: data
|
||||
steps:
|
||||
- tool_id: cat1
|
||||
label: first_cat
|
||||
first_cat:
|
||||
tool_id: cat1
|
||||
in:
|
||||
input1: outer_input
|
||||
- run:
|
||||
nested_workflow:
|
||||
run:
|
||||
class: GalaxyWorkflow
|
||||
inputs:
|
||||
inner_input: data
|
||||
@@ -133,16 +143,10 @@ steps:
|
||||
seed_source:
|
||||
seed_source_selector: set_seed
|
||||
seed: asdf
|
||||
label: nested_workflow
|
||||
in:
|
||||
inner_input: first_cat/out_file1
|
||||
|
||||
test_data:
|
||||
outer_input:
|
||||
value: 1.bed
|
||||
type: File
|
||||
""")
|
||||
workflow = self._get("workflows/%s/download" % workflow_id).json()
|
||||
""", client_convert=False)
|
||||
workflow = self.workflow_populator.download_workflow(workflow_id)
|
||||
by_label = self._steps_by_label(workflow)
|
||||
if "nested_workflow" not in by_label:
|
||||
template = "Workflow [%s] does not contain label 'nested_workflow'."
|
||||
@@ -173,50 +177,88 @@ test_data:
|
||||
# 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:
|
||||
outer_cat:
|
||||
tool_id: cat
|
||||
in:
|
||||
input1: outer_input
|
||||
nested_workflow_1:
|
||||
run: '#nested'
|
||||
in:
|
||||
inner_input: outer_cat/out_file1
|
||||
nested_workflow_2:
|
||||
run: '#nested'
|
||||
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
|
||||
steps:
|
||||
- label: test_input
|
||||
test_input:
|
||||
type: input
|
||||
- label: first_cat
|
||||
first_cat:
|
||||
tool_id: cat1
|
||||
state:
|
||||
input1:
|
||||
$link: test_input
|
||||
- label: the_pause
|
||||
the_pause:
|
||||
type: pause
|
||||
in:
|
||||
input: first_cat/out_file1
|
||||
- label: second_cat
|
||||
second_cat:
|
||||
tool_id: cat1
|
||||
in:
|
||||
input1: the_pause
|
||||
""")
|
||||
print(self._get("workflows/%s/download" % workflow_id).json())
|
||||
self.workflow_populator.dump_workflow(workflow_id)
|
||||
|
||||
def test_implicit_connections(self):
|
||||
workflow_id = self._upload_yaml_workflow("""
|
||||
class: GalaxyWorkflow
|
||||
inputs:
|
||||
test_input: data
|
||||
steps:
|
||||
- label: test_input
|
||||
type: input
|
||||
- label: first_cat
|
||||
first_cat:
|
||||
tool_id: cat1
|
||||
state:
|
||||
input1:
|
||||
$link: test_input
|
||||
- label: the_pause
|
||||
in:
|
||||
input1: test_input
|
||||
the_pause:
|
||||
type: pause
|
||||
connect:
|
||||
input:
|
||||
- first_cat#out_file1
|
||||
- label: second_cat
|
||||
in:
|
||||
input: first_cat/out_file1
|
||||
second_cat:
|
||||
tool_id: cat1
|
||||
state:
|
||||
input1:
|
||||
$link: the_pause
|
||||
- label: third_cat
|
||||
in:
|
||||
input1: the_pause
|
||||
third_cat:
|
||||
tool_id: cat1
|
||||
connect:
|
||||
$step: second_cat
|
||||
@@ -224,22 +266,21 @@ steps:
|
||||
input1:
|
||||
$link: test_input
|
||||
""")
|
||||
workflow = self._get("workflows/%s/download" % workflow_id).json()
|
||||
print(workflow)
|
||||
self.workflow_populator.dump_workflow(workflow_id)
|
||||
|
||||
@uses_test_history()
|
||||
def test_conditional_ints(self, history_id):
|
||||
self._run_jobs("""
|
||||
class: GalaxyWorkflow
|
||||
steps:
|
||||
- label: test_input
|
||||
test_input:
|
||||
tool_id: disambiguate_cond
|
||||
state:
|
||||
p3:
|
||||
use: true
|
||||
files:
|
||||
attach_files: false
|
||||
""", test_data={}, history_id=history_id)
|
||||
""", test_data={}, history_id=history_id, round_trip_format_conversion=True)
|
||||
content = self.dataset_populator.get_history_dataset_content(history_id)
|
||||
assert "no file specified" in content
|
||||
assert "7 7 4" in content
|
||||
@@ -247,7 +288,7 @@ steps:
|
||||
self._run_jobs("""
|
||||
class: GalaxyWorkflow
|
||||
steps:
|
||||
- label: test_input
|
||||
test_input:
|
||||
tool_id: disambiguate_cond
|
||||
state:
|
||||
p3:
|
||||
@@ -255,7 +296,7 @@ steps:
|
||||
p3v: 5
|
||||
files:
|
||||
attach_files: false
|
||||
""", test_data={}, history_id=history_id)
|
||||
""", test_data={}, history_id=history_id, round_trip_format_conversion=True)
|
||||
content = self.dataset_populator.get_history_dataset_content(history_id)
|
||||
assert "no file specified" in content
|
||||
assert "7 7 5" in content
|
||||
|
||||
@@ -201,6 +201,7 @@ def setup_galaxy_config(
|
||||
cleanup_job='onsuccess',
|
||||
data_manager_config_file=data_manager_config_file,
|
||||
enable_beta_tool_formats=True,
|
||||
enable_beta_workflow_format=True,
|
||||
expose_dataset_path=True,
|
||||
file_path=file_path,
|
||||
ftp_upload_purge=False,
|
||||
|
||||
+32
-4
@@ -547,8 +547,18 @@ class BaseWorkflowPopulator(object):
|
||||
return upload_response
|
||||
|
||||
def upload_yaml_workflow(self, has_yaml, **kwds):
|
||||
round_trip_conversion = kwds.get("round_trip_format_conversion", False)
|
||||
client_convert = kwds.pop("client_convert", not round_trip_conversion)
|
||||
kwds["convert"] = client_convert
|
||||
workflow = convert_and_import_workflow(has_yaml, galaxy_interface=self, **kwds)
|
||||
return workflow["id"]
|
||||
workflow_id = workflow["id"]
|
||||
if round_trip_conversion:
|
||||
workflow_yaml_wrapped = self.download_workflow(workflow_id, style="format2_wrapped_yaml")
|
||||
assert "yaml_content" in workflow_yaml_wrapped, workflow_yaml_wrapped
|
||||
round_trip_converted_content = workflow_yaml_wrapped["yaml_content"]
|
||||
workflow_id = self.upload_yaml_workflow(round_trip_converted_content, client_convert=False, round_trip_conversion=False)
|
||||
|
||||
return workflow_id
|
||||
|
||||
def wait_for_invocation(self, workflow_id, invocation_id, timeout=DEFAULT_TIMEOUT):
|
||||
url = "workflows/%s/usage/%s" % (workflow_id, invocation_id)
|
||||
@@ -578,7 +588,15 @@ class BaseWorkflowPopulator(object):
|
||||
else:
|
||||
return invocation_response
|
||||
|
||||
def run_workflow(self, has_workflow, test_data=None, history_id=None, wait=True, source_type=None, jobs_descriptions=None, expected_response=200, assert_ok=True):
|
||||
def download_workflow(self, workflow_id, style=None):
|
||||
params = {}
|
||||
if style is not None:
|
||||
params["style"] = style
|
||||
response = self._get("workflows/%s/download" % workflow_id, data=params)
|
||||
api_asserts.assert_status_code_is(response, 200)
|
||||
return response.json()
|
||||
|
||||
def run_workflow(self, has_workflow, test_data=None, history_id=None, wait=True, source_type=None, jobs_descriptions=None, expected_response=200, assert_ok=True, client_convert=None, round_trip_format_conversion=False, raw_yaml=False):
|
||||
"""High-level wrapper around workflow API, etc. to invoke format 2 workflows."""
|
||||
workflow_populator = self
|
||||
|
||||
@@ -588,7 +606,10 @@ class BaseWorkflowPopulator(object):
|
||||
content = open(filename, "r").read()
|
||||
return content
|
||||
|
||||
workflow_id = workflow_populator.upload_yaml_workflow(has_workflow, source_type=source_type)
|
||||
if client_convert is None:
|
||||
client_convert = not round_trip_format_conversion
|
||||
|
||||
workflow_id = workflow_populator.upload_yaml_workflow(has_workflow, source_type=source_type, client_convert=client_convert, round_trip_format_conversion=round_trip_format_conversion, raw_yaml=raw_yaml)
|
||||
|
||||
if test_data is None:
|
||||
if jobs_descriptions is None:
|
||||
@@ -636,6 +657,13 @@ class BaseWorkflowPopulator(object):
|
||||
workflow_request=workflow_request
|
||||
)
|
||||
|
||||
def dump_workflow(self, workflow_id, style=None):
|
||||
raw_workflow = self.download_workflow(workflow_id, style=style)
|
||||
if style == "format2_wrapped_yaml":
|
||||
print(raw_workflow["yaml_content"])
|
||||
else:
|
||||
print(json.dumps(raw_workflow, sort_keys=True, indent=2))
|
||||
|
||||
|
||||
RunJobsSummary = namedtuple('RunJobsSummary', ['history_id', 'workflow_id', 'invocation_id', 'inputs', 'jobs', 'invocation', 'workflow_request'])
|
||||
|
||||
@@ -662,7 +690,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()
|
||||
|
||||
|
||||
|
||||
@@ -2,10 +2,15 @@
|
||||
|
||||
WORKFLOW_SIMPLE_CAT_AND_RANDOM_LINES = """
|
||||
class: GalaxyWorkflow
|
||||
doc: |
|
||||
Simple workflow that no-op cats a file and then selects 10 random lines.
|
||||
inputs:
|
||||
- id: the_input
|
||||
the_input:
|
||||
type: data
|
||||
doc: input doc
|
||||
steps:
|
||||
- tool_id: cat1
|
||||
doc: cat doc
|
||||
in:
|
||||
input1: the_input
|
||||
- tool_id: cat1
|
||||
@@ -28,8 +33,8 @@ class: GalaxyWorkflow
|
||||
inputs:
|
||||
input1: data
|
||||
steps:
|
||||
- tool_id: cat
|
||||
label: first_cat
|
||||
first_cat:
|
||||
tool_id: cat
|
||||
in:
|
||||
input1: input1
|
||||
queries_0|input2: input1
|
||||
@@ -41,7 +46,8 @@ class: GalaxyWorkflow
|
||||
inputs:
|
||||
input1: data
|
||||
steps:
|
||||
- tool_id: multiple_versions
|
||||
mul_versions:
|
||||
tool_id: multiple_versions
|
||||
tool_version: "0.0.1"
|
||||
state:
|
||||
inttest: 8
|
||||
@@ -53,7 +59,8 @@ class: GalaxyWorkflow
|
||||
inputs:
|
||||
input1: data
|
||||
steps:
|
||||
- tool_id: multiple_versions
|
||||
mul_versions:
|
||||
tool_id: multiple_versions
|
||||
tool_version: "0.0.1"
|
||||
state:
|
||||
inttest: "moocow"
|
||||
@@ -65,11 +72,12 @@ class: GalaxyWorkflow
|
||||
inputs:
|
||||
text_input: data
|
||||
steps:
|
||||
- label: split_up
|
||||
split_up:
|
||||
tool_id: collection_creates_pair
|
||||
in:
|
||||
input1: text_input
|
||||
- tool_id: collection_paired_test
|
||||
paired:
|
||||
tool_id: collection_paired_test
|
||||
in:
|
||||
f1: split_up/paired_output
|
||||
test_data:
|
||||
@@ -83,25 +91,21 @@ test_data:
|
||||
|
||||
WORKFLOW_WITH_DYNAMIC_OUTPUT_COLLECTION = """
|
||||
class: GalaxyWorkflow
|
||||
inputs:
|
||||
text_input1: data
|
||||
text_input2: data
|
||||
steps:
|
||||
- label: text_input1
|
||||
type: input
|
||||
- label: text_input2
|
||||
type: input
|
||||
- label: cat_inputs
|
||||
cat_inputs:
|
||||
tool_id: cat1
|
||||
state:
|
||||
input1:
|
||||
$link: text_input1
|
||||
queries:
|
||||
- input2:
|
||||
$link: text_input2
|
||||
- label: split_up
|
||||
in:
|
||||
input1: text_input1
|
||||
queries_0|input2: text_input2
|
||||
split_up:
|
||||
tool_id: collection_split_on_column
|
||||
state:
|
||||
input1:
|
||||
$link: cat_inputs#out_file1
|
||||
- tool_id: cat_list
|
||||
in:
|
||||
input1: cat_inputs/out_file1
|
||||
cat_list:
|
||||
tool_id: cat_list
|
||||
in:
|
||||
input1: split_up/split_output
|
||||
test_data:
|
||||
@@ -121,8 +125,8 @@ inputs:
|
||||
type: collection
|
||||
collection_type: list
|
||||
steps:
|
||||
- tool_id: cat
|
||||
label: cat
|
||||
cat:
|
||||
tool_id: cat
|
||||
in:
|
||||
input1: input1
|
||||
"""
|
||||
@@ -152,7 +156,7 @@ class: GalaxyWorkflow
|
||||
inputs:
|
||||
input_c: collection
|
||||
steps:
|
||||
- label: apply
|
||||
apply:
|
||||
tool_id: __APPLY_RULES__
|
||||
state:
|
||||
input:
|
||||
@@ -166,8 +170,8 @@ steps:
|
||||
mapping:
|
||||
- type: list_identifiers
|
||||
columns: [0, 1]
|
||||
- tool_id: random_lines1
|
||||
label: random_lines
|
||||
random_lines:
|
||||
tool_id: random_lines1
|
||||
state:
|
||||
num_lines: 1
|
||||
input:
|
||||
@@ -191,7 +195,7 @@ class: GalaxyWorkflow
|
||||
inputs:
|
||||
input_c: collection
|
||||
steps:
|
||||
- label: apply
|
||||
apply:
|
||||
tool_id: __APPLY_RULES__
|
||||
state:
|
||||
input:
|
||||
@@ -205,8 +209,8 @@ steps:
|
||||
mapping:
|
||||
- type: list_identifiers
|
||||
columns: [0, 1]
|
||||
- tool_id: collection_creates_list
|
||||
label: copy_list
|
||||
copy_list:
|
||||
tool_id: collection_creates_list
|
||||
in:
|
||||
input1: apply/output
|
||||
test_data:
|
||||
@@ -228,11 +232,12 @@ outputs:
|
||||
outer_output:
|
||||
outputSource: second_cat/out_file1
|
||||
steps:
|
||||
- tool_id: cat1
|
||||
label: first_cat
|
||||
first_cat:
|
||||
tool_id: cat1
|
||||
in:
|
||||
input1: outer_input
|
||||
- run:
|
||||
nested_workflow:
|
||||
run:
|
||||
class: GalaxyWorkflow
|
||||
inputs:
|
||||
inner_input: data
|
||||
@@ -240,8 +245,8 @@ steps:
|
||||
workflow_output:
|
||||
outputSource: random_lines/out_file1
|
||||
steps:
|
||||
- tool_id: random_lines1
|
||||
label: random_lines
|
||||
random_lines:
|
||||
tool_id: random_lines1
|
||||
state:
|
||||
num_lines: 1
|
||||
input:
|
||||
@@ -249,17 +254,13 @@ steps:
|
||||
seed_source:
|
||||
seed_source_selector: set_seed
|
||||
seed: asdf
|
||||
label: nested_workflow
|
||||
in:
|
||||
inner_input: first_cat/out_file1
|
||||
- tool_id: cat1
|
||||
label: second_cat
|
||||
state:
|
||||
input1:
|
||||
$link: nested_workflow#workflow_output
|
||||
queries:
|
||||
- input2:
|
||||
$link: nested_workflow#workflow_output
|
||||
second_cat:
|
||||
tool_id: cat1
|
||||
in:
|
||||
input1: nested_workflow/workflow_output
|
||||
queries_0|input2: nested_workflow/workflow_output
|
||||
"""
|
||||
|
||||
|
||||
@@ -271,7 +272,8 @@ outputs:
|
||||
outer_output:
|
||||
outputSource: nested_workflow/workflow_output
|
||||
steps:
|
||||
- run:
|
||||
nested_workflow:
|
||||
run:
|
||||
class: GalaxyWorkflow
|
||||
inputs:
|
||||
inner_input: data
|
||||
@@ -289,7 +291,6 @@ steps:
|
||||
seed_source:
|
||||
seed_source_selector: set_seed
|
||||
seed: asdf
|
||||
label: nested_workflow
|
||||
in:
|
||||
inner_input: outer_input
|
||||
"""
|
||||
@@ -300,15 +301,16 @@ class: GalaxyWorkflow
|
||||
inputs:
|
||||
input1: data
|
||||
steps:
|
||||
- tool_id: cat1
|
||||
label: first_cat
|
||||
first_cat:
|
||||
tool_id: cat1
|
||||
outputs:
|
||||
out_file1:
|
||||
hide: true
|
||||
rename: "the new value"
|
||||
in:
|
||||
input1: input1
|
||||
- tool_id: cat1
|
||||
second_cat:
|
||||
tool_id: cat1
|
||||
in:
|
||||
input1: first_cat/out_file1
|
||||
"""
|
||||
@@ -319,7 +321,8 @@ class: GalaxyWorkflow
|
||||
inputs:
|
||||
input1: data
|
||||
steps:
|
||||
- tool_id: random_lines1
|
||||
random:
|
||||
tool_id: random_lines1
|
||||
runtime_inputs:
|
||||
- num_lines
|
||||
state:
|
||||
@@ -336,11 +339,12 @@ class: GalaxyWorkflow
|
||||
inputs:
|
||||
input1: data
|
||||
steps:
|
||||
- label: the_pause
|
||||
the_pause:
|
||||
type: pause
|
||||
in:
|
||||
input: input1
|
||||
- tool_id: random_lines1
|
||||
random:
|
||||
tool_id: random_lines1
|
||||
runtime_inputs:
|
||||
- num_lines
|
||||
state:
|
||||
@@ -356,8 +360,8 @@ class: GalaxyWorkflow
|
||||
inputs:
|
||||
input1: data
|
||||
steps:
|
||||
- tool_id: cat
|
||||
label: first_cat
|
||||
first_cat:
|
||||
tool_id: cat
|
||||
state:
|
||||
input1:
|
||||
$link: input1
|
||||
@@ -376,11 +380,10 @@ class: GalaxyWorkflow
|
||||
inputs:
|
||||
input1: data
|
||||
steps:
|
||||
- tool_id: cat
|
||||
label: first_cat
|
||||
state:
|
||||
input1:
|
||||
$link: input1
|
||||
first_cat:
|
||||
tool_id: cat
|
||||
in:
|
||||
input1: input1
|
||||
outputs:
|
||||
out_file1:
|
||||
rename: "${replaceme} suffix"
|
||||
@@ -394,7 +397,8 @@ outputs:
|
||||
outer_output:
|
||||
outputSource: nested_workflow/workflow_output
|
||||
steps:
|
||||
- run:
|
||||
nested_workflow:
|
||||
run:
|
||||
class: GalaxyWorkflow
|
||||
inputs:
|
||||
inner_input: data
|
||||
@@ -409,7 +413,6 @@ steps:
|
||||
outputs:
|
||||
out_file1:
|
||||
rename: "${replaceme} suffix"
|
||||
label: nested_workflow
|
||||
in:
|
||||
inner_input: outer_input
|
||||
"""
|
||||
@@ -422,12 +425,9 @@ outputs:
|
||||
wf_output_1:
|
||||
outputSource: first_cat/out_file1
|
||||
steps:
|
||||
- tool_id: cat1
|
||||
label: first_cat
|
||||
state:
|
||||
input1:
|
||||
$link: input1
|
||||
queries:
|
||||
- input2:
|
||||
$link: input1
|
||||
first_cat:
|
||||
tool_id: cat1
|
||||
in:
|
||||
input1: input1
|
||||
queries_0|input2: input1
|
||||
"""
|
||||
|
||||
Reference in New Issue
Block a user