From f48016048b3fcf5e1491b3967d2eadbf65d2d2eb Mon Sep 17 00:00:00 2001 From: pyranota <92104930+pyranota@users.noreply.github.com> Date: Thu, 11 Sep 2025 12:46:01 +0200 Subject: [PATCH] fix(apps): fix relative imports cache invalidation (#6564) * v0 Signed-off-by: pyranota * optimize relocks * make it work with relative relative imports Signed-off-by: pyranota * use fallback Signed-off-by: pyranota * remove dbg and todos Signed-off-by: pyranota * future proof a bit Signed-off-by: pyranota * cleanup Signed-off-by: pyranota * more cleanup Signed-off-by: pyranota * remove final TODO Signed-off-by: pyranota * do not use bytemuck Signed-off-by: pyranota * optimize hashing Signed-off-by: pyranota --------- Signed-off-by: pyranota --- ...153c43903f929ae5d62fbba12610f89c36d55.json | 2 +- ...420ef920d472286cf4c1f78452446a4bcb604.json | 22 ++ ...e359770b759e8661687b4503c4b692ecd245e.json | 23 ++ .../windmill-worker/src/worker_lockfiles.rs | 272 ++++++++++++++++-- 4 files changed, 287 insertions(+), 32 deletions(-) create mode 100644 backend/.sqlx/query-83232f2db5eb1b6fef744998e60420ef920d472286cf4c1f78452446a4bcb604.json create mode 100644 backend/.sqlx/query-ca15fe5d43f0e94f50408efe5c9e359770b759e8661687b4503c4b692ecd245e.json 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/.sqlx/query-83232f2db5eb1b6fef744998e60420ef920d472286cf4c1f78452446a4bcb604.json b/backend/.sqlx/query-83232f2db5eb1b6fef744998e60420ef920d472286cf4c1f78452446a4bcb604.json new file mode 100644 index 0000000000..27a5df6de9 --- /dev/null +++ b/backend/.sqlx/query-83232f2db5eb1b6fef744998e60420ef920d472286cf4c1f78452446a4bcb604.json @@ -0,0 +1,22 @@ +{ + "db_name": "PostgreSQL", + "query": "INSERT INTO app_version\n (app_id, value, created_by, raw_app)\n SELECT app_id, value, created_by, raw_app\n FROM app_version WHERE id = $1\n RETURNING id", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "id", + "type_info": "Int8" + } + ], + "parameters": { + "Left": [ + "Int8" + ] + }, + "nullable": [ + false + ] + }, + "hash": "83232f2db5eb1b6fef744998e60420ef920d472286cf4c1f78452446a4bcb604" +} diff --git a/backend/.sqlx/query-ca15fe5d43f0e94f50408efe5c9e359770b759e8661687b4503c4b692ecd245e.json b/backend/.sqlx/query-ca15fe5d43f0e94f50408efe5c9e359770b759e8661687b4503c4b692ecd245e.json new file mode 100644 index 0000000000..6812bfdfc7 --- /dev/null +++ b/backend/.sqlx/query-ca15fe5d43f0e94f50408efe5c9e359770b759e8661687b4503c4b692ecd245e.json @@ -0,0 +1,23 @@ +{ + "db_name": "PostgreSQL", + "query": "SELECT versions[array_upper(versions, 1)] FROM app WHERE path = $1 AND workspace_id = $2", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "versions", + "type_info": "Int8" + } + ], + "parameters": { + "Left": [ + "Text", + "Text" + ] + }, + "nullable": [ + null + ] + }, + "hash": "ca15fe5d43f0e94f50408efe5c9e359770b759e8661687b4503c4b692ecd245e" +} diff --git a/backend/windmill-worker/src/worker_lockfiles.rs b/backend/windmill-worker/src/worker_lockfiles.rs index 3ef64e5378..be0cc38cc3 100644 --- a/backend/windmill-worker/src/worker_lockfiles.rs +++ b/backend/windmill-worker/src/worker_lockfiles.rs @@ -6,8 +6,9 @@ use std::path::{Component, Path, PathBuf}; #[cfg(feature = "python")] use crate::ansible_executor::{get_git_repos_lock, AnsibleDependencyLocks}; use async_recursion::async_recursion; +use itertools::Itertools; use serde_json::value::RawValue; -use serde_json::{json, Value}; +use serde_json::{from_value, json, Value}; use sha2::Digest; use sqlx::types::Json; use uuid::Uuid; @@ -38,10 +39,12 @@ use windmill_parser_py_imports::parse_relative_imports; use windmill_parser_ts::parse_expr_for_imports; use windmill_queue::{append_logs, CanceledBy, MiniPulledJob, PushIsolationLevel}; +// TODO: To be removed in future versions lazy_static::lazy_static! { - // TODO: To be removed in future versions static ref WMDEBUG_NO_HASH_CHANGE_ON_DJ: bool = std::env::var("WMDEBUG_NO_HASH_CHANGE_ON_DJ").is_ok(); static ref WMDEBUG_NO_NEW_FLOW_VERSION_ON_DJ: bool = std::env::var("WMDEBUG_NO_NEW_FLOW_VERSION_ON_DJ").is_ok(); + static ref WMDEBUG_NO_NEW_APP_VERSION_ON_DJ: bool = std::env::var("WMDEBUG_NO_NEW_APP_VERSION_ON_DJ").is_ok(); + static ref WMDEBUG_NO_COMPONENTS_TO_RELOCK: bool = std::env::var("WMDEBUG_NO_COMPONENTS_TO_RELOCK").is_ok(); } use crate::common::OccupancyMetrics; @@ -668,6 +671,8 @@ pub async fn trigger_dependents_to_recompute_dependencies( let kind = s.importer_kind.clone().unwrap_or_default(); let job_payload = if kind == "script" { let r = + // TODO: Not sure if this is safe: + // might have race conditions in edge-cases get_latest_deployed_hash_for_path(None, db.clone(), w_id, s.importer_path.as_str()) .await; match r { @@ -780,6 +785,78 @@ pub async fn trigger_dependents_to_recompute_dependencies( continue; } } + } else if kind == "app" && !*WMDEBUG_NO_NEW_APP_VERSION_ON_DJ { + // Create transaction to make operation atomic. + let mut tx = db.begin().await?; + + args.insert( + "components_to_relock".to_string(), + to_raw_value(&s.importer_node_ids), + ); + + let r = sqlx::query_scalar!( + "SELECT versions[array_upper(versions, 1)] FROM app WHERE path = $1 AND workspace_id = $2", + s.importer_path, + w_id, + ).fetch_one(&mut *tx) + .await + .map_err(to_anyhow); + + match r { + // Get current version of current flow. + Ok(Some(cur_version)) => { + // NOTE: Temporary solution. See the usage for more details. + args.insert( + "triggered_by_relative_import".to_string(), + to_raw_value(&()), + ); + + let new_version = sqlx::query_scalar!( + "INSERT INTO app_version + (app_id, value, created_by, raw_app) + SELECT app_id, value, created_by, raw_app + FROM app_version WHERE id = $1 + RETURNING id", + cur_version + ) + .fetch_one(&mut *tx) + .await + .map_err(|e| { + error::Error::internal_err(format!( + "Error updating App due to App history insert: {e:#}" + )) + })?; + + // Commit the transaction. + // NOTE: + // We do not append app.versions with new version. + // We will do this in the end of the dependency job handler. + // Otherwise it might become a source of race-conditions. + tx.commit().await?; + JobPayload::AppDependencies { + path: s.importer_path.clone(), + // Point Dep Job to the new version. + // We do this since we want to assume old ones are immutable. + version: new_version, + } + } + Ok(None) => { + tracing::error!( + "no app version found for path {path}", + path = s.importer_path + ); + // Do not commit the transaction. It will be dropped and rollbacked + continue; + } + Err(err) => { + tracing::error!( + "error getting latest deployed app version for path {path}: {err}", + path = s.importer_path, + ); + // Do not commit the transaction. It will be dropped and rollbacked + continue; + } + } } else { tracing::error!( "unexpected importer kind: {kind} for path {path}", @@ -1489,6 +1566,34 @@ async fn lock_modules<'c>( Ok((new_flow_modules, tx, modified_ids, errors)) } +/// Parse relative imports in script and call db to get each scripts' hash. +async fn relative_imports_bytes<'a>( + e: impl sqlx::Executor<'a, Database = sqlx::Postgres>, + code: Option<&String>, + path: &str, + language: Option, +) -> Result> { + Ok( + if let Some(imports) = extract_relative_imports( + code.map(|s| s.as_str()).unwrap_or_default(), + path, + &language, + ) { + sqlx::query_scalar::<_, i64>( + "SELECT hash FROM script WHERE path = ANY($1) AND archived = false", + ) + .bind(imports) + .fetch_all(e) + .await? + .iter() + .flat_map(|h| h.to_le_bytes()) + .collect_vec() + } else { + vec![] + }, + ) +} + async fn insert_flow_node<'c>( mut tx: sqlx::Transaction<'c, sqlx::Postgres>, path: &str, @@ -1503,25 +1608,11 @@ async fn insert_flow_node<'c>( hasher.update(code.unwrap_or(&Default::default())); hasher.update(lock.unwrap_or(&Default::default())); hasher.update(flow.unwrap_or(&Default::default()).get()); - if !*WMDEBUG_NO_NEW_FLOW_VERSION_ON_DJ { - if let Some(imports) = extract_relative_imports( - code.map(|s| s.as_str()).unwrap_or_default(), - path, - &language, - ) { - // We also want to take into account hashes of relative imports. - // TODO: May be use bytemuck or cast it in different, more proper way. - hasher.update(&format!( - "{:?}", - sqlx::query_scalar::<_, i64>( - "SELECT hash FROM script WHERE path = ANY($1) AND archived = false" - ) - .bind(imports) - .fetch_all(&mut *tx) - .await? - )); - } + // We also want to take into account hashes of relative imports. + hasher.update( + relative_imports_bytes(&mut *tx, code, &format!("{path}/flow"), language).await?, + ); } format!("{:x}", hasher.finalize()) }; @@ -1546,11 +1637,14 @@ async fn insert_flow_node<'c>( Ok((tx, FlowNodeId(id))) } +// TODO: Clean up dependency map when moved/renamed? async fn insert_app_script( db: &sqlx::Pool, + path: &str, app: i64, code: String, lock: Option, + language: Option, ) -> Result { let code_sha256 = format!("{:x}", sha2::Sha256::digest(&code)); let hash = { @@ -1558,6 +1652,12 @@ async fn insert_app_script( hasher.update(app.to_le_bytes()); hasher.update(&code_sha256); hasher.update(lock.as_ref().unwrap_or(&Default::default())); + // We also want to take into account hashes of relative imports. + if !*WMDEBUG_NO_NEW_APP_VERSION_ON_DJ { + hasher.update( + relative_imports_bytes(db, Some(&code), &format!("{path}/app"), language).await?, + ); + } format!("{:x}", hasher.finalize()) }; @@ -1741,15 +1841,17 @@ async fn reduce_flow<'c>( Ok(tx) } -async fn reduce_app(db: &sqlx::Pool, value: &mut Value, app: i64) -> Result<()> { +async fn reduce_app( + db: &sqlx::Pool, + path: &str, + value: &mut Value, + app: i64, +) -> Result<()> { match value { Value::Object(object) => { if let Some(Value::Object(script)) = object.get_mut("inlineScript") { - if script - .get("language") - .and_then(|x| x.as_str()) - .is_some_and(|x| x == "frontend") - { + let language = script.get("language").cloned(); + if language == Some(Value::String("frontend".to_owned())) { return Ok(()); } // replace `content` with an empty string: @@ -1764,18 +1866,26 @@ async fn reduce_app(db: &sqlx::Pool, value: &mut Value, app: i64 Value::String(s) => Some(s), _ => None, }); - let id = insert_app_script(db, app, code, lock).await?; + let id = insert_app_script( + db, + path, + app, + code, + lock, + language.map(|v| from_value(v).ok()).flatten(), + ) + .await?; // insert the `id` into the `script` object: script.insert("id".to_string(), json!(id.0)); } else { for (_, value) in object { - Box::pin(reduce_app(db, value, app)).await?; + Box::pin(reduce_app(db, path, value, app)).await?; } } } Value::Array(array) => { for value in array { - Box::pin(reduce_app(db, value, app)).await?; + Box::pin(reduce_app(db, path, value, app)).await?; } } _ => {} @@ -1809,6 +1919,9 @@ async fn lock_modules_app( base_internal_url: &str, token: &str, occupancy_metrics: &mut OccupancyMetrics, + locks_to_reload: &Option>, + // Represents the closest container id + container_id: Option, ) -> Result { match value { Value::Object(mut m) => { @@ -1842,7 +1955,24 @@ async fn lock_modules_app( .unwrap_or_default() .to_string(); let mut logs = "".to_string(); - if v.get("lock") + + if let Some((l, id)) = locks_to_reload + .as_ref() + .zip(container_id.as_ref()) + // TODO: Remove fallback + .and_then(|e| { + if *WMDEBUG_NO_COMPONENTS_TO_RELOCK { + None + } else { + Some(e) + } + }) + { + if !l.contains(id) { + return Ok(Value::Object(m.clone())); + } + } else if v + .get("lock") .is_some_and(|x| !x.as_str().unwrap().trim().is_empty()) { if skip_creating_new_lock(&language, &content) { @@ -1875,6 +2005,46 @@ async fn lock_modules_app( match new_lock { Ok(new_lock) => { append_logs(&job.id, &job.workspace_id, logs, &db.into()).await; + + let mut tx = db.begin().await?; + + tx = clear_dependency_map_for_item( + &job_path, + &job.workspace_id, + "app", + tx, + &container_id, + ) + .await?; + + let relative_imports = extract_relative_imports( + &content, + &format!("{job_path}/app"), + &Some(language.clone()), + ); + + if let Some(relative_imports) = relative_imports { + let mut logs = "".to_string(); + logs.push_str( + format!("\n\n--- RELATIVE IMPORTS ---\n\n").as_str(), + ); + + tx = add_relative_imports_to_dependency_map( + &job_path, + &job.workspace_id, + relative_imports, + "app", + tx, + &mut logs, + container_id, + ) + .await?; + append_logs(&job.id, &job.workspace_id, logs, &db.into()) + .await; + } + + tx.commit().await?; + let anns = windmill_common::worker::TypeScriptAnnotations::parse( &content, @@ -1928,6 +2098,11 @@ async fn lock_modules_app( base_internal_url, token, occupancy_metrics, + locks_to_reload, + m.get("id") + .and_then(Value::as_str) + .map(str::to_owned) + .or(container_id.clone()), ) .await?, ); @@ -1951,6 +2126,8 @@ async fn lock_modules_app( base_internal_url, token, occupancy_metrics, + locks_to_reload, + container_id.clone(), ) .await?, ); @@ -1985,6 +2162,22 @@ pub async fn handle_app_dependency_job( .ok_or_else(|| Error::internal_err("App Dependency requires script hash".to_owned()))? .0; + let components_to_relock = job + .args + .as_ref() + .map(|x| { + x.get("components_to_relock") + .map(|v| serde_json::from_str::>(v.get()).ok()) + .flatten() + }) + .flatten(); + + let triggered_by_relative_import = job + .args + .as_ref() + .map(|x| x.get("triggered_by_relative_import").is_some()) + .unwrap_or_default(); + sqlx::query!( "DELETE FROM workspace_runnable_dependencies WHERE app_path = $1 AND workspace_id = $2", job_path, @@ -1998,6 +2191,7 @@ pub async fn handle_app_dependency_job( .await? .map(|record| (record.app_id, record.value)); + // TODO: Use transaction for entire segment? if let Some((app_id, value)) = record { let value = lock_modules_app( value, @@ -2012,12 +2206,14 @@ pub async fn handle_app_dependency_job( base_internal_url, token, occupancy_metrics, + &components_to_relock, + None, ) .await?; // Compute a lite version of the app value (w/ `inlineScript.{lock,code}`). let mut value_lite = value.clone(); - reduce_app(db, &mut value_lite, app_id).await?; + reduce_app(db, &job_path, &mut value_lite, app_id).await?; if let Value::Object(object) = &mut value_lite { object.insert("version".to_string(), json!(id)); } @@ -2049,6 +2245,20 @@ pub async fn handle_app_dependency_job( .execute(db) .await?; + // NOTE: Temporary solution. + // Ideally we do this for every job regardless whether it was triggered by relative import or by creation/update of the app. + // NOTE: For now is not solving any problem but at some point we will introduce latest version caching + // and when we do this will be last operation that will make new version appear as the latest and will trigger cache invalidation for all worker. + if triggered_by_relative_import { + sqlx::query!( + "UPDATE app SET versions = array_append(versions, $1::bigint) WHERE path = $2 AND workspace_id = $3", + id, + &job_path, + &job.workspace_id + ) + .execute(db) + .await?; + } let (deployment_message, parent_path) = get_deployment_msg_and_parent_path_from_args(job.args.clone());