Files
supabase/apps/studio/data/replication/create-destination-pipeline-mutation.ts
b2b150fa3c feat(pipelines): Add UI selector for choosing which tables to skip copy of (#47808)
## Summary

Adds initial-copy scoping to Pipelines in Studio. Users can copy all
existing rows, skip all initial copies, copy only selected publication
tables, or skip selected table copies. All publication tables continue
streaming new changes regardless of the initial-copy policy.

The policy now round-trips through create, edit, validation, and the
generated Management API contract. Initial-copy estimates and
table-restart confirmations use the same scope. Edit requests also
preserve redacted credentials and pipeline settings that Studio does not
own.

This completes the Studio layer of the [ETL API
change](https://github.com/supabase/etl/pull/897) and [Management API
change](https://github.com/supabase/platform/pull/35479).

## Screenshots

### Selector

<img width="1153" height="465" alt="image"
src="https://github.com/user-attachments/assets/bf615e82-ee61-4222-979d-a8695a957e82"
/>

### Select certain tables only

<img width="1153" height="465" alt="image"
src="https://github.com/user-attachments/assets/28adaa24-f239-4d1d-8fb8-fdb1988320cd"
/>

### Confirm copy costs

As the final step before the pipeline is created:

<img width="597" height="619" alt="image"
src="https://github.com/user-attachments/assets/a660bd87-bfb8-41c5-8099-4cdbdef943bf"
/>


### Policy-aware initial-copy estimate

#### Copy no table is selected

<img width="407" height="464" alt="image"
src="https://github.com/user-attachments/assets/99d859ec-2ec3-452a-ab69-11924a8db260"
/>

#### Some tables are selected

<img width="407" height="464" alt="image"
src="https://github.com/user-attachments/assets/e68aedf4-66bc-4372-98ef-0dd7fecef324"
/>


<!-- This is an auto-generated comment: release notes by coderabbit.ai
-->
## Summary by CodeRabbit

- **New Features**
- Added configurable “initial table copy” policies (copy/skip all and
copy/skip selected) during replication setup, including table-picker
behavior, pruning of stale selections, and updated restart/cost
estimates.
- **Bug Fixes**
- Improved restart flows to consistently use `schema.table` identity and
simplified “errored tables” targeting to match error-state tables.
- Reduced unnecessary loading by gating publication/table fetches to
when panels are visible; improved validation/toast handling when
publication tables are unavailable.
- **Tests**
- Added/expanded coverage for destination form submission, table-copy
selection, restart/cost dialogs, and copy-estimate summarization.
- **Style**
- Refreshed warning/label text for clearer configuration and
confirmation messaging.
<!-- end of auto-generated comment: release notes by coderabbit.ai -->

---------

Co-authored-by: Victor Farazdagi <simple.square@gmail.com>
Co-authored-by: Joshen Lim <joshenlimek@gmail.com>
2026-07-24 10:29:46 +03:00

335 lines
9.9 KiB
TypeScript

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 { handleError, post } from '@/data/fetchers'
import type { ResponseError, UseCustomMutationOptions } from '@/types'
export type { TableSyncCopyConfig } from '@/components/interfaces/Database/Replication/TableSyncCopy.utils'
export type DestinationConfig =
| { bigQuery: BigQueryDestinationConfig }
| { iceberg: IcebergDestinationConfig }
| { ducklake: DucklakeDestinationConfig }
| { snowflake: SnowflakeDestinationConfig }
| { clickHouse: ClickHouseDestinationConfig }
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 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
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,
project_ref: config.catalogProjectRef,
pool_size: config.poolSize,
metadata_schema: config.metadataSchema,
},
storage: {
type: 'supabase_storage' as const,
project_ref: config.storageProjectRef,
bucket: config.bucket,
...(config.path ? { path: config.path } : {}),
},
},
}
}
return {
ducklake: {
catalog_url: maybeOmitBlankSecret(config.catalogUrl, omitBlankSecrets),
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_region: config.s3Region,
s3_endpoint: config.s3Endpoint,
s3_url_style: config.s3UrlStyle,
s3_use_ssl: config.s3UseSsl,
metadata_schema: config.metadataSchema,
},
}
}
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'
}
export type BatchConfig = {
maxFillMs?: number
maxBytes?: number
memoryBudgetRatio?: number
}
export type PipelineConfig = {
publicationName: string
batch?: BatchConfig
maxTableSyncWorkers?: number
maxCopyConnectionsPerTable?: number
invalidatedSlotBehavior?: 'error' | 'recreate'
tableSyncCopy: TableSyncCopyConfig
}
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,
})
export type CreateDestinationPipelineParams = {
projectRef: string
destinationName: string
destinationConfig: DestinationConfig
sourceId: number
pipelineConfig: PipelineConfig
}
async function createDestinationPipeline(
{
projectRef,
destinationName: destinationName,
destinationConfig,
pipelineConfig,
sourceId,
}: CreateDestinationPipelineParams,
signal?: AbortSignal
) {
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 pipeline_config = buildPipelineApiConfig(pipelineConfig)
const { data, error } = await post('/platform/replication/{ref}/destinations-pipelines', {
params: { path: { ref: projectRef } },
body: {
source_id: sourceId,
destination_name: destinationName,
destination_config,
pipeline_config:
pipeline_config as components['schemas']['CreateReplicationDestinationPipelineBody']['pipeline_config'],
},
signal,
})
if (error) handleError(error)
return data
}
type CreateDestinationPipelineData = Awaited<ReturnType<typeof createDestinationPipeline>>
export const useCreateDestinationPipelineMutation = ({
onSuccess,
onError,
...options
}: Omit<
UseCustomMutationOptions<
CreateDestinationPipelineData,
ResponseError,
CreateDestinationPipelineParams
>,
'mutationFn'
> = {}) => {
const queryClient = useQueryClient()
return useMutation<CreateDestinationPipelineData, ResponseError, CreateDestinationPipelineParams>(
{
mutationFn: (vars) => createDestinationPipeline(vars),
async onSuccess(data, variables, context) {
const { projectRef } = variables
await Promise.all([
queryClient.invalidateQueries({ queryKey: replicationKeys.destinations(projectRef) }),
queryClient.invalidateQueries({ queryKey: replicationKeys.pipelines(projectRef) }),
])
await onSuccess?.(data, variables, context)
},
async onError(data, variables, context) {
if (onError === undefined) {
toast.error(`Failed to create destination or pipeline: ${data.message}`)
} else {
onError(data, variables, context)
}
},
...options,
}
)
}