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.

<!-- This is an auto-generated comment: release notes by coderabbit.ai
-->
## 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.
<!-- end of auto-generated comment: release notes by coderabbit.ai -->

---------

Co-authored-by: Joshen Lim <joshenlimek@gmail.com>
This commit is contained in:
Danny WhiteandJoshen Lim authored and GitHub committed 2026-09-07 15:02:09 +08:00
1 parent 9f5b5ea6a7
commit a351a36e9b
14 files changed
+449 -441

No files matched your search

@@ -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
@@ -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,
@@ -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'
@@ -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,
@@ -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
@@ -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' },
@@ -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
@@ -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,
@@ -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<CreateDestinationApiConfig, { big_query: unknown }>
type CreateDucklakeApiConfig = Extract<CreateDestinationApiConfig, { ducklake: unknown }>
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,
})
+95
View File
@@ -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<CreatePipelineApiConfig['table_sync_copy']>
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'
}
@@ -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<UpdateDestinationApiConfig, { big_query: unknown }>
type UpdateDucklakeApiConfig = Extract<UpdateDestinationApiConfig, { ducklake: unknown }>
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,
}
+34
View File
@@ -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,
})
@@ -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<ValidateDestinationResponse> {
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<ReturnType<typeof validateDestination>>
@@ -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<ReturnType<typeof validatePipeline>>