This commit is contained in:
dieriba
2025-11-12 17:17:31 +01:00
parent 727c55a397
commit 7dc02ab4d2
19 changed files with 420 additions and 84 deletions
+57
View File
@@ -8348,6 +8348,63 @@ paths:
items:
type: string
/w/{workspace}/jobs/queue/resume_suspended:
post:
summary: resume all suspended jobs with the given suspend number
operationId: resumeSuspendedJobs
tags:
- job
parameters:
- $ref: "#/components/parameters/WorkspaceId"
requestBody:
description: suspend number of the jobs to resume
required: true
content:
application/json:
schema:
type: object
properties:
suspend_number:
type: integer
description: the suspend number of jobs to resume
required:
- suspend_number
responses:
"200":
content:
application/json:
schema:
type: string
/w/{workspace}/jobs/queue/cancel_suspended:
post:
summary: cancel all suspended jobs with the given suspend number
operationId: cancelSuspendedJobs
tags:
- job
parameters:
- $ref: "#/components/parameters/WorkspaceId"
requestBody:
description: suspend number of the jobs to cancel
required: true
content:
application/json:
schema:
type: object
properties:
suspend_number:
type: integer
description: the suspend number of jobs to cancel
required:
- suspend_number
responses:
"200":
description: confirmation message
content:
application/json:
schema:
type: string
/w/{workspace}/jobs/completed/list:
get:
summary: list all completed jobs
+176
View File
@@ -15,6 +15,7 @@ use futures::{StreamExt, TryFutureExt};
use http::{HeaderMap, HeaderName};
use itertools::Itertools;
use quick_cache::sync::Cache;
use rand::SeedableRng;
use serde_json::value::RawValue;
use serde_json::Value;
use sqlx::Pool;
@@ -106,6 +107,56 @@ use windmill_queue::{
use crate::flow_conversations;
use windmill_common::flow_conversations::MessageType;
pub async fn generate_unique_suspend_number(
db: &Pool<sqlx::Postgres>,
workspace_id: &str,
) -> error::Result<i32> {
use rand::Rng;
let mut rng = rand::rngs::StdRng::from_os_rng();
let mut attempts = 0;
let max_attempts = 10;
loop {
attempts += 1;
if attempts > max_attempts {
return Err(error::Error::InternalErr(format!(
"Failed to generate unique suspend number after {} attempts",
max_attempts
)));
}
let suspend_number = rng.random_range(101..i32::MAX);
let exists = sqlx::query_scalar!(
"SELECT EXISTS(SELECT 1 FROM v2_job_queue WHERE workspace_id = $1 AND suspend = $2)",
workspace_id,
suspend_number
)
.fetch_one(db)
.await?
.unwrap_or(false);
if !exists {
return Ok(suspend_number);
}
}
}
pub async fn get_suspend_number_for_unactive_mode(
db: &Pool<sqlx::Postgres>,
workspace_id: &str,
active_mode: Option<bool>,
) -> error::Result<Option<i32>> {
if let Some(false) = active_mode {
Ok(Some(
generate_unique_suspend_number(db, workspace_id).await?,
))
} else {
Ok(None)
}
}
pub fn workspaced_service() -> Router {
let cors = CorsLayer::new()
.allow_methods([http::Method::GET, http::Method::POST])
@@ -237,6 +288,8 @@ pub fn workspaced_service() -> Router {
.route("/queue/position/:timestamp", get(get_queue_position))
.route("/queue/scheduled_for/:id", get(get_scheduled_for))
.route("/queue/cancel_selection", post(cancel_selection))
.route("/queue/resume_suspended", post(resume_suspended_jobs))
.route("/queue/cancel_suspended", post(cancel_suspended_jobs))
.route("/completed/count", get(count_completed_jobs))
.route("/completed/count_jobs", get(count_completed_jobs_detail))
.route(
@@ -2176,6 +2229,129 @@ async fn cancel_selection(
.await
}
#[derive(Deserialize)]
pub struct ResumeSuspendedJobsRequest {
suspend_number: i32,
}
async fn resume_suspended_jobs(
authed: ApiAuthed,
Extension(user_db): Extension<UserDB>,
Path(w_id): Path<String>,
Json(request): Json<ResumeSuspendedJobsRequest>,
) -> error::JsonResult<String> {
require_admin(authed.is_admin, &authed.username)?;
let mut tx = user_db.begin(&authed).await?;
sqlx::query!(
r#"
WITH jobs_to_resume AS (
SELECT
jq.id
FROM
v2_job_queue jq
INNER JOIN v2_job j ON j.id = jq.id
WHERE
j.workspace_id = $1 AND
jq.suspend = $2 AND
jq.canceled_by IS NULL
)
UPDATE
v2_job_queue
SET
suspend = 0,
scheduled_for = NOW()
WHERE
id IN (SELECT id FROM jobs_to_resume)
"#,
w_id,
request.suspend_number
)
.fetch_all(&mut *tx)
.await?;
audit_log(
&mut *tx,
&authed,
"jobs.resume_suspended",
ActionKind::Update,
&w_id,
Some(&format!("suspend_number:{}", request.suspend_number)),
None,
)
.await?;
tx.commit().await?;
Ok(Json(format!(
"Resumed all suspended workspace jobs for suspend number: {}",
request.suspend_number
)))
}
#[derive(Deserialize)]
pub struct CancelSuspendedJobsRequest {
suspend_number: i32,
}
async fn cancel_suspended_jobs(
authed: ApiAuthed,
Extension(user_db): Extension<UserDB>,
Path(w_id): Path<String>,
Json(request): Json<CancelSuspendedJobsRequest>,
) -> error::JsonResult<String> {
require_admin(authed.is_admin, &authed.username)?;
let mut tx = user_db.begin(&authed).await?;
sqlx::query!(
r#"
WITH jobs_to_cancel AS (
SELECT
jq.id
FROM
v2_job_queue jq
INNER JOIN v2_job j ON j.id = jq.id
WHERE
j.workspace_id = $1 AND
jq.suspend = $2 AND
jq.canceled_by IS NULL
)
UPDATE
v2_job_queue
SET
canceled_by = $3,
canceled_reason = 'Canceled all suspended jobs with suspend number'
WHERE
id IN (SELECT id FROM jobs_to_cancel)
"#,
w_id,
request.suspend_number,
authed.username
)
.fetch_all(&mut *tx)
.await?;
audit_log(
&mut *tx,
&authed,
"jobs.cancel_suspended",
ActionKind::Delete,
&w_id,
Some(&format!("suspend_number:{}", request.suspend_number)),
None,
)
.await?;
tx.commit().await?;
Ok(Json(format!(
"Canceled all suspended workspace jobs for suspend number: {}",
request.suspend_number
)))
}
async fn list_filtered_job_uuids(
authed: ApiAuthed,
Extension(db): Extension<DB>,
@@ -46,6 +46,7 @@ impl TriggerCrud for EmailTrigger {
_authed: &ApiAuthed,
_w_id: &str,
_trigger: TriggerData<Self::TriggerConfigRequest>,
_suspend_number: Option<i32>,
) -> Result<()> {
Err(Error::BadRequest(
"Email triggers are not available in open source version".to_string(),
@@ -60,6 +61,7 @@ impl TriggerCrud for EmailTrigger {
_workspace_id: &str,
_path: &str,
_trigger: TriggerData<Self::TriggerConfigRequest>,
_suspend_number: Option<i32>,
) -> Result<()> {
Err(Error::BadRequest(
"Email triggers are not available in open source version".to_string(),
+22 -5
View File
@@ -1,8 +1,7 @@
use crate::{
db::ApiAuthed,
triggers::{
trigger_helpers::get_suspend_number_for_unactive_mode, StandardTriggerQuery, TriggerData,
},
jobs::generate_unique_suspend_number,
triggers::{StandardTriggerQuery, TriggerData},
};
use async_trait::async_trait;
use serde::{de::DeserializeOwned, Deserialize, Serialize};
@@ -365,6 +364,20 @@ pub fn trigger_routes<T: TriggerCrud + 'static>() -> Router {
router
}
pub async fn get_suspend_number_for_inactive_mode(
db: &DB,
workspace_id: &str,
active_mode: Option<bool>,
) -> Result<Option<i32>> {
if let Some(false) = active_mode {
Ok(Some(
generate_unique_suspend_number(db, workspace_id).await?,
))
} else {
Ok(None)
}
}
async fn create_trigger<T: TriggerCrud>(
Extension(handler): Extension<Arc<T>>,
authed: ApiAuthed,
@@ -395,7 +408,9 @@ async fn create_trigger<T: TriggerCrud>(
let mut tx = user_db.begin(&authed).await?;
let new_path = new_trigger.base.path.clone();
let suspend_number = get_suspend_number_for_unactive_mode(new_trigger.base.active_mode);
let suspend_number =
get_suspend_number_for_inactive_mode(&db, &workspace_id, new_trigger.base.active_mode)
.await?;
handler
.create_trigger(
@@ -496,7 +511,9 @@ async fn update_trigger<T: TriggerCrud>(
let mut tx = user_db.begin(&authed).await?;
let new_path = edit_trigger.base.path.to_string();
let suspend_number = get_suspend_number_for_unactive_mode(edit_trigger.base.active_mode);
let suspend_number =
get_suspend_number_for_inactive_mode(&db, &workspace_id, edit_trigger.base.active_mode)
.await?;
handler
.update_trigger(
@@ -8,13 +8,14 @@ use crate::{
jobs::start_job_update_sse_stream,
resources::try_get_resource_from_db_as,
triggers::{
handler::get_suspend_number_for_inactive_mode,
http::{
refresh_routers, validate_authentication_method, HttpConfig, HttpConfigRequest,
RouteExists, ROUTE_PATH_KEY_RE, VALID_ROUTE_PATH_RE,
},
trigger_helpers::{
get_runnable_format, get_suspend_number_for_unactive_mode, trigger_runnable,
trigger_runnable_and_wait_for_result, trigger_runnable_inner, RunnableId,
get_runnable_format, trigger_runnable, trigger_runnable_and_wait_for_result,
trigger_runnable_inner, RunnableId,
},
Trigger, TriggerCrud, TriggerData,
},
@@ -283,7 +284,8 @@ pub async fn create_many_http_triggers(
for (new_http_trigger, route_path_key) in new_http_triggers.iter().zip(route_path_keys.iter()) {
let suspend_number =
get_suspend_number_for_unactive_mode(new_http_trigger.base.active_mode);
get_suspend_number_for_inactive_mode(&db, &w_id, new_http_trigger.base.active_mode)
.await?;
insert_new_trigger_into_db(
&authed,
&mut tx,
@@ -1052,6 +1054,35 @@ async fn route_job(
)
.map_err(|e| e.into_response())?;
if let Some(suspend_number) = trigger.suspend_number {
let _ = trigger_runnable(
&db,
Some(user_db),
authed,
&trigger.workspace_id,
&trigger.script_path,
trigger.is_flow,
args,
trigger.retry.as_ref(),
trigger.error_handler_path.as_deref(),
trigger.error_handler_args.as_ref(),
format!("http_trigger/{}", trigger.path),
None,
Some(suspend_number),
)
.await
.map_err(|e| e.into_response())?;
return Ok((
StatusCode::OK,
format!(
"Trigger: {} in inactive mode, incoming request has been queued",
&trigger.path
),
)
.into_response());
}
// Handle execution based on the execution mode
match trigger.request_type {
RequestType::SyncSse => {
+2 -1
View File
@@ -52,7 +52,7 @@ pub struct BaseTrigger {
pub email: String,
pub edited_at: DateTime<Utc>,
pub extra_perms: Option<serde_json::Value>,
pub suspend_number: Option<i32>
pub suspend_number: Option<i32>,
}
#[derive(Debug, FromRow, Clone, Serialize, Deserialize)]
@@ -117,6 +117,7 @@ pub struct BaseTriggerData {
pub is_flow: bool,
pub enabled: Option<bool>,
pub active_mode: Option<bool>,
pub process_queued_jobs: Option<bool>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
@@ -903,18 +903,6 @@ async fn trigger_script_with_retry_and_error_handler(
Ok((uuid, delete_after_use))
}
pub fn generate_trigger_suspend_number() -> i32 {
use rand::Rng;
let mut rng = rand::rng();
rng.random_range(101..i32::MAX)
}
pub fn get_suspend_number_for_unactive_mode(active_mode: Option<bool>) -> Option<i32> {
if let Some(false) = active_mode {
return Some(generate_trigger_suspend_number());
}
None
}
async fn trigger_runnable_with_suspend(
db: &DB,
@@ -365,7 +365,25 @@ impl Listener for WebsocketTrigger {
),
None => (None, None, None),
};
if let Some(ReturnMessageChannels { send_message_tx, mut killpill_rx }) = extra {
if suspend_number.is_some() || extra.is_none() {
trigger_runnable(
db,
None,
authed,
&workspace_id,
&script_path,
*is_flow,
args,
retry,
error_handler_path,
error_handler_args,
format!("websocket_trigger/{}", listening_trigger.path),
None,
*suspend_number,
)
.await?;
} else if let Some(ReturnMessageChannels { send_message_tx, mut killpill_rx }) = extra {
let db_ = db.clone();
let url = url.to_owned();
let script_path = script_path.to_owned();
@@ -415,23 +433,6 @@ impl Listener for WebsocketTrigger {
};
tokio::spawn(handle_response_f);
} else {
trigger_runnable(
db,
None,
authed,
&workspace_id,
&script_path,
*is_flow,
args,
retry,
error_handler_path,
error_handler_args,
format!("websocket_trigger/{}", listening_trigger.path),
None,
*suspend_number
)
.await?;
}
Ok(())
@@ -50,7 +50,7 @@
{#if isOpen}
<Portal name="always-mounted" {target}>
<div
class={'fixed top-0 bottom-0 left-0 right-0 transition-all overflow-auto z-[1100] bg-black bg-opacity-60 w-full h-full'}
class={'fixed top-0 bottom-0 left-0 right-0 transition-all overflow-auto z-[9999] bg-black bg-opacity-60 w-full h-full'}
transition:fadeFast|local
>
<div class="flex min-h-full items-center justify-center p-8">
@@ -7,23 +7,23 @@
import { workspaceStore } from '$lib/stores'
import RunRow from '../runs/RunRow.svelte'
import '../runs/runs-grid.css'
import { sendUserToast } from '$lib/toast'
type Props = {
active_mode: boolean
suspendNumber?: number
triggerPath?: string
onRunSuspendedJobs?: (data: { suspendNumber: number; jobIds: string[] }) => void
}
let { active_mode = $bindable(), suspendNumber, onRunSuspendedJobs }: Props = $props()
let { active_mode = $bindable(), suspendNumber }: Props = $props()
let wasInUnactiveMode = $state(!active_mode)
let wasInInactiveMode = $state(!active_mode)
let shouldShowModal = $state(false)
$effect(() => {
if (active_mode && wasInUnactiveMode && suspendNumber !== undefined) {
if (active_mode && wasInInactiveMode && suspendNumber !== undefined) {
shouldShowModal = true
} else if (!active_mode) {
wasInUnactiveMode = true
wasInInactiveMode = true
shouldShowModal = false
}
})
@@ -32,7 +32,6 @@
let error = $state<string | null>(null)
let processingAction = $state(false)
let workspace = $workspaceStore!
let selectedJobIds = $state<string[]>([])
let containerWidth = $state(1000)
$effect(() => {
if (shouldShowModal && suspendNumber !== undefined) {
@@ -63,59 +62,60 @@
}
async function runAllJobs() {
if (queuedJobs.length === 0) return
if (queuedJobs.length === 0 || !suspendNumber) return
processingAction = true
error = null
try {
if (onRunSuspendedJobs && suspendNumber !== undefined) {
onRunSuspendedJobs({ suspendNumber, jobIds: queuedJobs.map((j) => j.id) })
}
closeModal()
const resumedJobs = await JobService.resumeSuspendedJobs({
workspace,
requestBody: {
suspend_number: suspendNumber
}
})
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
if (queuedJobs.length === 0 || !suspendNumber) return
processingAction = true
error = null
try {
const jobIds = queuedJobs.map((job) => job.id)
await JobService.cancelSelection({
await JobService.cancelSuspendedJobs({
workspace,
requestBody: jobIds
requestBody: {
suspend_number: suspendNumber
}
})
closeModal()
sendUserToast(`Successfully canceled all jobs with suspend number: ${suspendNumber}`)
} catch (e) {
error = `Failed to discard jobs: ${e}`
console.error('Failed to discard jobs:', e)
} finally {
processingAction = false
closeModal()
}
}
function closeModal() {
wasInUnactiveMode = false
wasInInactiveMode = false
shouldShowModal = false
}
function cancelToggle() {
active_mode = false
closeModal()
}
</script>
<Toggle bind:checked={active_mode} options={{ right: 'Active', left: 'Unactive' }} />
<Toggle bind:checked={active_mode} options={{ right: 'Active', left: 'Inactive' }} />
{#if shouldShowModal}
<Modal2
@@ -124,7 +124,7 @@
target="#content"
fixedSize="lg"
>
<div class="flex flex-col gap-4 h-full">
<div class="flex w-full flex-col gap-4 h-full">
{#if loading}
<div class="flex items-center justify-center py-8">
<div
@@ -137,17 +137,22 @@
{error}
</div>
{:else if queuedJobs.length === 0}
<div class="text-center py-8 text-gray-500"> No queued jobs found for this trigger. </div>
<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 queued jobs found</div>
<div class="text-sm text-tertiary"
>This trigger has no 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 text-gray-900">Queued Jobs ({queuedJobs.length})</h3>
<h3 class="text-sm font-medium">Queued Jobs ({queuedJobs.length})</h3>
<p class="text-xs text-gray-500 mt-1">Click on any job to view details</p>
</div>
<!-- Job table container - matching runs page styling -->
<div class="divide-y h-full border min-w-[650px]" bind:clientWidth={containerWidth}>
<!-- Table header - using same bg-surface-secondary as runs page -->
<div
class="bg-surface-secondary sticky top-0 w-full py-2 pr-4 grid grid-runs-table-no-tag"
>
@@ -159,7 +164,6 @@
<div class=""></div>
</div>
<!-- Job rows - no background, let RunRow handle its own styling -->
<div class="h-full">
{#each queuedJobs as job}
<div class="flex flex-row items-center h-[42px] w-full">
@@ -169,7 +173,6 @@
showTag={false}
activeLabel={null}
on:select={() => {
// Handle job selection - navigate to job details or show modal
window.open(`/run/${job.id}?workspace=${workspace}`, '_blank')
}}
/>
@@ -188,10 +191,6 @@
{/if}
<div class="flex gap-2 pt-4 border-t">
<Button variant="border" size="sm" on:click={cancelToggle} disabled={processingAction}>
Cancel
</Button>
{#if !loading && !error && queuedJobs.length > 0}
<Button
variant="border"
@@ -24,6 +24,7 @@
import Tabs from '$lib/components/common/tabs/Tabs.svelte'
import Tab from '$lib/components/common/tabs/Tab.svelte'
import TriggerRetriesAndErrorHandler from '../TriggerRetriesAndErrorHandler.svelte'
import TriggerStateToggle from '../TriggerStateToggle.svelte'
import { saveEmailTriggerFromCfg } from './utils'
let {
@@ -65,6 +66,8 @@
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 suspend_number: number | undefined = $state()
// Component references
let drawer = $state<Drawer | undefined>(undefined)
let initialConfig: NewEmailTrigger | undefined = undefined
@@ -161,6 +164,8 @@
error_handler_args = cfg?.error_handler_args ?? {}
retry = cfg?.retry
errorHandlerSelected = getHandlerType(error_handler_path ?? '')
active_mode = cfg?.suspend_number ? false : true
suspend_number = cfg?.suspend_number
}
async function loadTrigger(defaultConfig?: Partial<EmailTrigger>): Promise<void> {
@@ -210,7 +215,8 @@
extra_perms: extraPerms,
error_handler_path,
error_handler_args,
retry
retry,
active_mode
}
return nCfg
@@ -296,6 +302,8 @@
</Section>
{/if}
<TriggerStateToggle suspendNumber={suspend_number} bind:active_mode />
<EmailTriggerEditorConfigSection
initialTriggerPath={initialPath}
bind:local_part
@@ -62,7 +62,7 @@
let error_handler_args: Record<string, any> = $state({})
let retry: Retry | undefined = $state()
let active_mode = $state(true)
let suspend_number: number | undefined
let suspend_number: number | undefined = $state(undefined)
let {
useDrawer = true,
description = undefined,
@@ -45,6 +45,7 @@
import Tabs from '$lib/components/common/tabs/Tabs.svelte'
import Tab from '$lib/components/common/tabs/Tab.svelte'
import TriggerRetriesAndErrorHandler from '../TriggerRetriesAndErrorHandler.svelte'
import TriggerStateToggle from '../TriggerStateToggle.svelte'
let {
useDrawer = true,
@@ -110,6 +111,8 @@
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 suspend_number: number | undefined = $state(undefined)
const isAdmin = $derived($userStore?.is_admin || $userStore?.is_super_admin)
const routeConfig = $derived.by(getRouteConfig)
const captureConfig = $derived.by(isEditor ? getCaptureConfig : () => ({}))
@@ -288,6 +291,8 @@
error_handler_args = cfg?.error_handler_args ?? {}
retry = cfg?.retry
errorHandlerSelected = getHandlerType(error_handler_path ?? '')
active_mode = cfg?.suspend_number ? false : true
suspend_number = cfg?.suspend_number
}
async function loadTrigger(defaultConfig?: Partial<HttpTrigger>): Promise<void> {
@@ -355,7 +360,8 @@
description: routeDescription,
error_handler_path,
error_handler_args,
retry
retry,
active_mode
}
return nCfg
@@ -590,6 +596,8 @@
{/if}
</div>
</Section>
<TriggerStateToggle suspendNumber={suspend_number} bind:active_mode />
{/if}
<RouteEditorConfigSection
@@ -19,6 +19,7 @@
import Tabs from '$lib/components/common/tabs/Tabs.svelte'
import Tab from '$lib/components/common/tabs/Tab.svelte'
import TriggerRetriesAndErrorHandler from '../TriggerRetriesAndErrorHandler.svelte'
import TriggerStateToggle from '../TriggerStateToggle.svelte'
interface Props {
useDrawer?: boolean
@@ -82,6 +83,8 @@
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 suspend_number: number | undefined = $state(undefined)
const isValid = $derived(
!!kafkaResourcePath &&
@@ -188,6 +191,8 @@
error_handler_args = cfg?.error_handler_args ?? {}
retry = cfg?.retry
errorHandlerSelected = getHandlerType(error_handler_path ?? '')
active_mode = cfg?.suspend_number ? false : true
suspend_number = cfg?.suspend_number
}
async function loadTrigger(defaultConfig?: Record<string, any>): Promise<void> {
@@ -215,7 +220,8 @@
extra_perms: extra_perms,
error_handler_path,
error_handler_args,
retry
retry,
active_mode
}
}
@@ -382,6 +388,8 @@
{/if}
</div>
</Section>
<TriggerStateToggle suspendNumber={suspend_number} bind:active_mode />
{/if}
<KafkaTriggersConfigSection
@@ -27,6 +27,7 @@
import Tabs from '$lib/components/common/tabs/Tabs.svelte'
import Tab from '$lib/components/common/tabs/Tab.svelte'
import TriggerRetriesAndErrorHandler from '../TriggerRetriesAndErrorHandler.svelte'
import TriggerStateToggle from '../TriggerStateToggle.svelte'
import Toggle from '$lib/components/Toggle.svelte'
import ToggleButtonGroup from '$lib/components/common/toggleButton-v2/ToggleButtonGroup.svelte'
import ToggleButton from '$lib/components/common/toggleButton-v2/ToggleButton.svelte'
@@ -96,6 +97,8 @@
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 suspend_number: number | undefined = $state(undefined)
let optionTabSelected: 'connection_options' | 'error_handler' | 'retries' =
$state('connection_options')
@@ -202,6 +205,8 @@
error_handler_args = cfg?.error_handler_args ?? {}
retry = cfg?.retry
errorHandlerSelected = getHandlerType(error_handler_path ?? '')
active_mode = cfg?.suspend_number ? false : true
suspend_number = cfg?.suspend_number
activateV5Options.topic_alias_maximum = Boolean(v5_config.topic_alias_maximum)
activateV5Options.session_expiry_interval = Boolean(v5_config.session_expiry_interval)
} catch (error) {
@@ -240,7 +245,8 @@
is_flow,
error_handler_path,
error_handler_args,
retry
retry,
active_mode
}
}
@@ -417,6 +423,8 @@
</Section>
{/if}
<TriggerStateToggle suspendNumber={suspend_number} bind:active_mode />
<MqttEditorConfigSection
bind:mqtt_resource_path
bind:subscribe_topics
@@ -19,6 +19,7 @@
import Tabs from '$lib/components/common/tabs/Tabs.svelte'
import Tab from '$lib/components/common/tabs/Tab.svelte'
import TriggerRetriesAndErrorHandler from '../TriggerRetriesAndErrorHandler.svelte'
import TriggerStateToggle from '../TriggerStateToggle.svelte'
interface Props {
useDrawer?: boolean
@@ -91,6 +92,8 @@
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 suspend_number: number | undefined = $state(undefined)
const saveDisabled = $derived(
pathError != '' || emptyString(script_path) || !can_write || !isValid
@@ -190,6 +193,8 @@
error_handler_args = cfg?.error_handler_args ?? {}
retry = cfg?.retry
errorHandlerSelected = getHandlerType(error_handler_path ?? '')
active_mode = cfg?.suspend_number ? false : true
suspend_number = cfg?.suspend_number
}
async function loadTrigger(defaultConfig?: Record<string, any>): Promise<void> {
@@ -218,7 +223,8 @@
use_jetstream: natsCfg.use_jetstream,
error_handler_path,
error_handler_args,
retry
retry,
active_mode
}
}
@@ -398,6 +404,8 @@
</Section>
{/if}
<TriggerStateToggle suspendNumber={suspend_number} bind:active_mode />
<NatsTriggersConfigSection
{path}
bind:natsResourcePath
@@ -35,6 +35,7 @@
import TestingBadge from '../testingBadge.svelte'
import { getHandlerType, handleConfigChange, type Trigger } from '../utils'
import TriggerRetriesAndErrorHandler from '../TriggerRetriesAndErrorHandler.svelte'
import TriggerStateToggle from '../TriggerStateToggle.svelte'
import { fade } from 'svelte/transition'
import MultiSelect from '$lib/components/select/MultiSelect.svelte'
import { safeSelectItems } from '$lib/components/select/utils.svelte'
@@ -116,6 +117,8 @@
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 suspend_number: number | undefined = $state(undefined)
const errorMessage = $derived.by(() => {
if (relations && relations.length > 0) {
@@ -299,7 +302,8 @@
: undefined,
error_handler_path,
error_handler_args,
retry
retry,
active_mode
}
return cfg
}
@@ -320,6 +324,8 @@
error_handler_args = cfg?.error_handler_args ?? {}
retry = cfg?.retry
errorHandlerSelected = getHandlerType(error_handler_path ?? '')
active_mode = cfg?.suspend_number ? false : true
suspend_number = cfg?.suspend_number
}
async function loadTrigger(defaultConfig?: Record<string, any>): Promise<void> {
@@ -575,6 +581,8 @@
</div>
</Section>
{/if}
<TriggerStateToggle suspendNumber={suspend_number} bind:active_mode />
<Section label="Database">
{#snippet badge()}
{#if isEditor}
@@ -24,6 +24,7 @@
import Tabs from '$lib/components/common/tabs/Tabs.svelte'
import Tab from '$lib/components/common/tabs/Tab.svelte'
import TriggerRetriesAndErrorHandler from '../TriggerRetriesAndErrorHandler.svelte'
import TriggerStateToggle from '../TriggerStateToggle.svelte'
interface Props {
useDrawer?: boolean
@@ -88,6 +89,8 @@
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 suspend_number: number | undefined = $state(undefined)
const sqsConfig = $derived.by(getSaveCfg)
const captureConfig = $derived.by(getCaptureConfig)
@@ -176,6 +179,8 @@
error_handler_args = cfg?.error_handler_args ?? {}
retry = cfg?.retry
errorHandlerSelected = getHandlerType(error_handler_path ?? '')
active_mode = cfg?.suspend_number ? false : true
suspend_number = cfg?.suspend_number
} catch (error) {
sendUserToast(`Could not load SQS trigger config: ${error.body}`, true)
}
@@ -210,7 +215,8 @@
enabled,
error_handler_path,
error_handler_args,
retry
retry,
active_mode
}
}
@@ -383,6 +389,8 @@
{/if}
</div>
</Section>
<TriggerStateToggle suspendNumber={suspend_number} bind:active_mode />
{/if}
<SqsTriggerEditorConfigSection
@@ -34,6 +34,7 @@
import Tabs from '$lib/components/common/tabs/Tabs.svelte'
import Tab from '$lib/components/common/tabs/Tab.svelte'
import TriggerRetriesAndErrorHandler from '../TriggerRetriesAndErrorHandler.svelte'
import TriggerStateToggle from '../TriggerStateToggle.svelte'
interface Props {
useDrawer?: boolean
@@ -105,6 +106,8 @@
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 suspend_number: number | undefined = $state(undefined)
const websocketCfg = $derived.by(getSaveCfg)
const captureConfig = $derived.by(isEditor ? getCaptureConfig : () => ({}))
@@ -219,6 +222,8 @@
error_handler_args = cfg?.error_handler_args ?? {}
retry = cfg?.retry
errorHandlerSelected = getHandlerType(error_handler_path ?? '')
active_mode = cfg?.suspend_number ? false : true
suspend_number = cfg?.suspend_number
}
function getSaveCfg() {
@@ -236,7 +241,8 @@
enabled,
error_handler_path,
error_handler_args,
retry
retry,
active_mode
}
}
@@ -490,6 +496,8 @@
/>
</Section>
<TriggerStateToggle suspendNumber={suspend_number} bind:active_mode />
<WebsocketEditorConfigSection
bind:url
bind:url_runnable_args