diff --git a/backend/windmill-worker/src/worker_flow.rs b/backend/windmill-worker/src/worker_flow.rs index 6c5abbfd7b..a4b65a589c 100644 --- a/backend/windmill-worker/src/worker_flow.rs +++ b/backend/windmill-worker/src/worker_flow.rs @@ -162,6 +162,11 @@ pub struct SkipIfStopped { pub args: Option>>>, } +#[derive(Deserialize)] +struct RecoveryObject { + recover: Option, +} + #[derive(sqlx::FromRow, Deserialize)] pub struct RowFlowStatus { pub flow_status: sqlx::types::Json>, @@ -758,7 +763,7 @@ pub async fn update_flow_status_after_job_completion_internal< let module = get_module(&flow_job, module_index); // tracing::error!( - // "UPDATE FLOW STATUS 3: {module:#?} {unrecoverable} {} {is_last_step} {success} {skip_error_handler}", flow_job.canceled + // "UPDATE FLOW STATUS 3: {module:#?} {unrecoverable} {} {is_last_step} {success} {skip_error_handler} is_failure_step {is_failure_step}", flow_job.canceled // ); let should_continue_flow = match success { @@ -874,7 +879,13 @@ pub async fn update_flow_status_after_job_completion_internal< save_in_cache(db, client, &flow_job, cached_res_path, &nresult).await; } - let success = success && !is_failure_step && !skip_error_handler; + fn result_has_recover_true(nresult: Arc>) -> bool { + let recover = serde_json::from_str::(nresult.get()); + return recover.map(|r| r.recover.unwrap_or(false)).unwrap_or(false); + } + let success = success + && (!is_failure_step || result_has_recover_true(nresult.clone())) + && !skip_error_handler; if success { add_completed_job( db, diff --git a/frontend/src/lib/components/flows/content/FlowInputs.svelte b/frontend/src/lib/components/flows/content/FlowInputs.svelte index 02a903231f..db5f9eac17 100644 --- a/frontend/src/lib/components/flows/content/FlowInputs.svelte +++ b/frontend/src/lib/components/flows/content/FlowInputs.svelte @@ -221,7 +221,6 @@ return } } - console.log(lang, kind) dispatch('new', { language: lang == 'docker' ? 'bash' : lang, kind, diff --git a/frontend/src/lib/components/flows/flowInfers.ts b/frontend/src/lib/components/flows/flowInfers.ts index 7c4e34eb7b..93bd366196 100644 --- a/frontend/src/lib/components/flows/flowInfers.ts +++ b/frontend/src/lib/components/flows/flowInfers.ts @@ -34,7 +34,7 @@ export async function loadSchemaFromModule(module: FlowModule): Promise<{ input_transforms = keys.reduce((accu, key) => { let nv = input_transforms[key] ?? - (module.id == 'failure' && ['message', 'name'].includes(key) + (module.id == 'failure' && ['message', 'name', 'step_id'].includes(key) ? { type: 'javascript', expr: `error.${key}` } : { type: 'static', diff --git a/frontend/src/lib/init_scripts/python_failure_module.ts b/frontend/src/lib/init_scripts/python_failure_module.ts index f6acb652f5..db8e1a7705 100644 --- a/frontend/src/lib/init_scripts/python_failure_module.ts +++ b/frontend/src/lib/init_scripts/python_failure_module.ts @@ -1,7 +1,8 @@ export default `import os -def main(message: str, name: str): - flow_id = os.environ.get("WM_FLOW_JOB_ID") +def main(message: str, name: str, step_id: str): + flow_id = os.environ.get("WM_ROOT_FLOW_JOB_ID") print("message", message) print("name", name) - return message, flow_id` + print("step_id", step_id) + return { "message": message, "flow_id": flow_id, "step_id": step_id, "recover": False }` diff --git a/frontend/src/lib/script_helpers.ts b/frontend/src/lib/script_helpers.ts index d450a8b451..326cea6d83 100644 --- a/frontend/src/lib/script_helpers.ts +++ b/frontend/src/lib/script_helpers.ts @@ -161,20 +161,22 @@ export async function main(x: string) { ` export const DENO_FAILURE_MODULE_CODE = ` -export async function main(message: string, name: string) { - const flow_id = Deno.env.get("WM_FLOW_JOB_ID") +export async function main(message: string, name: string, step_id: string) { + const flow_id = Deno.env.get("WM_ROOT_FLOW_JOB_ID") console.log("message", message) console.log("name",name) - return { message, flow_id } + console.log("step_id", step_id) + return { message, flow_id, step_id, recover: false } } ` export const BUN_FAILURE_MODULE_CODE = ` -export async function main(message: string, name: string) { - const flow_id = process.env.WM_FLOW_JOB_ID +export async function main(message: string, name: string, step_id: string) { + const flow_id = process.env.WM_ROOT_FLOW_JOB_ID console.log("message", message) console.log("name",name) - return { message, flow_id } + console.log("step_id", step_id) + return { message, flow_id, step_id, recover: false } } ` @@ -519,10 +521,10 @@ export function initialCode( return PYTHON_INIT_CODE_TRIGGER } else if (kind === 'approval') { return PYTHON_INIT_CODE_APPROVAL - } else if (subkind === 'flow') { - return PYTHON_INIT_CODE_CLEAR } else if (kind === 'failure') { return PYTHON_FAILURE_MODULE_CODE + } else if (subkind === 'flow') { + return PYTHON_INIT_CODE_CLEAR } else { return PYTHON_INIT_CODE }