From d3ae20be67ddf33bbd63816596bb4ef818879725 Mon Sep 17 00:00:00 2001 From: Diego Imbert Date: Fri, 18 Sep 2026 15:48:37 +0200 Subject: [PATCH] fix(datatables): lock the fork copy at finalization, and order cleanup's registry write after its settings row - Fork finalization holds the copy's database lock from its availability check to the commit. - Cleanup removes the registry entry in its own transaction, in a task of its own, instead of on a second connection; a rename migrates reservations after its settings rewrites, so both take the settings rows before the registry. Co-Authored-By: Claude Opus 5 (1M context) --- .../windmill-api-workspaces/src/workspaces.rs | 9 ++ .../src/workspaces_extra.rs | 135 ++++++++++-------- backend/windmill-common/src/lib.rs | 22 +-- 3 files changed, 101 insertions(+), 65 deletions(-) diff --git a/backend/windmill-api-workspaces/src/workspaces.rs b/backend/windmill-api-workspaces/src/workspaces.rs index 363a810dda..a93407f6f1 100644 --- a/backend/windmill-api-workspaces/src/workspaces.rs +++ b/backend/windmill-api-workspaces/src/workspaces.rs @@ -8284,6 +8284,15 @@ async fn apply_forked_datatable( })?, }; + if database.resource_type == DataTableCatalogResourceType::Instance { + // Held until the fork commits, as every save newly naming a database takes it: none may + // claim the copy between the check below and this fork's entry landing on it. + windmill_common::datatable_roles::lock_instance_databases_governance( + &mut **tx, + [fdt.new_dbname.as_str()], + ) + .await?; + } if database.resource_type == DataTableCatalogResourceType::Instance && !windmill_api_auth::is_super_admin_authed(db, authed).await? { diff --git a/backend/windmill-api-workspaces/src/workspaces_extra.rs b/backend/windmill-api-workspaces/src/workspaces_extra.rs index 605f82da38..91b4567507 100644 --- a/backend/windmill-api-workspaces/src/workspaces_extra.rs +++ b/backend/windmill-api-workspaces/src/workspaces_extra.rs @@ -110,8 +110,6 @@ pub(crate) async fn change_workspace_id( .execute(&mut *tx) .await?; - migrate_fork_reservations(&mut tx, &old_id, &rw.new_id).await?; - // Duplicate workspace settings (keep copy in old workspace for reference) info!("Duplicating workspace_settings table"); sqlx::query!( @@ -852,6 +850,10 @@ pub(crate) async fn change_workspace_id( } } + // After every workspace_settings write above: fork cleanup locks a settings row before the + // registry, so taking the registry first here would deadlock with it. + migrate_fork_reservations(&mut tx, &old_id, &rw.new_id).await?; + // Audit log in the same transaction as the workspace changes audit_log( &mut *tx, @@ -1455,68 +1457,87 @@ pub async fn drop_forked_datatable_databases( // The fork's own entry is what is going away; anything else still reaching the copy, // a child fork's pointer at this entry included, keeps it. The lock keeps a child fork // from gaining such a pointer before the drop. - let dropped = async { - let mut tx = db.begin().await?; - // The three locks a settings save takes, in its order: this workspace's data - // tables, its settings row, and the database itself. Without them a save could - // rename this entry, or point another one here, either side of the check below. - windmill_common::workspaces::lock_fork_datatables(&mut tx, &w_id).await?; - // The snapshot above was read unlocked: a save committing since could have - // repointed this entry, and the entry is removed below whatever it names by then. - let current = sqlx::query_scalar::<_, Option>( - "SELECT datatable->'datatables'->$2 FROM workspace_settings + // A task of its own, so a client going away cannot stop it between dropping the + // database and committing the entry's removal. + let dropped = tokio::spawn({ + let (db, w_id, dt_name, db_to_drop) = ( + db.clone(), + w_id.clone(), + dt_name.clone(), + db_to_drop.clone(), + ); + let resource_type = database.resource_type; + async move { + let mut tx = db.begin().await?; + // The three locks a settings save takes, in its order: this workspace's data + // tables, its settings row, and the database itself. Without them a save could + // rename this entry, or point another one here, either side of the check below. + windmill_common::workspaces::lock_fork_datatables(&mut tx, &w_id).await?; + // The snapshot above was read unlocked: a save committing since could have + // repointed this entry, and the entry is removed below whatever it names by then. + let current = sqlx::query_scalar::<_, Option>( + "SELECT datatable->'datatables'->$2 FROM workspace_settings WHERE workspace_id = $1 FOR UPDATE", - ) - .bind(&w_id) - .bind(dt_name) - .fetch_optional(&mut *tx) - .await? - .flatten() - .and_then(|v| serde_json::from_value::(v).ok()); - if !current.is_some_and(|dt| { - dt.forked_from.is_some() - && dt.database.is_some_and(|d| { - d.resource_type == database.resource_type - && &d.resource_path == db_to_drop - }) - }) { - return Err(Error::BadRequest( - "the data table changed while it was being cleaned up".to_string(), - )); - } - windmill_common::datatable_roles::lock_instance_databases_governance( - &mut tx, - [db_to_drop.as_str()], - ) - .await?; - let uses = windmill_common::workspaces::managed_database_uses( - &mut tx, - windmill_common::workspaces::DataTableCatalogResourceType::Instance, - db_to_drop, - Some((&w_id, dt_name)), - ) - .await?; - if !uses.is_empty() { - return Err(Error::BadRequest(format!( - "it is still used by {}", - uses.join(", ") - ))); - } - // The entry goes with the database: a fork this one is cloned into afterwards must - // not inherit a pointer at a data table whose database is gone. - sqlx::query( + ) + .bind(&w_id) + .bind(&dt_name) + .fetch_optional(&mut *tx) + .await? + .flatten() + .and_then(|v| serde_json::from_value::(v).ok()); + if !current.is_some_and(|dt| { + dt.forked_from.is_some() + && dt.database.is_some_and(|d| { + d.resource_type == resource_type && d.resource_path == db_to_drop + }) + }) { + return Err(Error::BadRequest( + "the data table changed while it was being cleaned up".to_string(), + )); + } + windmill_common::datatable_roles::lock_instance_databases_governance( + &mut tx, + [db_to_drop.as_str()], + ) + .await?; + let uses = windmill_common::workspaces::managed_database_uses( + &mut tx, + windmill_common::workspaces::DataTableCatalogResourceType::Instance, + &db_to_drop, + Some((w_id.as_str(), dt_name.as_str())), + ) + .await?; + if !uses.is_empty() { + return Err(Error::BadRequest(format!( + "it is still used by {}", + uses.join(", ") + ))); + } + // The entry goes with the database: a fork this one is cloned into afterwards must + // not inherit a pointer at a data table whose database is gone. + sqlx::query( "UPDATE workspace_settings SET datatable = datatable #- ARRAY['datatables', $2] WHERE workspace_id = $1", ) .bind(&w_id) - .bind(dt_name) + .bind(&dt_name) .execute(&mut *tx) .await?; - windmill_common::drop_custom_instance_database(&db, db_to_drop).await?; - tx.commit().await?; - Ok::<_, Error>(()) - } - .await; + windmill_common::drop_custom_instance_database_keep_entry(&db, &db_to_drop) + .await?; + sqlx::query( + "UPDATE global_settings SET value = value #- ARRAY['databases', $1] + WHERE name = 'custom_instance_pg_databases'", + ) + .bind(&db_to_drop) + .execute(&mut *tx) + .await?; + tx.commit().await?; + Ok::<_, Error>(()) + } + }) + .await + .unwrap_or_else(|e| Err(Error::internal_err(format!("cleanup task failed: {e}")))); if let Err(e) = dropped { errors.push(format!( "Could not drop instance database '{}' for datatable://{}: {}", diff --git a/backend/windmill-common/src/lib.rs b/backend/windmill-common/src/lib.rs index db52194730..5920dbf658 100644 --- a/backend/windmill-common/src/lib.rs +++ b/backend/windmill-common/src/lib.rs @@ -1476,6 +1476,20 @@ pub fn validate_dbname(dbname: &str) -> error::Result<()> { /// Drop a custom instance database: validate, terminate connections, DROP DATABASE, remove from global_settings. pub async fn drop_custom_instance_database(db: &DB, dbname: &str) -> error::Result<()> { + drop_custom_instance_database_keep_entry(db, dbname).await?; + sqlx::query!( + r#"UPDATE global_settings SET value = value #- ARRAY['databases', $1] WHERE name = 'custom_instance_pg_databases'"#, + dbname.trim() + ) + .execute(db) + .await?; + Ok(()) +} + +/// [`drop_custom_instance_database`] leaving its registry entry, for a caller holding row locks in +/// a transaction: the registry write has to go through that transaction, as waiting on another +/// connection for a lock the transaction's own peers hold is a deadlock Postgres cannot see. +pub async fn drop_custom_instance_database_keep_entry(db: &DB, dbname: &str) -> error::Result<()> { let dbname = dbname.trim(); validate_dbname(dbname)?; @@ -1521,14 +1535,6 @@ pub async fn drop_custom_instance_database(db: &DB, dbname: &str) -> error::Resu tracing::info!("Database '{}' does not exist, skipping drop", dbname); } - // Always remove from global_settings - sqlx::query!( - r#"UPDATE global_settings SET value = value #- ARRAY['databases', $1] WHERE name = 'custom_instance_pg_databases'"#, - dbname - ) - .execute(db) - .await?; - Ok(()) }