diff --git a/backend/tests/worker.rs b/backend/tests/worker.rs index d4fb929cba..c5facab96c 100644 --- a/backend/tests/worker.rs +++ b/backend/tests/worker.rs @@ -1,8 +1,9 @@ use futures::{stream, Stream}; use serde_json::json; -use sqlx::{postgres::PgListener, query_scalar, types::Uuid, Pool, Postgres, Transaction}; +use sqlx::{postgres::PgListener, types::Uuid, Pool, Postgres, Transaction}; use windmill_api::jobs::{CompletedJob, Job}; use windmill_common::{ + flow_status::{FlowStatus, FlowStatusModule}, flows::{FlowModule, FlowModuleValue, FlowValue, InputTransform}, scripts::ScriptLang, DEFAULT_SLEEP_QUEUE, @@ -107,6 +108,29 @@ impl ApiServer { } } +async fn _print_job(id: Uuid, db: &Pool) -> Result<(), anyhow::Error> { + tracing::info!( + "{:#?}", + get_job_by_id(db.begin().await?, "test-workspace", id) + .await? + .0 + ); + Ok(()) +} + +fn get_module(cjob: &CompletedJob, id: &str) -> Option { + cjob.flow_status.clone().and_then(|fs| { + find_module_in_vec( + serde_json::from_value::(fs).unwrap().modules, + id, + ) + }) +} + +fn find_module_in_vec(modules: Vec, id: &str) -> Option { + modules.into_iter().find(|s| s.id() == id) +} + mod suspend_resume { use futures::{Stream, StreamExt}; @@ -135,16 +159,6 @@ mod suspend_resume { } } - async fn _print_job(id: Uuid, db: &Pool) -> Result<(), anyhow::Error> { - tracing::info!( - "{:#?}", - get_job_by_id(db.begin().await?, "test-workspace", id) - .await? - .0 - ); - Ok(()) - } - fn flow() -> FlowValue { serde_json::from_value(serde_json::json!({ "modules": [{ @@ -273,7 +287,7 @@ mod suspend_resume { server.close().await.unwrap(); - let result = completed_job_result(flow, &db).await; + let result = completed_job(flow, &db).await.result.unwrap(); assert_eq!( json!({ @@ -310,7 +324,9 @@ mod suspend_resume { .arg("op", json!("cancel")) .arg("port", json!(port)) .run_until_complete(&db, port) - .await; + .await + .result + .unwrap(); server.close().await.unwrap(); @@ -374,7 +390,7 @@ mod suspend_resume { server.close().await.unwrap(); - let result = completed_job_result(flow, &db).await; + let result = completed_job(flow, &db).await.result.unwrap(); assert_eq!( json!({"error": "Job canceled: approval request disapproved by unknown" }), @@ -528,7 +544,9 @@ def main(last, port): .arg("items", json!(["unused", "unused", "unused"])) .arg("port", json!(server.addr.port())) .run_until_complete(&db, server.addr.port()) - .await; + .await + .result + .unwrap(); assert_eq!(server.close().await, attempts); assert_eq!(json!([3, 5, 7, 9]), result); @@ -555,7 +573,9 @@ def main(last, port): .arg("items", json!(["unused", "unused", "unused"])) .arg("port", json!(server.addr.port())) .run_until_complete(&db, server.addr.port()) - .await; + .await + .result + .unwrap(); assert_eq!(server.close().await, attempts); assert!(result["error"] @@ -587,7 +607,9 @@ def main(last, port): .arg("items", json!(["unused", "unused", "unused"])) .arg("port", json!(server.addr.port())) .run_until_complete(&db, server.addr.port()) - .await; + .await + .result + .unwrap(); assert_eq!(server.close().await, attempts); assert!(result["error"] @@ -604,8 +626,9 @@ def main(last, port): let value = serde_json::from_value(json!({ "modules": [{ - "input_transform": { "port": { "type": "javascript", "expr": "flow_input.port" } }, + "id": "a", "value": { + "input_transform": { "port": { "type": "javascript", "expr": "flow_input.port" } }, "type": "rawscript", "language": "python3", "content": r#" @@ -617,9 +640,9 @@ def main(port): "retry": { "constant": { "attempts": 1, "seconds": 0 } }, }], "failure_module": { - "input_transform": { "error": { "type": "javascript", "expr": "previous_result", }, - "port": { "type": "javascript", "expr": "flow_input.port" } }, "value": { + "input_transform": { "error": { "type": "javascript", "expr": "previous_result", }, + "port": { "type": "javascript", "expr": "flow_input.port" } }, "type": "rawscript", "language": "python3", "content": r#" @@ -643,10 +666,16 @@ def main(error, port): .into_iter() .unzip::<_, _, Vec<_>, Vec<_>>(); let server = Server::start(responses).await; - let result = RunJob::from(JobPayload::RawFlow { value, path: None }) + let cjob = RunJob::from(JobPayload::RawFlow { value, path: None }) .arg("port", json!(server.addr.port())) .run_until_complete(&db, server.addr.port()) .await; + let result = cjob.result.clone().unwrap(); + let failed_module = get_module(&cjob, "a").unwrap(); + match failed_module { + FlowStatusModule::Failure { .. } => {} + _ => panic!("expected failure module"), + } assert_eq!(server.close().await, attempts); assert_eq!( @@ -694,14 +723,18 @@ async fn test_iteration(db: Pool) { let result = RunJob::from(JobPayload::RawFlow { value: flow.clone(), path: None }) .arg("items", json!([])) .run_until_complete(&db, server.addr.port()) - .await; + .await + .result + .unwrap(); assert_eq!(result, serde_json::json!([])); /* Don't actually test that this does 257 jobs or that will take forever. */ let result = RunJob::from(JobPayload::RawFlow { value: flow.clone(), path: None }) .arg("items", json!((0..257).collect::>())) .run_until_complete(&db, server.addr.port()) - .await; + .await + .result + .unwrap(); assert!(matches!(result, serde_json::Value::Object(_))); assert!(result["error"] .as_str() @@ -751,11 +784,11 @@ impl RunJob { } /// push the job, spawn a worker, wait until the job is in completed_job - async fn run_until_complete(self, db: &Pool, port: u16) -> serde_json::Value { + async fn run_until_complete(self, db: &Pool, port: u16) -> CompletedJob { let uuid = self.push(db).await; let listener = listen_for_completed_jobs(db).await; in_test_worker(db, listener.find(&uuid), port).await; - completed_job_result(uuid, db).await + completed_job(uuid, db).await } } @@ -763,7 +796,7 @@ async fn run_job_in_new_worker_until_complete( db: &Pool, job: JobPayload, port: u16, -) -> serde_json::Value { +) -> CompletedJob { RunJob::from(job).run_until_complete(db, port).await } @@ -880,8 +913,8 @@ async fn listen_for_uuid_on( })) } -async fn completed_job_result(uuid: Uuid, db: &Pool) -> serde_json::Value { - query_scalar("SELECT result FROM completed_job WHERE id = $1") +async fn completed_job(uuid: Uuid, db: &Pool) -> CompletedJob { + sqlx::query_as::<_, CompletedJob>("SELECT * FROM completed_job WHERE id = $1") .bind(uuid) .fetch_one(db) .await @@ -974,7 +1007,10 @@ async fn test_deno_flow(db: Pool) { for i in 0..50 { println!("deno flow iteration: {}", i); - let result = run_job_in_new_worker_until_complete(&db, job.clone(), port).await; + let result = run_job_in_new_worker_until_complete(&db, job.clone(), port) + .await + .result + .unwrap(); assert_eq!(result, serde_json::json!([2, 4, 6]), "iteration: {}", i); } } @@ -1129,7 +1165,10 @@ async fn test_deno_flow_same_worker(db: Pool) { let job = JobPayload::RawFlow { value: flow, path: None }; - let result = run_job_in_new_worker_until_complete(&db, job.clone(), server.addr.port()).await; + let result = run_job_in_new_worker_until_complete(&db, job.clone(), server.addr.port()) + .await + .result + .unwrap(); assert_eq!( result, serde_json::json!("false 1,true 1,false 1,true 2,false 1,true 3,false 1,true 3") @@ -1181,7 +1220,10 @@ async fn test_flow_result_by_id(db: Pool) { .unwrap(); let job = JobPayload::RawFlow { value: flow, path: None }; - let result = run_job_in_new_worker_until_complete(&db, job.clone(), port).await; + let result = run_job_in_new_worker_until_complete(&db, job.clone(), port) + .await + .result + .unwrap(); assert_eq!(result, serde_json::json!([[42]])); } @@ -1195,8 +1237,9 @@ async fn test_stop_after_if(db: Pool) { let flow: FlowValue = serde_json::from_value(serde_json::json!({ "modules": [ { - "input_transforms": { "n": { "type": "javascript", "expr": "flow_input.n" } }, + "id": "a", "value": { + "input_transforms": { "n": { "type": "javascript", "expr": "flow_input.n" } }, "type": "rawscript", "language": "python3", "content": "def main(n): return n", @@ -1207,8 +1250,9 @@ async fn test_stop_after_if(db: Pool) { }, }, { - "input_transforms": { "n": { "type": "javascript", "expr": "previous_result" } }, + "id": "b", "value": { + "input_transforms": { "n": { "type": "javascript", "expr": "previous_result" } }, "type": "rawscript", "language": "python3", "content": "def main(n): return f'last step saw {n}'", @@ -1222,16 +1266,78 @@ async fn test_stop_after_if(db: Pool) { let result = RunJob::from(job.clone()) .arg("n", json!(123)) .run_until_complete(&db, port) - .await; + .await + .result + .unwrap(); assert_eq!(json!("last step saw 123"), result); - let result = RunJob::from(job.clone()) + let cjob = RunJob::from(job.clone()) .arg("n", json!(-123)) .run_until_complete(&db, port) .await; + + let result = cjob.result.unwrap(); assert_eq!(json!(-123), result); } +#[sqlx::test(fixtures("base"))] +async fn test_stop_after_if_nested(db: Pool) { + initialize_tracing().await; + // let server = ApiServer::start(db.clone()).await; + // let port = server.addr.port(); + + let port = 123; + let flow: FlowValue = serde_json::from_value(serde_json::json!({ + "modules": [ + { + "id": "a", + "value": { + "branches": [{"modules": [{ + "id": "b", + "value": { + "input_transforms": { "n": { "type": "javascript", "expr": "flow_input.n" } }, + "type": "rawscript", + "language": "python3", + "content": "def main(n): return n", + }, + "stop_after_if": { + "expr": "result < 0", + "skip_if_stopped": false, + }}]}], + "type": "branchall" + }, + }, + { + "id": "c", + "value": { + "input_transforms": { "n": { "type": "javascript", "expr": "previous_result" } }, + "type": "rawscript", + "language": "python3", + "content": "def main(n): return f'last step saw {n}'", + }, + }, + ], + })) + .unwrap(); + let job = JobPayload::RawFlow { value: flow, path: None }; + + let result = RunJob::from(job.clone()) + .arg("n", json!(123)) + .run_until_complete(&db, port) + .await + .result + .unwrap(); + assert_eq!(json!("last step saw [123]"), result); + + let cjob = RunJob::from(job.clone()) + .arg("n", json!(-123)) + .run_until_complete(&db, port) + .await; + + let result = cjob.result.unwrap(); + assert_eq!(json!([-123]), result); +} + #[sqlx::test(fixtures("base"))] async fn test_python_flow(db: Pool) { initialize_tracing().await; @@ -1282,7 +1388,9 @@ async fn test_python_flow(db: Pool) { JobPayload::RawFlow { value: flow.clone(), path: None }, port, ) - .await; + .await + .result + .unwrap(); assert_eq!(result, serde_json::json!([2, 4, 6]), "iteration: {i}"); } @@ -1315,7 +1423,9 @@ async fn test_python_flow_2(db: Pool) { JobPayload::RawFlow { value: flow.clone(), path: None }, port, ) - .await; + .await + .result + .unwrap(); assert_eq!(result, serde_json::json!("Hello"), "iteration: {i}"); } @@ -1344,7 +1454,9 @@ func main(derp string) (string, error) { })) .arg("derp", json!("world")) .run_until_complete(&db, port) - .await; + .await + .result + .unwrap(); assert_eq!(result, serde_json::json!("hello world")); } @@ -1363,7 +1475,10 @@ def main(): let job = JobPayload::Code(RawCode { content, path: None, language: ScriptLang::Python3 }); - let result = run_job_in_new_worker_until_complete(&db, job, port).await; + let result = run_job_in_new_worker_until_complete(&db, job, port) + .await + .result + .unwrap(); assert_eq!(result, serde_json::json!("hello world")); } @@ -1385,7 +1500,10 @@ def main(): let job = JobPayload::Code(RawCode { content, path: None, language: ScriptLang::Python3 }); - let result = run_job_in_new_worker_until_complete(&db, job, port).await; + let result = run_job_in_new_worker_until_complete(&db, job, port) + .await + .result + .unwrap(); assert_eq!(result, serde_json::json!(3)); } @@ -1406,7 +1524,10 @@ def main(): let job = JobPayload::Code(RawCode { content, path: None, language: ScriptLang::Python3 }); - let result = run_job_in_new_worker_until_complete(&db, job, port).await; + let result = run_job_in_new_worker_until_complete(&db, job, port) + .await + .result + .unwrap(); assert_eq!(result, serde_json::json!("test-workspace")); } @@ -1458,7 +1579,10 @@ async fn test_empty_loop(db: Pool) { .unwrap(); let flow = JobPayload::RawFlow { value: flow, path: None }; - let result = run_job_in_new_worker_until_complete(&db, flow, port).await; + let result = run_job_in_new_worker_until_complete(&db, flow, port) + .await + .result + .unwrap(); assert_eq!(result, serde_json::json!(0)); } @@ -1497,7 +1621,10 @@ async fn test_empty_loop_2(db: Pool) { .unwrap(); let flow = JobPayload::RawFlow { value: flow, path: None }; - let result = run_job_in_new_worker_until_complete(&db, flow, port).await; + let result = run_job_in_new_worker_until_complete(&db, flow, port) + .await + .result + .unwrap(); assert_eq!(result, serde_json::json!([])); } @@ -1548,7 +1675,10 @@ async fn test_step_after_loop(db: Pool) { .unwrap(); let flow = JobPayload::RawFlow { value: flow, path: None }; - let result = run_job_in_new_worker_until_complete(&db, flow, port).await; + let result = run_job_in_new_worker_until_complete(&db, flow, port) + .await + .result + .unwrap(); assert_eq!(result, serde_json::json!(9)); } @@ -1611,7 +1741,10 @@ async fn test_branchone_simple(db: Pool) { .unwrap(); let flow = JobPayload::RawFlow { value: flow, path: None }; - let result = run_job_in_new_worker_until_complete(&db, flow, port).await; + let result = run_job_in_new_worker_until_complete(&db, flow, port) + .await + .result + .unwrap(); assert_eq!(result, serde_json::json!([1, 2])); } @@ -1644,7 +1777,10 @@ async fn test_branchall_simple(db: Pool) { .unwrap(); let flow = JobPayload::RawFlow { value: flow, path: None }; - let result = run_job_in_new_worker_until_complete(&db, flow, port).await; + let result = run_job_in_new_worker_until_complete(&db, flow, port) + .await + .result + .unwrap(); assert_eq!(result, serde_json::json!([[1, 2], [1, 3]])); } @@ -1677,7 +1813,10 @@ async fn test_branchall_skip_failure(db: Pool) { .unwrap(); let flow = JobPayload::RawFlow { value: flow, path: None }; - let result = run_job_in_new_worker_until_complete(&db, flow, port).await; + let result = run_job_in_new_worker_until_complete(&db, flow, port) + .await + .result + .unwrap(); assert_eq!( result, @@ -1707,7 +1846,10 @@ async fn test_branchall_skip_failure(db: Pool) { .unwrap(); let flow = JobPayload::RawFlow { value: flow, path: None }; - let result = run_job_in_new_worker_until_complete(&db, flow, port).await; + let result = run_job_in_new_worker_until_complete(&db, flow, port) + .await + .result + .unwrap(); assert_eq!( result, @@ -1763,7 +1905,10 @@ async fn test_branchone_nested(db: Pool) { .unwrap(); let flow = JobPayload::RawFlow { value: flow, path: None }; - let result = run_job_in_new_worker_until_complete(&db, flow, port).await; + let result = run_job_in_new_worker_until_complete(&db, flow, port) + .await + .result + .unwrap(); assert_eq!(result, serde_json::json!([1, 2, 3])); } @@ -1813,7 +1958,10 @@ async fn test_branchall_nested(db: Pool) { .unwrap(); let flow = JobPayload::RawFlow { value: flow, path: None }; - let result = run_job_in_new_worker_until_complete(&db, flow, port).await; + let result = run_job_in_new_worker_until_complete(&db, flow, port) + .await + .result + .unwrap(); assert_eq!( result, @@ -1873,7 +2021,9 @@ async fn test_failure_module(db: Pool) { let result = RunJob::from(JobPayload::RawFlow { value: flow.clone(), path: None }) .arg("n", json!(0)) .run_until_complete(&db, port) - .await; + .await + .result + .unwrap(); assert!(result["from failure module"]["error"] .as_str() .unwrap() @@ -1882,7 +2032,9 @@ async fn test_failure_module(db: Pool) { let result = RunJob::from(JobPayload::RawFlow { value: flow.clone(), path: None }) .arg("n", json!(1)) .run_until_complete(&db, port) - .await; + .await + .result + .unwrap(); assert!(result["from failure module"]["error"] .as_str() .unwrap() @@ -1891,7 +2043,9 @@ async fn test_failure_module(db: Pool) { let result = RunJob::from(JobPayload::RawFlow { value: flow.clone(), path: None }) .arg("n", json!(2)) .run_until_complete(&db, port) - .await; + .await + .result + .unwrap(); assert!(result["from failure module"]["error"] .as_str() .unwrap() @@ -1900,6 +2054,8 @@ async fn test_failure_module(db: Pool) { let result = RunJob::from(JobPayload::RawFlow { value: flow.clone(), path: None }) .arg("n", json!(3)) .run_until_complete(&db, port) - .await; + .await + .result + .unwrap(); assert_eq!(json!({ "l": [0, 1, 2] }), result); } diff --git a/backend/windmill-api/src/jobs.rs b/backend/windmill-api/src/jobs.rs index 92bc126279..08121d601b 100644 --- a/backend/windmill-api/src/jobs.rs +++ b/backend/windmill-api/src/jobs.rs @@ -21,12 +21,12 @@ use sqlx::{query_scalar, types::Uuid, Postgres, Transaction}; use windmill_audit::{audit_log, ActionKind}; use windmill_common::{ error::{self, to_anyhow, Error}, + flow_status::{Approval, FlowStatus, FlowStatusModule}, flows::FlowValue, oauth2::HmacSha256, scripts::{ScriptHash, ScriptLang}, users::owner_to_token_owner, utils::{not_found_if_none, now_from_db, paginate, require_admin, Pagination, StripPath}, - worker_flow::{Approval, FlowStatus, FlowStatusModule}, }; use windmill_queue::{get_queued_job, push, JobKind, JobPayload, QueuedJob, RawCode}; @@ -208,32 +208,32 @@ pub async fn get_job_by_id<'c>( #[derive(Debug, sqlx::FromRow, Serialize)] pub struct CompletedJob { - workspace_id: String, - id: Uuid, - parent_job: Option, - created_by: String, - created_at: chrono::DateTime, - started_at: chrono::DateTime, - duration_ms: i32, - success: bool, - script_hash: Option, - script_path: Option, - args: Option, - result: Option, - logs: Option, - deleted: bool, - raw_code: Option, - canceled: bool, - canceled_by: Option, - canceled_reason: Option, - job_kind: JobKind, - schedule_path: Option, - permissioned_as: String, - flow_status: Option, - raw_flow: Option, - is_flow_step: bool, - language: Option, - is_skipped: bool, + pub workspace_id: String, + pub id: Uuid, + pub parent_job: Option, + pub created_by: String, + pub created_at: chrono::DateTime, + pub started_at: chrono::DateTime, + pub duration_ms: i32, + pub success: bool, + pub script_hash: Option, + pub script_path: Option, + pub args: Option, + pub result: Option, + pub logs: Option, + pub deleted: bool, + pub raw_code: Option, + pub canceled: bool, + pub canceled_by: Option, + pub canceled_reason: Option, + pub job_kind: JobKind, + pub schedule_path: Option, + pub permissioned_as: String, + pub flow_status: Option, + pub raw_flow: Option, + pub is_flow_step: bool, + pub language: Option, + pub is_skipped: bool, } #[derive(Deserialize, Clone, Copy)] diff --git a/backend/windmill-common/src/worker_flow.rs b/backend/windmill-common/src/flow_status.rs similarity index 94% rename from backend/windmill-common/src/worker_flow.rs rename to backend/windmill-common/src/flow_status.rs index c6d1f2b66d..8d9d739c16 100644 --- a/backend/windmill-common/src/worker_flow.rs +++ b/backend/windmill-common/src/flow_status.rs @@ -100,6 +100,8 @@ pub enum FlowStatusModule { flow_jobs: Option>, #[serde(skip_serializing_if = "Option::is_none")] branch_chosen: Option, + #[serde(default)] + #[serde(skip_serializing_if = "Vec::is_empty")] approvers: Vec, }, Failure { @@ -145,7 +147,13 @@ impl FlowStatus { .iter() .map(|m| FlowStatusModule::WaitingForPriorSteps { id: m.id.clone() }) .collect(), - failure_module: FlowStatusModule::WaitingForPriorSteps { id: "failure".to_string() }, + failure_module: FlowStatusModule::WaitingForPriorSteps { + id: f + .failure_module + .as_ref() + .map(|x| x.id.clone()) + .unwrap_or_else(|| "failure".to_string()), + }, retry: RetryStatus { fail_count: 0, previous_result: None, failed_jobs: vec![] }, } } diff --git a/backend/windmill-common/src/lib.rs b/backend/windmill-common/src/lib.rs index 9318acfc31..b9d27cdc0e 100644 --- a/backend/windmill-common/src/lib.rs +++ b/backend/windmill-common/src/lib.rs @@ -12,6 +12,7 @@ use error::Error; pub mod error; pub mod external_ip; +pub mod flow_status; pub mod flows; pub mod more_serde; pub mod oauth2; @@ -19,7 +20,6 @@ pub mod scripts; pub mod users; pub mod utils; pub mod variables; -pub mod worker_flow; #[cfg(feature = "tracing_init")] pub mod tracing_init; diff --git a/backend/windmill-queue/src/jobs.rs b/backend/windmill-queue/src/jobs.rs index 65e5189667..890746b040 100644 --- a/backend/windmill-queue/src/jobs.rs +++ b/backend/windmill-queue/src/jobs.rs @@ -16,10 +16,10 @@ 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}, flows::FlowValue, scripts::{get_full_hub_script_by_path, HubScript, ScriptHash, ScriptLang}, utils::StripPath, - worker_flow::{init_flow_status, FlowStatus, MAX_RETRY_ATTEMPTS, MAX_RETRY_INTERVAL}, }; lazy_static::lazy_static! { @@ -200,7 +200,7 @@ pub async fn delete_job( .map_err(|e| Error::InternalErr(format!("Error during deletion of job {job_id}: {e}")))? .unwrap_or(0) == 1; - tracing::debug!("Job {job_id} deletion was achieved with success: {job_removed}"); + tracing::debug!("Job {job_id} deleted: {job_removed}"); Ok(()) } diff --git a/backend/windmill-worker/src/js_eval.rs b/backend/windmill-worker/src/js_eval.rs index 21ff4915f6..716ed03789 100644 --- a/backend/windmill-worker/src/js_eval.rs +++ b/backend/windmill-worker/src/js_eval.rs @@ -238,7 +238,7 @@ async function resource(path) {{ }) .join(""), ); - tracing::debug!("{}", code); + // tracing::debug!("{}", code); let global = context.execute_script("", &code)?; let global = context.resolve_value(global).await?; diff --git a/backend/windmill-worker/src/worker.rs b/backend/windmill-worker/src/worker.rs index c2ba1dbb85..a0aa6f70a3 100644 --- a/backend/windmill-worker/src/worker.rs +++ b/backend/windmill-worker/src/worker.rs @@ -14,6 +14,7 @@ use uuid::Uuid; use windmill_common::{ error::{self, to_anyhow, Error}, scripts::{ScriptHash, ScriptLang}, + utils::rd_string, variables, }; use windmill_queue::{ @@ -57,12 +58,7 @@ pub async fn create_token_for_owner<'c>( username: &str, ) -> error::Result<(Transaction<'c, Postgres>, String)> { // TODO: Bad implementation. We should not have access to this DB here. - use rand::prelude::*; - let token: String = rand::thread_rng() - .sample_iter(&rand::distributions::Alphanumeric) - .take(30) - .map(char::from) - .collect(); + let token: String = rd_string(30); let is_super_admin = username.contains('@') && sqlx::query_scalar!( "SELECT super_admin FROM password WHERE email = $1", diff --git a/backend/windmill-worker/src/worker_flow.rs b/backend/windmill-worker/src/worker_flow.rs index aa9babec0c..c063f0e6ae 100644 --- a/backend/windmill-worker/src/worker_flow.rs +++ b/backend/windmill-worker/src/worker_flow.rs @@ -21,11 +21,11 @@ use tracing::instrument; use uuid::Uuid; use windmill_common::{ error::{self, to_anyhow, Error}, - flows::{FlowModule, FlowModuleValue, FlowValue, InputTransform, Retry, Suspend}, - worker_flow::{ + flow_status::{ Approval, BranchAllStatus, BranchChosen, FlowStatus, FlowStatusModule, RetryStatus, MAX_RETRY_ATTEMPTS, MAX_RETRY_INTERVAL, }, + flows::{FlowModule, FlowModuleValue, FlowValue, InputTransform, Retry, Suspend}, }; type DB = sqlx::Pool; @@ -81,6 +81,8 @@ pub async fn update_flow_status_after_job_completion( .and_then(|i| old_status.modules.get(i)) .unwrap_or(&old_status.failure_module); + tracing::debug!("UPDATE FLOW STATUS 2: {module_index:#?} {module_status:#?} {old_status:#?} "); + let skip_loop_failures = if matches!( module_status, FlowStatusModule::InProgress { iterator: Some(_), .. } @@ -105,7 +107,7 @@ pub async fn update_flow_status_after_job_completion( let (step_counter, new_status) = match module_status { FlowStatusModule::InProgress { - iterator: Some(windmill_common::worker_flow::Iterator { index, itered, .. }), + iterator: Some(windmill_common::flow_status::Iterator { index, itered, .. }), .. } if (*index + 1 < itered.len() && (success || skip_loop_failures)) => { (old_status.step, module_status.clone()) @@ -154,7 +156,37 @@ pub async fn update_flow_status_after_job_completion( .unwrap_or(true); let (stop_early, skip_if_stop_early) = if let Some(se) = stop_early_override { + sqlx::query!( + " + UPDATE queue + SET flow_status = JSONB_SET( + JSONB_SET(flow_status, ARRAY['modules', $1::TEXT], $2), + ARRAY['step'], $3) + WHERE id = $4 + ", + old_status.step.to_string(), + json!(new_status), + json!(step_counter), + flow + ) + .execute(&mut tx) + .await?; + (true, se) + } else if old_status.step >= old_status.modules.len() as i32 { + tracing::debug!("SET NEW STATUS: {new_status:#?} "); + sqlx::query!( + " + UPDATE queue + SET flow_status = JSONB_SET(flow_status, ARRAY['failure_module'], $1) + WHERE id = $2 + ", + json!(new_status), + flow + ) + .execute(&mut tx) + .await?; + (false, false) } else { let (stop_early_expr, skip_if_stop_early) = sqlx::query_as::< _, @@ -172,8 +204,8 @@ pub async fn update_flow_status_after_job_completion( ", ) .bind(old_status.step) - .bind(serde_json::json!(new_status)) - .bind(serde_json::json!(step_counter)) + .bind(json!(new_status)) + .bind(json!(step_counter)) .bind(flow) .fetch_one(&mut tx) .await @@ -280,7 +312,6 @@ pub async fn update_flow_status_after_job_completion( } else { "Flow job completed".to_string() }; - tracing::debug!("{skip_if_stop_early:?}"); if flow_job.canceled { add_completed_job_error( db, @@ -465,31 +496,54 @@ pub async fn update_flow_status_in_progress( job_in_progress: Uuid, ) -> error::Result<()> { let step = get_step_of_flow_status(db, flow).await?; - sqlx::query(&format!( - "UPDATE queue - SET flow_status = jsonb_set(jsonb_set(flow_status, '{{modules, {step}, job}}', $1), '{{modules, {step}, type}}', $2) - WHERE id = $3 AND workspace_id = $4", - )) - .bind(json!(job_in_progress.to_string())) - .bind(json!("InProgress")) - .bind(flow) - .bind(w_id) - .execute(db) - .await?; + if let Step::Step(step) = step { + sqlx::query(&format!( + "UPDATE queue + SET flow_status = jsonb_set(jsonb_set(flow_status, '{{modules, {step}, job}}', $1), '{{modules, {step}, type}}', $2) + WHERE id = $3 AND workspace_id = $4", + )) + .bind(json!(job_in_progress.to_string())) + .bind(json!("InProgress")) + .bind(flow) + .bind(w_id) + .execute(db) + .await?; + } else { + sqlx::query(&format!( + "UPDATE queue + SET flow_status = jsonb_set(jsonb_set(flow_status, '{{failure_module, job}}', $1), '{{failure_module, type}}', $2) + WHERE id = $3 AND workspace_id = $4", + )) + .bind(json!(job_in_progress.to_string())) + .bind(json!("InProgress")) + .bind(flow) + .bind(w_id) + .execute(db) + .await?; + } Ok(()) } +pub enum Step { + Step(i32), + FailureStep, +} #[instrument(level = "trace", skip_all)] -pub async fn get_step_of_flow_status(db: &DB, id: Uuid) -> error::Result { - let r = sqlx::query_scalar!( - "SELECT (flow_status->'step')::integer FROM queue WHERE id = $1", +pub async fn get_step_of_flow_status(db: &DB, id: Uuid) -> error::Result { + let r = sqlx::query!( + "SELECT (flow_status->'step')::integer as step, jsonb_array_length(flow_status->'modules') as len FROM queue WHERE id = $1", id ) .fetch_one(db) .await - .map_err(|e| Error::InternalErr(format!("fetching step flow status: {e}")))? - .ok_or_else(|| Error::InternalErr(format!("not found step")))?; - Ok(r) + .map_err(|e| Error::InternalErr(format!("fetching step flow status: {e}")))?; + if r.step < r.len { + Ok(Step::Step(r.step.ok_or_else(|| { + Error::InternalErr("step is null".to_string()) + })?)) + } else { + Ok(Step::FailureStep) + } } /// resumes should be in order of timestamp ascending, so that more recent are at the end @@ -520,7 +574,7 @@ async fn transform_input( let flow_input = flow_args.clone().unwrap_or_else(|| json!({})); let previous_result = flatten_previous_result(last_result.clone()); let context = vec![ - ("params".to_string(), serde_json::json!(mapped)), + ("params".to_string(), json!(mapped)), ("previous_result".to_string(), previous_result), ("flow_input".to_string(), flow_input), ( @@ -596,7 +650,7 @@ pub async fn handle_flow( &Uuid::nil(), flow_job.workspace_id.as_str(), true, - serde_json::json!({}), + json!({}), None, true, same_worker_tx, @@ -613,7 +667,7 @@ pub async fn handle_flow( serde_json::from_value::(flow_job.flow_status.clone().unwrap_or_default()) .with_context(|| format!("parse flow status {}", flow_job.id))?; - tracing::debug!("handle_flow: {:#?}", flow_job); + tracing::debug!("handle_flow: {:#?}", flow_job.flow_status); push_next_flow_job( flow_job, status, @@ -695,7 +749,7 @@ async fn push_next_flow_job( .unwrap_or_else(|| status.failure_module.clone()); tracing::debug!( - "push_next_flow_job: module: {:#?}, status: {:#?}", + "push_next_flow_job {i}: module: {:#?}, status: {:#?}", module.value, status_module ); @@ -1048,11 +1102,11 @@ async fn push_next_flow_job( FlowStatusModule::InProgress { job: uuid, - iterator: Some(windmill_common::worker_flow::Iterator { index, itered }), + iterator: Some(windmill_common::flow_status::Iterator { index, itered }), flow_jobs: Some(flow_jobs), branch_chosen: None, branchall: None, - id: module.id.clone(), + id: status_module.id(), } } NextStatus::NextBranchStep(NextBranch { mut flow_jobs, status, .. }) => { @@ -1064,7 +1118,7 @@ async fn push_next_flow_job( flow_jobs: Some(flow_jobs), branch_chosen: None, branchall: Some(status), - id: module.id.clone(), + id: status_module.id(), } } @@ -1074,23 +1128,30 @@ async fn push_next_flow_job( flow_jobs: None, branch_chosen: Some(branch), branchall: None, - id: module.id.clone(), + id: status_module.id(), }, NextStatus::NextStep => { - FlowStatusModule::WaitingForExecutor { id: module.id.clone(), job: uuid } + FlowStatusModule::WaitingForExecutor { id: status_module.id(), job: uuid } } }; - sqlx::query( + tracing::debug!("STATUS STEP: {:?} {i} {:#?}", status.step, new_status); + + let json_pointer = if i >= flow.modules.len() { + "'failure_module'" + } else { + "'modules', $1::TEXT" + }; + sqlx::query(&format!( " - UPDATE queue - SET flow_status = JSONB_SET( - JSONB_SET(flow_status, ARRAY['modules', $1::TEXT], $2), - ARRAY['step'], $3) - WHERE id = $4 - ", - ) - .bind(status.step) + UPDATE queue + SET flow_status = JSONB_SET( + JSONB_SET(flow_status, ARRAY[{json_pointer}], $2), + ARRAY['step'], $3) + WHERE id = $4 + " + )) + .bind(i as i32) .bind(json!(new_status)) .bind(json!(i)) .bind(flow_job.id) @@ -1308,7 +1369,7 @@ async fn compute_next_flow_transform<'c>( } FlowStatusModule::InProgress { - iterator: Some(windmill_common::worker_flow::Iterator { itered, index }), + iterator: Some(windmill_common::flow_status::Iterator { itered, index }), flow_jobs: Some(flow_jobs), .. } => {