mirror of
https://github.com/galaxyproject/galaxy.git
synced 2026-09-24 16:30:27 +08:00
Merge branch 'release_20.09' into dev
This commit is contained in:
@@ -3194,8 +3194,8 @@
|
||||
associated with scheduling workflows at the expense of increased
|
||||
total DB traffic because model objects are expunged from the SQL
|
||||
alchemy session between workflow invocation scheduling iterations.
|
||||
Set to -1 to disable any such maximum (the default).
|
||||
:Default: ``-1``
|
||||
Set to -1 to disable any such maximum.
|
||||
:Default: ``1000``
|
||||
:Type: int
|
||||
|
||||
|
||||
|
||||
@@ -1577,8 +1577,8 @@ galaxy:
|
||||
# scheduling workflows at the expense of increased total DB traffic
|
||||
# because model objects are expunged from the SQL alchemy session
|
||||
# between workflow invocation scheduling iterations. Set to -1 to
|
||||
# disable any such maximum (the default).
|
||||
#maximum_workflow_jobs_per_scheduling_iteration: -1
|
||||
# disable any such maximum.
|
||||
#maximum_workflow_jobs_per_scheduling_iteration: 1000
|
||||
|
||||
# Force serial scheduling of workflows within the context of a
|
||||
# particular history
|
||||
|
||||
@@ -42,6 +42,8 @@ class KubernetesJobRunner(AsynchronousJobRunner):
|
||||
"""
|
||||
runner_name = "KubernetesRunner"
|
||||
|
||||
LABEL_REGEX = re.compile("[^-A-Za-z0-9_.]")
|
||||
|
||||
def __init__(self, app, nworkers, **kwargs):
|
||||
# Check if pykube was importable, fail if not
|
||||
ensure_pykube()
|
||||
@@ -206,19 +208,19 @@ class KubernetesJobRunner(AsynchronousJobRunner):
|
||||
k8s_spec_template = {
|
||||
"metadata": {
|
||||
"labels": {
|
||||
"app.kubernetes.io/name": ajs.job_wrapper.tool.old_id,
|
||||
"app.kubernetes.io/name": self.LABEL_REGEX.sub("_", ajs.job_wrapper.tool.old_id),
|
||||
"app.kubernetes.io/instance": self.__produce_k8s_job_prefix(),
|
||||
"app.kubernetes.io/version": ajs.job_wrapper.tool.version,
|
||||
"app.kubernetes.io/version": self.LABEL_REGEX.sub("_", str(ajs.job_wrapper.tool.version)),
|
||||
"app.kubernetes.io/component": "tool",
|
||||
"app.kubernetes.io/part-of": "galaxy",
|
||||
"app.kubernetes.io/managed-by": "galaxy",
|
||||
"app.galaxyproject.org/job_id": ajs.job_wrapper.get_id_tag(),
|
||||
"app.galaxyproject.org/instance": self._galaxy_instance_id or "",
|
||||
"app.galaxyproject.org/handler": self.app.config.server_name,
|
||||
"app.galaxyproject.org/destination": ajs.job_wrapper.job_destination.id,
|
||||
"app.galaxyproject.org/job_id": self.LABEL_REGEX.sub("_", ajs.job_wrapper.get_id_tag()),
|
||||
"app.galaxyproject.org/handler": self.LABEL_REGEX.sub("_", self.app.config.server_name),
|
||||
"app.galaxyproject.org/destination": self.LABEL_REGEX.sub(
|
||||
"_", str(ajs.job_wrapper.job_destination.id))
|
||||
},
|
||||
"annotations": {
|
||||
"app.galaxyproject.org/tool_id": ajs.job_wrapper.tool.id,
|
||||
"app.galaxyproject.org/tool_id": ajs.job_wrapper.tool.id
|
||||
}
|
||||
},
|
||||
"spec": {
|
||||
|
||||
@@ -513,7 +513,8 @@ class XmlToolSource(ToolSource):
|
||||
|
||||
if tests_elem is not None:
|
||||
for i, test_elem in enumerate(tests_elem.findall("test")):
|
||||
tests.append(_test_elem_to_dict(test_elem, i))
|
||||
profile = self.parse_profile()
|
||||
tests.append(_test_elem_to_dict(test_elem, i, profile))
|
||||
|
||||
return rval
|
||||
|
||||
@@ -533,10 +534,10 @@ class XmlToolSource(ToolSource):
|
||||
return python_template_version
|
||||
|
||||
|
||||
def _test_elem_to_dict(test_elem, i):
|
||||
def _test_elem_to_dict(test_elem, i, profile=None):
|
||||
rval = dict(
|
||||
outputs=__parse_output_elems(test_elem),
|
||||
output_collections=__parse_output_collection_elems(test_elem),
|
||||
output_collections=__parse_output_collection_elems(test_elem, profile=profile),
|
||||
inputs=__parse_input_elems(test_elem, i),
|
||||
expect_num_outputs=test_elem.get("expect_num_outputs"),
|
||||
command=__parse_assert_list_from_elem(test_elem.find("assert_command")),
|
||||
@@ -579,36 +580,38 @@ def __parse_command_elem(test_elem):
|
||||
return __parse_assert_list_from_elem(assert_elem)
|
||||
|
||||
|
||||
def __parse_output_collection_elems(test_elem):
|
||||
def __parse_output_collection_elems(test_elem, profile=None):
|
||||
output_collections = []
|
||||
for output_collection_elem in test_elem.findall("output_collection"):
|
||||
output_collection_def = __parse_output_collection_elem(output_collection_elem)
|
||||
output_collection_def = __parse_output_collection_elem(output_collection_elem, profile=profile)
|
||||
output_collections.append(output_collection_def)
|
||||
return output_collections
|
||||
|
||||
|
||||
def __parse_output_collection_elem(output_collection_elem):
|
||||
def __parse_output_collection_elem(output_collection_elem, profile=None):
|
||||
attrib = dict(output_collection_elem.attrib)
|
||||
name = attrib.pop('name', None)
|
||||
if name is None:
|
||||
raise Exception("Test output collection does not have a 'name'")
|
||||
element_tests = __parse_element_tests(output_collection_elem)
|
||||
element_tests = __parse_element_tests(output_collection_elem, profile=profile)
|
||||
return TestCollectionOutputDef(name, attrib, element_tests).to_dict()
|
||||
|
||||
|
||||
def __parse_element_tests(parent_element):
|
||||
def __parse_element_tests(parent_element, profile=None):
|
||||
element_tests = {}
|
||||
for idx, element in enumerate(parent_element.findall("element")):
|
||||
element_attrib = dict(element.attrib)
|
||||
identifier = element_attrib.pop('name', None)
|
||||
if identifier is None:
|
||||
raise Exception("Test primary dataset does not have a 'identifier'")
|
||||
element_tests[identifier] = __parse_test_attributes(element, element_attrib, parse_elements=True)
|
||||
element_tests[identifier][1]["element_index"] = idx
|
||||
element_tests[identifier] = __parse_test_attributes(element, element_attrib, parse_elements=True, profile=profile)
|
||||
if profile and profile >= "20.09":
|
||||
element_tests[identifier][1]["expected_sort_order"] = idx
|
||||
|
||||
return element_tests
|
||||
|
||||
|
||||
def __parse_test_attributes(output_elem, attrib, parse_elements=False, parse_discovered_datasets=False):
|
||||
def __parse_test_attributes(output_elem, attrib, parse_elements=False, parse_discovered_datasets=False, profile=None):
|
||||
assert_list = __parse_assert_list(output_elem)
|
||||
|
||||
# Allow either file or value to specify a target file to compare result with
|
||||
@@ -638,7 +641,7 @@ def __parse_test_attributes(output_elem, attrib, parse_elements=False, parse_dis
|
||||
checksum = attrib.get("checksum", None)
|
||||
element_tests = {}
|
||||
if parse_elements:
|
||||
element_tests = __parse_element_tests(output_elem)
|
||||
element_tests = __parse_element_tests(output_elem, profile=profile)
|
||||
|
||||
primary_datasets = {}
|
||||
if parse_discovered_datasets:
|
||||
|
||||
@@ -724,36 +724,29 @@ def verify_collection(output_collection_def, data_collection, verify_dataset):
|
||||
message = template % (name, expected_element_count, actual_element_count)
|
||||
raise AssertionError(message)
|
||||
|
||||
def get_element(elements, id):
|
||||
for element in elements:
|
||||
if element["element_identifier"] == id:
|
||||
return element
|
||||
return False
|
||||
|
||||
def verify_elements(element_objects, element_tests):
|
||||
sorted_test_ids = [None] * len(element_tests)
|
||||
# sorted_test_ids = [None] * len(element_tests)
|
||||
expected_sort_order = []
|
||||
|
||||
eo_ids = [_["element_identifier"] for _ in element_objects]
|
||||
for element_identifier, element_test in element_tests.items():
|
||||
if isinstance(element_test, dict):
|
||||
element_outfile, element_attrib = None, element_test
|
||||
else:
|
||||
element_outfile, element_attrib = element_test
|
||||
sorted_test_ids[element_attrib["element_index"]] = element_identifier
|
||||
if 'expected_sort_order' in element_attrib:
|
||||
expected_sort_order.append(element_identifier)
|
||||
|
||||
i = 0
|
||||
for element_identifier in sorted_test_ids:
|
||||
element_test = element_tests[element_identifier]
|
||||
if isinstance(element_test, dict):
|
||||
element_outfile, element_attrib = None, element_test
|
||||
else:
|
||||
element_outfile, element_attrib = element_test
|
||||
|
||||
element = None
|
||||
while i < len(element_objects):
|
||||
if element_objects[i]["element_identifier"] == element_identifier:
|
||||
element = element_objects[i]
|
||||
i += 1
|
||||
break
|
||||
i += 1
|
||||
|
||||
if element is None:
|
||||
template = "Failed to find identifier '%s' of test collection %s in the tool generated collection elements %s (at the correct position)"
|
||||
eo_ids = [_["element_identifier"] for _ in element_objects]
|
||||
message = template % (element_identifier, sorted_test_ids,
|
||||
eo_ids)
|
||||
element = get_element(element_objects, element_identifier)
|
||||
if not element:
|
||||
template = "Failed to find identifier '%s' in the tool generated collection elements %s"
|
||||
message = template % (element_identifier, eo_ids)
|
||||
raise AssertionError(message)
|
||||
|
||||
element_type = element["element_type"]
|
||||
@@ -763,6 +756,21 @@ def verify_collection(output_collection_def, data_collection, verify_dataset):
|
||||
elements = element["object"]["elements"]
|
||||
verify_elements(elements, element_attrib.get("elements", {}))
|
||||
|
||||
if len(expected_sort_order) > 0:
|
||||
i = 0
|
||||
for element_identifier in expected_sort_order:
|
||||
element = None
|
||||
while i < len(element_objects):
|
||||
if element_objects[i]["element_identifier"] == element_identifier:
|
||||
element = element_objects[i]
|
||||
i += 1
|
||||
break
|
||||
i += 1
|
||||
if element is None:
|
||||
template = "Collection identifier '%s' found out of order, expected order of %s for the tool generated collection elements %s"
|
||||
message = template % (element_identifier, expected_sort_order, eo_ids)
|
||||
raise AssertionError(message)
|
||||
|
||||
verify_elements(data_collection["elements"], output_collection_def.element_tests)
|
||||
|
||||
|
||||
|
||||
@@ -1507,7 +1507,7 @@ Note that this tool uses ``assign_primary_output="true"`` for ``<discover_data_s
|
||||
<xs:annotation>
|
||||
<xs:documentation xml:lang="en"><
|
||||
|
||||
@@ -98,11 +98,13 @@ def execute(trans, tool, mapping_params, history, rerun_remap_job_id=None, colle
|
||||
break
|
||||
else:
|
||||
execute_single_job(execution_slice, completed_jobs[i])
|
||||
history = execution_slice.history or history
|
||||
jobs_executed += 1
|
||||
if execution_slice.datasets_to_persist:
|
||||
datasets_to_persist.extend(execution_slice.datasets_to_persist)
|
||||
|
||||
if datasets_to_persist:
|
||||
execution_slice.history.add_datasets(trans.sa_session, datasets_to_persist, set_hid=True, quota=False, flush=False)
|
||||
history.add_datasets(trans.sa_session, datasets_to_persist, set_hid=True, quota=False, flush=False)
|
||||
# a side effect of history.add_datasets is a commit within db_next_hid (even with flush=False).
|
||||
else:
|
||||
# Make sure collections, implicit jobs etc are flushed even if there are no precreated output datasets
|
||||
|
||||
@@ -2368,7 +2368,7 @@ mapping:
|
||||
|
||||
maximum_workflow_jobs_per_scheduling_iteration:
|
||||
type: int
|
||||
default: -1
|
||||
default: 1000
|
||||
required: false
|
||||
desc: |
|
||||
Specify a maximum number of jobs that any given workflow scheduling iteration can create.
|
||||
@@ -2376,7 +2376,7 @@ mapping:
|
||||
preventing other jobs from executing. This may also mitigate memory issues associated with
|
||||
scheduling workflows at the expense of increased total DB traffic because model objects
|
||||
are expunged from the SQL alchemy session between workflow invocation scheduling iterations.
|
||||
Set to -1 to disable any such maximum (the default).
|
||||
Set to -1 to disable any such maximum.
|
||||
|
||||
history_local_serial_workflow_scheduling:
|
||||
type: bool
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
<tool id="discover_sort_by" name="discover_sort_by" version="0.1.0">
|
||||
<tool id="discover_sort_by" name="discover_sort_by" version="0.1.0" profile="20.09">
|
||||
<command><![CDATA[
|
||||
for i in \$(seq 1 10);
|
||||
do
|
||||
|
||||
@@ -0,0 +1,42 @@
|
||||
<tool id="discover_sort_by_legacy_test" name="discover_sort_by_legacy_test" version="0.1.0" profile="20.05">
|
||||
<!-- same as discover_sort_by but with tests not ordered correctly, this would cause
|
||||
failure in profile 20.09 or newer.
|
||||
-->
|
||||
<command><![CDATA[
|
||||
for i in \$(seq 1 10);
|
||||
do
|
||||
echo "\$i" > \$i.txt;
|
||||
done
|
||||
]]></command>
|
||||
<inputs/>
|
||||
<outputs>
|
||||
<collection name="collection_numeric_name" type="list" label="num">
|
||||
<discover_datasets pattern="__name_and_ext__" sort_by="numeric_name"/>
|
||||
</collection>
|
||||
<collection name="collection_rev_numeric_name" type="list" label="num rev">
|
||||
<discover_datasets pattern="__name_and_ext__" sort_by="reverse_numeric_name"/>
|
||||
</collection>
|
||||
<collection name="collection_lexical_name" type="list" label="num">
|
||||
<discover_datasets pattern="__name_and_ext__" sort_by="lexical_name" />
|
||||
</collection>
|
||||
<data name="data_reverse_lexical_name">
|
||||
<discover_datasets pattern="__name_and_ext__" format="txt" assign_primary_output="true" sort_by="reverse_lexical_name" visible="true"/>
|
||||
</data>
|
||||
</outputs>
|
||||
<tests>
|
||||
<test expect_num_outputs="4">
|
||||
<param name="input1" value="tinywga.fam" />
|
||||
<output_collection name="collection_numeric_name" type="list" count="10">
|
||||
<element name="2">
|
||||
<assert_contents><has_text_matching expression="^.*$"/></assert_contents>
|
||||
</element>
|
||||
<element name="1">
|
||||
<assert_contents><has_text_matching expression="^.*$"/></assert_contents>
|
||||
</element>
|
||||
<element name="10">
|
||||
<assert_contents><has_text_matching expression="^.*$"/></assert_contents>
|
||||
</element>
|
||||
</output_collection>
|
||||
</test>
|
||||
</tests>
|
||||
</tool>
|
||||
@@ -24,7 +24,7 @@ class MaximumWorkflowInvocationDurationTestCase(integration_util.IntegrationTest
|
||||
def handle_galaxy_config_kwds(cls, config):
|
||||
config["maximum_workflow_invocation_duration"] = 20
|
||||
|
||||
def do_test(self):
|
||||
def test(self):
|
||||
workflow = self.workflow_populator.load_workflow_from_resource("test_workflow_pause")
|
||||
workflow_id = self.workflow_populator.create_workflow(workflow)
|
||||
history_id = self.dataset_populator.new_history()
|
||||
@@ -61,7 +61,7 @@ class MaximumWorkflowJobsPerSchedulingIterationTestCase(integration_util.Integra
|
||||
def handle_galaxy_config_kwds(cls, config):
|
||||
config["maximum_workflow_jobs_per_scheduling_iteration"] = 1
|
||||
|
||||
def do_test(self):
|
||||
def test(self):
|
||||
workflow_id = self.workflow_populator.upload_yaml_workflow("""
|
||||
class: GalaxyWorkflow
|
||||
steps:
|
||||
@@ -73,11 +73,11 @@ steps:
|
||||
- tool_id: collection_paired_test
|
||||
state:
|
||||
f1:
|
||||
$link: 1#paired_output
|
||||
$link: 1/paired_output
|
||||
- tool_id: cat_list
|
||||
state:
|
||||
input1:
|
||||
$link: 2#out1
|
||||
$link: 2/out1
|
||||
""")
|
||||
with self.dataset_populator.test_history() as history_id:
|
||||
hdca1 = self.dataset_collection_populator.create_list_in_history(history_id, contents=["a\nb\nc\nd\n", "e\nf\ng\nh\n"]).json()
|
||||
|
||||
Reference in New Issue
Block a user