feat: clean openflow spec v1 (#491)

* clean api 2

* the rest

* clean tests

* stop_after_if test

* unbox `modules: Vec<FlowModule>` 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 <somebody@froghat.ca>
This commit is contained in:
Ruben Fiszel
2022-09-02 01:26:39 +02:00
committed by GitHub
parent 79b4c18c9a
commit cf7209bdb9
18 changed files with 289 additions and 163 deletions
@@ -0,0 +1,2 @@
-- The corresponding migrate up isn't really reversible but should be
-- idempotent...
@@ -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;
+21 -12
View File
@@ -74,6 +74,11 @@ pub struct FlowValue {
pub modules: Vec<FlowModule>,
pub failure_module: Option<FlowModule>,
}
#[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<String, InputTransform>,
pub value: FlowModuleValue,
pub stop_after_if_expr: Option<String>,
pub skip_if_stopped: Option<bool>,
pub stop_after_if: Option<StopAfterIf>,
pub summary: Option<String>,
}
@@ -107,7 +111,7 @@ pub enum FlowModuleValue {
},
ForloopFlow {
iterator: InputTransform,
value: Box<FlowValue>,
modules: Vec<FlowModule>,
#[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,
}),
};
+128 -101
View File
@@ -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)",
},
}],
},
}],
}))
+7 -4
View File
@@ -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 { .. } => {
@@ -247,7 +247,7 @@
{/if}
<PageHeader title="Extra Params" primary={false} />
{#if !manual && resource_type != ''}
{#each extra_params as [k, v], i}
{#each extra_params as [k, v]}
<div class="flex flex-row max-w-md">
<input type="text" bind:value={k} />
<input type="text" bind:value={v} />
+1 -1
View File
@@ -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 })
}
+1 -1
View File
@@ -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
@@ -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
}
</script>
@@ -133,7 +133,7 @@
>
{#if open[i]}
<div class="border border-black p-2 bg-gray-50 divide-y">
<InputTransformsViewer inputTransforms={mod?.input_transform} />
<InputTransformsViewer inputTransforms={mod?.input_transforms} />
<div class="w-full h-full mt-6">
<iframe
style="height: 400px;"
@@ -157,13 +157,13 @@
Raw {mod?.value?.language} script {open[i] ? '(-)' : '(+)'}</button
>
{#if open[i]}
<div
transition:slide
class="border border-black p-2 bg-gray-50 w-full"
>
<InputTransformsViewer inputTransforms={mod?.input_transform} />
<InputTransformsViewer inputTransforms={mod?.input_transforms} />
<Highlight
language={mod?.value?.language == 'deno' ? typescript : python}
code={mod?.value?.content}
@@ -135,7 +135,7 @@
inputTransform={true}
importPath={String(indexes.join('-'))}
bind:pickableProperties={stepPropPicker.pickableProperties}
bind:args={mod.input_transform}
bind:args={mod.input_transforms}
bind:extraLib={stepPropPicker.extraLib}
/>
{/if}
@@ -41,7 +41,7 @@ export async function flowModulesToFlowState(flowModules: FlowModule[]): Promise
const value = flowModule.value
if (value.type === 'forloopflow') {
const childFlowModules = await Promise.all(
value.value.modules.map(async (module) => loadFlowModuleSchema(module))
value.modules.map(async (module) => loadFlowModuleSchema(module))
)
const loopFlowModule = await loadFlowModuleSchema(flowModule)
@@ -69,7 +69,7 @@ export function flowStateToFlow(flowState: FlowState, flow: Flow): Flow {
const fmv = flowModule.value
if (fmv.type === 'forloopflow' && childFlowModules && Array.isArray(childFlowModules)) {
fmv.value.modules = childFlowModules.map((cfm) => cfm.flowModule)
fmv.modules = childFlowModules.map((cfm) => cfm.flowModule)
flowModule.value = fmv
}
@@ -31,9 +31,9 @@ export function emptyFlowModuleSchema(): FlowModuleSchema {
export async function loadFlowModuleSchema(flowModule: FlowModule): Promise<FlowModuleSchema> {
try {
const { input_transform, schema } = await loadSchemaFromModule(flowModule)
const { input_transforms, schema } = await loadSchemaFromModule(flowModule)
flowModule.input_transform = input_transform
flowModule.input_transforms = input_transforms
return { flowModule, schema }
} catch (e) {
@@ -44,7 +44,7 @@ export async function loadFlowModuleSchema(flowModule: FlowModule): Promise<Flow
export async function pickScript(path: string): Promise<FlowModuleSchema> {
const flowModule: FlowModule = {
value: { type: 'script', path },
input_transform: {}
input_transforms: {}
}
return await loadFlowModuleSchema(flowModule)
@@ -61,7 +61,7 @@ export async function createInlineScriptModule({
const flowModule: FlowModule = {
value: { type: 'rawscript', content: code, language },
input_transform: {}
input_transforms: {}
}
return await loadFlowModuleSchema(flowModule)
@@ -71,13 +71,11 @@ export async function createLoop(): Promise<FlowModuleSchema> {
const loopFlowModule: FlowModule = {
value: {
type: 'forloopflow',
value: {
modules: []
},
modules: [],
iterator: { type: 'javascript', expr: 'result' },
skip_failures: true
},
input_transform: {}
input_transforms: {}
}
const { flowModule, schema } = await loadFlowModuleSchema(loopFlowModule)
@@ -109,7 +107,7 @@ export async function createInlineScriptModuleFromPath(path: string): Promise<Fl
content: content,
path
},
input_transform: {}
input_transforms: {}
}
}
+26 -1
View File
@@ -1,4 +1,4 @@
import type { Flow } from '$lib/gen'
import type { Flow, FlowModule, ForloopFlow, InputTransform } from '$lib/gen'
import { get, writable } from 'svelte/store'
import { flowStateStore, initFlowState } from './flowState'
@@ -7,8 +7,33 @@ export type FlowMode = 'push' | 'pull'
export const flowStore = writable<Flow>(undefined)
export function initFlow(flow: Flow) {
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.modules.forEach(migrateFlowModule)
flowVal.value = undefined
}
}
}
flowStore.set(flow)
initFlowState(flow)
function migrateFlowModule(mod: FlowModule) {
let modVal = mod as FlowModule & { input_transform?: Record<string, InputTransform>, 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.stop_after_if_expr) {
modVal.stop_after_if = { expr: modVal.stop_after_if_expr, skipped_if_stopped: modVal.skip_if_stopped}
delete modVal.stop_after_if_expr;
delete modVal.skip_if_stopped;
}
}
}
export async function copyFirstStepSchema() {
+10 -10
View File
@@ -3,13 +3,13 @@ import { JobService, type Flow, type FlowModule, type InputTransform, type Job }
import { inferArgs } from '$lib/infer'
import { loadSchema } from '$lib/scripts'
import { workspaceStore } from '$lib/stores'
import { emptySchema, schemaToObject } from '$lib/utils'
import { emptySchema } from '$lib/utils'
import { get } from 'svelte/store'
export function cleanInputs(flow: Flow | any): Flow {
const newFlow: Flow = JSON.parse(JSON.stringify(flow))
newFlow.value.modules.forEach((mod) => {
Object.values(mod.input_transform).forEach((inp) => {
Object.values(mod.input_transforms).forEach((inp) => {
// for now we use the value for dynamic expression when done in the static editor so we have to resort to this
if (inp.type == 'javascript') {
//@ts-ignore
@@ -58,7 +58,7 @@ export function scrollIntoView(element: any) {
}
export async function loadSchemaFromModule(module: FlowModule): Promise<{
input_transform: Record<string, InputTransform>
input_transforms: Record<string, InputTransform>
schema: Schema
}> {
const mod = module.value
@@ -72,20 +72,20 @@ export async function loadSchemaFromModule(module: FlowModule): Promise<{
schema = await loadSchema(mod.path!)
} else {
return {
input_transform: {},
input_transforms: {},
schema: emptySchema()
}
}
const keys = Object.keys(schema?.properties ?? {})
let input_transform = module.input_transform
let input_transforms = module.input_transforms
if (
JSON.stringify(keys.sort()) !== JSON.stringify(Object.keys(module.input_transform).sort())
JSON.stringify(keys.sort()) !== JSON.stringify(Object.keys(module.input_transforms).sort())
) {
input_transform = keys.reduce((accu, key) => {
let nv = module.input_transform[key] ?? {
input_transforms = keys.reduce((accu, key) => {
let nv = module.input_transforms[key] ?? {
type: 'static',
value: undefined
}
@@ -95,13 +95,13 @@ export async function loadSchemaFromModule(module: FlowModule): Promise<{
}
return {
input_transform: input_transform,
input_transforms: input_transforms,
schema: schema ?? emptySchema()
}
}
return {
input_transform: {},
input_transforms: {},
schema: emptySchema()
}
}
+3 -4
View File
@@ -54,9 +54,8 @@ export function displayDate(dateString: string | undefined): string {
if (date.toString() === 'Invalid Date') {
return ''
} else {
return `${date.getFullYear()}/${
date.getMonth() + 1
}/${date.getDate()} at ${date.toLocaleTimeString()}`
return `${date.getFullYear()}/${date.getMonth() + 1
}/${date.getDate()} at ${date.toLocaleTimeString()}`
}
}
@@ -144,7 +143,7 @@ export function emptySchema() {
export function emptyModule(): FlowModule {
return {
value: { type: 'script', path: '' },
input_transform: {}
input_transforms: {}
}
}
+1 -1
View File
@@ -149,7 +149,7 @@
sendUserToast(`Successfully archived script ${path}`)
}
async function viewCode(path) {
async function viewCode(path: string) {
codeViewerContent = (await getScriptByPath(path)).content
codeViewerPath = path
codeViewer.openModal()
+16 -9
View File
@@ -48,20 +48,25 @@ components:
FlowModule:
type: object
properties:
input_transform:
input_transforms:
type: object
additionalProperties:
$ref: "#/components/schemas/InputTransform"
value:
$ref: "#/components/schemas/FlowModuleValue"
stop_after_if_expr:
type: string
skip_if_stopped:
type: boolean
stop_after_if:
type: object
properties:
skipped_if_stopped:
type: boolean
expr:
type: string
required:
- expr
summary:
type: string
required:
- input_transform
- input_transforms
- value
InputTransform:
@@ -160,8 +165,10 @@ components:
ForloopFlow:
type: object
properties:
value:
$ref: "#/components/schemas/FlowValue"
modules:
type: array
items:
$ref: "#/components/schemas/FlowModule"
iterator:
$ref: "#/components/schemas/InputTransform"
skip_failures:
@@ -171,7 +178,7 @@ components:
enum:
- forloopflow
required:
- value
- modules
- iterator
- skip_failures
- type