Typing in workflows tests...

This commit is contained in:
John Chilton
2021-10-04 20:27:31 -04:00
parent 2770cd0bf4
commit 5483a4ffbd
2 changed files with 119 additions and 85 deletions
+95 -73
View File
@@ -2,6 +2,7 @@ import json
import os
import time
from json import dumps
from typing import Any, cast, Dict, Union
from uuid import uuid4
import pytest
@@ -12,6 +13,7 @@ from galaxy_test.base import rules_test_data
from galaxy_test.base.populators import (
DatasetCollectionPopulator,
DatasetPopulator,
RunJobsSummary,
skip_without_tool,
wait_on,
WorkflowPopulator
@@ -117,7 +119,7 @@ class BaseWorkflowsApiTestCase(ApiTestCase):
upload_response = self.workflow_populator.import_workflow(workflow, **kwds)
return upload_response
def _upload_yaml_workflow(self, has_yaml, **kwds):
def _upload_yaml_workflow(self, has_yaml, **kwds) -> str:
return self.workflow_populator.upload_yaml_workflow(has_yaml, **kwds)
def _setup_workflow_run(self, workflow=None, inputs_by='step_id', history_id=None, workflow_id=None):
@@ -182,12 +184,19 @@ class BaseWorkflowsApiTestCase(ApiTestCase):
invocation_details = invocation_details_response.json()
return invocation_details
def _run_jobs(self, has_workflow, history_id=None, **kwds):
def _run_jobs(self, has_workflow, history_id=None, **kwds) -> Union[Dict[str, Any], RunJobsSummary]:
if history_id is None:
history_id = self.history_id
return self.workflow_populator.run_workflow(has_workflow, history_id=history_id, **kwds)
def _run_workflow(self, has_workflow, history_id=None, **kwds) -> RunJobsSummary:
if history_id is None:
history_id = self.history_id
assert "expected_response" not in kwds
run_summary = self.workflow_populator.run_workflow(has_workflow, history_id=history_id, **kwds)
return cast(RunJobsSummary, run_summary)
def _history_jobs(self, history_id):
return self._get("jobs", {"history_id": history_id, "order_by": "create_time"}).json()
@@ -217,6 +226,8 @@ class BaseWorkflowsApiTestCase(ApiTestCase):
class ChangeDatatypeTestCase:
dataset_populator: DatasetPopulator
workflow_populator: WorkflowPopulator
def test_assign_column_pja(self):
with self.dataset_populator.test_history() as history_id:
@@ -933,9 +944,9 @@ steps:
self.assertEqual(invocation_response.json().get('err_msg'), "Workflow was not invoked; some required tools are not installed.")
@skip_without_tool("collection_creates_pair")
def test_workflow_run_output_collections(self):
def test_workflow_run_output_collections(self) -> None:
with self.dataset_populator.test_history() as history_id:
self._run_jobs(WORKFLOW_WITH_OUTPUT_COLLECTION, history_id=history_id, assert_ok=True, wait=True)
self._run_workflow(WORKFLOW_WITH_OUTPUT_COLLECTION, history_id=history_id)
self.assertEqual("a\nc\nb\nd\n", self.dataset_populator.get_history_dataset_content(history_id, hid=0))
@skip_without_tool("job_properties")
@@ -1033,7 +1044,7 @@ steps:
@skip_without_tool("identifier_collection")
def test_workflow_resume_with_mapped_over_input(self):
with self.dataset_populator.test_history() as history_id:
job_summary = self._run_jobs("""
self._run_workflow("""
class: GalaxyWorkflow
inputs:
input_datasets: collection
@@ -1058,8 +1069,7 @@ test_data:
- identifier: success
value: 1.fastq
type: File
""", history_id=history_id, assert_ok=False, wait=False)
self.wait_for_invocation_and_jobs(history_id, job_summary.workflow_id, job_summary.invocation_id, assert_ok=False)
""", history_id=history_id, assert_ok=False, wait=True)
history_contents = self.dataset_populator._get_contents_request(history_id=history_id).json()
first_input = history_contents[1]
assert first_input['history_content_type'] == 'dataset'
@@ -1087,7 +1097,7 @@ test_data:
def test_workflow_resume_with_mapped_over_collection_input(self):
# Test that replacement and resume also works if the failed job re-run works on a input DCE
with self.dataset_populator.test_history() as history_id:
job_summary = self._run_jobs("""
job_summary = self._run_workflow("""
class: GalaxyWorkflow
inputs:
input_collection: collection
@@ -1374,17 +1384,18 @@ test_data:
type: File
""", history_id=history_id)
def test_run_subworkflow_simple(self):
def test_run_subworkflow_simple(self) -> None:
with self.dataset_populator.test_history() as history_id:
run_response = self._run_jobs(WORKFLOW_NESTED_SIMPLE, test_data="""
summary = self._run_workflow(WORKFLOW_NESTED_SIMPLE, test_data="""
outer_input:
value: 1.bed
type: File
""", history_id=history_id)
invocation_id = summary.invocation_id
content = self.dataset_populator.get_history_dataset_content(history_id)
self.assertEqual("chrX\t152691446\t152691471\tCCDS14735.1_cds_0_0_chrX_152691447_f\t0\t+\nchrX\t152691446\t152691471\tCCDS14735.1_cds_0_0_chrX_152691447_f\t0\t+\n", content)
steps = self.workflow_populator.get_invocation(run_response.invocation_id)['steps']
steps = self.workflow_populator.get_invocation(invocation_id)['steps']
assert sum(1 for step in steps if step['subworkflow_invocation_id'] is None) == 3
subworkflow_invocation_id = [step['subworkflow_invocation_id'] for step in steps if step['subworkflow_invocation_id']][0]
subworkflow_invocation = self.workflow_populator.get_invocation(subworkflow_invocation_id)
@@ -1392,7 +1403,7 @@ outer_input:
assert [step for step in subworkflow_invocation['steps'] if step['workflow_step_label'] == 'inner_input']
assert [step for step in subworkflow_invocation['steps'] if step['workflow_step_label'] == 'random_lines']
bco = self.workflow_populator.get_biocompute_object(run_response.invocation_id)
bco = self.workflow_populator.get_biocompute_object(invocation_id)
self.workflow_populator.validate_biocompute_object(bco)
@skip_without_tool("random_lines1")
@@ -1485,7 +1496,7 @@ test_data:
value: 1.bed
type: File
"""
job_summary = self._run_jobs(workflow_run_description, history_id=history_id, wait=False)
job_summary = self._run_workflow(workflow_run_description, history_id=history_id, wait=False)
uploaded_workflow_id, invocation_id = job_summary.workflow_id, job_summary.invocation_id
# Wait for at least one scheduling step.
@@ -1518,8 +1529,9 @@ test_data:
value: 1.bed
type: File
"""
job_summary = self._run_jobs(workflow_text, test_data=test_data, history_id=history_id)
assert len(job_summary.jobs) == 4, "4 jobs expected, got %d jobs" % len(job_summary.jobs)
summary = self._run_workflow(workflow_text, test_data=test_data, history_id=history_id)
jobs = summary.jobs
assert len(jobs) == 4, "4 jobs expected, got %d jobs" % len(jobs)
content = self.dataset_populator.get_history_dataset_content(history_id)
self.assertEqual(
@@ -1632,7 +1644,7 @@ input_1:
type: File
"""
with self.dataset_populator.test_history() as history_id:
summary = self._run_jobs("""
summary = self._run_workflow("""
class: GalaxyWorkflow
inputs:
input_1: data
@@ -1659,7 +1671,7 @@ steps:
@skip_without_tool("cat")
def test_workflow_invocation_report_custom(self):
with self.dataset_populator.test_history() as history_id:
summary = self._run_jobs(
summary = self._run_workflow(
WORKFLOW_WITH_CUSTOM_REPORT_1,
test_data=WORKFLOW_WITH_CUSTOM_REPORT_1_TEST_DATA,
history_id=history_id
@@ -1683,7 +1695,7 @@ steps:
@skip_without_tool("cat1")
def test_export_invocation_bco(self):
with self.dataset_populator.test_history() as history_id:
summary = self._run_jobs(WORKFLOW_SIMPLE, test_data={"input1": "hello world"}, history_id=history_id)
summary = self._run_workflow(WORKFLOW_SIMPLE, test_data={"input1": "hello world"}, history_id=history_id)
invocation_id = summary.invocation_id
bco = self.workflow_populator.get_biocompute_object(invocation_id)
self.workflow_populator.validate_biocompute_object(bco)
@@ -1692,13 +1704,13 @@ steps:
@skip_without_tool("__APPLY_RULES__")
def test_workflow_run_apply_rules(self):
with self.dataset_populator.test_history() as history_id:
self._run_jobs(WORKFLOW_WITH_RULES_1, history_id=history_id, wait=True, assert_ok=True, round_trip_format_conversion=True)
self._run_workflow(WORKFLOW_WITH_RULES_1, history_id=history_id, wait=True, assert_ok=True, round_trip_format_conversion=True)
output_content = self.dataset_populator.get_history_collection_details(history_id, hid=6)
rules_test_data.check_example_2(output_content, self.dataset_populator)
def test_filter_failed_mapping(self):
with self.dataset_populator.test_history() as history_id:
summary = self._run_jobs("""
summary = self._run_workflow("""
class: GalaxyWorkflow
inputs:
input_c: collection
@@ -1750,7 +1762,7 @@ input_c:
def test_workflow_output_dataset(self):
with self.dataset_populator.test_history() as history_id:
summary = self._run_jobs(WORKFLOW_SIMPLE, test_data={"input1": "hello world"}, history_id=history_id)
summary = self._run_workflow(WORKFLOW_SIMPLE, test_data={"input1": "hello world"}, history_id=history_id)
workflow_id = summary.workflow_id
invocation_id = summary.invocation_id
invocation_response = self._get(f"workflows/{workflow_id}/invocations/{invocation_id}")
@@ -1765,7 +1777,25 @@ input_c:
@skip_without_tool("cat")
def test_workflow_output_dataset_collection(self):
with self.dataset_populator.test_history() as history_id:
summary = self._run_jobs("""
summary = self._run_workflow_with_output_collections(history_id)
workflow_id = summary.workflow_id
invocation_id = summary.invocation_id
invocation_response = self._get(f"workflows/{workflow_id}/invocations/{invocation_id}")
self._assert_status_code_is(invocation_response, 200)
invocation = invocation_response.json()
self._assert_has_keys(invocation, "id", "outputs", "output_collections")
assert len(invocation["output_collections"]) == 1
assert len(invocation["outputs"]) == 0
output_content = self.dataset_populator.get_history_collection_details(history_id, content_id=invocation["output_collections"]["wf_output_1"]["id"])
self._assert_has_keys(output_content, "id", "elements")
assert output_content["collection_type"] == "list"
elements = output_content["elements"]
assert len(elements) == 1
elements0 = elements[0]
assert elements0["element_identifier"] == "el1"
def _run_workflow_with_output_collections(self, history_id) -> RunJobsSummary:
summary = self._run_workflow("""
class: GalaxyWorkflow
inputs:
input1:
@@ -1788,25 +1818,10 @@ input1:
value: 1.fastq
type: File
""", history_id=history_id, round_trip_format_conversion=True)
workflow_id = summary.workflow_id
invocation_id = summary.invocation_id
invocation_response = self._get(f"workflows/{workflow_id}/invocations/{invocation_id}")
self._assert_status_code_is(invocation_response, 200)
invocation = invocation_response.json()
self._assert_has_keys(invocation, "id", "outputs", "output_collections")
assert len(invocation["output_collections"]) == 1
assert len(invocation["outputs"]) == 0
output_content = self.dataset_populator.get_history_collection_details(history_id, content_id=invocation["output_collections"]["wf_output_1"]["id"])
self._assert_has_keys(output_content, "id", "elements")
assert output_content["collection_type"] == "list"
elements = output_content["elements"]
assert len(elements) == 1
elements0 = elements[0]
assert elements0["element_identifier"] == "el1"
return summary
def test_workflow_input_as_output(self):
with self.dataset_populator.test_history() as history_id:
summary = self._run_jobs("""
def _run_workflow_with_inputs_as_outputs(self, history_id) -> RunJobsSummary:
summary = self._run_workflow("""
class: GalaxyWorkflow
inputs:
input1: data
@@ -1818,6 +1833,11 @@ outputs:
outputSource: text_input
steps: []
""", test_data={"input1": "hello world", "text_input": {"value": "A text variable", "type": "raw"}}, history_id=history_id)
return summary
def test_workflow_input_as_output(self):
with self.dataset_populator.test_history() as history_id:
summary = self._run_workflow_with_inputs_as_outputs(history_id)
workflow_id = summary.workflow_id
invocation_id = summary.invocation_id
invocation_response = self._get(f"workflows/{workflow_id}/invocations/{invocation_id}")
@@ -1834,7 +1854,7 @@ steps: []
def test_subworkflow_output_as_output(self):
with self.dataset_populator.test_history() as history_id:
summary = self._run_jobs("""
summary = self._run_workflow("""
class: GalaxyWorkflow
inputs:
input1: data
@@ -1868,7 +1888,7 @@ steps:
@skip_without_tool("cat")
def test_workflow_input_mapping(self):
with self.dataset_populator.test_history() as history_id:
summary = self._run_jobs("""
summary = self._run_workflow("""
class: GalaxyWorkflow
inputs:
input1: data
@@ -1910,7 +1930,7 @@ input1:
@skip_without_tool("collection_creates_pair")
def test_workflow_run_input_mapping_with_output_collections(self):
with self.dataset_populator.test_history() as history_id:
summary = self._run_jobs("""
summary = self._run_workflow("""
class: GalaxyWorkflow
inputs:
text_input: data
@@ -1988,7 +2008,7 @@ outer_input:
value: 1.fastq
type: File
"""
summary = self._run_jobs(WORKFLOW_NESTED_SIMPLE, test_data=test_data, history_id=history_id)
summary = self._run_workflow(WORKFLOW_NESTED_SIMPLE, test_data=test_data, history_id=history_id)
workflow_id = summary.workflow_id
invocation_id = summary.invocation_id
invocation_response = self._get(f"workflows/{workflow_id}/invocations/{invocation_id}")
@@ -2016,7 +2036,7 @@ outer_input:
# evaluation. Testing rescheduling and propagating connections within a subworkflow
# is handled by the next test case.
with self.dataset_populator.test_history() as history_id:
self._run_jobs("""
self._run_workflow("""
class: GalaxyWorkflow
inputs:
outer_input: data
@@ -2218,7 +2238,7 @@ input1:
@skip_without_tool("random_lines1")
def test_change_datatype_collection_map_over(self):
with self.dataset_populator.test_history() as history_id:
jobs_summary = self._run_jobs("""
jobs_summary = self._run_workflow("""
class: GalaxyWorkflow
inputs:
text_input1: collection
@@ -2244,7 +2264,7 @@ text_input1:
@skip_without_tool("collection_type_source_map_over")
def test_mapping_and_subcollection_mapping(self):
with self.dataset_populator.test_history() as history_id:
jobs_summary = self._run_jobs("""
jobs_summary = self._run_workflow("""
class: GalaxyWorkflow
inputs:
text_input1: collection
@@ -2266,7 +2286,7 @@ text_input1:
@skip_without_tool("random_lines1")
def test_empty_list_reduction(self):
with self.dataset_populator.test_history() as history_id:
self._run_jobs("""
self._run_workflow("""
class: GalaxyWorkflow
inputs:
input1: data
@@ -2494,7 +2514,7 @@ steps:
def test_run_with_implicit_connection(self):
with self.dataset_populator.test_history() as history_id:
run_summary = self._run_jobs("""
run_summary = self._run_workflow("""
class: GalaxyWorkflow
inputs:
test_input: data
@@ -2541,7 +2561,7 @@ steps:
def test_run_with_optional_data_specified_to_multi_data(self):
with self.dataset_populator.test_history() as history_id:
self._run_jobs(WORKFLOW_OPTIONAL_TRUE_INPUT_DATA, test_data="""
self._run_workflow(WORKFLOW_OPTIONAL_TRUE_INPUT_DATA, test_data="""
input1:
value: 1.bed
type: File
@@ -2596,7 +2616,7 @@ input1:
def test_run_with_validated_parameter_connection_optional(self):
with self.dataset_populator.test_history() as history_id:
run_summary = self._run_jobs("""
self._run_workflow("""
class: GalaxyWorkflow
inputs:
text_input: text
@@ -2612,7 +2632,6 @@ text_input:
value: "abd"
type: raw
""", history_id=history_id, wait=True, round_trip_format_conversion=True)
self.wait_for_invocation_and_jobs(history_id, run_summary.workflow_id, run_summary.invocation_id)
jobs = self._history_jobs(history_id)
assert len(jobs) == 1
@@ -2629,7 +2648,7 @@ data_input:
assert '(int_input) is not optional' in str(e)
failed = True
assert failed
run_response = self._run_jobs(WORKFLOW_PARAMETER_INPUT_INTEGER_REQUIRED, test_data="""
run_response = self._run_workflow(WORKFLOW_PARAMETER_INPUT_INTEGER_REQUIRED, test_data="""
data_input:
value: 1.bed
type: File
@@ -2637,13 +2656,13 @@ int_input:
value: 1
type: raw
""", history_id=history_id, wait=True, assert_ok=True)
self.dataset_populator.wait_for_history(history_id, assert_ok=True)
# self.dataset_populator.wait_for_history(history_id, assert_ok=True)
content = self.dataset_populator.get_history_dataset_content(history_id)
assert len(content.splitlines()) == 1, content
invocation = self.workflow_populator.get_invocation(run_response.invocation_id)
assert invocation['input_step_parameters']['int_input']['parameter_value'] == 1
run_response = self._run_jobs(WORKFLOW_PARAMETER_INPUT_INTEGER_OPTIONAL, test_data="""
run_response = self._run_workflow(WORKFLOW_PARAMETER_INPUT_INTEGER_OPTIONAL, test_data="""
data_input:
value: 1.bed
type: File
@@ -2656,7 +2675,7 @@ data_input:
with self.dataset_populator.test_history() as history_id:
workflow = self.workflow_populator.load_workflow_from_resource("test_subworkflow_with_integer_input")
workflow_id = self.workflow_populator.create_workflow(workflow)
hda = self.dataset_populator.new_dataset(history_id, content="1 2 3")
hda: dict = self.dataset_populator.new_dataset(history_id, content="1 2 3")
workflow_request = {
'history_id': history_id,
'inputs_by': 'name',
@@ -2885,7 +2904,7 @@ outer_input:
value: 1.bed
type: File
"""
run_jobs_summary = self._run_jobs(WORKFLOW_NESTED_SIMPLE, test_data=test_data, history_id=history_id_one)
run_jobs_summary = self._run_workflow(WORKFLOW_NESTED_SIMPLE, test_data=test_data, history_id=history_id_one)
workflow_id = run_jobs_summary.workflow_id
workflow_request = run_jobs_summary.workflow_request
# We copy the inputs to a new history and re-run the workflow
@@ -2974,7 +2993,7 @@ outer_input:
def test_empty_create(self):
response = self._post("workflows")
self._assert_status_code_is(response, 400)
self._assert_error_code_is(response, error_codes.USER_REQUEST_MISSING_PARAMETER)
self._assert_error_code_is(response, error_codes.error_codes_by_name["USER_REQUEST_MISSING_PARAMETER"])
def test_invalid_create_multiple_types(self):
data = {
@@ -2983,7 +3002,7 @@ outer_input:
}
response = self._post("workflows", data)
self._assert_status_code_is(response, 400)
self._assert_error_code_is(response, error_codes.USER_REQUEST_INVALID_PARAMETER)
self._assert_error_code_is(response, error_codes.error_codes_by_name["USER_REQUEST_INVALID_PARAMETER"])
@skip_without_tool("cat1")
def test_run_with_pja(self):
@@ -2999,7 +3018,7 @@ outer_input:
@skip_without_tool("hidden_param")
def test_hidden_param_in_workflow(self):
with self.dataset_populator.test_history() as history_id:
run_object = self._run_jobs("""
run_object = self._run_workflow("""
class: GalaxyWorkflow
steps:
step1:
@@ -3016,7 +3035,7 @@ steps:
@skip_without_tool("output_filter")
def test_optional_workflow_output(self):
with self.dataset_populator.test_history() as history_id:
run_object = self._run_jobs("""
run_object = self._run_workflow("""
class: GalaxyWorkflow
inputs: []
outputs:
@@ -3046,7 +3065,7 @@ input1:
- identifier: A
content: A
"""
run_object = self._run_jobs("""
run_object = self._run_workflow("""
class: GalaxyWorkflow
inputs:
input1:
@@ -4144,7 +4163,7 @@ input:
@skip_without_tool("random_lines1")
def test_run_replace_params_over_default_delayed(self):
with self.dataset_populator.test_history() as history_id:
run_summary = self._run_jobs("""
run_summary = self._run_workflow("""
class: GalaxyWorkflow
inputs:
input: data
@@ -4233,7 +4252,7 @@ input:
usage_details = self._invocation_details(workflow_id, invocation_id)
invocation_steps = usage_details["steps"]
invocation_input_step, invocation_tool_step = None, None
invocation_input_step, invocation_tool_step = {}, {}
for invocation_step in invocation_steps:
self._assert_has_keys(invocation_step, "workflow_step_id", "order_index", "id")
order_index = invocation_step["order_index"]
@@ -4275,6 +4294,9 @@ input:
assert invocation_tool_step is None
invocation_tool_step = invocation_step
assert invocation_input_step
assert invocation_tool_step
# Tool steps have non-null job_ids (deprecated though they may be)
assert invocation_input_step.get("job_id", None) is None
assert invocation_tool_step.get("job_id", None) is None
@@ -4302,7 +4324,7 @@ input:
def _run_mapping_workflow(self):
history_id = self.dataset_populator.new_history()
summary = self._run_jobs("""
summary = self._run_workflow("""
class: GalaxyWorkflow
inputs:
input_c: collection
@@ -4336,8 +4358,8 @@ input_c:
self._assert_status_code_is(response, 200)
assert len(response.json()) == 0
run_workflow_response = self.workflow_populator.invoke_workflow_raw(workflow_id, workflow_request, assert_ok=True)
run_workflow_response = run_workflow_response.json()
invocation_id = run_workflow_response['id']
run_workflow_dict = run_workflow_response.json()
invocation_id = run_workflow_dict['id']
usage_details_response = self._get(f"workflows/{other_id}/usage/{invocation_id}")
self._assert_status_code_is(usage_details_response, 200)
@@ -4350,8 +4372,8 @@ input_c:
self._assert_status_code_is(response, 200)
assert len(response.json()) == 0
run_workflow_response = self.workflow_populator.invoke_workflow_raw(workflow_id, workflow_request, assert_ok=True)
run_workflow_response = run_workflow_response.json()
invocation_id = run_workflow_response['id']
run_workflow_dict = run_workflow_response.json()
invocation_id = run_workflow_dict['id']
usage_details_response = self._get(f"workflows/{workflow_id}/usage/{invocation_id}")
self._assert_status_code_is(usage_details_response, 200)
@@ -4363,8 +4385,8 @@ input_c:
self._assert_status_code_is(response, 200)
assert len(response.json()) == 0
run_workflow_response = self.workflow_populator.invoke_workflow_raw(workflow_id, workflow_request, assert_ok=True)
run_workflow_response = run_workflow_response.json()
invocation_id = run_workflow_response['id']
run_workflow_dict = run_workflow_response.json()
invocation_id = run_workflow_dict['id']
with self._different_user():
usage_details_response = self._get(f"workflows/{workflow_id}/usage/{invocation_id}")
self._assert_status_code_is(usage_details_response, 403)
@@ -4386,13 +4408,13 @@ input_c:
f.write(WORKFLOW_NESTED_REPLACEMENT_PARAMETER)
import_response = self.workflow_populator.import_workflow_from_path_raw(workflow_path)
self._assert_status_code_is(import_response, 403)
self._assert_error_code_is(import_response, error_codes.ADMIN_REQUIRED)
self._assert_error_code_is(import_response, error_codes.error_codes_by_name["ADMIN_REQUIRED"])
path_as_uri = f"file://{workflow_path}"
import_data = dict(archive_source=path_as_uri)
import_response = self._post("workflows", data=import_data)
self._assert_status_code_is(import_response, 403)
self._assert_error_code_is(import_response, error_codes.ADMIN_REQUIRED)
self._assert_error_code_is(import_response, error_codes.error_codes_by_name["ADMIN_REQUIRED"])
def _invoke_paused_workflow(self, history_id):
workflow = self.workflow_populator.load_workflow_from_resource("test_workflow_pause")
+24 -12
View File
@@ -44,11 +44,16 @@ import random
import string
import unittest
from abc import ABCMeta, abstractmethod
from collections import namedtuple
from functools import wraps
from io import StringIO
from operator import itemgetter
from typing import Any, Callable, Dict, Optional
from typing import (
Any,
Callable,
Dict,
NamedTuple,
Optional,
)
import requests
import yaml
@@ -236,7 +241,7 @@ class BaseDatasetPopulator(BasePopulator):
Galaxy - implementations must implement _get, _post and _delete.
"""
def new_dataset(self, history_id: str, content=None, wait: bool = False, **kwds) -> str:
def new_dataset(self, history_id: str, content=None, wait: bool = False, **kwds) -> dict:
"""Create a new history dataset instance (HDA) and return its ID.
:returns: the HDA id of the new object
@@ -1101,7 +1106,14 @@ class BaseWorkflowPopulator(BasePopulator):
print(json.dumps(raw_workflow, sort_keys=True, indent=2))
RunJobsSummary = namedtuple('RunJobsSummary', ['history_id', 'workflow_id', 'invocation_id', 'inputs', 'jobs', 'invocation', 'workflow_request'])
class RunJobsSummary(NamedTuple):
history_id: str
workflow_id: str
invocation_id: str
inputs: dict
jobs: list
invocation: dict
workflow_request: dict
class WorkflowPopulator(GalaxyInteractorHttpMixin, BaseWorkflowPopulator, ImporterGalaxyInterface):
@@ -1113,7 +1125,7 @@ class WorkflowPopulator(GalaxyInteractorHttpMixin, BaseWorkflowPopulator, Import
# Required for ImporterGalaxyInterface interface - so we can recursively import
# nested workflows.
def import_workflow(self, workflow, **kwds):
def import_workflow(self, workflow, **kwds) -> Dict[str, Any]:
workflow_str = json.dumps(workflow, indent=4)
data = {
'workflow': workflow_str,
@@ -1123,7 +1135,7 @@ class WorkflowPopulator(GalaxyInteractorHttpMixin, BaseWorkflowPopulator, Import
assert upload_response.status_code == 200, upload_response.content
return upload_response.json()
def import_tool(self, tool):
def import_tool(self, tool) -> Dict[str, Any]:
""" Import a workflow via POST /api/workflows or
comparable interface into Galaxy.
"""
@@ -1131,7 +1143,7 @@ class WorkflowPopulator(GalaxyInteractorHttpMixin, BaseWorkflowPopulator, Import
assert upload_response.status_code == 200, upload_response
return upload_response.json()
def _import_tool_response(self, tool):
def _import_tool_response(self, tool) -> Response:
tool_str = json.dumps(tool, indent=4)
data = {
'representation': tool_str
@@ -1144,7 +1156,7 @@ class WorkflowPopulator(GalaxyInteractorHttpMixin, BaseWorkflowPopulator, Import
has_workflow = yaml.dump(workflow_dict)
return has_workflow
def _scale_workflow_dict(self, workflow_type="simple", **kwd):
def _scale_workflow_dict(self, workflow_type="simple", **kwd) -> Dict[str, Any]:
if workflow_type == "two_outputs":
return self._scale_workflow_dict_two_outputs(**kwd)
elif workflow_type == "wave_simple":
@@ -1152,7 +1164,7 @@ class WorkflowPopulator(GalaxyInteractorHttpMixin, BaseWorkflowPopulator, Import
else:
return self._scale_workflow_dict_simple(**kwd)
def _scale_workflow_dict_simple(self, **kwd):
def _scale_workflow_dict_simple(self, **kwd) -> Dict[str, Any]:
collection_size = kwd.get("collection_size", 2)
workflow_depth = kwd.get("workflow_depth", 3)
@@ -1174,7 +1186,7 @@ class WorkflowPopulator(GalaxyInteractorHttpMixin, BaseWorkflowPopulator, Import
}
return workflow_dict
def _scale_workflow_dict_two_outputs(self, **kwd):
def _scale_workflow_dict_two_outputs(self, **kwd) -> Dict[str, Any]:
collection_size = kwd.get("collection_size", 10)
workflow_depth = kwd.get("workflow_depth", 10)
@@ -1196,7 +1208,7 @@ class WorkflowPopulator(GalaxyInteractorHttpMixin, BaseWorkflowPopulator, Import
}
return workflow_dict
def _scale_workflow_dict_wave(self, **kwd):
def _scale_workflow_dict_wave(self, **kwd) -> Dict[str, Any]:
collection_size = kwd.get("collection_size", 10)
workflow_depth = kwd.get("workflow_depth", 10)
@@ -1222,7 +1234,7 @@ class WorkflowPopulator(GalaxyInteractorHttpMixin, BaseWorkflowPopulator, Import
return workflow_dict
@staticmethod
def _link(link, output_name=None):
def _link(link: str, output_name: Optional[str] = None) -> Dict[str, Any]:
if output_name is not None:
link = f"{str(link)}/{output_name}"
return {"$link": link}