From cf7209bdb92bc4f029224640ccdc5213e2c3cb98 Mon Sep 17 00:00:00 2001 From: Ruben Fiszel Date: Fri, 2 Sep 2022 01:26:39 +0200 Subject: [PATCH] feat: clean openflow spec v1 (#491) * clean api 2 * the rest * clean tests * stop_after_if test * unbox `modules: Vec` in `ForloopFlow` * migrate * initFlow stop_after_if_expr and skip_if_stopped * s/migrateInitTransform/migrateFlowModule I didn't read the name before... oops * sql migration for openflow changes * fix frontend migration code Co-authored-by: sqwishy --- .../20220901174503_openflow_changes.down.sql | 2 + .../20220901174503_openflow_changes.up.sql | 56 +++++ backend/src/flows.rs | 33 ++- backend/src/worker.rs | 229 ++++++++++-------- backend/src/worker_flow.rs | 11 +- frontend/src/lib/components/AppConnect.svelte | 2 +- frontend/src/lib/components/ArgInfo.svelte | 2 +- frontend/src/lib/components/ArgInput.svelte | 2 +- .../src/lib/components/FlowPreview.svelte | 6 +- frontend/src/lib/components/FlowViewer.svelte | 6 +- frontend/src/lib/components/ModuleStep.svelte | 2 +- .../src/lib/components/flows/flowState.ts | 4 +- .../lib/components/flows/flowStateUtils.ts | 16 +- .../src/lib/components/flows/flowStore.ts | 27 ++- frontend/src/lib/components/flows/utils.ts | 20 +- frontend/src/lib/utils.ts | 7 +- frontend/src/routes/scripts.svelte | 2 +- openflow.openapi.yaml | 25 +- 18 files changed, 289 insertions(+), 163 deletions(-) create mode 100644 backend/migrations/20220901174503_openflow_changes.down.sql create mode 100644 backend/migrations/20220901174503_openflow_changes.up.sql diff --git a/backend/migrations/20220901174503_openflow_changes.down.sql b/backend/migrations/20220901174503_openflow_changes.down.sql new file mode 100644 index 0000000000..8269131b66 --- /dev/null +++ b/backend/migrations/20220901174503_openflow_changes.down.sql @@ -0,0 +1,2 @@ +-- The corresponding migrate up isn't really reversible but should be +-- idempotent... diff --git a/backend/migrations/20220901174503_openflow_changes.up.sql b/backend/migrations/20220901174503_openflow_changes.up.sql new file mode 100644 index 0000000000..9de0ce8027 --- /dev/null +++ b/backend/migrations/20220901174503_openflow_changes.up.sql @@ -0,0 +1,56 @@ +-- https://github.com/windmill-labs/windmill/pull/491 +CREATE FUNCTION migrate_flow(flow jsonb) + RETURNS jsonb +AS $$ +DECLARE module jsonb; + i integer := 0; +BEGIN + if flow->'value'?'modules' THEN + flow = JSONB_SET(flow, ARRAY['modules'], flow->'value'->'modules') - 'value'; + END IF; + + FOR module IN SELECT JSONB_ARRAY_ELEMENTS(flow->'modules') LOOP + flow = JSONB_SET(flow, ARRAY['modules', i::text], migrate_flow_module(module)); + i = i + 1; + END LOOP; + + RETURN flow; +END; +$$ LANGUAGE plpgsql; + +CREATE FUNCTION migrate_flow_module(module jsonb) + RETURNS jsonb +AS $$ +BEGIN + IF module?'input_transform' AND module->'input_transform' != 'null'::jsonb THEN + module = JSONB_SET(module, ARRAY['input_transforms'], module->'input_transform') + - 'input_transform'; + END IF; + + IF module?'stop_after_if_expr' AND module->'stop_after_if_expr' != 'null'::jsonb THEN + IF NOT module?'stop_after_if' THEN + module = JSONB_SET(module, ARRAY['stop_after_if'], '{}'::jsonb); + END IF; + module = JSONB_SET(module, ARRAY['stop_after_if', 'expr'], module->'stop_after_if_expr') + - 'stop_after_if_expr'; + END IF; + + IF module?'skip_if_stopped' AND module->'skip_if_stopped' != 'null'::jsonb THEN + IF NOT module?'stop_after_if' THEN + module = JSONB_SET(module, ARRAY['stop_after_if'], '{}'::jsonb); + END IF; + module = JSONB_SET(module, ARRAY['stop_after_if', 'skip_if_stopped'], module->'skip_if_stopped') + - 'skip_if_stopped'; + END IF; + + if module->'value'->>'type' = 'forloopflow' THEN + module = JSONB_SET(module, ARRAY['value'], migrate_flow(module->'value')); + END IF; + + RETURN module; +END; +$$ LANGUAGE plpgsql; + +UPDATE flow SET value = migrate_flow(value); + +DROP FUNCTION migrate_flow_module, migrate_flow; diff --git a/backend/src/flows.rs b/backend/src/flows.rs index 904db755b0..f373fd0010 100644 --- a/backend/src/flows.rs +++ b/backend/src/flows.rs @@ -74,6 +74,11 @@ pub struct FlowValue { pub modules: Vec, pub failure_module: Option, } +#[derive(Deserialize, Serialize, Debug, Clone)] +pub struct StopAfterIf { + pub expr: String, + pub skip_if_stopped: bool, +} #[derive(Deserialize, Serialize, Debug, Clone)] pub struct FlowModule { @@ -81,8 +86,7 @@ pub struct FlowModule { #[serde(alias = "input_transform")] pub input_transforms: HashMap, pub value: FlowModuleValue, - pub stop_after_if_expr: Option, - pub skip_if_stopped: Option, + pub stop_after_if: Option, pub summary: Option, } @@ -107,7 +111,7 @@ pub enum FlowModuleValue { }, ForloopFlow { iterator: InputTransform, - value: Box, + modules: Vec, #[serde(default = "default_true")] skip_failures: bool, }, @@ -421,8 +425,7 @@ mod tests { FlowModule { input_transforms: hm, value: FlowModuleValue::Script { path: "test".to_string() }, - stop_after_if_expr: None, - skip_if_stopped: Some(false), + stop_after_if: None, summary: None, }, FlowModule { @@ -432,8 +435,10 @@ mod tests { language: crate::scripts::ScriptLang::Deno, path: None, }), - stop_after_if_expr: Some("foo = 'bar'".to_string()), - skip_if_stopped: None, + stop_after_if: Some(StopAfterIf { + expr: "foo = 'bar'".to_string(), + skip_if_stopped: false, + }), summary: None, }, FlowModule { @@ -444,19 +449,23 @@ mod tests { .into(), value: FlowModuleValue::ForloopFlow { iterator: InputTransform::Static { value: serde_json::json!([1, 2, 3]) }, - value: Box::new(FlowValue { modules: vec![], failure_module: None }), + modules: vec![], skip_failures: true, }, - stop_after_if_expr: Some("previous.isEmpty()".to_string()), - skip_if_stopped: None, + stop_after_if: Some(StopAfterIf { + expr: "previous.isEmpty()".to_string(), + skip_if_stopped: false, + }), summary: None, }, ], failure_module: Some(FlowModule { input_transforms: HashMap::new(), value: FlowModuleValue::Flow { path: "test".to_string() }, - stop_after_if_expr: Some("previous.isEmpty()".to_string()), - skip_if_stopped: None, + stop_after_if: Some(StopAfterIf { + expr: "previous.isEmpty()".to_string(), + skip_if_stopped: false, + }), summary: None, }), }; diff --git a/backend/src/worker.rs b/backend/src/worker.rs index e39f3e67ad..fcf428c193 100644 --- a/backend/src/worker.rs +++ b/backend/src/worker.rs @@ -1272,39 +1272,32 @@ mod tests { path: None, }), input_transforms: Default::default(), - stop_after_if_expr: Default::default(), - skip_if_stopped: Default::default(), + stop_after_if: Default::default(), summary: Default::default(), }, FlowModule { value: FlowModuleValue::ForloopFlow { iterator: InputTransform::Javascript { expr: "result".to_string() }, skip_failures: false, - value: FlowValue { - modules: vec![FlowModule { - value: FlowModuleValue::RawScript(RawCode { - language: ScriptLang::Deno, - content: doubles.to_string(), - path: None, - }), - input_transforms: [( - "n".to_string(), - InputTransform::Javascript { - expr: "previous_result.iter.value".to_string(), - }, - )] - .into(), - stop_after_if_expr: Default::default(), - skip_if_stopped: Default::default(), - summary: Default::default(), - }], - failure_module: Default::default(), - } - .into(), + modules: vec![FlowModule { + value: FlowModuleValue::RawScript(RawCode { + language: ScriptLang::Deno, + content: doubles.to_string(), + path: None, + }), + input_transforms: [( + "n".to_string(), + InputTransform::Javascript { + expr: "previous_result.iter.value".to_string(), + }, + )] + .into(), + stop_after_if: Default::default(), + summary: Default::default(), + }], }, input_transforms: Default::default(), - stop_after_if_expr: Default::default(), - skip_if_stopped: Default::default(), + stop_after_if: Default::default(), summary: Default::default(), }, ], @@ -1320,6 +1313,50 @@ mod tests { } } + #[sqlx::test(fixtures("base"))] + async fn test_stop_after_if(db: DB) { + initialize_tracing().await; + + let flow: FlowValue = serde_json::from_value(serde_json::json!({ + "modules": [ + { + "input_transforms": { "n": { "type": "javascript", "expr": "flow_input.n" } }, + "value": { + "type": "rawscript", + "language": "python3", + "content": "def main(n): return n", + }, + "stop_after_if": { + "expr": "result < 0", + "skip_if_stopped": false, + }, + }, + { + "input_transforms": { "n": { "type": "javascript", "expr": "previous_result" } }, + "value": { + "type": "rawscript", + "language": "python3", + "content": "def main(n): return f'last step saw {n}'", + }, + }, + ], + })) + .unwrap(); + let job = JobPayload::RawFlow { value: flow, path: None }; + + let result = RunJob::from(job.clone()) + .arg("n", json!(123)) + .wait_until_complete(&db) + .await; + assert_eq!(json!("last step saw 123"), result); + + let result = RunJob::from(job.clone()) + .arg("n", json!(-123)) + .wait_until_complete(&db) + .await; + assert_eq!(json!(-123), result); + } + #[sqlx::test(fixtures("base"))] async fn test_python_flow(db: DB) { initialize_tracing().await; @@ -1342,21 +1379,19 @@ mod tests { "type": "forloopflow", "iterator": { "type": "javascript", "expr": "result" }, "skip_failures": false, - "value": { - "modules": [{ - "value": { - "type": "rawscript", - "language": "python3", - "content": doubles, + "modules": [{ + "value": { + "type": "rawscript", + "language": "python3", + "content": doubles, + }, + "input_transform": { + "n": { + "type": "javascript", + "expr": "previous_result.iter.value", }, - "input_transform": { - "n": { - "type": "javascript", - "expr": "previous_result.iter.value", - }, - }, - }], - } + }, + }], }, }, ], @@ -1451,23 +1486,21 @@ def main(): "value": { "type": "forloopflow", "iterator": { "type": "static", "value": [] }, - "value": { - "modules": [ - { - "input_transform": { - "n": { - "type": "javascript", - "expr": "previous_result.iter.value", - }, + "modules": [ + { + "input_transform": { + "n": { + "type": "javascript", + "expr": "previous_result.iter.value", }, - "value": { - "type": "rawscript", - "language": "python3", - "content": "def main(n): return n", - }, - } - ], - } + }, + "value": { + "type": "rawscript", + "language": "python3", + "content": "def main(n): return n", + }, + } + ], }, }, { @@ -1503,23 +1536,21 @@ def main(): "value": { "type": "forloopflow", "iterator": { "type": "static", "value": [] }, - "value": { - "modules": [ - { - "input_transform": { - "n": { - "type": "javascript", - "expr": "previous_result.iter.value", - }, + "modules": [ + { + "input_transform": { + "n": { + "type": "javascript", + "expr": "previous_result.iter.value", }, - "value": { - "type": "rawscript", - "language": "python3", - "content": "def main(n): return n", - }, - } - ], - } + }, + "value": { + "type": "rawscript", + "language": "python3", + "content": "def main(n): return n", + }, + } + ], }, }, ], @@ -1542,23 +1573,21 @@ def main(): "value": { "type": "forloopflow", "iterator": { "type": "static", "value": [2,3,4] }, - "value": { - "modules": [ - { - "input_transform": { - "n": { - "type": "javascript", - "expr": "previous_result.iter.value", - }, + "modules": [ + { + "input_transform": { + "n": { + "type": "javascript", + "expr": "previous_result.iter.value", }, - "value": { - "type": "rawscript", - "language": "python3", - "content": "def main(n): return n", - } , - } - ], - } + }, + "value": { + "type": "rawscript", + "language": "python3", + "content": "def main(n): return n", + } , + } + ], }, }, { @@ -1675,21 +1704,19 @@ def main(): "type": "forloopflow", "iterator": { "type": "javascript", "expr": "result.items" }, "skip_failures": false, - "value": { - "modules": [{ - "input_transform": { - "n": { - "type": "javascript", - "expr": "previous_result.iter.value", - }, + "modules": [{ + "input_transform": { + "n": { + "type": "javascript", + "expr": "previous_result.iter.value", }, - "value": { - "type": "rawscript", - "language": "python3", - "content": "def main(n):\n if 1 < n:\n raise StopIteration(n)", - }, - }], - } + }, + "value": { + "type": "rawscript", + "language": "python3", + "content": "def main(n):\n if 1 < n:\n raise StopIteration(n)", + }, + }], }, }], })) diff --git a/backend/src/worker_flow.rs b/backend/src/worker_flow.rs index b15c40f884..4cdaf3ade3 100644 --- a/backend/src/worker_flow.rs +++ b/backend/src/worker_flow.rs @@ -147,8 +147,8 @@ pub async fn update_flow_status_after_job_completion( ARRAY['step'], $3) WHERE id = $4 RETURNING - (raw_flow->'modules'->$1->>'stop_after_if_expr'), - (raw_flow->'modules'->$1->>'skip_if_stopped')::bool + (raw_flow->'modules'->$1->'stop_after_if'->>'expr'), + (raw_flow->'modules'->$1->'stop_after_if'->>'skip_if_stopped')::bool ", ) .bind(old_status.step) @@ -647,8 +647,11 @@ async fn push_next_flow_job( } JobPayload::Code(raw_code) } - FlowModuleValue::ForloopFlow { value, .. } => JobPayload::RawFlow { - value: (**value).clone(), + FlowModuleValue::ForloopFlow { modules, .. } => JobPayload::RawFlow { + value: FlowValue { + modules: (*modules).clone(), + failure_module: flow.failure_module.clone(), + }, path: Some(format!("{}/{}", flow_job.script_path(), status.step)), }, a @ FlowModuleValue::Flow { .. } => { diff --git a/frontend/src/lib/components/AppConnect.svelte b/frontend/src/lib/components/AppConnect.svelte index 43f217259d..d6cb204b83 100644 --- a/frontend/src/lib/components/AppConnect.svelte +++ b/frontend/src/lib/components/AppConnect.svelte @@ -247,7 +247,7 @@ {/if} {#if !manual && resource_type != ''} - {#each extra_params as [k, v], i} + {#each extra_params as [k, v]}
diff --git a/frontend/src/lib/components/ArgInfo.svelte b/frontend/src/lib/components/ArgInfo.svelte index c806eaa09a..22f802127c 100644 --- a/frontend/src/lib/components/ArgInfo.svelte +++ b/frontend/src/lib/components/ArgInfo.svelte @@ -15,7 +15,7 @@ return typeof value === 'string' || value instanceof String } - async function getResource(path) { + async function getResource(path: string) { resource = await ResourceService.getResource({ workspace: $workspaceStore!, path }) } diff --git a/frontend/src/lib/components/ArgInput.svelte b/frontend/src/lib/components/ArgInput.svelte index 86f46f868c..d5ca8bd99d 100644 --- a/frontend/src/lib/components/ArgInput.svelte +++ b/frontend/src/lib/components/ArgInput.svelte @@ -20,7 +20,7 @@ import ResourcePicker from './ResourcePicker.svelte' import StringTypeNarrowing from './StringTypeNarrowing.svelte' import SchemaForm from './SchemaForm.svelte' - import type { Schema, SchemaProperty } from '$lib/common' + import type { SchemaProperty } from '$lib/common' export let label: string = '' export let value: any diff --git a/frontend/src/lib/components/FlowPreview.svelte b/frontend/src/lib/components/FlowPreview.svelte index d2977fc670..6fe7b41266 100644 --- a/frontend/src/lib/components/FlowPreview.svelte +++ b/frontend/src/lib/components/FlowPreview.svelte @@ -53,14 +53,14 @@ } function setInputTransformFromArgs(flow: Flow, args: any) { - let input_transform = {} + let input_transforms = {} Object.entries(args).forEach(([key, value]) => { - input_transform[key] = { + input_transforms[key] = { type: 'static', value: value } }) - flow.value.modules[0].input_transform = input_transform + flow.value.modules[0].input_transforms = input_transforms return flow } diff --git a/frontend/src/lib/components/FlowViewer.svelte b/frontend/src/lib/components/FlowViewer.svelte index bd9e4cec6a..b71cb08ba8 100644 --- a/frontend/src/lib/components/FlowViewer.svelte +++ b/frontend/src/lib/components/FlowViewer.svelte @@ -133,7 +133,7 @@ > {#if open[i]}
- +