diff --git a/backend/.sqlx/query-0c0f3909b80c35210fc64c685905308621f9135c2c45a2fa0531ea750387da1f.json b/backend/.sqlx/query-0c0f3909b80c35210fc64c685905308621f9135c2c45a2fa0531ea750387da1f.json new file mode 100644 index 0000000000..cec59be959 --- /dev/null +++ b/backend/.sqlx/query-0c0f3909b80c35210fc64c685905308621f9135c2c45a2fa0531ea750387da1f.json @@ -0,0 +1,25 @@ +{ + "db_name": "PostgreSQL", + "query": "\n SELECT \n CASE \n WHEN flow_version.id IS NOT NULL THEN\n (flow_version.value -> 'flow_env' -> $3) #> $4\n ELSE\n (root_job.raw_flow -> 'flow_env' -> $3) #> $4\n END AS \"flow_env: sqlx::types::Json>\"\n FROM \n v2_job current_job\n JOIN \n v2_job root_job ON root_job.id = COALESCE(current_job.root_job, current_job.flow_innermost_root_job, current_job.parent_job, current_job.id)\n AND root_job.workspace_id = current_job.workspace_id\n LEFT JOIN\n flow_version ON flow_version.id = root_job.runnable_id\n AND flow_version.path = root_job.runnable_path\n AND flow_version.workspace_id = root_job.workspace_id\n WHERE \n current_job.id = $1 AND \n current_job.workspace_id = $2", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "flow_env: sqlx::types::Json>", + "type_info": "Jsonb" + } + ], + "parameters": { + "Left": [ + "Uuid", + "Text", + "Text", + "TextArray" + ] + }, + "nullable": [ + null + ] + }, + "hash": "0c0f3909b80c35210fc64c685905308621f9135c2c45a2fa0531ea750387da1f" +} diff --git a/backend/.sqlx/query-5a219a2532517869578c4504ff3153c43903f929ae5d62fbba12610f89c36d55.json b/backend/.sqlx/query-5a219a2532517869578c4504ff3153c43903f929ae5d62fbba12610f89c36d55.json index 713ccb9dd3..36ddb8ab9f 100644 --- a/backend/.sqlx/query-5a219a2532517869578c4504ff3153c43903f929ae5d62fbba12610f89c36d55.json +++ b/backend/.sqlx/query-5a219a2532517869578c4504ff3153c43903f929ae5d62fbba12610f89c36d55.json @@ -15,7 +15,7 @@ ] }, "nullable": [ - null + true ] }, "hash": "5a219a2532517869578c4504ff3153c43903f929ae5d62fbba12610f89c36d55" diff --git a/backend/windmill-api/src/flows.rs b/backend/windmill-api/src/flows.rs index 0f4e8db038..9952ca1ddd 100644 --- a/backend/windmill-api/src/flows.rs +++ b/backend/windmill-api/src/flows.rs @@ -1634,6 +1634,7 @@ mod tests { early_return: None, concurrency_key: None, chat_input_enabled: None, + flow_env: None, debounce_key: None, debounce_delay_s: None, }; diff --git a/backend/windmill-api/src/jobs.rs b/backend/windmill-api/src/jobs.rs index da646e1311..d2de1a8f49 100644 --- a/backend/windmill-api/src/jobs.rs +++ b/backend/windmill-api/src/jobs.rs @@ -281,6 +281,10 @@ pub fn workspaced_service() -> Router { "/result_by_id/:job_id/:node_id", get(get_result_by_id).layer(cors.clone()), ) + .route( + "/flow_env_by_flow_job_id/:flow_job_id/:var_name", + get(get_flow_env_by_flow_job_id).layer(cors.clone()), + ) .route("/run/dependencies", post(run_dependencies_job)) .route("/run/flow_dependencies", post(run_flow_dependencies_job)) .route( @@ -376,6 +380,59 @@ async fn get_root_job( Ok(Json(res)) } +async fn get_flow_env_by_flow_job_id( + authed: ApiAuthed, + tokened: Tokened, + Extension(db): Extension, + Path((w_id, flow_job_id, var_name)): Path<(String, Uuid, String)>, + Query(JsonPath { json_path, .. }): Query, +) -> windmill_common::error::JsonResult> { + let flow_env = sqlx::query_scalar!( + r#" + SELECT + CASE + WHEN flow_version.id IS NOT NULL THEN + (flow_version.value -> 'flow_env' -> $3) #> $4 + ELSE + (root_job.raw_flow -> 'flow_env' -> $3) #> $4 + END AS "flow_env: sqlx::types::Json>" + FROM + v2_job current_job + JOIN + v2_job root_job ON root_job.id = COALESCE(current_job.root_job, current_job.flow_innermost_root_job, current_job.parent_job, current_job.id) + AND root_job.workspace_id = current_job.workspace_id + LEFT JOIN + flow_version ON flow_version.id = root_job.runnable_id + AND flow_version.path = root_job.runnable_path + AND flow_version.workspace_id = root_job.workspace_id + WHERE + current_job.id = $1 AND + current_job.workspace_id = $2"#, + flow_job_id, + w_id, + var_name, + json_path + .as_ref() + .map(|x| x.split(".").collect::>()) + .unwrap_or_default() as Vec<&str>, + ) + .fetch_optional(&db) + .await? + .map(|r| r.map(|x| x.0)) + .flatten() + .unwrap_or_else(|| to_raw_value(&serde_json::Value::Null)); + + log_job_view( + &db, + Some(&authed), + Some(&tokened.token), + &w_id, + &flow_job_id, + ) + .await?; + Ok(Json(flow_env)) +} + async fn compute_root_job_for_flow(db: &DB, w_id: &str, job_id: Uuid) -> error::Result { let root_job = sqlx::query_scalar!( r#"SELECT COALESCE(root_job, flow_innermost_root_job, parent_job, id) as "root_job!" FROM v2_job WHERE id = $1 AND workspace_id = $2"#, diff --git a/backend/windmill-common/src/client.rs b/backend/windmill-common/src/client.rs index b05fdfa2f7..146cc98e68 100644 --- a/backend/windmill-common/src/client.rs +++ b/backend/windmill-common/src/client.rs @@ -49,14 +49,7 @@ impl AuthedClient { "{}/api/w/{}/oidc/token/{}", self.base_internal_url, self.workspace, audience ); - let response = self.get(&url, vec![]).await?; - match response.status().as_u16() { - 200u16 => Ok(response - .json::() - .await - .context("decoding oidc token as json string")?), - _ => Err(anyhow::anyhow!(response.text().await.unwrap_or_default())), - } + make_basic_get_request(self, &url, None, Some("decoding oidc token as json string")).await } pub async fn get_resource_value(&self, path: &str) -> anyhow::Result { @@ -64,14 +57,7 @@ impl AuthedClient { "{}/api/w/{}/resources/get_value/{}", self.base_internal_url, self.workspace, path ); - let response = self.get(&url, vec![]).await?; - match response.status().as_u16() { - 200u16 => Ok(response - .json::() - .await - .context("decoding resource value as json")?), - _ => Err(anyhow::anyhow!(response.text().await.unwrap_or_default())), - } + make_basic_get_request(self, &url, None, Some("decoding resource value as json")).await } pub async fn get_variable_value(&self, path: &str) -> anyhow::Result { @@ -79,14 +65,7 @@ impl AuthedClient { "{}/api/w/{}/variables/get_value/{}", self.base_internal_url, self.workspace, path ); - let response = self.get(&url, vec![]).await?; - match response.status().as_u16() { - 200u16 => Ok(response - .json::() - .await - .context("decoding variable value as json")?), - _ => Err(anyhow::anyhow!(response.text().await.unwrap_or_default())), - } + make_basic_get_request(self, &url, None, Some("decoding variable value as json")).await } pub async fn get_resource_value_interpolated( @@ -121,19 +100,34 @@ impl AuthedClient { "{}/api/w/{}/jobs_u/completed/get_result/{}", self.base_internal_url, self.workspace, path ); - let query = if let Some(json_path) = json_path { - vec![("json_path", json_path)] - } else { - vec![] - }; - let response = self.get(&url, query).await?; - match response.status().as_u16() { - 200u16 => Ok(response - .json::() - .await - .context("decoding completed job result as json")?), - _ => Err(anyhow::anyhow!(response.text().await.unwrap_or_default())), - } + let query = query_from_json_path(json_path); + make_basic_get_request( + self, + &url, + Some(query), + Some("decoding completed job result as json"), + ) + .await + } + + pub async fn get_flow_env_by_flow_job_id( + &self, + root_job_id: &str, + var_name: &str, + json_path: Option, + ) -> anyhow::Result { + let url = format!( + "{}/api/w/{}/jobs/flow_env_by_flow_job_id/{}/{}", + self.base_internal_url, self.workspace, root_job_id, var_name + ); + let query = query_from_json_path(json_path); + make_basic_get_request( + self, + &url, + Some(query), + Some("decoding flow env variable as json"), + ) + .await } pub async fn get_result_by_id( @@ -146,19 +140,14 @@ impl AuthedClient { "{}/api/w/{}/jobs/result_by_id/{}/{}", self.base_internal_url, self.workspace, flow_job_id, node_id ); - let query = if let Some(json_path) = json_path { - vec![("json_path", json_path)] - } else { - vec![] - }; - let response = self.get(&url, query).await?; - match response.status().as_u16() { - 200u16 => Ok(response - .json::() - .await - .context("decoding result by id as json")?), - _ => Err(anyhow::anyhow!(response.text().await.unwrap_or_default())), - } + let query = query_from_json_path(json_path); + make_basic_get_request( + self, + &url, + Some(query), + Some("decoding result by id as json"), + ) + .await } pub async fn upload_s3_file( @@ -245,3 +234,33 @@ impl AuthedClient { } } } + +#[inline] +fn query_from_json_path(json_path: Option) -> Vec<(&'static str, String)> { + json_path + .map(|json_path| vec![("json_path", json_path)]) + .unwrap_or_else(|| Vec::new()) +} + +#[inline] +async fn make_basic_get_request( + client: &AuthedClient, + url: &str, + query: Option>, + context: Option<&'static str>, +) -> anyhow::Result { + let response = client + .get(&url, query.unwrap_or_else(|| Vec::new())) + .await?; + + match response.status().as_u16() { + 200u16 => { + let json_body = response + .json::() + .await + .context(context.unwrap_or("error decoding body as json"))?; + Ok(json_body) + } + _ => Err(anyhow::anyhow!(response.text().await.unwrap_or_default())), + } +} diff --git a/backend/windmill-common/src/flows.rs b/backend/windmill-common/src/flows.rs index ca1f08e8a2..929736eff5 100644 --- a/backend/windmill-common/src/flows.rs +++ b/backend/windmill-common/src/flows.rs @@ -192,6 +192,8 @@ pub struct FlowValue { pub priority: Option, #[serde(skip_serializing_if = "Option::is_none")] pub chat_input_enabled: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub flow_env: Option>> } impl FlowValue { diff --git a/backend/windmill-common/src/variables.rs b/backend/windmill-common/src/variables.rs index 9dd9edc64d..1edd292cbe 100644 --- a/backend/windmill-common/src/variables.rs +++ b/backend/windmill-common/src/variables.rs @@ -401,8 +401,7 @@ pub async fn get_reserved_variables( value, description: "Custom workspace environment variable".to_string(), is_custom: true, -}) -).collect() +})).collect() } async fn get_cached_workspace_envs(conn: &Connection, w_id: &str) -> Vec<(String, String)> { @@ -439,7 +438,11 @@ async fn get_cached_workspace_envs(conn: &Connection, w_id: &str) -> Vec<(String custom_envs } -pub async fn get_variable_or_self(path: String, db: &DB, w_id: &str) -> crate::error::Result { +pub async fn get_variable_or_self( + path: String, + db: &DB, + w_id: &str, +) -> crate::error::Result { if !path.starts_with("$var:") { return Ok(path); } diff --git a/backend/windmill-queue/src/jobs.rs b/backend/windmill-queue/src/jobs.rs index f55e0f08e9..a6370163d6 100644 --- a/backend/windmill-queue/src/jobs.rs +++ b/backend/windmill-queue/src/jobs.rs @@ -4329,6 +4329,7 @@ pub async fn push<'c, 'd>( skip_expr: None, preprocessor_module: None, chat_input_enabled: None, + flow_env: None, }; // this is a new flow being pushed, flow_status is set to flow_value: let flow_status: FlowStatus = FlowStatus::new(&flow_value); diff --git a/backend/windmill-worker/src/ai/tools.rs b/backend/windmill-worker/src/ai/tools.rs index 88469652bd..243a67fb4e 100644 --- a/backend/windmill-worker/src/ai/tools.rs +++ b/backend/windmill-worker/src/ai/tools.rs @@ -338,6 +338,7 @@ async fn execute_windmill_tool( transform, last_result.clone(), flow_inputs.clone(), + None, Some(ctx.client), ctx.id_context.as_ref(), ) diff --git a/backend/windmill-worker/src/common.rs b/backend/windmill-worker/src/common.rs index 892a3c8b7e..259b12ec84 100644 --- a/backend/windmill-worker/src/common.rs +++ b/backend/windmill-worker/src/common.rs @@ -1068,7 +1068,6 @@ pub fn build_http_client(timeout_duration: std::time::Duration) -> error::Result } pub fn get_root_job_id(job: &MiniPulledJob) -> uuid::Uuid { - // fallback to flow_innermost_root_job and parent_job as root_job is not set if equal to innermost root job or parent job job.root_job .or(job.flow_innermost_root_job) .or(job.parent_job) diff --git a/backend/windmill-worker/src/js_eval.rs b/backend/windmill-worker/src/js_eval.rs index c87b506686..ec28b6f509 100644 --- a/backend/windmill-worker/src/js_eval.rs +++ b/backend/windmill-worker/src/js_eval.rs @@ -165,10 +165,98 @@ impl NetPermissions for PermissionsContainer { #[cfg(feature = "deno_core")] pub struct OptAuthedClient(Option); +const FLOW_INPUT_PREFIX: &'static str = "flow_input"; +const ENV_KEY_PREFIX: &'static str = "flow_env"; +const DOT_PATTERN: &'static str = "."; +const START_BRACKET_PATTERN: &'static str = "[\""; +const END_BRACKET_PATTERN: &'static str = "\"]"; + +fn try_exact_property_access( + expr: &str, + flow_input: Option<&mappable_rc::Marc>>>, + flow_env: Option<&HashMap>>, +) -> Option> { + let obj = if expr.starts_with(FLOW_INPUT_PREFIX) { + Some(( + FLOW_INPUT_PREFIX, + flow_input.as_ref().map(|obj| obj.as_ref()), + )) + } else if expr.starts_with(ENV_KEY_PREFIX) { + Some((ENV_KEY_PREFIX, flow_env)) + } else { + None + }; + + if let Some((prefix, obj)) = obj { + let access_pattern_pos = prefix.len(); + let suffix = &expr[access_pattern_pos..]; + let maybe_key_name = if suffix.starts_with(DOT_PATTERN) { + let key_name_pos = DOT_PATTERN.len(); + Some(&expr[key_name_pos..]) + } else if suffix.starts_with(START_BRACKET_PATTERN) { + let key_name_pos = START_BRACKET_PATTERN.len(); + let suffix = &suffix[key_name_pos..]; + + let flow_arg_name = suffix + .ends_with(END_BRACKET_PATTERN) + .then(|| { + let start_key_name_pos = access_pattern_pos + key_name_pos; + let end_key_name_pos = expr.len() - END_BRACKET_PATTERN.len(); + &expr[start_key_name_pos..end_key_name_pos] + }) + .filter(|s| s.len() > 0); + flow_arg_name + } else { + None + }; + + if let Some(key_name) = maybe_key_name { + if let Some(key_value) = obj.and_then(|obj| obj.get(key_name)) { + return Some(key_value.clone()); + } + } + } + None +} + +async fn handle_full_regex( + captures: regex::Captures<'_>, + authed_client: &AuthedClient, + by_id: &IdContext, +) -> anyhow::Result> { + let obj_name = captures.get(1).unwrap().as_str(); + let obj_key = captures.get(2).unwrap().as_str(); + let idx_o = captures.get(3).map(|y| y.as_str()); + let rest = captures.get(4).map(|y| y.as_str()); + let query = if let Some(idx) = idx_o { + match rest { + Some(rest) => Some(format!("{}{}", idx, rest)), + None => Some(idx.to_string()), + } + } else { + rest.map(|x| x.trim_start_matches('.').to_string()) + }; + + let result = if obj_name == "results" { + authed_client + .get_result_by_id(&by_id.flow_job.to_string(), obj_key, query) + .await + } else if obj_name == "flow_env" { + authed_client + .get_flow_env_by_flow_job_id(&by_id.flow_job.to_string(), obj_key, query) + .await + } else { + unreachable!(); + }; + + return result; +} + pub async fn eval_timeout( expr: String, transform_context: HashMap>>, flow_input: Option>>>, + flow_env: Option<&HashMap>>, authed_client: Option<&AuthedClient>, by_id: Option<&IdContext>, #[allow(unused_variables)] ctx: Option>, @@ -180,21 +268,13 @@ pub async fn eval_timeout( expr, transform_context ); - for (k, v) in transform_context.iter() { - if k == &expr { - return Ok(v.as_ref().clone()); - } + + if let Some(value) = transform_context.get(&expr) { + return Ok(value.as_ref().to_owned()); } - if expr.starts_with("flow_input.") || expr.starts_with("flow_input[") { - if let Some(ref flow_input) = flow_input { - for (k, v) in flow_input.iter() { - if &format!("flow_input.{k}") == &expr || &format!("flow_input[\"{k}\"]") == &expr { - // tracing::error!("FLOW_INPUT"); - return Ok(v.clone()); - } - } - } + if let Some(value) = try_exact_property_access(&expr, flow_input.as_ref(), flow_env) { + return Ok(value); } let p_ids = by_id.map(|x| { @@ -219,24 +299,8 @@ pub async fn eval_timeout( } if let (Some(by_id), Some(authed_client)) = (by_id, authed_client) { - if let Some((id, idx_o, rest)) = RE_FULL.captures(&expr).map(|x| { - ( - x.get(1).unwrap().as_str(), - x.get(2).map(|y| y.as_str()), - x.get(3).map(|y| y.as_str()), - ) - }) { - let query = if let Some(idx) = idx_o { - match rest { - Some(rest) => Some(format!("{}{}", idx, rest)), - None => Some(idx.to_string()), - } - } else { - rest.map(|x| x.trim_start_matches('.').to_string()) - }; - return authed_client - .get_result_by_id(&by_id.flow_job.to_string(), id, query) - .await; + if let Some(captures) = RE_FULL.captures(&expr) { + return handle_full_regex(captures, authed_client, by_id).await; } } @@ -271,8 +335,8 @@ pub async fn eval_timeout( if by_id.is_some() && authed_client.is_some() { ops.push(op_get_result()); ops.push(op_get_id()); + ops.push(op_get_flow_env()); } - let ext = Extension { name: "js_eval", ops: ops.into(), ..Default::default() }; let exts = vec![ext]; // Use our snapshot to provision our new runtime @@ -307,13 +371,15 @@ pub async fn eval_timeout( let mut client = authed_client.clone(); if let Some(client) = client.as_mut() { client.force_client = Some( - configure_client(reqwest::ClientBuilder::new() - .user_agent("windmill/beta") - .danger_accept_invalid_certs( - std::env::var("ACCEPT_INVALID_CERTS").is_ok(), - )) - .build() - .unwrap(), + configure_client( + reqwest::ClientBuilder::new() + .user_agent("windmill/beta") + .danger_accept_invalid_certs( + std::env::var("ACCEPT_INVALID_CERTS").is_ok(), + ), + ) + .build() + .unwrap(), ); } op_state.put(OptAuthedClient(client)); @@ -323,7 +389,7 @@ pub async fn eval_timeout( .into_iter() .filter(|(a, _)| context_keys.contains(a)) .collect(), - }) + }); } sender @@ -379,10 +445,12 @@ fn replace_with_await(expr: String, fn_name: &str) -> String { } lazy_static! { static ref RE: Regex = - Regex::new(r#"(?m)(?Presults(?:\?)?(?:(?:\.[a-zA-Z_0-9]+)|(?:\[\".*?\"\])))"#).unwrap(); - static ref RE_FULL: Regex = - Regex::new(r"(?m)^results(?:\?)?\.([a-zA-Z_0-9]+)(?:\[(\d+)\])?((?:\.[a-zA-Z_0-9]+)+)?$") + Regex::new(r#"(?m)(?P(?:results|flow_env)(?:\?)?(?:(?:\.[a-zA-Z_0-9]+)|(?:\[\".*?\"\])))"#) .unwrap(); + static ref RE_FULL: Regex = Regex::new( + r"(?m)^(results|flow_env)(?:\?)?\.([a-zA-Z_0-9]+)(?:\[(\d+)\])?((?:\.[a-zA-Z_0-9]+)+)?$" + ) + .unwrap(); static ref RE_PROXY: Regex = Regex::new(r"^(https?)://(([^:@\s]+):([^:@\s]+)@)?([^:@\s]+)(:(\d+))?$").unwrap(); } @@ -453,6 +521,17 @@ const results = new Proxy({{}}, {{ }} }}); +async function flow_env_by_var_name(var_name) {{ + let root_job_id = "{}"; + return JSON.parse(await Deno.core.ops.op_get_flow_env(root_job_id, var_name, null)); +}} + +const flow_env = new Proxy({{}}, {{ + get: function(target, name, receiver) {{ + return flow_env_by_var_name(name); + }} +}}); + "#, by_id .steps_results @@ -469,6 +548,7 @@ const results = new Proxy({{}}, {{ .join(","), by_id.previous_id, by_id.flow_job, + by_id.flow_job ) } else { String::new() @@ -641,6 +721,33 @@ async fn op_resource( } } +#[cfg(feature = "deno_core")] +#[op2(async)] +#[string] +async fn op_get_flow_env( + op_state: Rc>, + #[string] root_job_id: String, + #[string] var_name: String, + #[string] json_path: Option, +) -> Result, deno_error::JsErrorBox> { + let client = op_state.borrow().borrow::().0.clone(); + if let Some(client) = client { + client + .get_flow_env_by_flow_job_id::>>( + &root_job_id, + &var_name, + json_path, + ) + .await + .map(|value| value.map(|val| val.get().to_string())) + .map_err(|e| deno_error::JsErrorBox::generic(e.to_string())) + } else { + Err(deno_error::JsErrorBox::generic( + "No client found in op state", + )) + } +} + #[cfg(feature = "deno_core")] pub struct TransformContext { pub envs: HashMap>>, @@ -1318,7 +1425,7 @@ multiline template`"; op_state.put(TransformContext { flow_input: None, envs: env.clone() }) } - let res = eval_timeout(code.to_string(), env, None, None, None, None).await?; + let res = eval_timeout(code.to_string(), env, None, None, None, None, None).await?; assert_eq!(res.get(), "2"); Ok(()) } diff --git a/backend/windmill-worker/src/worker.rs b/backend/windmill-worker/src/worker.rs index db568652e5..c457dc5343 100644 --- a/backend/windmill-worker/src/worker.rs +++ b/backend/windmill-worker/src/worker.rs @@ -877,7 +877,7 @@ pub fn start_interactive_worker_shell( token, precomputed_agent_info: precomputed_bundle, } = extract_job_and_perms(job, &conn).await; - + let authed_client = AuthedClient::new( base_internal_url.to_owned(), job.workspace_id.clone(), @@ -886,7 +886,7 @@ pub fn start_interactive_worker_shell( ); let arc_job = Arc::new(job); - + let _ = handle_queued_job( arc_job.clone(), raw_code, diff --git a/backend/windmill-worker/src/worker_flow.rs b/backend/windmill-worker/src/worker_flow.rs index f25e04d564..f9f8709f5c 100644 --- a/backend/windmill-worker/src/worker_flow.rs +++ b/backend/windmill-worker/src/worker_flow.rs @@ -245,6 +245,7 @@ async fn evaluate_stop_after_all_iters_if( let stop_early_after_all_iters = compute_bool_from_expr( &stop_after_all_iters_if.expr, Marc::new(args), + None, iters_result.clone(), None, None, @@ -427,17 +428,20 @@ pub async fn update_flow_status_after_job_completion_internal( ) .fetch_one(db) .await; + let args = + args.map(|flow_args| flow_args.map(|flow_args| flow_args.0).unwrap_or_default()); + args })); - let from_result_to_args = - |args: &Result>>>, sqlx::Error>| { - let args = args.as_ref().map_err(|e| { - Error::internal_err(format!("retrieval of args from state: {e:#}")) - })?; - Ok::<_, Error>(args.clone().unwrap_or_default().0) - }; + let from_result_to_args = |args: &Result>, sqlx::Error>| { + let args = args + .as_ref() + .map_err(|e| Error::internal_err(format!("retrieval of args from state: {e:#}")))?; + + Ok::<_, Error>(args.clone()) + }; let (mut stop_early, mut stop_early_err_msg, mut skip_if_stop_early, continue_on_error) = if stop_early_override.is_some() @@ -469,9 +473,11 @@ pub async fn update_flow_status_after_job_completion_internal( _ => None, }; let args = from_result_to_args(args.as_ref().await.get_ref())?; + compute_bool_from_expr( &expr, Marc::new(args), + None, result.clone(), all_iters, None, @@ -736,6 +742,7 @@ pub async fn update_flow_status_after_job_completion_internal( &mut stop_early_err_msg, &mut nresult, args, + ) .await?; } @@ -924,6 +931,7 @@ pub async fn update_flow_status_after_job_completion_internal( .and_then(|x| x.stop_after_all_iters_if.as_ref()) { let args = from_result_to_args(args.as_ref().await.get_ref())?; + evaluate_stop_after_all_iters_if( db, stop_after_all_iters_if, @@ -1016,6 +1024,7 @@ pub async fn update_flow_status_after_job_completion_internal( &old_status.retry, result.clone(), Marc::new(args), + None, Some(client), ) .await? @@ -1334,6 +1343,7 @@ pub async fn update_flow_status_after_job_completion_internal( &old_status.retry, result.clone(), Marc::new(args), + None, Some(client), ) .await? @@ -1824,6 +1834,7 @@ async fn evaluate_retry( status: &RetryStatus, result: Arc>, flow_args: Marc>>, + flow_env: Option<&HashMap>>, client: Option<&AuthedClient>, ) -> anyhow::Result> { if status.fail_count > MAX_RETRY_ATTEMPTS { @@ -1834,6 +1845,7 @@ async fn evaluate_retry( let should_retry = compute_bool_from_expr( &retry_if.expr, flow_args, + flow_env, result, None, None, @@ -1857,6 +1869,7 @@ async fn evaluate_retry( async fn compute_bool_from_expr( expr: &str, flow_args: Marc>>, + flow_env: Option<&HashMap>>, result: Arc>, all_iters: Option>>, by_id: Option<&IdContext>, @@ -1881,6 +1894,7 @@ async fn compute_bool_from_expr( format!("Boolean({expr})"), context, Some(flow_args), + flow_env, client, by_id, ctx, @@ -1906,6 +1920,7 @@ pub async fn evaluate_input_transform( transform: &InputTransform, last_result: Arc>, flow_args: Option>>>, + flow_env: Option<&HashMap>>, authed_client: Option<&AuthedClient>, by_id: Option<&IdContext>, ) -> error::Result @@ -1927,6 +1942,7 @@ where expr.to_string(), context, flow_args, + flow_env, authed_client, by_id, None, @@ -1955,6 +1971,7 @@ where #[instrument(level = "trace", skip_all)] async fn transform_input( flow_args: Marc>>, + flow_env: Option<&HashMap>>, last_result: Arc>, input_transforms: &HashMap, resumes: Arc>, @@ -1996,6 +2013,7 @@ async fn transform_input( expr.to_string(), env.clone(), Some(flow_args.clone()), + flow_env, Some(client), Some(by_id), None, @@ -2069,6 +2087,7 @@ pub async fn handle_flow( ); } } + let mut rec = PushNextFlowJobRec { flow_job: flow_job, status: status }; loop { let PushNextFlowJobRec { flow_job, status } = rec; @@ -2204,13 +2223,14 @@ async fn push_next_flow_job( // tracing::error!("status_module: {status_module:#?}"); let fj: mappable_rc::Marc = flow_job.clone().into(); - let arc_flow_job_args: Marc>> = Marc::map(fj, |x| { - if let Some(args) = &x.args { - &args.0 - } else { - &EHM - } - }); + let arc_flow_job_args: Marc>> = + Marc::map(fj, |x: &MiniPulledJob| { + if let Some(args) = &x.args { + &args.0 + } else { + &EHM + } + }); // if this is an empty module without preprocessor of if the module has already been completed, successfully, update the parent flow if (flow.modules.is_empty() && !step.is_preprocessor_step()) @@ -2300,6 +2320,7 @@ async fn push_next_flow_job( let skip = compute_bool_from_expr( &skip_expr, arc_flow_job_args.clone(), + flow.flow_env.as_ref(), Arc::new(to_raw_value(&json!("{}"))), None, None, @@ -2421,6 +2442,7 @@ async fn push_next_flow_job( expr.to_string(), context, Some(arc_flow_job_args.clone()), + flow.flow_env.as_ref(), None, None, None @@ -2681,6 +2703,7 @@ async fn push_next_flow_job( &input_transform, arc_last_job_result.clone(), Some(arc_flow_job_args.clone()), + flow.flow_env.as_ref(), Some(client), None, ) @@ -2718,6 +2741,7 @@ async fn push_next_flow_job( &status.retry, arc_last_job_result.clone(), arc_flow_job_args.clone(), + flow.flow_env.as_ref(), Some(client), ) .await? @@ -2807,6 +2831,7 @@ async fn push_next_flow_job( compute_bool_from_expr( &skip_if.expr, arc_flow_job_args.clone(), + flow.flow_env.as_ref(), arc_last_job_result.clone(), None, Some(&idcontext), @@ -2898,6 +2923,7 @@ async fn push_next_flow_job( }; transform_input( arc_flow_job_args.clone(), + flow.flow_env.as_ref(), arc_last_job_result.clone(), input_transforms, resumes.clone(), @@ -2924,6 +2950,7 @@ async fn push_next_flow_job( let next_flow_transform = compute_next_flow_transform( arc_flow_job_args.clone(), arc_last_job_result.clone(), + flow.flow_env.as_ref(), &flow_job, &flow, transform_context, @@ -3072,6 +3099,7 @@ async fn push_next_flow_job( .await?; let ti = transform_input( Marc::new(args), + flow.flow_env.as_ref(), arc_last_job_result.clone(), input_transforms, resumes.clone(), @@ -3122,6 +3150,7 @@ async fn push_next_flow_job( .await?; let ti = transform_input( Marc::new(hm), + flow.flow_env.as_ref(), arc_last_job_result.clone(), input_transforms, resumes.clone(), @@ -3225,6 +3254,7 @@ async fn push_next_flow_job( timeout_transform, arc_last_job_result.clone(), Some(arc_flow_job_args.clone()), + flow.flow_env.as_ref(), Some(client), Some(&ctx), ) @@ -3304,6 +3334,7 @@ async fn push_next_flow_job( parallelism_transform, arc_last_job_result.clone(), Some(arc_flow_job_args.clone()), + flow.flow_env.as_ref(), Some(client), Some(&ctx), ) @@ -3761,6 +3792,7 @@ pub fn get_path(flow_job: &MiniPulledJob, status: &FlowStatus, module: &FlowModu async fn compute_next_flow_transform( arc_flow_job_args: Marc>>, arc_last_job_result: Arc>, + flow_env: Option<&HashMap>>, flow_job: &MiniPulledJob, flow: &FlowValue, by_id: Option, @@ -3979,6 +4011,7 @@ async fn compute_next_flow_transform( resume, approvers, arc_flow_job_args, + flow_env, client, ¶llel, ) @@ -4074,6 +4107,7 @@ async fn compute_next_flow_transform( let pred = compute_bool_from_expr( &b.expr, arc_flow_job_args.clone(), + flow.flow_env.as_ref(), arc_last_job_result.clone(), None, Some(&idcontext), @@ -4338,6 +4372,7 @@ async fn next_forloop_status( resume: Arc>, approvers: Arc>, arc_flow_job_args: Marc>>, + flow_env: Option<&HashMap>>, client: &AuthedClient, parallel: &bool, ) -> Result { @@ -4370,6 +4405,7 @@ async fn next_forloop_status( expr.to_string(), context, Some(arc_flow_job_args), + flow_env, Some(client), Some(&by_id), None, @@ -4433,6 +4469,7 @@ async fn next_forloop_status( expr.to_string(), context, Some(arc_flow_job_args), + flow_env, Some(client), Some(&by_id), None, diff --git a/frontend/src/lib/components/InputTransformForm.svelte b/frontend/src/lib/components/InputTransformForm.svelte index 15d18f9e4b..7d3b2f464f 100644 --- a/frontend/src/lib/components/InputTransformForm.svelte +++ b/frontend/src/lib/components/InputTransformForm.svelte @@ -3,7 +3,8 @@ 'flow_input', 'results', 'resource', - 'variable' + 'variable', + 'flow_env' ]) diff --git a/frontend/src/lib/components/flows/content/FlowConstants.svelte b/frontend/src/lib/components/flows/content/FlowConstants.svelte deleted file mode 100644 index d9890281e3..0000000000 --- a/frontend/src/lib/components/flows/content/FlowConstants.svelte +++ /dev/null @@ -1,157 +0,0 @@ - - -
- - {#snippet header()} - - {/snippet} -
- This page centralizes the static inputs of every steps. It is aking to a file containing - all constants. Modifying a value here modifies it in the step input directly. It is - especially useful when forking a flow to get an overview of all the variables to parametrize - that are not exposed directly as flow inputs. - {#if Object.keys(resources).length > 0} - - The following resources are missing and the flow will not be fully runnable until they are - set. Add your own resources: - {#each Object.entries(resources) as [id, r]} - {#each r as resource} -
- {id} is missing a resource of type{' '} - {resource?.type} for the input{' '} - {resource?.argName} -
- {/each} - {/each} -
- {/if} - {#if steps.length == 0} -
- {#if flowStore.val.value.modules.length == 0} - - This flow has no steps. Add a step to see its static inputs. - - {:else} - - This flow has no steps with static inputs. Add a step with static inputs to see them - here. - - {/if} - {/if} - {#each steps as [_args, filter, m], index (m.id + index)} - {#if filter.length > 0} -
-

- {m.summary || m.value['path'] || 'Inline script'} - {m.id} -

- - -
- {/if} - {/each} -
-
-
diff --git a/frontend/src/lib/components/flows/content/FlowEditorPanel.svelte b/frontend/src/lib/components/flows/content/FlowEditorPanel.svelte index 0b9621ef39..65145fa0ac 100644 --- a/frontend/src/lib/components/flows/content/FlowEditorPanel.svelte +++ b/frontend/src/lib/components/flows/content/FlowEditorPanel.svelte @@ -5,7 +5,7 @@ import FlowSettings from './FlowSettings.svelte' import FlowInput from './FlowInput.svelte' import FlowFailureModule from './FlowFailureModule.svelte' - import FlowConstants from './FlowConstants.svelte' + import FlowEnvironmentVariables from './FlowEnvironmentVariables.svelte' import type { FlowModule, Flow, Job } from '$lib/gen' import FlowPreprocessorModule from './FlowPreprocessorModule.svelte' import type { TriggerContext } from '$lib/components/triggers' @@ -102,7 +102,7 @@ {:else if $selectedId === 'Result'} {:else if $selectedId === 'constants'} - + {:else if $selectedId === 'failure'} {:else if $selectedId === 'preprocessor'} diff --git a/frontend/src/lib/components/flows/content/FlowEnvironmentVariables.svelte b/frontend/src/lib/components/flows/content/FlowEnvironmentVariables.svelte new file mode 100644 index 0000000000..3391dcb311 --- /dev/null +++ b/frontend/src/lib/components/flows/content/FlowEnvironmentVariables.svelte @@ -0,0 +1,286 @@ + + +
+ +
+ + Flow envs can be referenced in any flow step input using the syntax{' '} + flow_env.VARIABLE_NAME or flow_env["VARIABLE_NAME"]. These + variables are available in the property picker and can be used in JavaScript expressions and + input bindings. You can choose between String or JSON types for each variable - JSON types + allow complex data structures. + + + {#if flowEnvEntries.length === 0} + + This flow has no flow env variables defined. Click "Add Variable" to create your first + flow env variable. + + {:else} +
+ {#each flowEnvEntries as entry (entry.id)} +
+
+
+ +
+ +
+
+ {/each} +
+ {/if} +
+ {#if !noEditor} + + {/if} +
+
+
+
diff --git a/frontend/src/lib/components/flows/map/FlowStickyNode.svelte b/frontend/src/lib/components/flows/map/FlowStickyNode.svelte index 65c68e443b..48e7445593 100644 --- a/frontend/src/lib/components/flows/map/FlowStickyNode.svelte +++ b/frontend/src/lib/components/flows/map/FlowStickyNode.svelte @@ -66,7 +66,7 @@ onClick={() => ($selectedId = 'constants')} /> {#snippet text()} - Static inputs + Environment Variables {/snippet} {/if} diff --git a/frontend/src/lib/components/flows/previousResults.ts b/frontend/src/lib/components/flows/previousResults.ts index 59fd466b80..1127b0c656 100644 --- a/frontend/src/lib/components/flows/previousResults.ts +++ b/frontend/src/lib/components/flows/previousResults.ts @@ -9,6 +9,7 @@ export type PickableProperties = { priorIds: Record previousId: string | undefined hasResume: boolean + flow_env?: Record } type StepPropPicker = { @@ -156,7 +157,8 @@ export function getFailureStepPropPicker(flowState: FlowState, flow: OpenFlow, a flow_input: schemaToObject(flow.schema as any, args), priorIds: priorIds, previousId: undefined, - hasResume: false + hasResume: false, + flow_env: flow.value.flow_env }, extraLib: ` /** @@ -178,6 +180,17 @@ declare const results = ${JSON.stringify(priorIds)} * flow input as an object */ declare const flow_input = ${JSON.stringify(flowInput)}; + +${ + flow.value.flow_env + ? ` +/** +* flow environment variables +*/ +declare const flow_env = ${JSON.stringify(flow.value.flow_env)}; +` + : '' +} ` } } @@ -218,7 +231,8 @@ export function getStepPropPicker( flow_input: flowInput, priorIds: priorIds, previousId: previousIds[0], - hasResume: previousModule?.suspend != undefined + hasResume: previousModule?.suspend != undefined, + flow_env: flow.value.flow_env } if (pickableProperties.hasResume) { @@ -230,7 +244,8 @@ export function getStepPropPicker( flowInput, priorIds, previousModule?.suspend != undefined, - previousModule?.id + previousModule?.id, + flow.value.flow_env ), pickableProperties } @@ -240,7 +255,8 @@ export function buildExtraLib( flowInput: Record, results: Record, resume: boolean, - previousId: string | undefined + previousId: string | undefined, + flowEnv?: Record ): string { return ` /** @@ -275,6 +291,17 @@ declare const results = ${JSON.stringify(results)}; */ declare const previous_result: ${previousId ? JSON.stringify(results[previousId]) : 'any'}; +${ + flowEnv + ? ` +/** + * flow environment variables + */ +declare const flow_env = ${JSON.stringify(flowEnv)}; +` + : '' +} + ${ resume ? ` diff --git a/frontend/src/lib/components/flows/propPicker/PropPickerWrapper.svelte b/frontend/src/lib/components/flows/propPicker/PropPickerWrapper.svelte index b5346566a2..5d4761d8d5 100644 --- a/frontend/src/lib/components/flows/propPicker/PropPickerWrapper.svelte +++ b/frontend/src/lib/components/flows/propPicker/PropPickerWrapper.svelte @@ -31,6 +31,7 @@ import type { PickableProperties } from '../previousResults' import AnimatedButton from '$lib/components/common/button/AnimatedButton.svelte' import type { PropPickerContext } from '$lib/components/prop_picker' + import type { FlowEditorContext } from '../types' interface Props { pickableProperties: PickableProperties | undefined @@ -69,6 +70,10 @@ const { flowPropPickerConfig } = getContext('PropPickerContext') flowPropPickerConfig.set(undefined) + + const { flowStore } = getContext('FlowEditorContext') + + let flow_env = $derived(pickableProperties?.flow_env || flowStore.val.value.flow_env) setContext('PropPickerWrapper', { propPickerConfig, inputMatches, @@ -156,6 +161,7 @@ | undefined = undefined let variables: Record = {} let resources: Record = {} let displayVariable = false let displayResources = false + let displayFlowEnv = false let allResultsCollapsed = true let collapsableInitialState: @@ -30,6 +32,7 @@ allResultsCollapsed: boolean displayVariable: boolean displayResources: boolean + displayFlowEnv: boolean } | undefined @@ -46,6 +49,7 @@ let flowInputsFiltered: any = pickableProperties.flow_input let resultByIdFiltered: any = pickableProperties.priorIds + let flowEnvFiltered: any = pickableProperties.flow_env let timeout: number | undefined function onSearch(search: string) { @@ -63,6 +67,9 @@ search === EMPTY_STRING ? pickableProperties.priorIds : keepByKey(pickableProperties.priorIds, search) + + flowEnvFiltered = + search === EMPTY_STRING ? pickableProperties.flow_env : keepByKey(pickableProperties.flow_env, search) }, 50) } @@ -98,6 +105,7 @@ if (search === EMPTY_STRING) { flowInputsFiltered = pickableProperties.flow_input resultByIdFiltered = pickableProperties.priorIds + flowEnvFiltered = pickableProperties.flow_env } filteringFlowInputsOrResult = '' return @@ -109,6 +117,9 @@ if (!$inputMatches?.some((match) => match.word === 'results')) { resultByIdFiltered = {} } + if (!$inputMatches?.some((match) => match.word === 'flow_env')) { + flowEnvFiltered = {} + } if ($inputMatches?.length == 1) { filteringFlowInputsOrResult = $inputMatches[0].value if ($inputMatches[0].word === 'flow_input') { @@ -125,6 +136,13 @@ if (Object.keys(filtered).length > 0) { resultByIdFiltered = filtered } + } else if ($inputMatches[0].word === 'flow_env') { + flowEnvFiltered = pickableProperties.flow_env + let [, ...nestedKeys] = $inputMatches[0].value.split('.') + let filtered = filterNestedObject(flowEnvFiltered, nestedKeys) + if (Object.keys(filtered).length > 0) { + flowEnvFiltered = filtered + } } } else { filteringFlowInputsOrResult = '' @@ -143,7 +161,12 @@ } if (!collapsableInitialState) { - collapsableInitialState = { allResultsCollapsed, displayVariable, displayResources } + collapsableInitialState = { + allResultsCollapsed, + displayVariable, + displayResources, + displayFlowEnv + } } if ($inputMatches[0].word === 'variable') { @@ -156,6 +179,10 @@ displayResources = true return } + if ($inputMatches[0].word === 'flow_env') { + displayFlowEnv = true + return + } if ($inputMatches[0].word === 'results') { allResultsCollapsed = false return @@ -166,7 +193,8 @@ if (!collapsableInitialState) { return } - ;({ allResultsCollapsed, displayVariable, displayResources } = collapsableInitialState) + ;({ allResultsCollapsed, displayVariable, displayResources, displayFlowEnv } = + collapsableInitialState) collapsableInitialState = undefined } @@ -183,6 +211,7 @@ if (prev && !filterActive) { flowInputsFiltered = pickableProperties.flow_input resultByIdFiltered = pickableProperties.priorIds + flowEnvFiltered = pickableProperties.flow_env } } @@ -192,7 +221,7 @@ await updateCollapsable() } - $: (search, $inputMatches, $propPickerConfig, pickableProperties, updateState()) + $: search, $inputMatches, $propPickerConfig, pickableProperties, updateState() onDestroy(() => { clearTimeout(timeout) @@ -400,6 +429,45 @@ {/if} {/if} + {#if flow_env && Object.keys(flow_env).length > 0 && (!filterActive || $inputMatches?.some((match) => match.word === 'flow_env'))} +
+ Flow Env Variables: + + {#if displayFlowEnv} + + + {:else} + + {/if} +
+ {/if} {/if} diff --git a/frontend/src/lib/components/propertyPicker/PropPickerResult.svelte b/frontend/src/lib/components/propertyPicker/PropPickerResult.svelte index cbb2f9c33c..7d56abeb05 100644 --- a/frontend/src/lib/components/propertyPicker/PropPickerResult.svelte +++ b/frontend/src/lib/components/propertyPicker/PropPickerResult.svelte @@ -5,6 +5,7 @@ export let result: any export let extraResults: any = undefined export let flow_input: any = undefined + export let flow_env: any = undefined
@@ -18,4 +19,10 @@
{/if} + {#if flow_env} + Flow Environment Variables +
+ +
+ {/if} diff --git a/openflow.openapi.yaml b/openflow.openapi.yaml index f332964c40..a9a4009211 100644 --- a/openflow.openapi.yaml +++ b/openflow.openapi.yaml @@ -62,6 +62,10 @@ components: type: string cache_ttl: type: number + flow_env: + type: object + additionalProperties: + type: string priority: type: number early_return: