diff --git a/client/src/components/Tool/ToolCard.vue b/client/src/components/Tool/ToolCard.vue index 21c858604c7..f5cae8303f3 100644 --- a/client/src/components/Tool/ToolCard.vue +++ b/client/src/components/Tool/ToolCard.vue @@ -117,8 +117,16 @@ const { isOnlyPreference } = useStorageLocationConfiguration(); const { currentUser, isAnonymous } = storeToRefs(useUserStore()); const { isLoaded: isConfigLoaded, config } = storeToRefs(useConfigStore()); const hasUser = computed(() => !isAnonymous.value); -const versions = computed(() => props.options.versions); -const showVersions = computed(() => props.options.versions?.length > 1); +const versions = computed(() => props.options.versions ?? []); +const hiddenVersions = computed(() => props.options.hidden_versions ?? []); +const visibleVersions = computed(() => { + const filtered = versions.value.filter((v) => !hiddenVersions.value.includes(v)); + if (props.version && !filtered.includes(props.version) && versions.value.includes(props.version)) { + filtered.push(props.version); + } + return filtered; +}); +const showVersions = computed(() => visibleVersions.value.length > 1); const storageLocationModalTitle = computed(() => { if (isOnlyPreference.value) { @@ -165,7 +173,7 @@ onBeforeMount(() => { { history_id: "historyId", targets: [ { - auto_decompress: false, + auto_decompress: true, destination: { type: "hdas" }, elements: [ { - auto_decompress: false, + auto_decompress: true, dbkey: "?", deferred: true, ext: "auto", @@ -236,7 +236,7 @@ describe("UploadQueue", () => { url: "http://test.me.0", }, { - auto_decompress: false, + auto_decompress: true, dbkey: "?", deferred: true, ext: "auto", @@ -247,7 +247,7 @@ describe("UploadQueue", () => { url: "http://test.me.1", }, { - auto_decompress: false, + auto_decompress: true, dbkey: "?", deferred: true, ext: "auto", diff --git a/client/src/utils/upload.test.ts b/client/src/utils/upload.test.ts index d3de5c6ad67..d72481591d3 100644 --- a/client/src/utils/upload.test.ts +++ b/client/src/utils/upload.test.ts @@ -427,6 +427,13 @@ describe("createUrlUploadItem", () => { expect(item.deferred).toBe(true); expect(item.dbkey).toBe("hg38"); }); + + test("trims surrounding whitespace from URL", () => { + const item = createUrlUploadItem(" http://example.com/data.bed\n", "historyId"); + + expect(item.url).toBe("http://example.com/data.bed"); + expect(item.name).toBe("data.bed"); + }); }); describe("parseContentToUploadItems", () => { @@ -462,6 +469,22 @@ describe("parseContentToUploadItems", () => { ); }); + test("throws on network URL with empty DNS labels", () => { + expect(() => parseContentToUploadItems("https://.../SRR1957099.fastq.gz", "historyId")).toThrow( + "Invalid URL: https://.../SRR1957099.fastq.gz", + ); + }); + + test("accepts Galaxy file-source URIs", () => { + const content = "gxfiles://myftp/file.txt\ndrs://example.org/abc\nzenodo://record/123"; + const items = parseContentToUploadItems(content, "historyId"); + + expect(items).toHaveLength(3); + expect((items[0] as { url: string }).url).toBe("gxfiles://myftp/file.txt"); + expect((items[1] as { url: string }).url).toBe("drs://example.org/abc"); + expect((items[2] as { url: string }).url).toBe("zenodo://record/123"); + }); + test("handles whitespace around URLs", () => { const items = parseContentToUploadItems(" http://example.com/file.txt \n ", "historyId"); @@ -537,6 +560,12 @@ describe("buildUploadPayload", () => { expect(() => buildUploadPayload(items)).toThrow("Invalid URL: not-a-valid-url"); }); + + test("rejects network URL with empty DNS labels", () => { + const items: ApiUploadItem[] = [createUrlUploadItem("https://.../SRR1957099.fastq.gz", "historyId")]; + + expect(() => buildUploadPayload(items)).toThrow("Invalid URL: https://.../SRR1957099.fastq.gz"); + }); }); // ============================================================================ diff --git a/client/src/utils/upload.ts b/client/src/utils/upload.ts index aa1a6fe28d1..e40f6c4fd61 100644 --- a/client/src/utils/upload.ts +++ b/client/src/utils/upload.ts @@ -58,7 +58,7 @@ import type { SupportedCollectionType, UploadCollectionConfig } from "@/composab import type { NewUploadItem } from "@/composables/upload/uploadItemTypes"; import { getAppRoot } from "@/onload/loadConfig"; import { errorMessageAsString } from "@/utils/simple-error"; -import { isUrl } from "@/utils/url"; +import { isUrl, isValidUrl } from "@/utils/url"; import { createTusUpload, type FileStream, type NamedBlob, type UploadableFile } from "./tusUpload"; @@ -127,6 +127,8 @@ interface UploadItemCommon { deferred: boolean; /** Optional hash values for verification */ hashes?: FetchDatasetHash[]; + /** Whether to auto-decompress the upload */ + auto_decompress: boolean; } /** Upload item from a local file */ @@ -349,6 +351,7 @@ export function createFileUploadItem( to_posix_lines: options.to_posix_lines ?? uploadItemDefaults.to_posix_lines, deferred: options.deferred ?? uploadItemDefaults.deferred, hashes: options.hashes, + auto_decompress: true, }; } @@ -385,6 +388,7 @@ export function createPastedUploadItem( to_posix_lines: options.to_posix_lines ?? uploadItemDefaults.to_posix_lines, deferred: options.deferred ?? uploadItemDefaults.deferred, hashes: options.hashes, + auto_decompress: true, }; } @@ -410,12 +414,13 @@ export function createUrlUploadItem( historyId: string, options: Partial> = {}, ): UrlUploadItem { + const trimmedUrl = url.trim(); // Extract filename from URL if not provided - const defaultName = url.split("/").pop()?.split("?")[0] || DEFAULT_FILE_NAME; + const defaultName = trimmedUrl.split("/").pop()?.split("?")[0] || DEFAULT_FILE_NAME; return { src: "url", - url, + url: trimmedUrl, historyId, name: options.name ?? defaultName, size: options.size ?? 0, @@ -425,6 +430,7 @@ export function createUrlUploadItem( to_posix_lines: options.to_posix_lines ?? uploadItemDefaults.to_posix_lines, deferred: options.deferred ?? uploadItemDefaults.deferred, hashes: options.hashes, + auto_decompress: true, }; } @@ -449,6 +455,7 @@ export function toApiUploadItem(item: NewUploadItem): ApiUploadItem { to_posix_lines: item.toPosixLines, deferred: item.deferred, hashes: item.hashes, + auto_decompress: true, }; switch (item.uploadMode) { @@ -510,7 +517,7 @@ export function parseContentToUploadItems( // If first line is a URL, treat all lines as URLs if (isUrl(firstLine)) { return lines.filter(Boolean).map((urlLine) => { - if (!isUrl(urlLine)) { + if (!isValidUrl(urlLine)) { throw new Error(`Invalid URL: ${urlLine}`); } return createUrlUploadItem(urlLine, historyId, options); @@ -536,7 +543,7 @@ function buildDataElement(item: ApiUploadItem): ApiDataElement { name: normalizeFileName(item.name), space_to_tab: item.space_to_tab, to_posix_lines: item.to_posix_lines, - auto_decompress: false, + auto_decompress: true, deferred: item.deferred, }; @@ -598,7 +605,7 @@ function validateItemContent(item: ApiUploadItem): void { if (!item.url || item.url.trim().length === 0) { throw new Error(`No URL for upload item: ${item.name}`); } - if (!isUrl(item.url)) { + if (!isValidUrl(item.url)) { throw new Error(`Invalid URL: ${item.url}`); } break; @@ -674,7 +681,7 @@ export function buildUploadPayload(items: ApiUploadItem[], options: BuildPayload history_id: historyId, targets: [ { - auto_decompress: false, + auto_decompress: true, destination: { type: "hdas" }, elements, }, diff --git a/client/src/utils/upload/itemMappers.ts b/client/src/utils/upload/itemMappers.ts index 27ed35b039f..0a434b81ff4 100644 --- a/client/src/utils/upload/itemMappers.ts +++ b/client/src/utils/upload/itemMappers.ts @@ -71,7 +71,7 @@ export function mapToPasteUrlUpload(item: PasteUrlItem, targetHistoryId: string) spaceToTab: item.spaceToTab, toPosixLines: item.toPosixLines, deferred: item.deferred, - url: item.url, + url: item.url.trim(), }; } diff --git a/client/src/utils/url.test.js b/client/src/utils/url.test.js index aef87258cf1..fe8deef1bd0 100644 --- a/client/src/utils/url.test.js +++ b/client/src/utils/url.test.js @@ -1,6 +1,6 @@ import { describe, expect, it, test } from "vitest"; -import { addSearchParams, isUrl } from "./url"; +import { addSearchParams, isUrl, isValidNetworkUrl, isValidUrl } from "./url"; describe("test url utilities", () => { it("adding parameters to url", async () => { @@ -14,4 +14,59 @@ describe("test url utilities", () => { expect(isUrl("ftp://")).toBeTruthy(); expect(isUrl("http://")).toBeTruthy(); }); + + test("network url validation accepts well-formed hosts", () => { + expect(isValidNetworkUrl("https://example.com/x")).toBe(true); + expect(isValidNetworkUrl("http://sub.example.org/path?q=1")).toBe(true); + expect(isValidNetworkUrl("ftp://ftp.example.org/a")).toBe(true); + }); + + test("network url validation rejects empty DNS labels", () => { + expect(isValidNetworkUrl("https://.../foo")).toBe(false); + expect(isValidNetworkUrl("https://./foo")).toBe(false); + expect(isValidNetworkUrl("https://..example.com/x")).toBe(false); + expect(isValidNetworkUrl("https:///foo")).toBe(false); + }); + + test("network url validation rejects non-network schemes", () => { + expect(isValidNetworkUrl("gxfiles://foo")).toBe(false); + expect(isValidNetworkUrl("drs://example.org/123")).toBe(false); + expect(isValidNetworkUrl("not-a-url")).toBe(false); + }); + + test("isValidUrl accepts custom Galaxy file-source schemes", () => { + expect(isValidUrl("gxfiles://myftp/file.txt")).toBe(true); + expect(isValidUrl("gxuserfiles://mysource/x")).toBe(true); + expect(isValidUrl("drs://example.org/abc")).toBe(true); + expect(isValidUrl("zenodo://record/123")).toBe(true); + expect(isValidUrl("invenio://record/123")).toBe(true); + expect(isValidUrl("dataverse://doi/x")).toBe(true); + expect(isValidUrl("base64://aGVsbG8=")).toBe(true); + }); + + test("isValidUrl accepts well-formed network urls", () => { + expect(isValidUrl("https://example.com/")).toBe(true); + expect(isValidUrl("https://example.com/file.txt")).toBe(true); + expect(isValidUrl("ftp://ftp.example.org/a")).toBe(true); + }); + + test("isValidUrl rejects network urls with empty DNS labels", () => { + expect(isValidUrl("https://.../foo")).toBe(false); + expect(isValidUrl("https://./foo")).toBe(false); + expect(isValidUrl("https:///foo")).toBe(false); + }); + + test("isValidUrl rejects unknown schemes and non-urls", () => { + expect(isValidUrl("")).toBe(false); + expect(isValidUrl("not-a-url")).toBe(false); + expect(isValidUrl("xyz://foo")).toBe(false); + }); + + test("isValidUrl tolerates surrounding whitespace", () => { + expect(isValidUrl(" https://example.com/x ")).toBe(true); + expect(isValidUrl("\thttps://example.com/x\n")).toBe(true); + expect(isValidUrl(" gxfiles://myftp/file.txt ")).toBe(true); + expect(isValidUrl(" ")).toBe(false); + expect(isValidUrl(" https://.../x ")).toBe(false); + }); }); diff --git a/client/src/utils/url.ts b/client/src/utils/url.ts index 2c7b4818ffc..12af02f09f0 100644 --- a/client/src/utils/url.ts +++ b/client/src/utils/url.ts @@ -38,6 +38,44 @@ export function isUrl(content: string): boolean { return URI_PREFIXES.some((prefix) => content.startsWith(prefix)); } +const NETWORK_SCHEMES = ["http://", "https://", "ftp://", "ftps://"]; + +export function isNetworkUrl(content: string): boolean { + return NETWORK_SCHEMES.some((prefix) => content.startsWith(prefix)); +} + +/** + * Structural validation for http(s)/ftp(s) URLs destined for the server + * fetch endpoint. Rejects empty DNS labels (leading/trailing dot, consecutive + * dots, or pure punctuation) that would otherwise crash server-side idna + * encoding and surface as an opaque 500. + */ +export function isValidNetworkUrl(content: string): boolean { + const prefix = NETWORK_SCHEMES.find((p) => content.startsWith(p)); + if (!prefix) { + return false; + } + const afterScheme = content.slice(prefix.length); + const authorityEnd = afterScheme.search(/[/?#]/); + const authority = authorityEnd === -1 ? afterScheme : afterScheme.slice(0, authorityEnd); + const hostAndPort = authority.includes("@") ? authority.slice(authority.lastIndexOf("@") + 1) : authority; + let host = hostAndPort; + if (!host.startsWith("[")) { + const portIdx = host.lastIndexOf(":"); + if (portIdx !== -1) { + host = host.slice(0, portIdx); + } + } + if (!host) { + return false; + } + // IPv6 literals are bracketed; they have no DNS labels to validate. + if (host.startsWith("[")) { + return true; + } + return !host.split(".").some((label) => label.length === 0); +} + export async function urlData({ url, headers, params, errorSimplify = true }: UrlDataOptions): Promise { try { headers = headers || {}; @@ -76,41 +114,24 @@ export interface UrlValidationResult { } /** - * Validates a URL string and returns detailed validation results. - * Checks for required fields, valid hostname, and meaningful content. - * - * @param url - The URL string to validate - * @returns Validation result with isValid flag and error message - * - * @example - * ```typescript - * const result = validateUrl("https://example.org/data.txt"); - * if (!result.isValid) { - * console.error(result.message); - * } - * ``` + * Validates a URL string for use in upload flows. Accepts any URI whose scheme + * appears in URI_PREFIXES (including Galaxy file-source schemes like + * gxfiles://, drs://, zenodo://, etc.). For http(s)/ftp(s) URLs it also applies + * a structural host check to reject empty DNS labels that would crash + * server-side idna encoding. */ export function validateUrl(url: string): UrlValidationResult { - if (!url.trim()) { + const trimmed = url.trim(); + if (!trimmed) { return { isValid: false, message: "URL is required" }; } - - try { - const urlObj = new URL(url); - - if (!urlObj.host) { - return { isValid: false, message: "URL must include a valid hostname" }; - } - - const hasContent = urlObj.pathname !== "/" || urlObj.search || urlObj.hash; - if (!hasContent) { - return { isValid: false, message: "URL should point to a specific file or resource" }; - } - - return { isValid: true, message: null }; - } catch { + if (!isUrl(trimmed)) { return { isValid: false, message: "URL format is invalid or missing protocol" }; } + if (isNetworkUrl(trimmed) && !isValidNetworkUrl(trimmed)) { + return { isValid: false, message: "URL must include a valid hostname" }; + } + return { isValid: true, message: null }; } /** diff --git a/client/visualizations.yml b/client/visualizations.yml index 4c72599c237..d687def32e6 100644 --- a/client/visualizations.yml +++ b/client/visualizations.yml @@ -71,7 +71,7 @@ nora: version: 1.2.4 openlayers: package: "@galaxyproject/openlayers" - version: 0.0.8 + version: 0.0.9 openseadragon: package: "@galaxyproject/openseadragon" version: 0.0.2 diff --git a/lib/galaxy/app/__init__.py b/lib/galaxy/app/__init__.py index 23a2c4332d8..e8b8ef1b941 100644 --- a/lib/galaxy/app/__init__.py +++ b/lib/galaxy/app/__init__.py @@ -70,6 +70,7 @@ from galaxy.managers.jobs import ( JobManager as JobQueryManager, JobSearch, ) +from galaxy.managers.landing import LandingRequestManager from galaxy.managers.libraries import LibraryManager from galaxy.managers.library_datasets import LibraryDatasetsManager from galaxy.managers.notification import NotificationManager @@ -367,7 +368,7 @@ class MinimalGalaxyApplication(BasicSharedApp, HaltableContainer, SentryClientMi self.config.tool_configs.append(self.config.migrated_tools_config) def _configure_toolbox(self): - self.citations_manager = CitationsManager(self) + self.citations_manager = self._register_singleton(CitationsManager, CitationsManager(self)) self.biotools_metadata_source = get_galaxy_biotools_metadata_source(self.config) self.dynamic_tool_manager = DynamicToolManager(self) @@ -647,6 +648,7 @@ class GalaxyManagerApplication(MinimalManagerApp, MinimalGalaxyApplication): self.dataset_collection_manager = self._register_singleton(DatasetCollectionManager) self.workflow_manager = self._register_singleton(WorkflowsManager) self.workflow_contents_manager = self._register_singleton(WorkflowContentsManager) + self.landing_request_manager = self._register_singleton(LandingRequestManager) self.library_folder_manager = self._register_singleton(FolderManager) self.library_manager = self._register_singleton(LibraryManager) self.library_datasets_manager = self._register_singleton(LibraryDatasetsManager) diff --git a/lib/galaxy/dependencies/dev-requirements.txt b/lib/galaxy/dependencies/dev-requirements.txt index 312146774ac..0af4545b5a4 100644 --- a/lib/galaxy/dependencies/dev-requirements.txt +++ b/lib/galaxy/dependencies/dev-requirements.txt @@ -1,5 +1,6 @@ # This file was autogenerated by uv via the following command: # uv export --frozen --no-annotate --no-hashes --only-group=dev +aiocop==1.1.4 aiohappyeyeballs==2.6.1 aiohttp==3.13.5 aiosignal==1.4.0 diff --git a/lib/galaxy/dependencies/pinned-test-requirements.txt b/lib/galaxy/dependencies/pinned-test-requirements.txt index 5d2f1c188ea..1f56f7ac90a 100644 --- a/lib/galaxy/dependencies/pinned-test-requirements.txt +++ b/lib/galaxy/dependencies/pinned-test-requirements.txt @@ -1,5 +1,6 @@ # This file was autogenerated by uv via the following command: # uv export --frozen --no-annotate --no-hashes --only-group=test +aiocop==1.1.4 aiohappyeyeballs==2.6.1 aiohttp==3.13.5 aiosignal==1.4.0 diff --git a/lib/galaxy/files/__init__.py b/lib/galaxy/files/__init__.py index cf21e294c03..a6a4bbef3e7 100644 --- a/lib/galaxy/files/__init__.py +++ b/lib/galaxy/files/__init__.py @@ -199,6 +199,7 @@ class ConfiguredFileSources: def get_file_source_path(self, uri): """Parse uri into a FileSource object and a path relative to its base.""" + uri = uri.strip() if "://" not in uri: raise exceptions.RequestParameterInvalidException(f"Invalid uri [{uri}]") file_source = self.find_best_match(uri) diff --git a/lib/galaxy/files/uris.py b/lib/galaxy/files/uris.py index d9455879a47..f565de8dad3 100644 --- a/lib/galaxy/files/uris.py +++ b/lib/galaxy/files/uris.py @@ -94,13 +94,15 @@ def split_port(parsed_url: str, url: str) -> tuple[str, int]: def validate_non_local(uri: str, ip_allowlist: list[IpAllowedListEntryT]) -> str: # If it doesn't look like a URL, ignore it. - if not (uri.lstrip().startswith("http://") or uri.lstrip().startswith("https://")): + if not (uri.strip().startswith("http://") or uri.strip().startswith("https://")): return uri - # Strip leading whitespace before passing url to urlparse() - url = uri.lstrip() + # Strip surrounding whitespace before passing url to urlparse() + url = uri.strip() # Extract hostname component parsed_url = urlparse(url).netloc + if not parsed_url: + raise RequestParameterInvalidException(f"Could not verify url '{url}'.") # If credentials are in this URL, we need to strip those. if parsed_url.count("@") > 0: # credentials. @@ -134,8 +136,10 @@ def validate_non_local(uri: str, ip_allowlist: list[IpAllowedListEntryT]) -> str # Call getaddrinfo to resolve hostname into tuples containing IPs. try: addrinfo = socket.getaddrinfo(parsed_url, port) - except socket.gaierror as e: - log.debug(f"Could not resolve url '{url}': {e}") + except (socket.gaierror, UnicodeError) as e: + # UnicodeError covers idna codec failures (e.g. empty DNS labels in hosts like '...' or '..example.com') + # which are not wrapped as socket.gaierror. + log.debug("Could not resolve url '%s': '%s'", url, e) raise RequestParameterInvalidException(f"Could not verify url '{url}'.") # Get the IP addresses that this entry resolves to (uniquely) # We drop: diff --git a/lib/galaxy/managers/collections.py b/lib/galaxy/managers/collections.py index bd3f6f2e59a..2a81db4ee3a 100644 --- a/lib/galaxy/managers/collections.py +++ b/lib/galaxy/managers/collections.py @@ -354,6 +354,7 @@ class DatasetCollectionManager: # else if elements is set, it better be an ordered dict! if elements is not self.ELEMENTS_UNINITIALIZED: + self._validate_nested_collection_elements(collection_type_description, elements) type_plugin = collection_type_description.rank_type_plugin() dataset_collection = builder.build_collection( type_plugin, elements, fields=fields, column_definitions=column_definitions, rows=rows @@ -577,6 +578,32 @@ class DatasetCollectionManager: session.commit() return dataset_collection_instance + def _validate_nested_collection_elements(self, collection_type_description, elements) -> None: + """For nested collection types (e.g. ``list:paired_or_unpaired``), verify that + every element value is a :class:`DatasetCollection` whose ``collection_type`` + matches the expected sub-collection type. Otherwise creating a + ``list:paired_or_unpaired`` from raw HDAs would silently produce a + structurally invalid collection that downstream tools cannot handle. + """ + if not collection_type_description.has_subcollections(): + return + if not isinstance(elements, dict): + return + expected_sub = collection_type_description.subcollection_type_description().collection_type + for identifier, value in elements.items(): + if isinstance(value, DatasetCollection) and value.collection_type == expected_sub: + continue + if isinstance(value, DatasetCollection): + actual = f"a collection of type '{value.collection_type}'" + elif getattr(value, "history_content_type", None) == "dataset": + actual = "a dataset" + else: + actual = type(value).__name__ + raise RequestParameterInvalidException( + f"Element '{identifier}' of collection type '{collection_type_description.collection_type}' " + f"must be a sub-collection of type '{expected_sub}', got {actual}." + ) + def __recursively_create_collections_for_identifiers( self, trans, element_identifiers, hide_source_items: bool, copy_elements: bool, history=None ): diff --git a/lib/galaxy/managers/hdas.py b/lib/galaxy/managers/hdas.py index 4e3c5ebd0f3..bc7095205f6 100644 --- a/lib/galaxy/managers/hdas.py +++ b/lib/galaxy/managers/hdas.py @@ -36,6 +36,7 @@ from galaxy import ( exceptions, model, ) +from galaxy.files import ProvidesFileSourcesUserContext from galaxy.managers import ( annotatable, base, @@ -69,7 +70,6 @@ from galaxy.schema.storage_cleaner import ( from galaxy.schema.tasks import ( MaterializeDatasetInstanceTaskRequest, PurgeDatasetsTaskRequest, - RequestUser, ) from galaxy.structured_app import ( MinimalManagerApp, @@ -81,6 +81,7 @@ from galaxy.tool_util_models.parameters import ( FileRequestUri, ) from galaxy.util.compression_utils import get_fileobj +from galaxy.work.context import WorkRequestContext if TYPE_CHECKING: from galaxy.model import LibraryDatasetDatasetAssociation @@ -175,15 +176,17 @@ class HDAManager( def materialize( self, request: MaterializeDatasetInstanceTaskRequest, session: Session, in_place: bool = False ) -> bool: - request_user: RequestUser = request.user + request_user = request.user + assert request_user.user_id + user = self.user_manager.by_id(request_user.user_id) + user_context = ProvidesFileSourcesUserContext(WorkRequestContext(app=self.app, user=user)) materializer = materializer_factory( True, # attached... object_store=self.app.object_store, file_sources=self.app.file_sources, sa_session=session, + user_context=user_context, ) - assert request_user.user_id - user = self.user_manager.by_id(request_user.user_id) if request.source == DatasetSourceType.hda: dataset_instance: Union[HistoryDatasetAssociation, LibraryDatasetDatasetAssociation] = self.get_accessible( request.content, user diff --git a/lib/galaxy/model/dataset_collections/matching.py b/lib/galaxy/model/dataset_collections/matching.py index 397f60fabcb..600c11892a8 100644 --- a/lib/galaxy/model/dataset_collections/matching.py +++ b/lib/galaxy/model/dataset_collections/matching.py @@ -3,6 +3,7 @@ from typing import Optional from galaxy import exceptions from galaxy.util import bunch from .structure import ( + get_collection, get_structure, leaf, ) @@ -32,11 +33,6 @@ class CollectionsToMatch: return self.collections.items() -def get_child_collection(item): - # item could be HDCA or DCE - return getattr(item, "child_collection", item.collection) - - class MatchingCollections: """Structure holding the result of matching a list of collections together. This class being different than the class above and being @@ -55,8 +51,12 @@ class MatchingCollections: self.action_tuples = {} self.when_values = None - def __attempt_add_to_linked_match(self, input_name, hdca, collection_type_description, subcollection_type): - structure = get_structure(hdca, collection_type_description, leaf_subcollection_type=subcollection_type) + def __attempt_add_to_linked_match( + self, input_name, hdca, child_collection, collection_type_description, subcollection_type + ): + structure = get_structure( + child_collection, collection_type_description, leaf_subcollection_type=subcollection_type + ) if not self.linked_structure: self.linked_structure = structure self.collections[input_name] = hdca @@ -69,7 +69,7 @@ class MatchingCollections: def slice_collections(self): self.linked_structure.when_values = self.when_values - return self.linked_structure.walk_collections(self.collections) + return self.linked_structure.walk_collections({k: get_collection(v) for k, v in self.collections.items()}) def subcollection_mapping_type(self, input_name): return self.subcollection_types[input_name] @@ -90,7 +90,7 @@ class MatchingCollections: def map_over_action_tuples(self, input_name): if input_name not in self.action_tuples: collection_instance = self.collections[input_name] - self.action_tuples[input_name] = get_child_collection(collection_instance).dataset_action_tuples + self.action_tuples[input_name] = get_collection(collection_instance).dataset_action_tuples return self.action_tuples[input_name] def is_mapped_over(self, input_name): @@ -104,17 +104,27 @@ class MatchingCollections: matching_collections = MatchingCollections() for input_key, to_match in sorted(collections_to_match.items()): hdca = to_match.hdca + # Resolve the contained collection: for an HDCA this is + # hdca.collection; for a DCE it is dce.child_collection + # (not dce.collection which is the *parent*). + # Both collection_type_description and get_structure must + # use the same collection so the type and elements agree. + child_collection = get_collection(hdca) collection_type_description = collection_type_descriptions.for_collection_type( - get_child_collection(hdca).collection_type + child_collection.collection_type ) subcollection_type = to_match.subcollection_type if to_match.linked: matching_collections.__attempt_add_to_linked_match( - input_key, hdca, collection_type_description, subcollection_type + input_key, hdca, child_collection, collection_type_description, subcollection_type ) else: - structure = get_structure(hdca, collection_type_description, leaf_subcollection_type=subcollection_type) + structure = get_structure( + child_collection, + collection_type_description, + leaf_subcollection_type=subcollection_type, + ) matching_collections.unlinked_structures.append(structure) return matching_collections diff --git a/lib/galaxy/model/dataset_collections/structure.py b/lib/galaxy/model/dataset_collections/structure.py index 37ef1471f0b..32366409bb5 100644 --- a/lib/galaxy/model/dataset_collections/structure.py +++ b/lib/galaxy/model/dataset_collections/structure.py @@ -1,10 +1,22 @@ """Module for reasoning about structure of and matching hierarchical collections of data.""" -from typing import TYPE_CHECKING +from typing import ( + Optional, + TYPE_CHECKING, + Union, +) + +from galaxy.model import DatasetCollectionElement if TYPE_CHECKING: + from galaxy.model import ( + DatasetCollection, + HistoryDatasetCollectionAssociation, + ) from .type_description import CollectionTypeDescription + CollectionLike = Union[DatasetCollectionElement, "HistoryDatasetCollectionAssociation"] + class Leaf: children_known = True @@ -100,8 +112,8 @@ class Tree(BaseTree): column_definitions=dataset_collection.column_definitions, ) - def walk_collections(self, hdca_dict): - return self._walk_collections(dict_map(lambda hdca: hdca.collection, hdca_dict)) + def walk_collections(self, collection_dict): + return self._walk_collections(collection_dict) def _walk_collections(self, collection_dict): for index, (_identifier, substructure) in enumerate(self.children): @@ -215,20 +227,55 @@ def dict_map(func, input_dict): return {k: func(v) for k, v in input_dict.items()} +def get_collection( + dataset_collection_instance: "CollectionLike", +) -> "DatasetCollection": + """Return the DatasetCollection contained by a collection instance. + + A DatasetCollectionElement has two collection references: + - ``collection``: the **parent** collection this element belongs to + - ``child_collection``: the nested collection this element *contains* + + An HDCA has one: + - ``collection``: the collection it wraps + + This helper returns the *contained* collection in both cases + (child_collection for DCE, collection for HDCA/adapters) and is + intended for callers that still hold a wrapper object and need a + DatasetCollection to pass to ``get_structure`` or ``walk_collections``. + """ + if ( + isinstance(dataset_collection_instance, DatasetCollectionElement) + and dataset_collection_instance.child_collection + ): + return dataset_collection_instance.child_collection + return dataset_collection_instance.collection + + def get_structure( - dataset_collection_instance, collection_type_description: "CollectionTypeDescription", leaf_subcollection_type=None + collection: "DatasetCollection", + collection_type_description: "CollectionTypeDescription", + leaf_subcollection_type: Optional[str] = None, ): + """Build a Tree (or UninitializedTree) describing a collection's shape. + + ``collection_type_description`` controls the depth of the tree: + elements below ``leaf_subcollection_type`` are treated as leaves. + """ if leaf_subcollection_type: - collection_type_description = collection_type_description.effective_collection_type_description( - leaf_subcollection_type - ) - if hasattr(dataset_collection_instance, "child_collection"): - collection_type_description = ( + if not collection_type_description.has_subcollections_of_type(leaf_subcollection_type): + # The described collection IS the leaf subcollection (no deeper + # structure to strip). Don't enumerate its elements; just record + # the type so multiply() can combine it with the mapping structure. + return UninitializedTree( collection_type_description.collection_type_description_factory.for_collection_type( leaf_subcollection_type ) ) - return UninitializedTree(collection_type_description) - - collection = dataset_collection_instance.collection + # Strip the leaf type from the description so it becomes a leaf + # in the resulting tree. E.g. "list:paired" with + # leaf_subcollection_type="paired" → description becomes "list". + collection_type_description = collection_type_description.effective_collection_type_description( + leaf_subcollection_type + ) return Tree.for_dataset_collection(collection, collection_type_description) diff --git a/lib/galaxy/model/dataset_collections/types/sample_sheet.py b/lib/galaxy/model/dataset_collections/types/sample_sheet.py index 3409ce70bda..987a269796d 100644 --- a/lib/galaxy/model/dataset_collections/types/sample_sheet.py +++ b/lib/galaxy/model/dataset_collections/types/sample_sheet.py @@ -1,12 +1,13 @@ -from typing import cast +from typing import ( + cast, + Optional, +) from galaxy.exceptions import RequestParameterMissingException from galaxy.model import DatasetCollectionElement +from galaxy.tool_util_models.sample_sheet import SampleSheetRow from . import BaseDatasetCollectionType -from .sample_sheet_util import ( - OptionalSampleSheetRows, - validate_row, -) +from .sample_sheet_util import validate_row class SampleSheetDatasetCollectionType(BaseDatasetCollectionType): @@ -15,18 +16,21 @@ class SampleSheetDatasetCollectionType(BaseDatasetCollectionType): collection_type = "sample_sheet" def generate_elements(self, dataset_instances, **kwds): - rows = cast(OptionalSampleSheetRows, kwds.get("rows", None)) + rows = cast(Optional[dict[str, Optional[SampleSheetRow]]], kwds.get("rows", None)) column_definitions = kwds.get("column_definitions", None) - if rows is None: - raise RequestParameterMissingException( - "Missing or null parameter 'rows' required for 'sample_sheet' collection types." - ) - if len(dataset_instances) != len(rows): - self._validation_failed("Supplied element do not match 'rows'.") + if not column_definitions: + rows = rows if rows is not None else dict.fromkeys(dataset_instances) + else: + if rows is None: + raise RequestParameterMissingException( + "Missing or null parameter 'rows' required for 'sample_sheet' collection types." + ) + if len(dataset_instances) != len(rows): + self._validation_failed("Supplied element do not match 'rows'.") all_element_identifiers = list(dataset_instances.keys()) for identifier, element in dataset_instances.items(): - columns = rows[identifier] + columns = rows.get(identifier) validate_row(columns, column_definitions, all_element_identifiers) association = DatasetCollectionElement( element=element, diff --git a/lib/galaxy/model/dataset_collections/types/sample_sheet_util.py b/lib/galaxy/model/dataset_collections/types/sample_sheet_util.py index c878d184673..8cb7d5d5fdd 100644 --- a/lib/galaxy/model/dataset_collections/types/sample_sheet_util.py +++ b/lib/galaxy/model/dataset_collections/types/sample_sheet_util.py @@ -95,9 +95,11 @@ def _validate_column_definition(column_definition: SampleSheetColumnDefinition): def validate_row( - row: SampleSheetRow, column_definitions: Optional[SampleSheetColumnDefinitions], element_identifiers: list[str] + row: Optional[SampleSheetRow], + column_definitions: Optional[SampleSheetColumnDefinitions], + element_identifiers: list[str], ): - if column_definitions is None: + if not column_definitions: return if row is None: raise RequestParameterInvalidException( diff --git a/lib/galaxy/model/deferred.py b/lib/galaxy/model/deferred.py index 532edb112a5..e7c9870051b 100644 --- a/lib/galaxy/model/deferred.py +++ b/lib/galaxy/model/deferred.py @@ -16,7 +16,10 @@ from galaxy.datatypes.sniff import ( stream_url_to_file, ) from galaxy.exceptions import ObjectAttributeInvalidException -from galaxy.files import ConfiguredFileSources +from galaxy.files import ( + ConfiguredFileSources, + OptionalUserContext, +) from galaxy.model import ( Dataset, DatasetCollection, @@ -79,18 +82,24 @@ class DatasetInstanceMaterializer: transient_path_mapper: Optional[TransientPathMapper] = None, file_sources: Optional[ConfiguredFileSources] = None, sa_session: Optional[Session] = None, + user_context: OptionalUserContext = None, ): """Constructor for DatasetInstanceMaterializer. If attached is true, these objects should be created in a supplied object store. If not, this class produces transient HDAs with external_filename and external_extra_files_path set. + + ``user_context`` is forwarded to file source operations so that access + controls (``requires_roles`` / ``requires_groups``) are enforced when + materializing from ``gxfiles://`` URIs. """ self._attached = attached self._transient_path_mapper = transient_path_mapper self._object_store_populator = object_store_populator self._file_sources = file_sources self._sa_session = sa_session + self._user_context = user_context self._previously_materialized: dict[int, HistoryDatasetAssociation] = {} def ensure_materialized( @@ -253,7 +262,7 @@ class DatasetInstanceMaterializer: source_uri = target_source.source_uri if source_uri is None: raise Exception("Cannot stream from dataset source without specified source_uri") - path = stream_url_to_file(source_uri, file_sources=self._file_sources) + path = stream_url_to_file(source_uri, file_sources=self._file_sources, user_context=self._user_context) if target_source.hashes: for source_hash in target_source.hashes: _validate_hash(path, source_hash, "downloaded file") @@ -386,6 +395,7 @@ def materializer_factory( transient_directory: Optional[str] = None, file_sources: Optional[ConfiguredFileSources] = None, sa_session: Optional[Session] = None, + user_context: OptionalUserContext = None, ) -> DatasetInstanceMaterializer: if object_store_populator is None and object_store is not None: object_store_populator = ObjectStorePopulator(object_store, None) @@ -397,6 +407,7 @@ def materializer_factory( transient_path_mapper=transient_path_mapper, file_sources=file_sources, sa_session=sa_session, + user_context=user_context, ) diff --git a/lib/galaxy/model/store/__init__.py b/lib/galaxy/model/store/__init__.py index 638dad3a097..2a7cb1386a5 100644 --- a/lib/galaxy/model/store/__init__.py +++ b/lib/galaxy/model/store/__init__.py @@ -3088,7 +3088,7 @@ def get_export_dataset_filename(name: str, ext: str, encoded_id: str, conversion """ Builds a filename for a dataset using its name an extension. """ - base = "".join(c in FILENAME_VALID_CHARS and c or "_" for c in name) + base = "".join(c in FILENAME_VALID_CHARS and c or "_" for c in name)[:150] if not conversion_key: return f"{base}_{encoded_id}.{ext}" else: diff --git a/lib/galaxy/model/store/_bco_convert_utils.py b/lib/galaxy/model/store/_bco_convert_utils.py index 6df2be1a922..28ac3c3491a 100644 --- a/lib/galaxy/model/store/_bco_convert_utils.py +++ b/lib/galaxy/model/store/_bco_convert_utils.py @@ -1,5 +1,3 @@ -import urllib.parse - from galaxy.model import ( Workflow, WorkflowStep, @@ -8,6 +6,7 @@ from galaxy.schema.bco import ( ExecutionDomainUri, SoftwarePrerequisite, ) +from galaxy.tool_util.identifiers import uri_safe_tool_id class SoftwarePrerequisiteTracker: @@ -25,12 +24,12 @@ class SoftwarePrerequisiteTracker: tool_version = step.tool_version assert tool_id self._recorded_tools.add(tool_id) - uri_safe_tool_id = urllib.parse.quote(tool_id) + safe_tool_id = uri_safe_tool_id(tool_id) if "repos/" in tool_id: # tool shed tool - give them a link... - uri = f"https://{uri_safe_tool_id}" + uri = f"https://{safe_tool_id}" else: - uri = f"gxstocktools://galaxyproject.org/{uri_safe_tool_id}" + uri = f"gxstocktools://galaxyproject.org/{safe_tool_id}" access_time = None # used to be uuid - but Pydanic validation... rightfully... disallows this software_prerequisite = SoftwarePrerequisite( diff --git a/lib/galaxy/tool_util/identifiers.py b/lib/galaxy/tool_util/identifiers.py new file mode 100644 index 00000000000..2e5932ccf6c --- /dev/null +++ b/lib/galaxy/tool_util/identifiers.py @@ -0,0 +1,8 @@ +from urllib.parse import quote + + +def uri_safe_tool_id(tool_id: str) -> str: + # Shed tool ids contain ``+`` in their version suffix (e.g. + # ``.../1.3.1+galaxy1``). In a query string the ``+`` decodes to a + # literal space unless percent-encoded as ``%2B``. + return quote(tool_id) diff --git a/lib/galaxy/tool_util/parser/xml.py b/lib/galaxy/tool_util/parser/xml.py index 5c5c69497b8..c925ecebd8c 100644 --- a/lib/galaxy/tool_util/parser/xml.py +++ b/lib/galaxy/tool_util/parser/xml.py @@ -967,7 +967,16 @@ def __parse_test_attributes( has_checksum = md5sum or checksum has_nested_tests = extra_files or element_tests or primary_datasets has_object = value_object is not VALUE_OBJECT_UNSET - if not (assert_list or file or metadata or has_checksum or has_nested_tests or has_object or has_count_assertions): + if not ( + assert_list + or file + or metadata + or ftype + or has_checksum + or has_nested_tests + or has_object + or has_count_assertions + ): raise Exception( "Test output defines nothing to check (e.g. must have a 'file' check against, assertions to check, metadata or checksum tests, etc...)" ) diff --git a/lib/galaxy/tools/__init__.py b/lib/galaxy/tools/__init__.py index 0e9dfad73fc..acbb0e1bcad 100644 --- a/lib/galaxy/tools/__init__.py +++ b/lib/galaxy/tools/__init__.py @@ -63,6 +63,7 @@ from galaxy.model import ( ToolRequest, ) from galaxy.model.dataset_collections.matching import MatchingCollections +from galaxy.model.dataset_collections.types.paired_or_unpaired import SINGLETON_IDENTIFIER from galaxy.model.dataset_collections.types.sample_sheet_workbook import _sample_sheet_to_list_collection_type from galaxy.objectstore import ObjectStorePopulator from galaxy.schema.credentials import CredentialsContext @@ -75,6 +76,7 @@ from galaxy.tool_util.deps import ( ) from galaxy.tool_util.deps.requirements import CredentialsRequirement from galaxy.tool_util.fetcher import ToolLocationFetcher +from galaxy.tool_util.identifiers import uri_safe_tool_id from galaxy.tool_util.loader import ( imported_macro_paths, raw_tool_xml_tree, @@ -195,7 +197,6 @@ from galaxy.tools.parameters.workflow_utils import workflow_building_modes from galaxy.tools.parameters.wrapped_json import json_wrap from galaxy.util import ( in_directory, - listify, Params, parse_xml_string, rst_to_html, @@ -208,7 +209,6 @@ from galaxy.util.bunch import Bunch from galaxy.util.compression_utils import get_fileobj_raw from galaxy.util.dictifiable import UsesDictVisibleKeys from galaxy.util.expressions import ExpressionContext -from galaxy.util.form_builder import SelectField from galaxy.util.json import ( safe_loads, swap_inf_nan, @@ -706,33 +706,6 @@ class ToolBox(AbstractToolBox): tool.name = tool.id return tool - def get_tool_components(self, tool_id, tool_version=None, get_loaded_tools_by_lineage=False, set_selected=False): - """ - Retrieve all loaded versions of a tool from the toolbox and return a select list enabling - selection of a different version, the list of the tool's loaded versions, and the specified tool. - """ - tool_version_select_field = None - tools = [] - tool = None - # Backwards compatibility for datasource tools that have default tool_id configured, but which - # are now using only GALAXY_URL. - tool_ids = listify(tool_id) - for tool_id in tool_ids: - if tool_id.endswith("/"): - # Some data sources send back redirects ending with `/`, this takes care of that case - tool_id = tool_id[:-1] - if get_loaded_tools_by_lineage: - tools = self.get_loaded_tools_by_lineage(tool_id) - else: - tools = self.get_tool(tool_id, tool_version=tool_version, get_all_versions=True) - if tools: - tool = self.get_tool(tool_id, tool_version=tool_version, get_all_versions=False) - assert tool - if len(tools) > 1: - tool_version_select_field = self.__build_tool_version_select_field(tools, tool.id, set_selected) - break - return tool_version_select_field, tools, tool - def _path_template_kwds(self): return { "model_tools_path": MODEL_TOOLS_PATH, @@ -782,20 +755,6 @@ class ToolBox(AbstractToolBox): stored = session.get_one(StoredWorkflow, id) return stored.latest_workflow - def __build_tool_version_select_field(self, tools, tool_id, set_selected): - """Build a SelectField whose options are the ids for the received list of tools.""" - options: list[tuple[str, str]] = [] - for tool in tools: - options.insert(0, (tool.version, tool.id)) - select_field = SelectField(name="tool_id") - for option_tup in options: - selected = set_selected and option_tup[1] == tool_id - if selected: - select_field.add_option(f"version {option_tup[0]}", option_tup[1], selected=True) - else: - select_field.add_option(f"version {option_tup[0]}", option_tup[1]) - return select_field - class DefaultToolState: """ @@ -1190,6 +1149,18 @@ class Tool(UsesDictVisibleKeys, MaybeToolParameterBundle): else: return [] + @property + def hidden_tool_versions(self): + if not self.lineage or not self.id: + return [] + versions_by_id = self.app.toolbox._tool_versions_by_id.get(self.id, {}) + hidden_versions = [] + for version in self.lineage.tool_versions: + tool = versions_by_id.get(version) + if tool and tool.hidden: + hidden_versions.append(version) + return hidden_versions + @property def is_latest_version(self): tool_versions = self.tool_versions @@ -3021,6 +2992,8 @@ class Tool(UsesDictVisibleKeys, MaybeToolParameterBundle): tool_dict["hidden"] = self.hidden tool_dict["is_workflow_compatible"] = self.is_workflow_compatible tool_dict["xrefs"] = self.xrefs + tool_dict["versions"] = self.tool_versions + tool_dict["hidden_versions"] = self.hidden_tool_versions if self.dynamic_tool: tool_dict["uuid"] = str(self.dynamic_tool.uuid) @@ -3174,6 +3147,7 @@ class Tool(UsesDictVisibleKeys, MaybeToolParameterBundle): "message": tool_message, "warnings": tool_warnings, "versions": self.tool_versions, + "hidden_versions": self.hidden_tool_versions, "requirements": [{"name": r.name, "version": r.version} for r in self.requirements], "credentials": [credential.to_dict() for credential in self.credentials] if self.credentials else [], "errors": state_errors, @@ -3275,9 +3249,8 @@ class Tool(UsesDictVisibleKeys, MaybeToolParameterBundle): return None message = "" try: - select_field, tools, tool = self.app.toolbox.get_tool_components( - tool_id, tool_version=tool_version, get_loaded_tools_by_lineage=False, set_selected=True - ) + tools = self.app.toolbox.get_tool(tool_id, tool_version=tool_version, get_all_versions=True) or [] + tool = self.app.toolbox.get_tool(tool_id, tool_version=tool_version) if tools else None if tool is None: raise exceptions.MessageException( f"This dataset was created by an obsolete tool ({tool_id}). Can't re-run." @@ -3542,8 +3515,10 @@ class DataSourceTool(OutputParameterJSONTool): return True def _build_GALAXY_URL_parameter(self): + assert self.id, "Tool id must be set to build GALAXY_URL parameter for data_source tool" return ToolParameter.build( - self, XML(f'') + self, + XML(f''), ) def parse_inputs(self, tool_source): @@ -3616,7 +3591,10 @@ class AsyncDataSourceTool(DataSourceTool): tool_type = "data_source_async" def _build_GALAXY_URL_parameter(self): - return ToolParameter.build(self, XML(f'')) + assert self.id, "Tool id must be set to build GALAXY_URL parameter for data_source_async tool" + return ToolParameter.build( + self, XML(f'') + ) class DataDestinationTool(Tool): @@ -3865,6 +3843,9 @@ class DatabaseOperationTool(Tool): if self.require_terminal_states and state in model.Dataset.non_ready_states: raise ToolInputsNotReadyException("An input dataset is pending.") + if state == model.Dataset.states.PAUSED and not self.require_terminal_or_paused_states: + raise ToolInputsNotReadyException(f"Input '{input_key}' is paused; the file is not yet available.") + if self.require_dataset_ok: if state != model.Dataset.states.OK: # TODO: frontend component should intercept and point to problematic input @@ -4093,14 +4074,26 @@ class SplitPairedAndUnpairedTool(DatabaseOperationTool): def _handle_unpaired(dce): element_identifier = dce.element_identifier - assert getattr(dce.element_object, "history_content_type", None) == "dataset" - copied_value = dce.element_object.copy(copy_tags=dce.element_object.tags, flush=False) + element_object = dce.element_object + # In list:paired_or_unpaired collections, unpaired elements are + # wrapped in a 1-element sub-collection. Unwrap to get the dataset. + if getattr(element_object, "elements", None): + inner_element = element_object.elements[0] + assert inner_element.element_identifier == SINGLETON_IDENTIFIER + element_object = inner_element.element_object + assert element_object.history_content_type == "dataset" + copied_value = element_object.copy(copy_tags=element_object.tags, flush=False) unpaired_dce_copies[element_identifier] = copied_value unpaired_dce_columns[element_identifier] = dce.columns def _handle_paired(dce): element_identifier = dce.element_identifier copied_value = dce.element_object.copy(flush=False) + # Normalize to 'paired' for list:paired output, since a + # list:paired_or_unpaired input may contain 2-element + # paired_or_unpaired sub-collections that are structurally + # equivalent to paired but carry the wider collection_type. + copied_value.collection_type = "paired" paired_dce_copies[element_identifier] = copied_value paired_datasets.append(copied_value.elements[0].element_object) paired_datasets.append(copied_value.elements[1].element_object) @@ -4114,7 +4107,8 @@ class SplitPairedAndUnpairedTool(DatabaseOperationTool): _handle_paired(element) elif collection_type == "list:paired_or_unpaired": for element in collection.elements: - if getattr(element.element_object, "history_content_type", None) == "dataset": + sub_collection = element.element_object + if len(sub_collection.elements) == 1: _handle_unpaired(element) else: _handle_paired(element) diff --git a/lib/galaxy/tools/evaluation.py b/lib/galaxy/tools/evaluation.py index 5a6e52d0328..f7c6f9171d3 100644 --- a/lib/galaxy/tools/evaluation.py +++ b/lib/galaxy/tools/evaluation.py @@ -20,6 +20,7 @@ from packaging.version import Version from galaxy import model from galaxy.authnz.util import provider_name_to_backend from galaxy.exceptions import RequestParameterInvalidException +from galaxy.files import ProvidesFileSourcesUserContext from galaxy.job_execution.compute_environment import ComputeEnvironment from galaxy.job_execution.datasets import DeferrableObjectsT from galaxy.job_execution.setup import ensure_configs_directory @@ -303,10 +304,14 @@ class ToolEvaluator: undeferred_objects: dict[str, DeferrableObjectsT] = {} transient_directory = os.path.join(job_working_directory, "inputs") safe_makedirs(transient_directory) + user_context = ProvidesFileSourcesUserContext( + WorkRequestContext(app=self.app, user=self._user, history=self._history) + ) dataset_materializer = materializer_factory( False, # unattached to a session. transient_directory=transient_directory, file_sources=self.app.file_sources, + user_context=user_context, ) for key, value in deferred_objects.items(): if isinstance(value, model.DatasetInstance): diff --git a/lib/galaxy/tools/execute.py b/lib/galaxy/tools/execute.py index a0704eabd8f..6df94dbd595 100644 --- a/lib/galaxy/tools/execute.py +++ b/lib/galaxy/tools/execute.py @@ -31,6 +31,7 @@ from galaxy.model import ( ) from galaxy.model.dataset_collections.matching import MatchingCollections from galaxy.model.dataset_collections.structure import ( + get_collection, get_structure, tool_output_to_structure, ) @@ -550,7 +551,9 @@ class ExecutionTracker: subcollection_mapping_type = collection_info.subcollection_mapping_type(qualified_name) return get_structure( - input_collection, collection_type_description, leaf_subcollection_type=subcollection_mapping_type + get_collection(input_collection), + collection_type_description, + leaf_subcollection_type=subcollection_mapping_type, ) def _structure_for_output(self, trans, tool_output): @@ -740,7 +743,9 @@ class ExecutionTracker: def walk_implicit_collections(self): collection_info = self.collection_info assert collection_info - return collection_info.structure.walk_collections(self.implicit_collections) + return collection_info.structure.walk_collections( + {k: get_collection(v) for k, v in self.implicit_collections.items()} + ) def new_execution_slices(self): if self.collection_info is None: diff --git a/lib/galaxy/tools/split_paired_and_unpaired.xml b/lib/galaxy/tools/split_paired_and_unpaired.xml index d1d05c88158..3b2ba48d881 100644 --- a/lib/galaxy/tools/split_paired_and_unpaired.xml +++ b/lib/galaxy/tools/split_paired_and_unpaired.xml @@ -76,14 +76,20 @@ - + - + + + + - + + + + diff --git a/lib/galaxy/web/framework/middleware/aiocop_integration.py b/lib/galaxy/web/framework/middleware/aiocop_integration.py new file mode 100644 index 00000000000..6892683ec49 --- /dev/null +++ b/lib/galaxy/web/framework/middleware/aiocop_integration.py @@ -0,0 +1,176 @@ +""" +aiocop integration — detect blocking I/O in async handlers. + +`aiocop `_ uses ``sys.audit`` hooks to catch +specific blocking syscalls (``socket.connect``, ``open``, ``subprocess.Popen``, +etc.) from within async tasks and report the exact call site. + +This module installs aiocop and adds an ASGI middleware that surfaces any +violations recorded during a request via an ``X-Aiocop-Violations`` response +header so test harnesses can fail requests that block the event loop. + +Activation is opt-in via the ``GALAXY_TEST_AIOCOP`` environment variable — +aiocop is a test-only dependency and is not imported in normal production +Galaxy servers. +""" + +import asyncio +import logging +import os +import threading +from contextvars import ContextVar +from typing import ( + Any, + Optional, +) + +from starlette.types import ( + ASGIApp, + Message, + Receive, + Scope, + Send, +) + +log = logging.getLogger(__name__) + +ENV_VAR = "GALAXY_TEST_AIOCOP" + +# High-severity threshold used by aiocop itself (THRESHOLD_HIGH == 50). +# Anything at or above is considered a test-failing violation. +HIGH_SEVERITY_SCORE = 50 + +_request_violations: ContextVar[Optional[list[dict[str, Any]]]] = ContextVar("aiocop_request_violations", default=None) + +_process_initialized = False + + +def aiocop_enabled() -> bool: + return os.environ.get(ENV_VAR, "").lower() in ("1", "true", "yes") + + +def _on_slow_task(event: Any) -> None: + if not event.blocking_events: + return + violations = _request_violations.get() + if violations is not None: + violations.extend(event.blocking_events) + log.error( + "aiocop detected blocking I/O on event loop (severity=%s, elapsed=%.1fms): %s", + event.severity_level, + event.elapsed_ms, + "; ".join(f"{e['event']} at {e['entry_point']}" for e in event.blocking_events), + ) + + +def install_aiocop(slow_task_threshold_ms: int = 50, trace_depth: int = 20) -> None: + """Activate aiocop monitoring on the currently running event loop. + + Must be called from the event-loop thread (e.g. a FastAPI ``startup`` + event). aiocop normally binds its audit hook to + :func:`threading.main_thread`, but Galaxy's test harness runs uvicorn's + event loop in a child thread — so we rebind the hook to the current + thread explicitly. + + The Galaxy integration test driver spins up a fresh event loop on a new + thread for every test module (see ``galaxy_test.driver.driver_util``). + aiocop has process-global state — ``sys.addaudithook`` cannot be undone + and ``patch_audit_functions`` wraps stdlib call sites in place — so we + do that setup exactly once. The per-loop pieces (loop patching, main + thread binding, slow-task callback) are re-armed on every call so the + Nth Galaxy instance is actually monitored. + """ + global _process_initialized + + import aiocop + from aiocop.core import ( + blocking_io as _bio, + slow_tasks as _slow_tasks, + ) + + if not _process_initialized: + aiocop.patch_audit_functions() + aiocop.start_blocking_io_detection(trace_depth=trace_depth) + _process_initialized = True + + # aiocop binds audit events to the process main thread; the Galaxy test + # server runs its event loop in a non-main thread, so pin aiocop to the + # thread that actually hosts the loop. + _bio._main_thread_id = threading.get_ident() + + # Re-arm detect_slow_tasks so the on-activate hook is fresh and the + # slow-task callback isn't a stale closure from a prior Galaxy instance. + # detect_slow_tasks() short-circuits on its module-global guard, so reset + # it and clear the callback list before calling again. + _slow_tasks._detect_slow_tasks_configured = False + aiocop.clear_slow_task_callbacks() + aiocop.detect_slow_tasks(threshold_ms=slow_task_threshold_ms, on_slow_task=_on_slow_task) + aiocop.activate() + log.info("aiocop blocking-I/O detection activated (threshold=%dms)", slow_task_threshold_ms) + + +class AiocopMiddleware: + """Attach per-request aiocop violations to an ``X-Aiocop-Violations`` header. + + The header is only emitted when at least one blocking event was captured. + Format: + + ``count=;severity=;first=@`` + + ``severity`` is the maximum severity weight among the captured events + (see ``aiocop.core.blocking_io.BLOCKING_EVENTS_DICT``). A value of + 50 or higher indicates a high-severity block that should fail tests. + + aiocop itself is installed on the first ``lifespan`` startup event so + it binds to the event loop actually serving traffic (Galaxy's test + harness runs uvicorn in a child thread — see + :func:`galaxy_test.driver.driver_util.launch_server`). + """ + + def __init__(self, app: ASGIApp): + self.app = app + self._installed = False + + async def __call__(self, scope: Scope, receive: Receive, send: Send) -> None: + if scope["type"] == "lifespan": + # Wrap lifespan startup to install aiocop as soon as the loop + # is running, before any request is served. + async def send_wrapper(message: Message) -> None: + if message.get("type") == "lifespan.startup.complete" and not self._installed: + install_aiocop() + self._installed = True + await send(message) + + return await self.app(scope, receive, send_wrapper) + + if scope["type"] != "http": + return await self.app(scope, receive, send) + + violations: list[dict[str, Any]] = [] + token = _request_violations.set(violations) + response_started = False + + async def send_with_header(message: Message) -> None: + nonlocal response_started + if message["type"] == "http.response.start" and not response_started: + response_started = True + # aiocop's slow-task callback fires when the task step that + # did the blocking I/O ends; yield once to let any pending + # callback populate `violations` before we emit headers. + await asyncio.sleep(0) + if violations: + max_severity = max(int(v.get("severity") or 0) for v in violations) + first = violations[0] + summary = ( + f"count={len(violations)};severity={max_severity};" + f"first={first['event']}@{first['entry_point']}" + ) + headers = list(message.get("headers", [])) + headers.append((b"x-aiocop-violations", summary.encode("latin-1"))) + message = {**message, "headers": headers} + await send(message) + + try: + await self.app(scope, receive, send_with_header) + finally: + _request_violations.reset(token) diff --git a/lib/galaxy/webapps/galaxy/api/tool_data.py b/lib/galaxy/webapps/galaxy/api/tool_data.py index 880e6420043..c7ae44f45ae 100644 --- a/lib/galaxy/webapps/galaxy/api/tool_data.py +++ b/lib/galaxy/webapps/galaxy/api/tool_data.py @@ -1,5 +1,7 @@ +from functools import partial from typing import Optional +import anyio from fastapi import ( Body, Path, @@ -117,7 +119,7 @@ class FastAPIToolData: field_name: str = ToolDataTableFieldName, ) -> ToolDataField: """Displays information about a data table field.""" - return self.tool_data_manager.show_field(table_name, field_name) + return await anyio.to_thread.run_sync(partial(self.tool_data_manager.show_field, table_name, field_name)) @router.get( "/api/tool_data/{table_name}/fields/{field_name}/files/{file_name}", diff --git a/lib/galaxy/webapps/galaxy/controllers/data_manager.py b/lib/galaxy/webapps/galaxy/controllers/data_manager.py index 93c134de60b..14a827d7f3f 100644 --- a/lib/galaxy/webapps/galaxy/controllers/data_manager.py +++ b/lib/galaxy/webapps/galaxy/controllers/data_manager.py @@ -8,6 +8,7 @@ from galaxy.model import ( DataManagerJobAssociation, Job, ) +from galaxy.tool_util.identifiers import uri_safe_tool_id from galaxy.util import ( nice_size, unicodify, @@ -34,7 +35,7 @@ class DataManager(BaseUIController): ): data_managers.append( { - "toolUrl": web.url_for(f"/?tool_id={data_manager.tool.id}"), + "toolUrl": web.url_for(f"/?tool_id={uri_safe_tool_id(data_manager.tool.id)}"), "id": data_manager_id, "name": data_manager.name, "description": data_manager.description.lower(), @@ -90,7 +91,7 @@ class DataManager(BaseUIController): "dataManager": { "name": data_manager.name, "description": data_manager.description.lower(), - "toolUrl": web.url_for(f"/?tool_id={data_manager.tool.id}"), + "toolUrl": web.url_for(f"/?tool_id={uri_safe_tool_id(data_manager.tool.id)}"), }, "jobs": jobs, "viewOnly": not_is_admin, @@ -155,7 +156,7 @@ class DataManager(BaseUIController): "id": data_manager_id, "name": data_manager.name, "description": data_manager.description.lower(), - "toolUrl": web.url_for(f"/?tool_id={data_manager.tool.id}"), + "toolUrl": web.url_for(f"/?tool_id={uri_safe_tool_id(data_manager.tool.id)}"), }, "hdaInfo": hda_info, "dataManagerOutput": data_manager_output, diff --git a/lib/galaxy/webapps/galaxy/controllers/tool_runner.py b/lib/galaxy/webapps/galaxy/controllers/tool_runner.py index 5f1fdac9300..b8d09529ab5 100644 --- a/lib/galaxy/webapps/galaxy/controllers/tool_runner.py +++ b/lib/galaxy/webapps/galaxy/controllers/tool_runner.py @@ -9,6 +9,7 @@ from markupsafe import escape import galaxy.util from galaxy import web from galaxy.model import HistoryDatasetAssociation +from galaxy.tool_util.identifiers import uri_safe_tool_id from galaxy.tools import DataSourceTool from galaxy.web import ( error, @@ -37,11 +38,10 @@ class ToolRunner(BaseUIController): """Catches the tool id and redirects as needed""" return self.index(trans, tool_id=tool_id, **kwd) - def __get_tool(self, tool_id, tool_version=None, get_loaded_tools_by_lineage=False, set_selected=False): - tool_version_select_field, tools, tool = self.get_toolbox().get_tool_components( - tool_id, tool_version, get_loaded_tools_by_lineage, set_selected - ) - return tool + def __get_tool(self, tool_id, tool_version=None): + # Some data sources send back redirects ending with `/`, this takes care of that case + tool_id = tool_id.rstrip("/") + return self.get_toolbox().get_tool(tool_id, tool_version=tool_version) @web.expose def index(self, trans, tool_id=None, from_noframe=None, **kwd): @@ -74,7 +74,7 @@ class ToolRunner(BaseUIController): return __tool_404__() # FIXME: Tool class should define behavior if tool.tool_type in ["default", "interactivetool"]: - return trans.response.send_redirect(url_for(f"/?tool_id={tool_id}")) + return trans.response.send_redirect(url_for(f"/?tool_id={uri_safe_tool_id(tool_id)}")) # execute tool without displaying form # (used for datasource tools, but note that data_source_async tools diff --git a/lib/galaxy/webapps/galaxy/fast_app.py b/lib/galaxy/webapps/galaxy/fast_app.py index 939e7e862b9..e169a2c88b5 100644 --- a/lib/galaxy/webapps/galaxy/fast_app.py +++ b/lib/galaxy/webapps/galaxy/fast_app.py @@ -167,6 +167,13 @@ def add_galaxy_middleware(app: FastAPI, gx_app): max_age=600, ) + from galaxy.web.framework.middleware.aiocop_integration import aiocop_enabled + + if aiocop_enabled(): + from galaxy.web.framework.middleware.aiocop_integration import AiocopMiddleware + + app.add_middleware(AiocopMiddleware) + def include_legacy_openapi(app, gx_app): if app.openapi_schema: diff --git a/lib/galaxy_test/api/test_dataset_collections.py b/lib/galaxy_test/api/test_dataset_collections.py index 7a505e7346f..0b21a16cc15 100644 --- a/lib/galaxy_test/api/test_dataset_collections.py +++ b/lib/galaxy_test/api/test_dataset_collections.py @@ -7,6 +7,7 @@ from urllib.parse import quote from galaxy.tool_util_models.sample_sheet import SampleSheetColumnDefinitions from galaxy.util import galaxy_root_path from galaxy_test.base.api_asserts import ( + assert_error_message_contains, assert_has_key, assert_object_id_error, assert_status_code_is, @@ -152,6 +153,30 @@ class TestDatasetCollectionsApi(ApiTestCase): returned_collections = dataset_collection["elements"] assert len(returned_collections) == 1, dataset_collection + def test_create_list_paired_or_unpaired_rejects_raw_hda_child(self, history_id): + element_identifiers = self.dataset_collection_populator.list_identifiers(history_id, contents=[("el1", "hi")]) + payload = dict( + instance_type="history", + history_id=history_id, + element_identifiers=element_identifiers, + collection_type="list:paired_or_unpaired", + ) + create_response = self._post("dataset_collections", payload, json=True) + assert_status_code_is(create_response, 400) + assert_error_message_contains(create_response, "sub-collection of type 'paired_or_unpaired'") + + def test_create_list_paired_rejects_raw_hda_child(self, history_id): + element_identifiers = self.dataset_collection_populator.list_identifiers(history_id, contents=[("el1", "hi")]) + payload = dict( + instance_type="history", + history_id=history_id, + element_identifiers=element_identifiers, + collection_type="list:paired", + ) + create_response = self._post("dataset_collections", payload, json=True) + assert_status_code_is(create_response, 400) + assert_error_message_contains(create_response, "sub-collection of type 'paired'") + def test_create_record(self, history_id): contents = [ ("condition", "1\t2\t3"), diff --git a/lib/galaxy_test/api/test_tool_execute.py b/lib/galaxy_test/api/test_tool_execute.py index 83eaadc4f29..c85f67aaa95 100644 --- a/lib/galaxy_test/api/test_tool_execute.py +++ b/lib/galaxy_test/api/test_tool_execute.py @@ -974,5 +974,5 @@ def test_map_over_dce_on_non_multiple_data_param( execute = required_tool.execute().with_inputs(inputs) execute.assert_has_n_jobs(2).assert_creates_n_implicit_collections(1) output_collection = execute.assert_creates_implicit_collection(0) - output_collection.assert_has_dataset_element("test0").with_contents_stripped("123") - output_collection.assert_has_dataset_element("test1").with_contents_stripped("456") + output_collection.assert_has_dataset_element("forward").with_contents_stripped("123") + output_collection.assert_has_dataset_element("reverse").with_contents_stripped("456") diff --git a/lib/galaxy_test/api/test_tools.py b/lib/galaxy_test/api/test_tools.py index 1f2f2e92366..68ac9bcad62 100644 --- a/lib/galaxy_test/api/test_tools.py +++ b/lib/galaxy_test/api/test_tools.py @@ -1737,6 +1737,28 @@ class TestToolsApi(ApiTestCase, TestsTools): # Return last version assert tool_info["version"] == "0.2" + @skip_without_tool("multiple_versions_hidden") + def test_show_lists_hidden_versions_separately(self): + tool_info = self._show_valid_tool("multiple_versions_hidden", tool_version="0.1") + assert tool_info["version"] == "0.1" + assert tool_info["versions"] == ["0.1", "0.2"] + assert tool_info["hidden_versions"] == ["0.1"] + + @skip_without_tool("multiple_versions_hidden") + def test_run_hidden_version(self): + with self.dataset_populator.test_history_for(self.test_run_hidden_version) as history_id: + outputs = self._run( + tool_id="multiple_versions_hidden", + history_id=history_id, + tool_version="0.1", + assert_ok=True, + wait_for_job=True, + ) + assert len(outputs["outputs"]) == 1 + output = outputs["outputs"][0] + output_content = self.dataset_populator.get_history_dataset_content(history_id, dataset=output) + assert output_content.strip() == "Hidden Version 0.1" + @skip_without_tool("cat1") def test_run_cat1_single_meta_wrapper(self): with self.dataset_populator.test_history_for(self.test_run_cat1_single_meta_wrapper) as history_id: @@ -2767,6 +2789,33 @@ class TestToolsApi(ApiTestCase, TestsTools): assert len(response_object["jobs"]) == 2 assert len(response_object["implicit_collections"]) == 1 + def test_can_map_over_dce_from_larger_list_paired(self): + """Regression: mapping a DCE should use the child collection structure, + not the parent list structure. Previously raised KeyError when the parent + list had more elements than the child pair collection.""" + with self.dataset_populator.test_history() as history_id: + pair_ids = [] + for _ in range(3): + pair_id = self.dataset_collection_populator.create_pair_in_history( + history_id, contents=["0", "0"], wait=True + ).json()["outputs"][0]["id"] + pair_ids.append(pair_id) + ok_hdca = self.dataset_collection_populator.create_list_from_pairs(history_id, pair_ids) + dce_id = ok_hdca.json()["elements"][0]["id"] + + inputs = { + "input1": { + "batch": True, + "values": [{"src": "dce", "id": dce_id, "map_over_type": None}], + }, + } + response = self._run_cat1(history_id, inputs=inputs) + self._assert_status_code_is(response, 200) + + response_object = response.json() + assert len(response_object["jobs"]) == 2 + assert len(response_object["implicit_collections"]) == 1 + @skip_without_tool("identifier_source") def test_default_identifier_source_map_over(self): with self.dataset_populator.test_history() as history_id: diff --git a/lib/galaxy_test/api/test_workflows.py b/lib/galaxy_test/api/test_workflows.py index 0ecf56e50d0..35fdbfcac96 100644 --- a/lib/galaxy_test/api/test_workflows.py +++ b/lib/galaxy_test/api/test_workflows.py @@ -6279,6 +6279,68 @@ test_data: assert message["workflow_step_id"] == 2 assert "Invalid new collection identifier" in message["details"] + @skip_without_tool("__RELABEL_FROM_FILE__") + @skip_without_tool("job_properties") + @skip_without_tool("cat1") + def test_relabel_from_file_with_paused_labels_does_not_crash(self, history_id): + # Regression test for https://github.com/galaxyproject/galaxy/issues/22520: + # a PAUSED HDA wired as the labels input to __RELABEL_FROM_FILE__ used to + # crash produce_outputs with FileNotFoundError. + summary = self._run_workflow( + """ +class: GalaxyWorkflow +inputs: + input_c: + type: collection + collection_type: list +steps: + failing_step: + tool_id: job_properties + state: + thebool: true + failbool: true + paused_labels: + tool_id: cat1 + in: + input1: failing_step/out_file1 + relabel: + tool_id: "__RELABEL_FROM_FILE__" + in: + input: input_c + how|labels: paused_labels/out_file1 +""", + test_data=""" +input_c: + collection_type: list + elements: + - identifier: i1 + content: "A" +""", + history_id=history_id, + wait=False, + assert_ok=False, + ) + invocation_id = summary.invocation_id + + def cat1_paused(): + for j in self._history_jobs(history_id): + if j.get("tool_id") == "cat1" and j.get("state") == "paused": + return True + return None + + try: + wait_on(cat1_paused, "cat1 job to be paused after upstream failure") + # Give the scheduler a couple of iterations to attempt the relabel step. + time.sleep(3) + + invocation = self.workflow_populator.get_invocation(invocation_id, step_details=True) + assert invocation["state"] != "failed", f"Relabel step crashed on paused input; invocation: {invocation}" + relabel_jobs = [j for j in self._history_jobs(history_id) if j.get("tool_id") == "__RELABEL_FROM_FILE__"] + for rj in relabel_jobs: + assert rj["state"] != "error", f"Relabel produced an errored job: {rj}" + finally: + self.workflow_populator.cancel_invocation(invocation_id) + @skip_without_tool("identifier_multiple") def test_invocation_map_over(self, history_id): summary = self._run_workflow( diff --git a/lib/galaxy_test/base/api.py b/lib/galaxy_test/base/api.py index 2dfcfce13ce..da0f014d1e4 100644 --- a/lib/galaxy_test/base/api.py +++ b/lib/galaxy_test/base/api.py @@ -100,11 +100,14 @@ class UsesApiTestCaseMixin: def tearDown(self): if os.environ.get("GALAXY_TEST_EXTERNAL") is None: - # Only kill running jobs after test for managed test instances - response = self.galaxy_interactor.get("jobs?state=running") + # Only kill running jobs after test for managed test instances. + # Use low-level methods to bypass aiocop blocking-I/O checks — + # tearDown cleanup should not fail the test for server-side + # blocking events. + response = self.galaxy_interactor._get("jobs?state=running") if response.ok: for job in response.json(): - self._delete(f"jobs/{job['id']}") + self.galaxy_interactor._delete(f"jobs/{job['id']}") def _api_url(self, path, params=None, use_key=None, use_admin_key=None): if not params: @@ -224,6 +227,37 @@ class UsesApiTestCaseMixin: _assert_has_key = _assert_has_keys +AIOCOP_ENABLED = os.environ.get("GALAXY_TEST_AIOCOP", "").lower() in ("1", "true", "yes") +# Matches aiocop's THRESHOLD_HIGH — high-severity events (socket.connect, +# subprocess, etc.) fail the request; lower-severity stat-like calls are +# logged but not treated as test failures. +AIOCOP_FAIL_SEVERITY = 50 + + +def _check_aiocop_violations(response: "Response") -> None: + """Fail the test if the server reported a high-severity aiocop violation. + + When the Galaxy test server has ``GALAXY_TEST_AIOCOP=1`` set, each + response includes an ``X-Aiocop-Violations`` header summarizing any + blocking-I/O events caught by the aiocop audit hook during that + request. See ``galaxy.web.framework.middleware.aiocop_integration``. + """ + header = response.headers.get("x-aiocop-violations") + if not header: + return + fields = dict(part.split("=", 1) for part in header.split(";") if "=" in part) + try: + severity = int(fields.get("severity", "0")) + except ValueError: + severity = 0 + if severity >= AIOCOP_FAIL_SEVERITY: + raise AssertionError( + f"aiocop detected high-severity blocking I/O (severity={severity}, " + f"count={fields.get('count', '?')}, first={fields.get('first', '?')}) " + f"during {response.request.method} {response.request.path_url}" + ) + + class ApiTestInteractor(BaseInteractor): """Specialized variant of the API interactor (originally developed for tool functional tests) for testing the API generally. @@ -239,26 +273,31 @@ class ApiTestInteractor(BaseInteractor): # directly to test API - instead of relying on higher-level constructs for # specific pieces of the API (the way it is done with the variant for tool) # testing. + def _maybe_check_aiocop(self, response: "Response") -> "Response": + if AIOCOP_ENABLED: + _check_aiocop_violations(response) + return response + def get(self, *args, **kwds): - return self._get(*args, **kwds) + return self._maybe_check_aiocop(self._get(*args, **kwds)) def head(self, *args, **kwds): - return self._head(*args, **kwds) + return self._maybe_check_aiocop(self._head(*args, **kwds)) def post(self, *args, **kwds): - return self._post(*args, **kwds) + return self._maybe_check_aiocop(self._post(*args, **kwds)) def options(self, *args, **kwds): - return self._options(*args, **kwds) + return self._maybe_check_aiocop(self._options(*args, **kwds)) def delete(self, *args, **kwds): - return self._delete(*args, **kwds) + return self._maybe_check_aiocop(self._delete(*args, **kwds)) def patch(self, *args, **kwds): - return self._patch(*args, **kwds) + return self._maybe_check_aiocop(self._patch(*args, **kwds)) def put(self, *args, **kwds): - return self._put(*args, **kwds) + return self._maybe_check_aiocop(self._put(*args, **kwds)) class AnonymousGalaxyInteractor(ApiTestInteractor): diff --git a/lib/galaxy_test/driver/driver_util.py b/lib/galaxy_test/driver/driver_util.py index 9aeb37cf0d6..ec48dc8c768 100644 --- a/lib/galaxy_test/driver/driver_util.py +++ b/lib/galaxy_test/driver/driver_util.py @@ -938,12 +938,16 @@ class GalaxyTestDriver(TestDriver): if handle_galaxy_config_kwds is not None: handle_galaxy_config_kwds(galaxy_config) - server_wrapper = launch_server( + launch_kwargs: dict[str, Any] = dict( app_factory=lambda: self.build_galaxy_app(galaxy_config), webapp_factory=lambda *args, **kwd: buildapp.app_factory(*args, wsgi_preflight=False, **kwd), galaxy_config=galaxy_config, config_object=config_object, ) + custom_init_fast_app = getattr(config_object, "init_fast_app", None) + if custom_init_fast_app is not None: + launch_kwargs["init_fast_app"] = custom_init_fast_app + server_wrapper = launch_server(**launch_kwargs) self.server_wrappers.append(server_wrapper) else: log.info(f"Functional tests will be run against test external Galaxy server {self.external_galaxy}") diff --git a/lib/galaxy_test/workflow/split_paired_and_unpaired.gxwf-tests.yml b/lib/galaxy_test/workflow/split_paired_and_unpaired.gxwf-tests.yml new file mode 100644 index 00000000000..161372786d3 --- /dev/null +++ b/lib/galaxy_test/workflow/split_paired_and_unpaired.gxwf-tests.yml @@ -0,0 +1,48 @@ +- doc: | + Test splitting a list:paired_or_unpaired collection with mixed paired + and unpaired elements into separate paired and unpaired output collections. + job: + input_collection: + class: Collection + collection_type: list:paired_or_unpaired + elements: + - identifier: el1 + class: Collection + type: paired_or_unpaired + elements: + - identifier: forward + class: File + contents: "forward content" + - identifier: reverse + class: File + contents: "reverse content" + - identifier: el2 + class: Collection + type: paired_or_unpaired + elements: + - identifier: unpaired + class: File + contents: "unpaired content" + outputs: + output_unpaired: + class: Collection + collection_type: list + elements: + el2: + asserts: + - that: has_text + text: "unpaired content" + output_paired: + class: Collection + collection_type: list:paired + elements: + el1: + elements: + forward: + asserts: + - that: has_text + text: "forward content" + reverse: + asserts: + - that: has_text + text: "reverse content" diff --git a/lib/galaxy_test/workflow/split_paired_and_unpaired.gxwf.yml b/lib/galaxy_test/workflow/split_paired_and_unpaired.gxwf.yml new file mode 100644 index 00000000000..96b4d319dc3 --- /dev/null +++ b/lib/galaxy_test/workflow/split_paired_and_unpaired.gxwf.yml @@ -0,0 +1,15 @@ +class: GalaxyWorkflow +inputs: + input_collection: + type: collection + collection_type: list:paired_or_unpaired +outputs: + output_unpaired: + outputSource: split/output_unpaired + output_paired: + outputSource: split/output_paired +steps: + split: + tool_id: '__SPLIT_PAIRED_AND_UNPAIRED__' + in: + input: input_collection diff --git a/packages/web_framework/setup.cfg b/packages/web_framework/setup.cfg index afbc21cdf3b..46175b12b24 100644 --- a/packages/web_framework/setup.cfg +++ b/packages/web_framework/setup.cfg @@ -40,6 +40,7 @@ install_requires = requests Routes SQLAlchemy>=2.0.37,<2.1,!=2.0.41 + starlette WebOb package_dir = =src diff --git a/pyproject.toml b/pyproject.toml index 45f5775d2ed..ee68a0fc1c0 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -122,6 +122,7 @@ Repository = "https://github.com/galaxyproject/galaxy" [dependency-groups] test = [ + "aiocop", "ase>=3.18.1", "axe-selenium-python", "boto3", diff --git a/run_tests.sh b/run_tests.sh index 0e3e5b75390..cca4a79e946 100755 --- a/run_tests.sh +++ b/run_tests.sh @@ -334,6 +334,10 @@ fi xunit_report_file="" structured_data_report_file="" structured_data_html=0 +# Detect blocking I/O in async handlers during tests via aiocop +# (https://github.com/Feverup/aiocop). Set to 0 to disable. +GALAXY_TEST_AIOCOP=${GALAXY_TEST_AIOCOP:-1} +export GALAXY_TEST_AIOCOP SKIP_CLIENT_BUILD=${GALAXY_SKIP_CLIENT_BUILD:-1} if [ "$SKIP_CLIENT_BUILD" = "1" ]; then skip_client_build="--skip-client-build" diff --git a/test/functional/tools/multiple_versions_hidden_v01.xml b/test/functional/tools/multiple_versions_hidden_v01.xml new file mode 100644 index 00000000000..57ffb32317f --- /dev/null +++ b/test/functional/tools/multiple_versions_hidden_v01.xml @@ -0,0 +1,11 @@ + + '$out_file1' + ]]> + + + + + + + diff --git a/test/functional/tools/multiple_versions_hidden_v02.xml b/test/functional/tools/multiple_versions_hidden_v02.xml new file mode 100644 index 00000000000..3e2d9620642 --- /dev/null +++ b/test/functional/tools/multiple_versions_hidden_v02.xml @@ -0,0 +1,11 @@ + + '$out_file1' + ]]> + + + + + + + diff --git a/test/functional/tools/sample_tool_conf.xml b/test/functional/tools/sample_tool_conf.xml index 91efcbd959c..c3a9bf6f92d 100644 --- a/test/functional/tools/sample_tool_conf.xml +++ b/test/functional/tools/sample_tool_conf.xml @@ -250,6 +250,8 @@ +