Merge pull request #22542 from nsoranzo/merge_dev_20260422

Merge release_26.0 into dev 20260422
This commit is contained in:
Marius van den Beek
2026-04-22 09:16:04 +02:00
committed by GitHub
58 changed files with 1205 additions and 191 deletions
+11 -3
View File
@@ -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(() => {
<ToolVersionsButton
v-if="showVersions"
:version="props.version"
:versions="versions"
:versions="visibleVersions"
@onChangeVersion="onChangeVersion" />
<ToolOptionsButton
:id="props.id"
+2
View File
@@ -41,6 +41,7 @@ export interface UploadRowModel {
targetHistoryId?: string;
toPosixLines: boolean;
id?: string;
autoDecompress: boolean;
}
export const defaultModel: UploadRowModel = {
@@ -62,4 +63,5 @@ export const defaultModel: UploadRowModel = {
spaceToTab: false,
status: "init",
toPosixLines: true,
autoDecompress: true,
};
+1
View File
@@ -193,6 +193,7 @@ export function useZipExplorer() {
ext: defaultModel.extension,
space_to_tab: defaultModel.spaceToTab,
to_posix_lines: defaultModel.toPosixLines,
auto_decompress: defaultModel.autoDecompress,
deferred: defaultModel.deferred ?? false,
};
uploadItems.push(uploadItem);
+4 -4
View File
@@ -221,11 +221,11 @@ describe("UploadQueue", () => {
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",
+29
View File
@@ -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");
});
});
// ============================================================================
+14 -7
View File
@@ -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<Omit<UrlUploadItem, "src" | "url" | "historyId">> = {},
): 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,
},
+1 -1
View File
@@ -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(),
};
}
+56 -1
View File
@@ -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);
});
});
+50 -29
View File
@@ -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<R>({ url, headers, params, errorSimplify = true }: UrlDataOptions): Promise<R> {
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 };
}
/**
+1 -1
View File
@@ -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
+3 -1
View File
@@ -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)
@@ -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
@@ -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
+1
View File
@@ -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)
+9 -5
View File
@@ -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:
+27
View File
@@ -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
):
+7 -4
View File
@@ -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
@@ -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
@@ -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)
@@ -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,
@@ -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(
+13 -2
View File
@@ -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,
)
+1 -1
View File
@@ -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:
+4 -5
View File
@@ -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(
+8
View File
@@ -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)
+10 -1
View File
@@ -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...)"
)
+45 -51
View File
@@ -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'<param name="GALAXY_URL" type="baseurl" value="/tool_runner?tool_id={self.id}" />')
self,
XML(f'<param name="GALAXY_URL" type="baseurl" value="/tool_runner?tool_id={uri_safe_tool_id(self.id)}" />'),
)
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'<param name="GALAXY_URL" type="baseurl" value="/async/{self.id}" />'))
assert self.id, "Tool id must be set to build GALAXY_URL parameter for data_source_async tool"
return ToolParameter.build(
self, XML(f'<param name="GALAXY_URL" type="baseurl" value="/async/{uri_safe_tool_id(self.id)}" />')
)
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)
+5
View File
@@ -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):
+7 -2
View File
@@ -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:
@@ -76,14 +76,20 @@
<param name="input">
<collection type="list:paired_or_unpaired">
<element name="el1">
<collection type="paired">
<collection type="paired_or_unpaired">
<element name="forward" value="simple_line.txt" />
<element name="reverse" value="simple_line_alternative.txt" />
</collection>
</element>
<element name="el2" value="simple_line.txt">
<element name="el2">
<collection type="paired_or_unpaired">
<element name="unpaired" value="simple_line.txt" />
</collection>
</element>
<element name="el3" value="simple_line_alternative.txt">
<element name="el3">
<collection type="paired_or_unpaired">
<element name="unpaired" value="simple_line_alternative.txt" />
</collection>
</element>
</collection>
</param>
@@ -0,0 +1,176 @@
"""
aiocop integration — detect blocking I/O in async handlers.
`aiocop <https://github.com/Feverup/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=<n>;severity=<max>;first=<event>@<entry_point>``
``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)
+3 -1
View File
@@ -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}",
@@ -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,
@@ -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
+7
View File
@@ -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:
@@ -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"),
+2 -2
View File
@@ -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")
+49
View File
@@ -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:
+62
View File
@@ -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(
+49 -10
View File
@@ -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):
+5 -1
View File
@@ -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}")
@@ -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"
@@ -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
+1
View File
@@ -40,6 +40,7 @@ install_requires =
requests
Routes
SQLAlchemy>=2.0.37,<2.1,!=2.0.41
starlette
WebOb
package_dir =
=src
+1
View File
@@ -122,6 +122,7 @@ Repository = "https://github.com/galaxyproject/galaxy"
[dependency-groups]
test = [
"aiocop",
"ase>=3.18.1",
"axe-selenium-python",
"boto3",
+4
View File
@@ -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"
@@ -0,0 +1,11 @@
<tool id="multiple_versions_hidden" name="multiple_versions_hidden" version="0.1">
<command><![CDATA[
echo "Hidden Version 0.1" > '$out_file1'
]]></command>
<inputs>
<param name="inttest" value="1" type="integer" />
</inputs>
<outputs>
<data name="out_file1" format="txt" />
</outputs>
</tool>
@@ -0,0 +1,11 @@
<tool id="multiple_versions_hidden" name="multiple_versions_hidden" version="0.2">
<command><![CDATA[
echo "Hidden Version 0.2" > '$out_file1'
]]></command>
<inputs>
<param name="inttest" value="1" type="integer" />
</inputs>
<outputs>
<data name="out_file1" format="txt" />
</outputs>
</tool>
@@ -250,6 +250,8 @@
<tool file="multiple_versions_v01.xml" />
<tool file="multiple_versions_v01galaxy6.xml" />
<tool file="multiple_versions_v02.xml" />
<tool file="multiple_versions_hidden_v01.xml" hidden="true" />
<tool file="multiple_versions_hidden_v02.xml" />
<section id="test_section_multi" name="Test Section with Multiple Versions">
<tool file="multiple_versions_v01.xml" />
<tool file="multiple_versions_v01galaxy6.xml" />
@@ -0,0 +1,83 @@
"""Integration tests for aiocop blocking-I/O detection.
Spins up a Galaxy server with aiocop enabled and registers test-only debug
endpoints that intentionally perform high-severity blocking I/O from inside
an async handler. The aiocop audit hook must catch the call and the test
interactor must fail the request via the ``X-Aiocop-Violations`` header.
The debug endpoints are **only** added to the test server — they never
exist in a production Galaxy instance.
"""
import socket
import pytest
from galaxy.webapps.galaxy.fast_app import initialize_fast_app
from galaxy_test.base.api import AIOCOP_ENABLED
from galaxy_test.driver import integration_util
def _init_fast_app_with_debug_routes(wsgi_webapp, gx_app):
"""Wrap the standard FastAPI initializer to add debug blocking routes.
Routes must be injected before the WSGI catch-all mount (``app.mount("/", ...)``)
that ``initialize_fast_app`` adds, otherwise Starlette never reaches them.
We pop the catch-all, register our routes, then re-append it.
"""
app = initialize_fast_app(wsgi_webapp, gx_app)
wsgi_catch_all = app.routes.pop()
@app.get("/api/debug/block")
async def block_event_loop() -> dict:
# socket.getaddrinfo is a high-severity blocking call in aiocop's
# severity table and will trigger the audit hook.
try:
socket.getaddrinfo("invalid.local.test.galaxy.example", 80)
except socket.gaierror:
pass
return {"status": "blocked"}
@app.get("/api/debug/ok")
async def ok() -> dict:
return {"status": "ok"}
app.routes.append(wsgi_catch_all)
return app
@pytest.mark.skipif(
not AIOCOP_ENABLED,
reason="GALAXY_TEST_AIOCOP not set",
)
class TestAiocopBlockingDetection(integration_util.IntegrationTestCase):
init_fast_app = staticmethod(_init_fast_app_with_debug_routes)
def test_ok_endpoint_does_not_block(self):
"""A trivial async endpoint must not trigger a high-severity aiocop violation.
Low-severity events (e.g. ``os.stat`` from Galaxy middleware) may
still appear in the ``X-Aiocop-Violations`` header; only events at
``AIOCOP_FAIL_SEVERITY`` or above fail the request.
"""
response = self._get("debug/ok")
self._assert_status_code_is(response, 200)
assert response.json()["status"] == "ok"
header = response.headers.get("x-aiocop-violations", "")
if header:
fields = dict(p.split("=", 1) for p in header.split(";") if "=" in p)
assert (
int(fields.get("severity", "0")) < 50
), f"OK endpoint triggered high-severity aiocop violation: {header!r}"
def test_blocking_endpoint_detected(self):
"""An endpoint that makes a blocking syscall must be caught by aiocop."""
with pytest.raises(AssertionError, match="aiocop detected high-severity blocking I/O"):
self._get("debug/block")
def test_violations_header_present_on_block(self):
"""The X-Aiocop-Violations header must be set when blocking I/O occurs."""
response = self.galaxy_interactor._get("debug/block")
header = response.headers.get("x-aiocop-violations", "")
assert "count=" in header and "severity=" in header, f"Missing aiocop violation summary: {header!r}"
+93 -1
View File
@@ -4,7 +4,11 @@ from sqlalchemy import select
from galaxy.model import Dataset
from galaxy_test.base import api_asserts
from galaxy_test.base.populators import DatasetPopulator
from galaxy_test.base.populators import (
DatasetPopulator,
WorkflowPopulator,
)
from galaxy_test.base.workflow_fixtures import WORKFLOW_SIMPLE_CAT_TWICE
from galaxy_test.driver import integration_util
from galaxy_test.driver.integration_setup import (
GROUP_A,
@@ -17,6 +21,7 @@ from galaxy_test.driver.integration_setup import (
class TestPosixFileSourceIntegration(PosixFileSourceSetup, integration_util.IntegrationTestCase):
dataset_populator: DatasetPopulator
framework_tool_and_types = True
required_role_expression = REQUIRED_ROLE_EXPRESSION
required_group_expression = REQUIRED_GROUP_EXPRESSION
@@ -24,6 +29,7 @@ class TestPosixFileSourceIntegration(PosixFileSourceSetup, integration_util.Inte
super().setUp()
self._write_file_fixtures()
self.dataset_populator = DatasetPopulator(self.galaxy_interactor)
self.workflow_populator = WorkflowPopulator(self.galaxy_interactor)
def test_plugin_config(self):
# Default user has required role but not required group, so cannot see plugin
@@ -68,6 +74,92 @@ class TestPosixFileSourceIntegration(PosixFileSourceSetup, integration_util.Inte
list_response = self.galaxy_interactor.get("remote_files", data)
self._assert_list_response_matches_fixtures(list_response)
def test_workflow_input_materialized_from_group_file_source(self):
# Ensure the default non-admin user is in the required group so the
# file source is accessible.
group_a_id = self._ensure_group(GROUP_A)
self._add_user_to_group(group_a_id, self.dataset_populator.user_id())
with self.dataset_populator.test_history() as history_id:
# Pass the workflow input as a URL at invocation time so the scheduler
# creates a deferred HDA and then materializes it from the
# group-restricted file source (requires_materialization path in
# galaxy.workflow.run_request).
workflow_id, response = self._invoke_gxfiles_posix_test_workflow(history_id)
invocation_id = response.json()["id"]
self.workflow_populator.wait_for_workflow(workflow_id, invocation_id, history_id, assert_ok=True)
# The input was materialized from gxfiles://posix_test/a; its fixture
# content is "a\n", so cat of it concatenated with itself is "a\na\n".
input_hda = self.dataset_populator.get_history_dataset_details(history_id, hid=1)
assert input_hda["state"] == "ok", input_hda
cat_output = self.dataset_populator.get_history_dataset_content(history_id, hid=2)
assert cat_output == "a\na\n", cat_output
def test_workflow_input_materialization_denied_without_group(self):
# Counterpart to test_workflow_input_materialized_from_group_file_source:
# a user who is not a member of any of the groups required by posix_test
# cannot see the file source, and any workflow invocation that tries to
# materialize a gxfiles://posix_test/... input must fail rather than
# silently leak data across the group boundary.
with self._different_user("wf_no_fs_access_user@bx.psu.edu"):
# The restricted plugin is not advertised to this user.
plugin_config_response = self.galaxy_interactor.get("remote_files/plugins")
api_asserts.assert_status_code_is_ok(plugin_config_response)
plugins = plugin_config_response.json()
assert all(p["uri_root"] != "gxfiles://posix_test" for p in plugins), plugins
# And direct access to the file source is forbidden.
list_response = self.galaxy_interactor.get("remote_files", {"target": "gxfiles://posix_test"})
self._assert_access_forbidden_response(list_response)
# Invoking a workflow with a gxfiles://posix_test/... input must not
# result in the file being materialized into this user's history.
with self.dataset_populator.test_history() as history_id:
workflow_id, response = self._invoke_gxfiles_posix_test_workflow(history_id)
invocation_id = response.json()["id"]
# Wait for the invocation to reach a terminal state; materialization
# should fail because the file source is not accessible to this user.
self.workflow_populator.wait_for_invocation(workflow_id, invocation_id, assert_ok=False)
invocation = self.workflow_populator.get_invocation(invocation_id)
assert invocation["state"] == "failed", invocation
# The scheduler records a `dataset_failed` message pointing at the
# HDA whose materialization failed. The ItemAccessibilityException
# from _check_user_access is surfaced on that HDA's `info` field.
messages = invocation["messages"]
dataset_failed = [m for m in messages if m["reason"] == "dataset_failed"]
assert dataset_failed, messages
failed_hda = self.dataset_populator.get_history_dataset_details(history_id, hid=1, assert_ok=False)
assert failed_hda["state"] == "error", failed_hda
assert "no access to file source" in failed_hda["misc_info"], failed_hda
def _invoke_gxfiles_posix_test_workflow(self, history_id: str):
"""Submit WORKFLOW_SIMPLE_CAT_TWICE with a single gxfiles://posix_test/a
URL input and return (workflow_id, raw invocation response)."""
workflow_id = self.workflow_populator.upload_yaml_workflow(WORKFLOW_SIMPLE_CAT_TWICE)
workflow_request = {
"history": f"hist_id={history_id}",
"inputs": {
"input1": {
"src": "url",
"url": "gxfiles://posix_test/a",
"ext": "txt",
"deferred": False,
},
},
"inputs_by": "name",
}
return workflow_id, self.workflow_populator.invoke_workflow_raw(workflow_id, workflow_request, assert_ok=True)
def _ensure_group(self, group_name: str) -> str:
"""Idempotent: return the id of an existing group with this name, or create it."""
groups_response = self._get("groups", admin=True)
api_asserts.assert_status_code_is_ok(groups_response)
for group in groups_response.json():
if group["name"] == group_name:
return group["id"]
return self._create_group(group_name)
def _create_group(self, group_name: str):
payload = {
"name": group_name,
@@ -0,0 +1,52 @@
import pytest
from galaxy.exceptions import (
RequestParameterInvalidException,
RequestParameterMissingException,
)
from galaxy.model import HistoryDatasetAssociation
from galaxy.model.dataset_collections.types.sample_sheet import SampleSheetDatasetCollectionType
def _make_instances(*identifiers):
return {identifier: HistoryDatasetAssociation(extension="txt") for identifier in identifiers}
def test_generate_elements_no_column_definitions_no_rows():
# Uploading a sample_sheet collection without any sample-sheet-specific
# columns (no column_definitions, no rows) should succeed and produce one
# element per identifier with empty columns.
instances = _make_instances("a", "b")
elements = list(SampleSheetDatasetCollectionType().generate_elements(instances))
assert [e.element_identifier for e in elements] == ["a", "b"]
assert all(e.columns is None for e in elements)
def test_generate_elements_empty_list_column_definitions_no_rows():
instances = _make_instances("a", "b")
elements = list(SampleSheetDatasetCollectionType().generate_elements(instances, column_definitions=[], rows=None))
assert [e.element_identifier for e in elements] == ["a", "b"]
assert all(e.columns is None for e in elements)
def test_generate_elements_rows_required_with_column_definitions():
instances = _make_instances("a", "b")
column_definitions = [{"type": "int", "name": "replicate", "optional": False}]
with pytest.raises(RequestParameterMissingException):
list(
SampleSheetDatasetCollectionType().generate_elements(
instances, column_definitions=column_definitions, rows=None
)
)
def test_generate_elements_rows_validated_with_column_definitions():
instances = _make_instances("a", "b")
column_definitions = [{"type": "int", "name": "replicate", "optional": False}]
rows = {"a": [1], "b": None}
with pytest.raises(RequestParameterInvalidException):
list(
SampleSheetDatasetCollectionType().generate_elements(
instances, column_definitions=column_definitions, rows=rows
)
)
@@ -33,6 +33,19 @@ def test_sample_sheet_validation_skipped_on_empty_definitions():
validate_row([0, 1], None) # just ensure no exception is thrown
def test_sample_sheet_validation_skipped_on_empty_list_definitions():
# Empty list of column definitions should also short-circuit validation,
# so a missing row is not treated as an error.
validate_row(None, [])
validate_row([0, 1], [])
def test_sample_sheet_validation_missing_row_when_definitions_declared():
# When column definitions exist, a missing row is still an error.
with pytest.raises(RequestParameterInvalidException):
validate_row(None, [{"type": "int", "name": "replicate", "optional": False}])
def test_sample_sheet_validation_number_columns():
with pytest.raises(RequestParameterInvalidException):
validate_row([0, 1], [{"type": "int", "name": "replicate number", "default_value": 0, "optional": False}])
@@ -12,7 +12,7 @@ factory = CollectionTypeDescriptionFactory(None)
def test_get_structure_simple():
paired_type_description = factory.for_collection_type("paired")
tree = get_structure(pair_instance(), paired_type_description)
tree = get_structure(pair_instance().collection, paired_type_description)
assert len(tree.children) == 2
assert tree.children[0][0] == "left" # why not forward?
assert tree.children[0][1].is_leaf
@@ -20,7 +20,7 @@ def test_get_structure_simple():
def test_get_structure_list_paired_over_paired():
paired_type_description = factory.for_collection_type("list:paired")
tree = get_structure(list_paired_instance(), paired_type_description, "paired")
tree = get_structure(list_paired_instance().collection, paired_type_description, "paired")
assert tree.collection_type_description.collection_type == "list"
assert len(tree.children) == 3
assert tree.children[0][0] == "data1"
@@ -29,7 +29,7 @@ def test_get_structure_list_paired_over_paired():
def test_get_structure_list_of_lists():
list_of_lists_type_description = factory.for_collection_type("list:list")
tree = get_structure(list_of_lists_instance(), list_of_lists_type_description)
tree = get_structure(list_of_lists_instance().collection, list_of_lists_type_description)
assert tree.collection_type_description.collection_type == "list:list"
assert len(tree.children) == 2
assert tree.children[0][0] == "outer1"
@@ -38,7 +38,7 @@ def test_get_structure_list_of_lists():
def test_get_structure_list_of_lists_over_list():
list_of_lists_type_description = factory.for_collection_type("list:list")
tree = get_structure(list_of_lists_instance(), list_of_lists_type_description, "list")
tree = get_structure(list_of_lists_instance().collection, list_of_lists_type_description, "list")
assert tree.collection_type_description.collection_type == "list"
assert len(tree.children) == 2
assert tree.children[0][0] == "outer1"
@@ -47,7 +47,7 @@ def test_get_structure_list_of_lists_over_list():
def test_get_structure_list_paired_or_unpaired():
list_pair_or_unpaired_description = factory.for_collection_type("list:paired_or_unpaired")
tree = get_structure(list_of_paired_and_unpaired_instance(), list_pair_or_unpaired_description)
tree = get_structure(list_of_paired_and_unpaired_instance().collection, list_pair_or_unpaired_description)
assert tree.collection_type_description.collection_type == "list:paired_or_unpaired"
assert len(tree.children) == 2
assert tree.children[0][0] == "el1"
@@ -57,7 +57,7 @@ def test_get_structure_list_paired_or_unpaired():
def test_get_structure_list_paired_or_unpaired_over_paired_or_unpaired():
list_pair_or_unpaired_description = factory.for_collection_type("list:paired_or_unpaired")
tree = get_structure(
list_of_paired_and_unpaired_instance(), list_pair_or_unpaired_description, "paired_or_unpaired"
list_of_paired_and_unpaired_instance().collection, list_pair_or_unpaired_description, "paired_or_unpaired"
)
assert tree.collection_type_description.collection_type == "list"
assert len(tree.children) == 2
@@ -67,7 +67,7 @@ def test_get_structure_list_paired_or_unpaired_over_paired_or_unpaired():
def test_get_structure_list_of_lists_over_single_datasests():
list_of_lists_type_description = factory.for_collection_type("list:list")
tree = get_structure(list_of_lists_instance(), list_of_lists_type_description, "single_datasets")
tree = get_structure(list_of_lists_instance().collection, list_of_lists_type_description, "single_datasets")
assert tree.collection_type_description.collection_type == "list:list"
assert len(tree.children) == 2
assert tree.children[0][0] == "outer1"
+13
View File
@@ -51,6 +51,19 @@ TEST_PATH_2_CONVERTED = TESTCASE_DIRECTORY / "2.txt"
DEFAULT_OBJECT_STORE_BY = "id"
def test_get_export_dataset_filename_truncates_long_name():
long_name = "https___example.com_" + "a" * 2000 + ".fastq.gz"
filename = store.get_export_dataset_filename(long_name, "fastqsanger.gz", "abcdef1234567890", conversion_key=None)
assert len(filename.encode("utf-8")) <= 255
assert filename.endswith("_abcdef1234567890.fastqsanger.gz")
filename_conv = store.get_export_dataset_filename(
long_name, "bam", "abcdef1234567890", conversion_key="0123456789abcdef"
)
assert len(filename_conv.encode("utf-8")) <= 255
assert filename_conv.endswith("_abcdef1234567890_conversion_0123456789abcdef.bam")
def test_import_export_history():
"""Test a simple job import/export after decompressing an archive (like history import/export tool)."""
app = _mock_app()
+7
View File
@@ -549,3 +549,10 @@ def test_score_url_match_requires_prefix():
# Embedded scheme must not match
assert file_source.score_url_match("gxfiles://test1http://evil.com/foo") == 0
assert file_source.score_url_match("http://evil.com/gxfiles://test1/a") == 0
def test_get_file_source_path_strips_whitespace():
file_sources = _configured_file_sources()
resolved = file_sources.get_file_source_path("\ngxfiles://test1/a\n")
assert resolved.file_source is not None
assert resolved.path == "/a"
+25
View File
@@ -1,8 +1,11 @@
import ipaddress
import pytest
from galaxy.exceptions import (
AdminRequiredException,
ConfigDoesNotAllowException,
RequestParameterInvalidException,
)
from galaxy.files.uris import (
validate_non_local,
@@ -39,6 +42,28 @@ def test_validate():
assert validates("file://foo/bar", True, [])
@pytest.mark.parametrize(
"bad_url",
[
"https://.../SRR1957099.fastq.gz",
"https://./foo",
"https://..example.com/x",
"https:///foo",
"http:// /foo",
],
)
def test_validate_non_local_rejects_malformed_host(bad_url):
with pytest.raises(RequestParameterInvalidException):
validate_non_local(bad_url, [])
def test_validate_non_local_strips_surrounding_whitespace():
# Leading whitespace was already tolerated; trailing whitespace used to
# slip through to getaddrinfo and raise an opaque error.
assert validates_as_non_local(" http://google.com/x ", [])
assert validates_as_non_local("\thttp://google.com/x\n", [])
def validates_as_non_local(uri: str, allow_list):
try:
validate_non_local(uri, allow_list)
@@ -168,6 +168,15 @@ class TestTestParsing(TestCase):
test_dicts = self._parse_tests()
self._verify_each(test_dicts[1].to_dict(), BIGWIG_TO_WIG_EXPECTATIONS)
def test_harmonize_bare_ftype_output_check(self):
if in_packages():
skip("tool not available when running from packages")
tool_path = os.path.join(galaxy_directory(), "lib", "galaxy", "tools", "harmonize_two_collections_list.xml")
self._init_tool_for_path(tool_path)
test_dicts = [td.to_dict() for td in self._parse_tests()]
for td in test_dicts:
assert not td.get("exception"), f"Test failed to parse: {td.get('exception')}"
def _verify_each(self, target_dict: dict, expectations: List[Any]):
exception = target_dict.get("exception")
assert not exception, f"Test failed to generate with exception {exception}"