From cc1d3f219320127d286cc6e2358b510a88bd43dd Mon Sep 17 00:00:00 2001 From: Guilhem Date: Mon, 4 Nov 2024 23:18:09 +0100 Subject: [PATCH] feat(frontend): improve and simplify scheduled poll flows (#4560) --- backend/tests/worker.rs | 8 + backend/windmill-api/src/flows.rs | 3 + backend/windmill-common/src/flows.rs | 7 + backend/windmill-queue/src/jobs.rs | 1 + .../windmill-worker/src/worker_lockfiles.rs | 3 + frontend/package.json | 5 + frontend/src/lib/components/Dev.svelte | 5 +- .../src/lib/components/FlowBuilder.svelte | 17 +- .../components/FlowStatusViewerInner.svelte | 1 + .../lib/components/RunPageSchedules.svelte | 164 +++++--------- .../src/lib/components/ScriptBuilder.svelte | 6 +- frontend/src/lib/components/Toggle.svelte | 6 +- .../lib/components/copilot/MetadataGen.svelte | 2 +- .../details/DetailPageDetailPanel.svelte | 6 +- .../details/DetailPageLayout.svelte | 12 +- .../details/DetailPageTriggerPanel.svelte | 116 +++++----- .../flows/content/FlowPathViewer.svelte | 9 +- .../lib/components/flows/flowStateUtils.ts | 13 +- .../flows/map/FlowModuleSchemaItem.svelte | 11 +- .../flows/map/FlowModuleSchemaMap.svelte | 22 +- .../flows/map/InsertModuleButton.svelte | 22 +- .../lib/components/flows/map/MapItem.svelte | 3 + .../components/flows/map/VirtualItem.svelte | 47 +--- .../flows/map/VirtualItemWrapper.svelte | 45 ++++ .../src/lib/components/flows/scheduleUtils.ts | 110 +++++++++- .../lib/components/graph/FlowGraphV2.svelte | 207 ++++++++++++++---- .../src/lib/components/graph/graphBuilder.ts | 96 +++++--- .../graph/renderers/edges/BaseEdge.svelte | 39 ---- .../renderers/nodes/ForLoopEndNode.svelte | 38 +++- .../renderers/nodes/ForLoopStartNode.svelte | 3 +- .../graph/renderers/nodes/InputNode.svelte | 3 +- .../graph/renderers/nodes/TriggersNode.svelte | 113 +++++++++- .../renderers/triggers/TriggerButton.svelte | 2 +- .../renderers/triggers/TriggersBadge.svelte | 107 +++------ .../renderers/triggers/TriggersWrapper.svelte | 54 +++-- .../components/icons/SchedulePollIcon.svelte | 36 +++ frontend/src/lib/components/icons/index.ts | 2 + frontend/src/lib/components/triggers.ts | 26 ++- .../triggers/ScheduledPollPanel.svelte | 41 ++++ .../components/triggers/TriggersEditor.svelte | 128 ++++++----- frontend/src/routes/flows/dev/+page.svelte | 5 +- frontend/static/create_action.png | Bin 0 -> 40254 bytes frontend/static/script-picker.png | Bin 0 -> 12652 bytes frontend/static/trigger_button.png | Bin 17525 -> 0 bytes frontend/tailwind.config.cjs | 3 +- openflow.openapi.yaml | 4 + 46 files changed, 1011 insertions(+), 540 deletions(-) create mode 100644 frontend/src/lib/components/flows/map/VirtualItemWrapper.svelte create mode 100644 frontend/src/lib/components/icons/SchedulePollIcon.svelte create mode 100644 frontend/src/lib/components/triggers/ScheduledPollPanel.svelte create mode 100644 frontend/static/create_action.png create mode 100644 frontend/static/script-picker.png delete mode 100644 frontend/static/trigger_button.png diff --git a/backend/tests/worker.rs b/backend/tests/worker.rs index a93c6eef1d..69482f3922 100644 --- a/backend/tests/worker.rs +++ b/backend/tests/worker.rs @@ -1112,6 +1112,7 @@ async fn test_deno_flow(db: Pool) { custom_concurrency_key: None, concurrent_limit: None, concurrency_time_window_s: None, + is_trigger: None, } .into(), stop_after_if: Default::default(), @@ -1153,6 +1154,7 @@ async fn test_deno_flow(db: Pool) { custom_concurrency_key: None, concurrent_limit: None, concurrency_time_window_s: None, + is_trigger: None, } .into(), stop_after_if: Default::default(), @@ -1276,6 +1278,8 @@ async fn test_deno_flow_same_worker(db: Pool) { custom_concurrency_key: None, concurrent_limit: None, concurrency_time_window_s: None, + is_trigger: None, + }.into(), stop_after_if: Default::default(), stop_after_all_iters_if: Default::default(), @@ -1327,6 +1331,7 @@ async fn test_deno_flow_same_worker(db: Pool) { custom_concurrency_key: None, concurrent_limit: None, concurrency_time_window_s: None, + is_trigger: None, }.into(), stop_after_if: Default::default(), stop_after_all_iters_if: Default::default(), @@ -1364,6 +1369,8 @@ async fn test_deno_flow_same_worker(db: Pool) { custom_concurrency_key: None, concurrent_limit: None, concurrency_time_window_s: None, + is_trigger: None, + }.into(), stop_after_if: Default::default(), stop_after_all_iters_if: Default::default(), @@ -1424,6 +1431,7 @@ async fn test_deno_flow_same_worker(db: Pool) { custom_concurrency_key: None, concurrent_limit: None, concurrency_time_window_s: None, + is_trigger: None, }.into(), stop_after_if: Default::default(), stop_after_all_iters_if: Default::default(), diff --git a/backend/windmill-api/src/flows.rs b/backend/windmill-api/src/flows.rs index 4fc9ce514f..f02f91dc41 100644 --- a/backend/windmill-api/src/flows.rs +++ b/backend/windmill-api/src/flows.rs @@ -1209,6 +1209,7 @@ mod tests { .into(), hash: None, tag_override: None, + is_trigger: None, }), stop_after_if: None, stop_after_all_iters_if: None, @@ -1236,6 +1237,7 @@ mod tests { custom_concurrency_key: None, concurrent_limit: None, concurrency_time_window_s: None, + is_trigger: None, }), stop_after_if: Some(StopAfterIf { expr: "foo = 'bar'".to_string(), @@ -1290,6 +1292,7 @@ mod tests { input_transforms: HashMap::new(), hash: None, tag_override: None, + is_trigger: None, } .into(), stop_after_if: Some(StopAfterIf { diff --git a/backend/windmill-common/src/flows.rs b/backend/windmill-common/src/flows.rs index 9b66cc65ac..d357a5ec72 100644 --- a/backend/windmill-common/src/flows.rs +++ b/backend/windmill-common/src/flows.rs @@ -426,6 +426,8 @@ pub enum FlowModuleValue { hash: Option, #[serde(skip_serializing_if = "Option::is_none")] tag_override: Option, + #[serde(skip_serializing_if = "Option::is_none")] + is_trigger: Option, }, Flow { #[serde(default)] @@ -474,6 +476,8 @@ pub enum FlowModuleValue { concurrent_limit: Option, #[serde(skip_serializing_if = "Option::is_none")] concurrency_time_window_s: Option, + #[serde(skip_serializing_if = "Option::is_none")] + is_trigger: Option, }, Identity, } @@ -505,6 +509,7 @@ struct UntaggedFlowModuleValue { custom_concurrency_key: Option, concurrent_limit: Option, concurrency_time_window_s: Option, + is_trigger: Option, } impl<'de> Deserialize<'de> for FlowModuleValue { @@ -522,6 +527,7 @@ impl<'de> Deserialize<'de> for FlowModuleValue { .ok_or_else(|| serde::de::Error::missing_field("path"))?, hash: untagged.hash, tag_override: untagged.tag_override, + is_trigger: untagged.is_trigger, }), "flow" => Ok(FlowModuleValue::Flow { input_transforms: untagged.input_transforms.unwrap_or_default(), @@ -574,6 +580,7 @@ impl<'de> Deserialize<'de> for FlowModuleValue { custom_concurrency_key: untagged.custom_concurrency_key, concurrent_limit: untagged.concurrent_limit, concurrency_time_window_s: untagged.concurrency_time_window_s, + is_trigger: untagged.is_trigger, }), "identity" => Ok(FlowModuleValue::Identity), other => Err(serde::de::Error::unknown_variant( diff --git a/backend/windmill-queue/src/jobs.rs b/backend/windmill-queue/src/jobs.rs index c0235c9b26..9b94c0f198 100644 --- a/backend/windmill-queue/src/jobs.rs +++ b/backend/windmill-queue/src/jobs.rs @@ -3522,6 +3522,7 @@ pub async fn push<'c, 'd, R: rsmq_async::RsmqConnection + Send + 'c>( path: path.clone(), hash: Some(hash), tag_override: tag_override, + is_trigger: None, }, ), stop_after_if: None, diff --git a/backend/windmill-worker/src/worker_lockfiles.rs b/backend/windmill-worker/src/worker_lockfiles.rs index f71c73f618..6c88294126 100644 --- a/backend/windmill-worker/src/worker_lockfiles.rs +++ b/backend/windmill-worker/src/worker_lockfiles.rs @@ -735,6 +735,7 @@ async fn lock_modules<'c>( custom_concurrency_key, concurrent_limit, concurrency_time_window_s, + is_trigger, } = e.get_value()? else { match e.get_value()? { @@ -970,6 +971,7 @@ async fn lock_modules<'c>( custom_concurrency_key, concurrent_limit, concurrency_time_window_s, + is_trigger, }); new_flow_modules.push(e); continue; @@ -992,6 +994,7 @@ async fn lock_modules<'c>( custom_concurrency_key, concurrent_limit, concurrency_time_window_s, + is_trigger, }); new_flow_modules.push(e); continue; diff --git a/frontend/package.json b/frontend/package.json index 0c22355ecc..a15b2e575c 100644 --- a/frontend/package.json +++ b/frontend/package.json @@ -182,6 +182,11 @@ "svelte": "./package/components/icons/WindmillIcon2.svelte", "default": "./package/components/icons/WindmillIcon2.svelte" }, + "./components/icons/SchedulePollIcon.svelte": { + "types": "./package/components/icons/SchedulePollIcon.d.ts", + "svelte": "./package/components/icons/SchedulePollIcon.svelte", + "default": "./package/components/icons/SchedulePollIcon.svelte" + }, "./components/IconedResourceType.svelte": { "types": "./package/components/IconedResourceType.svelte.d.ts", "svelte": "./package/components/IconedResourceType.svelte", diff --git a/frontend/src/lib/components/Dev.svelte b/frontend/src/lib/components/Dev.svelte index 5e3b52d867..bde2bd265f 100644 --- a/frontend/src/lib/components/Dev.svelte +++ b/frontend/src/lib/components/Dev.svelte @@ -476,7 +476,7 @@ const testStepStore = writable>({}) const selectedIdStore = writable('settings-metadata') const selectedTriggerStore = writable< - 'webhooks' | 'emails' | 'schedules' | 'cli' | 'routes' | 'websockets' + 'webhooks' | 'emails' | 'schedules' | 'cli' | 'routes' | 'websockets' | 'scheduledPoll' >('webhooks') const primaryScheduleStore = writable(undefined) @@ -484,7 +484,8 @@ setContext('TriggerContext', { primarySchedule: primaryScheduleStore, selectedTrigger: selectedTriggerStore, - triggersCount: triggersCount + triggersCount: triggersCount, + simplifiedPoll: writable(false) }) setContext('FlowEditorContext', { selectedId: selectedIdStore, diff --git a/frontend/src/lib/components/FlowBuilder.svelte b/frontend/src/lib/components/FlowBuilder.svelte index b04fbdf8dc..f51b7e8bc8 100644 --- a/frontend/src/lib/components/FlowBuilder.svelte +++ b/frontend/src/lib/components/FlowBuilder.svelte @@ -144,7 +144,7 @@ ? { schedule_count: 1, primary_schedule: { schedule: savedPrimarySchedule.cron } } : undefined ) - + const simplifiedPoll = writable(false) export function setPrimarySchedule(schedule: ScheduleTrigger | undefined | false) { primaryScheduleStore.set(schedule) loadTriggers() @@ -468,7 +468,7 @@ const selectedIdStore = writable(selectedId ?? 'settings-metadata') const selectedTriggerStore = writable< - 'webhooks' | 'emails' | 'schedules' | 'cli' | 'routes' | 'websockets' + 'webhooks' | 'emails' | 'schedules' | 'cli' | 'routes' | 'websockets' | 'scheduledPoll' >('webhooks') export function getSelectedId() { @@ -490,7 +490,14 @@ } function selectTrigger( - selectedTrigger: 'webhooks' | 'emails' | 'schedules' | 'cli' | 'routes' | 'websockets' + selectedTrigger: + | 'webhooks' + | 'emails' + | 'schedules' + | 'cli' + | 'routes' + | 'websockets' + | 'scheduledPoll' ) { selectedTriggerStore.set(selectedTrigger) } @@ -517,7 +524,8 @@ setContext('TriggerContext', { selectedTrigger: selectedTriggerStore, primarySchedule: primaryScheduleStore, - triggersCount + triggersCount, + simplifiedPoll }) async function loadTriggers() { @@ -876,7 +884,6 @@ await tick() select(module.id) await tick() - await tick() focusCopilot() let isFirstInLoop = false diff --git a/frontend/src/lib/components/FlowStatusViewerInner.svelte b/frontend/src/lib/components/FlowStatusViewerInner.svelte index e881dad539..ec6982a5df 100644 --- a/frontend/src/lib/components/FlowStatusViewerInner.svelte +++ b/frontend/src/lib/components/FlowStatusViewerInner.svelte @@ -1008,6 +1008,7 @@ modules={job.raw_flow?.modules ?? []} failureModule={job.raw_flow?.failure_module} preprocessorModule={job.raw_flow?.preprocessor_module} + allowSimplifiedPoll={false} />
('TriggerContext') let scheduleEditor: ScheduleEditor + let schedules: Writable = writable(undefined) + let initialPrimarySchedule: Writable = writable(undefined) - $: loadSchedules(false) || path - - let schedules: Schedule[] | undefined = undefined - let initialPrimarySchedule: ScheduleTrigger | false | undefined = undefined - async function loadSchedules(forceRefresh: boolean) { - if (!path || path == '') { - schedules = [] - if ($primarySchedule == undefined) { - $primarySchedule = false - } - initialPrimarySchedule = structuredClone($primarySchedule) - return - } - try { - const allSchedules = await ScheduleService.listSchedules({ - workspace: $workspaceStore ?? '', - path: path, - isFlow - }) - const primary = allSchedules.find((s) => s.path == path) - let remotePrimarySchedule = primary - ? { - summary: primary.summary, - args: primary.args ?? {}, - cron: primary.schedule, - timezone: primary.timezone, - enabled: primary.enabled - } - : false - if ($primarySchedule == undefined || forceRefresh) { - $primarySchedule = remotePrimarySchedule - } - initialPrimarySchedule = structuredClone(remotePrimarySchedule) - - $triggersCount = { - ...($triggersCount ?? {}), - schedule_count: allSchedules.length, - primary_schedule: $primarySchedule ? { schedule: $primarySchedule.cron } : undefined - } - schedules = allSchedules.filter((s) => s.path != path) - } catch (e) { - console.error('impossible to load schedules', e) - } + async function updateSchedules(forceRefresh: boolean) { + loadSchedules( + forceRefresh, + path, + isFlow, + schedules, + primarySchedule, + initialPrimarySchedule, + $workspaceStore ?? '', + triggersCount + ) } + $: updateSchedules(false) || path + async function save() { - const scheduleExists = - path != '' && - !newItem && - (await ScheduleService.existsSchedule({ - workspace: $workspaceStore!, - path - })) - if (scheduleExists) { - console.log('primary schedule exists') - if ($primarySchedule) { - await ScheduleService.updateSchedule({ - workspace: $workspaceStore!, - path, - requestBody: { - summary: $primarySchedule.summary, - args: $primarySchedule.args, - schedule: $primarySchedule.cron, - timezone: $primarySchedule.timezone - } - }) - sendUserToast(`Primary schedule updated`) - } else { - await ScheduleService.deleteSchedule({ workspace: $workspaceStore!, path }) - sendUserToast(`Primary schedule deleted`) - } - } else { - if ($primarySchedule) { - await ScheduleService.createSchedule({ - workspace: $workspaceStore!, - requestBody: { - path, - script_path: path, - is_flow: isFlow, - summary: $primarySchedule.summary, - args: $primarySchedule.args, - schedule: $primarySchedule.cron, - timezone: $primarySchedule.timezone, - enabled: $primarySchedule.enabled - } - }) - sendUserToast(`Primary schedule created`) - } - } - loadSchedules(true) + await saveSchedule(path, newItem, $workspaceStore ?? '', primarySchedule, isFlow) + updateSchedules(true) } { - loadSchedules(true) + updateSchedules(true) }} bind:this={scheduleEditor} /> @@ -204,7 +136,7 @@

Define a schedule frequency first

{/if} - {#if initialPrimarySchedule != false} + {#if $initialPrimarySchedule != false}
- - {#if initialPrimarySchedule != undefined && initialPrimarySchedule != false && !newItem} +
+ +
+ {#if $initialPrimarySchedule != undefined && $initialPrimarySchedule != false && !newItem} @@ -274,12 +208,12 @@ {/if}