Merge pull request #22912 from jmchilton/extract_bug_fixes

[26.1] Keep collection-operation outputs as workflow steps in extraction
This commit is contained in:
Marius van den Beek
2026-06-16 16:33:40 +02:00
committed by GitHub
3 changed files with 119 additions and 8 deletions
+10 -2
View File
@@ -713,13 +713,21 @@ def step_inputs_by_id(trans: ProvidesHistoryContext, job: Job) -> tuple[ToolInpu
def _original_hda(hda: HistoryDatasetAssociation) -> HistoryDatasetAssociation:
while hda.copied_from_history_dataset_association:
# Follow plain copies back to their source, but stop at anything with its own
# creating job: collection-operation tools (Extract Dataset, Filter, Relabel,
# ...) produce a copy *and* record a job, and those are real workflow steps -
# normalizing past them would attribute the output to its source and drop the
# operation from the extracted workflow.
while hda.copied_from_history_dataset_association and not hda.creating_job_associations:
hda = hda.copied_from_history_dataset_association
return hda
def _original_hdca(hdca: HistoryDatasetCollectionAssociation) -> HistoryDatasetCollectionAssociation:
while hdca.copied_from_history_dataset_collection_association:
# Same creating-job guard as _original_hda (see there for the rationale): a
# collection that records its own creating job - e.g. reimported with its job
# association - is a real step and must not normalize past copied_from.
while hdca.copied_from_history_dataset_collection_association and not hdca.creating_job_associations:
hdca = hdca.copied_from_history_dataset_collection_association
return hdca
@@ -11,6 +11,7 @@ from json import (
)
from typing import (
Any,
Optional,
TYPE_CHECKING,
)
@@ -29,6 +30,13 @@ if TYPE_CHECKING:
from requests import Response
def _connection_step_id(connection: Any) -> int:
# .ga format may yield a single dict or a list of one dict.
if isinstance(connection, list):
connection = connection[0]
return connection["id"]
class _ExtractionHelpersMixin:
"""Shared helpers for HID-based and ID-based workflow extraction tests."""
@@ -43,6 +51,53 @@ class _ExtractionHelpersMixin:
def _assert_status_code_is(self, response: "Response", expected_status_code: int) -> None: ...
def assert_steps_of_type(
self, workflow: dict[str, Any], step_type: str, expected_len: Optional[int] = None
) -> list[dict[str, Any]]: ...
def _setup_extract_dataset_then_cat(self, history_id):
"""Build a list, extract its first element, and feed the result to cat1.
The __EXTRACT_DATASET__ output is an HDA copied_from the source element
*and* carrying its own creating job - the shape that must stay a real
workflow step. Returns (input_hdca, extract_job_id, cat_job_id).
"""
hdca = self.dataset_collection_populator.create_list_in_history(
history_id, contents=["a\nb\n", "c\nd\n"], wait=True
).json()["outputs"][0]
extract_run = self.dataset_populator.run_tool(
tool_id="__EXTRACT_DATASET__",
inputs={"input": {"src": "hdca", "id": hdca["id"]}, "which|which_dataset": "first"},
history_id=history_id,
)
extract_job_id = extract_run["jobs"][0]["id"]
extracted = extract_run["outputs"][0]
self.dataset_populator.wait_for_history(history_id, assert_ok=True)
cat_run = self.dataset_populator.run_tool(
tool_id="cat1",
inputs={"input1": {"src": "hda", "id": extracted["id"]}},
history_id=history_id,
)
cat_job_id = cat_run["jobs"][0]["id"]
self.dataset_populator.wait_for_history(history_id, assert_ok=True)
return hdca, extract_job_id, cat_job_id
def _assert_extract_dataset_step_kept(self, downloaded):
"""Assert the Extract Dataset operation survived as its own tool step:
fed by the collection input and feeding the cat1 consumer. Normalizing
past copied_from would drop it and leave cat1 input-less.
"""
collection_step = self.assert_steps_of_type(downloaded, "data_collection_input", expected_len=1)[0]
tool_steps = self.assert_steps_of_type(downloaded, "tool", expected_len=2)
extract_step = next((s for s in tool_steps if s.get("tool_id") == "__EXTRACT_DATASET__"), None)
cat_step = next((s for s in tool_steps if s.get("tool_id") == "cat1"), None)
assert extract_step is not None, f"Extract Dataset step missing: {[s.get('tool_id') for s in tool_steps]}"
assert cat_step is not None, f"cat1 step missing: {[s.get('tool_id') for s in tool_steps]}"
extract_connections = extract_step["input_connections"]
cat_connections = cat_step["input_connections"]
assert _connection_step_id(extract_connections["input"]) == collection_step["id"], extract_connections
assert _connection_step_id(cat_connections["input1"]) == extract_step["id"], cat_connections
def _run_tool_get_collection_and_job_id(self, history_id, tool_id, inputs):
run = self.dataset_populator.run_tool(tool_id=tool_id, inputs=inputs, history_id=history_id)
implicit_hdca = run["implicit_collections"][0]
@@ -587,6 +642,24 @@ test_data:
tool_ids=tool_ids,
)
@skip_without_tool("__EXTRACT_DATASET__")
@skip_without_tool("cat1")
@summarize_instance_history_on_error
def test_extract_keeps_extract_dataset_operation_step(self, history_id):
"""Extract Dataset output carries copied_from (the source element) *and*
its own creating job. The summary must attribute the output to the
Extract Dataset job and keep it as a real step, not normalize past
copied_from back to the source element's creating job - which drops the
operation step and leaves the downstream consumer input-less.
"""
hdca, extract_job_id, cat_job_id = self._setup_extract_dataset_then_cat(history_id)
downloaded_workflow = self._extract_and_download_workflow(
history_id,
dataset_collection_ids=[hdca["hid"]],
job_ids=[extract_job_id, cat_job_id],
)
self._assert_extract_dataset_step_kept(downloaded_workflow)
def __run_random_lines_mapped_over_singleton(self, history_id):
hdca = self.dataset_collection_populator.create_list_in_history(history_id, contents=["1 2 3\n4 5 6"]).json()
hdca_id = hdca["id"]
@@ -745,15 +818,30 @@ class TestWorkflowExtractionByIdsApi(_ExtractionHelpersMixin, BaseWorkflowsApiTe
assert len(input_steps) == 1 and len(tool_steps) == 1
if expected_tool_id is not None:
assert tool_steps[0]["tool_id"] == expected_tool_id
connection = tool_steps[0]["input_connections"]["input1"]
# .ga format may yield a single dict or a list of one dict.
connection = connection[0] if isinstance(connection, list) else connection
assert connection["id"] == input_steps[0]["id"]
assert _connection_step_id(tool_steps[0]["input_connections"]["input1"]) == input_steps[0]["id"]
def _assert_extract_rejected(self, payload, allowed_codes):
response = self._post("workflows/extract", data=payload, json=True)
assert response.status_code in allowed_codes, response.text
@skip_without_tool("__EXTRACT_DATASET__")
@skip_without_tool("cat1")
@summarize_instance_history_on_error
def test_extract_keeps_extract_dataset_operation_step_by_ids(self, history_id):
"""ID-path sibling of the HID Extract Dataset operation-step test.
The by-ids closure normalizes copied_from symmetrically on both output
registration and input lookup, so this scenario already wires correctly
here - it is a regression guard ensuring the copied_from/creating-job
change keeps the Extract Dataset step connected, not a red->green proof.
"""
hdca, extract_job_id, cat_job_id = self._setup_extract_dataset_then_cat(history_id)
downloaded = self._extract_and_download_workflow_by_ids(
hdca_ids=[hdca["id"]],
job_ids=[extract_job_id, cat_job_id],
)
self._assert_extract_dataset_step_kept(downloaded)
@skip_without_tool("cat1")
@summarize_instance_history_on_error
def test_extract_with_hda_ids(self, history_id):
+17 -2
View File
@@ -42,10 +42,12 @@ class TestWorkflowExtractSummary(TestCase):
assert job_dict[hda3.job] == [("out3", hda3)]
def test_finds_original_job_if_copied(self):
# Passive copies (no creating job of their own) normalize back to the
# source, so the output is attributed to the source's creating job.
hda = MockHda()
derived_hda_1 = MockHda()
derived_hda_1 = MockHda(job=UNDEFINED_JOB)
derived_hda_1.copied_from_history_dataset_association = hda
derived_hda_2 = MockHda()
derived_hda_2 = MockHda(job=UNDEFINED_JOB)
derived_hda_2.copied_from_history_dataset_association = derived_hda_1
self.history.active_datasets.append(derived_hda_2)
job_dict, warnings = self._summarize()
@@ -53,6 +55,19 @@ class TestWorkflowExtractSummary(TestCase):
assert len(job_dict) == 1
assert job_dict[hda.job] == [("out1", derived_hda_2)]
def test_keeps_copy_with_own_creating_job(self):
# A copy that records its own creating job (e.g. Extract Dataset output)
# is a real step and must not normalize past copied_from to the source.
hda = MockHda()
derived_hda = MockHda()
derived_hda.copied_from_history_dataset_association = hda
self.history.active_datasets.append(derived_hda)
job_dict, warnings = self._summarize()
assert not warnings
assert len(job_dict) == 1
assert job_dict[derived_hda.job] == [("out1", derived_hda)]
assert hda.job not in job_dict
def test_fake_job_hda(self):
"""Fakes job if creating_job_associations is empty."""
hda = MockHda(job=UNDEFINED_JOB)