From ffbf3e1fcee4edcfe6746d3ca37679fd94b4909e Mon Sep 17 00:00:00 2001 From: Diego Imbert Date: Mon, 23 Mar 2026 11:55:30 +0100 Subject: [PATCH] refactor: remove datatable fork logic from workspace fork route Co-Authored-By: Claude Opus 4.5 --- .../windmill-api-workspaces/src/workspaces.rs | 336 +----------------- backend/windmill-api/openapi.yaml | 9 - .../CreateWorkspaceInner.svelte | 61 +--- 3 files changed, 4 insertions(+), 402 deletions(-) diff --git a/backend/windmill-api-workspaces/src/workspaces.rs b/backend/windmill-api-workspaces/src/workspaces.rs index 53ac0bf082..28c5d13cc0 100644 --- a/backend/windmill-api-workspaces/src/workspaces.rs +++ b/backend/windmill-api-workspaces/src/workspaces.rs @@ -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, - #[serde(default)] - datatable_behaviors: Option>, } #[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 { - // 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>, -) -> 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 = 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, @@ -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)) } diff --git a/backend/windmill-api/openapi.yaml b/backend/windmill-api/openapi.yaml index 959446b778..e19ff12b17 100644 --- a/backend/windmill-api/openapi.yaml +++ b/backend/windmill-api/openapi.yaml @@ -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 diff --git a/frontend/src/lib/components/workspaceSettings/CreateWorkspaceInner.svelte b/frontend/src/lib/components/workspaceSettings/CreateWorkspaceInner.svelte index 2144bbf713..2e53a08baa 100644 --- a/frontend/src/lib/components/workspaceSettings/CreateWorkspaceInner.svelte +++ b/frontend/src/lib/components/workspaceSettings/CreateWorkspaceInner.svelte @@ -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 = - $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 = {} - 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} - {#if isFork && allDatatables.current?.length} -