diff --git a/backend/tests/worker.rs b/backend/tests/worker.rs index b1de2ba9ab..048c1c4e7b 100644 --- a/backend/tests/worker.rs +++ b/backend/tests/worker.rs @@ -1001,7 +1001,6 @@ async fn test_deno_flow(db: Pool) { path: None, lock: None, }, - input_transforms: Default::default(), stop_after_if: Default::default(), summary: Default::default(), suspend: Default::default(), @@ -1029,7 +1028,6 @@ async fn test_deno_flow(db: Pool) { path: None, lock: None, }, - input_transforms: Default::default(), stop_after_if: Default::default(), summary: Default::default(), suspend: Default::default(), @@ -1037,7 +1035,6 @@ async fn test_deno_flow(db: Pool) { sleep: None, }], }, - input_transforms: Default::default(), stop_after_if: Default::default(), summary: Default::default(), suspend: Default::default(), @@ -1131,7 +1128,6 @@ async fn test_deno_flow_same_worker(db: Pool) { path: None, lock: None, }, - input_transforms: Default::default(), stop_after_if: Default::default(), summary: Default::default(), suspend: Default::default(), @@ -1147,25 +1143,24 @@ async fn test_deno_flow_same_worker(db: Pool) { modules: vec![ FlowModule { id: "d".to_string(), - input_transforms: [ - ( - "i".to_string(), - InputTransform::Javascript { - expr: "flow_input.iter.value".to_string(), - }, - ), - ( - "loop".to_string(), - InputTransform::Static { value: json!(true) }, - ), - ( - "path".to_string(), - InputTransform::Static { value: json!("inner.txt") }, - ), - ] - .into(), value: FlowModuleValue::RawScript { - input_transforms: [].into(), + input_transforms: [ + ( + "i".to_string(), + InputTransform::Javascript { + expr: "flow_input.iter.value".to_string(), + }, + ), + ( + "loop".to_string(), + InputTransform::Static { value: json!(true) }, + ), + ( + "path".to_string(), + InputTransform::Static { value: json!("inner.txt") }, + ), + ] + .into(), language: ScriptLang::Deno, content: write_file, path: None, @@ -1196,7 +1191,6 @@ async fn test_deno_flow_same_worker(db: Pool) { path: None, lock: None, }, - input_transforms: [].into(), stop_after_if: Default::default(), summary: Default::default(), suspend: Default::default(), @@ -1205,7 +1199,6 @@ async fn test_deno_flow_same_worker(db: Pool) { }, ], }, - input_transforms: Default::default(), stop_after_if: Default::default(), summary: Default::default(), suspend: Default::default(), @@ -1240,7 +1233,6 @@ async fn test_deno_flow_same_worker(db: Pool) { path: None, lock: None, }, - input_transforms: [].into(), stop_after_if: Default::default(), summary: Default::default(), suspend: Default::default(), diff --git a/backend/windmill-api/src/flows.rs b/backend/windmill-api/src/flows.rs index 16b8173bf5..51b018af75 100644 --- a/backend/windmill-api/src/flows.rs +++ b/backend/windmill-api/src/flows.rs @@ -552,7 +552,6 @@ mod tests { modules: vec![ FlowModule { id: "a".to_string(), - input_transforms: [].into(), value: FlowModuleValue::Script { path: "test".to_string(), input_transforms: [( @@ -570,7 +569,6 @@ mod tests { }, FlowModule { id: "b".to_string(), - input_transforms: HashMap::new(), value: FlowModuleValue::RawScript { input_transforms: HashMap::new(), content: "test".to_string(), @@ -589,7 +587,6 @@ mod tests { }, FlowModule { id: "c".to_string(), - input_transforms: HashMap::new(), value: FlowModuleValue::ForloopFlow { iterator: InputTransform::Static { value: serde_json::json!([1, 2, 3]) }, modules: vec![], @@ -608,7 +605,6 @@ mod tests { ], failure_module: Some(FlowModule { id: "d".to_string(), - input_transforms: HashMap::new(), value: FlowModuleValue::Script { path: "test".to_string(), input_transforms: HashMap::new(), @@ -695,31 +691,6 @@ mod tests { assert_eq!(dbg!(serde_json::json!(fv)), dbg!(expect)); } - #[test] - fn test_back_compat() { - /* renamed input_transform -> input_transforms but should deserialize old name */ - let s = r#" - { - "value": { - "type": "rawscript", - "content": "def main(n): return", - "language": "python3" - }, - "input_transform": { - "n": { - "expr": "flow_input.iter.value", - "type": "javascript" - } - } - } - "#; - let module: FlowModule = serde_json::from_str(s).unwrap(); - assert_eq!( - module.input_transforms["n"], - InputTransform::Javascript { expr: "flow_input.iter.value".to_string() } - ); - } - #[test] fn retry_serde() { assert_eq!(Retry::default(), serde_json::from_str(r#"{}"#).unwrap()); diff --git a/backend/windmill-common/src/flows.rs b/backend/windmill-common/src/flows.rs index 99e01f2c30..0e95a84ff2 100644 --- a/backend/windmill-common/src/flows.rs +++ b/backend/windmill-common/src/flows.rs @@ -6,9 +6,12 @@ * LICENSE-AGPL for a copy of the license. */ -use std::{collections::HashMap, time::Duration}; +use std::{ + collections::{BTreeMap, HashMap}, + time::Duration, +}; -use serde::{self, Deserialize, Serialize}; +use serde::{self, Deserialize, Serialize, Serializer}; use crate::{ more_serde::{ @@ -149,9 +152,6 @@ pub struct Suspend { pub struct FlowModule { #[serde(default = "default_id")] pub id: String, - #[serde(default)] - #[serde(alias = "input_transform")] - pub input_transforms: HashMap, pub value: FlowModuleValue, #[serde(skip_serializing_if = "Option::is_none")] pub stop_after_if: Option, @@ -245,7 +245,7 @@ pub enum FlowModuleValue { }, RawScript { #[serde(default)] - #[serde(alias = "input_transform")] + #[serde(alias = "input_transform", serialize_with = "ordered_map")] input_transforms: HashMap, content: String, #[serde(skip_serializing_if = "Option::is_none")] @@ -257,6 +257,14 @@ pub enum FlowModuleValue { Identity, } +fn ordered_map(value: &HashMap, serializer: S) -> Result +where + S: Serializer, +{ + let ordered: BTreeMap<_, _> = value.iter().collect(); + ordered.serialize(serializer) +} + #[derive(Deserialize)] pub struct ListFlowQuery { pub path_start: Option, diff --git a/backend/windmill-queue/src/jobs.rs b/backend/windmill-queue/src/jobs.rs index dd9ac1afa4..fc232554d6 100644 --- a/backend/windmill-queue/src/jobs.rs +++ b/backend/windmill-queue/src/jobs.rs @@ -489,7 +489,6 @@ pub async fn push<'c>( modules.push(FlowModule { id: format!("{}-v", flow.modules[flow.modules.len() - 1].id), value: FlowModuleValue::Identity, - input_transforms: HashMap::new(), stop_after_if: None, summary: Some( "Virtual module needed for suspend/sleep when last module".to_string(), diff --git a/backend/windmill-worker/src/worker_flow.rs b/backend/windmill-worker/src/worker_flow.rs index c398c96407..92f45d740a 100644 --- a/backend/windmill-worker/src/worker_flow.rs +++ b/backend/windmill-worker/src/worker_flow.rs @@ -1124,11 +1124,7 @@ async fn push_next_flow_job( transform_input( &flow_job.args, last_result.clone(), - if !input_transforms.is_empty() { - input_transforms - } else { - &module.input_transforms - }, + input_transforms, resume_messages.as_slice(), approvers, by_id, diff --git a/frontend/src/lib/components/flows/flowStore.ts b/frontend/src/lib/components/flows/flowStore.ts index 04af7cfedd..4eec3f0e9b 100644 --- a/frontend/src/lib/components/flows/flowStore.ts +++ b/frontend/src/lib/components/flows/flowStore.ts @@ -1,13 +1,11 @@ -import type { Flow, FlowModule, ForloopFlow, InputTransform } from '$lib/gen' +import type { Flow, FlowModule } from '$lib/gen' import { writable, type Writable } from 'svelte/store' import { initFlowState, type FlowState } from './flowState' -import { numberToChars } from './utils' export type FlowMode = 'push' | 'pull' export const importFlowStore = writable(undefined) - export function dfs(modules: FlowModule[], f: (x: FlowModule) => T): T[] { let result: T[] = [] for (const module of modules) { @@ -16,7 +14,15 @@ export function dfs(modules: FlowModule[], f: (x: FlowModule) => T): T[] { result = result.concat(dfs(module.value.modules, f)) } else if (module.value.type == 'branchone') { result = result.concat(f(module)) - result = result.concat(dfs(module.value.branches.map((b) => b.modules).flat().concat(module.value.default), f)) + result = result.concat( + dfs( + module.value.branches + .map((b) => b.modules) + .flat() + .concat(module.value.default), + f + ) + ) } else if (module.value.type == 'branchall') { result = result.concat(f(module)) result = result.concat(dfs(module.value.branches.map((b) => b.modules).flat(), f)) @@ -27,57 +33,13 @@ export function dfs(modules: FlowModule[], f: (x: FlowModule) => T): T[] { return result } - -export async function initFlow(flow: Flow, flowStore: Writable, flowStateStore: Writable) { - let counter = 0 - for (const mod of flow.value.modules) { - migrateFlowModule(mod) - let val = mod.value - if (val.type == 'forloopflow') { - let flowVal = val as ForloopFlow & { value?: { modules?: FlowModule[] } } - if (flowVal.value && flowVal.value.modules) { - flowVal.modules = flowVal.value.modules - flowVal.value = undefined - } - flowVal.modules.forEach(migrateFlowModule) - - } - } - +export async function initFlow( + flow: Flow, + flowStore: Writable, + flowStateStore: Writable +) { await initFlowState(flow, flowStateStore) flowStore.set(flow) - - function migrateFlowModule(mod: FlowModule) { - if (mod.id == undefined) { - mod.id = numberToChars(counter++) - } - let modVal = mod as FlowModule & { - input_transform?: Record - stop_after_if_expr?: string - skip_if_stopped?: boolean - } - if (modVal.input_transform) { - modVal.input_transforms = modVal.input_transform - delete modVal.input_transform - } - if ( - (modVal.input_transforms && modVal.value.type == 'script') || - modVal.value.type == 'rawscript' - ) { - if (modVal.input_transforms && Object.keys(modVal.input_transforms).length > 0) { - modVal.value.input_transforms = modVal.input_transforms - delete modVal.input_transforms - } - } - if (modVal.stop_after_if_expr) { - modVal.stop_after_if = { - expr: modVal.stop_after_if_expr, - skip_if_stopped: modVal.skip_if_stopped - } - delete modVal.stop_after_if_expr - delete modVal.skip_if_stopped - } - } } export async function copyFirstStepSchema(flowState: FlowState, flowStore: Writable) { diff --git a/frontend/src/lib/infer.ts b/frontend/src/lib/infer.ts index eabfa75286..9b27d8ceac 100644 --- a/frontend/src/lib/infer.ts +++ b/frontend/src/lib/infer.ts @@ -1,6 +1,7 @@ import { ScriptService, type MainArgSignature } from '$lib/gen' import { get, writable } from 'svelte/store' import type { Schema, SchemaProperty } from './common.js' +import { sortObject } from './utils.js' const loadSchemaLastRun = writable<[string | undefined, MainArgSignature | undefined]>(undefined) @@ -49,6 +50,7 @@ export async function inferArgs( } else { schema.properties[arg.name] = oldProperties[arg.name] } + schema.properties[arg.name] = sortObject(schema.properties[arg.name]) argSigToJsonSchemaType(arg.typ, schema.properties[arg.name]) schema.properties[arg.name].default = arg.default diff --git a/frontend/src/lib/utils.ts b/frontend/src/lib/utils.ts index 0b5bf80d88..9a3c8e6bd3 100644 --- a/frontend/src/lib/utils.ts +++ b/frontend/src/lib/utils.ts @@ -717,3 +717,12 @@ export function getModifierKey(): string { export function isValidHexColor(color: string): boolean { return /^#(([A-F0-9]{2}){3,4}|[A-F0-9]{3})$/i.test(color) } + +export function sortObject(o: T & object): T { + return Object.keys(o) + .sort() + .reduce((obj, key) => { + obj[key] = obj[key] + return obj + }, {}) as T +} diff --git a/openflow.openapi.yaml b/openflow.openapi.yaml index 8cdb94d1b7..c12a03aa48 100644 --- a/openflow.openapi.yaml +++ b/openflow.openapi.yaml @@ -75,11 +75,6 @@ components: properties: id: type: string - # to be removed in favor of raw/script input_transforms once migration is over - input_transforms: - type: object - additionalProperties: - $ref: "#/components/schemas/InputTransform" value: $ref: "#/components/schemas/FlowModuleValue" stop_after_if: @@ -327,7 +322,7 @@ components: type: boolean required: - type - + FlowStatus: type: object properties: @@ -344,7 +339,7 @@ components: properties: parent_module: type: string - + retry: type: object properties: