refactor: rename fork_pg_database to import_pg_database with source/target/override params

Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>
This commit is contained in:
Diego Imbert
2026-03-27 14:14:07 +01:00
parent 6374ce84f4
commit b090f87bdb
4 changed files with 85 additions and 75 deletions

View File

@@ -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<String>,
#[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<DB>,
Path(w_id): Path<String>,
Json(req): Json<ForkPgDatabaseRequest>,
Json(req): Json<ImportPgDatabaseRequest>,
) -> Result<String> {
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
))
}

View File

@@ -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:

View File

@@ -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
}
})

View File

@@ -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
}
})