mirror of
https://github.com/windmill-labs/windmill.git
synced 2026-10-04 16:02:17 +00:00
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) <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 5.5
parent
35ec16d127
commit
fa419af9b3
@@ -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/<wid>/…` 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<bool> {
|
||||
.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);
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user