From be251013775d91e59c733f43cb0c30b829c8eddb Mon Sep 17 00:00:00 2001 From: Ruben Fiszel Date: Thu, 18 Jul 2024 02:10:15 +0200 Subject: [PATCH] fix: add WM_SCHEDULED_FOR to contextual variables and early stop of flows --- backend/windmill-api/src/resources.rs | 1 + backend/windmill-api/src/variables.rs | 1 + backend/windmill-common/src/variables.rs | 12 ++++++++++++ backend/windmill-worker/src/bun_executor.rs | 1 + backend/windmill-worker/src/common.rs | 2 ++ backend/windmill-worker/src/deno_executor.rs | 1 + backend/windmill-worker/src/js_eval.rs | 19 +++++++++++++++++-- .../windmill-worker/src/python_executor.rs | 2 ++ backend/windmill-worker/src/worker_flow.rs | 13 +++++++++++++ .../flows/content/FlowSettings.svelte | 9 ++++++++- frontend/src/lib/deno_fetch.d.ts.txt | 1 + frontend/src/lib/process.d.ts.txt | 1 + 12 files changed, 60 insertions(+), 3 deletions(-) diff --git a/backend/windmill-api/src/resources.rs b/backend/windmill-api/src/resources.rs index b8b03d0259..fbc5f53d00 100644 --- a/backend/windmill-api/src/resources.rs +++ b/backend/windmill-api/src/resources.rs @@ -567,6 +567,7 @@ pub async fn transform_json_value<'c>( job.flow_step_id.clone(), job.root_job.map(|x| x.to_string()), None, + Some(job.scheduled_for.clone()), ) .await; diff --git a/backend/windmill-api/src/variables.rs b/backend/windmill-api/src/variables.rs index b6f1f0ecdc..62c95647ad 100644 --- a/backend/windmill-api/src/variables.rs +++ b/backend/windmill-api/src/variables.rs @@ -75,6 +75,7 @@ async fn list_contextual_variables( Some("c".to_string()), Some("017e0ad5-f499-73b6-5488-92a61c5196dd".to_string()), Some("eyJhbGciOiJIUzI1NiIsInR5cCI6IkpXVCJ9.eyJzdWIiOiIxMjM0NTY3ODkwIiwibmFtZSI6IkpvaG4gRG9lIiwiaWF0IjoxNTE2MjM5MDIyfQ.SflKxwRJSMeKKF2QT4fwpMeJf36POk6yJV_adQssw5c".to_string()), + Some(chrono::offset::Utc::now()) ) .await .to_vec(), diff --git a/backend/windmill-common/src/variables.rs b/backend/windmill-common/src/variables.rs index a908802f4b..8239b8c321 100644 --- a/backend/windmill-common/src/variables.rs +++ b/backend/windmill-common/src/variables.rs @@ -6,6 +6,7 @@ * LICENSE-AGPL for a copy of the license. */ +use chrono::Utc; use magic_crypt::{MagicCrypt256, MagicCryptError, MagicCryptTrait}; use serde::{Deserialize, Serialize}; @@ -147,6 +148,8 @@ pub async fn decrypt_value_with_mc( })?) } +pub const WM_SCHEDULED_FOR: &str = "WM_SCHEDULED_FOR"; + pub async fn get_reserved_variables( db: &DB, w_id: &str, @@ -162,6 +165,7 @@ pub async fn get_reserved_variables( step_id: Option, root_flow_id: Option, jwt_token: Option, + scheduled_for: Option>, ) -> Vec { let state_path = { let trigger = if schedule_path.is_some() { @@ -247,6 +251,14 @@ pub async fn get_reserved_variables( description: "Job id of the current script".to_string(), is_custom: false, }, + ContextualVariable { + name: WM_SCHEDULED_FOR.to_string(), + value: scheduled_for + .map(|ts| ts.to_string()) + .unwrap_or_else(|| "".to_string()), + description: "date-time in UTC (e.g: 2014-11-28T12:45:59.324310806Z) of when the job was scheduled".to_string(), + is_custom: false, + }, ContextualVariable { name: "WM_JOB_PATH".to_string(), value: path.unwrap_or_else(|| "".to_string()), diff --git a/backend/windmill-worker/src/bun_executor.rs b/backend/windmill-worker/src/bun_executor.rs index ae349e85fe..e93716e743 100644 --- a/backend/windmill-worker/src/bun_executor.rs +++ b/backend/windmill-worker/src/bun_executor.rs @@ -945,6 +945,7 @@ pub async fn start_worker( None, None, None, + None, ) .await; let context_envs = build_envs_map(context.to_vec()).await; diff --git a/backend/windmill-worker/src/common.rs b/backend/windmill-worker/src/common.rs index 741ac9e4e6..cfd00020b8 100644 --- a/backend/windmill-worker/src/common.rs +++ b/backend/windmill-worker/src/common.rs @@ -314,6 +314,7 @@ pub async fn transform_json_value( job.flow_step_id.clone(), job.root_job.clone().map(|x| x.to_string()), None, + Some(job.scheduled_for.clone()), ) .await; @@ -411,6 +412,7 @@ pub async fn get_reserved_variables( job.flow_step_id.clone(), job.root_job.clone().map(|x| x.to_string()), None, + Some(job.scheduled_for.clone()), ) .await .to_vec(); diff --git a/backend/windmill-worker/src/deno_executor.rs b/backend/windmill-worker/src/deno_executor.rs index dec1c78aca..54462644de 100644 --- a/backend/windmill-worker/src/deno_executor.rs +++ b/backend/windmill-worker/src/deno_executor.rs @@ -435,6 +435,7 @@ pub async fn start_worker( None, None, None, + None, ) .await; let context_envs = build_envs_map(context.to_vec()).await; diff --git a/backend/windmill-worker/src/js_eval.rs b/backend/windmill-worker/src/js_eval.rs index 0fc75d0209..591219de43 100644 --- a/backend/windmill-worker/src/js_eval.rs +++ b/backend/windmill-worker/src/js_eval.rs @@ -141,6 +141,7 @@ pub async fn eval_timeout( flow_input: Option>>>, authed_client: Option<&AuthedClient>, by_id: Option, + ctx: Option>, ) -> anyhow::Result> { let expr = expr.trim().to_string(); @@ -305,6 +306,7 @@ pub async fn eval_timeout( context_keys, by_id, has_client, + ctx, ))?; Ok(r) as anyhow::Result> @@ -367,8 +369,10 @@ async fn eval( transform_context: Vec, by_id: Option, has_client: bool, + ctx: Option>, ) -> anyhow::Result> { tracing::debug!("evaluating: {} {:#?}", expr, by_id); + let (api_code, by_id_code) = if has_client { let by_id_code = if let Some(by_id) = by_id { format!( @@ -440,11 +444,21 @@ async function resource(path) {{ } else { format!("return {expr}") }; + + let ctx_str = ctx + .map(|x| { + x.into_iter() + .map(|(k, v)| format!("let {} = \"{}\";", k, v)) + .join("\n") + }) + .unwrap_or_default(); let code = format!( r#" function get_from_env(name) {{ return JSON.parse(Deno.core.ops.op_get_context(name)); }} +{ctx_str} + {api_code} {} {} @@ -868,6 +882,7 @@ mod tests { vec!["params".to_string(), "value".to_string()], None, false, + None, ) .await?; assert_eq!(res.get(), "4"); @@ -882,7 +897,7 @@ return `my ${x} multiline template`"; let mut runtime = JsRuntime::new(RuntimeOptions::default()); - let res = eval(&mut runtime, code, env, None, false).await?; + let res = eval(&mut runtime, code, env, None, false, None).await?; assert_eq!(res.get(), "\"my 5\\nmultiline template\""); Ok(()) } @@ -908,7 +923,7 @@ multiline template`"; op_state.put(TransformContext { flow_input: None, envs: env.clone() }) } - let res = eval_timeout(code.to_string(), env, None, None, None).await?; + let res = eval_timeout(code.to_string(), env, None, None, None, None).await?; assert_eq!(res.get(), "2"); Ok(()) } diff --git a/backend/windmill-worker/src/python_executor.rs b/backend/windmill-worker/src/python_executor.rs index 72a482da3d..dd89632957 100644 --- a/backend/windmill-worker/src/python_executor.rs +++ b/backend/windmill-worker/src/python_executor.rs @@ -1049,6 +1049,7 @@ pub async fn start_worker( None, None, None, + None, ) .await .to_vec(); @@ -1167,6 +1168,7 @@ for line in sys.stdin: None, None, None, + None, ) .await; diff --git a/backend/windmill-worker/src/worker_flow.rs b/backend/windmill-worker/src/worker_flow.rs index 3787dac0fd..b839405b89 100644 --- a/backend/windmill-worker/src/worker_flow.rs +++ b/backend/windmill-worker/src/worker_flow.rs @@ -288,6 +288,7 @@ pub async fn update_flow_status_after_job_completion_internal< None, Some(client), None, + None, ) .await? } else { @@ -1174,6 +1175,7 @@ async fn compute_bool_from_expr( by_id: Option, client: Option<&AuthedClient>, resumes: Option<(Arc>, Arc>, Arc>)>, + ctx: Option>, ) -> error::Result { let mut context = HashMap::with_capacity(if resumes.is_some() { 7 } else { 3 }); context.insert("result".to_string(), result.clone()); @@ -1191,6 +1193,7 @@ async fn compute_bool_from_expr( Some(flow_args), client, by_id, + ctx, ) .await? .get() @@ -1300,6 +1303,7 @@ async fn transform_input( Some(flow_args.clone()), Some(client), Some(by_id.clone()), + None, ) .await .map_err(|e| { @@ -1552,6 +1556,10 @@ async fn push_next_flow_job None, Some(client), None, + Some(vec![( + windmill_common::variables::WM_SCHEDULED_FOR.to_string(), + flow_job.scheduled_for.to_string(), + )]), ) .await?; if skip { @@ -1667,6 +1675,7 @@ async fn push_next_flow_job Some(arc_flow_job_args.clone()), None, None, + None ) .await .map_err(|e| { @@ -1855,6 +1864,7 @@ async fn push_next_flow_job Some(arc_flow_job_args.clone()), None, None, + None, ) .await .map_err(|e| { @@ -2936,6 +2946,7 @@ async fn compute_next_flow_transform( Some(idcontext.clone()), Some(client), Some((resumes.clone(), resume.clone(), approvers.clone())), + None, ) .await?; @@ -3254,6 +3265,7 @@ async fn next_forloop_status( Some(arc_flow_job_args), Some(client), Some(by_id), + None, ) .await? } @@ -3314,6 +3326,7 @@ async fn next_forloop_status( Some(arc_flow_job_args), Some(client), Some(by_id), + None, ) .await? } diff --git a/frontend/src/lib/components/flows/content/FlowSettings.svelte b/frontend/src/lib/components/flows/content/FlowSettings.svelte index 707a925612..8644e836bf 100644 --- a/frontend/src/lib/components/flows/content/FlowSettings.svelte +++ b/frontend/src/lib/components/flows/content/FlowSettings.svelte @@ -431,8 +431,15 @@ class="small-editor" extraLib={`declare const flow_input = ${JSON.stringify( schemaToObject(asSchema($flowStore.schema), $previewArgs) - )};`} + )}; + declare const WM_SCHEDULED_FOR: string;`} /> +
+ You can use the variable `flow_input` to access the inputs of the flow.
The variable `WM_SCHEDULED_FOR` contains the time the flow was scheduled for + which you can use to stop early non fresh jobs: +
new Date().getTime() - new Date(WM_SCHEDULED_FOR).getTime() {'>'} X
+
{:else}