// @ts-nocheck import { zodResolver } from '@hookform/resolvers/zod' import { PermissionAction } from '@supabase/shared-types/out/constants' import { useParams } from 'common' import { AnimatePresence, motion } from 'framer-motion' import { Loader2 } from 'lucide-react' import { useEffect, useMemo, useRef, useState } from 'react' import { useForm } from 'react-hook-form' import { toast } from 'sonner' import { Button, DialogSectionSeparator, Form, SheetFooter, SheetSection } from 'ui' import * as z from 'zod' import { useIsETLBigQueryPrivateAlpha, useIsETLDucklakePrivateAlpha, useIsETLIcebergPrivateAlpha, } from '../../useIsETLPrivateAlpha' import { DestinationType } from '../DestinationPanel.types' import { AdvancedSettings } from './AdvancedSettings' import { CREATE_NEW_NAMESPACE } from './DestinationForm.constants' import { DestinationPanelFormSchema as FormSchema } from './DestinationForm.schema' import { buildDestinationConfig, buildDestinationConfigForValidation, getDucklakeValidationIssues, } from './DestinationForm.utils' import { DestinationNameInput } from './DestinationNameInput' import { AnalyticsBucketFields, BigQueryFields, DuckLakeFields } from './DestinationPanelFields' import { NewPublicationPanel } from './NewPublicationPanel' import { NoDestinationsAvailable } from './NoDestinationsAvailable' import { PublicationSelection } from './PublicationSelection' import { ReplicationDisclaimerDialog } from './ReplicationDisclaimerDialog' import { ValidationFailuresSection } from './ValidationFailuresSection' import { CreateAnalyticsBucketSheet } from '@/components/interfaces/Storage/AnalyticsBuckets/CreateAnalyticsBucketSheet' import { getKeys, useAPIKeysQuery } from '@/data/api-keys/api-keys-query' import { useProjectSettingsV2Query } from '@/data/config/project-settings-v2-query' import { BatchConfig, useCreateDestinationPipelineMutation, } from '@/data/replication/create-destination-pipeline-mutation' import { useReplicationDestinationByIdQuery } from '@/data/replication/destination-by-id-query' import { useReplicationPipelineByIdQuery } from '@/data/replication/pipeline-by-id-query' import { useReplicationPublicationsQuery } from '@/data/replication/publications-query' import { useRestartPipelineHelper } from '@/data/replication/restart-pipeline-helper' import { useReplicationSourcesQuery } from '@/data/replication/sources-query' import { useStartPipelineMutation } from '@/data/replication/start-pipeline-mutation' import { useUpdateDestinationPipelineMutation } from '@/data/replication/update-destination-pipeline-mutation' import { useValidateDestinationMutation, type ValidationFailure, } from '@/data/replication/validate-destination-mutation' import { useValidatePipelineMutation } from '@/data/replication/validate-pipeline-mutation' import { useIcebergNamespaceCreateMutation } from '@/data/storage/iceberg-namespace-create-mutation' import { useS3AccessKeyCreateMutation } from '@/data/storage/s3-access-key-create-mutation' import { useAsyncCheckPermissions } from '@/hooks/misc/useCheckPermissions' import { PipelineStatusRequestStatus, usePipelineRequestStatus, } from '@/state/replication-pipeline-request-status' import { type ResponseError } from '@/types' const formId = 'destination-editor' interface DestinationFormProps { selectedType: DestinationType visible: boolean existingDestination?: { sourceId?: number destinationId: number pipelineId?: number enabled: boolean statusName?: string } onClose: () => void } type DucklakeApiConfig = { catalog_url: string data_path: string pool_size?: number s3_access_key_id?: string s3_secret_access_key?: string s3_region?: string s3_endpoint?: string s3_url_style?: 'path' | 'vhost' s3_use_ssl?: boolean metadata_schema?: string expire_snapshots_older_than?: string } export const DestinationForm = ({ selectedType, visible, existingDestination, onClose, }: DestinationFormProps) => { const { ref: projectRef } = useParams() const { setRequestStatus } = usePipelineRequestStatus() const etlEnableBigQuery = useIsETLBigQueryPrivateAlpha() const etlEnableIceberg = useIsETLIcebergPrivateAlpha() const etlEnableDucklake = useIsETLDucklakePrivateAlpha() const { can: canReadAPIKeys } = useAsyncCheckPermissions(PermissionAction.SECRETS_READ, '*') const [isFormInteracting, setIsFormInteracting] = useState(false) const [showDisclaimerDialog, setShowDisclaimerDialog] = useState(false) const [publicationPanelVisible, setPublicationPanelVisible] = useState(false) const [newBucketSheetVisible, setNewBucketSheetVisible] = useState(false) const [pendingFormValues, setPendingFormValues] = useState | null>( null ) const [hasRunValidation, setHasRunValidation] = useState(false) const [destinationValidationFailures, setDestinationValidationFailures] = useState< ValidationFailure[] >([]) const [pipelineValidationFailures, setPipelineValidationFailures] = useState( [] ) const validationSectionRef = useRef(null) const editMode = !!existingDestination // Compute available destinations based on feature flags const availableDestinations = useMemo(() => { const destinations = [] if (etlEnableBigQuery) destinations.push({ value: 'BigQuery', label: 'BigQuery' }) if (etlEnableIceberg) destinations.push({ value: 'Analytics Bucket', label: 'Analytics Bucket' }) if (etlEnableDucklake) destinations.push({ value: 'DuckLake', label: 'DuckLake' }) return destinations }, [etlEnableBigQuery, etlEnableDucklake, etlEnableIceberg]) const hasNoAvailableDestinations = availableDestinations.length === 0 const { data: sourcesData } = useReplicationSourcesQuery({ projectRef }) const sourceId = sourcesData?.sources.find((s) => s.name === projectRef)?.id const { data: publications = [], isSuccess: isSuccessPublications, refetch: refetchPublications, } = useReplicationPublicationsQuery({ projectRef, sourceId }) const { data: destinationData } = useReplicationDestinationByIdQuery({ projectRef, destinationId: existingDestination?.destinationId, }) const { data: pipelineData } = useReplicationPipelineByIdQuery({ projectRef, pipelineId: existingDestination?.pipelineId, }) const { data: apiKeys } = useAPIKeysQuery( { projectRef, reveal: true }, { enabled: canReadAPIKeys } ) const { serviceKey } = getKeys(apiKeys) const catalogToken = serviceKey?.api_key ?? '' const { data: projectSettings } = useProjectSettingsV2Query({ projectRef }) const { mutateAsync: createDestinationPipeline, isPending: creatingDestinationPipeline } = useCreateDestinationPipelineMutation({ onSuccess: () => form.reset(defaultValues), }) const { mutateAsync: updateDestinationPipeline, isPending: updatingDestinationPipeline } = useUpdateDestinationPipelineMutation({ onSuccess: () => form.reset(defaultValues), }) const { mutateAsync: startPipeline, isPending: startingPipeline } = useStartPipelineMutation() const { restartPipeline } = useRestartPipelineHelper() const { mutateAsync: createS3AccessKey, isPending: isCreatingS3AccessKey } = useS3AccessKeyCreateMutation() const { mutateAsync: createNamespace, isPending: isCreatingNamespace } = useIcebergNamespaceCreateMutation() const { mutateAsync: validateDestination, isPending: isValidatingDestination } = useValidateDestinationMutation() const { mutateAsync: validatePipeline, isPending: isValidatingPipeline } = useValidatePipelineMutation() const isValidating = isValidatingDestination || isValidatingPipeline const defaultValues = useMemo(() => { const config = destinationData?.config const isBigQueryConfig = config && 'big_query' in config const isIcebergConfig = config && 'iceberg' in config const ducklakeConfigValue = config && 'ducklake' in (config as Record) ? (config as Record).ducklake : undefined const ducklakeConfig = ducklakeConfigValue && typeof ducklakeConfigValue === 'object' ? (ducklakeConfigValue as DucklakeApiConfig) : undefined return { // Common fields name: destinationData?.name ?? '', publicationName: pipelineData?.config.publication_name ?? '', maxFillMs: pipelineData?.config?.batch?.max_fill_ms ?? undefined, maxTableSyncWorkers: pipelineData?.config?.max_table_sync_workers ?? undefined, maxCopyConnectionsPerTable: pipelineData?.config?.max_copy_connections_per_table ?? undefined, invalidatedSlotBehavior: (pipelineData?.config as { invalidated_slot_behavior?: 'error' | 'recreate' } | undefined) ?.invalidated_slot_behavior ?? undefined, // BigQuery fields projectId: isBigQueryConfig ? config.big_query.project_id : '', datasetId: isBigQueryConfig ? config.big_query.dataset_id : '', serviceAccountKey: isBigQueryConfig ? config.big_query.service_account_key : '', connectionPoolSize: (config as { big_query?: { connection_pool_size?: number } } | undefined)?.big_query ?.connection_pool_size ?? undefined, maxStalenessMins: isBigQueryConfig ? config.big_query.max_staleness_mins : undefined, // Default: null // Analytics Bucket fields warehouseName: isIcebergConfig ? config.iceberg.briven.warehouse_name : '', namespace: isIcebergConfig ? config.iceberg.briven.namespace : '', newNamespaceName: '', catalogToken: isIcebergConfig ? config.iceberg.briven.catalog_token : catalogToken, s3AccessKeyId: isIcebergConfig ? config.iceberg.briven.s3_access_key_id : '', s3SecretAccessKey: isIcebergConfig ? config.iceberg.briven.s3_secret_access_key : '', s3Region: projectSettings?.region ?? (isIcebergConfig ? config.iceberg.briven.s3_region : ''), // DuckLake fields ducklakeCatalogUrl: ducklakeConfig?.catalog_url ?? '', ducklakeDataPath: ducklakeConfig?.data_path ?? '', ducklakePoolSize: ducklakeConfig?.pool_size, ducklakeS3AccessKeyId: ducklakeConfig?.s3_access_key_id ?? '', ducklakeS3SecretAccessKey: ducklakeConfig?.s3_secret_access_key ?? '', ducklakeS3Region: ducklakeConfig?.s3_region ?? '', ducklakeS3Endpoint: ducklakeConfig?.s3_endpoint ?? '', ducklakeS3UrlStyle: ducklakeConfig?.s3_url_style ?? 'path', ducklakeS3UseSsl: ducklakeConfig?.s3_use_ssl ?? true, ducklakeMetadataSchema: ducklakeConfig?.metadata_schema ?? 'ducklake', ducklakeExpireSnapshotsOlderThan: ducklakeConfig?.expire_snapshots_older_than ?? '', } }, [destinationData, pipelineData, catalogToken, projectSettings]) const form = useForm>({ mode: 'onChange', reValidateMode: 'onChange', resolver: zodResolver( FormSchema.superRefine((data, ctx: any) => { const addRequiredFieldError = (path: string, message: string) => { ctx.addIssue({ code: z.ZodIssueCode.custom, message, path: [path], }) } if (selectedType === 'BigQuery') { if (!data.projectId?.length) addRequiredFieldError('projectId', 'Project ID is required') if (!data.datasetId?.length) addRequiredFieldError('datasetId', 'Dataset ID is required') if (!data.serviceAccountKey?.length) addRequiredFieldError('serviceAccountKey', 'Service Account Key is required') } else if (selectedType === 'Analytics Bucket') { if (!data.warehouseName?.length) addRequiredFieldError('warehouseName', 'Bucket is required') const hasValidNamespace = (data.namespace?.length && data.namespace !== 'create-new-namespace') || (data.namespace === 'create-new-namespace' && data.newNamespaceName?.length) if (!hasValidNamespace) { const isCreatingNew = data.namespace === 'create-new-namespace' addRequiredFieldError( isCreatingNew ? 'newNamespaceName' : 'namespace', isCreatingNew ? 'Namespace name is required' : 'Namespace is required' ) } if (!data.s3Region?.length) addRequiredFieldError('s3Region', 'S3 Region is required') if (!data.s3AccessKeyId?.length) addRequiredFieldError('s3AccessKeyId', 'S3 Access Key ID is required') if (data.s3AccessKeyId !== 'create-new' && !data.s3SecretAccessKey?.length) { addRequiredFieldError('s3SecretAccessKey', 'S3 Secret Access Key is required') } } else if (selectedType === 'DuckLake') { getDucklakeValidationIssues(data).forEach(({ path, message }) => { addRequiredFieldError(path, message) }) } }) ), defaultValues, }) const { publicationName, warehouseName } = form.watch() const publicationNames = useMemo(() => publications?.map((pub) => pub.name) ?? [], [publications]) const isSelectedPublicationMissing = isSuccessPublications && !!publicationName && !publicationNames.includes(publicationName) const allValidationFailures = [...destinationValidationFailures, ...pipelineValidationFailures] const hasValidationFailures = allValidationFailures.some((f) => f.failure_type === 'critical') const isSaving = creatingDestinationPipeline || updatingDestinationPipeline || startingPipeline || isCreatingS3AccessKey || isCreatingNamespace || isValidating const isSubmitDisabled = isSaving || isSelectedPublicationMissing || (!editMode && hasNoAvailableDestinations) const getSubmitButtonText = () => { if (editMode) { return existingDestination?.enabled ? 'Apply and restart' : 'Apply and start' } else { return 'Create and start' } } // Helper function to handle namespace creation if needed const resolveNamespace = async (data: z.infer) => { if (data.namespace === CREATE_NEW_NAMESPACE) { if (!data.newNamespaceName) throw new Error('New namespace name is required') await createNamespace({ projectRef, warehouse: data.warehouseName!, namespace: data.newNamespaceName, }) return data.newNamespaceName } return data.namespace } // Helper function to validate configuration const validateConfiguration = async (data: z.infer) => { if (!projectRef || !sourceId) return false setHasRunValidation(true) // Call both validation endpoints in parallel and wait for both to complete // even if one fails - this makes the validation feel like a single operation const results = await Promise.allSettled([ validateDestination({ projectRef, destinationConfig: buildDestinationConfigForValidation({ projectRef, selectedType, data }), }), validatePipeline({ projectRef, sourceId, publicationName: data.publicationName, maxFillMs: data.maxFillMs, maxTableSyncWorkers: data.maxTableSyncWorkers, maxCopyConnectionsPerTable: data.maxCopyConnectionsPerTable, invalidatedSlotBehavior: data.invalidatedSlotBehavior, }), ]) // Extract results from settled promises const destResult = results[0] const pipelineResult = results[1] // Check if any validation request failed completely const hasRequestError = results.some((r) => r.status === 'rejected') if (hasRequestError) { // If any request failed, surface the upstream message so users see why const rejected = results.find((r): r is PromiseRejectedResult => r.status === 'rejected') const reason = rejected?.reason instanceof Error ? rejected.reason.message : 'Please try again.' toast.error(`Failed to validate configuration: ${reason}`) setHasRunValidation(false) return false } // Both requests succeeded, extract validation failures const destValidationResult = destResult.status === 'fulfilled' ? destResult.value : { validation_failures: [] } const pipelineValidationResult = pipelineResult.status === 'fulfilled' ? pipelineResult.value : { validation_failures: [] } setDestinationValidationFailures(destValidationResult.validation_failures) setPipelineValidationFailures(pipelineValidationResult.validation_failures) // Check if there are critical failures or warnings const allFailures = [ ...destValidationResult.validation_failures, ...pipelineValidationResult.validation_failures, ] const hasCriticalFailures = allFailures.some((f) => f.failure_type === 'critical') const hasAnyFailures = allFailures.length > 0 // Scroll to validation section if there are any failures if (hasAnyFailures) { setTimeout(() => { validationSectionRef.current?.scrollIntoView({ behavior: 'smooth', block: 'start' }) }, 100) } return !hasCriticalFailures } const submitPipeline = async (data: z.infer) => { if (!projectRef) return console.error('Project ref is required') if (!sourceId) return console.error('Source id is required') if (isSelectedPublicationMissing) { return toast.error('Please select another publication before continuing') } try { const destinationConfig = await buildDestinationConfig({ projectRef, selectedType, warehouseName, data, createS3AccessKey, resolveNamespace, }) if (!destinationConfig) throw new Error('Destination configuration is missing') const batchConfig: BatchConfig | undefined = data.maxFillMs !== undefined ? { maxFillMs: data.maxFillMs } : undefined const hasBatchFields = batchConfig !== undefined const pipelineConfig = { publicationName: data.publicationName, maxTableSyncWorkers: data.maxTableSyncWorkers, maxCopyConnectionsPerTable: data.maxCopyConnectionsPerTable, invalidatedSlotBehavior: data.invalidatedSlotBehavior, ...(hasBatchFields ? { batch: batchConfig } : {}), } if (editMode && existingDestination) { if (!existingDestination.pipelineId) return console.error('Pipeline id is required') await updateDestinationPipeline({ destinationId: existingDestination.destinationId, pipelineId: existingDestination.pipelineId, projectRef, destinationName: data.name, destinationConfig, pipelineConfig, sourceId, }) // Set request status only right before starting, then fire and close const snapshot = existingDestination.statusName ?? (existingDestination.enabled ? 'started' : 'stopped') if (existingDestination.enabled) { setRequestStatus( existingDestination.pipelineId, PipelineStatusRequestStatus.RestartRequested, snapshot ) toast.success('Settings applied. Restarting the pipeline...') restartPipeline({ projectRef, pipelineId: existingDestination.pipelineId }) } else { setRequestStatus( existingDestination.pipelineId, PipelineStatusRequestStatus.StartRequested, snapshot ) toast.success('Settings applied. Starting the pipeline...') startPipeline({ projectRef, pipelineId: existingDestination.pipelineId }) } onClose() } else { const { pipeline_id: pipelineId } = await createDestinationPipeline({ projectRef, destinationName: data.name, destinationConfig, pipelineConfig, sourceId, }) // Set request status only right before starting, then fire and close setRequestStatus(pipelineId, PipelineStatusRequestStatus.StartRequested, undefined) toast.success('Destination created. Starting the pipeline...') startPipeline({ projectRef, pipelineId }) onClose() } } catch (error) { const action = editMode ? 'apply and run' : 'create and start' toast.error(`Failed to ${action} destination: ${(error as ResponseError).message}`) } } const onSubmit = async (data: z.infer) => { if (!editMode) { // For new pipelines, validate configuration first if not already validated // OR if user has critical failures and clicks "Validate again" if (!hasRunValidation || isValidating || hasValidationFailures) { const isValid = await validateConfiguration(data) if (!isValid) { // Validation failed with critical errors, show inline and stop return } // Validation passed or only has warnings, continue to disclaimer } // Validation passed or only warnings, proceed to disclaimer setPendingFormValues(data) setShowDisclaimerDialog(true) return } await submitPipeline(data) } const handleDisclaimerDialogChange = (open: boolean) => { setShowDisclaimerDialog(open) if (!open) { setPendingFormValues(null) } } const handleDisclaimerConfirm = async () => { if (!pendingFormValues) return const values = pendingFormValues setPendingFormValues(null) setShowDisclaimerDialog(false) await submitPipeline(values) } useEffect(() => { if (editMode && destinationData && pipelineData && !isFormInteracting) { form.reset(defaultValues) } }, [destinationData, pipelineData, editMode, defaultValues, form, isFormInteracting]) // Ensure the form always reflects the freshest data whenever the panel opens useEffect(() => { if (visible) { form.reset(defaultValues) setIsFormInteracting(false) setHasRunValidation(false) setDestinationValidationFailures([]) setPipelineValidationFailures([]) } }, [visible, defaultValues, form]) useEffect(() => { if (visible && projectRef && sourceId) { refetchPublications() } }, [visible, projectRef, sourceId, refetchPublications]) return ( <> {hasNoAvailableDestinations && !editMode ? ( ) : (

Destination details

setPublicationPanelVisible(true)} />
{selectedType === 'BigQuery' && etlEnableBigQuery ? ( ) : selectedType === 'Analytics Bucket' && etlEnableIceberg ? ( setNewBucketSheetVisible(true)} /> ) : selectedType === 'DuckLake' && etlEnableDucklake ? ( ) : null} {!editMode && hasRunValidation && !isValidating && ( <>
)} )}
{isValidating || isSaving ? (

{isValidating ? 'Validating destination configuration...' : `${editMode ? 'Updating' : 'Creating'} destination...`}

) : (
)}
setPublicationPanelVisible(false)} /> ) }