suspended mode draft

This commit is contained in:
hugocasa
2025-11-27 19:34:02 +01:00
parent 04808ea94b
commit b95a8586b9
52 changed files with 842 additions and 528 deletions
@@ -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;
@@ -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;
@@ -0,0 +1 @@
-- Add down migration script here
@@ -0,0 +1,2 @@
-- Add up migration script here
ALTER TYPE JOB_KIND ADD VALUE IF NOT EXISTS 'unassigned';
@@ -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;
@@ -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;
+2 -1
View File
@@ -208,7 +208,8 @@ impl RunJob {
false,
None,
debounce_job_id_o,
None
None,
None,
)
.await
.expect("push has to succeed");
+7
View File
@@ -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();
+19 -39
View File
@@ -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
+6 -3
View File
@@ -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?;
+4 -2
View File
@@ -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?;
+15
View File
@@ -1757,6 +1757,7 @@ pub struct RunJobQuery {
pub skip_preprocessor: Option<bool>,
pub poll_delay_ms: Option<u64>,
pub memory_id: Option<Uuid>,
pub suspended_mode: Option<bool>,
}
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?;
+5 -6
View File
@@ -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))
@@ -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<sqlx::types::Json<HashMap<String, Box<RawValue>>>>,
}
pub async fn reassign_suspended_jobs(
authed: ApiAuthed,
Extension(db): Extension<DB>,
Extension(user_db): Extension<UserDB>,
Path((w_id, trigger_kind, trigger_path)): Path<(String, JobTriggerKind, String)>,
) -> error::Result<Json<String>> {
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<UserDB>,
Path((w_id, trigger_kind, trigger_path)): Path<(String, JobTriggerKind, String)>,
) -> error::Result<Json<String>> {
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)))
}
+69 -17
View File
@@ -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<Utc>,
pub error_handler_path: Option<String>,
pub error_handler_args:
Option<sqlx::types::Json<HashMap<String, serde_json::Value>>>,
pub retry: Option<sqlx::types::Json<windmill_common::flows::Retry>>,
}
#[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<TriggerForReassignment> {
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<bool> {
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
}
@@ -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
@@ -48,7 +48,7 @@ pub struct TriggerRoute {
error_handler_path: Option<String>,
error_handler_args: Option<sqlx::types::Json<HashMap<String, serde_json::Value>>>,
retry: Option<sqlx::types::Json<Retry>>,
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
@@ -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<T> {
pub script_path: String,
pub trigger_mode: bool,
pub error_handling: Option<TriggerErrorHandling>,
pub active_mode: Option<bool>,
pub suspended_mode: Option<bool>,
}
impl<T> ListeningTrigger<T> {
+4 -8
View File
@@ -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<Utc>,
pub extra_perms: Option<serde_json::Value>,
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<bool>,
pub active_mode: Option<bool>,
pub suspended_mode: Option<bool>,
}
#[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> = Utc.with_ymd_and_hms(9999, 12, 31, 23, 59, 59).unwrap();
}
@@ -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?;
@@ -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?;
@@ -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>) -> 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<Uuid>,
trigger: TriggerMetadata,
active_mode: Option<bool>,
suspended_mode: Option<bool>,
) -> Result<(Uuid, Option<bool>, Option<String>)> {
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<HashMap<String, serde_json::Value>>>,
trigger_path: String,
job_id: Option<Uuid>,
active_mode: bool,
suspended_mode: bool,
trigger: TriggerMetadata,
) -> Result<axum::response::Response> {
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<Uuid>,
trigger: TriggerMetadata,
scheduled_for: Option<DateTime<Utc>>,
suspended_mode: Option<bool>,
) -> Result<(Uuid, Option<bool>)> {
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<Uuid>,
trigger: TriggerMetadata,
scheduled_for: Option<DateTime<Utc>>,
suspended_mode: Option<bool>,
) -> Result<(Uuid, Option<bool>)> {
#[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?;
@@ -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?;
@@ -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?;
+1
View File
@@ -94,6 +94,7 @@ pub enum JobKind {
FlowNode,
AppScript,
AIAgent,
Unassigned,
}
impl JobKind {
+1
View File
@@ -85,6 +85,7 @@ lazy_static! {
Cache::new(1000);
}
#[derive(Debug, Clone)]
pub struct TriggerMetadata {
pub trigger_path: Option<String>,
pub trigger_kind: JobTriggerKind,
+17 -3
View File
@@ -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<Uuid>,
trigger: Option<TriggerMetadata>,
suspended_mode: Option<bool>,
) -> 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<JobTriggerKind>,
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 {
+1
View File
@@ -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?;
+2 -1
View File
@@ -471,7 +471,8 @@ async fn execute_windmill_tool(
true,
None,
None,
None
None,
None,
)
.await?;
+1 -2
View File
@@ -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 {
+2 -1
View File
@@ -3337,7 +3337,8 @@ async fn push_next_flow_job(
continue_with_runners,
None,
None,
None
None,
None,
)
.warn_after_seconds(2)
.await?;
@@ -671,7 +671,8 @@ pub async fn trigger_dependents_to_recompute_dependencies(
false,
None,
debounce_job_id_o,
None
None,
None,
)
.await?;
@@ -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 @@
}
</script>
<Toggle bind:checked={active_mode} options={{ right: 'Active', left: 'Inactive' }} />
{#if shouldShowModal}
<Modal2
bind:isOpen={shouldShowModal}
title="{hasMorePages
? `${queuedJobs.length}+ suspended`
: `${queuedJobs.length} suspended`} job{queuedJobs.length === 1 ? '' : 's'} for this trigger"
target="#content"
fixedSize="lg"
>
<div class="flex w-full flex-col gap-4 h-full">
{#if loading}
<div class="flex items-center justify-center py-8">
<div
class="animate-spin h-6 w-6 border-2 border-blue-500 border-t-transparent rounded-full"
></div>
<span class="ml-2">Loading queued jobs...</span>
</div>
{:else if error}
<div class="bg-red-50 border border-red-200 text-red-700 px-4 py-3 rounded">
{error}
</div>
{:else if queuedJobs.length === 0}
<div class="flex flex-col items-center w-full py-12 px-4">
<div class="text-center">
<div class="text-base font-medium text-secondary mb-2">No suspended jobs found</div>
<div class="text-sm text-tertiary"
>This trigger has no suspended jobs waiting to be processed.</div
>
</div>
</div>
{:else}
<div class="flex-1 overflow-auto">
<div class="mb-3">
<h3 class="text-sm font-medium">
Suspended Jobs {#if hasMorePages}(Page {currentPage}){:else}({queuedJobs.length}){/if}
</h3>
<p class="text-xs text-gray-500 mt-1">Click on any job to view details</p>
</div>
<div class="divide-y h-full border min-w-[650px]" bind:clientWidth={containerWidth}>
<div
class="bg-surface-secondary sticky top-0 w-full py-2 pr-4 grid grid-runs-table-no-tag"
>
<div class="text-2xs px-2 font-semibold">Status</div>
<div class="text-xs font-semibold">Started</div>
<div class="text-xs font-semibold">Duration</div>
<div class="text-xs font-semibold">Path</div>
<div class="text-xs font-semibold">Triggered by</div>
<div class=""></div>
</div>
<div class="h-full">
{#each queuedJobs as job}
<div class="flex flex-row items-center h-[42px] w-full">
<RunRow
{job}
{containerWidth}
showTag={false}
activeLabel={null}
on:select={() => {
window.open(`/run/${job.id}?workspace=${workspace}`, '_blank')
}}
/>
</div>
{/each}
</div>
</div>
</div>
{#if queuedJobs.length > 0 && (currentPage > 1 || hasMorePages)}
<div
class="w-full bg-surface border-t flex flex-row justify-between p-2 items-center gap-2"
>
<div class="flex flex-row gap-2 items-center">
<span class="text-xs text-secondary">
{queuedJobs.length}
{hasMorePages ? '+' : ''} suspended job{queuedJobs.length === 1 ? '' : 's'}
</span>
</div>
<div class="flex flex-row gap-3 items-center">
<div class="flex text-xs text-secondary">Page {currentPage}</div>
<Button
variant="subtle"
size="xs2"
startIcon={{ icon: ChevronLeft }}
on:click={prevPage}
disabled={currentPage === 1 || loading}
loading={isPreviousLoading}
>
Previous
</Button>
<Button
variant="subtle"
size="xs2"
endIcon={{ icon: ChevronRight }}
on:click={nextPage}
disabled={!hasMorePages || loading}
loading={isNextLoading}
>
Next
</Button>
</div>
</div>
{/if}
<div class="bg-blue-50 p-4 rounded-lg">
<p class="text-sm text-blue-700">
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'}?
</p>
</div>
{/if}
<div class="flex gap-2 pt-4 border-t">
{#if !loading && !error && queuedJobs.length > 0}
<Button
variant="border"
size="sm"
onClick={discardAllJobs}
disabled={processingAction}
color="red"
>
{processingAction ? 'Discarding...' : 'Discard All Jobs'}
</Button>
<Button
variant="contained"
size="sm"
onClick={runAllJobs}
disabled={processingAction}
color="green"
>
{processingAction ? 'Running...' : 'Run All Jobs'}
</Button>
{/if}
</div>
</div>
</Modal2>
{/if}
<Toggle bind:checked={suspended_mode} options={{ right: 'Active', left: 'Inactive' }} />
@@ -0,0 +1,289 @@
<script lang="ts">
import Modal2 from '../common/modal/Modal2.svelte'
import type { QueuedJob } from '$lib/gen/types.gen'
import Button from '../common/button/Button.svelte'
import { workspaceStore } from '$lib/stores'
import RunRow from '../runs/RunRow.svelte'
import '../runs/runs-grid.css'
import { JobService, TriggerService } from '$lib/gen'
import { sendUserToast } from '$lib/toast'
import { type JobTriggerType } from './utils'
import { ChevronLeft, ChevronRight } from 'lucide-svelte'
type Props = {
suspended_mode: boolean
triggerPath: string
jobTriggerKind: JobTriggerType
}
let { suspended_mode = $bindable(), jobTriggerKind, triggerPath }: Props = $props()
let wasInInactiveMode = $state(!suspended_mode)
let shouldShowModal = $state(false)
$effect(() => {
if (suspended_mode && wasInInactiveMode) {
shouldShowModal = true
} else if (!suspended_mode) {
wasInInactiveMode = true
shouldShowModal = false
}
})
let queuedJobs = $state<QueuedJob[]>([])
let loading = $state(false)
let error = $state<string | null>(null)
let processingAction = $state(false)
let isPreviousLoading = $state(false)
let isNextLoading = $state(false)
let workspace = $workspaceStore!
let containerWidth = $state(1000)
let currentPage = $state(1)
let perPage = $state(20)
let hasMorePages = $derived(queuedJobs.length === perPage)
$effect(() => {
if (shouldShowModal) {
fetchQueuedJobs()
}
})
async function fetchQueuedJobs(resetPage = false) {
if (resetPage) {
currentPage = 1
}
loading = true
error = null
try {
const allSuspendedJobs = await JobService.listQueue({
workspace,
triggerKind: jobTriggerKind,
jobKinds: 'unassigned',
triggerPath,
running: false,
perPage,
page: currentPage
})
queuedJobs = allSuspendedJobs
} catch (e) {
error = `Failed to fetch queued jobs: ${e}`
console.error('Failed to fetch queued jobs:', e)
} finally {
loading = false
}
}
function nextPage() {
if (hasMorePages && !loading) {
isNextLoading = true
currentPage += 1
fetchQueuedJobs().finally(() => {
isNextLoading = false
})
}
}
function prevPage() {
if (currentPage > 1 && !loading) {
isPreviousLoading = true
currentPage -= 1
fetchQueuedJobs().finally(() => {
isPreviousLoading = false
})
}
}
async function runAllJobs() {
if (queuedJobs.length === 0) return
processingAction = true
error = null
try {
const resumedJobs = await TriggerService.resumeSuspendedTriggerJobs({
workspace,
triggerKind: jobTriggerKind,
triggerPath
})
sendUserToast(resumedJobs)
} catch (e) {
error = `Failed to run jobs: ${e}`
console.error('Failed to run jobs:', e)
} finally {
processingAction = false
closeModal()
}
}
async function discardAllJobs() {
if (queuedJobs.length === 0) return
processingAction = true
error = null
try {
await TriggerService.cancelSuspendedTriggerJobs({
workspace,
triggerKind: jobTriggerKind,
triggerPath
})
sendUserToast(`Successfully canceled all jobs`)
} catch (e) {
error = `Failed to discard jobs: ${e}`
console.error('Failed to discard jobs:', e)
} finally {
processingAction = false
closeModal()
}
}
function closeModal() {
wasInInactiveMode = false
shouldShowModal = false
}
</script>
{#if shouldShowModal}
<Modal2
bind:isOpen={shouldShowModal}
title="{hasMorePages
? `${queuedJobs.length}+ suspended`
: `${queuedJobs.length} suspended`} job{queuedJobs.length === 1 ? '' : 's'} for this trigger"
target="#content"
fixedSize="lg"
>
<div class="flex w-full flex-col gap-4 h-full">
{#if loading}
<div class="flex items-center justify-center py-8">
<div
class="animate-spin h-6 w-6 border-2 border-blue-500 border-t-transparent rounded-full"
></div>
<span class="ml-2">Loading queued jobs...</span>
</div>
{:else if error}
<div class="bg-red-50 border border-red-200 text-red-700 px-4 py-3 rounded">
{error}
</div>
{:else if queuedJobs.length === 0}
<div class="flex flex-col items-center w-full py-12 px-4">
<div class="text-center">
<div class="text-base font-medium text-secondary mb-2">No suspended jobs found</div>
<div class="text-sm text-tertiary"
>This trigger has no suspended jobs waiting to be processed.</div
>
</div>
</div>
{:else}
<div class="flex-1 overflow-auto">
<div class="mb-3">
<h3 class="text-sm font-medium">
Suspended Jobs {#if hasMorePages}(Page {currentPage}){:else}({queuedJobs.length}){/if}
</h3>
<p class="text-xs text-gray-500 mt-1">Click on any job to view details</p>
</div>
<div class="divide-y h-full border min-w-[650px]" bind:clientWidth={containerWidth}>
<div
class="bg-surface-secondary sticky top-0 w-full py-2 pr-4 grid grid-runs-table-no-tag"
>
<div class="text-2xs px-2 font-semibold">Status</div>
<div class="text-xs font-semibold">Started</div>
<div class="text-xs font-semibold">Duration</div>
<div class="text-xs font-semibold">Path</div>
<div class="text-xs font-semibold">Triggered by</div>
<div class=""></div>
</div>
<div class="h-full">
{#each queuedJobs as job}
<div class="flex flex-row items-center h-[42px] w-full">
<RunRow
{job}
{containerWidth}
showTag={false}
activeLabel={null}
on:select={() => {
window.open(`/run/${job.id}?workspace=${workspace}`, '_blank')
}}
/>
</div>
{/each}
</div>
</div>
</div>
{#if queuedJobs.length > 0 && (currentPage > 1 || hasMorePages)}
<div
class="w-full bg-surface border-t flex flex-row justify-between p-2 items-center gap-2"
>
<div class="flex flex-row gap-2 items-center">
<span class="text-xs text-secondary">
{queuedJobs.length}
{hasMorePages ? '+' : ''} suspended job{queuedJobs.length === 1 ? '' : 's'}
</span>
</div>
<div class="flex flex-row gap-3 items-center">
<div class="flex text-xs text-secondary">Page {currentPage}</div>
<Button
variant="subtle"
size="xs2"
startIcon={{ icon: ChevronLeft }}
on:click={prevPage}
disabled={currentPage === 1 || loading}
loading={isPreviousLoading}
>
Previous
</Button>
<Button
variant="subtle"
size="xs2"
endIcon={{ icon: ChevronRight }}
on:click={nextPage}
disabled={!hasMorePages || loading}
loading={isNextLoading}
>
Next
</Button>
</div>
</div>
{/if}
<div class="bg-blue-50 p-4 rounded-lg">
<p class="text-sm text-blue-700">
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'}?
</p>
</div>
{/if}
<div class="flex gap-2 pt-4 border-t">
{#if !loading && !error && queuedJobs.length > 0}
<Button
variant="border"
size="sm"
onClick={discardAllJobs}
disabled={processingAction}
color="red"
>
{processingAction ? 'Discarding...' : 'Discard All Jobs'}
</Button>
<Button
variant="contained"
size="sm"
onClick={runAllJobs}
disabled={processingAction}
color="green"
>
{processingAction ? 'Running...' : 'Run All Jobs'}
</Button>
{/if}
</div>
</div>
</Modal2>
{/if}
@@ -67,7 +67,7 @@
let error_handler_args: Record<string, any> = $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<Drawer | undefined>(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<EmailTrigger>): Promise<void> {
@@ -218,7 +218,7 @@
error_handler_args,
retry,
enabled,
active_mode
suspended_mode
}
return nCfg
@@ -316,7 +316,7 @@
</Section>
{/if}
<TriggerActiveMode triggerPath={path} jobTriggerKind={'email'} bind:active_mode />
<TriggerActiveMode triggerPath={path} jobTriggerKind={'email'} bind:suspended_mode />
<EmailTriggerEditorConfigSection
initialTriggerPath={initialPath}
@@ -29,7 +29,7 @@ export async function saveEmailTriggerFromCfg(
workspaced_local_part: emailCfg.workspaced_local_part,
error_handler_path: emailCfg.error_handler_path,
error_handler_args: emailCfg.error_handler_path ? emailCfg.error_handler_args : undefined,
active_mode: emailCfg.active_mode,
suspended_mode: emailCfg.suspended_mode,
enabled: emailCfg.enabled,
retry: emailCfg.retry
}
@@ -61,7 +61,7 @@
let error_handler_path: string | undefined = $state()
let error_handler_args: Record<string, any> = $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 @@
</div>
</Section>
<TriggerActiveMode triggerPath={path} jobTriggerKind={'gcp'} bind:active_mode />
<TriggerActiveMode triggerPath={path} jobTriggerKind={'gcp'} bind:suspended_mode />
{/if}
<GcpTriggerEditorConfigSection
@@ -32,7 +32,7 @@ export async function saveGcpTriggerFromCfg(
is_flow: cfg.is_flow,
auto_acknowledge_msg: cfg.auto_acknowledge_msg,
ack_deadline: cfg.ack_deadline,
active_mode: cfg.active_mode,
suspended_mode: cfg.suspended_mode,
...errorHandlerAndRetries
}
if (edit) {
@@ -112,7 +112,7 @@
let deploymentLoading = $state(false)
let optionTabSelected: 'request_options' | 'error_handler' | 'retries' = $state('request_options')
let errorHandlerSelected: ErrorHandler = $state('slack')
let active_mode = $state(true)
let suspended_mode = $state(true)
const isAdmin = $derived($userStore?.is_admin || $userStore?.is_super_admin)
const routeConfig = $derived.by(getRouteConfig)
const captureConfig = $derived.by(isEditor ? getCaptureConfig : () => ({}))
@@ -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<HttpTrigger>): Promise<void> {
@@ -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 @@
</div>
</Section>
<TriggerActiveMode triggerPath={path} jobTriggerKind={'http'} bind:active_mode />
<TriggerActiveMode triggerPath={path} jobTriggerKind={'http'} bind:suspended_mode />
{/if}
<RouteEditorConfigSection
@@ -62,7 +62,7 @@ export async function saveHttpRouteFromCfg(
error_handler_args: routeCfg.error_handler_path ? routeCfg.error_handler_args : undefined,
retry: routeCfg.retry,
enabled: routeCfg.enabled,
active_mode: routeCfg.active_mode
suspended_mode: routeCfg.suspended_mode
}
try {
if (edit) {
@@ -83,7 +83,7 @@
let error_handler_path: string | undefined = $state()
let error_handler_args: Record<string, any> = $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<string, any>): Promise<void> {
@@ -219,7 +219,7 @@
error_handler_path,
error_handler_args,
retry,
active_mode
suspended_mode
}
}
@@ -387,7 +387,7 @@
</div>
</Section>
<TriggerActiveMode triggerPath={path} jobTriggerKind={'kafka'} bind:active_mode />
<TriggerActiveMode triggerPath={path} jobTriggerKind={'kafka'} bind:suspended_mode />
{/if}
<KafkaTriggersConfigSection
@@ -13,7 +13,7 @@ export async function saveKafkaTriggerFromCfg(
? {
error_handler_path: cfg.error_handler_path,
error_handler_args: cfg.error_handler_path ? cfg.error_handler_args : undefined,
retry: cfg.retry,
retry: cfg.retry
}
: {}
const requestBody: EditKafkaTrigger = {
@@ -23,7 +23,7 @@ export async function saveKafkaTriggerFromCfg(
kafka_resource_path: cfg.kafka_resource_path,
group_id: cfg.group_id,
topics: cfg.topics,
active_mode: cfg.active_mode,
suspended_mode: cfg.suspended_mode,
...errorHandlerAndRetries
}
try {
@@ -97,7 +97,7 @@
let error_handler_path: string | undefined = $state()
let error_handler_args: Record<string, any> = $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 @@
</Section>
{/if}
<TriggerActiveMode triggerPath={path} jobTriggerKind={'mqtt'} bind:active_mode />
<TriggerActiveMode triggerPath={path} jobTriggerKind={'mqtt'} bind:suspended_mode />
<MqttEditorConfigSection
bind:mqtt_resource_path
@@ -13,7 +13,7 @@ export async function saveMqttTriggerFromCfg(
? {
error_handler_path: cfg.error_handler_path,
error_handler_args: cfg.error_handler_path ? cfg.error_handler_args : undefined,
retry: cfg.retry,
retry: cfg.retry
}
: {}
const requestBody: EditMqttTrigger = {
@@ -27,7 +27,7 @@ export async function saveMqttTriggerFromCfg(
script_path: cfg.script_path,
enabled: cfg.enabled,
is_flow: cfg.is_flow,
active_mode: cfg.active_mode,
suspended_mode: cfg.suspended_mode,
...errorHandlerAndRetries
}
try {
@@ -92,7 +92,7 @@
let error_handler_path: string | undefined = $state()
let error_handler_args: Record<string, any> = $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<string, any>): Promise<void> {
@@ -222,7 +222,7 @@
error_handler_path,
error_handler_args,
retry,
active_mode
suspended_mode
}
}
@@ -402,7 +402,7 @@
</Section>
{/if}
<TriggerActiveMode triggerPath={path} jobTriggerKind={'nats'} bind:active_mode />
<TriggerActiveMode triggerPath={path} jobTriggerKind={'nats'} bind:suspended_mode />
<NatsTriggersConfigSection
{path}
@@ -13,7 +13,7 @@ export async function saveNatsTriggerFromCfg(
? {
error_handler_path: cfg.error_handler_path,
error_handler_args: cfg.error_handler_path ? cfg.error_handler_args : undefined,
retry: cfg.retry,
retry: cfg.retry
}
: {}
const requestBody: EditNatsTrigger = {
@@ -25,7 +25,7 @@ export async function saveNatsTriggerFromCfg(
consumer_name: cfg.consumer_name,
subjects: cfg.subjects,
use_jetstream: cfg.use_jetstream,
active_mode: cfg.active_mode,
suspended_mode: cfg.suspended_mode,
...errorHandlerAndRetries
}
try {
@@ -117,7 +117,7 @@
let error_handler_path: string | undefined = $state()
let error_handler_args: Record<string, any> = $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<string, any>): Promise<void> {
@@ -580,7 +580,7 @@
</Section>
{/if}
<TriggerActiveMode triggerPath={path} jobTriggerKind={'postgres'} bind:active_mode />
<TriggerActiveMode triggerPath={path} jobTriggerKind={'postgres'} bind:suspended_mode />
<Section label="Database">
{#snippet badge()}
{#if isEditor}
@@ -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) {
@@ -89,7 +89,7 @@
let error_handler_path: string | undefined = $state()
let error_handler_args: Record<string, any> = $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 @@
</div>
</Section>
<TriggerActiveMode triggerPath={path} jobTriggerKind={'sqs'} bind:active_mode />
<TriggerActiveMode triggerPath={path} jobTriggerKind={'sqs'} bind:suspended_mode />
{/if}
<SqsTriggerEditorConfigSection
@@ -13,7 +13,7 @@ export async function saveSqsTriggerFromCfg(
? {
error_handler_path: cfg.error_handler_path,
error_handler_args: cfg.error_handler_path ? cfg.error_handler_args : undefined,
retry: cfg.retry,
retry: cfg.retry
}
: {}
const requestBody: EditSqsTrigger = {
@@ -25,7 +25,7 @@ export async function saveSqsTriggerFromCfg(
message_attributes: cfg.message_attributes,
aws_auth_resource_type: cfg.aws_auth_resource_type,
enabled: cfg.enabled,
active_mode: cfg.active_mode,
suspended_mode: cfg.suspended_mode,
...errorHandlerAndRetries
}
try {
@@ -106,7 +106,7 @@
let error_handler_path: string | undefined = $state()
let error_handler_args: Record<string, any> = $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 @@
/>
</Section>
<TriggerActiveMode triggerPath={path} jobTriggerKind={'websocket'} bind:active_mode />
<TriggerActiveMode triggerPath={path} jobTriggerKind={'websocket'} bind:suspended_mode />
<WebsocketEditorConfigSection
bind:url
@@ -16,7 +16,7 @@ export async function saveWebsocketTriggerFromCfg(
error_handler_args: triggerCfg.error_handler_path
? triggerCfg.error_handler_args
: undefined,
retry: triggerCfg.retry,
retry: triggerCfg.retry
}
: {}
const requestBody: EditWebsocketTrigger = {
@@ -29,7 +29,7 @@ export async function saveWebsocketTriggerFromCfg(
url_runnable_args: triggerCfg.url_runnable_args,
can_return_message: triggerCfg.can_return_message,
can_return_error_result: triggerCfg.can_return_error_result,
active_mode: triggerCfg.active_mode,
suspended_mode: triggerCfg.suspended_mode,
...errorHandlerAndRetries
}
try {