fix: stop running deleted script versions from worker caches (#11456)

* fix: stop running deleted script versions from worker caches

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

* refactor: reuse script cache invalidate and test the deletion notify payload

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

* fix: keep a deleted copy from shadowing the same hash in another workspace

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

* fix: drop deletions missed while down and fills racing the eviction

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

* fix: batch deletion events, evict path caches and cover version pruning

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

* fix: evict the raw-import cache on both passes and split large deletion events

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

---------

Co-authored-by: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
This commit is contained in:
Ruben Fiszel
2026-10-01 13:37:00 +02:00
committed by GitHub
co-authored by Claude Opus 5.5
parent 840567fbe5
commit d2830cb4a0
10 changed files with 270 additions and 39 deletions
@@ -1,12 +1,17 @@
{
"db_name": "PostgreSQL",
"query": "DELETE FROM script WHERE workspace_id = $1 AND path = ANY($2) RETURNING path",
"query": "DELETE FROM script WHERE workspace_id = $1 AND path = ANY($2) RETURNING path, hash",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "path",
"type_info": "Varchar"
},
{
"ordinal": 1,
"name": "hash",
"type_info": "Int8"
}
],
"parameters": {
@@ -16,8 +21,9 @@
]
},
"nullable": [
false,
false
]
},
"hash": "819c4d78f10979a83529155c3f101f42f6958742bc68ab384944fa3c5cba094c"
"hash": "2d73fe76e45becbcae7a9d4384bbb2f9739cab232b8f3e357e27dc32bb35c3a8"
}
@@ -1,6 +1,6 @@
{
"db_name": "PostgreSQL",
"query": "SELECT content AS \"content!: String\",\n lock AS \"lock: String\", language AS \"language: Option<ScriptLang>\", envs AS \"envs: Vec<String>\", schema AS \"schema: String\", schema_validation AS \"schema_validation: bool\", codebase LIKE '%.tar' as use_tar, codebase LIKE '%.esm%' as is_esm, modules AS \"modules: serde_json::Value\" FROM script WHERE hash = $1 LIMIT 1",
"query": "SELECT content AS \"content!: String\",\n lock AS \"lock: String\", language AS \"language: Option<ScriptLang>\", envs AS \"envs: Vec<String>\", schema AS \"schema: String\", schema_validation AS \"schema_validation: bool\", codebase LIKE '%.tar' as use_tar, codebase LIKE '%.esm%' as is_esm, modules AS \"modules: serde_json::Value\", deleted FROM script WHERE hash = $1 ORDER BY deleted LIMIT 1",
"describe": {
"columns": [
{
@@ -80,6 +80,11 @@
"ordinal": 8,
"name": "modules: serde_json::Value",
"type_info": "Jsonb"
},
{
"ordinal": 9,
"name": "deleted",
"type_info": "Bool"
}
],
"parameters": {
@@ -96,8 +101,9 @@
false,
null,
null,
true
true,
false
]
},
"hash": "b44afcf1b9c047ac525638f0952c2cb01d65b1b46693331ac157dfca0dab6824"
"hash": "65e33a933ed8568c82f9b4e121b9ba4280123fe22a959d4d3d978ca767aeed29"
}
@@ -1,12 +1,12 @@
{
"db_name": "PostgreSQL",
"query": "DELETE FROM script WHERE path = $1 AND workspace_id = $2 RETURNING path",
"query": "DELETE FROM script WHERE path = $1 AND workspace_id = $2 RETURNING hash",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "path",
"type_info": "Varchar"
"name": "hash",
"type_info": "Int8"
}
],
"parameters": {
@@ -19,5 +19,5 @@
false
]
},
"hash": "0e7d95f4913e5775651971d741a3b5c1ef5dfe079be5325abe2866d39a7fe5fb"
"hash": "fa85a22c77844bb2c538b802b0984d4d69aecdd8a65db02e9495b0d6011c5c1f"
}
+11
View File
@@ -1232,6 +1232,11 @@ Windmill Community Edition {GIT_VERSION}
default_base_internal_url.clone()
};
// Deletions are delivered by `notify_event` to running processes only, so one that
// was down when a version was deleted would keep its code on disk: start from an
// empty script cache, refilled from the database.
windmill_common::cache::script::clear();
initial_load(
&conn,
killpill_tx.clone(),
@@ -1884,6 +1889,12 @@ async fn process_notify_event(
);
windmill_common::workspaces::PUBLIC_APP_RATE_LIMIT_CACHE.remove(payload);
}
windmill_common::SCRIPT_VERSION_DELETED_CHANNEL => {
match serde_json::from_str(payload) {
Ok(deleted) => windmill_api_scripts::scripts::evict_deleted_script_versions(deleted),
Err(e) => tracing::error!("Invalid script version deletion payload {payload}: {e}"),
}
}
"notify_runnable_version_change" => {
tracing::info!("Runnable version change detected: {}", payload);
match payload.split(':').collect::<Vec<&str>>().as_slice() {
+56
View File
@@ -0,0 +1,56 @@
//! Script data is cached by hash, and the hash leaves out the workspace, so forks and clones
//! share it: a deleted copy in one workspace must not stop the live copy in another.
use sqlx::{Pool, Postgres};
use windmill_common::{cache, scripts::ScriptHash, worker::Connection};
#[sqlx::test(fixtures("base"))]
async fn a_deleted_copy_does_not_shadow_a_live_one(db: Pool<Postgres>) {
const SHARED: i64 = 0x5de1_e7ed_0002;
sqlx::query("INSERT INTO workspace (id, name, owner) VALUES ('fork', 'fork', 'test-user')")
.execute(&db)
.await
.unwrap();
sqlx::query(
"INSERT INTO script (workspace_id, hash, path, summary, description, content,
created_by, language, lock, archived, deleted)
VALUES ('fork', $1, 'f/infra/tool', '', '', '', 'test-user', 'bash', '', true, true),
('test-workspace', $1, 'f/infra/tool', '', '', 'echo shared', 'test-user',
'bash', '', false, false)",
)
.bind(SHARED)
.execute(&db)
.await
.unwrap();
let conn = Connection::Sql(db.clone());
let (data, _) = cache::script::fetch(&conn, ScriptHash(SHARED))
.await
.unwrap();
assert_eq!(data.code, "echo shared");
}
#[sqlx::test(fixtures("base"))]
async fn a_large_deletion_is_split_across_events(db: Pool<Postgres>) {
let deleted = windmill_common::DeletedScriptVersions::new(
"test-workspace",
(0..501).map(|h| (format!("f/infra/s{}", h % 2), h)),
);
let mut conn = db.acquire().await.unwrap();
deleted.notify(&mut conn).await.unwrap();
let events: Vec<String> = sqlx::query_scalar(
"SELECT payload FROM notify_event WHERE channel = $1 ORDER BY id",
)
.bind(windmill_common::SCRIPT_VERSION_DELETED_CHANNEL)
.fetch_all(&db)
.await
.unwrap();
let events: Vec<windmill_common::DeletedScriptVersions> =
events.iter().map(|e| serde_json::from_str(e).unwrap()).collect();
assert_eq!(events.len(), 2);
let hashes: Vec<i64> = events.iter().flat_map(|e| e.hashes.clone()).collect();
let paths: Vec<String> = events.iter().flat_map(|e| e.paths.clone()).collect();
assert_eq!(hashes, deleted.hashes);
assert_eq!(paths, deleted.paths);
}
@@ -15,6 +15,17 @@ fn authed(builder: reqwest::RequestBuilder) -> reqwest::RequestBuilder {
builder.header("Authorization", "Bearer SECRET_TOKEN")
}
async fn last_script_deletion(db: &Pool<Postgres>) -> windmill_common::DeletedScriptVersions {
let payload: String = sqlx::query_scalar(
"SELECT payload FROM notify_event WHERE channel = $1 ORDER BY id DESC LIMIT 1",
)
.bind(windmill_common::SCRIPT_VERSION_DELETED_CHANNEL)
.fetch_one(db)
.await
.unwrap();
serde_json::from_str(&payload).unwrap()
}
async fn authed_get(port: u16, endpoint: &str, path: &str) -> reqwest::Response {
authed(client().get(script_url(port, endpoint, path)))
.send()
@@ -366,6 +377,12 @@ async fn test_script_endpoints(db: Pool<Postgres>) -> anyhow::Result<()> {
.await
.unwrap();
assert_eq!(resp.status(), 200);
// Every process evicts deleted versions from its script caches on this event.
let deleted = last_script_deletion(&db).await;
let another: windmill_common::scripts::ScriptHash =
serde_json::from_value(json!(another_hash))?;
assert_eq!(deleted.paths, ["u/test-user/another_script"]);
assert_eq!(deleted.hashes, [another.0]);
// --- delete_bulk ---
let resp = authed(client().delete(format!("{base}/delete_bulk")))
@@ -374,6 +391,9 @@ async fn test_script_endpoints(db: Pool<Postgres>) -> anyhow::Result<()> {
.await
.unwrap();
assert_eq!(resp.status(), 200, "delete_bulk: {}", resp.text().await?);
let deleted = last_script_deletion(&db).await;
assert_eq!(deleted.paths, ["u/test-user/test_script"]);
assert!(!deleted.hashes.is_empty(), "{deleted:?}");
let resp = authed_get(port, "exists/p", "u/test-user/test_script").await;
assert_eq!(resp.json::<bool>().await?, false);
+45 -16
View File
@@ -17,7 +17,7 @@ use windmill_common::{
utils::{BulkDeleteRequest, WithStarredInfoQuery, HTTP_CLIENT},
webhook::{WebhookMessage, WebhookShared},
workspaces::{check_deploy_rules, RuleCheckResult},
DB,
DeletedScriptVersions, DB,
};
use windmill_queue::schedule::clear_schedule;
@@ -818,6 +818,16 @@ fn invalidate_script_path_caches(w_id: &str, script_path: &str) {
RAW_SCRIPT_LATEST_HASH_CACHE.remove(&format!("{w_id}:{script_path}"));
}
/// [`DeletedScriptVersions::evict`], plus this crate's own path -> hash cache. Run by the
/// deleting process after its commit and by every process on the deletion event.
pub fn evict_deleted_script_versions(deleted: DeletedScriptVersions) {
deleted.evict_with(|deleted| {
for path in &deleted.paths {
RAW_SCRIPT_LATEST_HASH_CACHE.remove(&format!("{}:{path}", deleted.workspace_id));
}
});
}
/// What a script deploy still has to do once its transaction has committed.
enum PostCommitDeploy {
/// Everything, for a deploy with no dependency job to hand it to.
@@ -4161,6 +4171,8 @@ async fn delete_script_by_hash(
// the script was never a pipeline member.
clear_script_triggers(&mut *tx, &w_id, &script.path, AssetUsageKind::Script).await?;
clear_macro_registry(&mut *tx, &w_id, &script.path).await?;
let deleted = DeletedScriptVersions::new(&w_id, [(script.path.clone(), hash.0)]);
deleted.notify(&mut *tx).await?;
audit_log(
&mut *tx,
@@ -4173,6 +4185,7 @@ async fn delete_script_by_hash(
)
.await?;
tx.commit().await?;
evict_deleted_script_versions(deleted);
webhook.send_message(
w_id.clone(),
@@ -4240,14 +4253,24 @@ async fn delete_script_by_path(
.fetch_all(&mut *tx)
.await?;
let script = sqlx::query_scalar!(
"DELETE FROM script WHERE path = $1 AND workspace_id = $2 RETURNING path",
let deleted_hashes = sqlx::query_scalar!(
"DELETE FROM script WHERE path = $1 AND workspace_id = $2 RETURNING hash",
path,
w_id
)
.fetch_one(&mut *tx)
.fetch_all(&mut *tx)
.await
.map_err(|e| Error::internal_err(format!("deleting script by path {w_id}: {e:#}")))?;
if deleted_hashes.is_empty() {
return Err(Error::NotFound(format!(
"script {path} not found in {w_id}"
)));
}
let script = path.to_string();
let deleted = DeletedScriptVersions::new(
&w_id,
deleted_hashes.into_iter().map(|h| (script.clone(), h)),
);
// After the DELETE, never before: every dbt writer locks the `script` row
// first, so taking a sidecar ahead of it deadlocks one of the pair. The
@@ -4287,6 +4310,7 @@ async fn delete_script_by_path(
// the script was never a pipeline member.
clear_script_triggers(&mut *tx, &w_id, path, AssetUsageKind::Script).await?;
clear_macro_registry(&mut *tx, &w_id, path).await?;
deleted.notify(&mut *tx).await?;
if !query.keep_captures.unwrap_or(false) {
sqlx::query!(
@@ -4317,6 +4341,7 @@ async fn delete_script_by_path(
)
.await?;
tx.commit().await?;
evict_deleted_script_versions(deleted);
handle_deployment_metadata(
&authed.email,
@@ -4412,14 +4437,21 @@ async fn delete_scripts_bulk(
}
}
let mut deleted_paths = sqlx::query_scalar!(
"DELETE FROM script WHERE workspace_id = $1 AND path = ANY($2) RETURNING path",
w_id,
&request.paths
)
.fetch_all(&mut *tx)
.await
.map_err(|e| Error::internal_err(format!("deleting scripts in bulk {w_id}: {e:#}")))?;
let deleted = DeletedScriptVersions::new(
&w_id,
sqlx::query!(
"DELETE FROM script WHERE workspace_id = $1 AND path = ANY($2) RETURNING path, hash",
w_id,
&request.paths
)
.fetch_all(&mut *tx)
.await
.map_err(|e| Error::internal_err(format!("deleting scripts in bulk {w_id}: {e:#}")))?
.into_iter()
.map(|r| (r.path, r.hash)),
);
deleted.notify(&mut *tx).await?;
let deleted_paths = deleted.paths.clone();
// Same reason as the single-path delete, over every requested path rather
// than the deleted ones: a path that had no script left can still hold state.
@@ -4428,10 +4460,6 @@ async fn delete_scripts_bulk(
windmill_common::dbt_manifest::clear_dbt_editor_graphs(&mut tx, &w_id, p).await?;
}
// remove duplicates from deleted_paths
deleted_paths.sort();
deleted_paths.dedup();
sqlx::query!(
"DELETE FROM draft WHERE workspace_id = $1 AND path = ANY($2) AND typ = 'script'",
w_id,
@@ -4472,6 +4500,7 @@ async fn delete_scripts_bulk(
.await?;
tx.commit().await?;
evict_deleted_script_versions(deleted);
try_join_all(deleted_paths.iter().map(|path| {
handle_deployment_metadata(
@@ -14367,18 +14367,27 @@ async fn prune_versions(
let pruned = match req.resource_type.as_str() {
"scripts" => {
let result = sqlx::query(
"DELETE FROM script
WHERE workspace_id = $1 AND hash NOT IN (
SELECT DISTINCT ON (path) hash FROM script
WHERE workspace_id = $1 AND deleted = false
ORDER BY path, created_at DESC
)",
)
.bind(&w_id)
.execute(&db)
.await?;
result.rows_affected()
let mut tx = db.begin().await?;
let deleted = windmill_common::DeletedScriptVersions::new(
&w_id,
sqlx::query_as::<_, (String, i64)>(
"DELETE FROM script
WHERE workspace_id = $1 AND hash NOT IN (
SELECT DISTINCT ON (path) hash FROM script
WHERE workspace_id = $1 AND deleted = false
ORDER BY path, created_at DESC
)
RETURNING path, hash",
)
.bind(&w_id)
.fetch_all(&mut *tx)
.await?,
);
deleted.notify(&mut *tx).await?;
tx.commit().await?;
let pruned = deleted.hashes.len() as u64;
deleted.evict();
pruned
}
"flows" => {
let deleted = sqlx::query(
+11 -2
View File
@@ -687,8 +687,9 @@ pub mod script {
schema_validation AS \"schema_validation: bool\", \
codebase LIKE '%.tar' as use_tar, \
codebase LIKE '%.esm%' as is_esm, \
modules AS \"modules: serde_json::Value\" \
FROM script WHERE hash = $1 LIMIT 1",
modules AS \"modules: serde_json::Value\", \
deleted \
FROM script WHERE hash = $1 ORDER BY deleted LIMIT 1",
hash.0
)
.fetch_optional(db)
@@ -696,6 +697,14 @@ pub mod script {
.map_err(Into::into)
.and_then(unwrap_or_error(&loc, "Script", hash))
.and_then(|r| {
// The hash leaves out the workspace, so forks and clones share it, and a deleted
// copy keeps its row with the content wiped. Any live copy holds the same code;
// only when every copy is deleted is there nothing left to run.
if r.deleted {
return Err(error::Error::NotFound(format!(
"Script version {hash} was deleted"
)));
}
Ok(RawScript {
content: r.content,
lock: r.lock,
+85
View File
@@ -2479,6 +2479,91 @@ pub fn invalidate_deployed_script_hash_cache(w_id: &str, script_path: &str) {
DEPLOYED_SCRIPT_HASH_CACHE.remove(&(w_id.to_string(), script_path.to_string()));
}
pub const SCRIPT_VERSION_DELETED_CHANNEL: &str = "notify_script_version_deleted";
/// The payload of a [`SCRIPT_VERSION_DELETED_CHANNEL`] event. `paths` and `hashes` are the
/// deleted versions' paths and hashes as independent sets, not paired by position.
#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
pub struct DeletedScriptVersions {
pub workspace_id: String,
pub paths: Vec<String>,
pub hashes: Vec<i64>,
}
/// Bounds each event's payload: pruning a workspace deletes every past version of every path
/// in one call, and every process logs each event it receives.
const DELETED_SCRIPT_VERSIONS_PER_EVENT: usize = 500;
impl DeletedScriptVersions {
pub fn new(workspace_id: &str, deleted: impl IntoIterator<Item = (String, i64)>) -> Self {
let (mut paths, hashes): (Vec<String>, Vec<i64>) = deleted.into_iter().unzip();
paths.sort();
paths.dedup();
Self { workspace_id: workspace_id.to_string(), paths, hashes }
}
/// Tell every replica to drop these versions from its caches, in the transaction that
/// deletes them. Authorization is the caller's: only call it for versions the caller was
/// allowed to delete.
pub async fn notify(&self, db: &mut sqlx::PgConnection) -> error::Result<()> {
let per_event = DELETED_SCRIPT_VERSIONS_PER_EVENT;
let events = self.paths.len().max(self.hashes.len()).div_ceil(per_event);
for i in 0..events {
let range = i * per_event..(i + 1) * per_event;
let event = Self {
workspace_id: self.workspace_id.clone(),
paths: self
.paths
.get(range.start..range.end.min(self.paths.len()))
.unwrap_or_default()
.to_vec(),
hashes: self
.hashes
.get(range.start..range.end.min(self.hashes.len()))
.unwrap_or_default()
.to_vec(),
};
sqlx::query("INSERT INTO notify_event (channel, payload) VALUES ($1, $2)")
.bind(SCRIPT_VERSION_DELETED_CHANNEL)
.bind(serde_json::to_string(&event)?)
.execute(&mut *db)
.await?;
}
Ok(())
}
/// [`Self::evict_with`] for this crate's caches only.
pub fn evict(self) {
self.evict_with(|_| {});
}
/// Script data is cached by hash, memory and disk, with no expiry: without this, a
/// process that ran a version before its deletion keeps running that version's code,
/// and a path keeps resolving to its deleted latest version until its cache expires.
/// `also` drops the entries of caches other crates own, on both passes. Needs no
/// authorization: it only drops cache entries, refilled from the database. Must run
/// inside a Tokio runtime.
pub fn evict_with(self, also: impl Fn(&Self) + Send + 'static) {
let evict = move || {
for hash in &self.hashes {
cache::script::invalidate(ScriptHash(*hash));
DEPLOYED_SCRIPT_INFO_CACHE.remove(&(self.workspace_id.clone(), *hash));
}
for path in &self.paths {
invalidate_latest_script_hash_caches(&self.workspace_id, path);
}
also(&self);
};
evict();
// A fill that read a row before the deletion committed can still be writing it to
// the cache: a second pass, once such a fill has had time to finish, removes it.
spawn(async move {
tokio::time::sleep(std::time::Duration::from_secs(30)).await;
evict();
});
}
}
/// Same, for a new version row, which also moves the import-side answer (that one has no lock
/// predicate, so only a new row moves it).
pub fn invalidate_latest_script_hash_caches(w_id: &str, script_path: &str) {