From fa419af9b3f7d65a607bb8ea2bcea58223f743c4 Mon Sep 17 00:00:00 2001 From: Ruben Fiszel Date: Sun, 4 Oct 2026 08:55:30 +0200 Subject: [PATCH] fix: drop a fork's datatable database despite connected sessions (#11509) Claude-Session: https://claude.ai/code/session_019cpWeGy8jRMtWqvuRBpgu4 Co-authored-by: Claude Opus 5.5 (1M context) --- .../src/workspaces_extra.rs | 91 ++++++++++++++++++- 1 file changed, 86 insertions(+), 5 deletions(-) diff --git a/backend/windmill-api-workspaces/src/workspaces_extra.rs b/backend/windmill-api-workspaces/src/workspaces_extra.rs index 7e2958bb58..a0ff1f2370 100644 --- a/backend/windmill-api-workspaces/src/workspaces_extra.rs +++ b/backend/windmill-api-workspaces/src/workspaces_extra.rs @@ -1670,13 +1670,12 @@ pub async fn drop_forked_datatable_databases( match parent_pg.connect(Some(&db)).await { Ok((client, connection)) => { let join_handle = tokio::spawn(async move { connection.await }); - if let Err(e) = client - .execute(&format!("DROP DATABASE \"{}\"", db_to_drop), &[]) - .await - { + if let Err(e) = drop_database_with_sessions(&client, db_to_drop).await { errors.push(format!( "Could not drop database '{}' for datatable://{}: {}", - db_to_drop, dt_name, e + db_to_drop, + dt_name, + windmill_common::error::pg_error_message(&e) )); } drop(client); @@ -1695,6 +1694,44 @@ pub async fn drop_forked_datatable_databases( Ok(Json(errors)) } +/// Drops `dbname`, whose name the caller validated, along with the sessions still connected to +/// it: a worker keeps an idle connection to a database it ran a job on, and a plain +/// `DROP DATABASE` fails for as long as one exists. +async fn drop_database_with_sessions( + client: &tokio_postgres::Client, + dbname: &str, +) -> std::result::Result<(), tokio_postgres::Error> { + let version = client + .query_one("SELECT current_setting('server_version_num')::int", &[]) + .await? + .get::<_, i32>(0); + if version >= 130000 { + client + .execute(&format!("DROP DATABASE \"{dbname}\" WITH (FORCE)"), &[]) + .await?; + return Ok(()); + } + // No `FORCE` before 13. The drop below waits a few seconds for the sessions to exit, and + // reports the ones this role was not allowed to terminate. + if let Err(e) = client + .execute( + "SELECT pg_terminate_backend(pid) FROM pg_stat_activity + WHERE datname = $1 AND pid <> pg_backend_pid()", + &[&dbname], + ) + .await + { + tracing::warn!( + "Failed to terminate connections to '{dbname}': {}", + windmill_common::error::pg_error_message(&e) + ); + } + client + .execute(&format!("DROP DATABASE \"{dbname}\""), &[]) + .await?; + Ok(()) +} + /// Drop this fork workspace's ducklake namespaces: the `wm_fork_*` metadata schema in each /// lake's catalog database, plus (best effort) the fork's `__wm_forks//…` data files in /// the workspace storage. Driven by the `fork_ducklake_namespace` registry written at first @@ -2265,3 +2302,47 @@ async fn workspace_is_fork(db: &DB, w_id: &str) -> Result { .await? .unwrap_or(false)) } + +#[cfg(test)] +mod tests { + use super::drop_database_with_sessions; + + #[sqlx::test(migrations = false)] + async fn a_database_is_dropped_under_a_connected_session(pool: sqlx::PgPool) { + let mut config: tokio_postgres::Config = + std::env::var("DATABASE_URL").unwrap().parse().unwrap(); + let test_db = pool.connect_options().get_database().unwrap().to_string(); + let (client, connection) = config + .dbname(&test_db) + .connect(tokio_postgres::NoTls) + .await + .unwrap(); + tokio::spawn(connection); + // The test database's name already fills an identifier. + let held_db = format!("held_{}", &test_db[test_db.len().saturating_sub(40)..]); + client + .execute(&format!("CREATE DATABASE \"{held_db}\""), &[]) + .await + .unwrap(); + let (_session, connection) = config + .dbname(&held_db) + .connect(tokio_postgres::NoTls) + .await + .unwrap(); + tokio::spawn(connection); + + drop_database_with_sessions(&client, &held_db) + .await + .unwrap(); + + let left = client + .query_one( + "SELECT count(*) FROM pg_database WHERE datname = $1", + &[&held_db], + ) + .await + .unwrap() + .get::<_, i64>(0); + assert_eq!(left, 0); + } +}