From 7f2cbb40f67fc5f46cb02216a3f5b0e829891ea8 Mon Sep 17 00:00:00 2001 From: Ruben Fiszel Date: Sat, 13 Jul 2024 18:53:24 +0200 Subject: [PATCH] feat: improve flow status viewer (show branch chosen + all iterations in for loop) (#4074) * all * all * all * all * all --- ...ea026d8e94a3cb9cdcfa5fde19343acf81ecc.json | 17 + backend/substitute_ee_code.sh | 2 +- backend/windmill-common/src/flow_status.rs | 28 + backend/windmill-queue/src/jobs.rs | 20 +- backend/windmill-worker/src/worker_flow.rs | 191 ++++++- .../lib/components/FlowStatusViewer.svelte | 7 +- .../components/FlowStatusViewerInner.svelte | 502 ++++++++++++------ .../components/flows/map/FlowJobsMenu.svelte | 75 +++ .../flows/map/FlowModuleSchemaItem.svelte | 4 +- .../lib/components/flows/map/MapItem.svelte | 18 + .../components/flows/map/VirtualItem.svelte | 9 +- .../src/lib/components/graph/FlowGraph.svelte | 118 ++-- frontend/src/lib/components/graph/model.ts | 7 +- .../svelvet/container/views/GraphView.svelte | 2 + .../src/lib/components/runs/JobLoader.svelte | 1 + openflow.openapi.yaml | 4 + 16 files changed, 749 insertions(+), 256 deletions(-) create mode 100644 backend/.sqlx/query-3b4b62161a5197f37850c8c4197ea026d8e94a3cb9cdcfa5fde19343acf81ecc.json create mode 100644 frontend/src/lib/components/flows/map/FlowJobsMenu.svelte diff --git a/backend/.sqlx/query-3b4b62161a5197f37850c8c4197ea026d8e94a3cb9cdcfa5fde19343acf81ecc.json b/backend/.sqlx/query-3b4b62161a5197f37850c8c4197ea026d8e94a3cb9cdcfa5fde19343acf81ecc.json new file mode 100644 index 0000000000..367aae4a31 --- /dev/null +++ b/backend/.sqlx/query-3b4b62161a5197f37850c8c4197ea026d8e94a3cb9cdcfa5fde19343acf81ecc.json @@ -0,0 +1,17 @@ +{ + "db_name": "PostgreSQL", + "query": "UPDATE queue SET flow_status = JSONB_SET(flow_status, ARRAY['modules', $1::TEXT, 'flow_jobs_success', $3::TEXT], $4) WHERE id = $2", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "Text", + "Uuid", + "Text", + "Jsonb" + ] + }, + "nullable": [] + }, + "hash": "3b4b62161a5197f37850c8c4197ea026d8e94a3cb9cdcfa5fde19343acf81ecc" +} diff --git a/backend/substitute_ee_code.sh b/backend/substitute_ee_code.sh index d7f98bf0fe..60033cdc5d 100755 --- a/backend/substitute_ee_code.sh +++ b/backend/substitute_ee_code.sh @@ -71,7 +71,7 @@ if [ "$REVERT" == "YES" ]; then ce_file="${ee_file/${EE_CODE_DIR}/.}" ce_file="${root_dirpath}/backend/${ce_file}" if [ "$REVERT_PREVIOUS" == "YES" ]; then - git checkout HEAD@{5} ${ce_file} || true + git checkout HEAD@{9} ${ce_file} || true else git restore --staged ${ce_file} || true git restore ${ce_file} || true diff --git a/backend/windmill-common/src/flow_status.rs b/backend/windmill-common/src/flow_status.rs index 7fde7aa298..dd11b8e9af 100644 --- a/backend/windmill-common/src/flow_status.rs +++ b/backend/windmill-common/src/flow_status.rs @@ -119,6 +119,7 @@ struct UntaggedFlowStatusModule { job: Option, iterator: Option, flow_jobs: Option>, + flow_jobs_success: Option>>, branch_chosen: Option, branchall: Option, parallel: Option, @@ -150,6 +151,8 @@ pub enum FlowStatusModule { #[serde(skip_serializing_if = "Option::is_none")] flow_jobs: Option>, #[serde(skip_serializing_if = "Option::is_none")] + flow_jobs_success: Option>>, + #[serde(skip_serializing_if = "Option::is_none")] branch_chosen: Option, #[serde(skip_serializing_if = "Option::is_none")] branchall: Option, @@ -164,6 +167,8 @@ pub enum FlowStatusModule { #[serde(skip_serializing_if = "Option::is_none")] flow_jobs: Option>, #[serde(skip_serializing_if = "Option::is_none")] + flow_jobs_success: Option>>, + #[serde(skip_serializing_if = "Option::is_none")] branch_chosen: Option, #[serde(default)] #[serde(skip_serializing_if = "Vec::is_empty")] @@ -177,6 +182,8 @@ pub enum FlowStatusModule { #[serde(skip_serializing_if = "Option::is_none")] flow_jobs: Option>, #[serde(skip_serializing_if = "Option::is_none")] + flow_jobs_success: Option>>, + #[serde(skip_serializing_if = "Option::is_none")] branch_chosen: Option, #[serde(skip_serializing_if = "Vec::is_empty")] failed_retries: Vec, @@ -225,6 +232,7 @@ impl<'de> Deserialize<'de> for FlowStatusModule { .ok_or_else(|| serde::de::Error::missing_field("job"))?, iterator: untagged.iterator, flow_jobs: untagged.flow_jobs, + flow_jobs_success: untagged.flow_jobs_success, branch_chosen: untagged.branch_chosen, branchall: untagged.branchall, parallel: untagged.parallel.unwrap_or(false), @@ -238,6 +246,7 @@ impl<'de> Deserialize<'de> for FlowStatusModule { .job .ok_or_else(|| serde::de::Error::missing_field("job"))?, flow_jobs: untagged.flow_jobs, + flow_jobs_success: untagged.flow_jobs_success, branch_chosen: untagged.branch_chosen, approvers: untagged.approvers.unwrap_or_default(), failed_retries: untagged.failed_retries.unwrap_or_default(), @@ -250,6 +259,7 @@ impl<'de> Deserialize<'de> for FlowStatusModule { .job .ok_or_else(|| serde::de::Error::missing_field("job"))?, flow_jobs: untagged.flow_jobs, + flow_jobs_success: untagged.flow_jobs_success, branch_chosen: untagged.branch_chosen, failed_retries: untagged.failed_retries.unwrap_or_default(), }), @@ -295,6 +305,24 @@ impl FlowStatusModule { } } + pub fn branch_chosen(&self) -> Option { + match self { + FlowStatusModule::InProgress { branch_chosen, .. } => branch_chosen.clone(), + FlowStatusModule::Success { branch_chosen, .. } => branch_chosen.clone(), + FlowStatusModule::Failure { branch_chosen, .. } => branch_chosen.clone(), + _ => None, + } + } + + pub fn flow_jobs_success(&self) -> Option>> { + match self { + FlowStatusModule::InProgress { flow_jobs_success, .. } => flow_jobs_success.clone(), + FlowStatusModule::Success { flow_jobs_success, .. } => flow_jobs_success.clone(), + FlowStatusModule::Failure { flow_jobs_success, .. } => flow_jobs_success.clone(), + _ => None, + } + } + pub fn job_result(&self) -> Option { self.flow_jobs() .map(JobResult::ListJob) diff --git a/backend/windmill-queue/src/jobs.rs b/backend/windmill-queue/src/jobs.rs index 24c65642b0..5b3728ee09 100644 --- a/backend/windmill-queue/src/jobs.rs +++ b/backend/windmill-queue/src/jobs.rs @@ -4011,11 +4011,16 @@ async fn restarted_flows_resolution( } let mut new_flow_jobs = module.flow_jobs().unwrap_or_default(); new_flow_jobs.truncate(branch_or_iteration_n); + let mut new_flow_jobs_success = module.flow_jobs_success(); + if let Some(new_flow_jobs_success) = new_flow_jobs_success.as_mut() { + new_flow_jobs_success.truncate(branch_or_iteration_n); + } truncated_modules.push(FlowStatusModule::InProgress { id: module.id(), job: new_flow_jobs[new_flow_jobs.len() - 1], // set to last finished job from completed flow iterator: None, flow_jobs: Some(new_flow_jobs), + flow_jobs_success: new_flow_jobs_success, branch_chosen: None, branchall: Some(BranchAllStatus { branch: branch_or_iteration_n - 1, // Doing minus one here as this variable reflects the latest finished job in the iteration @@ -4043,7 +4048,10 @@ async fn restarted_flows_resolution( } let mut new_flow_jobs = module.flow_jobs().unwrap_or_default(); new_flow_jobs.truncate(branch_or_iteration_n); - + let mut new_flow_jobs_success = module.flow_jobs_success(); + if let Some(new_flow_jobs_success) = new_flow_jobs_success.as_mut() { + new_flow_jobs_success.truncate(branch_or_iteration_n); + } truncated_modules.push(FlowStatusModule::InProgress { id: module.id(), job: new_flow_jobs[new_flow_jobs.len() - 1], // set to last finished job from completed flow @@ -4052,6 +4060,7 @@ async fn restarted_flows_resolution( itered: vec![], // Setting itered to empty array here, such that input transforms will be re-computed by worker_flows }), flow_jobs: Some(new_flow_jobs), + flow_jobs_success: new_flow_jobs_success, branch_chosen: None, branchall: None, parallel: parallel, @@ -4074,14 +4083,7 @@ async fn restarted_flows_resolution( // else we simply "transfer" the module from the completed flow to the new one if it's a success step_n = step_n + 1; match module.clone() { - FlowStatusModule::Success { - id: _, - job: _, - flow_jobs: _, - branch_chosen: _, - approvers: _, - failed_retries: _, - } => Ok(truncated_modules.push(module)), + FlowStatusModule::Success { .. } => Ok(truncated_modules.push(module)), _ => Err(Error::InternalErr(format!( "Flow cannot be restarted from a non successful module", ))), diff --git a/backend/windmill-worker/src/worker_flow.rs b/backend/windmill-worker/src/worker_flow.rs index 64d37748c7..3787dac0fd 100644 --- a/backend/windmill-worker/src/worker_flow.rs +++ b/backend/windmill-worker/src/worker_flow.rs @@ -326,10 +326,22 @@ pub async fn update_flow_status_after_job_completion_internal< branchall, parallel, flow_jobs: Some(jobs), + flow_jobs_success, .. } if *parallel => { let (nindex, len) = match (iterator, branchall) { (Some(Iterator { itered, .. }), _) => { + set_success_in_flow_job_success( + flow_jobs_success, + jobs, + job_id_for_status, + &old_status, + flow, + success, + &mut tx, + ) + .await?; + let nindex = sqlx::query_scalar!( "UPDATE queue SET flow_status = JSONB_SET(flow_status, ARRAY['modules', $1::TEXT, 'iterator', 'index'], ((flow_status->'modules'->$1::int->'iterator'->>'index')::int + 1)::text::jsonb) @@ -353,6 +365,16 @@ pub async fn update_flow_status_after_job_completion_internal< (nindex, itered.len() as i32) } (_, Some(BranchAllStatus { len, .. })) => { + set_success_in_flow_job_success( + flow_jobs_success, + jobs, + job_id_for_status, + &old_status, + flow, + success, + &mut tx, + ) + .await?; let nindex = sqlx::query_scalar!( "UPDATE queue SET flow_status = JSONB_SET(flow_status, ARRAY['modules', $1::TEXT, 'branchall', 'branch'], ((flow_status->'modules'->$1::int->'branchall'->>'branch')::int + 1)::text::jsonb) @@ -376,6 +398,16 @@ pub async fn update_flow_status_after_job_completion_internal< )))?, }; if nindex == len { + let mut flow_jobs_success = flow_jobs_success.clone(); + if let Some(flow_job_success) = flow_jobs_success.as_mut() { + let position = jobs.iter().position(|x| x == job_id_for_status); + if let Some(position) = position { + if position < flow_job_success.len() { + flow_job_success[position] = Some(success); + } + } + } + let new_status = if skip_loop_failures || sqlx::query_scalar!( "SELECT success FROM completed_job WHERE id = ANY($1)", @@ -396,6 +428,7 @@ pub async fn update_flow_status_after_job_completion_internal< id: module_status.id(), job: job_id_for_status.clone(), flow_jobs: Some(jobs.clone()), + flow_jobs_success: flow_jobs_success.clone(), branch_chosen: None, approvers: vec![], failed_retries: vec![], @@ -406,6 +439,7 @@ pub async fn update_flow_status_after_job_completion_internal< id: module_status.id(), job: job_id_for_status.clone(), flow_jobs: Some(jobs.clone()), + flow_jobs_success: flow_jobs_success.clone(), branch_chosen: None, failed_retries: vec![], } @@ -478,18 +512,49 @@ pub async fn update_flow_status_after_job_completion_internal< } FlowStatusModule::InProgress { iterator: Some(windmill_common::flow_status::Iterator { index, itered, .. }), + flow_jobs_success, + flow_jobs, while_loop, .. } if (*while_loop || (*index + 1 < itered.len()) && (success || skip_loop_failures)) && !stop_early => { + if let Some(jobs) = flow_jobs { + set_success_in_flow_job_success( + flow_jobs_success, + jobs, + job_id_for_status, + &old_status, + flow, + success, + &mut tx, + ) + .await?; + } + (false, None) } FlowStatusModule::InProgress { branchall: Some(BranchAllStatus { branch, len, .. }), + flow_jobs_success, + flow_jobs, .. - } if branch.to_owned() < len - 1 && (success || skip_branch_failure) => (false, None), + } if branch.to_owned() < len - 1 && (success || skip_branch_failure) => { + if let Some(jobs) = flow_jobs { + set_success_in_flow_job_success( + flow_jobs_success, + jobs, + job_id_for_status, + &old_status, + flow, + success, + &mut tx, + ) + .await?; + } + (false, None) + } _ => { if stop_early && matches!( @@ -500,18 +565,21 @@ pub async fn update_flow_status_after_job_completion_internal< // if we're stopping early inside a loop, we just want to break the loop instead stop_early = false; } - let (flow_jobs, branch_chosen) = match module_status { - FlowStatusModule::InProgress { flow_jobs, branch_chosen, .. } => { - (flow_jobs.clone(), branch_chosen.clone()) + + let flow_jobs = module_status.flow_jobs(); + let branch_chosen = module_status.branch_chosen(); + let mut flow_jobs_success = module_status.flow_jobs_success(); + + if let (Some(flow_job_success), Some(flow_jobs)) = + (flow_jobs_success.as_mut(), flow_jobs.as_ref()) + { + let position = flow_jobs.iter().position(|x| x == job_id_for_status); + if let Some(position) = position { + if position < flow_job_success.len() { + flow_job_success[position] = Some(success); + } } - FlowStatusModule::Success { flow_jobs, branch_chosen, .. } => { - (flow_jobs.clone(), branch_chosen.clone()) - } - FlowStatusModule::Failure { flow_jobs, branch_chosen, .. } => { - (flow_jobs.clone(), branch_chosen.clone()) - } - _ => (None, None), - }; + } if success || (flow_jobs.is_some() && (skip_loop_failures || skip_branch_failure)) { success = true; ( @@ -520,6 +588,7 @@ pub async fn update_flow_status_after_job_completion_internal< id: module_status.id(), job: job_id_for_status.clone(), flow_jobs, + flow_jobs_success, branch_chosen, approvers: vec![], failed_retries: old_status.retry.failed_jobs.clone(), @@ -556,6 +625,7 @@ pub async fn update_flow_status_after_job_completion_internal< id: module_status.id(), job: job_id_for_status.clone(), flow_jobs, + flow_jobs_success, branch_chosen, failed_retries: old_status.retry.failed_jobs.clone(), }), @@ -903,6 +973,38 @@ pub async fn update_flow_status_after_job_completion_internal< } } +async fn set_success_in_flow_job_success<'c, R: rsmq_async::RsmqConnection + Send>( + flow_jobs_success: &Option>>, + flow_jobs: &Vec, + job_id_for_status: &Uuid, + old_status: &FlowStatus, + flow: Uuid, + success: bool, + tx: &mut QueueTransaction<'c, R>, +) -> error::Result<()> { + let flow_jobs_success = flow_jobs_success.clone(); + + if flow_jobs_success.is_some() { + let position = flow_jobs.iter().position(|x| x == job_id_for_status); + if let Some(position) = position { + sqlx::query!( + "UPDATE queue SET flow_status = JSONB_SET(flow_status, ARRAY['modules', $1::TEXT, 'flow_jobs_success', $3::TEXT], $4) WHERE id = $2", + old_status.step as i32, + flow, + position as i32, + json!(success) + ) + .execute(tx) + .await.map_err(|e| { + Error::InternalErr(format!( + "error while setting flow_jobs_success: {e:#}" + )) + })?; + } + } + Ok(()) +} + async fn retrieve_flow_jobs_results( db: &DB, w_id: &str, @@ -1975,6 +2077,7 @@ async fn push_next_flow_job id: status_module.id(), job: Uuid::nil(), flow_jobs: Some(vec![]), + flow_jobs_success: Some(vec![]), branch_chosen: None, approvers: vec![], failed_retries: vec![], @@ -2304,17 +2407,29 @@ async fn push_next_flow_job let first_uuid = uuids[0]; let new_status = match next_status { NextStatus::NextLoopIteration { - next: ForloopNextIteration { index, itered, mut flow_jobs, while_loop, .. }, + next: + ForloopNextIteration { + index, + itered, + mut flow_jobs, + while_loop, + mut flow_jobs_success, + .. + }, .. } => { let uuid = one_uuid?; flow_jobs.push(uuid); + if let Some(flow_jobs_success) = &mut flow_jobs_success { + flow_jobs_success.push(None); + } FlowStatusModule::InProgress { job: uuid, iterator: Some(windmill_common::flow_status::Iterator { index, itered }), flow_jobs: Some(flow_jobs), + flow_jobs_success, branch_chosen: None, branchall: None, id: status_module.id(), @@ -2325,6 +2440,7 @@ async fn push_next_flow_job NextStatus::AllFlowJobs { iterator, branchall, .. } => FlowStatusModule::InProgress { job: flow_job.id, iterator, + flow_jobs_success: Some(vec![None; uuids.len()]), flow_jobs: Some(uuids.clone()), branch_chosen: None, branchall, @@ -2332,14 +2448,22 @@ async fn push_next_flow_job parallel: true, while_loop: false, }, - NextStatus::NextBranchStep(NextBranch { mut flow_jobs, status, .. }) => { + NextStatus::NextBranchStep(NextBranch { + mut flow_jobs, + status, + mut flow_jobs_success, + .. + }) => { let uuid = one_uuid?; flow_jobs.push(uuid); - + if let Some(flow_jobs_success) = &mut flow_jobs_success { + flow_jobs_success.push(None); + } FlowStatusModule::InProgress { job: uuid, iterator: None, flow_jobs: Some(flow_jobs), + flow_jobs_success, branch_chosen: None, branchall: Some(status), id: status_module.id(), @@ -2352,6 +2476,7 @@ async fn push_next_flow_job job: one_uuid?, iterator: None, flow_jobs: None, + flow_jobs_success: None, branch_chosen: Some(branch), branchall: None, id: status_module.id(), @@ -2494,6 +2619,7 @@ struct ForloopNextIteration { index: usize, itered: Vec>, flow_jobs: Vec, + flow_jobs_success: Option>>, new_args: Iter, while_loop: bool, } @@ -2508,6 +2634,7 @@ enum ForLoopStatus { struct NextBranch { status: BranchAllStatus, flow_jobs: Vec, + flow_jobs_success: Option>>, } #[derive(Debug)] @@ -2656,11 +2783,13 @@ async fn compute_next_flow_transform( FlowModuleValue::WhileloopFlow { modules, .. } => { // if it's a simple single step flow, we will collapse it as an optimization and need to pass flow_input as an arg let is_simple = is_simple_modules(modules, flow); - let flow_jobs = match status_module { - FlowStatusModule::InProgress { flow_jobs: Some(flow_jobs), .. } => { - flow_jobs.clone() - } - _ => vec![], + let (flow_jobs, flow_jobs_success) = match status_module { + FlowStatusModule::InProgress { + flow_jobs: Some(flow_jobs), + flow_jobs_success, + .. + } => (flow_jobs.clone(), flow_jobs_success.clone()), + _ => (vec![], Some(vec![])), }; let next_loop_idx = flow_jobs.len(); next_loop_iteration( @@ -2669,7 +2798,8 @@ async fn compute_next_flow_transform( ForloopNextIteration { index: next_loop_idx, itered: vec![], - flow_jobs: flow_jobs.clone(), + flow_jobs: flow_jobs, + flow_jobs_success: flow_jobs_success, new_args: Iter { index: next_loop_idx as i32, value: windmill_common::worker::to_raw_value(&next_loop_idx), @@ -2870,7 +3000,7 @@ async fn compute_next_flow_transform( )) } FlowModuleValue::BranchAll { branches, parallel, .. } => { - let (branch_status, flow_jobs) = match status_module { + let (branch_status, flow_jobs, flow_jobs_success) = match status_module { FlowStatusModule::WaitingForPriorSteps { .. } | FlowStatusModule::WaitingForEvents { .. } | FlowStatusModule::WaitingForExecutor { .. } => { @@ -2928,16 +3058,22 @@ async fn compute_next_flow_transform( }, )); } else { - (BranchAllStatus { branch: 0, len: branches.len() }, vec![]) + ( + BranchAllStatus { branch: 0, len: branches.len() }, + vec![], + Some(vec![]), + ) } } FlowStatusModule::InProgress { branchall: Some(BranchAllStatus { branch, len }), flow_jobs: Some(flow_jobs), + flow_jobs_success, .. } if !*parallel => ( BranchAllStatus { branch: branch + 1, len: len.clone() }, flow_jobs.clone(), + flow_jobs_success.clone(), ), _ => Err(Error::BadRequest(format!( @@ -2986,7 +3122,11 @@ async fn compute_next_flow_transform( delete_after_use: delete_after_use, timeout: None, }), - NextStatus::NextBranchStep(NextBranch { status: branch_status, flow_jobs }), + NextStatus::NextBranchStep(NextBranch { + status: branch_status, + flow_jobs, + flow_jobs_success, + }), )) } } @@ -3136,6 +3276,7 @@ async fn next_forloop_status( index: 0, itered, flow_jobs: vec![], + flow_jobs_success: Some(vec![]), new_args: iter, while_loop: false, }) @@ -3147,6 +3288,7 @@ async fn next_forloop_status( FlowStatusModule::InProgress { iterator: Some(windmill_common::flow_status::Iterator { itered, index }), flow_jobs: Some(flow_jobs), + flow_jobs_success, .. } if !*parallel => { let itered_new = if itered.is_empty() { @@ -3195,6 +3337,7 @@ async fn next_forloop_status( index, itered: itered_new.clone(), flow_jobs: flow_jobs.clone(), + flow_jobs_success: flow_jobs_success.clone(), new_args: Iter { index: index as i32, value: next.to_owned() }, while_loop: false, }) diff --git a/frontend/src/lib/components/FlowStatusViewer.svelte b/frontend/src/lib/components/FlowStatusViewer.svelte index 379672cdbd..d9daaaddfa 100644 --- a/frontend/src/lib/components/FlowStatusViewer.svelte +++ b/frontend/src/lib/components/FlowStatusViewer.svelte @@ -46,11 +46,12 @@ { - if (detail.script_path != lastScriptPath && detail.script_path) { - lastScriptPath = detail.script_path + let { job } = detail + if (job.script_path != lastScriptPath && job.script_path) { + lastScriptPath = job.script_path loadOwner(lastScriptPath ?? '') } - dispatch('jobsLoaded', detail) + dispatch('jobsLoaded', job) }} globalDurationStatuses={[]} globalModuleStates={[]} diff --git a/frontend/src/lib/components/FlowStatusViewerInner.svelte b/frontend/src/lib/components/FlowStatusViewerInner.svelte index 6d836f07a3..66993d5f38 100644 --- a/frontend/src/lib/components/FlowStatusViewerInner.svelte +++ b/frontend/src/lib/components/FlowStatusViewerInner.svelte @@ -4,8 +4,6 @@ type Job, JobService, type FlowStatus, - type CompletedJob, - type QueuedJob, type FlowModuleValue, type FlowModule } from '$lib/gen' @@ -46,10 +44,10 @@ | { moduleId: string flowJobs: string[] + flowJobsSuccess: (boolean | undefined)[] length: number } | undefined = undefined - export let job: Job | undefined = undefined //only useful when forloops are optimized and the job doesn't contain the mod id anymore export let innerModule: FlowModuleValue | undefined = undefined @@ -62,37 +60,48 @@ export let globalModuleStates: Writable>[] export let globalDurationStatuses: Writable>[] + export let globalRefreshes: Record< + string, + (loopJob: { index: number; job: string }) => Promise + > = {} export let childFlow: boolean = false export let reducedPolling = false export let wideResults = false - let jobResults: any[] = [] - let jobFailures: boolean[] = [] + let jobResults: any[] = + flowJobIds?.flowJobs?.map((x, id) => `iter #${id + 1} not loaded by frontend yet`) ?? [] - let forloop_selected = '' let retry_selected = '' let timeout: NodeJS.Timeout let localModuleStates: Writable> = writable({}) let localDurationStatuses: Writable> = writable({}) - let lastSize = 0 - $: { - let len = (flowJobIds?.flowJobs ?? []).length - if (len != lastSize) { - updateForloop(len) - } - } + export let job: Job | undefined = undefined - function setModuleState(key: string, value: GraphModuleState) { - if (!deepEqual($localModuleStates[key], value)) { - // console.log('Setting module state', key, value) - $localModuleStates[key] = value + // let lastSize = 0 + // $: { + // let len = (flowJobIds?.flowJobs ?? []).length + // if (len != lastSize) { + // updateForloop(len) + // } + // } - globalModuleStates.forEach((s) => { + function setModuleState( + key: string, + value: Partial, + force?: boolean, + keepType?: boolean + ) { + let newValue = { ...($localModuleStates[key] ?? {}), ...value } + if (!deepEqual($localModuleStates[key], value) || force) { + ;[localModuleStates, ...globalModuleStates].forEach((s) => { s.update((x) => { - x[key] = value + if (keepType && (x[key]?.type == 'Success' || x[key]?.type == 'Failure')) { + newValue.type = x[key].type + } + x[key] = newValue return x }) }) @@ -105,6 +114,7 @@ globalDurationStatuses.forEach((s) => { s.update((x) => { x[key].byJob[id] = value + return x }) }) @@ -126,11 +136,6 @@ ) } - function updateForloop(len: number) { - forloop_selected = flowJobIds?.flowJobs[len - 1] ?? '' - lastSize = len - } - let innerModules: FlowStatusModule[] = [] function updateStatus(status: FlowStatus) { @@ -182,18 +187,76 @@ parent_module: mod['parent_module'], args: job?.args } - if (!deepEqual(newState, $localModuleStates[mod.id ?? ''])) { - setModuleState(mod.id ?? '', newState) - } + setModuleState(mod.id ?? '', newState) }) .catch((e) => { console.error(`Could not load inner module for job ${mod.job}`, e) }) + } else if ( + mod.flow_jobs && + (mod.type == 'Success' || mod.type == 'Failure') && + !['Success', 'Failure'].includes($localModuleStates?.[mod.id ?? '']?.type) + ) { + // console.log(mod.id, 'FOO') + setModuleState( + mod.id ?? '', + { + type: mod.type + }, + true + ) + } + if (mod.branch_chosen) { + setModuleState( + mod.id ?? '', + { + branchChosen: + mod.branch_chosen.type == 'default' ? 0 : (mod.branch_chosen.branch ?? 0) + 1 + }, + true + ) } }) } } + let recursiveRefresh: Record Promise> = {} + + export async function refresh( + root: boolean, + loopJob: { index: number; job: string } | undefined + ) { + let modId = flowJobIds?.moduleId + + if (!loopJob) { + loopJob = { + index: $localModuleStates[modId ?? '']?.selectedForloopIndex ?? 0, + job: $localModuleStates[modId ?? '']?.selectedForloop ?? '' + } + } + + let last = root ? undefined : flowJobIds?.flowJobs?.[flowJobIds?.flowJobs.length - 1] + + Object.entries(recursiveRefresh).forEach(([key, v]) => { + if (modId) { + if ((root && key == loopJob?.job) || key == last) { + v(false) + } else { + } + } else { + v(false) + } + }) + let njob = flowJobIds + ? root && modId + ? storedListJobs?.[loopJob.job] + : storedListJobs[flowJobIds.length - 1] + : job + if (njob) { + dispatch('jobsLoaded', { job: njob, force: true }) + } + } + let errorCount = 0 let notAnonynmous = false async function loadJobInProgress() { @@ -208,7 +271,7 @@ if (!deepEqual(job, newJob)) { job = newJob job?.flow_status && updateStatus(job?.flow_status) - dispatch('jobsLoaded', job) + dispatch('jobsLoaded', { job, force: false }) } errorCount = 0 notAnonynmous = false @@ -242,7 +305,7 @@ let common = { iteration_from: - $localDurationStatuses?.[modId]?.iteration_from ?? + // $localDurationStatuses?.[modId]?.iteration_from ?? Math.max(flowJobIds.flowJobs.length - 20, 0), iteration_total: $localDurationStatuses?.[modId]?.iteration_total ?? flowJobIds?.length } @@ -268,6 +331,20 @@ $: isListJob = flowJobIds != undefined && Array.isArray(flowJobIds?.flowJobs) + $: flowJobIds?.moduleId && onFlowJobFlowStatus() + + function onFlowJobFlowStatus() { + if (globalRefreshes) { + let modId = flowJobIds?.moduleId + if (modId) { + globalRefreshes[modId] = async (loopJob) => { + setIteration(loopJob.index, loopJob.job, false, modId ?? '') + refresh(true, loopJob) + } + } + } + } + onDestroy(() => { destroyed = true timeout && clearTimeout(timeout) @@ -283,7 +360,7 @@ } } - function onJobsLoaded(mod: FlowStatusModule, job: Job): void { + function onJobsLoaded(mod: FlowStatusModule, job: Job, force?: boolean): void { if (mod.id && (mod.flow_jobs ?? []).length == 0) { if (!childFlow) { if ($flowStateStore?.[mod.id]) { @@ -297,33 +374,42 @@ initializeByJob(mod.id) let started_at = job.started_at ? new Date(job.started_at).getTime() : undefined if (job.type == 'QueuedJob') { - setModuleState(mod.id, { - type: 'InProgress', - job_id: job.id, - logs: job.logs, - args: job.args, - started_at, - parent_module: mod['parent_module'] - }) + setModuleState( + mod.id, + { + type: 'InProgress', + job_id: job.id, + logs: job.logs, + args: job.args, + started_at, + parent_module: mod['parent_module'] + }, + force + ) setDurationStatusByJob(mod.id, job.id, { created_at: job.created_at ? new Date(job.created_at).getTime() : undefined, started_at }) } else { - setModuleState(mod.id, { - args: job.args, - type: job['success'] ? 'Success' : 'Failure', - logs: job.logs, - result: job['result'], - job_id: job.id, - parent_module: mod['parent_module'], - duration_ms: job['duration_ms'], - started_at: started_at, - iteration: mod.iterator?.itered?.length, - iteration_total: mod.iterator?.itered?.length, - retries: mod?.failed_retries?.length - // retries: $flowStateStore?.raw_flow - }) + setModuleState( + mod.id, + { + args: job.args, + type: job['success'] ? 'Success' : 'Failure', + logs: job.logs, + result: job['result'], + job_id: job.id, + parent_module: mod['parent_module'], + duration_ms: job['duration_ms'], + started_at: started_at, + flow_jobs: mod.flow_jobs, + flow_jobs_success: mod.flow_jobs_success, + iteration_total: mod.iterator?.itered?.length, + retries: mod?.failed_retries?.length + // retries: $flowStateStore?.raw_flow + }, + force + ) setDurationStatusByJob(mod.id, job.id, { created_at: job.created_at ? new Date(job.created_at).getTime() : undefined, started_at, @@ -333,13 +419,47 @@ } } - function innerJobLoaded( - jobLoaded: (QueuedJob & { type: 'QueuedJob' }) | (CompletedJob & { type: 'CompletedJob' }), - j: number - ) { - let modId = flowJobIds?.moduleId - + function setIteration(j: number, id: string, clicked: boolean, modId: string) { if (modId) { + if (!$localModuleStates?.[modId]) { + $localModuleStates[modId] = { + type: 'InProgress', + args: undefined + } + } + let state = $localModuleStates?.[modId] + + if (state) { + if (state.selectedForloop == id && clicked) { + setModuleState( + modId, + { + selectedForloop: undefined, + selectedForloopIndex: -1 + }, + false, + true + ) + } else { + setModuleState( + modId, + { + selectedForloop: id, + selectedForloopIndex: j + }, + false, + true + ) + clicked && refresh(true, undefined) + } + } + } + } + function innerJobLoaded(jobLoaded: Job, j: number, clicked: boolean, force: boolean) { + let modId = flowJobIds?.moduleId + if (modId) { + setIteration(j, jobLoaded.id, clicked, modId) + if ($flowStateStore && $flowStateStore?.[modId] == undefined) { $flowStateStore[modId] = { ...($flowStateStore[modId] ?? {}), @@ -358,10 +478,9 @@ } if (jobLoaded.type == 'QueuedJob') { jobResults[j] = 'Job in progress ...' - } else { + } else if (jobLoaded.type == 'CompletedJob') { $flowStateStore[modId].previewResult[j] = jobLoaded.result jobResults[j] = jobLoaded.result - jobFailures[j] = jobLoaded.success === false } } @@ -373,34 +492,47 @@ initializeByJob(modId) if (jobLoaded.type == 'QueuedJob') { - setModuleState(modId, { - type: 'InProgress', - started_at, - logs: jobLoaded.logs, - job_id, - args: jobLoaded.args, - iteration: flowJobIds?.flowJobs.length, - iteration_total: flowJobIds?.length, - duration_ms: undefined - }) - + if ($localModuleStates[modId]?.selectedForloopIndex == j) { + setModuleState( + modId, + { + started_at, + logs: jobLoaded.logs, + job_id, + args: jobLoaded.args, + flow_jobs: flowJobIds?.flowJobs, + flow_jobs_success: flowJobIds?.flowJobsSuccess, + iteration_total: flowJobIds?.length, + duration_ms: undefined + }, + force, + true + ) + } setDurationStatusByJob(modId, job_id, { created_at, started_at }) - } else { - setModuleState(modId, { - started_at, - args: jobLoaded.args, - type: jobLoaded.success ? 'Success' : 'Failure', - logs: 'All jobs completed', - result: jobResults, - job_id, - iteration: flowJobIds?.flowJobs.length, - iteration_total: flowJobIds?.length, - duration_ms: undefined, - isListJob: true - }) + } else if (jobLoaded.type == 'CompletedJob') { + if ($localModuleStates[modId]?.selectedForloopIndex == j) { + setModuleState( + modId, + { + started_at, + args: jobLoaded.args, + result: jobLoaded.result, + flow_jobs_results: jobResults, + job_id, + flow_jobs: flowJobIds?.flowJobs, + flow_jobs_success: flowJobIds?.flowJobsSuccess, + iteration_total: flowJobIds?.length, + duration_ms: undefined, + isListJob: true + }, + force, + true + ) + } setDurationStatusByJob(modId, job_id, { created_at, started_at, @@ -424,20 +556,6 @@ let rightColumnSelect: 'timeline' | 'node_status' | 'node_definition' | 'user_states' = 'timeline' - let slicedListJobIds: string[] = [] - - $: flowJobIds && !deepEqual(flowJobIds, lastFlowJobIds) && updateSlicedListJobIds() - let lastFlowJobIds: any = undefined - function updateSlicedListJobIds() { - lastFlowJobIds = flowJobIds - slicedListJobIds = - (flowJobIds?.flowJobs.length ?? 0) > 20 - ? flowJobIds?.flowJobs?.slice( - $localDurationStatuses[flowJobIds?.moduleId ?? '']?.iteration_from ?? 0 - ) ?? [] - : flowJobIds?.flowJobs ?? [] - } - function loadPreviousIters(lenToAdd: number) { let r = $localDurationStatuses[flowJobIds?.moduleId ?? ''] if (r.iteration_from) { @@ -445,11 +563,16 @@ $localDurationStatuses = $localDurationStatuses globalDurationStatuses.forEach((x) => x.update((x) => x)) } - jobResults = [...new Array(lenToAdd), ...jobResults] - updateSlicedListJobIds() + jobResults = [ + ...[...new Array(lenToAdd).keys()].map((x) => 'not computed or loaded yet'), + ...jobResults + ] + // updateSlicedListJobIds() } let stepDetail: FlowModule | string | undefined = undefined + + let storedListJobs: Record = {} {#if notAnonynmous} @@ -464,13 +587,11 @@
{/if} --> {#if isListJob} - {@const lenToAdd = Math.min( - 20, - $localDurationStatuses[flowJobIds?.moduleId ?? '']?.iteration_from ?? 0 - )} + {@const sliceFrom = $localDurationStatuses[flowJobIds?.moduleId ?? '']?.iteration_from ?? 0} + {@const lenToAdd = Math.min(20, sliceFrom)} {#if (flowJobIds?.flowJobs.length ?? 0) > 20 && lenToAdd > 0} - {@const allToAdd = (flowJobIds?.length ?? 0) - (slicedListJobIds.length ?? 0)} + {@const allToAdd = (flowJobIds?.length ?? 0) - sliceFrom}

For performance reasons, only the last 20 items are shown by default - {#if allToAdd > 0} + {#if allToAdd > 0 && allToAdd > lenToAdd} + {sliceFrom} - {#if allToAdd > 0} - + endIcon={{ + icon: ChevronDown, + classes: forloop_selected == loopJobId ? '!rotate-180' : '' + }} + > + + #{j + 1}: {loopJobId} + + {/if} -

- {/if} - {#each slicedListJobIds as loopJobId, j (loopJobId)} - {#if render} - - {/if} - - -
- 20} - {workspaceId} - jobId={loopJobId} - on:jobsLoaded={(e) => innerJobLoaded(e.detail, j)} - /> -
- {/each} + {#if j >= sliceFrom || forloop_selected == loopJobId} + +
+ 20} + {workspaceId} + jobId={loopJobId} + on:jobsLoaded={(e) => { + let { job, force } = e.detail + storedListJobs[j] = job + innerJobLoaded(job, j, false, force) + }} + /> +
+ {/if} + {/each} +
{:else if innerModules.length > 0}

    @@ -698,6 +811,8 @@
    { - onJobsLoaded(mod, e.detail) + let { force, job } = e.detail + onJobsLoaded(mod, job, force) }} /> + {:else if mod.flow_jobs?.length == 0 && mod.job == '00000000-0000-0000-0000-000000000000'} +
    no subflow (empty loop?)
    {:else} onJobsLoaded(mod, e.detail)} + on:jobsLoaded={(e) => { + let { job, force } = e.detail + onJobsLoaded(mod, job, force) + }} /> {/if} {:else} @@ -799,6 +925,14 @@ selectedNode = e.detail.id } }} + on:selectedIteration={(e) => { + let detail = e.detail + setModuleState(detail.moduleId, { + selectedForloop: detail.id, + selectedForloopIndex: detail.index + }) + globalRefreshes[detail.moduleId]?.({ job: detail.id, index: detail.index }) + }} modules={job.raw_flow?.modules ?? []} failureModule={job.raw_flow?.failure_module} /> @@ -853,6 +987,20 @@

    No arguments

    {/if} {:else if node} + {#if node.flow_jobs_results} + Result of step as collection of all subflows +
    +
    + +
    +
    + Selected subflow + {/if}
    {#if node.duration_ms} diff --git a/frontend/src/lib/components/flows/map/FlowJobsMenu.svelte b/frontend/src/lib/components/flows/map/FlowJobsMenu.svelte new file mode 100644 index 0000000000..216067f464 --- /dev/null +++ b/frontend/src/lib/components/flows/map/FlowJobsMenu.svelte @@ -0,0 +1,75 @@ + + + + +
    +
    + + {#each flowJobs ?? [] as id, idx (id)} + {#if filter == undefined || (idx + 1).toString().includes(filter.toString())} + + {/if} + {/each} +
    +
    +
    diff --git a/frontend/src/lib/components/flows/map/FlowModuleSchemaItem.svelte b/frontend/src/lib/components/flows/map/FlowModuleSchemaItem.svelte index 17de02070d..3c5c99d486 100644 --- a/frontend/src/lib/components/flows/map/FlowModuleSchemaItem.svelte +++ b/frontend/src/lib/components/flows/map/FlowModuleSchemaItem.svelte @@ -51,7 +51,7 @@
    {#if $$slots.icon} diff --git a/frontend/src/lib/components/flows/map/MapItem.svelte b/frontend/src/lib/components/flows/map/MapItem.svelte index 7032c8f4d8..fe8121be85 100644 --- a/frontend/src/lib/components/flows/map/MapItem.svelte +++ b/frontend/src/lib/components/flows/map/MapItem.svelte @@ -19,6 +19,7 @@ import { prettyLanguage } from '$lib/common' import { msToSec } from '$lib/utils' import BarsStaggered from '$lib/components/icons/BarsStaggered.svelte' + import FlowJobsMenu from './FlowJobsMenu.svelte' export let mod: FlowModule export let trigger: boolean @@ -33,6 +34,9 @@ export let disableAi: boolean = false export let wrapperId: string | undefined = undefined export let retries: number | undefined = undefined + export let flowJobs: + | { flowJobs: string[]; selected: number; flowJobsSuccess: (boolean | undefined)[] } + | undefined $: idx = modules.findIndex((m) => m.id === mod.id) @@ -48,6 +52,7 @@ select: string newBranch: { module: FlowModule } move: { module: FlowModule } | undefined + selectedIteration: { index: number; id: string } }>() $: itemProps = { @@ -123,6 +128,19 @@ {annotation}
    {/if} + {#if flowJobs && !insertable} +
    + { + dispatch('selectedIteration', e.detail) + }} + flowJobsSuccess={flowJobs.flowJobsSuccess} + flowJobs={flowJobs.flowJobs} + selected={flowJobs.selected} + index={idx} + /> +
    + {/if}
    {#if mod.value.type === 'forloopflow' || mod.value.type === 'whileloopflow'} diff --git a/frontend/src/lib/components/flows/map/VirtualItem.svelte b/frontend/src/lib/components/flows/map/VirtualItem.svelte index b29b45fc39..7003654d60 100644 --- a/frontend/src/lib/components/flows/map/VirtualItem.svelte +++ b/frontend/src/lib/components/flows/map/VirtualItem.svelte @@ -1,6 +1,6 @@