From b090f87bdbf382355b44a74af2e2d282a97fa9f9 Mon Sep 17 00:00:00 2001 From: Diego Imbert Date: Fri, 27 Mar 2026 14:14:07 +0100 Subject: [PATCH] refactor: rename fork_pg_database to import_pg_database with source/target/override params Co-Authored-By: Claude Opus 4.5 --- .../windmill-api-workspaces/src/workspaces.rs | 130 +++++++++--------- backend/windmill-api/openapi.yaml | 20 ++- .../src/lib/components/DBManagerDrawer.svelte | 4 +- .../ForkDatatableSection.svelte | 6 +- 4 files changed, 85 insertions(+), 75 deletions(-) diff --git a/backend/windmill-api-workspaces/src/workspaces.rs b/backend/windmill-api-workspaces/src/workspaces.rs index c440617483..ff44398958 100644 --- a/backend/windmill-api-workspaces/src/workspaces.rs +++ b/backend/windmill-api-workspaces/src/workspaces.rs @@ -151,7 +151,7 @@ pub fn workspaced_service() -> Router { post(reset_workspace_diffs), ) .route("/compare/:target_workspace_id", get(compare_workspaces)) - .route("/fork_pg_database", post(fork_pg_database)) + .route("/import_pg_database", post(import_pg_database)) .route("/export_pg_schema", post(export_pg_schema)) .route( "/get_datatable_full_schema", @@ -1480,17 +1480,21 @@ async fn pg_import_dump(target_db: &PgDatabase, dump_file: &DumpFile) -> Result< } #[derive(Deserialize)] -struct ForkPgDatabaseRequest { +struct ImportPgDatabaseRequest { source: String, - target_dbname: String, + target: String, + #[serde(default)] + target_dbname_override: Option, + #[serde(default)] + create_target_db: bool, fork_behavior: DataTableForkBehavior, } -async fn fork_pg_database( +async fn import_pg_database( ApiAuthed { is_admin, username, .. }: ApiAuthed, Extension(db): Extension, Path(w_id): Path, - Json(req): Json, + Json(req): Json, ) -> Result { if req.fork_behavior == DataTableForkBehavior::KeepOriginal { return Ok("No action needed for KeepOriginal behavior".to_string()); @@ -1500,75 +1504,73 @@ async fn fork_pg_database( require_admin(is_admin, &username)?; if *CLOUD_HOSTED { return Err(Error::BadRequest( - "Forking schema and data is not available on cloud".to_string(), + "Importing schema and data is not available on cloud".to_string(), )); } } let schema_only = req.fork_behavior == DataTableForkBehavior::SchemaOnly; let source_pg = resolve_pg_source(&db, &w_id, &req.source).await?; - let target_dbname = &req.target_dbname; + let mut target_pg = resolve_pg_source(&db, &w_id, &req.target).await?; - if req.source.starts_with("datatable://") { - // Instance datatable: create custom instance DB, then dump/import - windmill_common::create_custom_instance_database(&db, target_dbname, "datatable").await?; - - let target_pg = PgDatabase { - dbname: target_dbname.clone(), - ..PgDatabase::parse_uri(&get_database_url().await?.as_str().await)? - }; - let dump_file = pg_dump_database(&source_pg, schema_only).await?; - pg_import_dump(&target_pg, &dump_file).await?; - } else { - // Resource datatable: connect to the source server, create DB there, then dump/import - let (client, connection) = source_pg.connect().await?; - let join_handle = tokio::spawn(async move { connection.await }); - - let row = client - .query_one( - "SELECT EXISTS (SELECT 1 FROM pg_catalog.pg_database WHERE datname = $1)", - &[target_dbname], - ) - .await - .map_err(|e| { - Error::internal_err(format!("Failed to check database existence: {}", e)) - })?; - let db_exists: bool = row.get(0); - - if db_exists { - drop(client); - let _ = join_handle.await; - return Err(Error::BadRequest(format!( - "Database '{}' already exists on the resource server", - target_dbname - ))); - } - - client - .execute(&format!("CREATE DATABASE \"{}\"", target_dbname), &[]) - .await - .map_err(|e| { - Error::internal_err(format!( - "Failed to create database '{}': {}", - target_dbname, e - )) - })?; - - drop(client); - join_handle - .await - .map_err(|e| Error::internal_err(format!("join error: {}", e)))? - .map_err(|e| Error::internal_err(format!("tokio_postgres error: {}", e)))?; - - // Dump source → import into new DB on same server - let target_pg = PgDatabase { dbname: target_dbname.clone(), ..source_pg.clone() }; - let dump_file = pg_dump_database(&source_pg, schema_only).await?; - pg_import_dump(&target_pg, &dump_file).await?; + if let Some(ref override_dbname) = req.target_dbname_override { + target_pg.dbname = override_dbname.clone(); } + if req.create_target_db { + if req.target.starts_with("datatable://") { + // Instance: use shared create function + windmill_common::create_custom_instance_database(&db, &target_pg.dbname, "datatable") + .await?; + } else { + // Resource: CREATE DATABASE on the source server + let (client, connection) = source_pg.connect().await?; + let join_handle = tokio::spawn(async move { connection.await }); + + let row = client + .query_one( + "SELECT EXISTS (SELECT 1 FROM pg_catalog.pg_database WHERE datname = $1)", + &[&target_pg.dbname], + ) + .await + .map_err(|e| { + Error::internal_err(format!("Failed to check database existence: {}", e)) + })?; + let db_exists: bool = row.get(0); + + if db_exists { + drop(client); + let _ = join_handle.await; + return Err(Error::BadRequest(format!( + "Database '{}' already exists on the resource server", + target_pg.dbname + ))); + } + + client + .execute(&format!("CREATE DATABASE \"{}\"", &target_pg.dbname), &[]) + .await + .map_err(|e| { + Error::internal_err(format!( + "Failed to create database '{}': {}", + target_pg.dbname, e + )) + })?; + + drop(client); + join_handle + .await + .map_err(|e| Error::internal_err(format!("join error: {}", e)))? + .map_err(|e| Error::internal_err(format!("tokio_postgres error: {}", e)))?; + } + } + + let dump_file = pg_dump_database(&source_pg, schema_only).await?; + pg_import_dump(&target_pg, &dump_file).await?; + Ok(format!( - "Forked '{}' into new database '{}'", - req.source, target_dbname + "Imported from '{}' into '{}'", + req.source, target_pg.dbname )) } diff --git a/backend/windmill-api/openapi.yaml b/backend/windmill-api/openapi.yaml index deb4b4b08c..f7fa19a73d 100644 --- a/backend/windmill-api/openapi.yaml +++ b/backend/windmill-api/openapi.yaml @@ -3432,29 +3432,35 @@ paths: application/json: schema: {} - /w/{workspace}/workspaces/fork_pg_database: + /w/{workspace}/workspaces/import_pg_database: post: - summary: fork a PostgreSQL database from source to target - operationId: forkPgDatabase + summary: import a PostgreSQL database from source to target + operationId: importPgDatabase tags: - workspace parameters: - $ref: "#/components/parameters/WorkspaceId" requestBody: - description: Fork pg database request + description: Import pg database request required: true content: application/json: schema: type: object - required: [source, target_dbname, fork_behavior] + required: [source, target, fork_behavior] properties: source: type: string description: "Source database: 'datatable://name' or '$res:path'" - target_dbname: + target: type: string - description: "Name for the new database to create" + description: "Target database: 'datatable://name' or '$res:path'" + target_dbname_override: + type: string + description: "Override the target database name" + create_target_db: + type: boolean + description: "When true, CREATE DATABASE is run before importing" fork_behavior: type: string enum: diff --git a/frontend/src/lib/components/DBManagerDrawer.svelte b/frontend/src/lib/components/DBManagerDrawer.svelte index fd7d0a3dd5..1cbb996ff3 100644 --- a/frontend/src/lib/components/DBManagerDrawer.svelte +++ b/frontend/src/lib/components/DBManagerDrawer.svelte @@ -125,11 +125,11 @@ if (!target) return importLoading = true try { - await WorkspaceService.forkPgDatabase({ + await WorkspaceService.importPgDatabase({ workspace: $workspaceStore, requestBody: { source: toSourceIdentifier(importSource), - target_dbname: target.replace('datatable://', '').replace('$res:', ''), + target, fork_behavior: importBehavior } }) diff --git a/frontend/src/lib/components/workspaceSettings/ForkDatatableSection.svelte b/frontend/src/lib/components/workspaceSettings/ForkDatatableSection.svelte index 5b7c8cf1f3..d90297d04e 100644 --- a/frontend/src/lib/components/workspaceSettings/ForkDatatableSection.svelte +++ b/frontend/src/lib/components/workspaceSettings/ForkDatatableSection.svelte @@ -110,11 +110,13 @@ job.steps[stepIdx].status = 'running' try { - await WorkspaceService.forkPgDatabase({ + await WorkspaceService.importPgDatabase({ workspace: job._sourceWorkspace, requestBody: { source: `datatable://${job.name}`, - target_dbname: job._newDbName, + target: `datatable://${job.name}`, + target_dbname_override: job._newDbName, + create_target_db: true, fork_behavior: job.behavior } })