feat: improve flow status viewer (show branch chosen + all iterations in for loop) (#4074)

* all

* all

* all

* all

* all
This commit is contained in:
Ruben Fiszel
2024-07-13 18:53:24 +02:00
committed by GitHub
parent 83e06fa722
commit 7f2cbb40f6
16 changed files with 749 additions and 256 deletions
@@ -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"
}
+1 -1
View File
@@ -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
@@ -119,6 +119,7 @@ struct UntaggedFlowStatusModule {
job: Option<Uuid>,
iterator: Option<Iterator>,
flow_jobs: Option<Vec<Uuid>>,
flow_jobs_success: Option<Vec<Option<bool>>>,
branch_chosen: Option<BranchChosen>,
branchall: Option<BranchAllStatus>,
parallel: Option<bool>,
@@ -150,6 +151,8 @@ pub enum FlowStatusModule {
#[serde(skip_serializing_if = "Option::is_none")]
flow_jobs: Option<Vec<Uuid>>,
#[serde(skip_serializing_if = "Option::is_none")]
flow_jobs_success: Option<Vec<Option<bool>>>,
#[serde(skip_serializing_if = "Option::is_none")]
branch_chosen: Option<BranchChosen>,
#[serde(skip_serializing_if = "Option::is_none")]
branchall: Option<BranchAllStatus>,
@@ -164,6 +167,8 @@ pub enum FlowStatusModule {
#[serde(skip_serializing_if = "Option::is_none")]
flow_jobs: Option<Vec<Uuid>>,
#[serde(skip_serializing_if = "Option::is_none")]
flow_jobs_success: Option<Vec<Option<bool>>>,
#[serde(skip_serializing_if = "Option::is_none")]
branch_chosen: Option<BranchChosen>,
#[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<Vec<Uuid>>,
#[serde(skip_serializing_if = "Option::is_none")]
flow_jobs_success: Option<Vec<Option<bool>>>,
#[serde(skip_serializing_if = "Option::is_none")]
branch_chosen: Option<BranchChosen>,
#[serde(skip_serializing_if = "Vec::is_empty")]
failed_retries: Vec<Uuid>,
@@ -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<BranchChosen> {
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<Vec<Option<bool>>> {
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<JobResult> {
self.flow_jobs()
.map(JobResult::ListJob)
+11 -9
View File
@@ -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",
))),
+167 -24
View File
@@ -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<Vec<Option<bool>>>,
flow_jobs: &Vec<Uuid>,
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<R: rsmq_async::RsmqConnection + Send + Sync + Clone>
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<R: rsmq_async::RsmqConnection + Send + Sync + Clone>
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<R: rsmq_async::RsmqConnection + Send + Sync + Clone>
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<R: rsmq_async::RsmqConnection + Send + Sync + Clone>
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<R: rsmq_async::RsmqConnection + Send + Sync + Clone>
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<Box<RawValue>>,
flow_jobs: Vec<Uuid>,
flow_jobs_success: Option<Vec<Option<bool>>>,
new_args: Iter,
while_loop: bool,
}
@@ -2508,6 +2634,7 @@ enum ForLoopStatus {
struct NextBranch {
status: BranchAllStatus,
flow_jobs: Vec<Uuid>,
flow_jobs_success: Option<Vec<Option<bool>>>,
}
#[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,
})
@@ -46,11 +46,12 @@
<FlowStatusViewerInner
on:jobsLoaded={({ detail }) => {
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={[]}
@@ -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<Record<string, GraphModuleState>>[]
export let globalDurationStatuses: Writable<Record<string, DurationStatus>>[]
export let globalRefreshes: Record<
string,
(loopJob: { index: number; job: string }) => Promise<void>
> = {}
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<Record<string, GraphModuleState>> = writable({})
let localDurationStatuses: Writable<Record<string, DurationStatus>> = 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<GraphModuleState>,
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<string, (boolean) => Promise<void>> = {}
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<number, Job> = {}
</script>
{#if notAnonynmous}
@@ -464,13 +587,11 @@
<div class="h-8" />
{/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}
<p class="text-tertiary italic text-xs">
For performance reasons, only the last 20 items are shown by default <button
class="text-primary underline ml-4"
@@ -479,7 +600,8 @@
}}
>Load {lenToAdd} prior
</button>
{#if allToAdd > 0}
{#if allToAdd > 0 && allToAdd > lenToAdd}
{sliceFrom}
<button
class="text-primary underline ml-4"
on:click={() => {
@@ -570,79 +692,70 @@
{/if}
<div class="{selected != 'sequence' ? 'hidden' : ''} max-w-7xl mx-auto">
{#if isListJob}
{@const lenToAdd = Math.min(
20,
$localDurationStatuses[flowJobIds?.moduleId ?? '']?.iteration_from ?? 0
)}
{@const sliceFrom = $localDurationStatuses[flowJobIds?.moduleId ?? '']?.iteration_from ?? 0}
{@const forloop_selected =
$localModuleStates?.[flowJobIds?.moduleId ?? '']?.selectedForloop}
<h3 class="text-md leading-6 font-bold text-tertiary border-b mb-4">
Subflows: ({flowJobIds?.flowJobs.length} items)
Subflows ({flowJobIds?.flowJobs.length})
</h3>
{#if (flowJobIds?.flowJobs.length ?? 0) > 20 && lenToAdd > 0}
{@const allToAdd = (flowJobIds?.length ?? 0) - (slicedListJobIds.length ?? 0)}
<p class="text-tertiary italic text-xs">
For performance reasons, only the last 20 items are shown by default <button
class="text-primary underline ml-4"
on:click={() => {
loadPreviousIters(lenToAdd)
}}
>Load {lenToAdd} prior
</button>
{#if allToAdd > 0}
<button
class="text-primary underline ml-4"
on:click={() => {
loadPreviousIters(allToAdd)
<div class="overflow-auto max-h-1/2">
{forloop_selected}
{#each flowJobIds?.flowJobs ?? [] as loopJobId, j (loopJobId)}
{#if render}
<Button
variant={forloop_selected === loopJobId ? 'contained' : 'border'}
color={flowJobIds?.flowJobsSuccess?.[j] === false
? 'red'
: forloop_selected === loopJobId
? 'dark'
: 'light'}
btnClasses="w-full flex justify-start"
on:click={async () => {
let storedJob = storedListJobs[j]
if (!storedJob) {
storedJob = await JobService.getJob({
workspace: workspaceId ?? $workspaceStore ?? '',
id: loopJobId,
noLogs: true
})
storedListJobs[j] = storedJob
}
innerJobLoaded(storedJob, j, true, false)
}}
>Load {allToAdd} prior
</button>
endIcon={{
icon: ChevronDown,
classes: forloop_selected == loopJobId ? '!rotate-180' : ''
}}
>
<span class="truncate font-mono">
#{j + 1}: {loopJobId}
</span>
</Button>
{/if}
</p>
{/if}
{#each slicedListJobIds as loopJobId, j (loopJobId)}
{#if render}
<Button
variant={forloop_selected === loopJobId ? 'contained' : 'border'}
color={jobFailures[j] === true
? 'red'
: forloop_selected === loopJobId
? 'dark'
: 'light'}
btnClasses="w-full flex justify-start"
on:click={() => {
if (forloop_selected == loopJobId) {
forloop_selected = ''
} else {
forloop_selected = loopJobId
}
}}
endIcon={{
icon: ChevronDown,
classes: forloop_selected == loopJobId ? '!rotate-180' : ''
}}
>
<span class="truncate font-mono">
#{($localDurationStatuses[flowJobIds?.moduleId ?? '']?.iteration_from ?? 0) +
j +
1}: {loopJobId}
</span>
</Button>
{/if}
<!-- <LogId id={loopJobId} /> -->
<div class="border p-6" class:hidden={forloop_selected != loopJobId}>
<svelte:self
{childFlow}
globalModuleStates={[localModuleStates, ...globalModuleStates]}
globalDurationStatuses={[localDurationStatuses, ...globalDurationStatuses]}
render={forloop_selected == loopJobId && selected == 'sequence' && render}
reducedPolling={flowJobIds?.flowJobs.length && flowJobIds?.flowJobs.length > 20}
{workspaceId}
jobId={loopJobId}
on:jobsLoaded={(e) => innerJobLoaded(e.detail, j)}
/>
</div>
{/each}
{#if j >= sliceFrom || forloop_selected == loopJobId}
<!-- <LogId id={loopJobId} /> -->
<div class="border p-6" class:hidden={forloop_selected != loopJobId}>
<svelte:self
bind:refresh={recursiveRefresh[loopJobId]}
{globalRefreshes}
{childFlow}
job={storedListJobs[j]}
globalModuleStates={[localModuleStates, ...globalModuleStates]}
globalDurationStatuses={[localDurationStatuses, ...globalDurationStatuses]}
render={forloop_selected == loopJobId && selected == 'sequence' && render}
reducedPolling={flowJobIds?.flowJobs.length && flowJobIds?.flowJobs.length > 20}
{workspaceId}
jobId={loopJobId}
on:jobsLoaded={(e) => {
let { job, force } = e.detail
storedListJobs[j] = job
innerJobLoaded(job, j, false, force)
}}
/>
</div>
{/if}
{/each}
</div>
{:else if innerModules.length > 0}
<ul class="w-full">
<h3 class="text-md leading-6 font-bold text-primary border-b mb-4 py-2">
@@ -698,6 +811,8 @@
<!-- <LogId id={loopJobId} /> -->
<div class="border p-6" class:hidden={retry_selected != failedRetry}>
<svelte:self
{globalRefreshes}
bind:refresh={recursiveRefresh[failedRetry]}
{childFlow}
globalModuleStates={[localModuleStates, ...globalModuleStates]}
globalDurationStatuses={[localDurationStatuses, ...globalDurationStatuses]}
@@ -712,6 +827,8 @@
{#if ['InProgress', 'Success', 'Failure'].includes(mod.type)}
{#if job.raw_flow?.modules[i]?.value.type == 'flow'}
<svelte:self
{globalRefreshes}
bind:refresh={recursiveRefresh[mod.job ?? '']}
globalModuleStates={[]}
globalDurationStatuses={[]}
render={selected == 'sequence' && render}
@@ -719,11 +836,16 @@
jobId={mod.job}
childFlow
on:jobsLoaded={(e) => {
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'}
<div class="text-secondary">no subflow (empty loop?)</div>
{:else}
<svelte:self
{globalRefreshes}
bind:refresh={recursiveRefresh[mod.job ?? '']}
{childFlow}
globalModuleStates={[localModuleStates, ...globalModuleStates]}
globalDurationStatuses={[localDurationStatuses, ...globalDurationStatuses]}
@@ -735,10 +857,14 @@
? {
moduleId: mod.id,
flowJobs: mod.flow_jobs,
flowJobsSuccess: mod.flow_jobs_success,
length: mod.iterator?.itered?.length ?? mod.flow_jobs.length
}
: undefined}
on:jobsLoaded={(e) => 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 @@
<p class="p-2 text-secondary">No arguments</p>
{/if}
{:else if node}
{#if node.flow_jobs_results}
<span class="pl-1 text-tertiary"
>Result of step as collection of all subflows</span
>
<div class="p-2">
<div class="overflow-auto max-h-[200px]">
<DisplayResult
workspaceId={job?.workspace_id}
result={node.flow_jobs_results}
/>
</div>
</div>
<span class="pl-1 text-tertiary text-lg pt-4">Selected subflow</span>
{/if}
<div class="px-2 flex gap-2 min-w-0 w-full">
<ModuleStatus type={node.type} scheduled_for={node.scheduled_for} />
{#if node.duration_ms}
@@ -0,0 +1,75 @@
<script lang="ts">
import { Menu } from '$lib/components/common'
import { createEventDispatcher } from 'svelte'
import { ListFilter } from 'lucide-svelte'
const dispatch = createEventDispatcher()
export let open: boolean | undefined = undefined
export let index: number
export let flowJobs: string[] | undefined
export let flowJobsSuccess: (boolean | undefined)[] | undefined
export let selected: number
let filter: number | undefined = undefined
function onKeydown(event: KeyboardEvent) {
if (
event.key === 'Enter' &&
filter != undefined &&
flowJobs &&
filter < flowJobs.length &&
filter > 0
) {
event.preventDefault()
dispatch('selectedIteration', { index: filter - 1, id: flowJobs[filter - 1] })
}
}
</script>
<Menu
transitionDuration={0}
pointerDown
bind:show={open}
noMinW
placement="bottom-center"
let:close
>
<button
title="Pick an interation"
slot="trigger"
id={`flow-editor-iteration picker-${index}`}
type="button"
class=" text-xs bg-surface border-[1px] border-gray-300 dark:border-gray-500 focus:outline-none
hover:bg-surface-hover focus:ring-4 focus:ring-surface-selected font-medium rounded-sm w-[40px] gap-1 h-[20px]
flex items-center justify-center {flowJobsSuccess?.[selected] == false
? 'text-red-400'
: 'text-secondary'}"
>
#{selected == -1 ? '?' : selected + 1}
<ListFilter size={15} />
</button>
<div id="flow-editor-insert-module">
<div class="font-mono divide-y text-xs w-full text-secondary max-h-[200px] overflow-auto">
<input autofocus type="number" bind:value={filter} on:keydown={onKeydown} />
{#each flowJobs ?? [] as id, idx (id)}
{#if filter == undefined || (idx + 1).toString().includes(filter.toString())}
<button
class="w-full text-left py-1 pl-2 min-w-20 hover:bg-surface-hover whitespace-nowrap flex flex-row gap-2 items-center {flowJobsSuccess?.[
idx
] == false
? 'text-red-400'
: ''}"
on:pointerdown={() => {
close()
dispatch('selectedIteration', { index: idx, id })
}}
role="menuitem"
tabindex="-1"
>
#{idx + 1}
</button>
{/if}
{/each}
</div>
</div>
</Menu>
@@ -51,7 +51,7 @@
<div
class={classNames(
'w-full module flex rounded-sm cursor-pointer',
selected ? 'outline outline-offset-1 outline-2 outline-gray-600' : '',
selected ? 'outline outline-offset-1 outline-2 outline-gray-600 dark:outline-gray-400' : '',
'flex relative',
$copilotCurrentStepStore === id ? 'z-[901]' : ''
)}
@@ -142,7 +142,7 @@
<div
class="flex gap-1 justify-between items-center w-full overflow-hidden rounded-sm
border border-gray-400 p-2 text-2xs module text-primary"
border border-gray-400 dark:border-gray-600 p-2 text-2xs module text-primary"
>
{#if $$slots.icon}
<slot name="icon" />
@@ -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}
</div>
{/if}
{#if flowJobs && !insertable}
<div class="absolute z-10 right-8 -top-5">
<FlowJobsMenu
on:selectedIteration={(e) => {
dispatch('selectedIteration', e.detail)
}}
flowJobsSuccess={flowJobs.flowJobsSuccess}
flowJobs={flowJobs.flowJobs}
selected={flowJobs.selected}
index={idx}
/>
</div>
{/if}
<div class={moving == mod.id ? 'opacity-50' : ''}>
{#if mod.value.type === 'forloopflow' || mod.value.type === 'whileloopflow'}
@@ -1,6 +1,6 @@
<script lang="ts">
import { Badge } from '$lib/components/common'
import type { FlowModule } from '$lib/gen'
import type { FlowModule, FlowStatusModule } from '$lib/gen'
import { classNames } from '$lib/utils'
import { ClipboardCopy, ExternalLink, Wand2, X } from 'lucide-svelte'
import { createEventDispatcher, getContext } from 'svelte'
@@ -9,6 +9,7 @@
import { copilotInfo } from '$lib/stores'
import Menu from '$lib/components/common/menu/Menu.svelte'
import InsertTriggerButton from './InsertTriggerButton.svelte'
import { getStateColor } from '$lib/components/graph'
export let label: string
export let modules: FlowModule[] | undefined
@@ -23,6 +24,7 @@
export let center = true
export let disableAi: boolean = false
export let wrapperNode: FlowModule | undefined = undefined
export let borderStatus: FlowStatusModule['type'] | undefined = undefined
const dispatch = createEventDispatcher<{
insert: {
@@ -52,7 +54,7 @@
}
}}
type="button"
class="text-primary bg-surface border mx-[1px] border-gray-300 dark:border-gray-500 focus:outline-none hover:bg-surface-hover focus:ring-4 focus:ring-gray-200 font-medium rounded-full text-sm w-[25px] h-[25px] flex items-center justify-center"
class="text-primary bg-surface border mx-[1px] 'border-gray-300 dark:border-gray-500 focus:outline-none hover:bg-surface-hover focus:ring-4 focus:ring-gray-200 font-medium rounded-full text-sm w-[25px] h-[25px] flex items-center justify-center"
>
<X class="m-[5px]" size={15} />
</button>
@@ -80,9 +82,10 @@
id={`flow-editor-virtual-${label}`}
>
<div
style={borderStatus ? `border-color: ${getStateColor(borderStatus)}` : ''}
class="flex gap-1 justify-between {center
? 'items-center'
: 'items-baseline'} w-full overflow-hidden rounded-sm border p-2 text-2xs module text-primary border-gray-400"
: 'items-baseline'} w-full overflow-hidden rounded-sm border p-2 text-2xs module text-primary border-gray-400 dark:border-gray-600"
>
{#if $$slots.icon}
<slot name="icon" />
@@ -1,6 +1,6 @@
<script lang="ts">
import { sugiyama, dagStratify, decrossOpt, coordCenter } from 'd3-dag'
import { type FlowModule } from '../../gen'
import { type FlowModule, type FlowStatusModule } from '../../gen'
import {
NODE,
createIdGenerator,
@@ -123,7 +123,9 @@
true,
undefined,
undefined,
'Input'
'Input',
undefined,
success == undefined ? undefined : success ? 'Success' : 'Failure'
)
)
@@ -150,6 +152,8 @@
true,
undefined,
undefined,
undefined,
success == undefined ? undefined : success ? 'Success' : 'Failure',
undefined
)
)
@@ -269,7 +273,8 @@
insertableEnd,
false,
modules,
wrapper
wrapper,
undefined
)
}
@@ -291,19 +296,6 @@
return []
}
function getResultColor(): string {
const isDark = document.documentElement.classList.contains('dark')
switch (success) {
case true:
return getStateColor('Success')
case false:
return getStateColor('Failure')
default:
return isDark ? '#2e3440' : '#fff'
}
}
function flowModuleToNode(
parentIds: string[],
mod: FlowModule,
@@ -313,8 +305,15 @@
insertableEnd: boolean,
branchable: boolean,
modules: FlowModule[],
wrapper: FlowModule | undefined = undefined
wrapper: FlowModule | undefined,
flowJobs:
| { flowJobs: string[]; selected: number; flowJobsSuccess?: (boolean | undefined)[] }
| undefined
): Node {
let type = flowModuleStates?.[mod.id]?.type
if (!type && flowJobs) {
type = 'InProgress'
}
return {
type: 'node',
id: mod.id,
@@ -330,12 +329,13 @@
branchable,
retries: flowModuleStates?.[mod.id]?.retries,
duration_ms: flowModuleStates?.[mod.id]?.duration_ms,
bgColor: getStateColor(flowModuleStates?.[mod.id]?.type),
annotation,
bgColor: getStateColor(type),
annotation: annotation,
modules,
moving,
disableAi,
wrapperId: wrapper?.id
wrapperId: wrapper?.id,
flowJobs
},
cb: (e: string, detail: any) => {
if (e == 'delete') {
@@ -353,6 +353,8 @@
dispatch('newBranch', detail)
} else if (e == 'move') {
dispatch('move', { module: mod, modules })
} else if (e == 'selectedIteration') {
dispatch('selectedIteration', { ...detail, moduleId: mod.id })
}
}
}
@@ -373,6 +375,7 @@
parent: NestedNodes | string | undefined,
loopDepth: number
): Loop {
let state = flowModuleStates?.[module.id]
const loop: Loop = {
type: 'loop',
items: [
@@ -380,21 +383,30 @@
getParentIds(parent),
module,
undefined,
flowModuleStates?.[module.id]?.iteration
? 'Iteration ' +
flowModuleStates?.[module.id]?.iteration +
'/' +
(flowModuleStates?.[module.id]?.iteration_total ?? '?')
state?.flow_jobs
? 'iterations ' + state?.flow_jobs?.length + '/' + (state?.iteration_total ?? '?')
: '',
loopDepth,
false,
false,
modules
modules,
undefined,
state?.flow_jobs
? {
flowJobs: state?.flow_jobs,
selected: state?.selectedForloopIndex ?? -1,
flowJobsSuccess: state?.flow_jobs_success
}
: undefined
)
]
}
const innerModules = module.value.modules
let borderStatus: FlowStatusModule['type'] | undefined = undefined
let success = state?.flow_jobs_success?.[state?.selectedForloopIndex ?? 0]
if (success != undefined) {
borderStatus = success ? 'Success' : 'Failure'
}
loop.items.push(
createVirtualNode(
getParentIds(loop.items),
@@ -408,6 +420,8 @@
undefined,
undefined,
undefined,
undefined,
borderStatus,
true,
module
)
@@ -436,7 +450,9 @@
true,
undefined,
module.id,
undefined
undefined,
undefined,
flowModuleStates?.[module.id]?.type
)
)
return loop
@@ -460,7 +476,9 @@
loopDepth,
false,
true,
modules
modules,
undefined,
undefined
)
const bitems: NestedNodes[] = []
const branchParent = [node.id]
@@ -477,6 +495,8 @@
false,
undefined,
undefined,
undefined,
undefined,
undefined
)
])
@@ -485,6 +505,24 @@
branches.forEach(({ summary, modules, removable }, i) => {
const items: NestedNodes = []
let borderStatus: FlowStatusModule['type'] | undefined = undefined
if (module.value.type == 'branchall' || module.value.type == 'forloopflow') {
let flow_jobs_success = flowModuleStates?.[module.id]?.flow_jobs_success
if (!flow_jobs_success) {
borderStatus = 'WaitingForPriorSteps'
} else {
let status = flow_jobs_success?.[i]
if (status == undefined) {
borderStatus = 'WaitingForExecutor'
} else {
borderStatus = status ? 'Success' : 'Failure'
}
}
} else if (module.value.type == 'branchone') {
if (flowModuleStates?.[module.id]?.branchChosen == i) {
borderStatus = 'Success'
}
}
items.push(
createVirtualNode(
branchParent,
@@ -498,6 +536,8 @@
removable ? { module, index: i } : undefined,
undefined,
undefined,
undefined,
borderStatus,
false,
wrapper
)
@@ -533,7 +573,9 @@
true,
undefined,
module.id,
undefined
undefined,
undefined,
flowModuleStates?.[module.id]?.type
),
items: bitems
}
@@ -665,11 +707,19 @@
deleteBranch: { module: FlowModule; index: number } | undefined,
mid: string | undefined,
fixed_id: string | undefined,
module_status: FlowStatusModule['type'] | undefined,
borderStatus: FlowStatusModule['type'] | undefined,
center: boolean = true,
wrapperNode: FlowModule | undefined = undefined
): Node {
const id = fixed_id ?? -idGenerator.next().value - 2 + (offset ?? 0)
let bgColor
if (module_status) {
bgColor = getStateColor(module_status)
} else {
bgColor = document.documentElement.classList.contains('dark') ? '#2e3440' : '#dfe6ee'
}
return {
type: 'node',
id: id.toString(),
@@ -686,12 +736,8 @@
label,
insertable,
modules,
bgColor:
label == 'Result'
? getResultColor()
: document.documentElement.classList.contains('dark')
? '#2e3440'
: '#dfe6ee',
bgColor,
borderStatus,
selected: $selectedId == label,
index,
selectable,
+6 -1
View File
@@ -47,11 +47,16 @@ export type GraphModuleState = {
type: FlowStatusModule['type']
args: any
logs?: string
flow_jobs_results?: any
branchChosen?: number
result?: any
scheduled_for?: Date
job_id?: string
parent_module?: string
iteration?: number
selectedForloop?: string
selectedForloopIndex?: number
flow_jobs_success?: (boolean | undefined)[]
flow_jobs?: string[]
iteration_total?: number
retries?: number
duration_ms?: number
@@ -165,6 +165,8 @@
on:nodeInsert={(e) => node?.data?.custom?.cb?.('nodeInsert', e.detail)}
on:addBranch={(e) => node?.data?.custom?.cb?.('addBranch', e.detail)}
on:removeBranch={(e) => node?.data?.custom?.cb?.('removeBranch', e.detail)}
on:selectedIteration={(e) =>
node?.data?.custom?.cb?.('selectedIteration', e.detail)}
{...node.data.custom.props}
/>
</Node>
@@ -400,6 +400,7 @@
})
onDestroy(() => {
sync = false
if (intervalId) {
clearInterval(intervalId)
}
+4
View File
@@ -481,6 +481,10 @@ components:
type: array
items:
type: string
flow_jobs_success:
type: array
items:
type: boolean
branch_chosen:
type: object
properties: