mirror of
https://github.com/galaxyproject/galaxy.git
synced 2026-09-24 16:30:27 +08:00
Merge branch 'release_26.0' into dev
This commit is contained in:
@@ -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.
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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")
|
||||
|
||||
|
||||
|
||||
@@ -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."""
|
||||
|
||||
|
||||
@@ -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():
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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):
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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")
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -35,7 +35,7 @@ install_requires =
|
||||
galaxy-util
|
||||
bioblend
|
||||
cwltest>=2.5.20240906231108
|
||||
gxformat2
|
||||
gxformat2>=0.23.0
|
||||
pytest
|
||||
pytest-celery
|
||||
PyYAML
|
||||
|
||||
@@ -36,7 +36,7 @@ install_requires =
|
||||
galaxy-selenium
|
||||
galaxy-test-base
|
||||
galaxy-util
|
||||
gxformat2
|
||||
gxformat2>=0.23.0
|
||||
pytest
|
||||
PyYAML
|
||||
requests
|
||||
|
||||
@@ -45,7 +45,7 @@ install_requires =
|
||||
fastapi>=0.124.0
|
||||
gunicorn!=25.1.0
|
||||
httpx
|
||||
gxformat2
|
||||
gxformat2>=0.23.0
|
||||
Mako
|
||||
MarkupSafe
|
||||
Paste
|
||||
|
||||
+1
-1
@@ -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",
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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")
|
||||
|
||||
@@ -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():
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user