From 92c044af6a207d463ec9a66221a652c053979974 Mon Sep 17 00:00:00 2001 From: Diego Imbert Date: Fri, 18 Sep 2026 13:57:18 +0200 Subject: [PATCH] fix(datatables): hold the fork lock across a fork import, and carry the reservation inside the setup write Co-Authored-By: Claude Opus 5 (1M context) --- backend/windmill-api-settings/src/lib.rs | 33 ++++++++++--------- .../windmill-api-workspaces/src/workspaces.rs | 9 +++++ 2 files changed, 26 insertions(+), 16 deletions(-) diff --git a/backend/windmill-api-settings/src/lib.rs b/backend/windmill-api-settings/src/lib.rs index c666c6239d..1c0897d172 100644 --- a/backend/windmill-api-settings/src/lib.rs +++ b/backend/windmill-api-settings/src/lib.rs @@ -1767,16 +1767,6 @@ async fn setup_custom_instance_pg_database( // Before anything is recorded: the status written below replaces the registry entry, and with it // the workspace a fork copy is reserved for. require_super_admin(&db, &authed).await?; - // A re-run keeps the fork reservation: without it, the workspace the copy was made for could no - // longer import into it or finish its fork. - let workspace_id = sqlx::query_scalar::<_, Option>( - "SELECT value->'databases'->$1->>'workspace_id' FROM global_settings - WHERE name = 'custom_instance_pg_databases'", - ) - .bind(&dbname) - .fetch_optional(&db) - .await? - .flatten(); let mut logs = CustomInstanceDbLogs::default(); let result = setup_custom_instance_pg_database_inner(authed, &db, &dbname, &mut logs).await; let success = result.is_ok(); @@ -1787,14 +1777,25 @@ async fn setup_custom_instance_pg_database( error, tag: body.tag, used_by_workspaces: vec![], - workspace_id, + workspace_id: None, }; let status_json = serde_json::to_value(&status).map_err(to_anyhow)?; - // Save that the database was setup successfully - sqlx::query!( - r#"UPDATE global_settings SET value = jsonb_set(value, '{databases}', (COALESCE(value->'databases', '{}'::jsonb) || to_jsonb($1::json))) WHERE name = 'custom_instance_pg_databases'"#, - json!({ dbname: status_json }) - ).execute(&db).await?; + // The fork reservation is carried over inside the write, from whatever the row holds then: a + // rename migrating it while the setup above ran would otherwise be overwritten with the value + // this request started from, stranding the copy under the archived workspace. + let saved = sqlx::query_scalar::<_, serde_json::Value>( + r#"UPDATE global_settings SET value = jsonb_set(value, '{databases}', + COALESCE(value->'databases', '{}'::jsonb) + || jsonb_build_object($1::text, $2::jsonb || jsonb_build_object( + 'workspace_id', value->'databases'->$1::text->'workspace_id'))) + WHERE name = 'custom_instance_pg_databases' + RETURNING value->'databases'->$1::text"#, + ) + .bind(&dbname) + .bind(&status_json) + .fetch_one(&db) + .await?; + let status: CustomInstanceDb = serde_json::from_value(saved).map_err(to_anyhow)?; Ok(Json(status)) } diff --git a/backend/windmill-api-workspaces/src/workspaces.rs b/backend/windmill-api-workspaces/src/workspaces.rs index 322449acac..3d258b761b 100644 --- a/backend/windmill-api-workspaces/src/workspaces.rs +++ b/backend/windmill-api-workspaces/src/workspaces.rs @@ -3585,6 +3585,7 @@ async fn import_pg_database( } let schema_only = req.fork_behavior == DataTableForkBehavior::SchemaOnly; + let mut fork_lock: Option> = None; let source_pg = resolve_pg_source_checked(&db, &user_db, &authed, &w_id, &req.source).await?; let mut target_pg = resolve_pg_source_checked(&db, &user_db, &authed, &w_id, &req.target).await?; @@ -3598,8 +3599,13 @@ async fn import_pg_database( )); } if is_instance_datatable_source(&db, &w_id, &req.target).await? { + // Held until the restore is done, as fork finalization takes it: a fork must not + // commit this database while `psql` is still filling it. + let mut tx = db.begin().await?; + windmill_common::workspaces::lock_fork_datatables(&mut tx, &w_id).await?; windmill_common::ensure_fork_database_available_to(&db, override_dbname, &w_id) .await?; + fork_lock = Some(tx); } } target_pg.dbname = override_dbname.clone(); @@ -3619,6 +3625,9 @@ async fn import_pg_database( ) .await?; pg_import_dump(&target_pg, &dump_file).await?; + if let Some(tx) = fork_lock { + tx.commit().await?; + } Ok(format!( "Imported from '{}' into '{}'",