From 4929a50286300fa124af55e27bcc563de371b643 Mon Sep 17 00:00:00 2001 From: Ruben Fiszel Date: Sun, 16 Oct 2022 15:56:41 +0200 Subject: [PATCH] feat(backend): rework forloop flow job arg passing + reimplement branchone using flows --- backend/src/jobs.rs | 5 - backend/src/worker.rs | 16 +- backend/src/worker_flow.rs | 255 ++++++++++-------- .../components/flows/content/FlowLoop.svelte | 60 ----- 4 files changed, 151 insertions(+), 185 deletions(-) diff --git a/backend/src/jobs.rs b/backend/src/jobs.rs index 2dbc7d6239..59d6a3d941 100644 --- a/backend/src/jobs.rs +++ b/backend/src/jobs.rs @@ -1402,11 +1402,6 @@ pub async fn push<'c>( let is_running = same_worker; if let Some(flow) = raw_flow.as_ref() { same_worker = same_worker || flow.same_worker; - if flow.modules.len() == 0 { - Err(Error::BadRequest(format!( - "A flow needs at least one module to run" - )))?; - } for module in flow.modules.iter() { if let Some(retry) = &module.retry { diff --git a/backend/src/worker.rs b/backend/src/worker.rs index 7480d279e8..579f876edf 100644 --- a/backend/src/worker.rs +++ b/backend/src/worker.rs @@ -441,7 +441,7 @@ async fn handle_queued_job( match job.job_kind { JobKind::FlowPreview | JobKind::Flow => { let args = job.args.clone().unwrap_or(Value::Null); - handle_flow(&job, db, args, same_worker_tx).await?; + handle_flow(&job, db, args, same_worker_tx, worker_dir).await?; } _ => { let mut logs = "".to_string(); @@ -2642,7 +2642,7 @@ def main(): "value": { "type": "rawscript", "language": "deno", - "content": "export function main(array, i){ array.push(i); return array}", + "content": "export function main(array, i){ array.push(i); return array }", } }) } @@ -3527,13 +3527,13 @@ def main(error, port): base_internal_url: String::new(), base_url: String::new(), disable_nuser: std::env::var("DISABLE_NUSER") - .ok() - .and_then(|x| x.parse::().ok()) - .unwrap_or(false), + .ok() + .and_then(|x| x.parse::().ok()) + .unwrap_or(false), disable_nsjail: std::env::var("DISABLE_NSJAIL") - .ok() - .and_then(|x| x.parse::().ok()) - .unwrap_or(false), + .ok() + .and_then(|x| x.parse::().ok()) + .unwrap_or(false), keep_job_dir: std::env::var("KEEP_JOB_DIR") .ok() .and_then(|x| x.parse::().ok()) diff --git a/backend/src/worker_flow.rs b/backend/src/worker_flow.rs index c8e4b26693..f0e7bcaec4 100644 --- a/backend/src/worker_flow.rs +++ b/backend/src/worker_flow.rs @@ -51,7 +51,6 @@ pub struct RetryStatus { pub struct Iterator { pub index: usize, pub itered: Vec, - pub args: Map, } #[derive(Serialize, Deserialize, Debug, Clone)] @@ -176,19 +175,21 @@ pub async fn update_flow_status_after_job_completion( (old_status.step, module.clone()) } _ => { - let forloop_jobs = match module { - FlowStatusModule::InProgress { forloop_jobs, .. } => forloop_jobs.clone(), - _ => None, + let (forloop_jobs, branch_chosen) = match module { + FlowStatusModule::InProgress { forloop_jobs, branch_chosen, .. } => { + (forloop_jobs.clone(), branch_chosen.clone()) + } + _ => (None, None), }; if success || (forloop_jobs.is_some() && skip_loop_failures) { ( old_status.step + 1, - FlowStatusModule::Success { job: job.id, forloop_jobs, branch_chosen: None }, + FlowStatusModule::Success { job: job.id, forloop_jobs, branch_chosen }, ) } else { ( old_status.step, - FlowStatusModule::Failure { job: job.id, forloop_jobs, branch_chosen: None }, + FlowStatusModule::Failure { job: job.id, forloop_jobs, branch_chosen }, ) } } @@ -312,7 +313,15 @@ pub async fn update_flow_status_after_job_completion( .await?; true } else { - match handle_flow(&flow_job, db, result.clone(), same_worker_tx.clone()).await { + match handle_flow( + &flow_job, + db, + result.clone(), + same_worker_tx.clone(), + worker_dir, + ) + .await + { Err(err) => { let _ = add_completed_job_error( db, @@ -488,8 +497,10 @@ async fn transform_input( match val { InputTransform::Static { value: _ } => (), InputTransform::Javascript { expr } => { - let previous_result = last_result.clone(); let flow_input = flow_args.clone().unwrap_or_else(|| json!({})); + println!("FLOW_INPUT 1 {:#?}", flow_input); + + let previous_result = flatten_previous_result(last_result.clone()); let context = vec![ ("params".to_string(), serde_json::json!(mapped)), ("previous_result".to_string(), previous_result), @@ -522,12 +533,32 @@ async fn transform_input( Ok(mapped) } +fn flatten_previous_result(last_result: serde_json::Value) -> serde_json::Value { + if last_result.is_object() + && last_result + .as_object() + .unwrap() + .contains_key("previous_result") + { + println!("FLOW_INPUT {:#?}", last_result); + last_result + .as_object() + .unwrap() + .get("previous_result") + .unwrap() + .clone() + } else { + last_result.clone() + } +} + #[instrument(level = "trace", skip_all)] pub async fn handle_flow( flow_job: &QueuedJob, db: &sqlx::Pool, last_result: serde_json::Value, same_worker_tx: Sender, + worker_dir: &str, ) -> anyhow::Result<()> { let value = flow_job .raw_flow @@ -536,6 +567,23 @@ pub async fn handle_flow( .to_owned(); let flow = serde_json::from_value::(value.to_owned())?; + if flow.modules.is_empty() { + let fake_job = QueuedJob { parent_job: Some(flow_job.id), ..flow_job.clone() }; + update_flow_status_after_job_completion( + db, + &fake_job, + true, + serde_json::json!({}), + None, + true, + same_worker_tx, + worker_dir, + false, + ) + .await?; + return Ok(()); + } + let status: FlowStatus = serde_json::from_value::(flow_job.flow_status.clone().unwrap_or_default()) .with_context(|| format!("parse flow status {}", flow_job.id))?; @@ -550,7 +598,7 @@ pub async fn handle_flow( async fn push_next_flow_job( flow_job: &QueuedJob, mut status: FlowStatus, - mut flow: FlowValue, + flow: FlowValue, db: &sqlx::Pool, mut last_result: serde_json::Value, same_worker_tx: Sender, @@ -823,29 +871,35 @@ async fn push_next_flow_job( _ => (), } - /* Don't evaluate `module.input_transforms` after iteration has begun. Instead, args are - * carried through the Iterator by the InProgress variant. */ - - #[rustfmt::skip] - let compute_input_transform = !( matches!(&module.value, FlowModuleValue::ForloopFlow { .. }) - && matches!(&status_module, FlowStatusModule::InProgress { .. })); - let mut transform_context: Option<(String, Vec)> = None; - let mut args = if compute_input_transform { - transform_context = Some(get_transform_context(&db, &flow_job, &status).await?); - let (token, steps) = transform_context.as_ref().unwrap(); - transform_input( - &flow_job.args, - last_result.clone(), - &module.input_transforms, - &flow_job.workspace_id, - &token, - steps.to_vec(), - resume_messages.as_slice(), - ) - .await? - } else { - Map::new() + let mut args = match module.value { + FlowModuleValue::Script { .. } | FlowModuleValue::RawScript { .. } => { + transform_context = Some(get_transform_context(&db, &flow_job, &status).await?); + let (token, steps) = transform_context.as_ref().unwrap(); + transform_input( + &flow_job.args, + last_result.clone(), + &module.input_transforms, + &flow_job.workspace_id, + &token, + steps.to_vec(), + resume_messages.as_slice(), + ) + .await? + } + _ => { + /* embedded flow input is augmented with embedding flow input */ + if let Some(value) = &flow_job.args { + value + .as_object() + .ok_or_else(|| { + Error::BadRequest(format!("Expected an object value, found: {value:?}")) + })? + .clone() + } else { + Map::new() + } + } }; let next_flow_transform = compute_next_flow_transform( @@ -857,60 +911,17 @@ async fn push_next_flow_job( &status, &status_module, last_result.clone(), - &mut args, ) .await?; let (job_payload, next_status) = match next_flow_transform { NextFlowTransform::Continue(job_payload, next_state) => (job_payload, next_state), - NextFlowTransform::BranchChosen(branch, modules) => { - let added_flow_statuses = vec![FlowStatusModule::WaitingForPriorSteps; modules.len()]; - - let index = (status.step + 1) as usize; - status.modules.splice(index..index, added_flow_statuses); - flow.modules.splice(index..index, modules); - - let mut tx = db.begin().await?; - sqlx::query!( - " - UPDATE queue - SET flow_status = JSONB_SET(flow_status, ARRAY['modules'], $1), - raw_flow = JSONB_SET(raw_flow, ARRAY['modules'], $2) - WHERE id = $3 - ", - json!(status.modules), - json!(flow.modules), - flow_job.id - ) - .execute(&mut tx) - .await?; - - return jump_to_next_step( - status.step, - i, - &flow_job.id, - flow, - tx, - &db, - FlowStatusModule::Success { - job: flow_job.id, - forloop_jobs: None, - branch_chosen: Some(branch), - }, - last_result, - same_worker_tx, - ) - .await; - } NextFlowTransform::EmptyIterator => { - let tx = db.begin().await?; - return jump_to_next_step( status.step, i, &flow_job.id, flow.clone(), - tx, &db, FlowStatusModule::Success { job: flow_job.id, @@ -927,6 +938,19 @@ async fn push_next_flow_job( let continue_on_same_worker = flow.same_worker && module.suspend.is_none() && module.sleep.is_none(); + match &next_status { + NextStatus::NextLoopIteration(NextIteration { new_args, .. }) => { + args.extend(new_args.clone()) + } + NextStatus::BranchChosen(_) => { + args.insert( + "previous_result".to_string(), + flatten_previous_result(last_result), + ); + } + _ => (), + }; + /* Finally, push the job into the queue */ let tx = db.begin().await?; @@ -946,17 +970,25 @@ async fn push_next_flow_job( .await?; let new_status = match next_status { - NextStatus::NextLoopIteration(NextIteration { index, itered, mut forloop_jobs }) => { + NextStatus::NextLoopIteration(NextIteration { + index, itered, mut forloop_jobs, .. + }) => { forloop_jobs.push(uuid); FlowStatusModule::InProgress { job: uuid, - iterator: Some(Iterator { index, itered, args }), + iterator: Some(Iterator { index, itered }), forloop_jobs: Some(forloop_jobs), branch_chosen: None, } } - _ => FlowStatusModule::WaitingForExecutor { job: uuid }, + NextStatus::BranchChosen(branch) => FlowStatusModule::InProgress { + job: uuid, + iterator: None, + forloop_jobs: None, + branch_chosen: Some(branch), + }, + NextStatus::NextStep => FlowStatusModule::WaitingForExecutor { job: uuid }, }; sqlx::query( @@ -983,17 +1015,18 @@ async fn push_next_flow_job( return Ok(()); } -async fn jump_to_next_step<'c>( +async fn jump_to_next_step( status_step: i32, i: usize, job_id: &Uuid, flow: FlowValue, - mut tx: sqlx::Transaction<'c, sqlx::Postgres>, db: &DB, status_module: FlowStatusModule, last_result: serde_json::Value, same_worker_tx: Sender, ) -> anyhow::Result<()> { + let mut tx = db.begin().await?; + let next_step = i .checked_add(1) .filter(|i| (..flow.modules.len()).contains(i)); @@ -1049,6 +1082,7 @@ struct NextIteration { index: usize, itered: Vec, forloop_jobs: Vec, + new_args: Map, } enum LoopStatus { @@ -1058,13 +1092,13 @@ enum LoopStatus { enum NextStatus { NextStep, + BranchChosen(BranchChosen), NextLoopIteration(NextIteration), } enum NextFlowTransform { EmptyIterator, Continue(JobPayload, NextStatus), - BranchChosen(BranchChosen, Vec), } async fn compute_next_flow_transform( @@ -1076,7 +1110,7 @@ async fn compute_next_flow_transform( status: &FlowStatus, status_module: &FlowStatusModule, last_result: serde_json::Value, - args: &mut Map, + // args: &mut Map, ) -> error::Result { match &module.value { FlowModuleValue::Script { path: script_path } => Ok(NextFlowTransform::Continue( @@ -1096,6 +1130,8 @@ async fn compute_next_flow_transform( } /* forloop modules are expected set `iter: { value: Value, index: usize }` as job arguments */ FlowModuleValue::ForloopFlow { modules, iterator, .. } => { + let new_args: &mut Map = &mut Map::new(); + let next_loop_status = match status_module { FlowStatusModule::WaitingForPriorSteps => { let (token, steps) = if let Some(x) = transform_context { @@ -1126,12 +1162,13 @@ async fn compute_next_flow_transform( })?; if let Some(first) = itered.first() { - args.insert("iter".to_string(), json!({ "index": 0, "value": first })); + new_args.insert("iter".to_string(), json!({ "index": 0, "value": first })); LoopStatus::NextIteration(NextIteration { index: 0, itered, forloop_jobs: vec![], + new_args: new_args.clone(), }) } else { LoopStatus::EmptyIterator @@ -1139,7 +1176,7 @@ async fn compute_next_flow_transform( } FlowStatusModule::InProgress { - iterator: Some(Iterator { itered, index, args: iterator_args }), + iterator: Some(Iterator { itered, index }), forloop_jobs: Some(forloop_jobs), .. } => { @@ -1154,13 +1191,13 @@ async fn compute_next_flow_transform( format!("could not iterate index {index} of {itered:?}") })?; - args.extend(iterator_args.clone()); - args.insert("iter".to_string(), json!({ "index": index, "value": next })); + new_args.insert("iter".to_string(), json!({ "index": index, "value": next })); LoopStatus::NextIteration(NextIteration { index, itered: itered.clone(), forloop_jobs: forloop_jobs.clone(), + new_args: new_args.clone(), }) } @@ -1171,31 +1208,17 @@ async fn compute_next_flow_transform( match next_loop_status { LoopStatus::EmptyIterator => Ok(NextFlowTransform::EmptyIterator), - LoopStatus::NextIteration(ns) => { - /* embedded flow input is augmented with embedding flow input */ - if let Some(value) = &flow_job.args { - value - .as_object() - .ok_or_else(|| { - Error::BadRequest(format!( - "Expected an object value, found: {value:?}" - )) - }) - .map(|map| args.extend(map.clone()))?; - } - - Ok(NextFlowTransform::Continue( - JobPayload::RawFlow { - value: FlowValue { - modules: (*modules).clone(), - failure_module: flow.failure_module.clone(), - same_worker: flow.same_worker, - }, - path: Some(format!("{}/loop-{}", flow_job.script_path(), status.step)), + LoopStatus::NextIteration(ns) => Ok(NextFlowTransform::Continue( + JobPayload::RawFlow { + value: FlowValue { + modules: (*modules).clone(), + failure_module: flow.failure_module.clone(), + same_worker: flow.same_worker, }, - NextStatus::NextLoopIteration(ns), - )) - } + path: Some(format!("{}/loop-{}", flow_job.script_path(), status.step)), + }, + NextStatus::NextLoopIteration(ns), + )), } } FlowModuleValue::BranchOne { branches, default, .. } => { @@ -1224,9 +1247,17 @@ async fn compute_next_flow_transform( &default }; - // match inner_flow_transform {} - - Ok(NextFlowTransform::BranchChosen(branch, modules.clone())) + Ok(NextFlowTransform::Continue( + JobPayload::RawFlow { + value: FlowValue { + modules: (*modules).clone(), + failure_module: flow.failure_module.clone(), + same_worker: flow.same_worker, + }, + path: Some(format!("{}/loop-{}", flow_job.script_path(), status.step)), + }, + NextStatus::BranchChosen(branch), + )) } FlowModuleValue::BranchAll { branches: _branches, .. } => { todo!() diff --git a/frontend/src/lib/components/flows/content/FlowLoop.svelte b/frontend/src/lib/components/flows/content/FlowLoop.svelte index efbc0a1572..da269c47fb 100644 --- a/frontend/src/lib/components/flows/content/FlowLoop.svelte +++ b/frontend/src/lib/components/flows/content/FlowLoop.svelte @@ -23,12 +23,8 @@ export let index: number let editor: SimpleEditor | undefined = undefined - let monacos: { [id: string]: SimpleEditor } = {} - let selected: string = 'retries' - let inputTransformName = '' - $: mod = $flowStore.value.modules[index] $: pickableProperties = getStepPropPicker( @@ -90,62 +86,6 @@ right: 'Skip failures' }} /> - Pass specific flow context as loop flow input -
- - {#each Object.keys(mod.input_transforms) as key} -
- {key} - - -
-
- {#if mod.input_transforms[key].type == 'javascript'} - { - monacos[key]?.insertAtCursor(detail) - }} - > - - - {:else} -
- {/each} {/if}