mirror of
https://github.com/galaxyproject/galaxy.git
synced 2026-09-24 16:30:27 +08:00
Merge branch 'release_26.0' into dev
This commit is contained in:
@@ -1,5 +1,6 @@
|
||||
import createClient from "openapi-fetch";
|
||||
|
||||
import { pendingRequestsMiddleware } from "@/api/client/pendingRequestsMiddleware";
|
||||
import { createRateLimiterMiddleware } from "@/api/client/rateLimiter";
|
||||
import type { GalaxyApiPaths } from "@/api/schema";
|
||||
import { getAppRoot } from "@/onload/loadConfig";
|
||||
@@ -12,6 +13,9 @@ function getBaseUrl() {
|
||||
function apiClientFactory() {
|
||||
const client = createClient<GalaxyApiPaths>({ baseUrl: getBaseUrl() });
|
||||
|
||||
// Registered first so aborted requests bypass the rate-limiter queue.
|
||||
client.use(pendingRequestsMiddleware);
|
||||
|
||||
// TODO: Adjust based on server limits (maybe this goes in Galaxy config?)
|
||||
client.use(
|
||||
createRateLimiterMiddleware({
|
||||
|
||||
@@ -0,0 +1,24 @@
|
||||
import type { Middleware } from "openapi-fetch";
|
||||
|
||||
import { getPendingAbortSignal, SKIP_PENDING_REQUESTS_HEADER } from "@/api/pendingRequests";
|
||||
|
||||
/**
|
||||
* Attaches the shared pending-requests signal to every ``GalaxyApi`` request.
|
||||
* The ``openapi-fetch`` client uses native ``fetch()`` so the axios
|
||||
* interceptor does not apply; without this middleware, login/register
|
||||
* navigations cannot cancel in-flight ``/api/...`` calls and a late
|
||||
* anonymous-cookie response can clobber the authenticated ``galaxysession``
|
||||
* cookie. See ``client/src/api/pendingRequests.ts`` for the race.
|
||||
*/
|
||||
export const pendingRequestsMiddleware: Middleware = {
|
||||
async onRequest({ request }) {
|
||||
if (request.headers.has(SKIP_PENDING_REQUESTS_HEADER)) {
|
||||
const headers = new Headers(request.headers);
|
||||
headers.delete(SKIP_PENDING_REQUESTS_HEADER);
|
||||
return new Request(request, { headers });
|
||||
}
|
||||
const shared = getPendingAbortSignal();
|
||||
const signal = typeof AbortSignal.any === "function" ? AbortSignal.any([request.signal, shared]) : shared;
|
||||
return new Request(request, { signal });
|
||||
},
|
||||
};
|
||||
@@ -0,0 +1,58 @@
|
||||
/**
|
||||
* Shared ``AbortController`` used by both the axios interceptor and the
|
||||
* ``openapi-fetch`` middleware so we can cancel every in-flight request in
|
||||
* one shot right before a login/register navigation.
|
||||
*
|
||||
* Why this exists: when the server processes a request that carries the old
|
||||
* anonymous ``galaxysession`` cookie *after* ``handle_user_login`` has marked
|
||||
* that session ``is_valid=False``, it creates a fresh anonymous session and
|
||||
* responds with ``Set-Cookie: galaxysession=<new>``. Delivered into the cookie
|
||||
* jar after the login response but before the new page loads, it overwrites
|
||||
* the just-issued authenticated cookie — so the new page loads anonymous and
|
||||
* ``wait_for_logged_in`` times out in selenium. Aborting the TCP connection
|
||||
* before the response headers are parsed prevents that ``Set-Cookie`` from
|
||||
* ever applying.
|
||||
*/
|
||||
import axios, { type InternalAxiosRequestConfig } from "axios";
|
||||
|
||||
let activeController = new AbortController();
|
||||
|
||||
/** Explicit opt-out header that the login/register POST itself sets. */
|
||||
export const SKIP_PENDING_REQUESTS_HEADER = "x-galaxy-skip-pending-abort";
|
||||
|
||||
/**
|
||||
* The signal every request should ride on by default. Read lazily so the
|
||||
* ``openapi-fetch`` middleware picks up the fresh signal after each
|
||||
* ``cancelPendingRequests()`` rotation.
|
||||
*/
|
||||
export function getPendingAbortSignal(): AbortSignal {
|
||||
return activeController.signal;
|
||||
}
|
||||
|
||||
/**
|
||||
* Install a request interceptor that attaches the shared signal to every
|
||||
* outgoing axios request that didn't set one itself. Call once at app boot.
|
||||
*/
|
||||
export function installPendingRequestsInterceptor() {
|
||||
axios.interceptors.request.use((config: InternalAxiosRequestConfig) => {
|
||||
if (config.signal !== undefined) {
|
||||
return config;
|
||||
}
|
||||
if (config.headers?.[SKIP_PENDING_REQUESTS_HEADER]) {
|
||||
delete config.headers[SKIP_PENDING_REQUESTS_HEADER];
|
||||
return config;
|
||||
}
|
||||
config.signal = activeController.signal;
|
||||
return config;
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* Abort every request that is using the shared signal (both axios via the
|
||||
* interceptor and ``openapi-fetch`` via its middleware) and install a fresh
|
||||
* controller so subsequent requests can still go out.
|
||||
*/
|
||||
export function cancelPendingRequests() {
|
||||
activeController.abort();
|
||||
activeController = new AbortController();
|
||||
}
|
||||
@@ -38,7 +38,7 @@ const isCurrentTargetHistory = computed(() => {
|
||||
v-else
|
||||
v-g-tooltip.hover
|
||||
data-description="not current history indicator"
|
||||
class="text-warning"
|
||||
class="text-warning text-nowrap"
|
||||
title="This history is not your currently active history. You can click the link to switch to it.">
|
||||
(not current)
|
||||
</span>
|
||||
|
||||
@@ -15,6 +15,8 @@ import {
|
||||
import { computed, ref } from "vue";
|
||||
import { useRouter } from "vue-router/composables";
|
||||
|
||||
import { SKIP_PENDING_REQUESTS_HEADER } from "@/api/pendingRequests";
|
||||
import { discardActiveConnectionsBeforeAuthNavigation } from "@/composables/useAuthNavigation";
|
||||
import localize from "@/utils/localization";
|
||||
import { withPrefix } from "@/utils/redirect";
|
||||
import { errorMessageAsString } from "@/utils/simple-error";
|
||||
@@ -86,13 +88,22 @@ async function submitLogin() {
|
||||
redirect = props.redirect ?? null;
|
||||
}
|
||||
|
||||
// Stop polling and abort in-flight axios/GalaxyApi before sending the
|
||||
// login POST — otherwise a late anonymous-cookie response can overwrite
|
||||
// the authenticated cookie we're about to receive.
|
||||
discardActiveConnectionsBeforeAuthNavigation();
|
||||
|
||||
try {
|
||||
const response = await axios.post(withPrefix("/user/login"), {
|
||||
login: login.value,
|
||||
password: password.value,
|
||||
redirect: redirect,
|
||||
session_csrf_token: props.sessionCsrfToken,
|
||||
});
|
||||
const response = await axios.post(
|
||||
withPrefix("/user/login"),
|
||||
{
|
||||
login: login.value,
|
||||
password: password.value,
|
||||
redirect: redirect,
|
||||
session_csrf_token: props.sessionCsrfToken,
|
||||
},
|
||||
{ headers: { [SKIP_PENDING_REQUESTS_HEADER]: "1" } },
|
||||
);
|
||||
|
||||
if (response.data.message && response.data.status) {
|
||||
alert(response.data.message);
|
||||
|
||||
@@ -43,6 +43,9 @@ const props = withDefaults(defineProps<Props>(), {
|
||||
position: absolute;
|
||||
text-align: center;
|
||||
width: 100%;
|
||||
white-space: nowrap;
|
||||
overflow: hidden;
|
||||
text-overflow: ellipsis;
|
||||
}
|
||||
.progress-container {
|
||||
position: relative;
|
||||
|
||||
@@ -14,8 +14,10 @@ import {
|
||||
} from "bootstrap-vue";
|
||||
import { computed, type Ref, ref } from "vue";
|
||||
|
||||
import { SKIP_PENDING_REQUESTS_HEADER } from "@/api/pendingRequests";
|
||||
import { getOIDCIdpsWithRegistration, type OIDCConfig } from "@/components/User/ExternalIdentities/ExternalIDHelper";
|
||||
import { Toast } from "@/composables/toast";
|
||||
import { discardActiveConnectionsBeforeAuthNavigation } from "@/composables/useAuthNavigation";
|
||||
import localize from "@/utils/localization";
|
||||
import { withPrefix } from "@/utils/redirect";
|
||||
import { errorMessageAsString } from "@/utils/simple-error";
|
||||
@@ -72,15 +74,24 @@ const registerColumnDisplay = computed(() => Boolean(props.termsUrl));
|
||||
async function submit() {
|
||||
disableCreate.value = true;
|
||||
|
||||
// Stop polling and abort in-flight axios/GalaxyApi before sending the
|
||||
// register POST — otherwise a late anonymous-cookie response can overwrite
|
||||
// the authenticated cookie we're about to receive.
|
||||
discardActiveConnectionsBeforeAuthNavigation();
|
||||
|
||||
try {
|
||||
const response = await axios.post(withPrefix("/user/create"), {
|
||||
email: email.value,
|
||||
username: username.value,
|
||||
password: password.value,
|
||||
confirm: confirm.value,
|
||||
subscribe: subscribe.value,
|
||||
session_csrf_token: props.sessionCsrfToken,
|
||||
});
|
||||
const response = await axios.post(
|
||||
withPrefix("/user/create"),
|
||||
{
|
||||
email: email.value,
|
||||
username: username.value,
|
||||
password: password.value,
|
||||
confirm: confirm.value,
|
||||
subscribe: subscribe.value,
|
||||
session_csrf_token: props.sessionCsrfToken,
|
||||
},
|
||||
{ headers: { [SKIP_PENDING_REQUESTS_HEADER]: "1" } },
|
||||
);
|
||||
|
||||
if (response.data.message && response.data.status) {
|
||||
Toast.info(response.data.message);
|
||||
|
||||
@@ -42,7 +42,7 @@ const workflowTags = computed(() => {
|
||||
<template>
|
||||
<div v-if="workflow" class="pb-2 pl-2">
|
||||
<div class="d-flex justify-content-between align-items-center">
|
||||
<div>
|
||||
<div class="annotation-left">
|
||||
<i v-if="timeElapsed" data-description="workflow annotation time info">
|
||||
<FontAwesomeIcon :icon="faClock" class="mr-1" />
|
||||
<span v-localize>
|
||||
@@ -52,9 +52,11 @@ const workflowTags = computed(() => {
|
||||
</i>
|
||||
<TargetHistoryLink v-if="props.invocationCreateTime" :target-history-id="props.historyId" />
|
||||
</div>
|
||||
<slot name="middle-content" />
|
||||
<div class="d-flex align-items-center">
|
||||
<div class="d-flex flex-column align-items-end mr-2 flex-gapy-1">
|
||||
<div class="annotation-middle">
|
||||
<slot name="middle-content" />
|
||||
</div>
|
||||
<div class="annotation-right">
|
||||
<div class="annotation-right-content">
|
||||
<WorkflowIndicators :workflow="workflow" published-view no-edit-time />
|
||||
<WorkflowInvocationsCount v-if="owned" class="mr-1" :workflow="workflow" />
|
||||
</div>
|
||||
@@ -69,16 +71,66 @@ const workflowTags = computed(() => {
|
||||
</template>
|
||||
|
||||
<style scoped lang="scss">
|
||||
.history-link-wrapper {
|
||||
max-width: 300px;
|
||||
// Left column: 35% of the width
|
||||
.annotation-left {
|
||||
flex: 1 1 0;
|
||||
max-width: 35%;
|
||||
min-width: 0;
|
||||
overflow: hidden;
|
||||
|
||||
&:deep(.history-link) {
|
||||
.history-link-click {
|
||||
overflow: hidden;
|
||||
white-space: nowrap;
|
||||
text-overflow: ellipsis;
|
||||
display: block;
|
||||
}
|
||||
// Ensure neither the history name nor "(current)" escape the column.
|
||||
:deep(.history-link-wrapper) {
|
||||
min-width: 0;
|
||||
overflow: hidden;
|
||||
}
|
||||
|
||||
// SwitchToHistoryLink root div must be able to shrink to 0 so the
|
||||
// "(current)" label stays visible as long as there is any room.
|
||||
:deep(.history-link-wrapper > div) {
|
||||
min-width: 0;
|
||||
overflow: hidden;
|
||||
}
|
||||
|
||||
:deep(.history-link) {
|
||||
min-width: 0;
|
||||
}
|
||||
|
||||
:deep(.history-link-click) {
|
||||
overflow: hidden;
|
||||
white-space: nowrap;
|
||||
text-overflow: ellipsis;
|
||||
display: block;
|
||||
}
|
||||
}
|
||||
|
||||
// Middle column: 30% of the width
|
||||
.annotation-middle {
|
||||
flex: 1 1 0;
|
||||
max-width: 30%;
|
||||
min-width: 0;
|
||||
}
|
||||
|
||||
// Right column: 35% of the width
|
||||
.annotation-right {
|
||||
flex: 1 1 0;
|
||||
max-width: 35%;
|
||||
min-width: 0;
|
||||
overflow: hidden;
|
||||
justify-content: flex-end;
|
||||
|
||||
.annotation-right-content {
|
||||
display: flex;
|
||||
flex-direction: column;
|
||||
align-items: flex-end;
|
||||
gap: 0.25rem;
|
||||
min-width: 0;
|
||||
overflow: hidden;
|
||||
}
|
||||
|
||||
// Ensure creator badges stay within the column instead of overflowing
|
||||
:deep(.workflow-indicators) {
|
||||
flex-wrap: wrap;
|
||||
justify-content: flex-end;
|
||||
}
|
||||
}
|
||||
</style>
|
||||
|
||||
@@ -0,0 +1,25 @@
|
||||
import { cancelPendingRequests } from "@/api/pendingRequests";
|
||||
import { useEntryPointStore } from "@/stores/entryPointStore";
|
||||
import { useHistoryStore } from "@/stores/historyStore";
|
||||
import { useNotificationsStore } from "@/stores/notificationsStore";
|
||||
|
||||
/**
|
||||
* Tear down every long-lived connection and cancel every in-flight request
|
||||
* (both axios and ``openapi-fetch``/GalaxyApi) so that nothing issued under
|
||||
* the old anonymous ``galaxysession`` cookie can land on the server after
|
||||
* ``handle_user_login`` has invalidated it. See
|
||||
* ``client/src/api/pendingRequests.ts`` for the race this guards against.
|
||||
*
|
||||
* Call this synchronously as the first step of a login or registration
|
||||
* submit, before the authenticating POST goes out. The shared abort
|
||||
* controller is rotated, so the login/register POST (issued right after)
|
||||
* will use a fresh signal and is not affected.
|
||||
*/
|
||||
export function discardActiveConnectionsBeforeAuthNavigation() {
|
||||
// Stop polling watchers first so they can't kick off new fetches, then
|
||||
// abort any requests still in flight via the shared AbortController.
|
||||
useHistoryStore().stopWatchingHistory();
|
||||
useEntryPointStore().stopWatchingEntryPoints();
|
||||
useNotificationsStore().stopWatchingNotifications();
|
||||
cancelPendingRequests();
|
||||
}
|
||||
@@ -2,6 +2,7 @@
|
||||
import { createPinia, PiniaVuePlugin } from "pinia";
|
||||
import Vue from "vue";
|
||||
|
||||
import { installPendingRequestsInterceptor } from "@/api/pendingRequests";
|
||||
import { initGalaxyInstance } from "@/app";
|
||||
import { initSentry } from "@/app/addons/sentry";
|
||||
import { initWebhooks } from "@/app/addons/webhooks";
|
||||
@@ -13,6 +14,12 @@ import App from "./App.vue";
|
||||
Vue.use(PiniaVuePlugin);
|
||||
const pinia = createPinia();
|
||||
|
||||
// Attach the shared AbortController signal to every outgoing axios request
|
||||
// so we can cancel in-flight anonymous-cookie requests before login/register
|
||||
// navigates — otherwise their late ``Set-Cookie: galaxysession=<anon>`` can
|
||||
// clobber the authenticated cookie.
|
||||
installPendingRequestsInterceptor();
|
||||
|
||||
window.addEventListener("load", async () => {
|
||||
// Create Galaxy object
|
||||
const Galaxy = await initGalaxyInstance();
|
||||
|
||||
@@ -23,10 +23,11 @@ interface EntryPoint {
|
||||
}
|
||||
|
||||
export const useEntryPointStore = defineStore("entryPointStore", () => {
|
||||
const { startWatchingResource: startWatchingEntryPoints } = useResourceWatcher(fetchEntryPoints, {
|
||||
shortPollingInterval: ACTIVE_POLLING_INTERVAL,
|
||||
enableBackgroundPolling: false, // No need to poll in the background
|
||||
});
|
||||
const { startWatchingResource: startWatchingEntryPoints, stopWatchingResource: stopWatchingEntryPoints } =
|
||||
useResourceWatcher(fetchEntryPoints, {
|
||||
shortPollingInterval: ACTIVE_POLLING_INTERVAL,
|
||||
enableBackgroundPolling: false, // No need to poll in the background
|
||||
});
|
||||
|
||||
const entryPoints = ref<EntryPoint[]>([]);
|
||||
|
||||
@@ -92,5 +93,6 @@ export const useEntryPointStore = defineStore("entryPointStore", () => {
|
||||
updateEntryPoints,
|
||||
removeEntryPoint,
|
||||
startWatchingEntryPoints,
|
||||
stopWatchingEntryPoints,
|
||||
};
|
||||
});
|
||||
|
||||
@@ -391,13 +391,14 @@ export const useHistoryStore = defineStore("historyStore", () => {
|
||||
return watchHistorySuppliedApp(app);
|
||||
}
|
||||
|
||||
const { startWatchingResource: startWatchingHistory, isWatchingResource: isWatchingHistory } = useResourceWatcher(
|
||||
watchHistory,
|
||||
{
|
||||
shortPollingInterval: ACTIVE_POLLING_INTERVAL,
|
||||
longPollingInterval: INACTIVE_POLLING_INTERVAL,
|
||||
},
|
||||
);
|
||||
const {
|
||||
startWatchingResource: startWatchingHistory,
|
||||
stopWatchingResource: stopWatchingHistory,
|
||||
isWatchingResource: isWatchingHistory,
|
||||
} = useResourceWatcher(watchHistory, {
|
||||
shortPollingInterval: ACTIVE_POLLING_INTERVAL,
|
||||
longPollingInterval: INACTIVE_POLLING_INTERVAL,
|
||||
});
|
||||
|
||||
async function loadHistoryById(historyId: string) {
|
||||
if (!isLoadingHistory.has(historyId)) {
|
||||
@@ -525,6 +526,7 @@ export const useHistoryStore = defineStore("historyStore", () => {
|
||||
restoreHistories,
|
||||
handleTotalCountChange,
|
||||
startWatchingHistory,
|
||||
stopWatchingHistory,
|
||||
isWatchingHistory,
|
||||
loadCurrentHistory,
|
||||
loadCurrentHistoryId,
|
||||
|
||||
@@ -13,10 +13,11 @@ const ACTIVE_POLLING_INTERVAL = 30000; // 30 seconds
|
||||
const INACTIVE_POLLING_INTERVAL = ACTIVE_POLLING_INTERVAL * 20; // 10 minutes
|
||||
|
||||
export const useNotificationsStore = defineStore("notificationsStore", () => {
|
||||
const { startWatchingResource: startWatchingNotifications } = useResourceWatcher(getNotificationStatus, {
|
||||
shortPollingInterval: ACTIVE_POLLING_INTERVAL,
|
||||
longPollingInterval: INACTIVE_POLLING_INTERVAL,
|
||||
});
|
||||
const { startWatchingResource: startWatchingNotifications, stopWatchingResource: stopWatchingNotifications } =
|
||||
useResourceWatcher(getNotificationStatus, {
|
||||
shortPollingInterval: ACTIVE_POLLING_INTERVAL,
|
||||
longPollingInterval: INACTIVE_POLLING_INTERVAL,
|
||||
});
|
||||
const broadcastsStore = useBroadcastsStore();
|
||||
|
||||
const totalUnreadCount = ref<number>(0);
|
||||
@@ -106,5 +107,6 @@ export const useNotificationsStore = defineStore("notificationsStore", () => {
|
||||
updateNotification,
|
||||
updateBatchNotification,
|
||||
startWatchingNotifications,
|
||||
stopWatchingNotifications,
|
||||
};
|
||||
});
|
||||
|
||||
@@ -638,10 +638,7 @@ class DatasetCollectionManager:
|
||||
|
||||
new_elements = {}
|
||||
for key, element in elements.items():
|
||||
if isinstance(element, DatasetCollection):
|
||||
continue
|
||||
|
||||
if element.get("src") != "new_collection":
|
||||
if not isinstance(element, dict) or element.get("src") != "new_collection":
|
||||
continue
|
||||
|
||||
# element is a dict with src new_collection and
|
||||
|
||||
@@ -240,15 +240,11 @@ class NotificationManager:
|
||||
return self._is_subscribed_to_category(category_settings)
|
||||
|
||||
def _send_via_channels(self, notification: Notification, user: User, channel_settings: NotificationChannelSettings):
|
||||
channels = channel_settings.model_fields_set
|
||||
for channel in channels:
|
||||
if channel not in self.channel_plugins:
|
||||
continue # Skip unsupported channels
|
||||
for channel, plugin in self.channel_plugins.items():
|
||||
user_opted_out = getattr(channel_settings, channel, False) is False
|
||||
if user_opted_out and not self._is_urgent(notification):
|
||||
continue # Skip sending to opted-out users unless it's an urgent notification
|
||||
try:
|
||||
plugin = self.channel_plugins[channel]
|
||||
plugin.send(notification, user)
|
||||
except Exception as e:
|
||||
log.error(
|
||||
|
||||
@@ -248,7 +248,10 @@ class TourGenerator:
|
||||
dataset = self._test.inputs[tour_id][0]
|
||||
step_msg = f"Select dataset: <b>{hid}: {dataset}</b>"
|
||||
else:
|
||||
case_params = ", ".join(self._test.inputs[tour_id])
|
||||
case_params = ", ".join(
|
||||
("Yes" if v is True else "No" if v is False else str(v))
|
||||
for v in self._test.inputs[tour_id]
|
||||
)
|
||||
step_msg = f"Select parameter(s): <b>{case_params}</b>"
|
||||
cond_case_steps.append(
|
||||
TourStep(
|
||||
|
||||
@@ -37,7 +37,6 @@ from sqlalchemy import (
|
||||
true,
|
||||
)
|
||||
from sqlalchemy.orm import (
|
||||
aliased,
|
||||
joinedload,
|
||||
subqueryload,
|
||||
)
|
||||
@@ -76,10 +75,10 @@ from galaxy.model import (
|
||||
)
|
||||
from galaxy.model.base import ensure_object_added_to_session
|
||||
from galaxy.model.index_filter_util import (
|
||||
append_user_filter,
|
||||
raw_text_column_filter,
|
||||
tag_filter,
|
||||
tag_exists_filter,
|
||||
text_column_filter,
|
||||
user_exists_filter,
|
||||
)
|
||||
from galaxy.model.item_attrs import UsesAnnotations
|
||||
from galaxy.schema.invocation import InvocationCancellationUserRequest
|
||||
@@ -105,6 +104,7 @@ from galaxy.util.json import (
|
||||
)
|
||||
from galaxy.util.sanitize_html import sanitize_html
|
||||
from galaxy.util.search import (
|
||||
filter_terms,
|
||||
FilteredTerm,
|
||||
parse_filters_structured,
|
||||
RawTextTerm,
|
||||
@@ -215,13 +215,16 @@ class WorkflowsManager(sharable.SharableModelManager[model.StoredWorkflow], dele
|
||||
stmt = stmt.where(StoredWorkflow.hidden == (true() if show_hidden else false()))
|
||||
if payload.search:
|
||||
search_query = payload.search
|
||||
parsed_search = parse_filters_structured(search_query, INDEX_SEARCH_FILTERS)
|
||||
parsed_search = filter_terms(parse_filters_structured(search_query, INDEX_SEARCH_FILTERS))
|
||||
|
||||
def w_tag_filter(term_text: str, quoted: bool):
|
||||
nonlocal stmt
|
||||
alias = aliased(StoredWorkflowTagAssociation)
|
||||
stmt = stmt.outerjoin(StoredWorkflow.tags.of_type(alias))
|
||||
return tag_filter(alias, term_text, quoted)
|
||||
def w_tag_exists(term_text: str, quoted: bool):
|
||||
return tag_exists_filter(
|
||||
StoredWorkflowTagAssociation,
|
||||
StoredWorkflowTagAssociation.stored_workflow_id,
|
||||
StoredWorkflow.id,
|
||||
term_text,
|
||||
quoted,
|
||||
)
|
||||
|
||||
def name_filter(term):
|
||||
return text_column_filter(StoredWorkflow.name, term)
|
||||
@@ -231,12 +234,11 @@ class WorkflowsManager(sharable.SharableModelManager[model.StoredWorkflow], dele
|
||||
key = term.filter
|
||||
q = term.text
|
||||
if key == "tag":
|
||||
tf = w_tag_filter(term.text, term.quoted)
|
||||
stmt = stmt.where(tf)
|
||||
stmt = stmt.where(w_tag_exists(term.text, term.quoted))
|
||||
elif key == "name":
|
||||
stmt = stmt.where(name_filter(term))
|
||||
elif key == "user":
|
||||
stmt = append_user_filter(stmt, StoredWorkflow, term)
|
||||
stmt = stmt.where(user_exists_filter(StoredWorkflow.user_id, term.text))
|
||||
elif key == "is":
|
||||
if q == "published":
|
||||
stmt = stmt.where(StoredWorkflow.published == true())
|
||||
@@ -264,15 +266,12 @@ class WorkflowsManager(sharable.SharableModelManager[model.StoredWorkflow], dele
|
||||
model.StoredWorkflowMenuEntry.stored_workflow_id == StoredWorkflow.id,
|
||||
).where(model.StoredWorkflowMenuEntry.user_id == user.id)
|
||||
elif isinstance(term, RawTextTerm):
|
||||
tf = w_tag_filter(term.text, False)
|
||||
alias = aliased(User)
|
||||
stmt = stmt.outerjoin(StoredWorkflow.user.of_type(alias))
|
||||
stmt = stmt.where(
|
||||
raw_text_column_filter(
|
||||
[
|
||||
StoredWorkflow.name,
|
||||
tf,
|
||||
alias.username,
|
||||
w_tag_exists(term.text, False),
|
||||
user_exists_filter(StoredWorkflow.user_id, term.text),
|
||||
],
|
||||
term,
|
||||
)
|
||||
|
||||
@@ -7,6 +7,7 @@ from typing import (
|
||||
from sqlalchemy import (
|
||||
and_,
|
||||
or_,
|
||||
select,
|
||||
)
|
||||
from sqlalchemy.orm import (
|
||||
aliased,
|
||||
@@ -68,3 +69,32 @@ def append_user_filter(query, model_class, term: FilteredTerm):
|
||||
query = query.outerjoin(model_class.user.of_type(alias))
|
||||
query = query.filter(text_column_filter(alias.username, term))
|
||||
return query
|
||||
|
||||
|
||||
def tag_exists_filter(association_model_class, fk_column, parent_id_column, term_text, quoted: bool = False):
|
||||
"""Correlated EXISTS subquery that matches any tag on the parent row against term_text.
|
||||
|
||||
Prefer this over adding a per-term outer join on the tag-association table: each
|
||||
extra outer join multiplies rows (forcing an expensive DISTINCT) and in free-text
|
||||
search N whitespace-separated terms produce N such joins.
|
||||
"""
|
||||
return (
|
||||
select(1)
|
||||
.select_from(association_model_class)
|
||||
.where(fk_column == parent_id_column)
|
||||
.where(tag_filter(association_model_class, term_text, quoted))
|
||||
.correlate_except(association_model_class)
|
||||
.exists()
|
||||
)
|
||||
|
||||
|
||||
def user_exists_filter(owner_id_column, term_text: str):
|
||||
"""Correlated EXISTS subquery that matches the owning user's username."""
|
||||
return (
|
||||
select(1)
|
||||
.select_from(model.User)
|
||||
.where(model.User.id == owner_id_column)
|
||||
.where(model.User.username.ilike(f"%{term_text}%"))
|
||||
.correlate_except(model.User)
|
||||
.exists()
|
||||
)
|
||||
|
||||
@@ -1,4 +1,7 @@
|
||||
from datetime import datetime
|
||||
from datetime import (
|
||||
datetime,
|
||||
timezone,
|
||||
)
|
||||
|
||||
# NOTE REGARDING TIMESTAMPS:
|
||||
# It is currently difficult to have the timestamps calculated by the
|
||||
@@ -7,7 +10,12 @@ from datetime import datetime
|
||||
# relies on the client's clock being set correctly, so if clustering
|
||||
# web servers, use a time server to ensure synchronization
|
||||
|
||||
# Return the current time in UTC without any timezone information
|
||||
now = datetime.utcnow
|
||||
|
||||
def now():
|
||||
"""
|
||||
Return the current time in UTC without any timezone information.
|
||||
"""
|
||||
return datetime.now(timezone.utc).replace(tzinfo=None)
|
||||
|
||||
|
||||
__all__ = ("now",)
|
||||
|
||||
@@ -12,6 +12,12 @@ KeyedQueryT = Tuple[str, str]
|
||||
ParseFilterResultT = Tuple[Optional[List["FilteredTerm"]], Optional[str]]
|
||||
QUOTE_PATTERN = re.compile(r"\'(.*?)\'")
|
||||
|
||||
# Defaults for `filter_terms` used by index-search callers. A whitespace-rich
|
||||
# query turns into one WHERE clause (and, pre-trigram-index, one seq scan per
|
||||
# matching table) per raw term, so both floors are there to bound query cost.
|
||||
DEFAULT_MIN_RAW_TERM_LENGTH = 4
|
||||
DEFAULT_MAX_RAW_TERMS = 7
|
||||
|
||||
|
||||
def parse_filters(search_term: str, filters: Optional[Dict[str, str]] = None) -> ParseFilterResultT:
|
||||
"""Support github-like filters for narrowing the results.
|
||||
@@ -110,7 +116,39 @@ class ParsedSearch:
|
||||
return None if len(self.filter_terms) == 0 else self.filter_terms, " ".join([t.text for t in self.text_terms])
|
||||
|
||||
|
||||
def filter_terms(
|
||||
parsed: "ParsedSearch",
|
||||
min_raw_term_length: int = DEFAULT_MIN_RAW_TERM_LENGTH,
|
||||
max_raw_terms: Optional[int] = DEFAULT_MAX_RAW_TERMS,
|
||||
) -> "ParsedSearch":
|
||||
"""Return a new ParsedSearch with short / excess raw text terms dropped.
|
||||
|
||||
Raw (unquoted, non-keyed) terms shorter than ``min_raw_term_length`` are
|
||||
dropped, and the surviving raw terms are capped at ``max_raw_terms``.
|
||||
Filtered terms (``key:value``) and quoted raw terms ('foo bar') are
|
||||
always kept — those are explicit user intent.
|
||||
"""
|
||||
out = ParsedSearch()
|
||||
raw_kept = 0
|
||||
for term in parsed.terms:
|
||||
if isinstance(term, RawTextTerm) and not term.quoted:
|
||||
if len(term.text) < min_raw_term_length:
|
||||
continue
|
||||
if max_raw_terms is not None and raw_kept >= max_raw_terms:
|
||||
continue
|
||||
raw_kept += 1
|
||||
out.add_unfiltered_text(term.text, term.quoted)
|
||||
elif isinstance(term, RawTextTerm):
|
||||
out.add_unfiltered_text(term.text, term.quoted)
|
||||
else:
|
||||
out.add_keyed_term(term.filter, term.text, term.quoted)
|
||||
return out
|
||||
|
||||
|
||||
__all__ = (
|
||||
"DEFAULT_MAX_RAW_TERMS",
|
||||
"DEFAULT_MIN_RAW_TERM_LENGTH",
|
||||
"filter_terms",
|
||||
"parse_filters",
|
||||
"parse_filters_structured",
|
||||
)
|
||||
|
||||
@@ -228,10 +228,15 @@ class UrlBuilder:
|
||||
url = f"{url}?{urlencode(query_params)}"
|
||||
return url
|
||||
except NoMatchFound:
|
||||
# Fallback to legacy url_for
|
||||
# Fallback to legacy WSGI url_for for routes not registered with FastAPI
|
||||
if query_params:
|
||||
path_params.update(query_params)
|
||||
return web.url_for(name, **path_params)
|
||||
url = web.url_for(name, **path_params)
|
||||
if qualified and not url.startswith(("http://", "https://")):
|
||||
# routes.url_for has no thread-local request_config in an ASGI
|
||||
# request, so qualify the URL using the FastAPI request base_url.
|
||||
url = str(self.request.base_url).rstrip("/") + url
|
||||
return url
|
||||
|
||||
def _url_path_for(self, name: str, **path_params) -> str:
|
||||
"""O(1) route lookup using the app's pre-built name index.
|
||||
|
||||
@@ -8,6 +8,7 @@ from markupsafe import escape
|
||||
|
||||
import galaxy.util
|
||||
from galaxy import web
|
||||
from galaxy.exceptions import RequestParameterInvalidException
|
||||
from galaxy.model import HistoryDatasetAssociation
|
||||
from galaxy.tool_util.identifiers import uri_safe_tool_id
|
||||
from galaxy.tools import DataSourceTool
|
||||
@@ -39,6 +40,14 @@ class ToolRunner(BaseUIController):
|
||||
return self.index(trans, tool_id=tool_id, **kwd)
|
||||
|
||||
def __get_tool(self, tool_id, tool_version=None):
|
||||
# webob's params.mixed() returns a list when a form/query key is repeated
|
||||
# (some data sources redirect back with tool_id duplicated); accept that
|
||||
# case only when every value agrees.
|
||||
if isinstance(tool_id, list):
|
||||
unique_ids = set(tool_id)
|
||||
if len(unique_ids) != 1:
|
||||
raise RequestParameterInvalidException(f"Conflicting tool_id values supplied: {tool_id!r}")
|
||||
tool_id = unique_ids.pop()
|
||||
# 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)
|
||||
|
||||
@@ -17,6 +17,7 @@ from galaxy import model
|
||||
from galaxy.exceptions import HandlerAssignmentError
|
||||
from galaxy.jobs.handler import InvocationGrabber
|
||||
from galaxy.model.base import check_database_connection
|
||||
from galaxy.model.orm.now import now
|
||||
from galaxy.schema.invocation import (
|
||||
FailureReason,
|
||||
InvocationFailureDatasetFailed,
|
||||
@@ -350,12 +351,12 @@ class WorkflowRequestMonitor(Monitors):
|
||||
invocation_step_update_time := invocation.get_last_workflow_invocation_step_update_time()
|
||||
):
|
||||
do_schedule = invocation_step_update_time > last_schedule_time
|
||||
if not do_schedule and (datetime.now() - last_schedule_time) > self.timedelta:
|
||||
if not do_schedule and (now() - last_schedule_time) > self.timedelta:
|
||||
# If we haven't scheduled in a while, schedule anyway.
|
||||
log.debug(
|
||||
"Scheduling workflow invocation [%s] after %s seconds without scheduling.",
|
||||
invocation.id,
|
||||
(datetime.now() - last_schedule_time).total_seconds(),
|
||||
(now() - last_schedule_time).total_seconds(),
|
||||
)
|
||||
do_schedule = True
|
||||
return do_schedule
|
||||
@@ -459,7 +460,7 @@ class WorkflowRequestMonitor(Monitors):
|
||||
if i.active and i.id < workflow_invocation.id:
|
||||
return False
|
||||
if self.ready_to_schedule_more(workflow_invocation):
|
||||
self.update_time_tracking_dict[invocation_id] = datetime.now()
|
||||
self.update_time_tracking_dict[invocation_id] = now()
|
||||
workflow_scheduler.schedule(workflow_invocation)
|
||||
log.debug("Workflow invocation [%s] scheduled", invocation_id)
|
||||
except Exception:
|
||||
|
||||
@@ -23,6 +23,7 @@ from galaxy.util import galaxy_root_path
|
||||
from galaxy_test.base import rules_test_data
|
||||
from galaxy_test.base.api_asserts import (
|
||||
assert_error_code_is,
|
||||
assert_error_message_contains,
|
||||
assert_file_looks_like_xlsx,
|
||||
assert_has_keys,
|
||||
assert_status_code_is,
|
||||
@@ -850,6 +851,24 @@ class TestToolsApi(ApiTestCase, TestsTools):
|
||||
)
|
||||
assert zipped_hdca["collection_type"] == "list:paired"
|
||||
|
||||
@skip_without_tool("__MERGE_COLLECTION__")
|
||||
def test_merge_collection_rejects_structurally_invalid_inputs(self):
|
||||
with self.dataset_populator.test_history(require_new=False) as history_id:
|
||||
list_paired_id = self.dataset_collection_populator.create_list_of_pairs_in_history(
|
||||
history_id, wait=True
|
||||
).json()["outputs"][0]["id"]
|
||||
plain_list_id = self.dataset_collection_populator.create_list_in_history(
|
||||
history_id, contents=["a", "b"], wait=True
|
||||
).json()["outputs"][0]["id"]
|
||||
self.dataset_populator.wait_for_history(history_id, assert_ok=True)
|
||||
inputs = {
|
||||
"inputs_0|input": {"src": "hdca", "id": list_paired_id},
|
||||
"inputs_1|input": {"src": "hdca", "id": plain_list_id},
|
||||
}
|
||||
response = self._run("__MERGE_COLLECTION__", history_id, inputs, assert_ok=False)
|
||||
assert_status_code_is(response, 400)
|
||||
assert_error_message_contains(response, "must be a sub-collection")
|
||||
|
||||
@skip_without_tool("__EXTRACT_DATASET__")
|
||||
@skip_without_tool("cat_data_and_sleep")
|
||||
def test_database_operation_tool_with_pending_inputs(self):
|
||||
|
||||
@@ -690,6 +690,17 @@ steps:
|
||||
assert workflow_id_2 not in index_ids
|
||||
assert workflow_id_3 not in index_ids
|
||||
|
||||
def test_index_search_many_terms(self):
|
||||
# Regression: a whitespace-rich search string used to add one outer join
|
||||
# on stored_workflow_tag_association and one on galaxy_user per term,
|
||||
# producing an unusably expensive query for long searches.
|
||||
name = f"Copy of Genomic Assembly and analysis - RDH shared by user {uuid4()}"
|
||||
workflow_id = self.workflow_populator.simple_workflow(name)
|
||||
self.workflow_populator.set_tags(workflow_id, [f"manyterms-{uuid4()}"])
|
||||
search = "Copy of Genomic Assembly and analysis - RDH shared by user"
|
||||
index_ids = self.workflow_populator.index_ids(search=search)
|
||||
assert workflow_id in index_ids
|
||||
|
||||
def test_search_casing(self):
|
||||
name1, name2 = (
|
||||
self.dataset_populator.get_random_name().upper(),
|
||||
|
||||
@@ -61,3 +61,24 @@ class TestToolDescribingTours(SeleniumTestCase):
|
||||
self.tool_form_execute()
|
||||
self.history_panel_wait_for_hid_ok(2)
|
||||
self.screenshot("tool_describing_tour_3_after_execute")
|
||||
|
||||
@selenium_test
|
||||
def test_generate_tour_boolean_conditional(self):
|
||||
self.tool_open("gx_conditional_boolean")
|
||||
self.tool_form_generate_tour()
|
||||
popover_component = self.components.tour.popover._
|
||||
popover_component.wait_for_visible()
|
||||
|
||||
# Intro step: advance to the outer conditional step.
|
||||
popover_component.next.wait_for_and_click()
|
||||
self.sleep_for(self.wait_types.UX_RENDER)
|
||||
# Advance to the inner boolean_parameter case step.
|
||||
popover_component.next.wait_for_and_click()
|
||||
self.sleep_for(self.wait_types.UX_RENDER)
|
||||
|
||||
# tests[0] specifies boolean_parameter="true" → tour should render "Yes".
|
||||
text = popover_component.content.wait_for_visible().text
|
||||
assert "Yes" in text, text
|
||||
|
||||
popover_component.end.wait_for_and_click()
|
||||
popover_component.wait_for_absent_or_hidden()
|
||||
|
||||
@@ -16,6 +16,19 @@
|
||||
</conditional>
|
||||
</inputs>
|
||||
<tests>
|
||||
<test>
|
||||
<conditional name="conditional_parameter">
|
||||
<param name="test_parameter" value="false" />
|
||||
<param name="boolean_parameter" value="true" />
|
||||
</conditional>
|
||||
<expand macro="assert_output">
|
||||
<has_line line="test: false" />
|
||||
</expand>
|
||||
<expand macro="assert_inputs_json">
|
||||
<has_json_property_with_value property="test_parameter" value="false" />
|
||||
<has_json_property_with_value property="boolean_parameter" value="true" />
|
||||
</expand>
|
||||
</test>
|
||||
<test>
|
||||
<expand macro="assert_output">
|
||||
<has_line line="test: false" />
|
||||
|
||||
@@ -1,3 +1,6 @@
|
||||
import json
|
||||
import os
|
||||
import re
|
||||
from typing import (
|
||||
ClassVar,
|
||||
)
|
||||
@@ -91,3 +94,31 @@ class TestUnexposedUsersIntegration(UsersIntegrationCase):
|
||||
# And the current user has all fields, so no limited fields.
|
||||
expected_limited_user_keys = set()
|
||||
expected_regular_user_list_count = 1
|
||||
|
||||
|
||||
class TestAdminResendActivationEmail(integration_util.IntegrationTestCase):
|
||||
email_directory: ClassVar[str]
|
||||
|
||||
@classmethod
|
||||
def handle_galaxy_config_kwds(cls, config):
|
||||
super().handle_galaxy_config_kwds(config)
|
||||
cls.email_directory = cls._test_driver.mkdtemp()
|
||||
config["user_activation_on"] = True
|
||||
config["activation_grace_period"] = 3
|
||||
config["email_from"] = "galaxy-noreply@example.com"
|
||||
config["smtp_server"] = f"mock_emails_to_path://{cls.email_directory}/email.json"
|
||||
|
||||
def test_resend_activation_includes_qualified_link(self):
|
||||
user = self._setup_user("resend-activation@test.gx")
|
||||
response = self._post(f"users/{user['id']}/send_activation_email", admin=True)
|
||||
self._assert_status_code_is_ok(response)
|
||||
|
||||
with open(os.path.join(self.email_directory, "email.json")) as f:
|
||||
email = json.loads(f.read())
|
||||
assert email["to"] == "resend-activation@test.gx"
|
||||
assert email["subject"] == "Galaxy Account Activation"
|
||||
match = re.search(r"(https?://[^/\s]+/user/activate\?[^\s]+)", email["body"])
|
||||
assert match, f"No qualified activation link found in email body:\n{email['body']}"
|
||||
link = match.group(1)
|
||||
assert "activation_token=" in link
|
||||
assert "email=" in link
|
||||
|
||||
@@ -6,7 +6,10 @@ from typing import (
|
||||
Any,
|
||||
Optional,
|
||||
)
|
||||
from unittest.mock import patch
|
||||
from unittest.mock import (
|
||||
MagicMock,
|
||||
patch,
|
||||
)
|
||||
|
||||
import pytest
|
||||
|
||||
@@ -395,6 +398,25 @@ class TestUserNotifications(NotificationManagerBaseTestCase):
|
||||
user_notifications = self.notification_manager.get_user_notifications(user)
|
||||
assert len(user_notifications) == 1
|
||||
|
||||
def test_send_via_channels_uses_all_channel_fields(self):
|
||||
user = self._create_test_user()
|
||||
notification_data = NotificationCreateData(**self._default_test_notification_data())
|
||||
notification = self.notification_manager._create_notification_model(notification_data)
|
||||
|
||||
push_plugin = MagicMock()
|
||||
email_plugin = MagicMock()
|
||||
self.notification_manager.channel_plugins = {
|
||||
"push": push_plugin,
|
||||
"email": email_plugin,
|
||||
}
|
||||
|
||||
# `email` remains at its default value (True) and should still be considered.
|
||||
channel_settings = NotificationChannelSettings(push=False)
|
||||
self.notification_manager._send_via_channels(notification, user, channel_settings)
|
||||
|
||||
push_plugin.send.assert_not_called()
|
||||
email_plugin.send.assert_called_once_with(notification, user)
|
||||
|
||||
|
||||
class TestUserNotificationsWithTasks(NotificationManagerBaseTestCaseWithTasks):
|
||||
|
||||
|
||||
@@ -1,6 +1,9 @@
|
||||
from galaxy.util.search import (
|
||||
filter_terms,
|
||||
FilteredTerm,
|
||||
parse_filters,
|
||||
parse_filters_structured,
|
||||
RawTextTerm,
|
||||
)
|
||||
|
||||
|
||||
@@ -94,3 +97,55 @@ def test_parse_filters_structured():
|
||||
assert text_terms[1].quoted is True
|
||||
assert text_terms[2].text == "foo"
|
||||
assert text_terms[2].quoted is False
|
||||
|
||||
|
||||
def test_filter_terms_drops_short_raw_terms():
|
||||
parsed = parse_filters_structured("Copy of Genomic Assembly and analysis", {})
|
||||
filtered = filter_terms(parsed, min_raw_term_length=4, max_raw_terms=None)
|
||||
kept = [t.text for t in filtered.terms]
|
||||
assert kept == ["Copy", "Genomic", "Assembly", "analysis"]
|
||||
|
||||
|
||||
def test_filter_terms_preserves_quoted_raw_terms():
|
||||
parsed = parse_filters_structured("'ab' Copy 'de'", {})
|
||||
filtered = filter_terms(parsed, min_raw_term_length=4, max_raw_terms=None)
|
||||
assert [(t.text, t.quoted) for t in filtered.terms] == [
|
||||
("ab", True),
|
||||
("Copy", False),
|
||||
("de", True),
|
||||
]
|
||||
|
||||
|
||||
def test_filter_terms_preserves_filtered_terms_of_any_length():
|
||||
parsed = parse_filters_structured("tag:ab user:cd Copy", {"tag": "tag", "user": "user"})
|
||||
filtered = filter_terms(parsed, min_raw_term_length=4, max_raw_terms=None)
|
||||
kinds = [(t.__class__.__name__, t.text) for t in filtered.terms]
|
||||
# Both filtered terms are preserved even though their text is shorter than 4;
|
||||
# "Copy" (4) is preserved too.
|
||||
assert ("FilteredTerm", "ab") in kinds
|
||||
assert ("FilteredTerm", "cd") in kinds
|
||||
assert ("RawTextTerm", "Copy") in kinds
|
||||
|
||||
|
||||
def test_filter_terms_caps_raw_terms_only():
|
||||
parsed = parse_filters_structured(
|
||||
"tag:foo Copy Genomic Assembly analysis shared user nedflanders extra1 extra2",
|
||||
{"tag": "tag"},
|
||||
)
|
||||
filtered = filter_terms(parsed, min_raw_term_length=4, max_raw_terms=3)
|
||||
raw = [t.text for t in filtered.terms if isinstance(t, RawTextTerm)]
|
||||
filt = [t.text for t in filtered.terms if isinstance(t, FilteredTerm)]
|
||||
assert raw == ["Copy", "Genomic", "Assembly"]
|
||||
assert filt == ["foo"]
|
||||
|
||||
|
||||
def test_filter_terms_defaults():
|
||||
# Default behaviour: 4-char min length, 7-term cap on raw terms.
|
||||
parsed = parse_filters_structured(
|
||||
"Copy of Genomic Assembly and analysis - RDH shared by user nedflanders",
|
||||
{},
|
||||
)
|
||||
filtered = filter_terms(parsed)
|
||||
kept = [t.text for t in filtered.terms]
|
||||
# of, and, -, RDH, by dropped for length; everything else survives.
|
||||
assert kept == ["Copy", "Genomic", "Assembly", "analysis", "shared", "user", "nedflanders"]
|
||||
|
||||
Reference in New Issue
Block a user