diff --git a/backend/.sqlx/query-caeb49629b8673c1f1c84a6e40c3e2d2c3bc3fdbde530a0a6b6fd68a22b867c3.json b/backend/.sqlx/query-4c529e29a0eb084a92f4e952506d701592bc707fd8a8cb3311bf146be6078f26.json similarity index 64% rename from backend/.sqlx/query-caeb49629b8673c1f1c84a6e40c3e2d2c3bc3fdbde530a0a6b6fd68a22b867c3.json rename to backend/.sqlx/query-4c529e29a0eb084a92f4e952506d701592bc707fd8a8cb3311bf146be6078f26.json index 95315372d0..3f345bc907 100644 --- a/backend/.sqlx/query-caeb49629b8673c1f1c84a6e40c3e2d2c3bc3fdbde530a0a6b6fd68a22b867c3.json +++ b/backend/.sqlx/query-4c529e29a0eb084a92f4e952506d701592bc707fd8a8cb3311bf146be6078f26.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "SELECT id, running, is_flow_step FROM queue WHERE scheduled_for < now() AND workspace_id = $1 AND schedule_path IS NULL", + "query": "SELECT id, is_flow_step, running FROM queue WHERE id = ANY($1) AND schedule_path IS NULL", "describe": { "columns": [ { @@ -10,25 +10,25 @@ }, { "ordinal": 1, - "name": "running", + "name": "is_flow_step", "type_info": "Bool" }, { "ordinal": 2, - "name": "is_flow_step", + "name": "running", "type_info": "Bool" } ], "parameters": { "Left": [ - "Text" + "UuidArray" ] }, "nullable": [ false, - false, - true + true, + false ] }, - "hash": "caeb49629b8673c1f1c84a6e40c3e2d2c3bc3fdbde530a0a6b6fd68a22b867c3" + "hash": "4c529e29a0eb084a92f4e952506d701592bc707fd8a8cb3311bf146be6078f26" } diff --git a/backend/.sqlx/query-661f472ff3860983322162420457f5033b9c9afc344d9c3e385ba20a3ad2197a.json b/backend/.sqlx/query-661f472ff3860983322162420457f5033b9c9afc344d9c3e385ba20a3ad2197a.json index 75b8108281..1fa370e682 100644 --- a/backend/.sqlx/query-661f472ff3860983322162420457f5033b9c9afc344d9c3e385ba20a3ad2197a.json +++ b/backend/.sqlx/query-661f472ff3860983322162420457f5033b9c9afc344d9c3e385ba20a3ad2197a.json @@ -5,7 +5,7 @@ "columns": [ { "ordinal": 0, - "name": "?column?", + "name": "bool", "type_info": "Bool" } ], diff --git a/backend/.sqlx/query-b3dbdfb50ee8118bdaed3164b210cb549a34b96554ae1872355b90304f5dcb76.json b/backend/.sqlx/query-b3dbdfb50ee8118bdaed3164b210cb549a34b96554ae1872355b90304f5dcb76.json index 99269c9851..54e94cfb8f 100644 --- a/backend/.sqlx/query-b3dbdfb50ee8118bdaed3164b210cb549a34b96554ae1872355b90304f5dcb76.json +++ b/backend/.sqlx/query-b3dbdfb50ee8118bdaed3164b210cb549a34b96554ae1872355b90304f5dcb76.json @@ -18,8 +18,8 @@ "Left": [] }, "nullable": [ - true, - false + false, + true ] }, "hash": "b3dbdfb50ee8118bdaed3164b210cb549a34b96554ae1872355b90304f5dcb76" diff --git a/backend/windmill-api/openapi.yaml b/backend/windmill-api/openapi.yaml index 8e95629715..e154ffe960 100644 --- a/backend/windmill-api/openapi.yaml +++ b/backend/windmill-api/openapi.yaml @@ -5418,14 +5418,75 @@ paths: required: - database_length - /w/{workspace}/jobs/queue/cancel_all: - post: - summary: cancel all jobs - operationId: cancelAll + /w/{workspace}/jobs/queue/list_filtered_uuids: + get: + summary: get the ids of all jobs matching the given filters + operationId: listFilteredUuids tags: - job parameters: - $ref: "#/components/parameters/WorkspaceId" + - $ref: "#/components/parameters/OrderDesc" + - $ref: "#/components/parameters/CreatedBy" + - $ref: "#/components/parameters/ParentJob" + - $ref: "#/components/parameters/ScriptExactPath" + - $ref: "#/components/parameters/ScriptStartPath" + - $ref: "#/components/parameters/SchedulePath" + - $ref: "#/components/parameters/ScriptExactHash" + - $ref: "#/components/parameters/StartedBefore" + - $ref: "#/components/parameters/StartedAfter" + - $ref: "#/components/parameters/Success" + - $ref: "#/components/parameters/ScheduledForBeforeNow" + - $ref: "#/components/parameters/JobKinds" + - $ref: "#/components/parameters/Suspended" + - $ref: "#/components/parameters/Running" + - $ref: "#/components/parameters/ArgsFilter" + - $ref: "#/components/parameters/ResultFilter" + - $ref: "#/components/parameters/Tag" + - $ref: "#/components/parameters/Page" + - $ref: "#/components/parameters/PerPage" + - name: concurrency_key + in: query + required: false + schema: + type: string + - name: all_workspaces + description: get jobs from all workspaces (only valid if request come from the `admins` workspace) + in: query + schema: + type: boolean + - name: is_not_schedule + description: is not a scheduled job + in: query + schema: + type: boolean + responses: + "200": + description: uuids of jobs + content: + application/json: + schema: + type: array + items: + type: string + + /w/{workspace}/jobs/queue/cancel_selection: + post: + summary: cancel jobs based on the given uuids + operationId: cancelSelection + tags: + - job + parameters: + - $ref: "#/components/parameters/WorkspaceId" + requestBody: + description: uuids of the jobs to cancel + required: true + content: + application/json: + schema: + type: array + items: + type: string responses: "200": description: uuids of canceled jobs diff --git a/backend/windmill-api/src/concurrency_groups.rs b/backend/windmill-api/src/concurrency_groups.rs index cfd1cbde2e..40099b2bb6 100644 --- a/backend/windmill-api/src/concurrency_groups.rs +++ b/backend/windmill-api/src/concurrency_groups.rs @@ -120,11 +120,10 @@ struct ObscuredJob { } #[derive(Deserialize)] struct ExtendedJobsParams { - concurrency_key: Option, row_limit: Option, } -fn join_concurrency_key<'c>(concurrency_key: Option<&String>, mut sqlb: SqlBuilder) -> SqlBuilder { +pub fn join_concurrency_key<'c>(concurrency_key: Option<&String>, mut sqlb: SqlBuilder) -> SqlBuilder { if let Some(key) = concurrency_key { sqlb.join("concurrency_key") .on_eq("id", "concurrency_key.job_id") @@ -151,7 +150,6 @@ async fn get_concurrent_intervals( } let row_limit = iq.row_limit.unwrap_or(1000); - let concurrency_key = iq.concurrency_key; let lq = ListCompletedQuery { order_desc: Some(true), ..lq }; let lqc = lq.clone(); @@ -177,10 +175,10 @@ async fn get_concurrent_intervals( .limit(row_limit) .clone(); - sqlb_q = join_concurrency_key(concurrency_key.as_ref(), sqlb_q); - sqlb_c = join_concurrency_key(concurrency_key.as_ref(), sqlb_c); - sqlb_q_user = join_concurrency_key(concurrency_key.as_ref(), sqlb_q_user); - sqlb_c_user = join_concurrency_key(concurrency_key.as_ref(), sqlb_c_user); + sqlb_q = join_concurrency_key(lq.concurrency_key.as_ref(), sqlb_q); + sqlb_c = join_concurrency_key(lq.concurrency_key.as_ref(), sqlb_c); + sqlb_q_user = join_concurrency_key(lq.concurrency_key.as_ref(), sqlb_q_user); + sqlb_c_user = join_concurrency_key(lq.concurrency_key.as_ref(), sqlb_c_user); let should_fetch_obscured_jobs = match lq { ListCompletedQuery { @@ -211,7 +209,8 @@ async fn get_concurrent_intervals( job_kinds: _, is_flow_step: _, all_workspaces: _, - } => concurrency_key.is_some(), + concurrency_key: Some(_), + } => true, _ => false, }; diff --git a/backend/windmill-api/src/jobs.rs b/backend/windmill-api/src/jobs.rs index 5a9e4edd17..fba8b8623d 100644 --- a/backend/windmill-api/src/jobs.rs +++ b/backend/windmill-api/src/jobs.rs @@ -28,6 +28,7 @@ use windmill_common::worker::TMP_DIR; use windmill_common::scripts::PREVIEW_IS_CODEBASE_HASH; use windmill_common::variables::get_workspace_key; +use crate::concurrency_groups::join_concurrency_key; use crate::add_webhook_allowed_origin; use crate::db::ApiAuthed; @@ -191,7 +192,8 @@ pub fn workspaced_service() -> Router { ) .route("/queue/list", get(list_queue_jobs)) .route("/queue/count", get(count_queue_jobs)) - .route("/queue/cancel_all", post(cancel_all)) + .route("/queue/list_filtered_uuids", get(list_filtered_uuids)) + .route("/queue/cancel_selection", post(cancel_selection)) .route("/completed/count", get(count_completed_jobs)) .route( "/completed/list", @@ -994,6 +996,7 @@ pub struct ListQueueQuery { pub is_flow_step: Option, pub has_null_parent: Option, pub is_not_schedule: Option, + pub concurrency_key: Option, } impl From for ListQueueQuery { @@ -1022,6 +1025,7 @@ impl From for ListQueueQuery { is_flow_step: lcq.is_flow_step, has_null_parent: lcq.has_null_parent, is_not_schedule: lcq.is_not_schedule, + concurrency_key: lcq.concurrency_key, } } } @@ -1217,23 +1221,20 @@ async fn list_queue_jobs( Ok(Json(jobs)) } -async fn cancel_all( - authed: ApiAuthed, - Extension(db): Extension, - Extension(rsmq): Extension>, +#[derive(Deserialize, FromRow)] +struct JobToCancel { + id: Uuid, + is_flow_step: Option, + running: bool, +} - Path(w_id): Path, +async fn cancel_jobs( + jobs: Vec, + db: &DB, + username: &str, + w_id: &str, + rsmq: Option, ) -> error::JsonResult> { - require_admin(authed.is_admin, &authed.username)?; - - let jobs = sqlx::query!( - "SELECT id, running, is_flow_step FROM queue WHERE scheduled_for < now() AND workspace_id = $1 AND schedule_path IS NULL", - w_id, - ) - .fetch_all(&db) - .await?; - - let username = authed.username; let mut uuids = vec![]; for j in jobs.iter() { let r = sqlx::query!( @@ -1241,7 +1242,7 @@ async fn cancel_all( username, j.id, ) - .fetch_optional(&db) + .fetch_optional(db) .await; if r.as_ref().is_ok_and(|x| x.is_some()) { @@ -1254,7 +1255,7 @@ async fn cancel_all( if let Some(job_running) = job_running { append_logs( &j.id, - w_id.clone(), + w_id, format!("canceled by {username}: cancel_all"), db.clone(), ) @@ -1286,6 +1287,62 @@ async fn cancel_all( Ok(Json(uuids)) } +async fn cancel_selection( + authed: ApiAuthed, + Extension(db): Extension, + Extension(user_db): Extension, + Extension(rsmq): Extension>, + + Path(w_id): Path, + Json(jobs): Json>, +) -> error::JsonResult> { + require_admin(authed.is_admin, &authed.username)?; + + let mut tx = user_db.begin(&authed).await?; + let jobs_to_cancel = sqlx::query_as!( + JobToCancel, + "SELECT id, is_flow_step, running FROM queue WHERE id = ANY($1) AND schedule_path IS NULL", + &jobs + ) + .fetch_all(&mut *tx) + .await?; + tx.commit().await?; + + cancel_jobs( + jobs_to_cancel, + &db, + authed.username.as_str(), + w_id.as_str(), + rsmq, + ) + .await +} + +async fn list_filtered_uuids( + authed: ApiAuthed, + Extension(db): Extension, + + Path(w_id): Path, + Query(lq): Query, +) -> error::JsonResult> { + require_admin(authed.is_admin, &authed.username)?; + + let mut sqlb = SqlBuilder::select_from("queue") + .fields(&["id"]) + .clone(); + + sqlb = join_concurrency_key(lq.concurrency_key.as_ref(), sqlb); + + sqlb.and_where_is_null("schedule_path"); + + sqlb = filter_list_queue_query(sqlb, &lq, w_id.as_str(), false); + + let sql = sqlb.query()?; + let jobs = sqlx::query_scalar(sql.as_str()).fetch_all(&db).await?; + + Ok(Json(jobs)) +} + #[derive(Serialize, Debug, FromRow)] struct QueueStats { database_length: i64, @@ -4326,6 +4383,7 @@ pub struct ListCompletedQuery { pub has_null_parent: Option, pub label: Option, pub is_not_schedule: Option, + pub concurrency_key: Option, } async fn list_completed_jobs( diff --git a/frontend/src/lib/components/common/modal/Modal.svelte b/frontend/src/lib/components/common/modal/Modal.svelte index 37e63642fa..718a52000a 100644 --- a/frontend/src/lib/components/common/modal/Modal.svelte +++ b/frontend/src/lib/components/common/modal/Modal.svelte @@ -10,6 +10,7 @@ let c: string = '' export { c as class } export let style = '' + export let cancelText: string | undefined = undefined const dispatch = createEventDispatcher() @@ -87,7 +88,9 @@ color="light" size="sm" > - Cancel Escape + {cancelText ?? 'Cancel'}Escape diff --git a/frontend/src/lib/components/common/waitTimeWarning/WaitTimeWarning.svelte b/frontend/src/lib/components/common/waitTimeWarning/WaitTimeWarning.svelte index 070f1312e8..7e70d51e0d 100644 --- a/frontend/src/lib/components/common/waitTimeWarning/WaitTimeWarning.svelte +++ b/frontend/src/lib/components/common/waitTimeWarning/WaitTimeWarning.svelte @@ -36,7 +36,7 @@ - +
{#if self_wait_time_ms != undefined}
@@ -66,7 +66,7 @@
{/if}
In a healthy queue, jobs are expected to start in under 50ms.
- + {#if variant === 'icon'} {:else if variant === 'badge'} diff --git a/frontend/src/lib/components/runs/RunRow.svelte b/frontend/src/lib/components/runs/RunRow.svelte index 5510ba1fef..93b8c58f69 100644 --- a/frontend/src/lib/components/runs/RunRow.svelte +++ b/frontend/src/lib/components/runs/RunRow.svelte @@ -31,12 +31,17 @@ export let containerWidth: number = 0 export let containsLabel: boolean = false export let activeLabel: string | null + export let isSelectingJobsToCancel: boolean = false let scheduleEditor: ScheduleEditor $: isExternal = job && job.id === '-' let triggeredByWidth: number = 0 + + function isJobCancelable(j: Job): boolean { + return j.type === 'QueuedJob' && !j.schedule_path + } @@ -52,10 +57,17 @@ )} style="width: {containerWidth}px" on:click={() => { - dispatch('select') + if (!isSelectingJobsToCancel || isJobCancelable(job)) { + dispatch('select') + } }} >
+ {#if isSelectingJobsToCancel && isJobCancelable(job)} +
+ +
+ {/if} {#if isExternal} diff --git a/frontend/src/lib/components/runs/RunsTable.svelte b/frontend/src/lib/components/runs/RunsTable.svelte index c08294260a..9f979cfa05 100644 --- a/frontend/src/lib/components/runs/RunsTable.svelte +++ b/frontend/src/lib/components/runs/RunsTable.svelte @@ -13,6 +13,7 @@ export let externalJobs: Job[] = [] export let omittedObscuredJobs: boolean export let showExternalJobs: boolean = false + export let isSelectingJobsToCancel: boolean = false export let selectedIds: string[] = [] export let selectedWorkspace: string | undefined = undefined export let activeLabel: string | null = null @@ -171,13 +172,14 @@ >
- {:else if $workspaceStore !== 'admins' && omittedObscuredJobs} + {:else if $workspaceStore !== 'admins' && omittedObscuredJobs}
{jobs && jobCountString(jobs.length)} - + - Too specific filtering may have caused the omission of obscured jobs. This is done for security reasons. To see obscured jobs, try removing some filters. + Too specific filtering may have caused the omission of obscured jobs. This is done for + security reasons. To see obscured jobs, try removing some filters.
@@ -214,10 +216,21 @@ {containsLabel} job={jobOrDate.job} selected={jobOrDate.job.id !== '-' && selectedIds.includes(jobOrDate.job.id)} + {isSelectingJobsToCancel} on:select={() => { - selectedWorkspace = jobOrDate.job.workspace_id - selectedIds = [jobOrDate.job.id] - dispatch('select') + const jobId = jobOrDate.job.id + if (isSelectingJobsToCancel) { + if (selectedIds.includes(jobOrDate.job.id)) { + selectedIds = selectedIds.filter((id) => id != jobId) + } else { + selectedIds.push(jobId) + selectedIds = selectedIds + } + } else { + selectedWorkspace = jobOrDate.job.workspace_id + selectedIds = [jobOrDate.job.id] + dispatch('select') + } }} {activeLabel} on:filterByLabel diff --git a/frontend/src/routes/(root)/(logged)/runs/[...path]/+page.svelte b/frontend/src/routes/(root)/(logged)/runs/[...path]/+page.svelte index 60460390dd..828be2f3bf 100644 --- a/frontend/src/routes/(root)/(logged)/runs/[...path]/+page.svelte +++ b/frontend/src/routes/(root)/(logged)/runs/[...path]/+page.svelte @@ -31,10 +31,11 @@ import { twMerge } from 'tailwind-merge' import ManuelDatePicker from '$lib/components/runs/ManuelDatePicker.svelte' import JobLoader from '$lib/components/runs/JobLoader.svelte' - import { AlertTriangle, Calendar, Clock } from 'lucide-svelte' + import { AlertTriangle, Calendar, Check, ChevronDown, Clock, X } from 'lucide-svelte' import ConcurrentJobsChart from '$lib/components/ConcurrentJobsChart.svelte' import ToggleButtonGroup from '$lib/components/common/toggleButton-v2/ToggleButtonGroup.svelte' import ToggleButton from '$lib/components/common/toggleButton-v2/ToggleButton.svelte' + import DropdownV2 from '$lib/components/DropdownV2.svelte' let jobs: Job[] | undefined let selectedIds: string[] = [] @@ -97,7 +98,9 @@ let selectedManualDate = 0 let autoRefresh: boolean = true let runDrawer: Drawer - let cancelAllJobs = false + let isCancelingVisibleJobs = false + let isCancelingFilteredJobs = false + let innerWidth = window.innerWidth let jobLoader: JobLoader | undefined = undefined let externalJobs: Job[] | undefined = undefined @@ -248,6 +251,8 @@ completedJobs = undefined selectedManualDate = 0 selectedIds = [] + jobIdsToCancel = [] + isSelectingJobsToCancel = false selectedWorkspace = undefined jobLoader?.loadJobs(minTs, maxTs, true) } @@ -327,6 +332,60 @@ } } + let jobIdsToCancel: string[] = [] + let isSelectingJobsToCancel = false + let fetchingFilteredJobs = false + let selectedFiltersString: string | undefined = undefined + + async function cancelVisibleJobs() { + isSelectingJobsToCancel = true + selectedIds = jobs?.filter(isJobCancelable).map((j) => j.id) ?? [] + if (selectedIds.length === 0 ) { + sendUserToast("There are no visible jobs that can be canceled", true) + } + } + async function cancelFilteredJobs() { + isCancelingFilteredJobs = true + fetchingFilteredJobs = true + const selectedFilters = { + workspace: $workspaceStore ?? '', + startedBefore: maxTs, + startedAfter: minTs, + schedulePath, + scriptPathExact: path === null || path === '' ? undefined : path, + createdBy: user === null || user === '' ? undefined : user, + scriptPathStart: folder === null || folder === '' ? undefined : `f/${folder}/`, + jobKinds, + success: success == 'success' ? true : success == 'failure' ? false : undefined, + running: success == 'running' ? true : undefined, + isNotSchedule: showSchedules == false ? true : undefined, + scheduledForBeforeNow: showFutureJobs == false ? true : undefined, + args: + argFilter && argFilter != '{}' && argFilter != '' && argError == '' + ? argFilter + : undefined, + result: + resultFilter && resultFilter != '{}' && resultFilter != '' && resultError == '' + ? resultFilter + : undefined, + allWorkspaces: allWorkspaces ? true : undefined, + concurrencyKey: concurrencyKey ?? undefined + } + + selectedFiltersString = JSON.stringify(selectedFilters, null, 4) + jobIdsToCancel = await JobService.listFilteredUuids(selectedFilters) + fetchingFilteredJobs = false + } + + async function cancelSelectedJobs() { + jobIdsToCancel = selectedIds + isCancelingVisibleJobs = true + } + + function isJobCancelable(j: Job): boolean { + return j.type === 'QueuedJob' && !j.schedule_path + } + const warnJobLimitMsg = 'The exact number of concurrent job at the beginning of the time range may be incorrect as only the last 1000 jobs are taken into account: a job that was started earlier than this limit will not be taken into account' @@ -334,6 +393,10 @@ graph === 'ConcurrencyChart' && extendedJobs !== undefined && extendedJobs.jobs.length + extendedJobs.obscured_jobs.length >= 1000 + + $: if (selectedIds.length === 0) { + isSelectingJobsToCancel = false + } { - cancelAllJobs = false - let uuids = await JobService.cancelAll({ workspace: $workspaceStore ?? '' }) + isCancelingFilteredJobs = false + let uuids = await JobService.cancelSelection({ + workspace: $workspaceStore ?? '', + requestBody: jobIdsToCancel + }) + jobIdsToCancel = [] + selectedIds = [] + jobLoader?.loadJobs(minTs, maxTs, true, true) + sendUserToast(`Canceled ${uuids.length} jobs`) + }} + loading={fetchingFilteredJobs} + on:canceled={() => { + isCancelingFilteredJobs = false + }} +> +
{selectedFiltersString}
+
+ + { + isCancelingVisibleJobs = false + let uuids = await JobService.cancelSelection({ + workspace: $workspaceStore ?? '', + requestBody: jobIdsToCancel + }) + jobIdsToCancel = [] + selectedIds = [] jobLoader?.loadJobs(minTs, maxTs, true, true) sendUserToast(`Canceled ${uuids.length} jobs`) }} on:canceled={() => { - cancelAllJobs = false + isCancelingVisibleJobs = false }} /> @@ -492,14 +583,65 @@
- +
+ {#if isSelectingJobsToCancel} +
+
+ {:else if !$userStore?.is_admin && !$superadmin} + + +
+ Cancel jobs +
+
+
+ {:else} + + +
+ Cancel jobs + +
+
+
+ {/if} +
{/if} - +
+ {#if isSelectingJobsToCancel} +
+
+ {:else if !$userStore?.is_admin && !$superadmin} + + +
+ Cancel jobs +
+
+
+ {:else} + + +
+ Cancel jobs + +
+
+
+ {/if} +
@@ -834,10 +1028,11 @@ externalJobs={externalJobs ?? []} omittedObscuredJobs={extendedJobs?.omitted_obscured_jobs ?? false} showExternalJobs={!graphIsRunsChart} + {isSelectingJobsToCancel} bind:selectedIds bind:selectedWorkspace on:select={() => { - runDrawer.openDrawer() + if (!isSelectingJobsToCancel) runDrawer.openDrawer() }} on:filterByPath={filterByPath} on:filterByUser={filterByUser}