diff --git a/backend/tests/worker.rs b/backend/tests/worker.rs index c5facab96c..f298eadab8 100644 --- a/backend/tests/worker.rs +++ b/backend/tests/worker.rs @@ -1546,13 +1546,13 @@ async fn test_empty_loop(db: Pool) { "iterator": { "type": "static", "value": [] }, "modules": [ { - "input_transform": { - "n": { - "type": "javascript", - "expr": "previous_result.iter.value", - }, - }, "value": { + "input_transform": { + "n": { + "type": "javascript", + "expr": "previous_result.iter.value", + }, + }, "type": "rawscript", "language": "python3", "content": "def main(n): return n", @@ -1562,13 +1562,13 @@ async fn test_empty_loop(db: Pool) { }, }, { - "input_transform": { - "items": { - "type": "javascript", - "expr": "previous_result", - }, - }, "value": { + "input_transform": { + "items": { + "type": "javascript", + "expr": "previous_result", + }, + }, "type": "rawscript", "language": "python3", "content": "def main(items): return sum(items)", @@ -1587,6 +1587,45 @@ async fn test_empty_loop(db: Pool) { assert_eq!(result, serde_json::json!(0)); } +#[sqlx::test(fixtures("base"))] +async fn test_invalid_first_step(db: Pool) { + initialize_tracing().await; + let server = ApiServer::start(db.clone()).await; + let port = server.addr.port(); + + let flow: FlowValue = serde_json::from_value(serde_json::json!({ + "modules": [ + { + "value": { + "type": "forloopflow", + "iterator": { "type": "javascript", "expr": "flow_input" }, + "modules": [ + { + "value": { + "type": "identity", + }, + } + ], + }, + }, + { + "value": { + "type": "identity", + }, + }, + ], + })) + .unwrap(); + + let flow = JobPayload::RawFlow { value: flow, path: None }; + let job = run_job_in_new_worker_until_complete(&db, flow, port).await; + + assert_eq!( + job.result.unwrap(), + serde_json::json!({"error":"Expected an array value, found: {}"}) + ); +} + #[sqlx::test(fixtures("base"))] async fn test_empty_loop_2(db: Pool) { initialize_tracing().await; diff --git a/backend/windmill-common/src/flow_status.rs b/backend/windmill-common/src/flow_status.rs index 8d9d739c16..fe590ca291 100644 --- a/backend/windmill-common/src/flow_status.rs +++ b/backend/windmill-common/src/flow_status.rs @@ -164,7 +164,3 @@ impl FlowStatus { self.modules.get(i) } } - -pub fn init_flow_status(f: &FlowValue) -> FlowStatus { - FlowStatus::new(f) -} diff --git a/backend/windmill-queue/src/jobs.rs b/backend/windmill-queue/src/jobs.rs index 890746b040..9914e86746 100644 --- a/backend/windmill-queue/src/jobs.rs +++ b/backend/windmill-queue/src/jobs.rs @@ -16,7 +16,7 @@ use uuid::Uuid; use windmill_audit::{audit_log, ActionKind}; use windmill_common::{ error::{self, to_anyhow, Error}, - flow_status::{init_flow_status, FlowStatus, MAX_RETRY_ATTEMPTS, MAX_RETRY_INTERVAL}, + flow_status::{FlowStatus, MAX_RETRY_ATTEMPTS, MAX_RETRY_INTERVAL}, flows::FlowValue, scripts::{get_full_hub_script_by_path, HubScript, ScriptHash, ScriptLang}, utils::StripPath, @@ -392,7 +392,7 @@ pub async fn push<'c>( } } - let flow_status = raw_flow.as_ref().map(init_flow_status); + let flow_status = raw_flow.as_ref().map(FlowStatus::new); let uuid = sqlx::query_scalar!( "INSERT INTO queue (workspace_id, id, running, parent_job, created_by, permissioned_as, scheduled_for, diff --git a/backend/windmill-worker/src/jobs.rs b/backend/windmill-worker/src/jobs.rs index 42084331b4..77dea01606 100644 --- a/backend/windmill-worker/src/jobs.rs +++ b/backend/windmill-worker/src/jobs.rs @@ -6,6 +6,7 @@ * LICENSE-AGPL for a copy of the license. */ +use serde_json::{Map, Value}; use sqlx::{Pool, Postgres, Transaction}; use tracing::instrument; use uuid::Uuid; @@ -22,11 +23,8 @@ pub async fn add_completed_job_error( metrics: Option, ) -> Result<(Uuid, serde_json::Map), Error> { metrics.map(|m| m.worker_execution_failed.inc()); - let mut output_map = serde_json::Map::new(); - output_map.insert( - "error".to_string(), - serde_json::Value::String(e.to_string()), - ); + let mut output_map = Map::new(); + error_to_result(&mut output_map, &e); let a = add_completed_job( db, client, @@ -40,6 +38,16 @@ pub async fn add_completed_job_error( Ok((a, output_map)) } +pub fn error_to_result( + output_map: &mut Map, + err: &E, +) { + output_map.insert( + "error".to_string(), + serde_json::Value::String(err.to_string()), + ); +} + #[instrument(level = "trace", skip_all)] pub async fn add_completed_job( db: &Pool, diff --git a/backend/windmill-worker/src/worker.rs b/backend/windmill-worker/src/worker.rs index 5630676529..9730044e95 100644 --- a/backend/windmill-worker/src/worker.rs +++ b/backend/windmill-worker/src/worker.rs @@ -42,7 +42,7 @@ use futures::{ use async_recursion::async_recursion; use crate::{ - jobs::{add_completed_job, add_completed_job_error}, + jobs::{add_completed_job, add_completed_job_error, error_to_result}, worker_flow::{ handle_flow, update_flow_status_after_job_completion, update_flow_status_in_progress, }, @@ -467,24 +467,16 @@ async fn handle_job_error( keep_job_dir: bool, base_internal_url: &str, ) { - let m = add_completed_job_error( - db, - client, - &job, - "Unexpected error during job execution:\n".to_string(), - &err, - metrics.clone(), - ) - .await - .map(|(_, m)| m) - .unwrap_or_else(|_| Map::new()); - if job.is_flow_step || job.job_kind == JobKind::FlowPreview || job.job_kind == JobKind::Flow { let (flow, job_status_to_update) = if let Some(parent_job_id) = job.parent_job { (parent_job_id, job.id) } else { (job.id, Uuid::nil()) }; + println!("job {:#?} failed {} {}", job, flow, job_status_to_update); + + let mut output_map = serde_json::Map::new(); + error_to_result(&mut output_map, &err); let updated_flow = update_flow_status_after_job_completion( db, client, @@ -492,7 +484,7 @@ async fn handle_job_error( &job_status_to_update, &job.workspace_id, false, - serde_json::Value::Object(m), + serde_json::Value::Object(output_map), metrics.clone(), unrecoverable, same_worker_tx, @@ -503,6 +495,8 @@ async fn handle_job_error( ) .await; if let Err(err) = updated_flow { + println!("error updating flow status: {}", err); + if let Some(parent_job_id) = job.parent_job { if let Ok(mut tx) = db.begin().await { if let Ok(Some(parent_job)) = @@ -514,7 +508,7 @@ async fn handle_job_error( &parent_job, format!("Unexpected error during flow job error handling:\n{err}"), err, - metrics, + metrics.clone(), ) .await; } @@ -522,6 +516,18 @@ async fn handle_job_error( } } } + add_completed_job_error( + db, + client, + &job, + format!("Unexpected error during job execution:\n{err}"), + &err, + metrics, + ) + .await + .map(|(_, m)| m) + .unwrap_or_else(|_| Map::new()); + tracing::error!(job_id = %job.id, err = err.alt(), "error handling job: {} {} {}", job.id, job.workspace_id, job.created_by); } diff --git a/backend/windmill-worker/src/worker_flow.rs b/backend/windmill-worker/src/worker_flow.rs index c35e0983af..23d8cbea0c 100644 --- a/backend/windmill-worker/src/worker_flow.rs +++ b/backend/windmill-worker/src/worker_flow.rs @@ -1346,11 +1346,13 @@ async fn compute_next_flow_transform<'c>( tx = tx_new; res }; + let flow_input = flow_job.args.clone().unwrap_or_else(|| json!({})); /* Iterator is an InputTransform, evaluate it into an array. */ let itered = evaluate_with( iterator.clone(), || { vec![ + ("flow_input".to_string(), flow_input), ("result".to_string(), last_result.clone()), ("previous_result".to_string(), last_result.clone()), ]