diff --git a/backend/windmill-common/src/flow_status.rs b/backend/windmill-common/src/flow_status.rs index fe5c1467e1..f873208f70 100644 --- a/backend/windmill-common/src/flow_status.rs +++ b/backend/windmill-common/src/flow_status.rs @@ -48,10 +48,23 @@ pub struct FlowStatus { pub trait FlowStatusGetter { fn get_raw_flow_status(&self) -> Option<&sqlx::types::Json>>; } + +#[macro_export] +macro_rules! impl_flow_status_getter { + ($struct_name:ident) => { + impl FlowStatusGetter for $struct_name { + fn get_raw_flow_status( + &self, + ) -> Option<&sqlx::types::Json>> { + self.flow_status.as_ref() + } + } + }; +} + pub trait ParsedFlowStatusGetter { fn parse_flow_status(&self) -> Option; } - impl ParsedFlowStatusGetter for I { fn parse_flow_status(&self) -> Option { self.get_raw_flow_status() diff --git a/backend/windmill-common/src/flows.rs b/backend/windmill-common/src/flows.rs index 05efc6e136..790cd27e90 100644 --- a/backend/windmill-common/src/flows.rs +++ b/backend/windmill-common/src/flows.rs @@ -127,6 +127,20 @@ pub struct FlowValue { pub trait FlowValueGetter { fn get_raw_flow_value(&self) -> Option<&sqlx::types::Json>>; } + +#[macro_export] +macro_rules! impl_flow_value_getter { + ($struct_name:ident) => { + impl FlowValueGetter for $struct_name { + fn get_raw_flow_value( + &self, + ) -> Option<&sqlx::types::Json>> { + self.raw_flow.as_ref() + } + } + }; +} + pub trait ParsedFlowValueGetter { fn parse_raw_flow(&self) -> Option; } diff --git a/backend/windmill-common/src/jobs.rs b/backend/windmill-common/src/jobs.rs index 71a0fe8423..fa61792190 100644 --- a/backend/windmill-common/src/jobs.rs +++ b/backend/windmill-common/src/jobs.rs @@ -17,7 +17,7 @@ use crate::{ error::{self, to_anyhow, Error}, flow_status::{FlowStatusGetter, RestartedFrom}, flows::{FlowValue, FlowValueGetter, Retry}, - get_latest_deployed_hash_for_path, + get_latest_deployed_hash_for_path, impl_flow_status_getter, impl_flow_value_getter, scripts::{ScriptHash, ScriptLang}, worker::{to_raw_value, TMP_DIR}, }; @@ -138,17 +138,8 @@ impl QueuedJob { } } -impl FlowValueGetter for QueuedJob { - fn get_raw_flow_value(&self) -> Option<&sqlx::types::Json>> { - self.raw_flow.as_ref() - } -} - -impl FlowStatusGetter for QueuedJob { - fn get_raw_flow_status(&self) -> Option<&sqlx::types::Json>> { - self.flow_status.as_ref() - } -} +impl_flow_status_getter!(QueuedJob); +impl_flow_value_getter!(QueuedJob); impl Default for QueuedJob { fn default() -> Self { @@ -261,22 +252,12 @@ impl CompletedJob { pub fn json_result(&self) -> Option { self.result .as_ref() - .map(|r| serde_json::from_str(r.get()).ok()) - .flatten() + .and_then(|r| serde_json::from_str(r.get()).ok()) } } -impl FlowValueGetter for CompletedJob { - fn get_raw_flow_value(&self) -> Option<&sqlx::types::Json>> { - self.raw_flow.as_ref() - } -} - -impl FlowStatusGetter for CompletedJob { - fn get_raw_flow_status(&self) -> Option<&sqlx::types::Json>> { - self.flow_status.as_ref() - } -} +impl_flow_status_getter!(CompletedJob); +impl_flow_value_getter!(CompletedJob); #[derive(sqlx::FromRow)] pub struct BranchResults { diff --git a/backend/windmill-queue/src/jobs.rs b/backend/windmill-queue/src/jobs.rs index eaa421fb3c..54f82f1958 100644 --- a/backend/windmill-queue/src/jobs.rs +++ b/backend/windmill-queue/src/jobs.rs @@ -6,8 +6,6 @@ * LICENSE-AGPL for a copy of the license. */ -use std::{borrow::Borrow, collections::HashMap, sync::Arc, vec}; - use anyhow::Context; use async_recursion::async_recursion; use axum::{ @@ -31,18 +29,21 @@ use serde_json::{json, value::RawValue}; use sqlx::{types::Json, FromRow, Pool, Postgres, Transaction}; #[cfg(feature = "benchmark")] use std::time::Instant; +use std::{borrow::Borrow, collections::HashMap, sync::Arc, vec}; use tokio::{sync::RwLock, time::sleep}; use tracing::{instrument, Instrument}; use ulid::Ulid; use uuid::Uuid; use windmill_audit::audit_ee::{audit_log, AuditAuthor}; use windmill_audit::ActionKind; +use windmill_common::flows::FlowValueGetter; use windmill_common::{ add_time, auth::{fetch_authed_from_permissioned_as, permissioned_as_to_username}, db::{Authed, UserDB}, error::{self, to_anyhow, Error}, + fetch_one_with_fallback, flow_status::{ BranchAllStatus, FlowCleanupModule, FlowStatus, FlowStatusGetter, FlowStatusModule, FlowStatusModuleWParent, Iterator, JobResult, ParsedFlowStatusGetter, RestartedFrom, @@ -52,6 +53,7 @@ use windmill_common::{ add_virtual_items_if_necessary, FlowModule, FlowModuleValue, FlowValue, InputTransform, ParsedFlowValueGetter, }, + impl_flow_status_getter, impl_flow_value_getter, jobs::{ get_payload_tag_from_prefixed_path, CompletedJob, JobKind, JobPayload, QueuedJob, RawCode, ENTRYPOINT_OVERRIDE, PREPROCESSOR_FAKE_ENTRYPOINT, @@ -2502,26 +2504,28 @@ async fn extract_result_from_job_result( match job_result { JobResult::ListJob(job_ids) => match json_path { Some(json_path) => { - let mut parts = json_path.split("."); + let mut parts = json_path.split('.'); - let Some(idx) = parts.next().and_then(|x| x.parse::().ok()) else { + let Some(ref idx) = parts.next().and_then(|x| x.parse::().ok()) else { return Ok(to_raw_value(&serde_json::Value::Null)); }; - let Some(job_id) = job_ids.get(idx).cloned() else { + let Some(job_id) = job_ids.get(*idx) else { return Ok(to_raw_value(&serde_json::Value::Null)); }; - // ici - Ok(sqlx::query_as::<_, ResultR>( - "SELECT result #> $3 as result FROM completed_job WHERE id = $1 AND workspace_id = $2", - ) - .bind(job_id) - .bind(w_id) - .bind( - parts.map(|x| x.to_string()).collect::>() - ) - .fetch_optional(db) - .await? + let parts = parts.map(|x| x.to_string()).collect_vec(); + + Ok(fetch_one_with_fallback!( + db, + query_as, + fetch_optional, + ResultR, + "SELECT result #> $3 as result FROM {} WHERE id = $1 AND workspace_id = $2", + "completed_jobs_result" || "completed_job", + *job_id, + w_id, + &parts + )? .and_then(|r| r.result.map(|x| x.0)) .unwrap_or_else(|| to_raw_value(&serde_json::Value::Null))) } @@ -2536,6 +2540,7 @@ async fn extract_result_from_job_result( .into_iter() .filter_map(|x| x.result.map(|y| (x.id, y))) .collect::>>>(); + let result = job_ids .into_iter() .map(|id| { @@ -2543,25 +2548,32 @@ async fn extract_result_from_job_result( .map(|x| x.0.clone()) .unwrap_or_else(|| to_raw_value(&serde_json::Value::Null)) }) - .collect::>(); + .collect_vec(); Ok(to_raw_value(&result)) } }, // ici - JobResult::SingleJob(x) => Ok(sqlx::query_as::<_, ResultR>( - "SELECT result #> $3 as result FROM completed_job WHERE id = $1 AND workspace_id = $2", - ) - .bind(x) - .bind(w_id) - .bind( - json_path + JobResult::SingleJob(x) => { + let path = json_path .map(|x| x.split(".").map(|x| x.to_string()).collect::>()) - .unwrap_or_default(), - ) - .fetch_optional(db) - .await? - .and_then(|r| r.result.map(|x| x.0)) - .unwrap_or_else(|| to_raw_value(&serde_json::Value::Null))), + .unwrap_or_default(); + + let res = fetch_one_with_fallback!( + db, + query_as, + fetch_optional, + ResultR, + "SELECT result #> $3 as result FROM {} WHERE id = $1 AND workspace_id = $2", + "completed_jobs_result" || "completed_job", + x, + w_id, + &path + )? + .and_then(|r| r.result.map(|x| x.0)) + .unwrap_or_else(|| to_raw_value(&serde_json::Value::Null)); + + Ok(res) + } } } @@ -4123,11 +4135,8 @@ pub async fn push<'c, 'd, R: rsmq_async::RsmqConnection + Send + 'c>( } pub fn canceled_job_to_result(job: &QueuedJob) -> serde_json::Value { - let reason = job - .canceled_reason - .as_deref() - .unwrap_or_else(|| "no reason given"); - let canceler = job.canceled_by.as_deref().unwrap_or_else(|| "unknown"); + let reason = job.canceled_reason.as_deref().unwrap_or("no reason given"); + let canceler = job.canceled_by.as_deref().unwrap_or("unknown"); serde_json::json!({"message": format!("Job canceled: {reason} by {canceler}"), "name": "Canceled", "reason": reason, "canceler": canceler}) } @@ -4150,8 +4159,19 @@ async fn restarted_flows_resolution( ), Error, > { - let completed_job = sqlx::query_as::<_, CompletedJob>( - "SELECT *, null as labels FROM completed_job WHERE id = $1 and workspace_id = $2", + #[derive(Debug, sqlx::FromRow)] + struct QueryResults { + pub script_path: Option, + pub priority: Option, + pub raw_flow: Option>>, + pub flow_status: Option>>, + } + + impl_flow_status_getter!(QueryResults); + impl_flow_value_getter!(QueryResults); + + let completed_job = sqlx::query_as::<_, QueryResults>( + "SELECT script_path, priority, raw_flow, flow_status FROM completed_job WHERE id = $1 and workspace_id = $2", ) .bind(completed_flow_id) .bind(workspace_id) @@ -4184,10 +4204,9 @@ async fn restarted_flows_resolution( if flow_value_if_any .clone() .map(|fv| { - fv.modules + !fv.modules .iter() - .find(|flow_value_module| flow_value_module.id == module.id()) - .is_none() + .any(|flow_value_module| flow_value_module.id == module.id()) }) .unwrap_or(false) { @@ -4303,7 +4322,7 @@ async fn restarted_flows_resolution( truncated_modules.push(FlowStatusModule::WaitingForPriorSteps { id: module.id() }); } else { // else we simply "transfer" the module from the completed flow to the new one if it's a success - step_n = step_n + 1; + step_n += 1; match module.clone() { FlowStatusModule::Success { .. } => Ok(truncated_modules.push(module)), _ => Err(Error::InternalErr(format!( diff --git a/backend/windmill-worker/src/worker_flow.rs b/backend/windmill-worker/src/worker_flow.rs index f1f15c7d17..fb29481120 100644 --- a/backend/windmill-worker/src/worker_flow.rs +++ b/backend/windmill-worker/src/worker_flow.rs @@ -1220,7 +1220,7 @@ fn get_module(flow_job: &QueuedJob, module_step: &Step) -> Option { if let Some(raw_flow) = raw_flow { match module_step { Step::PreprocessorStep => raw_flow.preprocessor_module.map(|x| *x.clone()), - Step::Step(i) => raw_flow.modules.get(*i).map(|x| x.clone()), + Step::Step(i) => raw_flow.modules.get(*i).cloned(), Step::FailureStep => raw_flow.failure_module.map(|x| *x.clone()), } } else {