From a351a36e9b35b53f8dbbbb33d6e8e2ac3a3b9144 Mon Sep 17 00:00:00 2001 From: Danny White <3104761+dnywh@users.noreply.github.com> Date: Mon, 7 Sep 2026 17:02:09 +1000 Subject: [PATCH] refactor(studio): centralise replication payload builders (#49842) ## What kind of change does this PR introduce? Studio data-layer refactor. ## What is the current behavior? Pipeline creation, editing, and validation build similar destination and pipeline payloads separately. The duplicated mappings rely on type assertions and can drift between actions. ## What is the new behavior? Uses shared typed builders for create, update, and validation payloads across the existing destinations. Update payloads continue to omit blank secrets, while create payloads preserve their current values. This PR does not add table partitioning configuration. ## To test This is a data-layer refactor. No visible behaviour should change. 1. Open **Database > Replication** and click **Start a new pipeline**. 2. Select **BigQuery**, or any other enabled destination. 3. Edit a few non-secret fields and expand **Advanced settings**. 4. Confirm the form remains usable and no runtime errors appear. Create, update, validation, and secret-handling behaviour is covered by the focused tests and CI. Deploy previews and fresh local projects do not have the existing destinations or credentials needed to exercise those paths manually. ## Summary by CodeRabbit * **Bug Fixes** * Improved replication destination configuration handling during creation, updates, and validation. * Applied consistent configuration mapping across supported destination types. * Ensured blank secret values are omitted during updates while retained when creating destinations. * Standardized table synchronization defaults when no specific setting is provided. * **Tests** * Added coverage for BigQuery configuration mapping and secret handling. * Updated DuckLake tests for destination updates. --------- Co-authored-by: Joshen Lim --- .../Replication/BatchRestartDialog.tsx | 3 +- .../DestinationForm/DestinationForm.utils.ts | 10 +- .../DestinationForm/PipelineCostDialog.tsx | 2 +- .../DestinationForm/useDestinationForm.ts | 6 +- .../Replication/RestartTableDialog.tsx | 7 +- .../Replication/TableSyncCopy.utils.test.ts | 2 +- .../Replication/TableSyncCopy.utils.ts | 6 +- ...eate-destination-pipeline-mutation.test.ts | 61 +++- .../create-destination-pipeline-mutation.ts | 298 ++++++------------ apps/studio/data/replication/types.ts | 95 ++++++ .../update-destination-pipeline-mutation.ts | 227 ++++++++----- apps/studio/data/replication/utils.ts | 34 ++ .../validate-destination-mutation.ts | 114 +------ .../replication/validate-pipeline-mutation.ts | 25 +- 14 files changed, 449 insertions(+), 441 deletions(-) create mode 100644 apps/studio/data/replication/types.ts diff --git a/apps/studio/components/interfaces/Database/Replication/BatchRestartDialog.tsx b/apps/studio/components/interfaces/Database/Replication/BatchRestartDialog.tsx index 5c9b209f0fe..1c9948cb6e0 100644 --- a/apps/studio/components/interfaces/Database/Replication/BatchRestartDialog.tsx +++ b/apps/studio/components/interfaces/Database/Replication/BatchRestartDialog.tsx @@ -14,9 +14,10 @@ import { import { PipelineStatusName } from './Replication.constants' import { RestartCostEstimate } from './RestartCostEstimate' -import { getTableCopyTargets, type TableSyncCopyConfig } from './TableSyncCopy.utils' +import { getTableCopyTargets } from './TableSyncCopy.utils' import { ReplicationPipelineTableStatus } from '@/data/replication/pipeline-replication-status-query' import { useRollbackTablesMutation } from '@/data/replication/rollback-tables-mutation' +import type { TableSyncCopyConfig } from '@/data/replication/types' interface BatchRestartDialogProps { open: boolean diff --git a/apps/studio/components/interfaces/Database/Replication/DestinationPanel/DestinationForm/DestinationForm.utils.ts b/apps/studio/components/interfaces/Database/Replication/DestinationPanel/DestinationForm/DestinationForm.utils.ts index d36e4287aab..3ed2cb9a00b 100644 --- a/apps/studio/components/interfaces/Database/Replication/DestinationPanel/DestinationForm/DestinationForm.utils.ts +++ b/apps/studio/components/interfaces/Database/Replication/DestinationPanel/DestinationForm/DestinationForm.utils.ts @@ -24,7 +24,10 @@ import { } from './DuckLake/DuckLake.constants' import { type DucklakeApiConfig } from './DuckLake/DuckLake.utils' import { type SnowflakeApiConfig } from './Snowflake/Snowflake.utils' -import { +import { type ReplicationDestinationByIdData } from '@/data/replication/destination-by-id-query' +import { type ReplicationPipelineByIdData } from '@/data/replication/pipeline-by-id-query' +import { type ReplicationPublication } from '@/data/replication/publications-query' +import type { BatchConfig, BigQueryDestinationConfig, ClickHouseDestinationConfig, @@ -35,10 +38,7 @@ import { IcebergDestinationConfig, SnowflakeDestinationConfig, TableSyncCopyConfig, -} from '@/data/replication/create-destination-pipeline-mutation' -import { type ReplicationDestinationByIdData } from '@/data/replication/destination-by-id-query' -import { type ReplicationPipelineByIdData } from '@/data/replication/pipeline-by-id-query' -import { type ReplicationPublication } from '@/data/replication/publications-query' +} from '@/data/replication/types' import { type ValidationFailure } from '@/data/replication/validate-destination-mutation' import { type CreateS3AccessKeyCredentialVariables, diff --git a/apps/studio/components/interfaces/Database/Replication/DestinationPanel/DestinationForm/PipelineCostDialog.tsx b/apps/studio/components/interfaces/Database/Replication/DestinationPanel/DestinationForm/PipelineCostDialog.tsx index b472551c66f..199b1d6ee85 100644 --- a/apps/studio/components/interfaces/Database/Replication/DestinationPanel/DestinationForm/PipelineCostDialog.tsx +++ b/apps/studio/components/interfaces/Database/Replication/DestinationPanel/DestinationForm/PipelineCostDialog.tsx @@ -25,11 +25,11 @@ import { getTableCopyTargets, summarizeTableCopyEstimate, type ReplicationTableIdentity, - type TableSyncCopyConfig, } from '@/components/interfaces/Database/Replication/TableSyncCopy.utils' import { InlineLink } from '@/components/ui/InlineLink' import { useReplicationCostEstimateQuery } from '@/data/replication/cost-estimate-query' import { useReplicationSourceId } from '@/data/replication/sources-query' +import type { TableSyncCopyConfig } from '@/data/replication/types' import { useLatest } from '@/hooks/misc/useLatest' import { DOCS_URL } from '@/lib/constants' import { formatBytes, formatCurrency } from '@/lib/helpers' diff --git a/apps/studio/components/interfaces/Database/Replication/DestinationPanel/DestinationForm/useDestinationForm.ts b/apps/studio/components/interfaces/Database/Replication/DestinationPanel/DestinationForm/useDestinationForm.ts index a89977494cc..5fca795b9a5 100644 --- a/apps/studio/components/interfaces/Database/Replication/DestinationPanel/DestinationForm/useDestinationForm.ts +++ b/apps/studio/components/interfaces/Database/Replication/DestinationPanel/DestinationForm/useDestinationForm.ts @@ -12,13 +12,11 @@ import { buildDestinationConfigForValidation, buildTableSyncCopyConfig, } from './DestinationForm.utils' -import { - useCreateDestinationPipelineMutation, - type BatchConfig, -} from '@/data/replication/create-destination-pipeline-mutation' +import { useCreateDestinationPipelineMutation } from '@/data/replication/create-destination-pipeline-mutation' import type { ReplicationPipelineByIdData } from '@/data/replication/pipeline-by-id-query' import { useReplicationSourcesQuery } from '@/data/replication/sources-query' import { useStartPipelineMutation } from '@/data/replication/start-pipeline-mutation' +import { type BatchConfig } from '@/data/replication/types' import { useUpdateDestinationPipelineMutation } from '@/data/replication/update-destination-pipeline-mutation' import { useValidateDestinationMutation, diff --git a/apps/studio/components/interfaces/Database/Replication/RestartTableDialog.tsx b/apps/studio/components/interfaces/Database/Replication/RestartTableDialog.tsx index 59ce94d5f69..771e6eb034c 100644 --- a/apps/studio/components/interfaces/Database/Replication/RestartTableDialog.tsx +++ b/apps/studio/components/interfaces/Database/Replication/RestartTableDialog.tsx @@ -13,12 +13,9 @@ import { import { PipelineStatusName } from './Replication.constants' import { RestartCostEstimate } from './RestartCostEstimate' -import { - shouldCopyTable, - type ReplicationTableIdentity, - type TableSyncCopyConfig, -} from './TableSyncCopy.utils' +import { shouldCopyTable, type ReplicationTableIdentity } from './TableSyncCopy.utils' import { useRollbackTablesMutation } from '@/data/replication/rollback-tables-mutation' +import type { TableSyncCopyConfig } from '@/data/replication/types' interface RestartTableDialogProps { open: boolean diff --git a/apps/studio/components/interfaces/Database/Replication/TableSyncCopy.utils.test.ts b/apps/studio/components/interfaces/Database/Replication/TableSyncCopy.utils.test.ts index 023afa18441..83f0748fd1a 100644 --- a/apps/studio/components/interfaces/Database/Replication/TableSyncCopy.utils.test.ts +++ b/apps/studio/components/interfaces/Database/Replication/TableSyncCopy.utils.test.ts @@ -4,8 +4,8 @@ import { getTableCopyTargets, shouldCopyTable, summarizeTableCopyEstimate, - type TableSyncCopyConfig, } from './TableSyncCopy.utils' +import type { TableSyncCopyConfig } from '@/data/replication/types' const tables = [ { id: 101, schema: 'public', name: 'orders' }, diff --git a/apps/studio/components/interfaces/Database/Replication/TableSyncCopy.utils.ts b/apps/studio/components/interfaces/Database/Replication/TableSyncCopy.utils.ts index 71fe5b9f0e1..6e004d88238 100644 --- a/apps/studio/components/interfaces/Database/Replication/TableSyncCopy.utils.ts +++ b/apps/studio/components/interfaces/Database/Replication/TableSyncCopy.utils.ts @@ -1,8 +1,4 @@ -export type TableSyncCopyConfig = - | { type: 'include_all_tables' } - | { type: 'skip_all_tables' } - | { type: 'include_tables'; table_ids: number[] } - | { type: 'skip_tables'; table_ids: number[] } +import type { TableSyncCopyConfig } from '@/data/replication/types' export type ReplicationTableIdentity = { id: number diff --git a/apps/studio/data/replication/create-destination-pipeline-mutation.test.ts b/apps/studio/data/replication/create-destination-pipeline-mutation.test.ts index 5264a361700..d112d7cc56a 100644 --- a/apps/studio/data/replication/create-destination-pipeline-mutation.test.ts +++ b/apps/studio/data/replication/create-destination-pipeline-mutation.test.ts @@ -1,9 +1,14 @@ import { describe, expect, it } from 'vitest' import { + buildBigQueryApiConfig, buildDucklakeApiConfig, - buildPipelineApiConfig, } from './create-destination-pipeline-mutation' +import { + buildBigQueryUpdateApiConfig, + buildDucklakeUpdateApiConfig, +} from './update-destination-pipeline-mutation' +import { buildPipelineApiConfig } from './utils' describe('buildPipelineApiConfig', () => { it('maps selective initial-copy configuration to the ETL API shape', () => { @@ -27,6 +32,33 @@ describe('buildPipelineApiConfig', () => { }) }) +describe('buildBigQueryApiConfig', () => { + const baseConfig = { + projectId: 'my-project', + datasetId: 'analytics', + serviceAccountKey: '{}', + } + + it('maps the destination config to the API shape', () => { + expect(buildBigQueryApiConfig(baseConfig)).toEqual({ + big_query: { + project_id: 'my-project', + dataset_id: 'analytics', + service_account_key: '{}', + connection_pool_size: undefined, + max_staleness_mins: undefined, + }, + }) + }) + + it('omits blank service_account_key on update, but not on create', () => { + const config = { ...baseConfig, serviceAccountKey: '' } + + expect(buildBigQueryApiConfig(config).big_query.service_account_key).toBe('') + expect(buildBigQueryUpdateApiConfig(config).big_query.service_account_key).toBeUndefined() + }) +}) + describe('buildDucklakeApiConfig', () => { it('maps a "Use Supabase" config with catalog-level pool size + metadata schema', () => { expect( @@ -106,21 +138,18 @@ describe('buildDucklakeApiConfig', () => { it('omits blank custom secret fields when requested', () => { expect( - buildDucklakeApiConfig( - { - catalogUrl: ' ', - dataPath: 's3://bucket/path', - poolSize: 4, - s3AccessKeyId: '', - s3SecretAccessKey: '\n', - s3Region: 'eu-west-1', - s3Endpoint: 's3.example.com', - s3UrlStyle: 'path', - s3UseSsl: true, - metadataSchema: 'ducklake', - }, - { omitBlankSecrets: true } - ) + buildDucklakeUpdateApiConfig({ + catalogUrl: ' ', + dataPath: 's3://bucket/path', + poolSize: 4, + s3AccessKeyId: '', + s3SecretAccessKey: '\n', + s3Region: 'eu-west-1', + s3Endpoint: 's3.example.com', + s3UrlStyle: 'path', + s3UseSsl: true, + metadataSchema: 'ducklake', + }) ).toEqual({ ducklake: { catalog_url: undefined, diff --git a/apps/studio/data/replication/create-destination-pipeline-mutation.ts b/apps/studio/data/replication/create-destination-pipeline-mutation.ts index 4ada0784bd8..b583ab6fd77 100644 --- a/apps/studio/data/replication/create-destination-pipeline-mutation.ts +++ b/apps/studio/data/replication/create-destination-pipeline-mutation.ts @@ -2,102 +2,54 @@ import { useMutation, useQueryClient } from '@tanstack/react-query' import type { components } from 'api-types' import { toast } from 'sonner' -import { optionalSecret } from './destination-secret-utils' import { replicationKeys } from './keys' -import type { TableSyncCopyConfig } from '@/components/interfaces/Database/Replication/TableSyncCopy.utils' +import type { + BigQueryDestinationConfig, + DestinationConfig, + DucklakeDestinationConfig, + PipelineConfig, +} from './types' +import { buildPipelineApiConfig, isDucklakeSupabaseConfig } from './utils' import { handleError, post } from '@/data/fetchers' import type { ResponseError, UseCustomMutationOptions } from '@/types' -export type { TableSyncCopyConfig } from '@/components/interfaces/Database/Replication/TableSyncCopy.utils' +type CreateDestinationPipelineBody = + components['schemas']['CreateReplicationDestinationPipelineBody'] +type CreateDestinationApiConfig = CreateDestinationPipelineBody['destination_config'] -export type DestinationConfig = - | { bigQuery: BigQueryDestinationConfig } - | { iceberg: IcebergDestinationConfig } - | { ducklake: DucklakeDestinationConfig } - | { snowflake: SnowflakeDestinationConfig } - | { clickHouse: ClickHouseDestinationConfig } +type CreateBigQueryApiConfig = Extract +type CreateDucklakeApiConfig = Extract -export type BigQueryDestinationConfig = { - projectId: string - datasetId: string - serviceAccountKey: string - connectionPoolSize?: number - maxStalenessMins?: number -} - -export type IcebergDestinationConfig = { - projectRef: string - warehouseName: string - namespace?: string - catalogToken: string - s3AccessKeyId: string - s3SecretAccessKey: string - s3Region: string -} - -// "Custom parameters" DuckLake: caller provides the PostgreSQL catalog URL and the -// S3-compatible storage credentials directly. -export type DucklakeManualDestinationConfig = { - catalogUrl: string - dataPath: string - poolSize?: number - s3AccessKeyId: string - s3SecretAccessKey: string - s3Region: string - s3Endpoint: string - s3UrlStyle?: 'path' | 'vhost' - s3UseSsl?: boolean - metadataSchema?: string -} - -// "Use Supabase" DuckLake: caller provides Supabase project refs and a bucket; the platform -// API resolves these into a catalog URL + provisioned S3 credentials before persisting. -export type DucklakeSupabaseDestinationConfig = { - catalogProjectRef: string - storageProjectRef: string - bucket: string - path?: string - poolSize?: number - metadataSchema?: string -} - -export type DucklakeDestinationConfig = - | DucklakeManualDestinationConfig - | DucklakeSupabaseDestinationConfig - -function isDucklakeSupabaseConfig( - config: DucklakeDestinationConfig -): config is DucklakeSupabaseDestinationConfig { - return 'catalogProjectRef' in config -} - -const maybeOmitBlankSecret = (value: string | undefined, omitBlankSecrets: boolean) => { - if (omitBlankSecrets) return optionalSecret(value) - - return value +// Maps the studio-side BigQuery config to the snake_case `{ big_query: ... }` payload accepted +// by the platform API. Shared by the create and validate mutations. +export function buildBigQueryApiConfig(config: BigQueryDestinationConfig): CreateBigQueryApiConfig { + return { + big_query: { + project_id: config.projectId, + dataset_id: config.datasetId, + service_account_key: config.serviceAccountKey, + connection_pool_size: config.connectionPoolSize, + max_staleness_mins: config.maxStalenessMins, + }, + } } // Maps the studio-side DuckLake config to the snake_case `{ ducklake: ... }` payload accepted // by the platform API. Shared by the create / update / validate mutations. -export function buildDucklakeApiConfig( - config: DucklakeDestinationConfig, - options: { omitBlankSecrets?: boolean } = {} -) { - const omitBlankSecrets = options.omitBlankSecrets ?? false - +export function buildDucklakeApiConfig(config: DucklakeDestinationConfig): CreateDucklakeApiConfig { if (isDucklakeSupabaseConfig(config)) { return { ducklake: { // pool_size / metadata_schema live on the catalog so they apply to the selected // Supabase Postgres catalog (the API resolves catalog-level values over top-level). catalog: { - type: 'supabase_project' as const, + type: 'supabase_project', project_ref: config.catalogProjectRef, pool_size: config.poolSize, metadata_schema: config.metadataSchema, }, storage: { - type: 'supabase_storage' as const, + type: 'supabase_storage', project_ref: config.storageProjectRef, bucket: config.bucket, ...(config.path ? { path: config.path } : {}), @@ -108,11 +60,11 @@ export function buildDucklakeApiConfig( return { ducklake: { - catalog_url: maybeOmitBlankSecret(config.catalogUrl, omitBlankSecrets), + catalog_url: config.catalogUrl, data_path: config.dataPath, pool_size: config.poolSize, - s3_access_key_id: maybeOmitBlankSecret(config.s3AccessKeyId, omitBlankSecrets), - s3_secret_access_key: maybeOmitBlankSecret(config.s3SecretAccessKey, omitBlankSecrets), + s3_access_key_id: config.s3AccessKeyId, + s3_secret_access_key: config.s3SecretAccessKey, s3_region: config.s3Region, s3_endpoint: config.s3Endpoint, s3_url_style: config.s3UrlStyle, @@ -122,60 +74,70 @@ export function buildDucklakeApiConfig( } } -export type SnowflakeDestinationConfig = { - accountId: string - user: string - privateKey: string - privateKeyPassphrase?: string - database: string - schema: string - role?: string -} +export const buildCreateDestinationApiConfig = ( + destinationConfig: DestinationConfig +): CreateDestinationApiConfig => { + if ('bigQuery' in destinationConfig) { + return buildBigQueryApiConfig(destinationConfig.bigQuery) + } -export type ClickHouseDestinationConfig = { - url: string - user: string - password?: string - database: string - engine?: 'merge_tree' | 'replacing_merge_tree' -} + if ('iceberg' in destinationConfig) { + const { + projectRef, + namespace, + warehouseName, + catalogToken, + s3AccessKeyId, + s3SecretAccessKey, + s3Region, + } = destinationConfig.iceberg -export type BatchConfig = { - maxFillMs?: number - maxBytes?: number - memoryBudgetRatio?: number -} + return { + iceberg: { + supabase: { + namespace, + project_ref: projectRef, + warehouse_name: warehouseName, + catalog_token: catalogToken, + s3_access_key_id: s3AccessKeyId, + s3_secret_access_key: s3SecretAccessKey, + s3_region: s3Region, + }, + }, + } + } -export type PipelineConfig = { - publicationName: string - batch?: BatchConfig - maxTableSyncWorkers?: number - maxCopyConnectionsPerTable?: number - invalidatedSlotBehavior?: 'error' | 'recreate' - tableSyncCopy: TableSyncCopyConfig -} + if ('ducklake' in destinationConfig) { + return buildDucklakeApiConfig(destinationConfig.ducklake) + } -export const buildPipelineApiConfig = ({ - publicationName, - batch, - maxTableSyncWorkers, - maxCopyConnectionsPerTable, - invalidatedSlotBehavior, - tableSyncCopy, -}: PipelineConfig) => ({ - publication_name: publicationName, - max_table_sync_workers: maxTableSyncWorkers, - max_copy_connections_per_table: maxCopyConnectionsPerTable, - invalidated_slot_behavior: invalidatedSlotBehavior, - table_sync_copy: tableSyncCopy, - batch: batch - ? { - max_fill_ms: batch.maxFillMs, - max_bytes: batch.maxBytes, - memory_budget_ratio: batch.memoryBudgetRatio, - } - : undefined, -}) + if ('snowflake' in destinationConfig) { + const { accountId, user, privateKey, privateKeyPassphrase, database, schema, role } = + destinationConfig.snowflake + + return { + snowflake: { + account_id: accountId, + user, + private_key: privateKey, + private_key_passphrase: privateKeyPassphrase, + database, + schema, + role, + }, + } + } + + if ('clickHouse' in destinationConfig) { + const { url, user, password, database, engine } = destinationConfig.clickHouse + + return { clickhouse: { url, user, password, database, engine } } + } + + throw new Error( + 'Invalid destination config: must specify bigQuery, iceberg, ducklake, snowflake, or clickHouse' + ) +} export type CreateDestinationPipelineParams = { projectRef: string @@ -197,82 +159,7 @@ async function createDestinationPipeline( ) { if (!projectRef) throw new Error('projectRef is required') - // Build destination_config based on the type - let destination_config: components['schemas']['CreateReplicationDestinationPipelineBody']['destination_config'] - - if ('bigQuery' in destinationConfig) { - const { projectId, datasetId, serviceAccountKey, connectionPoolSize, maxStalenessMins } = - destinationConfig.bigQuery - - destination_config = { - big_query: { - project_id: projectId, - dataset_id: datasetId, - service_account_key: serviceAccountKey, - connection_pool_size: connectionPoolSize, - max_staleness_mins: maxStalenessMins, - }, - } as components['schemas']['CreateReplicationDestinationPipelineBody']['destination_config'] - } else if ('iceberg' in destinationConfig) { - const { - projectRef: icebergProjectRef, - namespace, - warehouseName, - catalogToken, - s3AccessKeyId, - s3SecretAccessKey, - s3Region, - } = destinationConfig.iceberg - - destination_config = { - iceberg: { - supabase: { - namespace, - project_ref: icebergProjectRef, - warehouse_name: warehouseName, - catalog_token: catalogToken, - s3_access_key_id: s3AccessKeyId, - s3_secret_access_key: s3SecretAccessKey, - s3_region: s3Region, - }, - }, - } - } else if ('ducklake' in destinationConfig) { - destination_config = buildDucklakeApiConfig( - destinationConfig.ducklake - ) as components['schemas']['CreateReplicationDestinationPipelineBody']['destination_config'] - } else if ('snowflake' in destinationConfig) { - const { accountId, user, privateKey, privateKeyPassphrase, database, schema, role } = - destinationConfig.snowflake - - destination_config = { - snowflake: { - account_id: accountId, - user, - private_key: privateKey, - private_key_passphrase: privateKeyPassphrase, - database, - schema, - role, - }, - } as components['schemas']['CreateReplicationDestinationPipelineBody']['destination_config'] - } else if ('clickHouse' in destinationConfig) { - const { url, user, password, database, engine } = destinationConfig.clickHouse - - destination_config = { - clickhouse: { - url, - user, - password, - database, - engine, - }, - } as components['schemas']['CreateReplicationDestinationPipelineBody']['destination_config'] - } else { - throw new Error( - 'Invalid destination config: must specify bigQuery, iceberg, ducklake, snowflake, or clickHouse' - ) - } + const destination_config = buildCreateDestinationApiConfig(destinationConfig) const pipeline_config = buildPipelineApiConfig(pipelineConfig) @@ -282,8 +169,7 @@ async function createDestinationPipeline( source_id: sourceId, destination_name: destinationName, destination_config, - pipeline_config: - pipeline_config as components['schemas']['CreateReplicationDestinationPipelineBody']['pipeline_config'], + pipeline_config, }, signal, }) diff --git a/apps/studio/data/replication/types.ts b/apps/studio/data/replication/types.ts new file mode 100644 index 00000000000..7b6d875de52 --- /dev/null +++ b/apps/studio/data/replication/types.ts @@ -0,0 +1,95 @@ +import { components } from 'api-types' + +type CreateDestinationPipelineBody = + components['schemas']['CreateReplicationDestinationPipelineBody'] +export type CreatePipelineApiConfig = CreateDestinationPipelineBody['pipeline_config'] + +export type BatchConfig = { + maxFillMs?: number + maxBytes?: number + memoryBudgetRatio?: number +} + +export type TableSyncCopyConfig = NonNullable + +export type PipelineConfig = { + publicationName: string + batch?: BatchConfig + maxTableSyncWorkers?: number + maxCopyConnectionsPerTable?: number + invalidatedSlotBehavior?: 'error' | 'recreate' + tableSyncCopy: TableSyncCopyConfig +} + +export type DestinationConfig = + | { bigQuery: BigQueryDestinationConfig } + | { iceberg: IcebergDestinationConfig } + | { ducklake: DucklakeDestinationConfig } + | { snowflake: SnowflakeDestinationConfig } + | { clickHouse: ClickHouseDestinationConfig } + +// "Custom parameters" DuckLake: caller provides the PostgreSQL catalog URL and the +// S3-compatible storage credentials directly. +export type DucklakeManualDestinationConfig = { + catalogUrl: string + dataPath: string + poolSize?: number + s3AccessKeyId: string + s3SecretAccessKey: string + s3Region: string + s3Endpoint: string + s3UrlStyle?: 'path' | 'vhost' + s3UseSsl?: boolean + metadataSchema?: string +} + +// "Use Supabase" DuckLake: caller provides Supabase project refs and a bucket; the platform +// API resolves these into a catalog URL + provisioned S3 credentials before persisting. +export type DucklakeSupabaseDestinationConfig = { + catalogProjectRef: string + storageProjectRef: string + bucket: string + path?: string + poolSize?: number + metadataSchema?: string +} + +export type DucklakeDestinationConfig = + | DucklakeManualDestinationConfig + | DucklakeSupabaseDestinationConfig + +export type BigQueryDestinationConfig = { + projectId: string + datasetId: string + serviceAccountKey: string + connectionPoolSize?: number + maxStalenessMins?: number +} + +export type IcebergDestinationConfig = { + projectRef: string + warehouseName: string + namespace?: string + catalogToken: string + s3AccessKeyId: string + s3SecretAccessKey: string + s3Region: string +} + +export type SnowflakeDestinationConfig = { + accountId: string + user: string + privateKey: string + privateKeyPassphrase?: string + database: string + schema: string + role?: string +} + +export type ClickHouseDestinationConfig = { + url: string + user: string + password?: string + database: string + engine?: 'merge_tree' | 'replacing_merge_tree' +} diff --git a/apps/studio/data/replication/update-destination-pipeline-mutation.ts b/apps/studio/data/replication/update-destination-pipeline-mutation.ts index 8356324f701..033a315584f 100644 --- a/apps/studio/data/replication/update-destination-pipeline-mutation.ts +++ b/apps/studio/data/replication/update-destination-pipeline-mutation.ts @@ -1,18 +1,151 @@ import { useMutation, useQueryClient } from '@tanstack/react-query' -import type { components } from 'api-types' +import { components } from 'api-types' import { toast } from 'sonner' -import { - buildDucklakeApiConfig, - buildPipelineApiConfig, - DestinationConfig, - PipelineConfig, -} from './create-destination-pipeline-mutation' import { optionalSecret } from './destination-secret-utils' import { replicationKeys } from './keys' +import type { + BigQueryDestinationConfig, + DestinationConfig, + DucklakeDestinationConfig, + PipelineConfig, +} from './types' +import { buildPipelineApiConfig, isDucklakeSupabaseConfig } from './utils' import { handleError, post } from '@/data/fetchers' import type { ResponseError, UseCustomMutationOptions } from '@/types' +type UpdateDestinationPipelineBody = + components['schemas']['UpdateReplicationDestinationPipelineBody'] +type UpdateDestinationApiConfig = UpdateDestinationPipelineBody['destination_config'] + +type UpdateBigQueryApiConfig = Extract +type UpdateDucklakeApiConfig = Extract + +export function buildBigQueryUpdateApiConfig( + config: BigQueryDestinationConfig +): UpdateBigQueryApiConfig { + return { + big_query: { + project_id: config.projectId, + dataset_id: config.datasetId, + service_account_key: optionalSecret(config.serviceAccountKey), + connection_pool_size: config.connectionPoolSize, + max_staleness_mins: config.maxStalenessMins, + }, + } +} + +export function buildDucklakeUpdateApiConfig( + config: DucklakeDestinationConfig +): UpdateDucklakeApiConfig { + if (isDucklakeSupabaseConfig(config)) { + return { + ducklake: { + catalog: { + type: 'supabase_project', + project_ref: config.catalogProjectRef, + pool_size: config.poolSize, + metadata_schema: config.metadataSchema, + }, + storage: { + type: 'supabase_storage', + project_ref: config.storageProjectRef, + bucket: config.bucket, + ...(config.path ? { path: config.path } : {}), + }, + }, + } + } + + return { + ducklake: { + catalog_url: optionalSecret(config.catalogUrl), + data_path: config.dataPath, + pool_size: config.poolSize, + s3_access_key_id: optionalSecret(config.s3AccessKeyId), + s3_secret_access_key: optionalSecret(config.s3SecretAccessKey), + s3_region: config.s3Region, + s3_endpoint: config.s3Endpoint, + s3_url_style: config.s3UrlStyle, + s3_use_ssl: config.s3UseSsl, + metadata_schema: config.metadataSchema, + }, + } +} + +export const buildUpdateDestinationApiConfig = ( + destinationConfig: DestinationConfig +): UpdateDestinationApiConfig => { + if ('bigQuery' in destinationConfig) { + return buildBigQueryUpdateApiConfig(destinationConfig.bigQuery) + } + + if ('iceberg' in destinationConfig) { + const { + projectRef, + warehouseName, + namespace, + catalogToken, + s3AccessKeyId, + s3SecretAccessKey, + s3Region, + } = destinationConfig.iceberg + + return { + iceberg: { + supabase: { + project_ref: projectRef, + warehouse_name: warehouseName, + namespace, + catalog_token: optionalSecret(catalogToken), + s3_access_key_id: optionalSecret(s3AccessKeyId), + s3_secret_access_key: optionalSecret(s3SecretAccessKey), + s3_region: s3Region, + }, + }, + } + } + + if ('ducklake' in destinationConfig) { + return buildDucklakeUpdateApiConfig(destinationConfig.ducklake) + } + + if ('snowflake' in destinationConfig) { + const { accountId, user, privateKey, privateKeyPassphrase, database, schema, role } = + destinationConfig.snowflake + + return { + snowflake: { + account_id: accountId, + user, + private_key: optionalSecret(privateKey), + private_key_passphrase: optionalSecret(privateKeyPassphrase), + database, + schema, + role, + }, + } + } + + if ('clickHouse' in destinationConfig) { + const { url, user, password, database, engine } = destinationConfig.clickHouse + + return { + clickhouse: { + url, + user, + password: optionalSecret(password), + database, + engine, + }, + } + } + + throw new Error( + 'Invalid destination config: must specify bigQuery, iceberg, ducklake, snowflake, or clickHouse' + ) +} + export type UpdateDestinationPipelineParams = { destinationId: number pipelineId: number @@ -23,11 +156,6 @@ export type UpdateDestinationPipelineParams = { pipelineConfig: PipelineConfig } -type UpdateDestinationPipelineBody = - components['schemas']['UpdateReplicationDestinationPipelineBody'] -type UpdateDestinationConfig = UpdateDestinationPipelineBody['destination_config'] -type UpdatePipelineConfig = UpdateDestinationPipelineBody['pipeline_config'] - async function updateDestinationPipeline( { destinationId: destinationId, @@ -42,78 +170,7 @@ async function updateDestinationPipeline( ) { if (!projectRef) throw new Error('projectRef is required') - // Build destination_config based on the type - let destination_config: UpdateDestinationConfig - - if ('bigQuery' in destinationConfig) { - const { projectId, datasetId, serviceAccountKey, connectionPoolSize, maxStalenessMins } = - destinationConfig.bigQuery - destination_config = { - big_query: { - project_id: projectId, - dataset_id: datasetId, - service_account_key: optionalSecret(serviceAccountKey), - connection_pool_size: connectionPoolSize, - max_staleness_mins: maxStalenessMins, - }, - } as UpdateDestinationConfig - } else if ('iceberg' in destinationConfig) { - const { - projectRef: icebergProjectRef, - warehouseName, - namespace, - catalogToken, - s3AccessKeyId, - s3SecretAccessKey, - s3Region, - } = destinationConfig.iceberg - destination_config = { - iceberg: { - supabase: { - project_ref: icebergProjectRef, - warehouse_name: warehouseName, - namespace: namespace, - catalog_token: optionalSecret(catalogToken), - s3_access_key_id: optionalSecret(s3AccessKeyId), - s3_secret_access_key: optionalSecret(s3SecretAccessKey), - s3_region: s3Region, - }, - }, - } as UpdateDestinationConfig - } else if ('ducklake' in destinationConfig) { - destination_config = buildDucklakeApiConfig(destinationConfig.ducklake, { - omitBlankSecrets: true, - }) as UpdateDestinationConfig - } else if ('snowflake' in destinationConfig) { - const { accountId, user, privateKey, privateKeyPassphrase, database, schema, role } = - destinationConfig.snowflake - destination_config = { - snowflake: { - account_id: accountId, - user, - private_key: optionalSecret(privateKey), - private_key_passphrase: optionalSecret(privateKeyPassphrase), - database, - schema, - role, - }, - } as UpdateDestinationConfig - } else if ('clickHouse' in destinationConfig) { - const { url, user, password, database, engine } = destinationConfig.clickHouse - destination_config = { - clickhouse: { - url, - user, - password: optionalSecret(password), - database, - engine, - }, - } as UpdateDestinationConfig - } else { - throw new Error( - 'Invalid destination config: must specify bigQuery, iceberg, ducklake, snowflake, or clickHouse' - ) - } + const destination_config = buildUpdateDestinationApiConfig(destinationConfig) const pipeline_config = buildPipelineApiConfig(pipelineConfig) @@ -125,7 +182,7 @@ async function updateDestinationPipeline( destination_config, source_id: sourceId, destination_name: destinationName, - pipeline_config: pipeline_config as UpdatePipelineConfig, + pipeline_config, }, signal, } diff --git a/apps/studio/data/replication/utils.ts b/apps/studio/data/replication/utils.ts index 31289a682e7..6ad7c043055 100644 --- a/apps/studio/data/replication/utils.ts +++ b/apps/studio/data/replication/utils.ts @@ -1,3 +1,9 @@ +import type { + CreatePipelineApiConfig, + DucklakeDestinationConfig, + DucklakeSupabaseDestinationConfig, + PipelineConfig, +} from './types' import { MAX_RETRY_FAILURE_COUNT } from '@/data/query-client' import { ResponseError } from '@/types' @@ -36,3 +42,31 @@ export const checkReplicationFeatureFlagRetry = ( return false } + +export function isDucklakeSupabaseConfig( + config: DucklakeDestinationConfig +): config is DucklakeSupabaseDestinationConfig { + return 'catalogProjectRef' in config +} + +export const buildPipelineApiConfig = ({ + publicationName, + batch, + maxTableSyncWorkers, + maxCopyConnectionsPerTable, + invalidatedSlotBehavior, + tableSyncCopy, +}: PipelineConfig): CreatePipelineApiConfig => ({ + publication_name: publicationName, + max_table_sync_workers: maxTableSyncWorkers, + max_copy_connections_per_table: maxCopyConnectionsPerTable, + invalidated_slot_behavior: invalidatedSlotBehavior, + table_sync_copy: tableSyncCopy, + batch: batch + ? { + max_fill_ms: batch.maxFillMs, + max_bytes: batch.maxBytes, + memory_budget_ratio: batch.memoryBudgetRatio, + } + : undefined, +}) diff --git a/apps/studio/data/replication/validate-destination-mutation.ts b/apps/studio/data/replication/validate-destination-mutation.ts index 3e263ac79b3..0d1fc23d5c9 100644 --- a/apps/studio/data/replication/validate-destination-mutation.ts +++ b/apps/studio/data/replication/validate-destination-mutation.ts @@ -1,11 +1,9 @@ import { useMutation } from '@tanstack/react-query' import type { components } from 'api-types' -import { - buildDucklakeApiConfig, - DestinationConfig, - TableSyncCopyConfig, -} from './create-destination-pipeline-mutation' +import { buildCreateDestinationApiConfig } from './create-destination-pipeline-mutation' +import type { DestinationConfig, TableSyncCopyConfig } from './types' +import { buildPipelineApiConfig } from './utils' import { handleError, post } from '@/data/fetchers' import type { ResponseError, UseCustomMutationOptions } from '@/types' @@ -40,108 +38,28 @@ async function validateDestination( ): Promise { if (!projectRef) throw new Error('projectRef is required') - // Build destination_config based on the type - let config: components['schemas']['ValidateReplicationDestinationBody']['config'] - - if ('bigQuery' in destinationConfig) { - const { projectId, datasetId, serviceAccountKey, connectionPoolSize, maxStalenessMins } = - destinationConfig.bigQuery - - config = { - big_query: { - project_id: projectId, - dataset_id: datasetId, - service_account_key: serviceAccountKey, - connection_pool_size: connectionPoolSize, - max_staleness_mins: maxStalenessMins, - }, - } as components['schemas']['ValidateReplicationDestinationBody']['config'] - } else if ('iceberg' in destinationConfig) { - const { - projectRef: icebergProjectRef, - namespace, - warehouseName, - catalogToken, - s3AccessKeyId, - s3SecretAccessKey, - s3Region, - } = destinationConfig.iceberg - - config = { - iceberg: { - supabase: { - namespace, - project_ref: icebergProjectRef, - warehouse_name: warehouseName, - catalog_token: catalogToken, - s3_access_key_id: s3AccessKeyId, - s3_secret_access_key: s3SecretAccessKey, - s3_region: s3Region, - }, - }, - } - } else if ('ducklake' in destinationConfig) { - config = buildDucklakeApiConfig( - destinationConfig.ducklake - ) as components['schemas']['ValidateReplicationDestinationBody']['config'] - } else if ('snowflake' in destinationConfig) { - const { accountId, user, privateKey, privateKeyPassphrase, database, schema, role } = - destinationConfig.snowflake - - config = { - snowflake: { - account_id: accountId, - user, - private_key: privateKey, - private_key_passphrase: privateKeyPassphrase, - database, - schema, - role, - }, - } as components['schemas']['ValidateReplicationDestinationBody']['config'] - } else if ('clickHouse' in destinationConfig) { - const { url, user, password, database, engine } = destinationConfig.clickHouse - - config = { - clickhouse: { - url, - user, - password, - database, - engine, - }, - } as components['schemas']['ValidateReplicationDestinationBody']['config'] - } else { - throw new Error( - 'Invalid destination config: must specify bigQuery, iceberg, ducklake, snowflake, or clickHouse' - ) - } - - const batchConfig = maxFillMs !== undefined ? { max_fill_ms: maxFillMs } : undefined - const pipelineConfig = - publicationName === undefined - ? undefined - : { - publication_name: publicationName, - max_table_sync_workers: maxTableSyncWorkers, - max_copy_connections_per_table: maxCopyConnectionsPerTable, - invalidated_slot_behavior: invalidatedSlotBehavior, - table_sync_copy: tableSyncCopy, - batch: batchConfig, - } - const { data, error } = await post('/platform/replication/{ref}/destinations/validate', { params: { path: { ref: projectRef } }, body: { - config, + config: buildCreateDestinationApiConfig(destinationConfig), source_id: sourceId, - pipeline_config: pipelineConfig, + pipeline_config: + publicationName === undefined + ? undefined + : buildPipelineApiConfig({ + publicationName, + maxTableSyncWorkers, + maxCopyConnectionsPerTable, + invalidatedSlotBehavior, + tableSyncCopy: tableSyncCopy ?? { type: 'include_all_tables' }, + batch: maxFillMs === undefined ? undefined : { maxFillMs }, + }), }, signal, }) if (error) handleError(error) - return data as ValidateDestinationResponse + return data } type ValidateDestinationData = Awaited> diff --git a/apps/studio/data/replication/validate-pipeline-mutation.ts b/apps/studio/data/replication/validate-pipeline-mutation.ts index 81fc461046a..92019fe2484 100644 --- a/apps/studio/data/replication/validate-pipeline-mutation.ts +++ b/apps/studio/data/replication/validate-pipeline-mutation.ts @@ -1,7 +1,8 @@ import { useMutation } from '@tanstack/react-query' import { components } from 'api-types' -import type { TableSyncCopyConfig } from './create-destination-pipeline-mutation' +import { type TableSyncCopyConfig } from './types' +import { buildPipelineApiConfig } from './utils' import { handleError, post } from '@/data/fetchers' import type { ResponseError, UseCustomMutationOptions } from '@/types' @@ -33,28 +34,24 @@ async function validatePipeline( if (!projectRef) throw new Error('projectRef is required') if (!sourceId) throw new Error('sourceId is required') - const batchConfig = maxFillMs !== undefined ? { max_fill_ms: maxFillMs } : undefined - - const config = { - publication_name: publicationName, - max_table_sync_workers: maxTableSyncWorkers, - max_copy_connections_per_table: maxCopyConnectionsPerTable, - invalidated_slot_behavior: invalidatedSlotBehavior, - table_sync_copy: tableSyncCopy, - batch: batchConfig, - } - const { data, error } = await post('/platform/replication/{ref}/pipelines/validate', { params: { path: { ref: projectRef } }, body: { source_id: sourceId, - config: config as components['schemas']['ValidateReplicationPipelineBody']['config'], + config: buildPipelineApiConfig({ + publicationName, + maxTableSyncWorkers, + maxCopyConnectionsPerTable, + invalidatedSlotBehavior, + tableSyncCopy, + batch: maxFillMs === undefined ? undefined : { maxFillMs }, + }), }, signal, }) if (error) handleError(error) - return data as ValidatePipelineResponse + return data } type ValidatePipelineData = Awaited>