diff --git a/backend/.sqlx/query-819c4d78f10979a83529155c3f101f42f6958742bc68ab384944fa3c5cba094c.json b/backend/.sqlx/query-2d73fe76e45becbcae7a9d4384bbb2f9739cab232b8f3e357e27dc32bb35c3a8.json similarity index 62% rename from backend/.sqlx/query-819c4d78f10979a83529155c3f101f42f6958742bc68ab384944fa3c5cba094c.json rename to backend/.sqlx/query-2d73fe76e45becbcae7a9d4384bbb2f9739cab232b8f3e357e27dc32bb35c3a8.json index b25b301a04..23b00521d2 100644 --- a/backend/.sqlx/query-819c4d78f10979a83529155c3f101f42f6958742bc68ab384944fa3c5cba094c.json +++ b/backend/.sqlx/query-2d73fe76e45becbcae7a9d4384bbb2f9739cab232b8f3e357e27dc32bb35c3a8.json @@ -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" } diff --git a/backend/.sqlx/query-b44afcf1b9c047ac525638f0952c2cb01d65b1b46693331ac157dfca0dab6824.json b/backend/.sqlx/query-65e33a933ed8568c82f9b4e121b9ba4280123fe22a959d4d3d978ca767aeed29.json similarity index 88% rename from backend/.sqlx/query-b44afcf1b9c047ac525638f0952c2cb01d65b1b46693331ac157dfca0dab6824.json rename to backend/.sqlx/query-65e33a933ed8568c82f9b4e121b9ba4280123fe22a959d4d3d978ca767aeed29.json index 9187cc29b7..8e059f92ed 100644 --- a/backend/.sqlx/query-b44afcf1b9c047ac525638f0952c2cb01d65b1b46693331ac157dfca0dab6824.json +++ b/backend/.sqlx/query-65e33a933ed8568c82f9b4e121b9ba4280123fe22a959d4d3d978ca767aeed29.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "SELECT content AS \"content!: String\",\n lock AS \"lock: String\", language AS \"language: Option\", envs AS \"envs: Vec\", 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\", envs AS \"envs: Vec\", 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" } diff --git a/backend/.sqlx/query-0e7d95f4913e5775651971d741a3b5c1ef5dfe079be5325abe2866d39a7fe5fb.json b/backend/.sqlx/query-fa85a22c77844bb2c538b802b0984d4d69aecdd8a65db02e9495b0d6011c5c1f.json similarity index 65% rename from backend/.sqlx/query-0e7d95f4913e5775651971d741a3b5c1ef5dfe079be5325abe2866d39a7fe5fb.json rename to backend/.sqlx/query-fa85a22c77844bb2c538b802b0984d4d69aecdd8a65db02e9495b0d6011c5c1f.json index 9548fa6997..6c54ede027 100644 --- a/backend/.sqlx/query-0e7d95f4913e5775651971d741a3b5c1ef5dfe079be5325abe2866d39a7fe5fb.json +++ b/backend/.sqlx/query-fa85a22c77844bb2c538b802b0984d4d69aecdd8a65db02e9495b0d6011c5c1f.json @@ -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" } diff --git a/backend/src/main.rs b/backend/src/main.rs index 257ad77d1d..23978d7e62 100644 --- a/backend/src/main.rs +++ b/backend/src/main.rs @@ -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::>().as_slice() { diff --git a/backend/tests/deleted_script_version.rs b/backend/tests/deleted_script_version.rs new file mode 100644 index 0000000000..fdc49043b9 --- /dev/null +++ b/backend/tests/deleted_script_version.rs @@ -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) { + 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) { + 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 = 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 = + events.iter().map(|e| serde_json::from_str(e).unwrap()).collect(); + assert_eq!(events.len(), 2); + let hashes: Vec = events.iter().flat_map(|e| e.hashes.clone()).collect(); + let paths: Vec = events.iter().flat_map(|e| e.paths.clone()).collect(); + assert_eq!(hashes, deleted.hashes); + assert_eq!(paths, deleted.paths); +} diff --git a/backend/windmill-api-integration-tests/tests/scripts.rs b/backend/windmill-api-integration-tests/tests/scripts.rs index 3a146add27..d6afc9f302 100644 --- a/backend/windmill-api-integration-tests/tests/scripts.rs +++ b/backend/windmill-api-integration-tests/tests/scripts.rs @@ -15,6 +15,17 @@ fn authed(builder: reqwest::RequestBuilder) -> reqwest::RequestBuilder { builder.header("Authorization", "Bearer SECRET_TOKEN") } +async fn last_script_deletion(db: &Pool) -> 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) -> 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) -> 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::().await?, false); diff --git a/backend/windmill-api-scripts/src/scripts.rs b/backend/windmill-api-scripts/src/scripts.rs index 80a175dd99..bce67755bb 100644 --- a/backend/windmill-api-scripts/src/scripts.rs +++ b/backend/windmill-api-scripts/src/scripts.rs @@ -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( diff --git a/backend/windmill-api-workspaces/src/workspaces.rs b/backend/windmill-api-workspaces/src/workspaces.rs index 8c891bf0cf..81fead981b 100644 --- a/backend/windmill-api-workspaces/src/workspaces.rs +++ b/backend/windmill-api-workspaces/src/workspaces.rs @@ -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( diff --git a/backend/windmill-common/src/cache.rs b/backend/windmill-common/src/cache.rs index 3c385da79f..1518e8e8e3 100644 --- a/backend/windmill-common/src/cache.rs +++ b/backend/windmill-common/src/cache.rs @@ -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, diff --git a/backend/windmill-common/src/lib.rs b/backend/windmill-common/src/lib.rs index 0a1b0974d5..c5c9bde0d3 100644 --- a/backend/windmill-common/src/lib.rs +++ b/backend/windmill-common/src/lib.rs @@ -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, + pub hashes: Vec, +} + +/// 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) -> Self { + let (mut paths, hashes): (Vec, Vec) = 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) {