fix: add validation to drop_custom_instance_database and use source db for CREATE/DROP

Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>
This commit is contained in:
Diego Imbert
2026-03-27 10:46:55 +01:00
co-authored by Claude Opus 4.5
parent fd60dd11f4
commit bef700b5cf
4 changed files with 48 additions and 48 deletions
+1 -30
View File
@@ -1013,36 +1013,7 @@ async fn drop_custom_instance_pg_database(
) -> Result<String> {
require_super_admin(&db, &authed.email).await?;
let dbname = dbname.trim();
if dbname.is_empty() {
return Err(error::Error::BadRequest(
"Database name cannot be empty".to_string(),
));
}
let wmill_pg_creds = PgDatabase::parse_uri(&get_database_url().await?.as_str().await)?;
if wmill_pg_creds.dbname.trim().eq_ignore_ascii_case(dbname) {
return Err(error::Error::BadRequest(
"Cannot drop the main Windmill database".to_string(),
));
}
let db_exists = sqlx::query_scalar!(
"SELECT EXISTS (SELECT 1 FROM pg_catalog.pg_database WHERE datname = $1)",
dbname
)
.fetch_one(&db)
.await?
.unwrap_or(false);
if !db_exists {
return Err(error::Error::NotFound(format!(
"Database '{}' does not exist",
dbname
)));
}
windmill_common::drop_custom_instance_database(&db, dbname).await;
windmill_common::drop_custom_instance_database(&db, &dbname).await?;
Ok(format!("Database '{}' dropped successfully", dbname))
}
@@ -1527,8 +1527,9 @@ async fn fork_pg_database(
}
// Optionally create the target database before importing
// Connect to the source database (which we know exists) to run CREATE DATABASE
if req.create_target_db {
let admin_pg = PgDatabase { dbname: "postgres".to_string(), ..target_pg.clone() };
let admin_pg = PgDatabase { dbname: source_pg.dbname.clone(), ..target_pg.clone() };
let (client, connection) = admin_pg.connect().await?;
let join_handle = tokio::spawn(async move { connection.await });
@@ -899,8 +899,11 @@ async fn drop_forked_datatable_databases(
if dt.database.resource_type
== windmill_common::workspaces::DataTableCatalogResourceType::Instance
{
windmill_common::drop_custom_instance_database(db, db_to_drop).await;
if let Err(e) = windmill_common::drop_custom_instance_database(db, db_to_drop).await {
tracing::error!("Failed to drop instance database '{}': {}", db_to_drop, e);
}
} else if let Some(original_resource) = forked_from.get("original_resource") {
// Connect to the original resource's database to run DROP on the forked db
let pg = match serde_json::from_value::<PgDatabase>(original_resource.clone()) {
Ok(pg) => pg,
Err(e) => {
+41 -16
View File
@@ -552,9 +552,37 @@ impl PgDatabase {
}
}
/// Drop a custom instance database: terminate connections, DROP DATABASE, remove from global_settings.
/// Non-fatal variant that logs errors instead of returning them.
pub async fn drop_custom_instance_database(db: &DB, dbname: &str) {
/// 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<()> {
let dbname = dbname.trim();
if dbname.is_empty() {
return Err(error::Error::BadRequest(
"Database name cannot be empty".to_string(),
));
}
let wmill_pg_creds = PgDatabase::parse_uri(&get_database_url().await?.as_str().await)?;
if wmill_pg_creds.dbname.trim().eq_ignore_ascii_case(dbname) {
return Err(error::Error::BadRequest(
"Cannot drop the main Windmill database".to_string(),
));
}
let db_exists = sqlx::query_scalar!(
"SELECT EXISTS (SELECT 1 FROM pg_catalog.pg_database WHERE datname = $1)",
dbname
)
.fetch_one(db)
.await?
.unwrap_or(false);
if !db_exists {
return Err(error::Error::NotFound(format!(
"Database '{}' does not exist",
dbname
)));
}
// Terminate active connections
if let Err(e) = sqlx::query(&format!(
"SELECT pg_terminate_backend(pid) FROM pg_stat_activity WHERE datname = '{}' AND pid <> pg_backend_pid()",
@@ -567,27 +595,24 @@ pub async fn drop_custom_instance_database(db: &DB, dbname: &str) {
}
// Drop the database
match sqlx::query(&format!("DROP DATABASE IF EXISTS \"{}\"", dbname))
sqlx::query(&format!("DROP DATABASE IF EXISTS \"{}\"", dbname))
.execute(db)
.await
{
Ok(_) => tracing::info!("Dropped instance database '{}'", dbname),
Err(e) => {
tracing::error!("Failed to drop instance database '{}': {}", dbname, e);
return;
}
}
.map_err(|e| {
error::Error::internal_err(format!("Failed to drop database '{}': {}", dbname, e))
})?;
tracing::info!("Dropped instance database '{}'", dbname);
// Remove from global_settings
if let Err(e) = sqlx::query!(
sqlx::query!(
r#"UPDATE global_settings SET value = value #- ARRAY['databases', $1] WHERE name = 'custom_instance_pg_databases'"#,
dbname
)
.execute(db)
.await
{
tracing::error!("Failed to remove '{}' from global_settings: {}", dbname, e);
}
.await?;
Ok(())
}
#[derive(Clone)]