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>