refactor: remove datatable fork logic from workspace fork route

Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>
This commit is contained in:
Diego Imbert
2026-03-23 11:55:30 +01:00
parent 50587fc954
commit ffbf3e1fce
3 changed files with 4 additions and 402 deletions

View File

@@ -41,18 +41,17 @@ use windmill_common::workspaces::GitRepositorySettings;
use windmill_common::workspaces::WorkspaceDeploymentUISettings;
use windmill_common::workspaces::{
check_user_against_rule, get_datatable_resource_from_db_unchecked, transform_json_unchecked,
DataTable, DataTableCatalogResourceType, DataTableDatabase, DataTableForkBehavior,
ProtectionRuleKind, ProtectionRules, ProtectionRuleset, RuleCheckResult,
WorkspaceGitSyncSettings,
DataTable, DataTableCatalogResourceType, DataTableForkBehavior, ProtectionRuleKind,
ProtectionRules, ProtectionRuleset, RuleCheckResult, WorkspaceGitSyncSettings,
};
use windmill_common::workspaces::{Ducklake, DucklakeCatalogResourceType};
use windmill_common::PgDatabase;
use windmill_common::{
error::{Error, JsonResult, Result},
global_settings::AUTOMATE_USERNAME_CREATION_SETTING,
oauth2::WORKSPACE_SLACK_BOT_TOKEN_PATH,
utils::{paginate, rd_string, require_admin, Pagination},
};
use windmill_common::{get_database_url, PgDatabase};
use windmill_dep_map::scoped_dependency_map::{
DependencyDependent, DependencyMap, ScopedDependencyMap,
};
@@ -346,8 +345,6 @@ struct CreateWorkspaceFork {
id: String,
name: String,
color: Option<String>,
#[serde(default)]
datatable_behaviors: Option<HashMap<String, DataTableForkBehavior>>,
}
#[derive(Deserialize)]
@@ -1512,324 +1509,6 @@ async fn export_pg_schema(
pg_dump_database(&pg, true).await
}
/// Core logic for forking a single datatable: creates a new instance DB,
/// dumps the source schema (and optionally data), imports into the new DB,
/// and updates the target workspace's datatable config.
async fn fork_datatable(
db: &DB,
w_id: &str,
source_datatable_name: &str,
new_datatable_name: &str,
new_custom_instance_database_name: &str,
include_data: bool,
) -> Result<String> {
// Resolve the source datatable to get the original instance DB name
let pg_db =
resolve_pg_source(db, w_id, &format!("datatable://{}", source_datatable_name)).await?;
// Interpolate $current_name with the original database name
let original_dbname = &pg_db.dbname;
let new_dbname = new_custom_instance_database_name.replace("$current_name", original_dbname);
// Create the new custom instance database
let wmill_pg_creds = PgDatabase::parse_uri(&get_database_url().await?.as_str().await)?;
let new_pg_creds = PgDatabase { dbname: new_dbname.clone(), ..wmill_pg_creds };
let db_exists = sqlx::query_scalar!(
"SELECT EXISTS (SELECT 1 FROM pg_catalog.pg_database WHERE datname = $1)",
&new_dbname
)
.fetch_one(db)
.await?
.unwrap_or(false);
if !db_exists {
sqlx::query(&format!("CREATE DATABASE \"{}\"", &new_dbname))
.execute(db)
.await?;
}
// Grant permissions to custom_instance_user on the new database
let (client, connection) = new_pg_creds.connect().await?;
let join_handle = tokio::spawn(async move { connection.await });
if let Err(e) = client
.batch_execute(&format!(
"GRANT CONNECT ON DATABASE \"{new_dbname}\" TO custom_instance_user;
GRANT USAGE ON SCHEMA public TO custom_instance_user;
GRANT CREATE ON SCHEMA public TO custom_instance_user;
GRANT CREATE ON DATABASE \"{new_dbname}\" TO custom_instance_user;
ALTER DEFAULT PRIVILEGES IN SCHEMA public
GRANT SELECT, INSERT, UPDATE, DELETE ON TABLES TO custom_instance_user;
ALTER ROLE custom_instance_user CREATEROLE;
ALTER ROLE custom_instance_user REPLICATION;"
))
.await
{
tracing::warn!(
"Failed to grant permissions to custom_instance_user on '{}': {}. Continuing with fork.",
new_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)))?;
// Register the new database in global_settings
let status_json = serde_json::json!({
"logs": {
"super_admin": "OK",
"database_credentials": "OK",
"valid_dbname": "OK",
"created_database": if db_exists { "SKIP" } else { "OK" },
"db_connect": "OK",
"grant_permissions": "OK"
},
"success": true,
"error": null,
"tag": "datatable"
});
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'"#,
serde_json::json!({ &new_dbname: status_json })
)
.execute(db)
.await?;
// Export the schema (and optionally data) from the source BEFORE updating the config
let dump = pg_dump_database(&pg_db, !include_data).await?;
// Update the forked workspace's datatable config to point to the new database
let new_datatable = DataTable {
database: DataTableDatabase {
resource_type: DataTableCatalogResourceType::Instance,
resource_path: new_dbname.clone(),
},
non_diffable: true,
forked_from: Some(original_dbname.clone()),
};
let datatable_value: serde_json::Value =
serde_json::to_value(&new_datatable).map_err(|err| Error::internal_err(err.to_string()))?;
sqlx::query!(
r#"UPDATE workspace_settings
SET datatable = jsonb_set(
COALESCE(datatable, '{"datatables":{}}'::jsonb),
ARRAY['datatables', $1],
$2
)
WHERE workspace_id = $3"#,
new_datatable_name,
datatable_value,
w_id
)
.execute(db)
.await?;
// Import the dumped schema into the new database
pg_import_dump(&new_pg_creds, &dump).await?;
Ok(format!(
"Forked datatable '{}' as '{}' with new database '{}'",
source_datatable_name, new_datatable_name, new_dbname
))
}
/// Fork a resource-type datatable by creating a new database on the same server,
/// dumping/importing data, and updating the resource in the forked workspace.
async fn fork_resource_datatable(
db: &DB,
source_workspace_id: &str,
target_workspace_id: &str,
datatable_name: &str,
resource_path: &str,
include_data: bool,
) -> Result<()> {
// Resolve the source PG credentials from the resource
let source_pg =
resolve_pg_source(db, source_workspace_id, &format!("$res:{}", resource_path)).await?;
let original_dbname = &source_pg.dbname;
// Generate the new database name
let new_dbname = format!(
"{}__{}_{}",
target_workspace_id.replace('-', "_"),
datatable_name,
original_dbname
);
// Connect to the server (using the source credentials but targeting the default/postgres db)
// to create the new database
let admin_pg = PgDatabase { dbname: "postgres".to_string(), ..source_pg.clone() };
let (client, connection) = admin_pg.connect().await?;
let join_handle = tokio::spawn(async move { connection.await });
// Check if the database already exists
let row = client
.query_one(
"SELECT EXISTS (SELECT 1 FROM pg_catalog.pg_database WHERE datname = $1)",
&[&new_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 {
client
.execute(&format!("CREATE DATABASE \"{}\"", &new_dbname), &[])
.await
.map_err(|e| {
Error::internal_err(format!("Failed to create database '{}': {}", new_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)))?;
// Export the schema (and optionally data) from the source
let dump = pg_dump_database(&source_pg, !include_data).await?;
// Import the dump into the new database
let new_pg = PgDatabase { dbname: new_dbname.clone(), ..source_pg.clone() };
pg_import_dump(&new_pg, &dump).await?;
// Update the resource in the forked workspace: set dbname and non_diffable
sqlx::query!(
r#"UPDATE resource
SET value = jsonb_set(COALESCE(value, '{}'::jsonb), '{dbname}', to_jsonb($1::text)),
non_diffable = true
WHERE workspace_id = $2 AND path = $3"#,
new_dbname,
target_workspace_id,
resource_path
)
.execute(db)
.await?;
// Tag the datatable config as non_diffable with forked_from info
let forked_from_value = format!("$res:{}", resource_path);
let tag_json = serde_json::json!({
"nonDiffable": true,
"forkedFrom": forked_from_value
});
sqlx::query!(
r#"UPDATE workspace_settings
SET datatable = jsonb_set(
COALESCE(datatable, '{"datatables":{}}'::jsonb),
ARRAY['datatables', $1],
COALESCE(datatable->'datatables'->$1, '{}'::jsonb) || $2
)
WHERE workspace_id = $3"#,
datatable_name,
tag_json,
target_workspace_id
)
.execute(db)
.await?;
tracing::info!(
"Forked resource datatable '{}': created database '{}' from '{}'",
datatable_name,
new_dbname,
original_dbname
);
Ok(())
}
/// Fork all datatables from a source workspace into a target workspace.
/// Uses per-datatable behaviors from the provided map, defaulting to KeepOriginal.
async fn fork_all_datatables(
db: &DB,
source_workspace_id: &str,
target_workspace_id: &str,
datatable_behaviors: &Option<HashMap<String, DataTableForkBehavior>>,
) -> Result<()> {
if !target_workspace_id.starts_with("wm-fork") {
return Err(Error::BadRequest(
"Target workspace is not a fork".to_string(),
));
}
let datatable_config = sqlx::query_scalar!(
"SELECT datatable FROM workspace_settings WHERE workspace_id = $1",
source_workspace_id
)
.fetch_one(db)
.await?;
let datatables: HashMap<String, DataTable> = match datatable_config {
Some(config) => {
let settings: DataTableSettings = serde_json::from_value(config)
.unwrap_or(DataTableSettings { datatables: HashMap::new() });
settings.datatables
}
None => return Ok(()),
};
for (name, dt) in &datatables {
let behavior = datatable_behaviors
.as_ref()
.and_then(|m| m.get(name).copied())
.unwrap_or(DataTableForkBehavior::KeepOriginal);
if behavior == DataTableForkBehavior::KeepOriginal {
continue;
}
let include_data = behavior == DataTableForkBehavior::SchemaAndData && !*CLOUD_HOSTED;
if dt.database.resource_type == DataTableCatalogResourceType::Instance {
let new_db_name = format!("{}__$current_name", target_workspace_id.replace('-', "_"));
if let Err(e) = fork_datatable(
db,
target_workspace_id,
name,
name,
&new_db_name,
include_data,
)
.await
{
tracing::error!(
"Failed to fork instance datatable '{}' from '{}' to '{}': {}",
name,
source_workspace_id,
target_workspace_id,
e
);
}
} else {
// Resource-type datatable (Postgresql)
if let Err(e) = fork_resource_datatable(
db,
source_workspace_id,
target_workspace_id,
name,
&dt.database.resource_path,
include_data,
)
.await
{
tracing::error!(
"Failed to fork resource datatable '{}' from '{}' to '{}': {}",
name,
source_workspace_id,
target_workspace_id,
e
);
}
}
}
Ok(())
}
async fn edit_ducklake_config(
authed: ApiAuthed,
Extension(db): Extension<DB>,
@@ -4037,15 +3716,6 @@ async fn create_workspace_fork(
.await?;
tx.commit().await?;
// Fork datatables after the transaction commits (creates external databases)
fork_all_datatables(
&db,
&parent_workspace_id,
&forked_id,
&nw.datatable_behaviors,
)
.await?;
Ok(format!("Created forked workspace {}", &forked_id))
}

View File

@@ -22890,15 +22890,6 @@ components:
type: string
color:
type: string
datatable_behaviors:
type: object
description: "Per-datatable fork behavior. Keys are datatable names, values are fork behaviors."
additionalProperties:
type: string
enum:
- schema_only
- schema_and_data
- keep_original
required:
- id
- name

View File

@@ -32,9 +32,7 @@
import { jobManager } from '$lib/services/JobManager'
import Alert from '../common/alert/Alert.svelte'
import { base } from '$lib/base'
import { resource } from 'runed'
import Label from '../Label.svelte'
import Select from '../select/Select.svelte'
interface Props {
isFork?: boolean
@@ -53,16 +51,6 @@
let codeCompletionEnabled = $state(true)
let checking = $state(false)
let allDatatables = resource([], async () =>
$workspaceStore
? WorkspaceService.listDataTables({
workspace: $workspaceStore ?? ''
})
: undefined
)
let datatableBehaviors: Record<string, 'schema_only' | 'schema_and_data' | 'keep_original'> =
$state({})
let workspaceColor: string | undefined = $state(undefined)
let colorEnabled = $state(false)
@@ -189,20 +177,13 @@
return
}
// Build datatable behaviors map, defaulting to keep_original
const behaviors: Record<string, 'schema_only' | 'schema_and_data' | 'keep_original'> = {}
for (const dt of allDatatables.current ?? []) {
behaviors[dt] = datatableBehaviors[dt] ?? 'keep_original'
}
try {
await WorkspaceService.createWorkspaceFork({
workspace: $workspaceStore!,
requestBody: {
id: prefixed_id,
name,
color: colorEnabled && workspaceColor ? workspaceColor : undefined,
datatable_behaviors: behaviors
color: colorEnabled && workspaceColor ? workspaceColor : undefined
}
})
} catch (e) {
@@ -486,46 +467,6 @@
{/if}
</div>
</Label>
{#if isFork && allDatatables.current?.length}
<Label label="Data table behavior">
<span class="text-xs text-secondary">
Fork datatables to avoid changing production data while working
</span>
<div class="border rounded-md divide-y">
{#each allDatatables.current ?? [] as dt}
<div class="flex items-center gap-2 justify-between px-4 py-1.5">
<span class="text-xs">{dt}</span>
<Select
dropdownClass="max-w-96"
bind:value={
() => datatableBehaviors[dt] ?? 'keep_original',
(v) => (datatableBehaviors[dt] = v)
}
items={[
{
value: 'schema_only',
label: 'Clone schema only',
subtitle: 'Create a fork of the datatable with the same schema but no data'
},
{
value: 'schema_and_data',
label: 'Clone schema and data',
subtitle:
'Copy the database and all rows from every table. This may take a long time and use significant storage depending on the size of the source.'
},
{
value: 'keep_original',
label: 'Keep original',
subtitle:
'Share the same datatable. Edits to this datatable will also happen on the original workspace.'
}
]}
/>
</div>
{/each}
</div>
</Label>
{/if}
{#if !automateUsernameCreation}
<Label label="Your username in that workspace">
<TextInput