From b95a8586b967c26cde60359f2cbdb7dd658783c4 Mon Sep 17 00:00:00 2001 From: hugocasa Date: Thu, 27 Nov 2025 19:34:02 +0100 Subject: [PATCH] suspended mode draft --- ...200000_add_queue_mode_to_triggers.down.sql | 27 -- ...06200000_add_queue_mode_to_triggers.up.sql | 28 -- ...0251127103823_add_unassigned_kind.down.sql | 1 + .../20251127103823_add_unassigned_kind.up.sql | 2 + ...00_add_suspended_mode_to_triggers.down.sql | 27 ++ ...4100_add_suspended_mode_to_triggers.up.sql | 28 ++ backend/tests/common/mod.rs | 3 +- backend/tests/relative_imports.rs | 7 + backend/windmill-api/openapi.yaml | 58 +-- backend/windmill-api/src/apps.rs | 9 +- backend/windmill-api/src/flows.rs | 6 +- backend/windmill-api/src/jobs.rs | 15 + backend/windmill-api/src/scripts.rs | 11 +- .../src/triggers/global_handler.rs | 365 +++++++++++------- backend/windmill-api/src/triggers/handler.rs | 86 ++++- .../windmill-api/src/triggers/http/handler.rs | 18 +- backend/windmill-api/src/triggers/http/mod.rs | 4 +- backend/windmill-api/src/triggers/listener.rs | 12 +- backend/windmill-api/src/triggers/mod.rs | 12 +- .../windmill-api/src/triggers/mqtt/handler.rs | 8 +- .../src/triggers/postgres/handler.rs | 8 +- .../src/triggers/trigger_helpers.rs | 35 +- .../src/triggers/websocket/handler.rs | 8 +- .../src/triggers/websocket/listener.rs | 8 +- backend/windmill-common/src/jobs.rs | 1 + backend/windmill-common/src/triggers.rs | 1 + backend/windmill-queue/src/jobs.rs | 20 +- backend/windmill-queue/src/schedule.rs | 1 + backend/windmill-worker/src/ai/tools.rs | 3 +- backend/windmill-worker/src/worker.rs | 3 +- backend/windmill-worker/src/worker_flow.rs | 3 +- .../windmill-worker/src/worker_lockfiles.rs | 3 +- .../triggers/TriggerActiveMode.svelte | 157 +------- .../triggers/TriggerSuspendedJobsModal.svelte | 289 ++++++++++++++ .../email/EmailTriggerEditorInner.svelte | 8 +- .../lib/components/triggers/email/utils.ts | 2 +- .../triggers/gcp/GcpTriggerEditorInner.svelte | 8 +- .../src/lib/components/triggers/gcp/utils.ts | 2 +- .../triggers/http/RouteEditorInner.svelte | 9 +- .../src/lib/components/triggers/http/utils.ts | 2 +- .../kafka/KafkaTriggerEditorInner.svelte | 8 +- .../lib/components/triggers/kafka/utils.ts | 4 +- .../mqtt/MqttTriggerEditorInner.svelte | 8 +- .../src/lib/components/triggers/mqtt/utils.ts | 4 +- .../nats/NatsTriggerEditorInner.svelte | 8 +- .../src/lib/components/triggers/nats/utils.ts | 4 +- .../PostgresTriggerEditorInner.svelte | 8 +- .../lib/components/triggers/postgres/utils.ts | 4 +- .../triggers/sqs/SqsTriggerEditorInner.svelte | 8 +- .../src/lib/components/triggers/sqs/utils.ts | 4 +- .../WebsocketTriggerEditorInner.svelte | 8 +- .../components/triggers/websocket/utils.ts | 4 +- 52 files changed, 842 insertions(+), 528 deletions(-) delete mode 100644 backend/migrations/20251106200000_add_queue_mode_to_triggers.down.sql delete mode 100644 backend/migrations/20251106200000_add_queue_mode_to_triggers.up.sql create mode 100644 backend/migrations/20251127103823_add_unassigned_kind.down.sql create mode 100644 backend/migrations/20251127103823_add_unassigned_kind.up.sql create mode 100644 backend/migrations/20251127104100_add_suspended_mode_to_triggers.down.sql create mode 100644 backend/migrations/20251127104100_add_suspended_mode_to_triggers.up.sql create mode 100644 frontend/src/lib/components/triggers/TriggerSuspendedJobsModal.svelte diff --git a/backend/migrations/20251106200000_add_queue_mode_to_triggers.down.sql b/backend/migrations/20251106200000_add_queue_mode_to_triggers.down.sql deleted file mode 100644 index 0a587265c2..0000000000 --- a/backend/migrations/20251106200000_add_queue_mode_to_triggers.down.sql +++ /dev/null @@ -1,27 +0,0 @@ --- Add down migration script here -ALTER TABLE websocket_trigger -DROP COLUMN IF EXISTS active_mode; - -ALTER TABLE sqs_trigger -DROP COLUMN IF EXISTS active_mode; - -ALTER TABLE postgres_trigger -DROP COLUMN IF EXISTS active_mode; - -ALTER TABLE nats_trigger -DROP COLUMN IF EXISTS active_mode; - -ALTER TABLE mqtt_trigger -DROP COLUMN IF EXISTS active_mode; - -ALTER TABLE kafka_trigger -DROP COLUMN IF EXISTS active_mode; - -ALTER TABLE http_trigger -DROP COLUMN IF EXISTS active_mode; - -ALTER TABLE gcp_trigger -DROP COLUMN IF EXISTS active_mode; - -ALTER TABLE email_trigger -DROP COLUMN IF EXISTS active_mode; \ No newline at end of file diff --git a/backend/migrations/20251106200000_add_queue_mode_to_triggers.up.sql b/backend/migrations/20251106200000_add_queue_mode_to_triggers.up.sql deleted file mode 100644 index e51716a178..0000000000 --- a/backend/migrations/20251106200000_add_queue_mode_to_triggers.up.sql +++ /dev/null @@ -1,28 +0,0 @@ --- Add up migration script here -ALTER TABLE gcp_trigger -ADD COLUMN active_mode BOOLEAN NOT NULL DEFAULT TRUE; - -ALTER TABLE http_trigger -ADD COLUMN active_mode BOOLEAN NOT NULL DEFAULT TRUE; - -ALTER TABLE kafka_trigger -ADD COLUMN active_mode BOOLEAN NOT NULL DEFAULT TRUE; - -ALTER TABLE mqtt_trigger -ADD COLUMN active_mode BOOLEAN NOT NULL DEFAULT TRUE; - -ALTER TABLE nats_trigger -ADD COLUMN active_mode BOOLEAN NOT NULL DEFAULT TRUE; - -ALTER TABLE postgres_trigger -ADD COLUMN active_mode BOOLEAN NOT NULL DEFAULT TRUE; - -ALTER TABLE sqs_trigger -ADD COLUMN active_mode BOOLEAN NOT NULL DEFAULT TRUE; - -ALTER TABLE websocket_trigger -ADD COLUMN active_mode BOOLEAN NOT NULL DEFAULT TRUE; - -ALTER TABLE email_trigger -ADD COLUMN active_mode BOOLEAN NOT NULL DEFAULT TRUE; - diff --git a/backend/migrations/20251127103823_add_unassigned_kind.down.sql b/backend/migrations/20251127103823_add_unassigned_kind.down.sql new file mode 100644 index 0000000000..d2f607c5b8 --- /dev/null +++ b/backend/migrations/20251127103823_add_unassigned_kind.down.sql @@ -0,0 +1 @@ +-- Add down migration script here diff --git a/backend/migrations/20251127103823_add_unassigned_kind.up.sql b/backend/migrations/20251127103823_add_unassigned_kind.up.sql new file mode 100644 index 0000000000..95f750731c --- /dev/null +++ b/backend/migrations/20251127103823_add_unassigned_kind.up.sql @@ -0,0 +1,2 @@ +-- Add up migration script here +ALTER TYPE JOB_KIND ADD VALUE IF NOT EXISTS 'unassigned'; diff --git a/backend/migrations/20251127104100_add_suspended_mode_to_triggers.down.sql b/backend/migrations/20251127104100_add_suspended_mode_to_triggers.down.sql new file mode 100644 index 0000000000..fb32a57a9a --- /dev/null +++ b/backend/migrations/20251127104100_add_suspended_mode_to_triggers.down.sql @@ -0,0 +1,27 @@ +-- Add down migration script here +ALTER TABLE websocket_trigger +DROP COLUMN IF EXISTS suspended_mode; + +ALTER TABLE sqs_trigger +DROP COLUMN IF EXISTS suspended_mode; + +ALTER TABLE postgres_trigger +DROP COLUMN IF EXISTS suspended_mode; + +ALTER TABLE nats_trigger +DROP COLUMN IF EXISTS suspended_mode; + +ALTER TABLE mqtt_trigger +DROP COLUMN IF EXISTS suspended_mode; + +ALTER TABLE kafka_trigger +DROP COLUMN IF EXISTS suspended_mode; + +ALTER TABLE http_trigger +DROP COLUMN IF EXISTS suspended_mode; + +ALTER TABLE gcp_trigger +DROP COLUMN IF EXISTS suspended_mode; + +ALTER TABLE email_trigger +DROP COLUMN IF EXISTS suspended_mode; \ No newline at end of file diff --git a/backend/migrations/20251127104100_add_suspended_mode_to_triggers.up.sql b/backend/migrations/20251127104100_add_suspended_mode_to_triggers.up.sql new file mode 100644 index 0000000000..994eb83804 --- /dev/null +++ b/backend/migrations/20251127104100_add_suspended_mode_to_triggers.up.sql @@ -0,0 +1,28 @@ +-- Add up migration script here +ALTER TABLE gcp_trigger +ADD COLUMN suspended_mode BOOLEAN NOT NULL DEFAULT FALSE; + +ALTER TABLE http_trigger +ADD COLUMN suspended_mode BOOLEAN NOT NULL DEFAULT FALSE; + +ALTER TABLE kafka_trigger +ADD COLUMN suspended_mode BOOLEAN NOT NULL DEFAULT FALSE; + +ALTER TABLE mqtt_trigger +ADD COLUMN suspended_mode BOOLEAN NOT NULL DEFAULT FALSE; + +ALTER TABLE nats_trigger +ADD COLUMN suspended_mode BOOLEAN NOT NULL DEFAULT FALSE; + +ALTER TABLE postgres_trigger +ADD COLUMN suspended_mode BOOLEAN NOT NULL DEFAULT FALSE; + +ALTER TABLE sqs_trigger +ADD COLUMN suspended_mode BOOLEAN NOT NULL DEFAULT FALSE; + +ALTER TABLE websocket_trigger +ADD COLUMN suspended_mode BOOLEAN NOT NULL DEFAULT FALSE; + +ALTER TABLE email_trigger +ADD COLUMN suspended_mode BOOLEAN NOT NULL DEFAULT FALSE; + diff --git a/backend/tests/common/mod.rs b/backend/tests/common/mod.rs index cfdfe4124f..5555861822 100644 --- a/backend/tests/common/mod.rs +++ b/backend/tests/common/mod.rs @@ -208,7 +208,8 @@ impl RunJob { false, None, debounce_job_id_o, - None + None, + None, ) .await .expect("push has to succeed"); diff --git a/backend/tests/relative_imports.rs b/backend/tests/relative_imports.rs index 6279bc107f..18c18a18c7 100644 --- a/backend/tests/relative_imports.rs +++ b/backend/tests/relative_imports.rs @@ -865,6 +865,7 @@ def main(): None, None, None, + None, ) .await .unwrap(); @@ -1025,6 +1026,7 @@ def main(): None, debounce_job_id_o, None, + None, ) .await .unwrap(); @@ -1204,6 +1206,7 @@ def main(): debounce_job_id_o, None, None, + None, ) .await .unwrap(); @@ -1713,6 +1716,7 @@ WHERE None, None, None, + None, ) .await .unwrap(); @@ -1853,6 +1857,7 @@ WHERE None, None, None, + None, ) .await .unwrap(); @@ -2284,6 +2289,7 @@ WHERE None, None, None, + None, ) .await .unwrap(); @@ -2410,6 +2416,7 @@ WHERE // None, // None, // None, + // None, // ) // .await // .unwrap(); diff --git a/backend/windmill-api/openapi.yaml b/backend/windmill-api/openapi.yaml index 867b0a790d..228883b2a1 100644 --- a/backend/windmill-api/openapi.yaml +++ b/backend/windmill-api/openapi.yaml @@ -17386,7 +17386,7 @@ components: type: boolean enabled: type: boolean - active_mode: + suspended_mode: type: boolean required: - path @@ -17398,7 +17398,7 @@ components: - edited_at - is_flow - enabled - - active_mode + - suspended_mode AuthenticationMethod: type: string @@ -17634,10 +17634,8 @@ components: $ref: "#/components/schemas/ScriptArgs" retry: $ref: "../../openflow.openapi.yaml#/components/schemas/Retry" - active_mode: + suspended_mode: type: boolean - default: true - description: If set to false, each incoming event will be suspend job until ran manually or set it to true required: - path @@ -17699,9 +17697,8 @@ components: $ref: "#/components/schemas/ScriptArgs" retry: $ref: "../../openflow.openapi.yaml#/components/schemas/Retry" - active_mode: + suspended_mode: type: boolean - default: true required: - path - script_path @@ -17832,9 +17829,8 @@ components: $ref: "#/components/schemas/ScriptArgs" retry: $ref: "../../openflow.openapi.yaml#/components/schemas/Retry" - active_mode: + suspended_mode: type: boolean - default: true required: - path @@ -17883,9 +17879,8 @@ components: $ref: "#/components/schemas/ScriptArgs" retry: $ref: "../../openflow.openapi.yaml#/components/schemas/Retry" - active_mode: + suspended_mode: type: boolean - default: true required: - path @@ -18023,9 +18018,8 @@ components: $ref: "#/components/schemas/ScriptArgs" retry: $ref: "../../openflow.openapi.yaml#/components/schemas/Retry" - active_mode: + suspended_mode: type: boolean - default: true required: - path - script_path @@ -18064,9 +18058,8 @@ components: $ref: "#/components/schemas/ScriptArgs" retry: $ref: "../../openflow.openapi.yaml#/components/schemas/Retry" - active_mode: + suspended_mode: type: boolean - default: true required: - path - script_path @@ -18175,7 +18168,7 @@ components: $ref: "#/components/schemas/ScriptArgs" retry: $ref: "../../openflow.openapi.yaml#/components/schemas/Retry" - active_mode: + suspended_mode: type: boolean required: - path @@ -18310,10 +18303,8 @@ components: $ref: "#/components/schemas/ScriptArgs" retry: $ref: "../../openflow.openapi.yaml#/components/schemas/Retry" - active_mode: + suspended_mode: type: boolean - default: true - description: If false, queue jobs with suspend functionality instead of immediate execution required: - queue_url - aws_resource_path @@ -18349,9 +18340,8 @@ components: $ref: "#/components/schemas/ScriptArgs" retry: $ref: "../../openflow.openapi.yaml#/components/schemas/Retry" - active_mode: + suspended_mode: type: boolean - default: true required: - queue_url - aws_resource_path @@ -18491,9 +18481,8 @@ components: $ref: "#/components/schemas/ScriptArgs" retry: $ref: "../../openflow.openapi.yaml#/components/schemas/Retry" - active_mode: + suspended_mode: type: boolean - default: true required: - path - script_path @@ -18526,9 +18515,8 @@ components: $ref: "#/components/schemas/ScriptArgs" retry: $ref: "../../openflow.openapi.yaml#/components/schemas/Retry" - active_mode: + suspended_modes: type: boolean - default: true required: - path - script_path @@ -18595,9 +18583,8 @@ components: $ref: "#/components/schemas/ScriptArgs" retry: $ref: "../../openflow.openapi.yaml#/components/schemas/Retry" - active_mode: + suspended_mode: type: boolean - default: true required: - path @@ -18628,9 +18615,8 @@ components: type: string error_handler_args: $ref: "#/components/schemas/ScriptArgs" - active_mode: + suspended_mode: type: boolean - default: true retry: $ref: "../../openflow.openapi.yaml#/components/schemas/Retry" @@ -18707,9 +18693,8 @@ components: $ref: "#/components/schemas/ScriptArgs" retry: $ref: "../../openflow.openapi.yaml#/components/schemas/Retry" - active_mode: + suspended_mode: type: boolean - default: true required: - path @@ -18746,9 +18731,8 @@ components: $ref: "#/components/schemas/ScriptArgs" retry: $ref: "../../openflow.openapi.yaml#/components/schemas/Retry" - active_mode: + suspended_mode: type: boolean - default: true required: - path - script_path @@ -18797,10 +18781,8 @@ components: $ref: "../../openflow.openapi.yaml#/components/schemas/Retry" enabled: type: boolean - active_mode: + suspended_mode: type: boolean - default: true - description: If set to false, each incoming event will be suspend job until ran manually or set it to true required: - path @@ -18827,10 +18809,8 @@ components: $ref: "#/components/schemas/ScriptArgs" retry: $ref: "../../openflow.openapi.yaml#/components/schemas/Retry" - active_mode: + suspended_mode: type: boolean - default: true - description: If set to false, each incoming event will be suspend job until ran manually or set it to true required: - path - script_path diff --git a/backend/windmill-api/src/apps.rs b/backend/windmill-api/src/apps.rs index 60e73b3990..b0781930d5 100644 --- a/backend/windmill-api/src/apps.rs +++ b/backend/windmill-api/src/apps.rs @@ -1242,7 +1242,8 @@ async fn create_app_internal<'a>( false, None, None, - None + None, + None, ) .await?; tracing::info!("Pushed app dependency job {}", dependency_job_uuid); @@ -1632,7 +1633,8 @@ async fn update_app_internal<'a>( false, None, None, - None + None, + None, ) .await?; tracing::info!("Pushed app dependency job {}", dependency_job_uuid); @@ -1960,7 +1962,8 @@ async fn execute_component( false, end_user_email, None, - None + None, + None, ) .await?; diff --git a/backend/windmill-api/src/flows.rs b/backend/windmill-api/src/flows.rs index 48eb2e58c0..9b8c8b3930 100644 --- a/backend/windmill-api/src/flows.rs +++ b/backend/windmill-api/src/flows.rs @@ -566,7 +566,8 @@ async fn create_flow( false, None, None, - None + None, + None, ) .await?; @@ -1032,7 +1033,8 @@ async fn update_flow( false, None, None, - None + None, + None, ) .await?; diff --git a/backend/windmill-api/src/jobs.rs b/backend/windmill-api/src/jobs.rs index e1a885e9c1..c8c4b3c555 100644 --- a/backend/windmill-api/src/jobs.rs +++ b/backend/windmill-api/src/jobs.rs @@ -1757,6 +1757,7 @@ pub struct RunJobQuery { pub skip_preprocessor: Option, pub poll_delay_ms: Option, pub memory_id: Option, + pub suspended_mode: Option, } impl RunJobQuery { @@ -4143,6 +4144,7 @@ pub async fn run_flow( None, None, trigger, + run_query.suspended_mode, ) .await?; @@ -4398,6 +4400,7 @@ pub async fn restart_flow( None, None, None, + run_query.suspended_mode, ) .await?; tx.commit().await?; @@ -4518,6 +4521,7 @@ pub async fn push_script_job_by_path_into_queue( None, None, trigger, + run_query.suspended_mode, ) .await?; tx.commit().await?; @@ -4676,6 +4680,7 @@ pub async fn run_workflow_as_code( None, None, None, + None, ) .await?; @@ -5214,6 +5219,7 @@ pub async fn run_wait_result_job_by_path_get( None, None, None, + run_query.suspended_mode, ) .await?; tx.commit().await?; @@ -5359,6 +5365,7 @@ pub async fn run_wait_result_script_by_path_internal( None, None, None, + run_query.suspended_mode, ) .await?; tx.commit().await?; @@ -5482,6 +5489,7 @@ pub async fn run_wait_result_script_by_hash( None, None, None, + run_query.suspended_mode, ) .await?; tx.commit().await?; @@ -5956,6 +5964,7 @@ async fn run_preview_script( None, None, None, + None, ) .await?; tx.commit().await?; @@ -6077,6 +6086,7 @@ async fn run_bundle_preview_script( None, None, None, + None, ) .await?; job_id = Some(uuid); @@ -6217,6 +6227,7 @@ async fn run_dependencies_job( None, None, None, + None, ) .await?; tx.commit().await?; @@ -6287,6 +6298,7 @@ async fn run_flow_dependencies_job( None, None, None, + None, ) .await?; tx.commit().await?; @@ -6641,6 +6653,7 @@ async fn run_preview_flow_job( None, None, None, + None, ) .await?; @@ -6839,6 +6852,7 @@ async fn run_dynamic_select( None, None, None, + None, ) .await?; tx.commit().await?; @@ -6983,6 +6997,7 @@ pub async fn run_job_by_hash_inner( None, None, trigger, + run_query.suspended_mode, ) .await?; tx.commit().await?; diff --git a/backend/windmill-api/src/scripts.rs b/backend/windmill-api/src/scripts.rs index 98fad0cc37..ec92ed66bd 100644 --- a/backend/windmill-api/src/scripts.rs +++ b/backend/windmill-api/src/scripts.rs @@ -41,11 +41,11 @@ use windmill_audit::ActionKind; use windmill_worker::{process_relative_imports, scoped_dependency_map::ScopedDependencyMap}; use windmill_common::{ - assets::{AssetUsageKind, AssetWithAltAccessType, clear_asset_usage, insert_asset_usage}, + assets::{clear_asset_usage, insert_asset_usage, AssetUsageKind, AssetWithAltAccessType}, error::to_anyhow, s3_helpers::upload_artifact_to_store, scripts::hash_script, - utils::{WarnAfterExt, paginate_without_limits}, + utils::{paginate_without_limits, WarnAfterExt}, worker::{CLOUD_HOSTED, MIN_VERSION_SUPPORTS_DEBOUNCING}, }; @@ -60,9 +60,7 @@ use windmill_common::{ ScriptHistory, ScriptHistoryUpdate, ScriptKind, ScriptLang, ScriptWithStarred, }, users::username_to_permissioned_as, - utils::{ - not_found_if_none, query_elems_from_hub, require_admin, Pagination, StripPath, - }, + utils::{not_found_if_none, query_elems_from_hub, require_admin, Pagination, StripPath}, worker::to_raw_value, HUB_BASE_URL, }; @@ -1027,7 +1025,8 @@ async fn create_script_internal<'c>( false, None, None, - None + None, + None, ) .await?; Ok((hash, new_tx, None)) diff --git a/backend/windmill-api/src/triggers/global_handler.rs b/backend/windmill-api/src/triggers/global_handler.rs index ae76b1d740..831982c4fa 100644 --- a/backend/windmill-api/src/triggers/global_handler.rs +++ b/backend/windmill-api/src/triggers/global_handler.rs @@ -1,149 +1,254 @@ -use crate::{db::ApiAuthed, triggers::INACTIVE_TRIGGER_SCHEDULED_FOR_DATE}; +use crate::{ + db::{ApiAuthed, DB}, + triggers::{trigger_helpers::trigger_runnable_inner, TriggerForReassignment, Trigger, TriggerCrud}, +}; + +#[cfg(feature = "http_trigger")] +use crate::triggers::http::{handler::HttpTrigger, HttpConfig}; + +#[cfg(feature = "mqtt_trigger")] +use crate::triggers::mqtt::{MqttConfig, MqttTrigger}; + +#[cfg(feature = "postgres_trigger")] +use crate::triggers::postgres::{PostgresConfig, PostgresTrigger}; + +#[cfg(feature = "websocket")] +use crate::triggers::websocket::{WebsocketConfig, WebsocketTrigger}; + +#[cfg(all(feature = "smtp", feature = "enterprise", feature = "private"))] +use crate::triggers::email::{EmailConfig, EmailTrigger}; + +#[cfg(all(feature = "gcp_trigger", feature = "enterprise", feature = "private"))] +use crate::triggers::gcp::{GcpConfig, GcpTrigger}; + +#[cfg(all(feature = "kafka", feature = "enterprise", feature = "private"))] +use crate::triggers::kafka::{KafkaConfig, KafkaTrigger}; + +#[cfg(all(feature = "nats", feature = "enterprise", feature = "private"))] +use crate::triggers::nats::{NatsConfig, NatsTrigger}; + +#[cfg(all(feature = "sqs_trigger", feature = "enterprise", feature = "private"))] +use crate::triggers::sqs::{SqsConfig, SqsTrigger}; use axum::{ extract::{Extension, Path}, response::Json, }; -use windmill_audit::{audit_oss::audit_log, ActionKind}; -use windmill_common::{db::UserDB, error, jobs::JobTriggerKind, utils::require_admin}; +use serde_json::value::RawValue; +use std::collections::HashMap; +use uuid::Uuid; +use windmill_common::{db::UserDB, error, jobs::JobTriggerKind, triggers::TriggerMetadata}; -pub async fn resume_suspended_trigger_jobs( +struct JobWithArgs { + id: Uuid, + args: Option>>>, +} + +pub async fn reassign_suspended_jobs( authed: ApiAuthed, + Extension(db): Extension, Extension(user_db): Extension, Path((w_id, trigger_kind, trigger_path)): Path<(String, JobTriggerKind, String)>, ) -> error::Result> { - require_admin(authed.is_admin, &authed.username)?; + let mut tx = user_db.clone().begin(&authed).await?; - let mut tx = user_db.begin(&authed).await?; + let trigger: TriggerForReassignment = match trigger_kind { + JobTriggerKind::Websocket => { + #[cfg(feature = "websocket")] + { + WebsocketTrigger + .get_trigger_for_reassignment(&mut *tx, &w_id, &trigger_path) + .await? + } + #[cfg(not(feature = "websocket"))] + { + return Err(error::Error::BadRequest( + "Websocket triggers are not enabled in this build".to_string(), + )); + } + } + JobTriggerKind::Http => { + #[cfg(feature = "http_trigger")] + { + HttpTrigger + .get_trigger_for_reassignment(&mut *tx, &w_id, &trigger_path) + .await? + } + #[cfg(not(feature = "http_trigger"))] + { + return Err(error::Error::BadRequest( + "HTTP triggers are not enabled in this build".to_string(), + )); + } + } + JobTriggerKind::Mqtt => { + #[cfg(feature = "mqtt_trigger")] + { + MqttTrigger + .get_trigger_for_reassignment(&mut *tx, &w_id, &trigger_path) + .await? + } + #[cfg(not(feature = "mqtt_trigger"))] + { + return Err(error::Error::BadRequest( + "MQTT triggers are not enabled in this build".to_string(), + )); + } + } + JobTriggerKind::Postgres => { + #[cfg(feature = "postgres_trigger")] + { + PostgresTrigger + .get_trigger_for_reassignment(&mut *tx, &w_id, &trigger_path) + .await? + } + #[cfg(not(feature = "postgres_trigger"))] + { + return Err(error::Error::BadRequest( + "Postgres triggers are not enabled in this build".to_string(), + )); + } + } + JobTriggerKind::Kafka => { + #[cfg(all(feature = "kafka", feature = "enterprise", feature = "private"))] + { + KafkaTrigger + .get_trigger_for_reassignment(&mut *tx, &w_id, &trigger_path) + .await? + } + #[cfg(not(all(feature = "kafka", feature = "enterprise", feature = "private")))] + { + return Err(error::Error::BadRequest( + "Kafka triggers are not enabled in this build".to_string(), + )); + } + } + JobTriggerKind::Email => { + #[cfg(all(feature = "smtp", feature = "enterprise", feature = "private"))] + { + EmailTrigger + .get_trigger_for_reassignment(&mut *tx, &w_id, &trigger_path) + .await? + } + #[cfg(not(all(feature = "smtp", feature = "enterprise", feature = "private")))] + { + return Err(error::Error::BadRequest( + "Email triggers are not enabled in this build".to_string(), + )); + } + } + JobTriggerKind::Nats => { + #[cfg(all(feature = "nats", feature = "enterprise", feature = "private"))] + { + NatsTrigger + .get_trigger_for_reassignment(&mut *tx, &w_id, &trigger_path) + .await? + } + #[cfg(not(all(feature = "nats", feature = "enterprise", feature = "private")))] + { + return Err(error::Error::BadRequest( + "NATS triggers are not enabled in this build".to_string(), + )); + } + } + JobTriggerKind::Sqs => { + #[cfg(all(feature = "sqs_trigger", feature = "enterprise", feature = "private"))] + { + SqsTrigger + .get_trigger_for_reassignment(&mut *tx, &w_id, &trigger_path) + .await? + } + #[cfg(not(all(feature = "sqs_trigger", feature = "enterprise", feature = "private")))] + { + return Err(error::Error::BadRequest( + "SQS triggers are not enabled in this build".to_string(), + )); + } + } + JobTriggerKind::Gcp => { + #[cfg(all(feature = "gcp_trigger", feature = "enterprise", feature = "private"))] + { + GcpTrigger + .get_trigger_for_reassignment(&mut *tx, &w_id, &trigger_path) + .await? + } + #[cfg(not(all(feature = "gcp_trigger", feature = "enterprise", feature = "private")))] + { + return Err(error::Error::BadRequest( + "GCP triggers are not enabled in this build".to_string(), + )); + } + } + JobTriggerKind::Webhook | JobTriggerKind::Schedule => { + return Err(error::Error::BadRequest( + "Webhook and Schedule triggers do not support job reassignment".to_string(), + )); + } + }; - // Use the date constant to identify suspended jobs - // This date (9999-12-31 23:59:59) is used as a marker for suspended jobs - let scheduled_for = INACTIVE_TRIGGER_SCHEDULED_FOR_DATE.clone(); - let result = sqlx::query!( - r#" - UPDATE - v2_job_queue - SET - scheduled_for = now() - FROM - v2_job - WHERE - v2_job_queue.id = v2_job.id AND - v2_job_queue.running is FALSE AND - v2_job_queue.scheduled_for = $1 AND - v2_job_queue.workspace_id = $2 AND - v2_job.trigger_kind = $3 AND - v2_job.trigger = $4 - "#, - scheduled_for, + let jobs = sqlx::query_as!(JobWithArgs, + "SELECT id, args as \"args: _\" FROM v2_job WHERE workspace_id = $1 AND kind = 'unassigned'::JOB_KIND AND trigger_kind = $2 AND trigger = $3", w_id, trigger_kind as _, - trigger_path - ) - .execute(&mut *tx) - .await?; + trigger_path, + ).fetch_all(&mut *tx).await?; - let count = result.rows_affected(); + let trigger_metadata = TriggerMetadata::new(Some(trigger_path.clone()), trigger_kind); - let trigger_kind_str = format!("{:?}", trigger_kind); - let count_str = count.to_string(); - audit_log( - &mut *tx, - &authed, - "triggers.bulk_resume", - ActionKind::Execute, - &w_id, - Some(&trigger_path), - Some( - [ - ("trigger_kind", trigger_kind_str.as_str()), - ("jobs_count", count_str.as_str()), - ] - .iter() - .cloned() - .collect(), - ), - ) - .await?; + let l = jobs.len(); + + for job in jobs { + trigger_runnable_inner( + &db, + Some(user_db.clone()), + authed.clone(), + &w_id, + &trigger.script_path, + trigger.is_flow, + windmill_queue::PushArgsOwned { + extra: None, + args: job.args.map(|a| a.0).unwrap_or_default(), + }, + trigger.retry.as_ref(), + trigger.error_handler_path.as_deref(), + trigger.error_handler_args.as_ref(), + trigger_path.clone(), + None, + trigger_metadata.clone(), + None, + ) + .await?; + + // Delete the unassigned job from all related tables + sqlx::query!("DELETE FROM v2_job_queue WHERE id = $1", job.id) + .execute(&mut *tx) + .await?; + + sqlx::query!("DELETE FROM v2_job_runtime WHERE id = $1", job.id) + .execute(&mut *tx) + .await?; + + sqlx::query!("DELETE FROM job_perms WHERE job_id = $1", job.id) + .execute(&mut *tx) + .await?; + + sqlx::query!("DELETE FROM concurrency_key WHERE job_id = $1", job.id) + .execute(&mut *tx) + .await?; + + sqlx::query!("DELETE FROM debounce_key WHERE job_id = $1", job.id) + .execute(&mut *tx) + .await?; + + sqlx::query!("DELETE FROM debounce_stale_data WHERE job_id = $1", job.id) + .execute(&mut *tx) + .await?; + + sqlx::query!("DELETE FROM v2_job WHERE id = $1", job.id) + .execute(&mut *tx) + .await?; + } tx.commit().await?; - let message = format!( - "Successfully resumed {} suspended job{} for trigger at path: {}", - count, - if count == 1 { "" } else { "s" }, - &trigger_path - ); - Ok(Json(message)) -} - -pub async fn cancel_suspended_trigger_jobs( - authed: ApiAuthed, - Extension(user_db): Extension, - Path((w_id, trigger_kind, trigger_path)): Path<(String, JobTriggerKind, String)>, -) -> error::Result> { - require_admin(authed.is_admin, &authed.username)?; - - let mut tx = user_db.begin(&authed).await?; - - let scheduled_for = INACTIVE_TRIGGER_SCHEDULED_FOR_DATE.clone(); - - let result = sqlx::query!( - r#" - UPDATE - v2_job_queue - SET - canceled_by = $1, - canceled_reason = 'cancelled by trigger bulk operation', - scheduled_for = now() - FROM - v2_job - WHERE - v2_job_queue.id = v2_job.id AND - v2_job_queue.workspace_id = $2 AND - v2_job_queue.running is FALSE AND - v2_job_queue.scheduled_for = $3 AND - v2_job.trigger_kind = $4 AND - v2_job.trigger = $5 - "#, - authed.username, - w_id, - scheduled_for, - trigger_kind as _, - trigger_path - ) - .execute(&mut *tx) - .await?; - - let count = result.rows_affected(); - - let trigger_kind_str = format!("{:?}", trigger_kind); - let count_str = count.to_string(); - audit_log( - &mut *tx, - &authed, - "triggers.bulk_cancel", - ActionKind::Delete, - &w_id, - Some(&trigger_path), - Some( - [ - ("trigger_kind", trigger_kind_str.as_str()), - ("jobs_count", count_str.as_str()), - ] - .iter() - .cloned() - .collect(), - ), - ) - .await?; - - tx.commit().await?; - - let message = format!( - "Successfully cancelled {} suspended job{} for trigger at path: {}", - count, - if count == 1 { "" } else { "s" }, - &trigger_path - ); - Ok(Json(message)) + Ok(Json(format!("Reassigned {} jobs", l))) } diff --git a/backend/windmill-api/src/triggers/handler.rs b/backend/windmill-api/src/triggers/handler.rs index c46b9202c3..76a2de6755 100644 --- a/backend/windmill-api/src/triggers/handler.rs +++ b/backend/windmill-api/src/triggers/handler.rs @@ -3,10 +3,11 @@ use crate::{ triggers::{StandardTriggerQuery, TriggerData}, }; use async_trait::async_trait; +use chrono::{DateTime, Utc}; use serde::{de::DeserializeOwned, Deserialize, Serialize}; use sql_builder::{bind::Bind, SqlBuilder}; use sqlx::{FromRow, PgConnection}; -use std::fmt::Debug; +use std::{collections::HashMap, fmt::Debug}; use windmill_common::{ db::UserDB, error::{Error, JsonResult, Result}, @@ -28,6 +29,19 @@ use windmill_git_sync::handle_deployment_metadata; use crate::utils::check_scopes; +#[derive(FromRow)] +pub struct TriggerForReassignment { + pub script_path: String, + pub is_flow: bool, + pub edited_by: String, + pub email: String, + pub edited_at: DateTime, + pub error_handler_path: Option, + pub error_handler_args: + Option>>, + pub retry: Option>, +} + #[async_trait] pub trait TriggerCrud: Send + Sync + 'static { type Trigger: Serialize @@ -144,7 +158,7 @@ pub trait TriggerCrud: Send + Sync + 'static { "edited_at", "extra_perms", "enabled", - "active_mode" + "suspended_mode", ]; if Self::SUPPORTS_SERVER_STATE { @@ -175,6 +189,44 @@ pub trait TriggerCrud: Send + Sync + 'static { .ok_or_else(|| Error::NotFound(format!("Trigger not found at path: {}", path))) } + async fn get_trigger_for_reassignment( + &self, + tx: &mut PgConnection, + workspace_id: &str, + path: &str, + ) -> Result { + let fields = vec![ + "script_path", + "is_flow", + "edited_by", + "email", + "edited_at", + "error_handler_path", + "error_handler_args", + "retry", + ]; + + let sql = format!( + r#"SELECT + {} + FROM + {} + WHERE + workspace_id = $1 AND + path = $2 + "#, + fields.join(", "), + Self::TABLE_NAME + ); + + sqlx::query_as(&sql) + .bind(workspace_id) + .bind(path) + .fetch_optional(&mut *tx) + .await? + .ok_or_else(|| Error::NotFound(format!("Trigger not found at path: {}", path))) + } + async fn exists(&self, db: &DB, workspace_id: &str, path: &str) -> Result { let exists = sqlx::query_scalar(&format!( "SELECT EXISTS(SELECT 1 FROM {} WHERE workspace_id = $1 AND path = $2)", @@ -323,7 +375,7 @@ pub trait TriggerCrud: Send + Sync + 'static { "edited_at", "extra_perms", "enabled", - "active_mode", + "suspended_mode", ]; if Self::SUPPORTS_SERVER_STATE { @@ -773,21 +825,21 @@ pub fn generate_trigger_routers() -> Router { ); } - { - use crate::triggers::global_handler::{ - cancel_suspended_trigger_jobs, resume_suspended_trigger_jobs, - }; + // { + // use crate::triggers::global_handler::{ + // cancel_suspended_trigger_jobs, resume_suspended_trigger_jobs, + // }; - router = router - .route( - "/trigger/:trigger_kind/resume_suspended_trigger_job/*trigger_path", - post(resume_suspended_trigger_jobs), - ) - .route( - "/trigger/:trigger_kind/cancel_suspended_trigger_job/*trigger_path", - post(cancel_suspended_trigger_jobs), - ); - } + // router = router + // .route( + // "/trigger/:trigger_kind/resume_suspended_trigger_job/*trigger_path", + // post(resume_suspended_trigger_jobs), + // ) + // .route( + // "/trigger/:trigger_kind/cancel_suspended_trigger_job/*trigger_path", + // post(cancel_suspended_trigger_jobs), + // ); + // } router } diff --git a/backend/windmill-api/src/triggers/http/handler.rs b/backend/windmill-api/src/triggers/http/handler.rs index 74008159f6..0bbb2bc076 100644 --- a/backend/windmill-api/src/triggers/http/handler.rs +++ b/backend/windmill-api/src/triggers/http/handler.rs @@ -209,7 +209,7 @@ pub async fn insert_new_trigger_into_db( error_handler_path, error_handler_args, retry, - active_mode + suspended_mode ) VALUES ( $1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16, $17, $18, $19, now(), $20, $21, $22, $23, $24 @@ -238,7 +238,7 @@ pub async fn insert_new_trigger_into_db( trigger.error_handling.error_handler_path, trigger.error_handling.error_handler_args as _, trigger.error_handling.retry as _, - trigger.base.active_mode.unwrap_or(true) + trigger.base.suspended_mode.unwrap_or(true) ) .execute(&mut *tx) .await?; @@ -500,7 +500,7 @@ impl TriggerCrud for HttpTrigger { error_handler_path = $20, error_handler_args = $21, retry = $22, - active_mode = $23 + suspended_mode = $23 WHERE workspace_id = $24 AND path = $25 @@ -527,7 +527,7 @@ impl TriggerCrud for HttpTrigger { trigger.error_handling.error_handler_path, trigger.error_handling.error_handler_args as _, trigger.error_handling.retry as _, - trigger.base.active_mode.unwrap_or(true), + trigger.base.suspended_mode.unwrap_or(true), workspace_id, path, ) @@ -561,7 +561,7 @@ impl TriggerCrud for HttpTrigger { error_handler_path = $17, error_handler_args = $18, retry = $19, - active_mode = $20 + suspended_mode = $20 WHERE workspace_id = $21 AND path = $22 @@ -585,7 +585,7 @@ impl TriggerCrud for HttpTrigger { trigger.error_handling.error_handler_path, trigger.error_handling.error_handler_args as _, trigger.error_handling.retry as _, - trigger.base.active_mode.unwrap_or(true), + trigger.base.suspended_mode.unwrap_or(true), workspace_id, path, ) @@ -1050,7 +1050,7 @@ async fn route_job( .map_err(|e| e.into_response())?; let trigger_info = TriggerMetadata::new(Some(trigger.path.clone()), JobTriggerKind::Http); - if !trigger.active_mode { + if !trigger.suspended_mode { let _ = trigger_runnable( &db, Some(user_db), @@ -1064,7 +1064,7 @@ async fn route_job( trigger.error_handler_args.as_ref(), format!("http_trigger/{}", trigger.path), None, - trigger.active_mode, + trigger.suspended_mode, trigger_info, ) .await @@ -1157,7 +1157,7 @@ async fn route_job( trigger.error_handler_args.as_ref(), format!("http_trigger/{}", trigger.path), None, - trigger.active_mode, + trigger.suspended_mode, trigger_info, ) .await diff --git a/backend/windmill-api/src/triggers/http/mod.rs b/backend/windmill-api/src/triggers/http/mod.rs index 8dcc43b754..4f9d049ca0 100644 --- a/backend/windmill-api/src/triggers/http/mod.rs +++ b/backend/windmill-api/src/triggers/http/mod.rs @@ -48,7 +48,7 @@ pub struct TriggerRoute { error_handler_path: Option, error_handler_args: Option>>, retry: Option>, - active_mode: bool, + suspended_mode: bool, } pub struct RoutersCache { @@ -259,7 +259,7 @@ pub async fn refresh_routers(db: &DB) -> Result<(bool, RwLockReadGuard<'_, Route error_handler_path, error_handler_args as "error_handler_args: _", retry as "retry: _", - active_mode + suspended_mode FROM http_trigger WHERE diff --git a/backend/windmill-api/src/triggers/listener.rs b/backend/windmill-api/src/triggers/listener.rs index 38f3a6f770..dfafa47f29 100644 --- a/backend/windmill-api/src/triggers/listener.rs +++ b/backend/windmill-api/src/triggers/listener.rs @@ -22,7 +22,7 @@ use tokio::sync::RwLock; use windmill_common::{ error::{Error, Result}, jobs::JobTriggerKind, - triggers::{TriggerMetadata, TriggerKind}, + triggers::{TriggerKind, TriggerMetadata}, utils::report_critical_error, DB, INSTANCE_NAME, }; @@ -71,7 +71,7 @@ pub trait Listener: TriggerCrud + TriggerJobArgs { "error_handler_path", "error_handler_args", "retry", - "active_mode" + "suspended_mode", ]; fields.extend_from_slice(Self::ADDITIONAL_SELECT_FIELDS); @@ -105,7 +105,7 @@ pub trait Listener: TriggerCrud + TriggerJobArgs { trigger_config: trigger.config, error_handling: Some(trigger.error_handling), trigger_mode: true, - active_mode: Some(trigger.base.active_mode), + suspended_mode: Some(trigger.base.suspended_mode), }) .collect_vec(); @@ -155,7 +155,7 @@ pub trait Listener: TriggerCrud + TriggerJobArgs { trigger_mode: false, is_flow: capture.is_flow, error_handling: None, - active_mode: None, + suspended_mode: None, }) .collect_vec(); @@ -516,7 +516,7 @@ pub trait Listener: TriggerCrud + TriggerJobArgs { error_handler_args, format!("{}_trigger/{}", Self::TRIGGER_KIND, listening_trigger.path), None, - listening_trigger.active_mode.unwrap_or(false), + listening_trigger.suspended_mode.unwrap_or(false), TriggerMetadata::new(Some(listening_trigger.path.clone()), Self::JOB_TRIGGER_KIND), ) .await?; @@ -812,7 +812,7 @@ pub struct ListeningTrigger { pub script_path: String, pub trigger_mode: bool, pub error_handling: Option, - pub active_mode: Option, + pub suspended_mode: Option, } impl ListeningTrigger { diff --git a/backend/windmill-api/src/triggers/mod.rs b/backend/windmill-api/src/triggers/mod.rs index 7d7cfee9b3..f716eebab1 100644 --- a/backend/windmill-api/src/triggers/mod.rs +++ b/backend/windmill-api/src/triggers/mod.rs @@ -30,14 +30,14 @@ pub mod sqs; #[cfg(feature = "websocket")] pub mod websocket; +pub mod global_handler; mod handler; mod listener; pub mod trigger_helpers; -pub mod global_handler; #[allow(unused)] pub(crate) use handler::TriggerCrud; -pub use handler::{generate_trigger_routers, get_triggers_count_internal, TriggersCount}; +pub use handler::{generate_trigger_routers, get_triggers_count_internal, TriggerForReassignment, TriggersCount}; pub use listener::start_all_listeners; #[allow(unused)] pub(crate) use listener::Listener; @@ -62,7 +62,7 @@ pub struct BaseTrigger { pub email: String, pub edited_at: DateTime, pub extra_perms: Option, - pub active_mode: bool, + pub suspended_mode: bool, } #[derive(Debug, FromRow, Clone, Serialize, Deserialize)] @@ -125,7 +125,7 @@ pub struct BaseTriggerData { pub script_path: String, pub is_flow: bool, pub enabled: Option, - pub active_mode: Option, + pub suspended_mode: Option, } #[derive(Debug, Clone, Serialize, Deserialize)] @@ -157,7 +157,3 @@ impl Default for StandardTriggerQuery { Self { page: Some(0), per_page: Some(100), path: None, path_start: None, is_flow: None } } } - -lazy_static::lazy_static! { - pub static ref INACTIVE_TRIGGER_SCHEDULED_FOR_DATE: DateTime = Utc.with_ymd_and_hms(9999, 12, 31, 23, 59, 59).unwrap(); -} diff --git a/backend/windmill-api/src/triggers/mqtt/handler.rs b/backend/windmill-api/src/triggers/mqtt/handler.rs index 310ad875bb..bc62f2c2a5 100644 --- a/backend/windmill-api/src/triggers/mqtt/handler.rs +++ b/backend/windmill-api/src/triggers/mqtt/handler.rs @@ -101,7 +101,7 @@ impl TriggerCrud for MqttTrigger { error_handler_path, error_handler_args, retry, - active_mode + suspended_mode ) VALUES ( $1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16, $17 @@ -122,7 +122,7 @@ impl TriggerCrud for MqttTrigger { trigger.error_handling.error_handler_path, trigger.error_handling.error_handler_args as _, trigger.error_handling.retry as _, - trigger.base.active_mode.unwrap_or(true) + trigger.base.suspended_mode.unwrap_or(true) ) .execute(tx) .await?; @@ -171,7 +171,7 @@ impl TriggerCrud for MqttTrigger { error_handler_path = $14, error_handler_args = $15, retry = $16, - active_mode = $17 + suspended_mode = $17 WHERE workspace_id = $12 AND path = $13 @@ -192,7 +192,7 @@ impl TriggerCrud for MqttTrigger { trigger.error_handling.error_handler_path, trigger.error_handling.error_handler_args as _, trigger.error_handling.retry as _, - trigger.base.active_mode.unwrap_or(true) + trigger.base.suspended_mode.unwrap_or(true) ) .execute(tx) .await?; diff --git a/backend/windmill-api/src/triggers/postgres/handler.rs b/backend/windmill-api/src/triggers/postgres/handler.rs index 73f5187cd5..089452be31 100644 --- a/backend/windmill-api/src/triggers/postgres/handler.rs +++ b/backend/windmill-api/src/triggers/postgres/handler.rs @@ -126,7 +126,7 @@ impl TriggerCrud for PostgresTrigger { error_handler_path, error_handler_args, retry, - active_mode + suspended_mode ) VALUES ( $1, $2, $3, $4, $5, $6, $7, $8, $9, $10, now(), $11, $12, $13, $14 ) @@ -144,7 +144,7 @@ impl TriggerCrud for PostgresTrigger { trigger.error_handling.error_handler_path, trigger.error_handling.error_handler_args as _, trigger.error_handling.retry as _, - trigger.base.active_mode.unwrap_or(true) + trigger.base.suspended_mode.unwrap_or(true) ) .execute(tx) .await?; @@ -229,7 +229,7 @@ impl TriggerCrud for PostgresTrigger { error_handler_path = $11, error_handler_args = $12, retry = $13, - active_mode = $14 + suspended_mode = $14 WHERE workspace_id = $9 AND path = $10 "#, @@ -246,7 +246,7 @@ impl TriggerCrud for PostgresTrigger { trigger.error_handling.error_handler_path, trigger.error_handling.error_handler_args as _, trigger.error_handling.retry as _, - trigger.base.active_mode.unwrap_or(true) + trigger.base.suspended_mode.unwrap_or(true) ) .execute(tx) .await?; diff --git a/backend/windmill-api/src/triggers/trigger_helpers.rs b/backend/windmill-api/src/triggers/trigger_helpers.rs index dd997e06a9..e8d50aa692 100644 --- a/backend/windmill-api/src/triggers/trigger_helpers.rs +++ b/backend/windmill-api/src/triggers/trigger_helpers.rs @@ -16,7 +16,7 @@ use windmill_common::{ jobs::{get_has_preprocessor_from_content_and_lang, script_path_to_payload, JobPayload}, scripts::{get_full_hub_script_by_path, ScriptHash, ScriptLang}, triggers::{ - HubOrWorkspaceId, RunnableFormat, RunnableFormatVersion, TriggerMetadata, TriggerKind, + HubOrWorkspaceId, RunnableFormat, RunnableFormatVersion, TriggerKind, TriggerMetadata, RUNNABLE_FORMAT_VERSION_CACHE, }, users::username_to_permissioned_as, @@ -25,12 +25,6 @@ use windmill_common::{ }; use windmill_queue::{push, PushArgs, PushArgsOwned, PushIsolationLevel}; -/// Helper function to check if triggers should use queue mode. -/// Queue mode suspends jobs by scheduling them for a far future date. -fn is_queue_mode(active_mode: Option) -> bool { - matches!(active_mode, Some(false)) -} - #[cfg(feature = "enterprise")] use crate::jobs::check_license_key_valid; use crate::{ @@ -40,7 +34,6 @@ use crate::{ push_flow_job_by_path_into_queue, push_script_job_by_path_into_queue, result_to_response, run_wait_result_internal, RunJobQuery, }, - triggers::INACTIVE_TRIGGER_SCHEDULED_FOR_DATE, utils::check_scopes, HTTP_CLIENT, }; @@ -526,7 +519,7 @@ pub async fn trigger_runnable_inner( trigger_path: String, job_id: Option, trigger: TriggerMetadata, - active_mode: Option, + suspended_mode: Option, ) -> Result<(Uuid, Option, Option)> { let error_handler_args = error_handler_args.map(|args| { let args = args @@ -538,13 +531,8 @@ pub async fn trigger_runnable_inner( }); let user_db = user_db.unwrap_or_else(|| UserDB::new(db.clone())); - let scheduled_for = if is_queue_mode(active_mode) { - Some(INACTIVE_TRIGGER_SCHEDULED_FOR_DATE.clone()) - } else { - None - }; let (uuid, delete_after_use, early_return) = if is_flow { - let run_query = RunJobQuery { job_id, scheduled_for, ..Default::default() }; + let run_query = RunJobQuery { job_id, suspended_mode, ..Default::default() }; let path = StripPath(runnable_path.to_string()); let (uuid, early_return) = push_flow_job_by_path_into_queue( authed, @@ -572,7 +560,7 @@ pub async fn trigger_runnable_inner( trigger_path, job_id, trigger, - scheduled_for, + suspended_mode, ) .await?; (uuid, delete_after_use, None) @@ -595,7 +583,7 @@ pub async fn trigger_runnable( error_handler_args: Option<&sqlx::types::Json>>, trigger_path: String, job_id: Option, - active_mode: bool, + suspended_mode: bool, trigger: TriggerMetadata, ) -> Result { let uuid = trigger_runnable_inner( @@ -612,7 +600,7 @@ pub async fn trigger_runnable( trigger_path, job_id, trigger, - Some(active_mode), + Some(suspended_mode), ) .await? .0; @@ -768,10 +756,10 @@ async fn trigger_script_internal( trigger_path: String, job_id: Option, trigger: TriggerMetadata, - scheduled_for: Option>, + suspended_mode: Option, ) -> Result<(Uuid, Option)> { if retry.is_none() && error_handler_path.is_none() { - let run_query = RunJobQuery { job_id, scheduled_for, ..Default::default() }; + let run_query = RunJobQuery { job_id, suspended_mode, ..Default::default() }; let path = StripPath(script_path.to_string()); push_script_job_by_path_into_queue( authed, @@ -798,7 +786,7 @@ async fn trigger_script_internal( trigger_path, job_id, trigger, - scheduled_for, + suspended_mode, ) .await } @@ -817,7 +805,7 @@ async fn trigger_script_with_retry_and_error_handler( trigger_path: String, job_id: Option, trigger: TriggerMetadata, - scheduled_for: Option>, + suspended_mode: Option, ) -> Result<(Uuid, Option)> { #[cfg(feature = "enterprise")] check_license_key_valid().await?; @@ -911,7 +899,7 @@ async fn trigger_script_with_retry_and_error_handler( email, permissioned_as, authed.token_prefix.as_deref(), - scheduled_for, + None, None, None, None, @@ -930,6 +918,7 @@ async fn trigger_script_with_retry_and_error_handler( None, None, Some(trigger), + suspended_mode, ) .await?; tx.commit().await?; diff --git a/backend/windmill-api/src/triggers/websocket/handler.rs b/backend/windmill-api/src/triggers/websocket/handler.rs index 2f23c654a9..c2df1fec10 100644 --- a/backend/windmill-api/src/triggers/websocket/handler.rs +++ b/backend/windmill-api/src/triggers/websocket/handler.rs @@ -112,7 +112,7 @@ impl TriggerCrud for WebsocketTrigger { error_handler_path, error_handler_args, retry, - active_mode + suspended_mode ) VALUES ( $1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, now(), $14, $15, $16, $17 ) @@ -136,7 +136,7 @@ impl TriggerCrud for WebsocketTrigger { trigger.error_handling.error_handler_path, trigger.error_handling.error_handler_args as _, trigger.error_handling.retry as _, - trigger.base.active_mode.unwrap_or(true) + trigger.base.suspended_mode.unwrap_or(true) ) .execute(&mut *tx) .await?; @@ -189,7 +189,7 @@ impl TriggerCrud for WebsocketTrigger { error_handler_path = $14, error_handler_args = $15, retry = $16, - active_mode = $17 + suspended_mode = $17 WHERE workspace_id = $12 AND path = $13 ", @@ -213,7 +213,7 @@ impl TriggerCrud for WebsocketTrigger { trigger.error_handling.error_handler_path, trigger.error_handling.error_handler_args as _, trigger.error_handling.retry as _, - trigger.base.active_mode.unwrap_or(true) + trigger.base.suspended_mode.unwrap_or(true) ) .execute(&mut *tx) .await?; diff --git a/backend/windmill-api/src/triggers/websocket/listener.rs b/backend/windmill-api/src/triggers/websocket/listener.rs index b1261fc30e..b293c2c1f5 100644 --- a/backend/windmill-api/src/triggers/websocket/listener.rs +++ b/backend/windmill-api/src/triggers/websocket/listener.rs @@ -334,7 +334,7 @@ impl Listener for WebsocketTrigger { trigger_config, script_path, error_handling, - active_mode, + suspended_mode, .. } = listening_trigger; @@ -367,9 +367,9 @@ impl Listener for WebsocketTrigger { ), None => (None, None, None), }; - let active_mode = active_mode.unwrap_or(false); + let suspended_mode = suspended_mode.unwrap_or(false); let trigger = TriggerMetadata::new(Some(path.to_owned()), Self::JOB_TRIGGER_KIND); - if active_mode || extra.is_none() { + if suspended_mode || extra.is_none() { trigger_runnable( db, None, @@ -383,7 +383,7 @@ impl Listener for WebsocketTrigger { error_handler_args, format!("websocket_trigger/{}", listening_trigger.path), None, - active_mode, + suspended_mode, trigger, ) .await?; diff --git a/backend/windmill-common/src/jobs.rs b/backend/windmill-common/src/jobs.rs index f32b3447c8..5f58a6b151 100644 --- a/backend/windmill-common/src/jobs.rs +++ b/backend/windmill-common/src/jobs.rs @@ -94,6 +94,7 @@ pub enum JobKind { FlowNode, AppScript, AIAgent, + Unassigned, } impl JobKind { diff --git a/backend/windmill-common/src/triggers.rs b/backend/windmill-common/src/triggers.rs index d3f0e1d98f..51c5ea41a2 100644 --- a/backend/windmill-common/src/triggers.rs +++ b/backend/windmill-common/src/triggers.rs @@ -85,6 +85,7 @@ lazy_static! { Cache::new(1000); } +#[derive(Debug, Clone)] pub struct TriggerMetadata { pub trigger_path: Option, pub trigger_kind: JobTriggerKind, diff --git a/backend/windmill-queue/src/jobs.rs b/backend/windmill-queue/src/jobs.rs index 78120b6e42..637fe3165c 100644 --- a/backend/windmill-queue/src/jobs.rs +++ b/backend/windmill-queue/src/jobs.rs @@ -479,6 +479,7 @@ pub async fn push_init_job<'c>( None, None, None, + None, ) .await?; inner_tx.commit().await?; @@ -538,6 +539,7 @@ pub async fn push_periodic_bash_job<'c>( None, None, None, + None, ) .await?; inner_tx.commit().await?; @@ -1384,6 +1386,7 @@ async fn restart_job_if_perpetual_inner( None, None, None, + None, ) .await?; tx.commit().await?; @@ -1935,6 +1938,7 @@ pub async fn push_error_handler<'a, 'c, T: Serialize + Send + Sync>( None, None, None, + None, ) .await?; tx.commit().await?; @@ -3859,6 +3863,7 @@ pub async fn push<'c, 'd>( // NOTE: Only works with dependency jobs triggered by relative imports debounce_job_id_o: Option, trigger: Option, + suspended_mode: Option, ) -> Result<(Uuid, Transaction<'c, Postgres>), Error> { #[cfg(feature = "cloud")] if *CLOUD_HOSTED { @@ -4045,7 +4050,7 @@ pub async fn push<'c, 'd>( script_hash, script_path, raw_code_tuple, - job_kind, + mut job_kind, raw_flow, flow_status, language, @@ -5274,6 +5279,12 @@ pub async fn push<'c, 'd>( root_job }; + let (job_kind, suspend, suspend_until) = if suspended_mode.unwrap_or(false) { + (JobKind::Unassigned, Some(1), Some(Utc::now() + chrono::Duration::days(30))) + } else { + (job_kind, None, None) + }; + sqlx::query!( "WITH inserted_job AS ( INSERT INTO v2_job (id, workspace_id, raw_code, raw_lock, raw_flow, tag, parent_job, @@ -5294,8 +5305,8 @@ pub async fn push<'c, 'd>( ON CONFLICT (job_id) DO UPDATE SET email = EXCLUDED.email, username = EXCLUDED.username, is_admin = EXCLUDED.is_admin, is_operator = EXCLUDED.is_operator, folders = EXCLUDED.folders, groups = EXCLUDED.groups, workspace_id = EXCLUDED.workspace_id, end_user_email = EXCLUDED.end_user_email ) INSERT INTO v2_job_queue - (workspace_id, id, running, scheduled_for, started_at, tag, priority) - VALUES ($2, $1, $28, COALESCE($29, now()), CASE WHEN $27 OR $40 THEN now() END, $30, $31)", + (workspace_id, id, running, scheduled_for, started_at, tag, priority, suspend, suspend_until) + VALUES ($2, $1, $28, COALESCE($29, now()), CASE WHEN $27 OR $40 THEN now() END, $30, $31, $42, $43)", job_id, workspace_id, raw_code, @@ -5341,6 +5352,8 @@ pub async fn push<'c, 'd>( trigger_kind as Option, running, end_user_email, + suspend, + suspend_until, ) .execute(&mut *tx) .warn_after_seconds(1) @@ -5413,6 +5426,7 @@ pub async fn push<'c, 'd>( JobKind::FlowNode => "jobs.run.flow_node", JobKind::AppScript => "jobs.run.app_script", JobKind::AIAgent => "jobs.run.ai_agent", + JobKind::Unassigned => "jobs.run.unassigned", }; let audit_author = if format!("u/{user}") != permissioned_as && user != permissioned_as { diff --git a/backend/windmill-queue/src/schedule.rs b/backend/windmill-queue/src/schedule.rs index 4a658a3682..b517283b62 100644 --- a/backend/windmill-queue/src/schedule.rs +++ b/backend/windmill-queue/src/schedule.rs @@ -506,6 +506,7 @@ pub async fn push_scheduled_job<'c>( Some(schedule.path.clone()), JobTriggerKind::Schedule, )), + None, ) .warn_after_seconds_with_sql(1, "push in push_scheduled_job".to_string()) .await?; diff --git a/backend/windmill-worker/src/ai/tools.rs b/backend/windmill-worker/src/ai/tools.rs index 89ccc58e55..808c24b22e 100644 --- a/backend/windmill-worker/src/ai/tools.rs +++ b/backend/windmill-worker/src/ai/tools.rs @@ -471,7 +471,8 @@ async fn execute_windmill_tool( true, None, None, - None + None, + None, ) .await?; diff --git a/backend/windmill-worker/src/worker.rs b/backend/windmill-worker/src/worker.rs index f1f2f5cef2..8d7caeda55 100644 --- a/backend/windmill-worker/src/worker.rs +++ b/backend/windmill-worker/src/worker.rs @@ -1017,8 +1017,6 @@ pub fn start_interactive_worker_shell( } }; - - match pulled_job { Ok(Some(job)) => { tracing::debug!(worker = %worker_name, hostname = %hostname, "started handling of job {}", job.id); @@ -3152,6 +3150,7 @@ async fn try_validate_schema( JobKind::Noop => 13, JobKind::FlowNode => 14, JobKind::AIAgent => 15, + JobKind::Unassigned => 16, }; let sv = match job.runnable_id { diff --git a/backend/windmill-worker/src/worker_flow.rs b/backend/windmill-worker/src/worker_flow.rs index 39cb479192..6a6fd7c8b1 100644 --- a/backend/windmill-worker/src/worker_flow.rs +++ b/backend/windmill-worker/src/worker_flow.rs @@ -3337,7 +3337,8 @@ async fn push_next_flow_job( continue_with_runners, None, None, - None + None, + None, ) .warn_after_seconds(2) .await?; diff --git a/backend/windmill-worker/src/worker_lockfiles.rs b/backend/windmill-worker/src/worker_lockfiles.rs index 54a5584e4f..0034aad108 100644 --- a/backend/windmill-worker/src/worker_lockfiles.rs +++ b/backend/windmill-worker/src/worker_lockfiles.rs @@ -671,7 +671,8 @@ pub async fn trigger_dependents_to_recompute_dependencies( false, None, debounce_job_id_o, - None + None, + None, ) .await?; diff --git a/frontend/src/lib/components/triggers/TriggerActiveMode.svelte b/frontend/src/lib/components/triggers/TriggerActiveMode.svelte index bf0eb7fdf1..3715cf4b24 100644 --- a/frontend/src/lib/components/triggers/TriggerActiveMode.svelte +++ b/frontend/src/lib/components/triggers/TriggerActiveMode.svelte @@ -11,20 +11,20 @@ import { type JobTriggerType } from './utils' import { ChevronLeft, ChevronRight } from 'lucide-svelte' type Props = { - active_mode: boolean + suspended_mode: boolean triggerPath: string jobTriggerKind: JobTriggerType } - let { active_mode = $bindable(), jobTriggerKind, triggerPath }: Props = $props() + let { suspended_mode = $bindable(), jobTriggerKind, triggerPath }: Props = $props() - let wasInInactiveMode = $state(!active_mode) + let wasInInactiveMode = $state(!suspended_mode) let shouldShowModal = $state(false) $effect(() => { - if (active_mode && wasInInactiveMode) { + if (suspended_mode && wasInInactiveMode) { shouldShowModal = true - } else if (!active_mode) { + } else if (!suspended_mode) { wasInInactiveMode = true shouldShowModal = false } @@ -143,149 +143,4 @@ } - - -{#if shouldShowModal} - -
- {#if loading} -
-
- Loading queued jobs... -
- {:else if error} -
- {error} -
- {:else if queuedJobs.length === 0} -
-
-
No suspended jobs found
-
This trigger has no suspended jobs waiting to be processed.
-
-
- {:else} -
-
-

- Suspended Jobs {#if hasMorePages}(Page {currentPage}){:else}({queuedJobs.length}){/if} -

-

Click on any job to view details

-
- -
-
-
Status
-
Started
-
Duration
-
Path
-
Triggered by
-
-
- -
- {#each queuedJobs as job} -
- { - window.open(`/run/${job.id}?workspace=${workspace}`, '_blank') - }} - /> -
- {/each} -
-
-
- - {#if queuedJobs.length > 0 && (currentPage > 1 || hasMorePages)} -
-
- - {queuedJobs.length} - {hasMorePages ? '+' : ''} suspended job{queuedJobs.length === 1 ? '' : 's'} - -
- -
-
Page {currentPage}
- - - - -
-
- {/if} - -
-

- You are switching this trigger from inactive to active mode. What would you like to do - with the {hasMorePages - ? `${queuedJobs.length}+ suspended` - : `${queuedJobs.length} suspended`} job{queuedJobs.length === 1 ? '' : 's'}? -

-
- {/if} - -
- {#if !loading && !error && queuedJobs.length > 0} - - - - {/if} -
-
-
-{/if} + diff --git a/frontend/src/lib/components/triggers/TriggerSuspendedJobsModal.svelte b/frontend/src/lib/components/triggers/TriggerSuspendedJobsModal.svelte new file mode 100644 index 0000000000..2979537b2b --- /dev/null +++ b/frontend/src/lib/components/triggers/TriggerSuspendedJobsModal.svelte @@ -0,0 +1,289 @@ + + +{#if shouldShowModal} + +
+ {#if loading} +
+
+ Loading queued jobs... +
+ {:else if error} +
+ {error} +
+ {:else if queuedJobs.length === 0} +
+
+
No suspended jobs found
+
This trigger has no suspended jobs waiting to be processed.
+
+
+ {:else} +
+
+

+ Suspended Jobs {#if hasMorePages}(Page {currentPage}){:else}({queuedJobs.length}){/if} +

+

Click on any job to view details

+
+ +
+
+
Status
+
Started
+
Duration
+
Path
+
Triggered by
+
+
+ +
+ {#each queuedJobs as job} +
+ { + window.open(`/run/${job.id}?workspace=${workspace}`, '_blank') + }} + /> +
+ {/each} +
+
+
+ + {#if queuedJobs.length > 0 && (currentPage > 1 || hasMorePages)} +
+
+ + {queuedJobs.length} + {hasMorePages ? '+' : ''} suspended job{queuedJobs.length === 1 ? '' : 's'} + +
+ +
+
Page {currentPage}
+ + + + +
+
+ {/if} + +
+

+ You are switching this trigger from inactive to active mode. What would you like to do + with the {hasMorePages + ? `${queuedJobs.length}+ suspended` + : `${queuedJobs.length} suspended`} job{queuedJobs.length === 1 ? '' : 's'}? +

+
+ {/if} + +
+ {#if !loading && !error && queuedJobs.length > 0} + + + + {/if} +
+
+
+{/if} diff --git a/frontend/src/lib/components/triggers/email/EmailTriggerEditorInner.svelte b/frontend/src/lib/components/triggers/email/EmailTriggerEditorInner.svelte index ffc43d8e80..7ef99e6031 100644 --- a/frontend/src/lib/components/triggers/email/EmailTriggerEditorInner.svelte +++ b/frontend/src/lib/components/triggers/email/EmailTriggerEditorInner.svelte @@ -67,7 +67,7 @@ let error_handler_args: Record = $state({}) let retry: Retry | undefined = $state() let enabled = $state(false) - let active_mode = $state(true) + let suspended_mode = $state(true) // Component references let drawer = $state(undefined) let initialConfig: NewEmailTrigger | undefined = undefined @@ -166,7 +166,7 @@ retry = cfg?.retry errorHandlerSelected = getHandlerType(error_handler_path ?? '') enabled = cfg?.enabled ?? false - active_mode = cfg?.active_mode ?? true + suspended_mode = cfg?.suspended_mode ?? true } async function loadTrigger(defaultConfig?: Partial): Promise { @@ -218,7 +218,7 @@ error_handler_args, retry, enabled, - active_mode + suspended_mode } return nCfg @@ -316,7 +316,7 @@ {/if} - + = $state({}) let retry: Retry | undefined = $state() - let active_mode = $state(true) + let suspended_mode = $state(true) let { useDrawer = true, description = undefined, @@ -190,7 +190,7 @@ can_write = canWrite(cfg?.path, cfg?.extra_perms, $userStore) error_handler_path = cfg?.error_handler_path error_handler_args = cfg?.error_handler_args ?? {} - active_mode = cfg?.active_mode ?? true + suspended_mode = cfg?.suspended_mode ?? true retry = cfg?.retry auto_acknowledge_msg = cfg?.auto_acknowledge_msg ?? true ack_deadline = cfg?.ack_deadline @@ -219,7 +219,7 @@ function getGcpConfig() { return { - active_mode, + suspended_mode, gcp_resource_path, subscription_mode, subscription_id, @@ -396,7 +396,7 @@ - + {/if} ({})) @@ -293,7 +293,7 @@ error_handler_args = cfg?.error_handler_args ?? {} retry = cfg?.retry errorHandlerSelected = getHandlerType(error_handler_path ?? '') - active_mode = cfg?.active_mode ?? true + suspended_mode = cfg?.suspended_mode ?? true } async function loadTrigger(defaultConfig?: Partial): Promise { @@ -363,7 +363,7 @@ error_handler_path, error_handler_args, retry, - active_mode + suspended_mode } return nCfg @@ -381,7 +381,6 @@ } } - // Update config for captures function getCaptureConfig() { const newCaptureConfig = { @@ -612,7 +611,7 @@ - + {/if} = $state({}) let retry: Retry | undefined = $state() - let active_mode = $state(true) + let suspended_mode = $state(true) const isValid = $derived( !!kafkaResourcePath && @@ -190,7 +190,7 @@ error_handler_args = cfg?.error_handler_args ?? {} retry = cfg?.retry errorHandlerSelected = getHandlerType(error_handler_path ?? '') - active_mode = cfg?.active_mode ?? true + suspended_mode = cfg?.suspended_mode ?? true } async function loadTrigger(defaultConfig?: Record): Promise { @@ -219,7 +219,7 @@ error_handler_path, error_handler_args, retry, - active_mode + suspended_mode } } @@ -387,7 +387,7 @@ - + {/if} = $state({}) let retry: Retry | undefined = $state() - let active_mode = $state(true) + let suspended_mode = $state(true) let optionTabSelected: 'connection_options' | 'error_handler' | 'retries' = $state('connection_options') @@ -204,7 +204,7 @@ error_handler_args = cfg?.error_handler_args ?? {} retry = cfg?.retry errorHandlerSelected = getHandlerType(error_handler_path ?? '') - active_mode = cfg?.active_mode ?? true + suspended_mode = cfg?.suspended_mode ?? true activateV5Options.topic_alias_maximum = Boolean(v5_config.topic_alias_maximum) activateV5Options.session_expiry_interval = Boolean(v5_config.session_expiry_interval) } catch (error) { @@ -244,7 +244,7 @@ error_handler_path, error_handler_args, retry, - active_mode + suspended_mode } } @@ -421,7 +421,7 @@ {/if} - + = $state({}) let retry: Retry | undefined = $state() - let active_mode = $state(true) + let suspended_mode = $state(true) const saveDisabled = $derived( pathError != '' || emptyString(script_path) || !can_write || !isValid @@ -192,7 +192,7 @@ error_handler_args = cfg?.error_handler_args ?? {} retry = cfg?.retry errorHandlerSelected = getHandlerType(error_handler_path ?? '') - active_mode = cfg?.active_mode ?? true + suspended_mode = cfg?.suspended_mode ?? true } async function loadTrigger(defaultConfig?: Record): Promise { @@ -222,7 +222,7 @@ error_handler_path, error_handler_args, retry, - active_mode + suspended_mode } } @@ -402,7 +402,7 @@ {/if} - + = $state({}) let retry: Retry | undefined = $state() - let active_mode = $state(true) + let suspended_mode = $state(true) const errorMessage = $derived.by(() => { if (relations && relations.length > 0) { @@ -302,7 +302,7 @@ error_handler_path, error_handler_args, retry, - active_mode + suspended_mode } return cfg } @@ -323,7 +323,7 @@ error_handler_args = cfg?.error_handler_args ?? {} retry = cfg?.retry errorHandlerSelected = getHandlerType(error_handler_path ?? '') - active_mode = cfg?.active_mode ?? true + suspended_mode = cfg?.suspended_mode ?? true } async function loadTrigger(defaultConfig?: Record): Promise { @@ -580,7 +580,7 @@ {/if} - +
{#snippet badge()} {#if isEditor} diff --git a/frontend/src/lib/components/triggers/postgres/utils.ts b/frontend/src/lib/components/triggers/postgres/utils.ts index 641b6ad692..2c29b0e004 100644 --- a/frontend/src/lib/components/triggers/postgres/utils.ts +++ b/frontend/src/lib/components/triggers/postgres/utils.ts @@ -114,7 +114,7 @@ export async function savePostgresTriggerFromCfg( ? { error_handler_path: config.error_handler_path, error_handler_args: config.error_handler_path ? config.error_handler_args : undefined, - retry: config.retry, + retry: config.retry } : {} const requestBody: EditPostgresTrigger = { @@ -126,7 +126,7 @@ export async function savePostgresTriggerFromCfg( publication_name: config.publication_name, publication: config.publication, enabled: config.enabled, - active_mode: config.active_mode, + suspended_mode: config.suspended_mode, ...errorHandlerAndRetries } if (edit) { diff --git a/frontend/src/lib/components/triggers/sqs/SqsTriggerEditorInner.svelte b/frontend/src/lib/components/triggers/sqs/SqsTriggerEditorInner.svelte index 408f9004f3..95842ee422 100644 --- a/frontend/src/lib/components/triggers/sqs/SqsTriggerEditorInner.svelte +++ b/frontend/src/lib/components/triggers/sqs/SqsTriggerEditorInner.svelte @@ -89,7 +89,7 @@ let error_handler_path: string | undefined = $state() let error_handler_args: Record = $state({}) let retry: Retry | undefined = $state() - let active_mode = $state(true) + let suspended_mode = $state(true) const sqsConfig = $derived.by(getSaveCfg) const captureConfig = $derived.by(getCaptureConfig) @@ -178,7 +178,7 @@ error_handler_args = cfg?.error_handler_args ?? {} retry = cfg?.retry errorHandlerSelected = getHandlerType(error_handler_path ?? '') - active_mode = cfg?.active_mode ?? true + suspended_mode = cfg?.suspended_mode ?? true } catch (error) { sendUserToast(`Could not load SQS trigger config: ${error.body}`, true) } @@ -214,7 +214,7 @@ error_handler_path, error_handler_args, retry, - active_mode + suspended_mode } } @@ -388,7 +388,7 @@
- + {/if} = $state({}) let retry: Retry | undefined = $state() - let active_mode = $state(true) + let suspended_mode = $state(true) const websocketCfg = $derived.by(getSaveCfg) const captureConfig = $derived.by(isEditor ? getCaptureConfig : () => ({})) @@ -221,7 +221,7 @@ error_handler_args = cfg?.error_handler_args ?? {} retry = cfg?.retry errorHandlerSelected = getHandlerType(error_handler_path ?? '') - active_mode = cfg?.active_mode ?? true + suspended_mode = cfg?.suspended_mode ?? true } function getSaveCfg() { @@ -240,7 +240,7 @@ error_handler_path, error_handler_args, retry, - active_mode + suspended_mode } } @@ -494,7 +494,7 @@ /> - +