From 58e928f02a476986182801225aa1445a93648a1e Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Tue, 17 Mar 2026 12:53:23 +0100 Subject: [PATCH 01/13] Fix duplicate collection elements in metadata outputs_populated When re-populating dynamic output collections, existing elements were not cleared before adding new ones, causing duplicate entries to accumulate in metadata/outputs_populated/collections_attrs.txt and then propagate to the database. Two fixes: - In collect_dynamic_outputs: clear elements and reset element_count before re-populating a dynamic collection - In _import_collection_instances: delete existing DB elements before materializing new ones from import attrs (allow_edit path) --- lib/galaxy/job_execution/output_collect.py | 3 +++ lib/galaxy/model/store/__init__.py | 5 +++++ 2 files changed, 8 insertions(+) diff --git a/lib/galaxy/job_execution/output_collect.py b/lib/galaxy/job_execution/output_collect.py index bea31b54daf..703cdaae339 100644 --- a/lib/galaxy/job_execution/output_collect.py +++ b/lib/galaxy/job_execution/output_collect.py @@ -189,6 +189,9 @@ def collect_dynamic_outputs( # We are adding dynamic collections, which may be precreated, but their actually state is still new! collection.populated_state = collection.populated_states.NEW + # Clear any existing elements to avoid duplicates when re-populating + collection.elements.clear() + collection.element_count = None try: collection_builder = builder.BoundCollectionBuilder(collection) diff --git a/lib/galaxy/model/store/__init__.py b/lib/galaxy/model/store/__init__.py index 65d9ae780d3..22d314d4745 100644 --- a/lib/galaxy/model/store/__init__.py +++ b/lib/galaxy/model/store/__init__.py @@ -899,6 +899,11 @@ class ModelImportStore(metaclass=abc.ABCMeta): for attribute in attributes: if attribute in collection_attrs: setattr(dc, attribute, collection_attrs.get(attribute)) + # Clear existing elements to avoid duplicates when re-importing + if "elements" in collection_attrs: + for element in list(dc.elements): + self.sa_session.delete(element) + dc.elements.clear() materialize_elements(dc) else: # create collection From e51262fd00324f93ae525a8e3d5e063d25ad66b1 Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Tue, 17 Mar 2026 12:53:24 +0100 Subject: [PATCH 02/13] Fix mapped dynamic collection outputs missing with extended metadata MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Commit 4172b9b5 ("Fix AttributeError serializing implicit output collections during job prep") excluded all output_dataset_collections from io_dicts when exclude_implicit_outputs=True. This broke mapped dynamic collection outputs because their per-job DCs were no longer serialized to outputs_new, so set_metadata.py could not discover and populate them. Selectively include output_dataset_collections where the name is not in out_data — this excludes shared DCs for mapped dataset outputs (which have N precreated elements with uninitialized sentinels) while including per-job DCs for mapped collection outputs. Add integration tests verifying no duplicate collection elements in both non-mapped and mapped dynamic collection outputs with extended metadata. --- lib/galaxy/model/__init__.py | 17 +- ...st_extended_metadata_duplicate_elements.py | 210 ++++++++++++++++++ test/unit/data/model/test_model.py | 28 ++- 3 files changed, 246 insertions(+), 9 deletions(-) create mode 100644 test/integration/test_extended_metadata_duplicate_elements.py diff --git a/lib/galaxy/model/__init__.py b/lib/galaxy/model/__init__.py index 5f506daee13..3d53d6a401a 100644 --- a/lib/galaxy/model/__init__.py +++ b/lib/galaxy/model/__init__.py @@ -1839,8 +1839,21 @@ class Job(Base, JobLike, UsesCreateAndUpdateTime, Dictifiable, Serializable): if obj.name not in out_data: out_collections[obj.name] = obj.dataset_collection_instance # else this is a mapped over output - if not exclude_implicit_outputs: - out_collections.update([(obj.name, obj.dataset_collection) for obj in self.output_dataset_collections]) + if exclude_implicit_outputs: + # Include implicit output dataset collections only when they represent + # a tool's collection output (name not in out_data). Exclude shared DCs + # for mapped dataset outputs (name in out_data) which have N precreated + # elements where only the current job's element is initialized. + for implicit_obj in self.output_dataset_collections: + if implicit_obj.name not in out_data: + out_collections[implicit_obj.name] = implicit_obj.dataset_collection + else: + out_collections.update( + [ + (implicit_obj.name, implicit_obj.dataset_collection) + for implicit_obj in self.output_dataset_collections + ] + ) return IoDicts(inp_data, out_data, out_collections) # TODO: Add accessors for members defined in SQL Alchemy for the Job table and diff --git a/test/integration/test_extended_metadata_duplicate_elements.py b/test/integration/test_extended_metadata_duplicate_elements.py new file mode 100644 index 00000000000..82ad0981bc4 --- /dev/null +++ b/test/integration/test_extended_metadata_duplicate_elements.py @@ -0,0 +1,210 @@ +"""Integration test verifying that extended metadata does not create duplicate collection elements.""" + +from sqlalchemy import ( + func, + select, +) + +from galaxy.model import ( + DatasetCollectionElement, + HistoryDatasetCollectionAssociation, +) +from galaxy_test.base.populators import ( + DatasetCollectionPopulator, + DatasetPopulator, +) +from galaxy_test.driver.integration_util import IntegrationTestCase + + +class TestExtendedMetadataDuplicateElements(IntegrationTestCase): + dataset_populator: DatasetPopulator + framework_tool_and_types = True + + @classmethod + def handle_galaxy_config_kwds(cls, config): + super().handle_galaxy_config_kwds(config) + config["metadata_strategy"] = "extended" + config["retry_metadata_internally"] = False + + def setUp(self): + super().setUp() + self.dataset_populator = DatasetPopulator(self.galaxy_interactor) + self.dataset_collection_populator = DatasetCollectionPopulator(self.galaxy_interactor) + + def test_no_duplicate_elements_in_dynamic_list_output(self, history_id): + """Run a tool with dynamic collection output and verify no duplicate elements.""" + response = self.dataset_populator.run_tool( + "collection_creates_dynamic_list_of_pairs", + {"foo": "bar"}, + history_id, + ) + job_api_id = response["jobs"][0]["id"] + self.dataset_populator.wait_for_job(job_api_id, assert_ok=True) + + output_collections = response["output_collections"] + hdca_details = self.dataset_populator.get_history_collection_details( + history_id, content_id=output_collections[0]["id"] + ) + + sa_session = self._app.model.session + hdca_id = self._app.security.decode_id(hdca_details["id"]) + hdca = sa_session.get(HistoryDatasetCollectionAssociation, hdca_id) + assert hdca is not None + dc = hdca.collection + + # Verify exactly 3 outer elements (samp1, samp2, samp3) - no duplicates + outer_count = sa_session.scalar( + select(func.count()).where(DatasetCollectionElement.dataset_collection_id == dc.id) + ) + assert outer_count == 3, f"Expected 3 outer elements but found {outer_count}. Duplicate elements detected!" + + # Verify each inner pair has exactly 2 elements (forward, reverse) + for element in dc.elements: + assert element.child_collection is not None + inner_count = sa_session.scalar( + select(func.count()).where( + DatasetCollectionElement.dataset_collection_id == element.child_collection.id + ) + ) + assert inner_count == 2, ( + f"Expected 2 inner elements for '{element.element_identifier}' " + f"but found {inner_count}. Duplicate elements detected!" + ) + + def test_no_duplicate_elements_in_mapped_dynamic_collection(self, history_id): + """Map a tool with dynamic collection output over a list and verify no duplicate elements.""" + # Create a list of 2 tabular datasets + fetch_response = self.dataset_collection_populator.create_list_in_history( + history_id, + contents=["101\t1\n101\t2\n105\t3\n", "201\t10\n201\t20\n205\t30\n"], + ext="tabular", + wait=True, + ) + self.dataset_populator.wait_for_history(history_id, assert_ok=True) + hdca_id = fetch_response.json()["output_collections"][0]["id"] + + # Map collection_split_on_column over the list + inputs = { + "input1": {"batch": True, "values": [{"src": "hdca", "id": hdca_id}]}, + } + response = self.dataset_populator.run_tool( + "collection_split_on_column", + inputs, + history_id, + ) + + # Wait for all mapping jobs to complete + for job in response["jobs"]: + self.dataset_populator.wait_for_job(job["id"], assert_ok=True) + + # The implicit output should be a list:list collection + implicit_collections = response["implicit_collections"] + assert len(implicit_collections) == 1 + implicit_hdca_details = self.dataset_populator.get_history_collection_details( + history_id, content_id=implicit_collections[0]["id"] + ) + + sa_session = self._app.model.session + implicit_hdca_id = self._app.security.decode_id(implicit_hdca_details["id"]) + implicit_hdca = sa_session.get(HistoryDatasetCollectionAssociation, implicit_hdca_id) + assert implicit_hdca is not None + outer_dc = implicit_hdca.collection + + # The outer collection should have exactly 2 elements (one per input) + outer_count = sa_session.scalar( + select(func.count()).where(DatasetCollectionElement.dataset_collection_id == outer_dc.id) + ) + assert ( + outer_count == 2 + ), f"Expected 2 outer elements but found {outer_count}. Duplicate elements detected in mapped output!" + + # Each inner collection should have no duplicate elements + for element in outer_dc.elements: + inner_dc = element.child_collection + assert inner_dc is not None, f"Inner collection for '{element.element_identifier}' is None" + inner_count = sa_session.scalar( + select(func.count()).where(DatasetCollectionElement.dataset_collection_id == inner_dc.id) + ) + # Each input has 2 unique first-column values, so 2 split files + assert ( + inner_count is not None and inner_count > 0 + ), f"Inner collection for '{element.element_identifier}' has no elements" + # Check for duplicates: element_identifiers should be unique + inner_elements = sa_session.scalars( + select(DatasetCollectionElement).where(DatasetCollectionElement.dataset_collection_id == inner_dc.id) + ).all() + identifiers = [e.element_identifier for e in inner_elements] + assert len(identifiers) == len(set(identifiers)), ( + f"Duplicate element identifiers found in inner collection " + f"for '{element.element_identifier}': {identifiers}" + ) + + def test_no_duplicate_elements_in_mapped_static_collection(self, history_id): + """Map a tool with a static collection output over a list and verify no duplicate elements. + + This covers the case where a tool declares a static collection output + (e.g. paired with forward/reverse) and is mapped over a list input, + producing a list:paired output. The per-job DC contains fully + initialized elements and must be included in io_dicts for metadata + serialization. + """ + hdca = self.dataset_collection_populator.create_list_in_history( + history_id, + contents=["line1\nline2\nline3\nline4\n", "lineA\nlineB\nlineC\nlineD\n"], + ext="txt", + wait=True, + ) + self.dataset_populator.wait_for_history(history_id, assert_ok=True) + hdca_id = hdca.json()["output_collections"][0]["id"] + + inputs = { + "input1": {"batch": True, "values": [{"src": "hdca", "id": hdca_id}]}, + } + response = self.dataset_populator.run_tool( + "collection_creates_pair", + inputs, + history_id, + ) + + for job in response["jobs"]: + self.dataset_populator.wait_for_job(job["id"], assert_ok=True) + + implicit_collections = response["implicit_collections"] + assert len(implicit_collections) == 1 + implicit_hdca_details = self.dataset_populator.get_history_collection_details( + history_id, content_id=implicit_collections[0]["id"] + ) + + sa_session = self._app.model.session + implicit_hdca_id = self._app.security.decode_id(implicit_hdca_details["id"]) + implicit_hdca = sa_session.get(HistoryDatasetCollectionAssociation, implicit_hdca_id) + assert implicit_hdca is not None + outer_dc = implicit_hdca.collection + + # The outer collection should have exactly 2 elements (one per input) + outer_count = sa_session.scalar( + select(func.count()).where(DatasetCollectionElement.dataset_collection_id == outer_dc.id) + ) + assert ( + outer_count == 2 + ), f"Expected 2 outer elements but found {outer_count}. Duplicate elements detected in mapped output!" + + # Each inner pair should have exactly 2 elements (forward, reverse) with no duplicates + for element in outer_dc.elements: + inner_dc = element.child_collection + assert inner_dc is not None, f"Inner collection for '{element.element_identifier}' is None" + inner_count = sa_session.scalar( + select(func.count()).where(DatasetCollectionElement.dataset_collection_id == inner_dc.id) + ) + assert inner_count == 2, ( + f"Expected 2 inner elements for '{element.element_identifier}' " + f"but found {inner_count}. Duplicate elements detected!" + ) + inner_elements = sa_session.scalars( + select(DatasetCollectionElement).where(DatasetCollectionElement.dataset_collection_id == inner_dc.id) + ).all() + identifiers = [e.element_identifier for e in inner_elements] + assert set(identifiers) == { + "forward", + "reverse", + }, f"Expected forward/reverse but got {identifiers} for '{element.element_identifier}'" diff --git a/test/unit/data/model/test_model.py b/test/unit/data/model/test_model.py index 4ec6c459af8..36ebbbffa56 100644 --- a/test/unit/data/model/test_model.py +++ b/test/unit/data/model/test_model.py @@ -23,22 +23,36 @@ def test_permitted_actions(): def test_io_dicts_excludes_implicit_output_collections(): """Regression test for https://github.com/galaxyproject/galaxy/issues/22015 - When a tool with an explicit collection output is mapped over a list, - each job gets a JobToImplicitOutputDatasetCollectionAssociation pointing - to a DatasetCollection with precreated (unpopulated) elements. - io_dicts(exclude_implicit_outputs=True) must exclude these to avoid - crashes during metadata serialization on unpopulated elements. + When a tool with a dataset output is mapped over a list, each job gets + both a JobToOutputDatasetAssociation and a + JobToImplicitOutputDatasetCollectionAssociation with the same name. + The implicit DC has precreated (unpopulated) elements; only the current + job's element is initialized. io_dicts(exclude_implicit_outputs=True) + must exclude these shared DCs (name in out_data) to avoid crashes + during metadata serialization, but must include implicit DCs for + collection outputs (name not in out_data) so set_metadata.py can + discover and populate them. """ job = model.Job() dc = model.DatasetCollection(collection_type="paired") assoc = model.JobToImplicitOutputDatasetCollectionAssociation(name="paired_output", dataset_collection=dc) job.output_dataset_collections.append(assoc) - # With exclude_implicit_outputs=True, output_dataset_collections must be excluded + # When the name is NOT in out_data (collection output), the implicit DC + # should be included even with exclude_implicit_outputs=True + io = job.io_dicts(exclude_implicit_outputs=True) + assert "paired_output" in io.out_collections + assert io.out_collections["paired_output"] is dc + + # Now simulate a mapped dataset output: same name in both out_data and + # output_dataset_collections. The shared DC must be excluded. + hda = model.HistoryDatasetAssociation() + out_assoc = model.JobToOutputDatasetAssociation(name="paired_output", dataset=hda) + job.output_datasets.append(out_assoc) + io = job.io_dicts(exclude_implicit_outputs=True) assert "paired_output" not in io.out_collections # With exclude_implicit_outputs=False (default), they should be included io = job.io_dicts(exclude_implicit_outputs=False) assert "paired_output" in io.out_collections - assert io.out_collections["paired_output"] is dc From 9e2d5f5a94aeb95d4e85f238bc75712bb0f25e57 Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Tue, 17 Mar 2026 16:36:04 +0100 Subject: [PATCH 03/13] Move tests into TestExtendedMetadataIntegration --- test/integration/test_extended_metadata.py | 189 ++++++++++++++++ ...st_extended_metadata_duplicate_elements.py | 210 ------------------ 2 files changed, 189 insertions(+), 210 deletions(-) delete mode 100644 test/integration/test_extended_metadata_duplicate_elements.py diff --git a/test/integration/test_extended_metadata.py b/test/integration/test_extended_metadata.py index b6a5d7ae192..91896daa18f 100644 --- a/test/integration/test_extended_metadata.py +++ b/test/integration/test_extended_metadata.py @@ -1,6 +1,16 @@ """Integration tests for the Pulsar embedded runner.""" +from sqlalchemy import ( + func, + select, +) + +from galaxy.model import ( + DatasetCollectionElement, + HistoryDatasetCollectionAssociation, +) from galaxy_test.base.populators import ( + DatasetCollectionPopulator, DatasetPopulator, LibraryPopulator, ) @@ -56,6 +66,7 @@ class TestExtendedMetadataIntegration(integration_util.IntegrationTestCase): def setUp(self): super().setUp() self.dataset_populator = DatasetPopulator(self.galaxy_interactor) + self.dataset_collection_populator = DatasetCollectionPopulator(self.galaxy_interactor) self.library_populator = LibraryPopulator(self.galaxy_interactor) @classmethod @@ -100,6 +111,184 @@ class TestExtendedMetadataIntegration(integration_util.IntegrationTestCase): assert dataset["file_ext"] == "bed", dataset assert dataset["created_from_basename"] == "4.bed" + def test_no_duplicate_elements_in_dynamic_list_output(self, history_id): + """Run a tool with dynamic collection output and verify no duplicate elements.""" + response = self.dataset_populator.run_tool( + "collection_creates_dynamic_list_of_pairs", + {"foo": "bar"}, + history_id, + ) + job_api_id = response["jobs"][0]["id"] + self.dataset_populator.wait_for_job(job_api_id, assert_ok=True) + + output_collections = response["output_collections"] + hdca_details = self.dataset_populator.get_history_collection_details( + history_id, content_id=output_collections[0]["id"] + ) + + sa_session = self._app.model.session + hdca_id = self._app.security.decode_id(hdca_details["id"]) + hdca = sa_session.get(HistoryDatasetCollectionAssociation, hdca_id) + assert hdca is not None + dc = hdca.collection + + # Verify exactly 3 outer elements (samp1, samp2, samp3) - no duplicates + outer_count = sa_session.scalar( + select(func.count()).where(DatasetCollectionElement.dataset_collection_id == dc.id) + ) + assert outer_count == 3, f"Expected 3 outer elements but found {outer_count}. Duplicate elements detected!" + + # Verify each inner pair has exactly 2 elements (forward, reverse) + for element in dc.elements: + assert element.child_collection is not None + inner_count = sa_session.scalar( + select(func.count()).where( + DatasetCollectionElement.dataset_collection_id == element.child_collection.id + ) + ) + assert inner_count == 2, ( + f"Expected 2 inner elements for '{element.element_identifier}' " + f"but found {inner_count}. Duplicate elements detected!" + ) + + def test_no_duplicate_elements_in_mapped_dynamic_collection(self, history_id): + """Map a tool with dynamic collection output over a list and verify no duplicate elements.""" + # Create a list of 2 tabular datasets + fetch_response = self.dataset_collection_populator.create_list_in_history( + history_id, + contents=["101\t1\n101\t2\n105\t3\n", "201\t10\n201\t20\n205\t30\n"], + ext="tabular", + wait=True, + ) + self.dataset_populator.wait_for_history(history_id, assert_ok=True) + hdca_id = fetch_response.json()["output_collections"][0]["id"] + + # Map collection_split_on_column over the list + inputs = { + "input1": {"batch": True, "values": [{"src": "hdca", "id": hdca_id}]}, + } + response = self.dataset_populator.run_tool( + "collection_split_on_column", + inputs, + history_id, + ) + + # Wait for all mapping jobs to complete + for job in response["jobs"]: + self.dataset_populator.wait_for_job(job["id"], assert_ok=True) + + # The implicit output should be a list:list collection + implicit_collections = response["implicit_collections"] + assert len(implicit_collections) == 1 + implicit_hdca_details = self.dataset_populator.get_history_collection_details( + history_id, content_id=implicit_collections[0]["id"] + ) + + sa_session = self._app.model.session + implicit_hdca_id = self._app.security.decode_id(implicit_hdca_details["id"]) + implicit_hdca = sa_session.get(HistoryDatasetCollectionAssociation, implicit_hdca_id) + assert implicit_hdca is not None + outer_dc = implicit_hdca.collection + + # The outer collection should have exactly 2 elements (one per input) + outer_count = sa_session.scalar( + select(func.count()).where(DatasetCollectionElement.dataset_collection_id == outer_dc.id) + ) + assert ( + outer_count == 2 + ), f"Expected 2 outer elements but found {outer_count}. Duplicate elements detected in mapped output!" + + # Each inner collection should have no duplicate elements + for element in outer_dc.elements: + inner_dc = element.child_collection + assert inner_dc is not None, f"Inner collection for '{element.element_identifier}' is None" + inner_count = sa_session.scalar( + select(func.count()).where(DatasetCollectionElement.dataset_collection_id == inner_dc.id) + ) + # Each input has 2 unique first-column values, so 2 split files + assert ( + inner_count is not None and inner_count > 0 + ), f"Inner collection for '{element.element_identifier}' has no elements" + # Check for duplicates: element_identifiers should be unique + inner_elements = sa_session.scalars( + select(DatasetCollectionElement).where(DatasetCollectionElement.dataset_collection_id == inner_dc.id) + ).all() + identifiers = [e.element_identifier for e in inner_elements] + assert len(identifiers) == len(set(identifiers)), ( + f"Duplicate element identifiers found in inner collection " + f"for '{element.element_identifier}': {identifiers}" + ) + + def test_no_duplicate_elements_in_mapped_static_collection(self, history_id): + """Map a tool with a static collection output over a list and verify no duplicate elements. + + This covers the case where a tool declares a static collection output + (e.g. paired with forward/reverse) and is mapped over a list input, + producing a list:paired output. The per-job DC contains fully + initialized elements and must be included in io_dicts for metadata + serialization. + """ + hdca = self.dataset_collection_populator.create_list_in_history( + history_id, + contents=["line1\nline2\nline3\nline4\n", "lineA\nlineB\nlineC\nlineD\n"], + ext="txt", + wait=True, + ) + self.dataset_populator.wait_for_history(history_id, assert_ok=True) + hdca_id = hdca.json()["output_collections"][0]["id"] + + inputs = { + "input1": {"batch": True, "values": [{"src": "hdca", "id": hdca_id}]}, + } + response = self.dataset_populator.run_tool( + "collection_creates_pair", + inputs, + history_id, + ) + + for job in response["jobs"]: + self.dataset_populator.wait_for_job(job["id"], assert_ok=True) + + implicit_collections = response["implicit_collections"] + assert len(implicit_collections) == 1 + implicit_hdca_details = self.dataset_populator.get_history_collection_details( + history_id, content_id=implicit_collections[0]["id"] + ) + + sa_session = self._app.model.session + implicit_hdca_id = self._app.security.decode_id(implicit_hdca_details["id"]) + implicit_hdca = sa_session.get(HistoryDatasetCollectionAssociation, implicit_hdca_id) + assert implicit_hdca is not None + outer_dc = implicit_hdca.collection + + # The outer collection should have exactly 2 elements (one per input) + outer_count = sa_session.scalar( + select(func.count()).where(DatasetCollectionElement.dataset_collection_id == outer_dc.id) + ) + assert ( + outer_count == 2 + ), f"Expected 2 outer elements but found {outer_count}. Duplicate elements detected in mapped output!" + + # Each inner pair should have exactly 2 elements (forward, reverse) with no duplicates + for element in outer_dc.elements: + inner_dc = element.child_collection + assert inner_dc is not None, f"Inner collection for '{element.element_identifier}' is None" + inner_count = sa_session.scalar( + select(func.count()).where(DatasetCollectionElement.dataset_collection_id == inner_dc.id) + ) + assert inner_count == 2, ( + f"Expected 2 inner elements for '{element.element_identifier}' " + f"but found {inner_count}. Duplicate elements detected!" + ) + inner_elements = sa_session.scalars( + select(DatasetCollectionElement).where(DatasetCollectionElement.dataset_collection_id == inner_dc.id) + ).all() + identifiers = [e.element_identifier for e in inner_elements] + assert set(identifiers) == { + "forward", + "reverse", + }, f"Expected forward/reverse but got {identifiers} for '{element.element_identifier}'" + def test_purge_while_job_running(self): # pass extra_sleep, since templating the command line will fail if the output # is deleted before remote_tool_eval runs. diff --git a/test/integration/test_extended_metadata_duplicate_elements.py b/test/integration/test_extended_metadata_duplicate_elements.py deleted file mode 100644 index 82ad0981bc4..00000000000 --- a/test/integration/test_extended_metadata_duplicate_elements.py +++ /dev/null @@ -1,210 +0,0 @@ -"""Integration test verifying that extended metadata does not create duplicate collection elements.""" - -from sqlalchemy import ( - func, - select, -) - -from galaxy.model import ( - DatasetCollectionElement, - HistoryDatasetCollectionAssociation, -) -from galaxy_test.base.populators import ( - DatasetCollectionPopulator, - DatasetPopulator, -) -from galaxy_test.driver.integration_util import IntegrationTestCase - - -class TestExtendedMetadataDuplicateElements(IntegrationTestCase): - dataset_populator: DatasetPopulator - framework_tool_and_types = True - - @classmethod - def handle_galaxy_config_kwds(cls, config): - super().handle_galaxy_config_kwds(config) - config["metadata_strategy"] = "extended" - config["retry_metadata_internally"] = False - - def setUp(self): - super().setUp() - self.dataset_populator = DatasetPopulator(self.galaxy_interactor) - self.dataset_collection_populator = DatasetCollectionPopulator(self.galaxy_interactor) - - def test_no_duplicate_elements_in_dynamic_list_output(self, history_id): - """Run a tool with dynamic collection output and verify no duplicate elements.""" - response = self.dataset_populator.run_tool( - "collection_creates_dynamic_list_of_pairs", - {"foo": "bar"}, - history_id, - ) - job_api_id = response["jobs"][0]["id"] - self.dataset_populator.wait_for_job(job_api_id, assert_ok=True) - - output_collections = response["output_collections"] - hdca_details = self.dataset_populator.get_history_collection_details( - history_id, content_id=output_collections[0]["id"] - ) - - sa_session = self._app.model.session - hdca_id = self._app.security.decode_id(hdca_details["id"]) - hdca = sa_session.get(HistoryDatasetCollectionAssociation, hdca_id) - assert hdca is not None - dc = hdca.collection - - # Verify exactly 3 outer elements (samp1, samp2, samp3) - no duplicates - outer_count = sa_session.scalar( - select(func.count()).where(DatasetCollectionElement.dataset_collection_id == dc.id) - ) - assert outer_count == 3, f"Expected 3 outer elements but found {outer_count}. Duplicate elements detected!" - - # Verify each inner pair has exactly 2 elements (forward, reverse) - for element in dc.elements: - assert element.child_collection is not None - inner_count = sa_session.scalar( - select(func.count()).where( - DatasetCollectionElement.dataset_collection_id == element.child_collection.id - ) - ) - assert inner_count == 2, ( - f"Expected 2 inner elements for '{element.element_identifier}' " - f"but found {inner_count}. Duplicate elements detected!" - ) - - def test_no_duplicate_elements_in_mapped_dynamic_collection(self, history_id): - """Map a tool with dynamic collection output over a list and verify no duplicate elements.""" - # Create a list of 2 tabular datasets - fetch_response = self.dataset_collection_populator.create_list_in_history( - history_id, - contents=["101\t1\n101\t2\n105\t3\n", "201\t10\n201\t20\n205\t30\n"], - ext="tabular", - wait=True, - ) - self.dataset_populator.wait_for_history(history_id, assert_ok=True) - hdca_id = fetch_response.json()["output_collections"][0]["id"] - - # Map collection_split_on_column over the list - inputs = { - "input1": {"batch": True, "values": [{"src": "hdca", "id": hdca_id}]}, - } - response = self.dataset_populator.run_tool( - "collection_split_on_column", - inputs, - history_id, - ) - - # Wait for all mapping jobs to complete - for job in response["jobs"]: - self.dataset_populator.wait_for_job(job["id"], assert_ok=True) - - # The implicit output should be a list:list collection - implicit_collections = response["implicit_collections"] - assert len(implicit_collections) == 1 - implicit_hdca_details = self.dataset_populator.get_history_collection_details( - history_id, content_id=implicit_collections[0]["id"] - ) - - sa_session = self._app.model.session - implicit_hdca_id = self._app.security.decode_id(implicit_hdca_details["id"]) - implicit_hdca = sa_session.get(HistoryDatasetCollectionAssociation, implicit_hdca_id) - assert implicit_hdca is not None - outer_dc = implicit_hdca.collection - - # The outer collection should have exactly 2 elements (one per input) - outer_count = sa_session.scalar( - select(func.count()).where(DatasetCollectionElement.dataset_collection_id == outer_dc.id) - ) - assert ( - outer_count == 2 - ), f"Expected 2 outer elements but found {outer_count}. Duplicate elements detected in mapped output!" - - # Each inner collection should have no duplicate elements - for element in outer_dc.elements: - inner_dc = element.child_collection - assert inner_dc is not None, f"Inner collection for '{element.element_identifier}' is None" - inner_count = sa_session.scalar( - select(func.count()).where(DatasetCollectionElement.dataset_collection_id == inner_dc.id) - ) - # Each input has 2 unique first-column values, so 2 split files - assert ( - inner_count is not None and inner_count > 0 - ), f"Inner collection for '{element.element_identifier}' has no elements" - # Check for duplicates: element_identifiers should be unique - inner_elements = sa_session.scalars( - select(DatasetCollectionElement).where(DatasetCollectionElement.dataset_collection_id == inner_dc.id) - ).all() - identifiers = [e.element_identifier for e in inner_elements] - assert len(identifiers) == len(set(identifiers)), ( - f"Duplicate element identifiers found in inner collection " - f"for '{element.element_identifier}': {identifiers}" - ) - - def test_no_duplicate_elements_in_mapped_static_collection(self, history_id): - """Map a tool with a static collection output over a list and verify no duplicate elements. - - This covers the case where a tool declares a static collection output - (e.g. paired with forward/reverse) and is mapped over a list input, - producing a list:paired output. The per-job DC contains fully - initialized elements and must be included in io_dicts for metadata - serialization. - """ - hdca = self.dataset_collection_populator.create_list_in_history( - history_id, - contents=["line1\nline2\nline3\nline4\n", "lineA\nlineB\nlineC\nlineD\n"], - ext="txt", - wait=True, - ) - self.dataset_populator.wait_for_history(history_id, assert_ok=True) - hdca_id = hdca.json()["output_collections"][0]["id"] - - inputs = { - "input1": {"batch": True, "values": [{"src": "hdca", "id": hdca_id}]}, - } - response = self.dataset_populator.run_tool( - "collection_creates_pair", - inputs, - history_id, - ) - - for job in response["jobs"]: - self.dataset_populator.wait_for_job(job["id"], assert_ok=True) - - implicit_collections = response["implicit_collections"] - assert len(implicit_collections) == 1 - implicit_hdca_details = self.dataset_populator.get_history_collection_details( - history_id, content_id=implicit_collections[0]["id"] - ) - - sa_session = self._app.model.session - implicit_hdca_id = self._app.security.decode_id(implicit_hdca_details["id"]) - implicit_hdca = sa_session.get(HistoryDatasetCollectionAssociation, implicit_hdca_id) - assert implicit_hdca is not None - outer_dc = implicit_hdca.collection - - # The outer collection should have exactly 2 elements (one per input) - outer_count = sa_session.scalar( - select(func.count()).where(DatasetCollectionElement.dataset_collection_id == outer_dc.id) - ) - assert ( - outer_count == 2 - ), f"Expected 2 outer elements but found {outer_count}. Duplicate elements detected in mapped output!" - - # Each inner pair should have exactly 2 elements (forward, reverse) with no duplicates - for element in outer_dc.elements: - inner_dc = element.child_collection - assert inner_dc is not None, f"Inner collection for '{element.element_identifier}' is None" - inner_count = sa_session.scalar( - select(func.count()).where(DatasetCollectionElement.dataset_collection_id == inner_dc.id) - ) - assert inner_count == 2, ( - f"Expected 2 inner elements for '{element.element_identifier}' " - f"but found {inner_count}. Duplicate elements detected!" - ) - inner_elements = sa_session.scalars( - select(DatasetCollectionElement).where(DatasetCollectionElement.dataset_collection_id == inner_dc.id) - ).all() - identifiers = [e.element_identifier for e in inner_elements] - assert set(identifiers) == { - "forward", - "reverse", - }, f"Expected forward/reverse but got {identifiers} for '{element.element_identifier}'" From 9d56b75d866ed1aeae4a798d89839a2c0ccf2c3b Mon Sep 17 00:00:00 2001 From: davelopez <46503462+davelopez@users.noreply.github.com> Date: Wed, 18 Mar 2026 17:43:36 +0100 Subject: [PATCH 04/13] Bump tiffviewer visualization to v0.0.4 Fixes resolution issue on macOS and zoom drifting --- client/visualizations.yml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/client/visualizations.yml b/client/visualizations.yml index f6ddbbdab5f..eaadeac63b8 100644 --- a/client/visualizations.yml +++ b/client/visualizations.yml @@ -110,7 +110,7 @@ tabulator: version: 0.0.9 tiffviewer: package: "@galaxyproject/tiffviewer" - version: 0.0.3 + version: 0.0.4 ts_visjs: package: "@galaxyproject/ts_visjs" version: 0.0.6 From ecc5d1c41fd5c12b05c7087bea741d93042643da Mon Sep 17 00:00:00 2001 From: John Chilton Date: Wed, 18 Mar 2026 21:21:58 -0400 Subject: [PATCH 05/13] Update gxformat2 dependency to 0.23.0 release. Co-Authored-By: Claude Opus 4.6 (1M context) --- lib/galaxy/dependencies/pinned-requirements.txt | 2 +- pyproject.toml | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/lib/galaxy/dependencies/pinned-requirements.txt b/lib/galaxy/dependencies/pinned-requirements.txt index 784a6a33dff..d83edca597f 100644 --- a/lib/galaxy/dependencies/pinned-requirements.txt +++ b/lib/galaxy/dependencies/pinned-requirements.txt @@ -104,7 +104,7 @@ groq==1.0.0 grpcio==1.76.0 grpcio-status==1.76.0 gunicorn==24.1.1 -gxformat2==0.21.0 +gxformat2==0.23.0 h11==0.16.0 h5grove==2.3.0 h5py==3.15.1 diff --git a/pyproject.toml b/pyproject.toml index 6cb1cdc602a..7ae497514d8 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -43,7 +43,7 @@ dependencies = [ "future>=1.0.0", # Python 3.12 support "gravity>=1.2.0", # Python 3.14 support "gunicorn", - "gxformat2>=0.21.0", + "gxformat2>=0.23.0", "h5grove>=1.2.1", "h5py>=3.12", # Python 3.13 support "httpx", From 761efd4496edb8cf3b298c13cc5e61d96c9aa9a8 Mon Sep 17 00:00:00 2001 From: John Chilton Date: Wed, 18 Mar 2026 21:22:21 -0400 Subject: [PATCH 06/13] Save workflow comments when converting to and from gxformat2. --- lib/galaxy_test/api/test_workflows.py | 79 +++++++++++++++++++++++++++ 1 file changed, 79 insertions(+) diff --git a/lib/galaxy_test/api/test_workflows.py b/lib/galaxy_test/api/test_workflows.py index ab572a654dd..ba2de188db9 100644 --- a/lib/galaxy_test/api/test_workflows.py +++ b/lib/galaxy_test/api/test_workflows.py @@ -1553,6 +1553,85 @@ steps: downloaded_workflow = self._download_workflow(uploaded_workflow_id, style="format2") assert downloaded_workflow["class"] == "GalaxyWorkflow" + def test_export_format2_comments(self): + # Build a native workflow with comments + workflow = self.workflow_populator.load_workflow("test_comments_format2") + workflow["comments"] = [ + { + "id": 0, + "type": "text", + "position": [100, 200], + "size": [200, 50], + "color": "blue", + "data": {"text": "Check adapters", "size": 2, "bold": True}, + }, + { + "id": 1, + "type": "markdown", + "position": [300, 50], + "size": [400, 300], + "color": "none", + "data": {"text": "# Preprocessing\nQuality filtering."}, + }, + { + "id": 2, + "type": "frame", + "position": [50, 50], + "size": [700, 500], + "color": "green", + "data": {"title": "Main Steps"}, + "child_steps": [0], + "child_comments": [0, 1], + }, + { + "id": 3, + "type": "freehand", + "position": [600, 100], + "size": [50, 50], + "color": "red", + "data": {"thickness": 3, "line": [[610, 110], [620, 120], [630, 130]]}, + }, + ] + uploaded_workflow_id = self.workflow_populator.create_workflow(workflow) + + # Verify native export preserves comments + native = self._download_workflow(uploaded_workflow_id) + assert "comments" in native + assert len(native["comments"]) == 4 + native_text = [c for c in native["comments"] if c["type"] == "text"][0] + assert native_text["data"]["text"] == "Check adapters" + + # Verify format2 export converts comments + fmt2 = self._download_workflow(uploaded_workflow_id, style="format2") + assert "comments" in fmt2 + comments = fmt2["comments"] + assert isinstance(comments, list) + assert len(comments) == 4 + + # Text comment: data fields should be flattened to top level + text_c = [c for c in comments if c["type"] == "text"][0] + assert text_c["text"] == "Check adapters" + assert text_c["text_size"] == 2 + assert text_c["bold"] is True + assert "data" not in text_c + + # Markdown comment + md_c = [c for c in comments if c["type"] == "markdown"][0] + assert "# Preprocessing" in md_c["text"] + + # Frame: child_steps/child_comments renamed to contains_steps/contains_comments + frame_c = [c for c in comments if c["type"] == "frame"][0] + assert frame_c["title"] == "Main Steps" + assert "contains_steps" in frame_c + assert "contains_comments" in frame_c + assert "child_steps" not in frame_c + assert "child_comments" not in frame_c + + # Freehand + fh_c = [c for c in comments if c["type"] == "freehand"][0] + assert fh_c["thickness"] == 3 + assert len(fh_c["line"]) == 3 + def test_export_editor(self): uploaded_workflow_id = self.workflow_populator.simple_workflow("test_for_export") downloaded_workflow = self._download_workflow(uploaded_workflow_id, style="editor") From 0c2feefed5a74143686a91230e3ee3ba8d26def9 Mon Sep 17 00:00:00 2001 From: John Chilton Date: Thu, 19 Mar 2026 06:26:03 -0400 Subject: [PATCH 07/13] Include expression tools in stock_tool_paths. MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit stock_tool_paths() missed tools/expression_tools/ (param_value_from_file, pick_value) — these are Galaxy stock tools but live outside lib/galaxy/tools/. Co-Authored-By: Claude Opus 4.6 (1M context) --- lib/galaxy/tools/stock.py | 1 + test/unit/app/tools/test_stock.py | 1 + 2 files changed, 2 insertions(+) diff --git a/lib/galaxy/tools/stock.py b/lib/galaxy/tools/stock.py index eead325e644..a1a6804ca4d 100644 --- a/lib/galaxy/tools/stock.py +++ b/lib/galaxy/tools/stock.py @@ -16,6 +16,7 @@ from galaxy.util.resources import files def stock_tool_paths(): yield from _walk_directory_for_tools(files(galaxy.tools)) yield from _walk_directory_for_tools(Path(galaxy_directory()) / "test" / "functional" / "tools") + yield from _walk_directory_for_tools(Path(galaxy_directory()) / "tools") def stock_tool_sources(): diff --git a/test/unit/app/tools/test_stock.py b/test/unit/app/tools/test_stock.py index 25aeca87531..3918bd72b90 100644 --- a/test/unit/app/tools/test_stock.py +++ b/test/unit/app/tools/test_stock.py @@ -9,6 +9,7 @@ def test_stock_tool_paths(): assert "merge_collection.xml" in file_names assert "meme.xml" in file_names assert "output_auto_format.xml" in file_names + assert "parse_values_from_file.xml" in file_names def test_stock_tool_sources(): From cc6a708bb02caff481cbedebe6e589891ab44521 Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Thu, 19 Mar 2026 11:59:58 +0100 Subject: [PATCH 08/13] Add batch celery task for history dataset purging When a user purges a large history, each dataset was individually dispatched as a celery task, which could block celery workers for days with histories containing millions of datasets. Replace the per-HDA task dispatch with a single `purge_history_datasets` celery task that processes all HDAs in the history at once: marks them as purged, adjusts quotas, and removes files from storage. --- lib/galaxy/celery/tasks.py | 55 +++++++++++++++++++++++++ lib/galaxy/managers/histories.py | 14 +++++-- lib/galaxy/schema/tasks.py | 4 ++ test/integration/test_purge_datasets.py | 43 ++++++++++++++++++- 4 files changed, 112 insertions(+), 4 deletions(-) diff --git a/lib/galaxy/celery/tasks.py b/lib/galaxy/celery/tasks.py index c674af43a49..4a3f3128133 100644 --- a/lib/galaxy/celery/tasks.py +++ b/lib/galaxy/celery/tasks.py @@ -11,8 +11,11 @@ from typing import ( ) from sqlalchemy import ( + and_, exists, + false, select, + update, ) from galaxy import model @@ -58,6 +61,7 @@ from galaxy.schema.tasks import ( MaterializeDatasetInstanceTaskRequest, PrepareDatasetCollectionDownload, PurgeDatasetsTaskRequest, + PurgeHistoryDatasetsTaskRequest, QueueJobs, SetupHistoryExportJob, TOOL_SOURCE_CLASS, @@ -117,6 +121,57 @@ def purge_datasets( dataset_manager.purge_datasets(request) +@galaxy_task(ignore_result=True, action="purge all datasets in a history") +def purge_history_datasets( + sa_session: galaxy_scoped_session, + dataset_manager: DatasetManager, + object_store: BaseObjectStore, + request: PurgeHistoryDatasetsTaskRequest, + task_user_id: Optional[int] = None, +): + """Batch purge all HDAs in a history in a single task. + + Bulk-marks all unpurged HDAs as deleted/purged, recalculates user quota, + and removes underlying dataset files from the object store. + """ + history = sa_session.get(model.History, request.history_id) + if not history: + log.error(f"Purge history datasets task failed, history {request.history_id} not found") + return + # Collect dataset IDs before the bulk update + dataset_id_stmt = ( + select(model.HistoryDatasetAssociation.dataset_id) + .where( + and_( + model.HistoryDatasetAssociation.history_id == request.history_id, + model.HistoryDatasetAssociation.purged == false(), + ) + ) + .distinct() + ) + dataset_ids = list(sa_session.scalars(dataset_id_stmt)) + if not dataset_ids: + return + # Bulk mark all unpurged HDAs as deleted and purged + sa_session.execute( + update(model.HistoryDatasetAssociation) + .where( + and_( + model.HistoryDatasetAssociation.history_id == request.history_id, + model.HistoryDatasetAssociation.purged == false(), + ) + ) + .values(deleted=True, purged=True) + ) + sa_session.commit() + # Recalculate user disk usage from scratch + user = history.user + if user: + user.calculate_and_set_disk_usage(object_store) + # Remove underlying dataset files from object store + dataset_manager.purge_datasets(PurgeDatasetsTaskRequest(dataset_ids=dataset_ids)) + + @galaxy_task(ignore_result=True, action="materializing dataset instance") def materialize( hda_manager: HDAManager, diff --git a/lib/galaxy/managers/histories.py b/lib/galaxy/managers/histories.py index cc9aa850ded..6f8bd1200e9 100644 --- a/lib/galaxy/managers/histories.py +++ b/lib/galaxy/managers/histories.py @@ -79,6 +79,7 @@ from galaxy.schema.storage_cleaner import ( StoredItem, StoredItemOrderBy, ) +from galaxy.schema.tasks import PurgeHistoryDatasetsTaskRequest from galaxy.security.validate_user_input import validate_preferred_object_store_id from galaxy.structured_app import MinimalManagerApp from galaxy.util.search import ( @@ -292,9 +293,16 @@ class HistoryManager(sharable.SharableModelManager[model.History], deletable.Pur self.error_unless_mutable(item) self.hda_manager.dataset_manager.error_unless_dataset_purge_allowed() # First purge all the datasets - for hda in item.datasets: - if not hda.purged: - self.hda_manager.purge(hda, flush=True, **kwargs) + if self.app.config.enable_celery_tasks: + from galaxy.celery.tasks import purge_history_datasets + + request = PurgeHistoryDatasetsTaskRequest(history_id=item.id) + user = item.user + purge_history_datasets.delay(request=request, task_user_id=user.id if user else None) + else: + for hda in item.datasets: + if not hda.purged: + self.hda_manager.purge(hda, flush=True, **kwargs) # Now mark the history as purged super().purge(item, flush=flush, **kwargs) diff --git a/lib/galaxy/schema/tasks.py b/lib/galaxy/schema/tasks.py index 8c6ef67e7f8..4ae0e7f10dc 100644 --- a/lib/galaxy/schema/tasks.py +++ b/lib/galaxy/schema/tasks.py @@ -139,6 +139,10 @@ class PurgeDatasetsTaskRequest(Model): dataset_ids: list[int] +class PurgeHistoryDatasetsTaskRequest(Model): + history_id: int + + class TaskState(str, Enum): """Enum representing the possible states of a task.""" diff --git a/test/integration/test_purge_datasets.py b/test/integration/test_purge_datasets.py index 363245397c1..8f0855ab8e8 100644 --- a/test/integration/test_purge_datasets.py +++ b/test/integration/test_purge_datasets.py @@ -4,7 +4,10 @@ from typing import ( Optional, ) -from galaxy_test.base.populators import DatasetPopulator +from galaxy_test.base.populators import ( + DatasetPopulator, + wait_on, +) from galaxy_test.driver import integration_util @@ -76,9 +79,47 @@ class TestPurgeDatasetsIntegration(integration_util.IntegrationTestCase): purge_result = purge_response.json() assert purge_result["success_count"] == 1 + def test_purge_history_removes_underlying_datasets_from_disk(self): + """Test that purging a history purges all its datasets and removes files from disk.""" + hda1 = self.dataset_populator.new_dataset(self.test_history_id, wait=True) + hda2 = self.dataset_populator.new_dataset(self.test_history_id, wait=True) + hda1_id = hda1["id"] + hda2_id = hda2["id"] + + # Ensure dataset files exist on disk + dataset_file1 = self._get_underlying_dataset_on_disk(hda1_id) + dataset_file2 = self._get_underlying_dataset_on_disk(hda2_id) + assert self._file_exists_on_disk(dataset_file1) + assert self._file_exists_on_disk(dataset_file2) + + # Purge the entire history + purge_response = self._delete(f"histories/{self.test_history_id}", data={"purge": True}, json=True) + self._assert_status_code_is_ok(purge_response) + + # Verify history is purged + history_response = self._get(f"histories/{self.test_history_id}").json() + assert history_response["purged"] + assert history_response["deleted"] + + # Verify HDAs are marked as purged + self.dataset_populator.wait_for_purge(self.test_history_id, hda1_id) + self.dataset_populator.wait_for_purge(self.test_history_id, hda2_id) + + # Wait for underlying dataset files to be removed from disk. + # With batched history purging, HDAs are marked purged synchronously + # but the actual file deletion happens via a batched celery task. + self._wait_for_file_deleted(dataset_file1) + self._wait_for_file_deleted(dataset_file2) + def _get_underlying_dataset_on_disk(self, hda_id: str) -> Optional[str]: detailed_response = self._get(f"datasets/{hda_id}", admin=True).json() return detailed_response.get("file_name") def _file_exists_on_disk(self, filename: Optional[str]) -> bool: return os.path.isfile(filename) if filename else False + + def _wait_for_file_deleted(self, filename: Optional[str], timeout: int = 10): + def _check(): + return True if not self._file_exists_on_disk(filename) else None + + wait_on(_check, f"file {filename} to be deleted", timeout=timeout) From 63b63e7d6f146d23f3ce2f2a07fc32846df283b7 Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Thu, 19 Mar 2026 14:43:04 +0100 Subject: [PATCH 09/13] Don't fail on probable delete race Fixes https://github.com/galaxyproject/galaxy/issues/22182 --- lib/galaxy/model/__init__.py | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/lib/galaxy/model/__init__.py b/lib/galaxy/model/__init__.py index 58345f05681..750a154a162 100644 --- a/lib/galaxy/model/__init__.py +++ b/lib/galaxy/model/__init__.py @@ -4737,14 +4737,16 @@ class Dataset(Base, StorableObject, Serializable): def full_delete(self): """Remove the file and extra files, marks deleted and purged""" - # os.unlink( self.file_name ) try: self.object_store.delete(self) except galaxy.exceptions.ObjectNotFound: pass if (rel_path := self._extra_files_rel_path) is not None: if self.object_store.exists(self, extra_dir=rel_path, dir_only=True): - self.object_store.delete(self, entire_dir=True, extra_dir=rel_path, dir_only=True) + try: + self.object_store.delete(self, entire_dir=True, extra_dir=rel_path, dir_only=True) + except galaxy.exceptions.ObjectNotFound: + pass # TODO: purge metadata files self.deleted = True self.purged = True From f79692760ee33ebb4f6dfb8db7d8eba4bd50f38b Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Thu, 19 Mar 2026 14:30:59 +0100 Subject: [PATCH 10/13] Return async task result --- client/src/api/schema/schema.ts | 30 ++++++++++++++++ lib/galaxy/celery/tasks.py | 2 +- lib/galaxy/managers/histories.py | 4 ++- lib/galaxy/schema/schema.py | 12 +++++++ .../webapps/galaxy/services/histories.py | 8 +++-- test/integration/test_purge_datasets.py | 36 +++++++++---------- 6 files changed, 68 insertions(+), 24 deletions(-) diff --git a/client/src/api/schema/schema.ts b/client/src/api/schema/schema.ts index 710506133cb..de8de9ba3bc 100644 --- a/client/src/api/schema/schema.ts +++ b/client/src/api/schema/schema.ts @@ -7463,6 +7463,11 @@ export interface components { * @description Whether this resource is currently publicly available to all users. */ published: boolean; + /** + * Purge Task + * @description Summary of the async task purging datasets in this history. Only present when purge is performed via a background task. + */ + purge_task?: components["schemas"]["AsyncTaskResultSummary"] | null; /** * Purged * @description Whether this item has been permanently removed. @@ -7586,6 +7591,11 @@ export interface components { * @description Whether this resource is currently publicly available to all users. */ published: boolean; + /** + * Purge Task + * @description Summary of the async task purging datasets in this history. Only present when purge is performed via a background task. + */ + purge_task?: components["schemas"]["AsyncTaskResultSummary"] | null; /** * Purged * @description Whether this item has been permanently removed. @@ -9974,6 +9984,11 @@ export interface components { * @description Whether this resource is currently publicly available to all users. */ published?: boolean | null; + /** + * Purge Task + * @description Summary of the async task purging datasets in this history. Only present when purge is performed via a background task. + */ + purge_task?: components["schemas"]["AsyncTaskResultSummary"] | null; /** * Purged * @description Whether this item has been permanently removed. @@ -10213,6 +10228,11 @@ export interface components { * @description Whether this resource is currently publicly available to all users. */ published?: boolean | null; + /** + * Purge Task + * @description Summary of the async task purging datasets in this history. Only present when purge is performed via a background task. + */ + purge_task?: components["schemas"]["AsyncTaskResultSummary"] | null; /** * Purged * @description Whether this item has been permanently removed. @@ -15492,6 +15512,11 @@ export interface components { * @description Whether this resource is currently publicly available to all users. */ published: boolean; + /** + * Purge Task + * @description Summary of the async task purging datasets in this history. Only present when purge is performed via a background task. + */ + purge_task?: components["schemas"]["AsyncTaskResultSummary"] | null; /** * Purged * @description Whether this item has been permanently removed. @@ -15613,6 +15638,11 @@ export interface components { * @description Whether this resource is currently publicly available to all users. */ published: boolean; + /** + * Purge Task + * @description Summary of the async task purging datasets in this history. Only present when purge is performed via a background task. + */ + purge_task?: components["schemas"]["AsyncTaskResultSummary"] | null; /** * Purged * @description Whether this item has been permanently removed. diff --git a/lib/galaxy/celery/tasks.py b/lib/galaxy/celery/tasks.py index 4a3f3128133..61833ba2be6 100644 --- a/lib/galaxy/celery/tasks.py +++ b/lib/galaxy/celery/tasks.py @@ -121,7 +121,7 @@ def purge_datasets( dataset_manager.purge_datasets(request) -@galaxy_task(ignore_result=True, action="purge all datasets in a history") +@galaxy_task(action="purge all datasets in a history") def purge_history_datasets( sa_session: galaxy_scoped_session, dataset_manager: DatasetManager, diff --git a/lib/galaxy/managers/histories.py b/lib/galaxy/managers/histories.py index 6f8bd1200e9..053d456e511 100644 --- a/lib/galaxy/managers/histories.py +++ b/lib/galaxy/managers/histories.py @@ -298,14 +298,16 @@ class HistoryManager(sharable.SharableModelManager[model.History], deletable.Pur request = PurgeHistoryDatasetsTaskRequest(history_id=item.id) user = item.user - purge_history_datasets.delay(request=request, task_user_id=user.id if user else None) + result = purge_history_datasets.delay(request=request, task_user_id=user.id if user else None) else: + result = None for hda in item.datasets: if not hda.purged: self.hda_manager.purge(hda, flush=True, **kwargs) # Now mark the history as purged super().purge(item, flush=flush, **kwargs) + return result # .... current # TODO: make something to bypass the anon user + current history permissions issue diff --git a/lib/galaxy/schema/schema.py b/lib/galaxy/schema/schema.py index 24cbf0e4d92..6a1c45d5f52 100644 --- a/lib/galaxy/schema/schema.py +++ b/lib/galaxy/schema/schema.py @@ -1435,6 +1435,11 @@ class HistorySummary(Model, WithModelClass): tags: TagCollection update_time: datetime = UpdateTimeField preferred_object_store_id: Optional[str] = PreferredObjectStoreIdField + purge_task: Optional["AsyncTaskResultSummary"] = Field( + None, + title="Purge Task", + description="Summary of the async task purging datasets in this history. Only present when purge is performed via a background task.", + ) class HistoryActiveContentCounts(Model): @@ -3997,6 +4002,13 @@ class AsyncTaskResultSummary(Model): ) +HistorySummary.model_rebuild() +HistoryDetailed.model_rebuild() +CustomHistoryView.model_rebuild() +ArchivedHistorySummary.model_rebuild() +ArchivedHistoryDetailed.model_rebuild() +CustomArchivedHistoryView.model_rebuild() + ToolRequestIdField = Field(title="ID", description="Encoded ID of the role") diff --git a/lib/galaxy/webapps/galaxy/services/histories.py b/lib/galaxy/webapps/galaxy/services/histories.py index 3d6a4fb0056..55549c44f21 100644 --- a/lib/galaxy/webapps/galaxy/services/histories.py +++ b/lib/galaxy/webapps/galaxy/services/histories.py @@ -454,10 +454,14 @@ class HistoriesService(ServiceBase, ConsumesModelStores, ServesExportStores): """ history = self.manager.get_mutable(history_id, trans.user, current_history=trans.history) if purge: - self.manager.purge(history) + result = self.manager.purge(history) else: + result = None self.manager.delete(history) - return self._serialize_history(trans, history, serialization_params) + rval = self._serialize_history(trans, history, serialization_params) + if result is not None: + rval["purge_task"] = async_task_summary(result) + return rval def undelete( self, diff --git a/test/integration/test_purge_datasets.py b/test/integration/test_purge_datasets.py index 8f0855ab8e8..753ec4420fd 100644 --- a/test/integration/test_purge_datasets.py +++ b/test/integration/test_purge_datasets.py @@ -4,10 +4,7 @@ from typing import ( Optional, ) -from galaxy_test.base.populators import ( - DatasetPopulator, - wait_on, -) +from galaxy_test.base.populators import DatasetPopulator from galaxy_test.driver import integration_util @@ -95,21 +92,26 @@ class TestPurgeDatasetsIntegration(integration_util.IntegrationTestCase): # Purge the entire history purge_response = self._delete(f"histories/{self.test_history_id}", data={"purge": True}, json=True) self._assert_status_code_is_ok(purge_response) + purge_result = purge_response.json() # Verify history is purged - history_response = self._get(f"histories/{self.test_history_id}").json() - assert history_response["purged"] - assert history_response["deleted"] + assert purge_result["purged"] + assert purge_result["deleted"] - # Verify HDAs are marked as purged + # Verify the response contains a purge_task with an id + assert "purge_task" in purge_result + purge_task = purge_result["purge_task"] + assert "id" in purge_task + purge_task_id = purge_task["id"] + + # Wait for the celery task to complete via the tasks API + self.dataset_populator.wait_on_task_id(purge_task_id) + + # After task completion, HDAs should be purged and files deleted self.dataset_populator.wait_for_purge(self.test_history_id, hda1_id) self.dataset_populator.wait_for_purge(self.test_history_id, hda2_id) - - # Wait for underlying dataset files to be removed from disk. - # With batched history purging, HDAs are marked purged synchronously - # but the actual file deletion happens via a batched celery task. - self._wait_for_file_deleted(dataset_file1) - self._wait_for_file_deleted(dataset_file2) + assert not self._file_exists_on_disk(dataset_file1) + assert not self._file_exists_on_disk(dataset_file2) def _get_underlying_dataset_on_disk(self, hda_id: str) -> Optional[str]: detailed_response = self._get(f"datasets/{hda_id}", admin=True).json() @@ -117,9 +119,3 @@ class TestPurgeDatasetsIntegration(integration_util.IntegrationTestCase): def _file_exists_on_disk(self, filename: Optional[str]) -> bool: return os.path.isfile(filename) if filename else False - - def _wait_for_file_deleted(self, filename: Optional[str], timeout: int = 10): - def _check(): - return True if not self._file_exists_on_disk(filename) else None - - wait_on(_check, f"file {filename} to be deleted", timeout=timeout) From c1d0065fb7fb5fb5e2ce0e0de250a9eefd4541db Mon Sep 17 00:00:00 2001 From: Nicola Soranzo Date: Mon, 23 Mar 2026 00:31:40 +0000 Subject: [PATCH 11/13] Sync gxformat2 in packages with pin in pyproject.toml --- packages/app/setup.cfg | 2 +- packages/test_base/setup.cfg | 2 +- packages/test_selenium/setup.cfg | 2 +- packages/web_apps/setup.cfg | 2 +- 4 files changed, 4 insertions(+), 4 deletions(-) diff --git a/packages/app/setup.cfg b/packages/app/setup.cfg index 8b5be408e75..757e32ab5f7 100644 --- a/packages/app/setup.cfg +++ b/packages/app/setup.cfg @@ -51,7 +51,7 @@ install_requires = cloudauthz==0.6.0 cwl-utils dparse - gxformat2 + gxformat2>=0.23.0 kombu>=5.3 lagom lxml!=4.2.2 diff --git a/packages/test_base/setup.cfg b/packages/test_base/setup.cfg index 6763eeacd71..81d38b817af 100644 --- a/packages/test_base/setup.cfg +++ b/packages/test_base/setup.cfg @@ -35,7 +35,7 @@ install_requires = galaxy-util bioblend cwltest>=2.5.20240906231108 - gxformat2 + gxformat2>=0.23.0 pytest pytest-celery PyYAML diff --git a/packages/test_selenium/setup.cfg b/packages/test_selenium/setup.cfg index d78fbf2e036..649a2b6d4dd 100644 --- a/packages/test_selenium/setup.cfg +++ b/packages/test_selenium/setup.cfg @@ -36,7 +36,7 @@ install_requires = galaxy-selenium galaxy-test-base galaxy-util - gxformat2 + gxformat2>=0.23.0 pytest PyYAML requests diff --git a/packages/web_apps/setup.cfg b/packages/web_apps/setup.cfg index 971ef27e1c2..0ca4bb912a0 100644 --- a/packages/web_apps/setup.cfg +++ b/packages/web_apps/setup.cfg @@ -45,7 +45,7 @@ install_requires = fastapi>=0.124.0 gunicorn httpx - gxformat2 + gxformat2>=0.23.0 Mako MarkupSafe Paste From 2d27004b23c52fdcd14bec0c6c4d69102208d948 Mon Sep 17 00:00:00 2001 From: Nicola Soranzo Date: Mon, 23 Mar 2026 00:32:56 +0000 Subject: [PATCH 12/13] Update fastmcp requirement to 3.0.2 and update its dependencies. This is needed to drop the requirement on redis 7.1.0, which conflicts with our conditional requirement `redis>=5.3.0,<6` , which is needed for compatibility with celery 5.6.2 . Fix https://github.com/galaxyproject/galaxy/issues/22201 . --- .../dependencies/pinned-requirements.txt | 22 +++++-------------- 1 file changed, 6 insertions(+), 16 deletions(-) diff --git a/lib/galaxy/dependencies/pinned-requirements.txt b/lib/galaxy/dependencies/pinned-requirements.txt index d83edca597f..e3b54823bb4 100644 --- a/lib/galaxy/dependencies/pinned-requirements.txt +++ b/lib/galaxy/dependencies/pinned-requirements.txt @@ -4,6 +4,7 @@ a2wsgi==1.10.10 adal==1.2.7 ag-ui-protocol==0.1.10 aiobotocore==3.1.1 +aiofile==3.9.0 aiofiles==25.1.0 aiohappyeyeballs==2.6.1 aiohttp==3.13.3 @@ -19,7 +20,7 @@ apispec==6.9.0 appdirs==1.4.4 arcp==0.2.1 argcomplete==3.6.3 -async-timeout==5.0.1 ; python_full_version < '3.11.3' +async-timeout==5.0.1 ; python_full_version < '3.11' attmap==0.13.2 attrs==25.4.0 authlib==1.6.6 @@ -41,6 +42,7 @@ botocore==1.42.30 bx-python==0.14.0 cachecontrol==0.14.4 cachetools==6.2.4 +caio==0.9.25 celery==5.6.2 certifi==2026.1.4 cffi==2.0.0 ; implementation_name == 'pypy' or platform_python_implementation != 'PyPy' @@ -52,7 +54,6 @@ click-plugins==1.1.1.2 click-repl==0.3.0 cloudauthz==0.6.0 cloudbridge==3.2.0 -cloudpickle==3.1.2 cohere==5.20.2 ; sys_platform != 'emscripten' colorama==0.4.6 coloredlogs==15.0.1 @@ -67,7 +68,6 @@ defusedxml==0.7.1 deprecated==1.3.1 deprecation==2.1.0 dictobj==0.4 -diskcache==5.6.3 distro==1.9.0 dnspython==2.8.0 docopt==0.6.2 @@ -81,10 +81,9 @@ et-xmlfile==2.0.0 eval-type-backport==0.3.1 exceptiongroup==1.3.1 executing==2.2.1 -fakeredis==2.33.0 fastapi==0.128.0 fastavro==1.12.1 ; sys_platform != 'emscripten' -fastmcp==2.14.4 +fastmcp==3.0.2 filelock==3.20.3 fissix==24.4.24 frozenlist==1.8.0 @@ -140,7 +139,6 @@ legacy-cgi==2.6.4 ; python_full_version >= '3.13' limits==5.6.0 logfire==4.19.0 logfire-api==4.19.0 -lupa==2.6 lxml==6.0.2 mako==1.3.10 markdown==3.10.1 @@ -173,7 +171,6 @@ openpyxl==3.1.5 opentelemetry-api==1.39.1 opentelemetry-exporter-otlp-proto-common==1.39.1 opentelemetry-exporter-otlp-proto-http==1.39.1 -opentelemetry-exporter-prometheus==0.60b1 opentelemetry-instrumentation==0.60b1 opentelemetry-instrumentation-httpx==0.60b1 opentelemetry-proto==1.39.1 @@ -188,11 +185,9 @@ parsley==1.3 paste==3.10.1 pastedeploy==3.1.0 pathable==0.4.4 -pathvalidate==3.3.1 pebble==5.1.3 pillow==12.1.0 platformdirs==4.5.1 -prometheus-client==0.24.1 prompt-toolkit==3.0.52 propcache==0.4.1 proto-plus==1.27.0 @@ -200,8 +195,7 @@ protobuf==6.33.4 prov==1.5.1 psutil==7.2.1 pulsar-galaxy-lib==0.15.14 -py-key-value-aio==0.3.0 -py-key-value-shared==0.3.0 +py-key-value-aio==0.4.4 pyasn1==0.6.2 pyasn1-modules==0.4.2 pycparser==3.0 ; (implementation_name != 'PyPy' and platform_python_implementation != 'PyPy') or (implementation_name == 'pypy' and platform_python_implementation == 'PyPy') @@ -215,7 +209,6 @@ pydantic-graph==1.44.0 pydantic-settings==2.12.0 pydantic-tes==0.2.0 pydicom==3.0.1 -pydocket==0.16.6 pydot==4.0.1 pyeventsystem==0.1.0 pyfaidx==0.9.0.3 @@ -231,7 +224,6 @@ pyreadline3==3.5.4 ; sys_platform == 'win32' pysam==0.23.3 python-dateutil==2.9.0.post0 python-dotenv==1.2.1 -python-json-logger==4.0.0 python-magic==0.4.27 python-multipart==0.0.22 python-slugify==8.0.4 @@ -242,7 +234,6 @@ pywin32-ctypes==0.2.3 ; sys_platform == 'win32' pyyaml==6.0.3 pyzmq==27.1.0 rdflib==7.5.0 -redis==7.1.0 referencing==0.36.2 refgenconf==0.12.2 regex==2026.1.15 @@ -265,7 +256,6 @@ s3transfer==0.16.0 schema-salad==8.9.20251102115403 secretstorage==3.5.0 ; sys_platform == 'linux' setuptools==80.10.1 -shellingham==1.5.4 six==1.17.0 slowapi==0.1.9 sniffio==1.3.1 @@ -292,7 +282,6 @@ tornado==6.5.4 tqdm==4.67.1 tuspy==1.1.0 tuspyserver==4.2.3 -typer==0.21.1 types-protobuf==6.32.1.20251210 types-requests==2.32.4.20260107 ; sys_platform != 'emscripten' typing-extensions==4.15.0 @@ -304,6 +293,7 @@ urllib3==2.6.3 uvicorn==0.40.0 uvloop==0.22.1 vine==5.1.0 +watchfiles==1.1.1 wcwidth==0.3.2 webencodings==0.5.1 webob==1.8.9 From 63954cb42e5d8da7c4ccc7ec01f240a6efcc5253 Mon Sep 17 00:00:00 2001 From: Nicola Soranzo Date: Mon, 23 Mar 2026 01:32:52 +0000 Subject: [PATCH 13/13] Drop use of deprecated ``method`` parameter of starlette ``FileResponse`` Removed in starlette 1.0, see https://github.com/Kludex/starlette/pull/3147 Fix the following mypy errors in package tests: ``` galaxy/webapps/base/api.py:94: error: Too many arguments for "__init__" of "FileResponse" [call-arg] super().__init__( ^~~~~~~~~~~~~~~~~ galaxy/webapps/base/api.py:95: error: Argument 8 to "__init__" of "FileResponse" has incompatible type "str | None"; expected "str" [arg-type] ..., media_type, background, filename, stat_result, method, content_dispo... ^~~~~~ ``` --- lib/galaxy/webapps/base/api.py | 10 ++++++++-- lib/galaxy/webapps/galaxy/api/datasets.py | 4 ++-- 2 files changed, 10 insertions(+), 4 deletions(-) diff --git a/lib/galaxy/webapps/base/api.py b/lib/galaxy/webapps/base/api.py index 1cc03d05141..01127f21fbb 100644 --- a/lib/galaxy/webapps/base/api.py +++ b/lib/galaxy/webapps/base/api.py @@ -88,11 +88,17 @@ class GalaxyFileResponse(FileResponse): background: Optional["BackgroundTask"] = None, filename: Optional[str] = None, stat_result: Optional[os.stat_result] = None, - method: Optional[str] = None, content_disposition_type: str = "attachment", ) -> None: super().__init__( - path, status_code, headers, media_type, background, filename, stat_result, method, content_disposition_type + path=path, + status_code=status_code, + headers=headers, + media_type=media_type, + background=background, + filename=filename, + stat_result=stat_result, + content_disposition_type=content_disposition_type, ) self.headers["accept-ranges"] = "bytes" self.xsendfile = self.nginx_x_accel_redirect_base or self.apache_xsendfile diff --git a/lib/galaxy/webapps/galaxy/api/datasets.py b/lib/galaxy/webapps/galaxy/api/datasets.py index b207a3ca350..35488d66eed 100644 --- a/lib/galaxy/webapps/galaxy/api/datasets.py +++ b/lib/galaxy/webapps/galaxy/api/datasets.py @@ -290,7 +290,7 @@ class FastAPIDatasets: assert isinstance(display_data, IOBase) file_name = getattr(display_data, "name", None) assert file_name - return GalaxyFileResponse(file_name, headers=headers, method=request.method) + return GalaxyFileResponse(file_name, headers=headers) @router.get( "/api/histories/{history_id}/contents/{history_content_id}/display", @@ -374,7 +374,7 @@ class FastAPIDatasets: if isinstance(display_data, IOBase): file_name = getattr(display_data, "name", None) if file_name: - return GalaxyFileResponse(file_name, headers=headers, method=request.method) + return GalaxyFileResponse(file_name, headers=headers) elif isinstance(display_data, ZipstreamWrapper): return StreamingResponse(display_data.response(), headers=headers) elif isinstance(display_data, bytes):