diff --git a/client/src/api/schema/schema.ts b/client/src/api/schema/schema.ts index d4f35364efb..1c508ca678f 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/client/visualizations.yml b/client/visualizations.yml index 6a635622e5c..9db3485a4f7 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 diff --git a/lib/galaxy/celery/tasks.py b/lib/galaxy/celery/tasks.py index 03152c46dd7..e8a7dfb4958 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(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/dependencies/pinned-requirements.txt b/lib/galaxy/dependencies/pinned-requirements.txt index 23382f28b9e..62e0c0af2d4 100644 --- a/lib/galaxy/dependencies/pinned-requirements.txt +++ b/lib/galaxy/dependencies/pinned-requirements.txt @@ -102,7 +102,7 @@ groq==1.1.1 grpcio==1.78.0 grpcio-status==1.78.0 gunicorn==25.0.3 -gxformat2==0.22.0 +gxformat2==0.23.0 h11==0.16.0 h5grove==3.0.0 h5py==3.16.0 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/managers/histories.py b/lib/galaxy/managers/histories.py index cc9aa850ded..053d456e511 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,12 +293,21 @@ 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 + 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/model/__init__.py b/lib/galaxy/model/__init__.py index bbc84c5c494..05756dc5295 100644 --- a/lib/galaxy/model/__init__.py +++ b/lib/galaxy/model/__init__.py @@ -1845,8 +1845,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 @@ -4730,14 +4743,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 diff --git a/lib/galaxy/model/store/__init__.py b/lib/galaxy/model/store/__init__.py index 56ddceae710..0d840342c25 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 diff --git a/lib/galaxy/schema/schema.py b/lib/galaxy/schema/schema.py index 2dbab68380a..85785b8abab 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): @@ -3995,6 +4000,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/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/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/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): 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/lib/galaxy_test/api/test_workflows.py b/lib/galaxy_test/api/test_workflows.py index 00fd842e866..f253ba5c996 100644 --- a/lib/galaxy_test/api/test_workflows.py +++ b/lib/galaxy_test/api/test_workflows.py @@ -1945,6 +1945,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") diff --git a/packages/app/setup.cfg b/packages/app/setup.cfg index 8bd63f4718e..aa633cf0d9a 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 a26989aaa36..64efc4ce46f 100644 --- a/packages/web_apps/setup.cfg +++ b/packages/web_apps/setup.cfg @@ -45,7 +45,7 @@ install_requires = fastapi>=0.124.0 gunicorn!=25.1.0 httpx - gxformat2 + gxformat2>=0.23.0 Mako MarkupSafe Paste diff --git a/pyproject.toml b/pyproject.toml index 5a833ee907b..63778e6bdb3 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!=25.1.0", # https://github.com/benoitc/gunicorn/discussions/3509 - "gxformat2>=0.22.0", + "gxformat2>=0.23.0", "h5grove>=1.2.1", "h5py>=3.12", # Python 3.13 support "httpx", 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_purge_datasets.py b/test/integration/test_purge_datasets.py index 363245397c1..753ec4420fd 100644 --- a/test/integration/test_purge_datasets.py +++ b/test/integration/test_purge_datasets.py @@ -76,6 +76,43 @@ 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) + purge_result = purge_response.json() + + # Verify history is purged + assert purge_result["purged"] + assert purge_result["deleted"] + + # 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) + 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() return detailed_response.get("file_name") 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(): 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