fix(apps): fix relative imports cache invalidation (#6564)

* v0

Signed-off-by: pyranota <pyra@duck.com>

* optimize relocks

* make it work with relative relative imports

Signed-off-by: pyranota <pyra@duck.com>

* use fallback

Signed-off-by: pyranota <pyra@duck.com>

* remove dbg and todos

Signed-off-by: pyranota <pyra@duck.com>

* future proof a bit

Signed-off-by: pyranota <pyra@duck.com>

* cleanup

Signed-off-by: pyranota <pyra@duck.com>

* more cleanup

Signed-off-by: pyranota <pyra@duck.com>

* remove final TODO

Signed-off-by: pyranota <pyra@duck.com>

* do not use bytemuck

Signed-off-by: pyranota <pyra@duck.com>

* optimize hashing

Signed-off-by: pyranota <pyra@duck.com>

---------

Signed-off-by: pyranota <pyra@duck.com>
This commit is contained in:
pyranota
2025-09-11 12:46:01 +02:00
committed by GitHub
parent 9ab3c1eb37
commit f48016048b
4 changed files with 287 additions and 32 deletions
@@ -15,7 +15,7 @@
]
},
"nullable": [
null
true
]
},
"hash": "5a219a2532517869578c4504ff3153c43903f929ae5d62fbba12610f89c36d55"
@@ -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"
}
@@ -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"
}
+241 -31
View File
@@ -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<ScriptLang>,
) -> Result<Vec<u8>> {
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<sqlx::Postgres>,
path: &str,
app: i64,
code: String,
lock: Option<String>,
language: Option<ScriptLang>,
) -> Result<AppScriptId> {
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<sqlx::Postgres>, value: &mut Value, app: i64) -> Result<()> {
async fn reduce_app(
db: &sqlx::Pool<sqlx::Postgres>,
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<sqlx::Postgres>, 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<Vec<String>>,
// Represents the closest container id
container_id: Option<String>,
) -> Result<Value> {
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::<Vec<String>>(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());