From 23ceb9072dd39b7cce613aae7af3c78c0dea1054 Mon Sep 17 00:00:00 2001 From: Raminder Singh Date: Wed, 23 Apr 2025 10:52:46 +0530 Subject: [PATCH] Add replication UI in Studio (#30090) * add dummy sinks and pipelines pages * update api types * show empty sources state * show empty replication state at /replication * create source when enable replication button is clicked * improve replication page when replication is enabled * replace sources page with publications page * publications table * show publications in table * create publication wip * show toast error instead of throwing an exception * user can now delete a publication * show empty sinks page * create and list sinks * add ui to delete a sink * show pipelines on the pipelines page * add ui to create and delete pipelines * get pipeline status wip * show pipeline status wip * show correct label on action buttons * start and stop pipelines * remove a couple of console.logs * fix error when deleting a pipeline * only consider replication enabled when a source with name = ref is present * add source and sink names * correct colspan for 'no pipelines' row * hide 'supabase_realtime' publication on ui * move filtering to fetch query * show sink name in ui * show source/sink names on pipelines page * fix start/stop status shown on ui * fix prettier formatting * update api types * extract pipeline action button as a separate component * fix a crashing page * fixed publications page crash * update to match with changes in api * add new replication page under database * hide replication page behind feature flag * update types * update api types * show destinations empty state * add destinations table * factor out components from Destinations table * show status dot * move pipeline fetch query to parent component * add ability to enable/disable a pipeline * show loader when starting or stopping a pipeline * fix a bug in which loading & empty states were shown together * fix a bug in which error & empty states were shown together * wrap in default layout * add new destination panel * add type field * fix a forwardRef error * fix layout * create destination * delete destination * add ability to create or delete destinations with pipelines * create source if missing * show only a single error * add an enable switch * new layout * add subsections * comment out unused code * show enabled switch only in the header * close panel when destination is created * disable buttons when api requests in flight * reduce panel size * remove commented out code * treat max size and max fill secs as numbers * use drop down to show publications * simpler vertical layout * add separators * add form validation * remove publications drop down padding * add new publication button * hide advanced settings behind an accordion * add some margin between icon and text * show publications panel on clicking new publication button * add header to new publication panel * fix validation not running for publication drop down * create publication in the new publication panel * add table selector in new publication panel * update api types * remove old code * update platform.d.ts * update navigation bar utils * remove a redirect from replication page to publications page * ask user for confirmation before deleting destination * edit destination panel * edit destination panel values fixed * bug fixes * fix prettier formatting * enable/disable pipeline after editing * rename snake_case params to camelCase * loading button when editing * remove merge markers * update api types * add max_staleness parameter in sinks for bigquery * add read replicas flow diagram to replication page * remove an unused import * Revert "add read replicas flow diagram to replication page" This reverts commit 8852d7847b457885603dba786141a8aaf8e99350. * add panel to warn users about additional cost before creating a destination * hide replication page contents behind a feature flag * fix merge conflicts * styling changes * revert static flag * styling updates * fixes * fix switch * copy * fix layout --------- Co-authored-by: Saxon Fletcher --- .../Replication/DeleteDestination.tsx | 41 ++ .../Database/Replication/DestinationPanel.tsx | 519 ++++++++++++++++++ .../Database/Replication/DestinationRow.tsx | 206 +++++++ .../Database/Replication/Destinations.tsx | 138 +++++ .../Replication/NewPublicationPanel.tsx | 179 ++++++ .../Database/Replication/PipelineStatus.tsx | 57 ++ .../Replication/PublicationsComboBox.tsx | 132 +++++ .../Database/Replication/RowMenu.tsx | 72 +++ .../layouts/DatabaseLayout/DatabaseLayout.tsx | 3 + .../DatabaseLayout/DatabaseMenu.utils.tsx | 14 +- .../replication/create-pipeline-mutation.ts | 77 +++ .../create-publication-mutation.ts | 66 +++ .../data/replication/create-sink-mutation.ts | 75 +++ .../replication/create-source-mutation.ts | 56 ++ .../replication/delete-pipeline-mutation.ts | 60 ++ .../delete-publication-mutation.ts | 64 +++ .../data/replication/delete-sink-mutation.ts | 54 ++ apps/studio/data/replication/keys.ts | 15 + .../data/replication/pipeline-by-id-query.ts | 42 ++ .../data/replication/pipeline-status-query.ts | 45 ++ .../data/replication/pipelines-query.ts | 39 ++ .../data/replication/publications-query.ts | 47 ++ .../data/replication/sink-by-id-query.ts | 42 ++ apps/studio/data/replication/sinks-query.ts | 33 ++ apps/studio/data/replication/sources-query.ts | 36 ++ .../replication/start-pipeline-mutation.ts | 60 ++ .../replication/stop-pipeline-mutation.ts | 57 ++ apps/studio/data/replication/tables-query.ts | 40 ++ .../replication/update-pipeline-mutation.ts | 79 +++ .../data/replication/update-sink-mutation.ts | 77 +++ apps/studio/next.config.js | 5 - .../project/[ref]/database/replication.tsx | 42 ++ 32 files changed, 2466 insertions(+), 6 deletions(-) create mode 100644 apps/studio/components/interfaces/Database/Replication/DeleteDestination.tsx create mode 100644 apps/studio/components/interfaces/Database/Replication/DestinationPanel.tsx create mode 100644 apps/studio/components/interfaces/Database/Replication/DestinationRow.tsx create mode 100644 apps/studio/components/interfaces/Database/Replication/Destinations.tsx create mode 100644 apps/studio/components/interfaces/Database/Replication/NewPublicationPanel.tsx create mode 100644 apps/studio/components/interfaces/Database/Replication/PipelineStatus.tsx create mode 100644 apps/studio/components/interfaces/Database/Replication/PublicationsComboBox.tsx create mode 100644 apps/studio/components/interfaces/Database/Replication/RowMenu.tsx create mode 100644 apps/studio/data/replication/create-pipeline-mutation.ts create mode 100644 apps/studio/data/replication/create-publication-mutation.ts create mode 100644 apps/studio/data/replication/create-sink-mutation.ts create mode 100644 apps/studio/data/replication/create-source-mutation.ts create mode 100644 apps/studio/data/replication/delete-pipeline-mutation.ts create mode 100644 apps/studio/data/replication/delete-publication-mutation.ts create mode 100644 apps/studio/data/replication/delete-sink-mutation.ts create mode 100644 apps/studio/data/replication/keys.ts create mode 100644 apps/studio/data/replication/pipeline-by-id-query.ts create mode 100644 apps/studio/data/replication/pipeline-status-query.ts create mode 100644 apps/studio/data/replication/pipelines-query.ts create mode 100644 apps/studio/data/replication/publications-query.ts create mode 100644 apps/studio/data/replication/sink-by-id-query.ts create mode 100644 apps/studio/data/replication/sinks-query.ts create mode 100644 apps/studio/data/replication/sources-query.ts create mode 100644 apps/studio/data/replication/start-pipeline-mutation.ts create mode 100644 apps/studio/data/replication/stop-pipeline-mutation.ts create mode 100644 apps/studio/data/replication/tables-query.ts create mode 100644 apps/studio/data/replication/update-pipeline-mutation.ts create mode 100644 apps/studio/data/replication/update-sink-mutation.ts create mode 100644 apps/studio/pages/project/[ref]/database/replication.tsx diff --git a/apps/studio/components/interfaces/Database/Replication/DeleteDestination.tsx b/apps/studio/components/interfaces/Database/Replication/DeleteDestination.tsx new file mode 100644 index 00000000000..4782500013a --- /dev/null +++ b/apps/studio/components/interfaces/Database/Replication/DeleteDestination.tsx @@ -0,0 +1,41 @@ +import TextConfirmModal from 'ui-patterns/Dialogs/TextConfirmModal' + +interface DeleteDestinationProps { + visible: boolean + setVisible: (value: boolean) => void + onDelete: () => void + isLoading: boolean + name: string +} + +const DeleteDestination = ({ + visible, + setVisible, + onDelete, + isLoading, + name, +}: DeleteDestinationProps) => { + return ( + <> + setVisible(!visible)} + onConfirm={onDelete} + title="Delete this destination" + loading={isLoading} + confirmLabel={`Delete destination`} + confirmPlaceholder="Type in name of destination" + confirmString={name ?? 'Unknown'} + text={ + <> + This will delete the destination{' '} + + } + alert={{ title: 'You cannot recover this destination once deleted.' }} + /> + + ) +} + +export default DeleteDestination diff --git a/apps/studio/components/interfaces/Database/Replication/DestinationPanel.tsx b/apps/studio/components/interfaces/Database/Replication/DestinationPanel.tsx new file mode 100644 index 00000000000..8acd2964420 --- /dev/null +++ b/apps/studio/components/interfaces/Database/Replication/DestinationPanel.tsx @@ -0,0 +1,519 @@ +import { zodResolver } from '@hookform/resolvers/zod' +import { useParams } from 'common' +import { useCreatePipelineMutation } from 'data/replication/create-pipeline-mutation' +import { useCreateSinkMutation } from 'data/replication/create-sink-mutation' +import { useCreateSourceMutation } from 'data/replication/create-source-mutation' +import { useReplicationPublicationsQuery } from 'data/replication/publications-query' +import { useStartPipelineMutation } from 'data/replication/start-pipeline-mutation' +import { useUpdateSinkMutation } from 'data/replication/update-sink-mutation' +import { useUpdatePipelineMutation } from 'data/replication/update-pipeline-mutation' +import { useForm } from 'react-hook-form' +import { toast } from 'sonner' +import { + Accordion_Shadcn_, + AccordionContent_Shadcn_, + AccordionItem_Shadcn_, + AccordionTrigger_Shadcn_, + Alert_Shadcn_, + AlertDescription_Shadcn_, + AlertTitle_Shadcn_, + Button, + Form_Shadcn_, + FormControl_Shadcn_, + FormField_Shadcn_, + Input_Shadcn_, + Select_Shadcn_, + SelectContent_Shadcn_, + SelectGroup_Shadcn_, + SelectItem_Shadcn_, + SelectTrigger_Shadcn_, + Sheet, + SheetContent, + SheetDescription, + SheetFooter, + SheetHeader, + SheetSection, + SheetTitle, + Switch, + TextArea_Shadcn_, + WarningIcon, + Label_Shadcn_ as Label, +} from 'ui' +import * as z from 'zod' +import PublicationsComboBox from './PublicationsComboBox' +import NewPublicationPanel from './NewPublicationPanel' +import { useState, useMemo, useEffect } from 'react' +import { useReplicationSinkByIdQuery } from 'data/replication/sink-by-id-query' +import { useReplicationPipelineByIdQuery } from 'data/replication/pipeline-by-id-query' +import { useStopPipelineMutation } from 'data/replication/stop-pipeline-mutation' +import { FormItemLayout } from 'ui-patterns/form/FormItemLayout/FormItemLayout' + +interface DestinationPanelProps { + visible: boolean + sourceId: number | undefined + onClose: () => void + existingDestination?: { + sourceId?: number + sinkId: number + pipelineId?: number + enabled: boolean + } +} + +const DestinationPanel = ({ + visible, + sourceId, + onClose, + existingDestination, +}: DestinationPanelProps) => { + const { ref: projectRef } = useParams() + const [publicationPanelVisible, setPublicationPanelVisible] = useState(false) + const { mutateAsync: createSource, isLoading: creatingSource } = useCreateSourceMutation() + const { mutateAsync: createSink, isLoading: creatingSink } = useCreateSinkMutation() + const { mutateAsync: createPipeline, isLoading: creatingPipeline } = useCreatePipelineMutation() + const { mutateAsync: startPipeline, isLoading: startingPipeline } = useStartPipelineMutation() + const { mutateAsync: stopPipeline, isLoading: stoppingPipeline } = useStopPipelineMutation() + const { mutateAsync: updateSink, isLoading: updatingSink } = useUpdateSinkMutation() + const { mutateAsync: updatePipeline, isLoading: updatingPipeline } = useUpdatePipelineMutation() + const { data: publications, isLoading: loadingPublications } = useReplicationPublicationsQuery({ + projectRef, + sourceId, + }) + + const { data: sinkData } = useReplicationSinkByIdQuery({ + projectRef, + sinkId: existingDestination?.sinkId, + }) + + const { data: pipelineData } = useReplicationPipelineByIdQuery({ + projectRef, + pipelineId: existingDestination?.pipelineId, + }) + + const isCreating = creatingSource || creatingSink || creatingPipeline || startingPipeline + const isUpdating = updatingSink || updatingPipeline || stoppingPipeline || startingPipeline + const isSubmitting = isCreating || isUpdating + const editMode = !!existingDestination + + const formId = 'destination-editor' + const types = ['BigQuery'] as const + const TypeEnum = z.enum(types) + const FormSchema = z.object({ + type: TypeEnum, + name: z.string().min(1, 'Name is required'), + projectId: z.string().min(1, 'Project id is required'), + datasetId: z.string().min(1, 'Dataset id is required'), + serviceAccountKey: z.string().min(1, 'Service account key is required'), + publicationName: z.string().min(1, 'Publication is required'), + maxSize: z.number().min(1, 'Max Size must be greater than 0').int(), + maxFillSecs: z.number().min(1, 'Max Fill seconds should be greater than 0').int(), + maxStalenessMins: z.number().nonnegative(), + enabled: z.boolean(), + }) + const defaultValues = useMemo( + () => ({ + type: TypeEnum.enum.BigQuery, + name: sinkData?.name ?? '', + projectId: sinkData?.config?.big_query?.project_id ?? '', + datasetId: sinkData?.config?.big_query?.dataset_id ?? '', + serviceAccountKey: sinkData?.config?.big_query?.service_account_key ?? '', + publicationName: pipelineData?.publication_name ?? '', + maxSize: pipelineData?.config?.config?.max_size ?? 1000, + maxFillSecs: pipelineData?.config?.config?.max_fill_secs ?? 10, + maxStalenessMins: sinkData?.config?.big_query?.max_staleness_mins ?? 5, + enabled: existingDestination?.enabled ?? true, + }), + [sinkData, pipelineData, existingDestination] + ) + const form = useForm>({ + mode: 'onBlur', + reValidateMode: 'onBlur', + resolver: zodResolver(FormSchema), + defaultValues, + }) + const onSubmit = async (data: z.infer) => { + if (!projectRef) return console.error('Project ref is required') + try { + if (editMode && existingDestination) { + if (!sourceId) { + console.error('Source id is required') + return + } + if (!existingDestination.pipelineId) { + console.error('Pipeline id is required') + return + } + // Update existing destination + await updateSink({ + projectRef, + sinkId: existingDestination.sinkId, + sinkName: data.name, + projectId: data.projectId, + datasetId: data.datasetId, + serviceAccountKey: data.serviceAccountKey, + maxStalenessMins: data.maxStalenessMins, + }) + + if (existingDestination.pipelineId) { + await updatePipeline({ + projectRef, + pipelineId: existingDestination.pipelineId, + sourceId, + sinkId: existingDestination.sinkId, + publicationName: data.publicationName, + config: { config: { maxSize: data.maxSize, maxFillSecs: data.maxFillSecs } }, + }) + } + if (data.enabled) { + await startPipeline({ projectRef, pipelineId: existingDestination.pipelineId }) + } else { + await stopPipeline({ projectRef, pipelineId: existingDestination.pipelineId }) + } + + toast.success('Successfully updated destination') + } else { + // Create new destination + if (!sourceId) { + console.error('Source id is required') + return + } + const { id: sinkId } = await createSink({ + projectRef, + sinkName: data.name, + projectId: data.projectId, + datasetId: data.datasetId, + serviceAccountKey: data.serviceAccountKey, + maxStalenessMins: data.maxStalenessMins, + }) + const { id: pipelineId } = await createPipeline({ + projectRef, + sourceId, + sinkId, + publicationName: data.publicationName, + config: { config: { maxSize: data.maxSize, maxFillSecs: data.maxFillSecs } }, + }) + if (data.enabled) { + await startPipeline({ projectRef, pipelineId }) + } + toast.success('Successfully created destination') + } + onClose() + } catch (error) { + toast.error(`Failed to ${editMode ? 'update' : 'create'} destination`) + } + } + const onEnableReplication = async () => { + if (!projectRef) return console.error('Project ref is required') + await createSource({ projectRef }) + } + + const { enabled } = form.watch() + + useEffect(() => { + if (editMode && sinkData && pipelineData) { + form.reset(defaultValues) + } + }, [sinkData, pipelineData, editMode, defaultValues, form]) + + return ( + <> + {sourceId ? ( + <> + + +
+ +
+ {editMode ? 'Edit Destination' : 'New Destination'} + + {editMode ? null : 'Send data to a new destination'} + +
+
+ { + form.setValue('enabled', checked) + }} + /> + +
+
+ + +
+ ( + + + + + + )} + /> +

What data to send

+ + ( + + + pub.name) || []} + loading={loadingPublications} + field={field} + onNewPublicationClick={() => setPublicationPanelVisible(true)} + /> + + + )} + /> +

Where to send that data

+ + ( + + + + {field.value} + + + + BigQuery + + + + + + + )} + > + + ( + + + + + + )} + /> + ( + + + + + + )} + /> + ( + + + + + + )} + /> + + + + + Advanced Settings + + + ( + + + + + + )} + /> + ( + + + + + + )} + /> + ( + + + + + + )} + /> + + + +
+ ( + + + + + + )} + /> +
+ +
+
+ + + + +
+
+
+ setPublicationPanelVisible(false)} + /> + + ) : ( + <> + + +
+ + New Destination + + + + + + {/* Pricing to be decided yet */} + Enabling replication will cost additional $xx.xx + + + +
+ +
+
+
+
+ + + +
+
+
+ + )} + + ) +} + +export default DestinationPanel diff --git a/apps/studio/components/interfaces/Database/Replication/DestinationRow.tsx b/apps/studio/components/interfaces/Database/Replication/DestinationRow.tsx new file mode 100644 index 00000000000..465a35cba47 --- /dev/null +++ b/apps/studio/components/interfaces/Database/Replication/DestinationRow.tsx @@ -0,0 +1,206 @@ +import Table from 'components/to-be-cleaned/Table' +import AlertError from 'components/ui/AlertError' +import { ReplicationPipelinesData } from 'data/replication/pipelines-query' +import { ResponseError } from 'types' +import ShimmeringLoader from 'ui-patterns/ShimmeringLoader' +import RowMenu from './RowMenu' +import PipelineStatus from './PipelineStatus' +import { useParams } from 'common' +import { useReplicationPipelineStatusQuery } from 'data/replication/pipeline-status-query' +import { useState } from 'react' +import { toast } from 'sonner' +import { useStartPipelineMutation } from 'data/replication/start-pipeline-mutation' +import { useStopPipelineMutation } from 'data/replication/stop-pipeline-mutation' +import { useDeleteSinkMutation } from 'data/replication/delete-sink-mutation' +import { useDeletePipelineMutation } from 'data/replication/delete-pipeline-mutation' +import DeleteDestination from './DeleteDestination' +import DestinationPanel from './DestinationPanel' + +export type Pipeline = ReplicationPipelinesData['pipelines'][0] + +interface DestinationRowProps { + sourceId: number | undefined + sinkId: number + sinkName: string + type: string + pipeline: Pipeline | undefined + error: ResponseError | null + isLoading: boolean + isError: boolean + isSuccess: boolean +} + +const DestinationRow = ({ + sourceId, + sinkId, + sinkName, + type, + pipeline, + error: pipelineError, + isLoading: isPipelineLoading, + isError: isPipelineError, + isSuccess: isPipelineSuccess, +}: DestinationRowProps) => { + const { ref: projectRef } = useParams() + const [refetchInterval, setRefetchInterval] = useState(false) + const [showDeleteDestinationForm, setShowDeleteDestinationForm] = useState(false) + const [showEditDestinationPanel, setShowEditDestinationPanel] = useState(false) + + const { + data: pipelineStatusData, + error: pipelineStatusError, + isLoading: isPipelineStatusLoading, + isError: isPipelineStatusError, + isSuccess: isPipelineStatusSuccess, + } = useReplicationPipelineStatusQuery( + { + projectRef, + pipelineId: pipeline?.id, + }, + { refetchInterval } + ) + const [requestStatus, setRequestStatus] = useState< + 'None' | 'EnableRequested' | 'DisableRequested' + >('None') + const { mutateAsync: startPipeline } = useStartPipelineMutation() + const { mutateAsync: stopPipeline } = useStopPipelineMutation() + const pipelineStatus = pipelineStatusData?.status + if ( + (requestStatus === 'EnableRequested' && pipelineStatus === 'Started') || + (requestStatus === 'DisableRequested' && pipelineStatus === 'Stopped') + ) { + setRefetchInterval(false) + setRequestStatus('None') + } + + const onEnableClick = async () => { + if (!projectRef) { + console.error('Project ref is required') + return + } + if (!pipeline) { + toast.error('No pipeline found') + return + } + + try { + await startPipeline({ projectRef, pipelineId: pipeline.id }) + } catch (error) { + toast.error('Failed to enable destination') + } + setRequestStatus('EnableRequested') + setRefetchInterval(5000) + } + const onDisableClick = async () => { + if (!projectRef) { + console.error('Project ref is required') + return + } + if (!pipeline) { + toast.error('No pipeline found') + return + } + + try { + await stopPipeline({ projectRef, pipelineId: pipeline.id }) + } catch (error) { + toast.error('Failed to disable destination') + } + setRequestStatus('DisableRequested') + setRefetchInterval(5000) + } + const { mutateAsync: deleteSink } = useDeleteSinkMutation({}) + const { mutateAsync: deletePipeline } = useDeletePipelineMutation({ + onSuccess: (_res: any) => { + toast.success('Successfully deleted destination') + }, + }) + + const onDeleteClick = async () => { + if (!projectRef) { + console.error('Project ref is required') + return + } + if (!pipeline) { + toast.error('No pipeline found') + return + } + + try { + await stopPipeline({ projectRef, pipelineId: pipeline.id }) + await deletePipeline({ projectRef, pipelineId: pipeline.id }) + await deleteSink({ projectRef, sinkId }) + } catch (error) { + toast.error('Failed to delete destination') + } + } + + return ( + <> + {isPipelineError && ( + + )} + {isPipelineSuccess && ( + + + {isPipelineLoading ? : sinkName} + + {isPipelineLoading ? : type} + + {isPipelineLoading || !pipeline ? ( + + ) : ( + + )} + + + {isPipelineLoading || !pipeline ? ( + + ) : ( + pipeline.publication_name + )} + + + setShowDeleteDestinationForm(true)} + onEditClick={() => setShowEditDestinationPanel(true)} + > + + + )} + + setShowEditDestinationPanel(false)} + sourceId={sourceId} + existingDestination={{ + sourceId, + sinkId, + pipelineId: pipeline?.id, + enabled: pipelineStatusData?.status === 'Started', + }} + /> + + ) +} + +export default DestinationRow diff --git a/apps/studio/components/interfaces/Database/Replication/Destinations.tsx b/apps/studio/components/interfaces/Database/Replication/Destinations.tsx new file mode 100644 index 00000000000..a48f109fa11 --- /dev/null +++ b/apps/studio/components/interfaces/Database/Replication/Destinations.tsx @@ -0,0 +1,138 @@ +import { useParams } from 'common' +import Table from 'components/to-be-cleaned/Table' +import AlertError from 'components/ui/AlertError' +import { ButtonTooltip } from 'components/ui/ButtonTooltip' +import { useReplicationSinksQuery } from 'data/replication/sinks-query' +import { Plus } from 'lucide-react' +import { Button, cn } from 'ui' +import { GenericSkeletonLoader } from 'ui-patterns' +import DestinationRow from './DestinationRow' +import { useReplicationPipelinesQuery } from 'data/replication/pipelines-query' +import { FormHeader } from 'components/ui/Forms/FormHeader' +import { useState } from 'react' +import NewDestinationPanel from './DestinationPanel' +import { useReplicationSourcesQuery } from 'data/replication/sources-query' +import { ScaffoldSection, ScaffoldSectionTitle } from 'components/layouts/Scaffold' + +const Destinations = () => { + const [showNewDestinationPanel, setShowNewDestinationPanel] = useState(false) + const { ref: projectRef } = useParams() + + const { + data: sourcesData, + error: sourcesError, + isLoading: isSourcesLoading, + isError: isSourcesError, + isSuccess: isSourcesSuccess, + } = useReplicationSourcesQuery({ + projectRef, + }) + + let sourceId = sourcesData?.sources.find((s) => s.name === projectRef)?.id + + const { + data: sinksData, + error: sinksError, + isLoading: isSinksLoading, + isError: isSinksError, + isSuccess: isSinksSuccess, + } = useReplicationSinksQuery({ + projectRef, + }) + + const { + data: pipelinesData, + error: pipelinesError, + isLoading: isPipelinesLoading, + isError: isPipelinesError, + isSuccess: isPipelinesSuccess, + } = useReplicationPipelinesQuery({ + projectRef, + }) + + const anySinks = isSinksSuccess && sinksData.sinks.length > 0 + + return ( + <> + +
+ Destinations + +
+ {(isSourcesLoading || isSinksLoading) && } + + {(isSourcesError || isSinksError) && ( + + )} + + {anySinks ? ( + Name, + Type, + Status, + Publication, + , + ]} + body={sinksData.sinks.map((sink) => { + const pipeline = pipelinesData?.pipelines.find((p) => p.sink_id === sink.id) + return ( + + ) + })} + >
+ ) : ( + !isSourcesLoading && + !isSinksLoading && + !isSourcesError && + !isSinksError && ( +
+

Send data to your first destination

+

+ Use destinations to improve performance or run analysis on your data via + integrations like BigQuery +

+ +
+ ) + )} +
+ + setShowNewDestinationPanel(false)} + > + + ) +} + +export default Destinations diff --git a/apps/studio/components/interfaces/Database/Replication/NewPublicationPanel.tsx b/apps/studio/components/interfaces/Database/Replication/NewPublicationPanel.tsx new file mode 100644 index 00000000000..2d06e06b4c0 --- /dev/null +++ b/apps/studio/components/interfaces/Database/Replication/NewPublicationPanel.tsx @@ -0,0 +1,179 @@ +import { zodResolver } from '@hookform/resolvers/zod' +import { useParams } from 'common' +import { useCreatePublicationMutation } from 'data/replication/create-publication-mutation' +import { useReplicationTablesQuery } from 'data/replication/tables-query' +import { X } from 'lucide-react' +import { useForm } from 'react-hook-form' +import { toast } from 'sonner' +import { + Sheet, + SheetContent, + SheetHeader, + SheetTitle, + SheetClose, + cn, + Button, + SheetFooter, + SheetSection, + Form_Shadcn_, + FormLabel_Shadcn_, + FormField_Shadcn_, + FormItem_Shadcn_, + FormControl_Shadcn_, + Input_Shadcn_, + FormMessage_Shadcn_, + Card, + CardContent, + SheetDescription, +} from 'ui' +import { MultiSelector } from 'ui-patterns/multi-select' +import { FormItemLayout } from 'ui-patterns/form/FormItemLayout/FormItemLayout' +import { z } from 'zod' + +interface NewPublicationPanelProps { + visible: boolean + sourceId?: number + onClose: () => void +} + +const NewPublicationPanel = ({ visible, sourceId, onClose }: NewPublicationPanelProps) => { + const { ref: projectRef } = useParams() + const { mutateAsync: createPublication, isLoading: creatingPublication } = + useCreatePublicationMutation() + const { data: tables } = useReplicationTablesQuery({ + projectRef, + sourceId, + }) + const formId = 'publication-editor' + const FormSchema = z.object({ + name: z.string().min(1, 'Name is required'), + tables: z.array(z.string()).min(1, 'At least one table is required'), + }) + const defaultValues = { + name: '', + tables: [], + } + const form = useForm>({ + mode: 'onBlur', + reValidateMode: 'onBlur', + resolver: zodResolver(FormSchema), + defaultValues, + }) + + const onSubmit = async (data: z.infer) => { + if (!projectRef) return console.error('Project ref is required') + if (!sourceId) return console.error('Source id is required') + try { + await createPublication({ + projectRef, + sourceId, + name: data.name, + tables: data.tables.map((table) => { + const [schema, name] = table.split('.') + return { schema, name } + }), + }) + toast.success('Successfully created publication') + onClose() + } catch (error) { + toast.error('Failed to create publication') + } + form.reset(defaultValues) + } + + return ( + <> + + +
+ +
+
+ New Publication + + Create a new publication to share table changes for replication + +
+ + + Close + +
+
+ + +
+ ( + + + + + + )} + /> + ( + + + + + + + + + {tables?.tables.map((table) => ( + + {`${table.schema}.${table.name}`} + + ))} + + + + + + )} + /> + +
+
+ + + + +
+
+
+ + ) +} + +export default NewPublicationPanel diff --git a/apps/studio/components/interfaces/Database/Replication/PipelineStatus.tsx b/apps/studio/components/interfaces/Database/Replication/PipelineStatus.tsx new file mode 100644 index 00000000000..2569c32cd14 --- /dev/null +++ b/apps/studio/components/interfaces/Database/Replication/PipelineStatus.tsx @@ -0,0 +1,57 @@ +import AlertError from 'components/ui/AlertError' +import ShimmeringLoader from 'ui-patterns/ShimmeringLoader' +import { cn } from 'ui' +import { ResponseError } from 'types' +import { Loader2 } from 'lucide-react' + +interface PipelineStatusProps { + pipelineStatus: string | undefined + error: ResponseError | null + isLoading: boolean + isError: boolean + isSuccess: boolean + requestStatus: 'None' | 'EnableRequested' | 'DisableRequested' +} + +const PipelineStatus = ({ + pipelineStatus, + error, + isLoading, + isError, + isSuccess, + requestStatus, +}: PipelineStatusProps) => { + const pipelineEnabled = pipelineStatus === 'Stopped' ? false : true + const requestInFlight = requestStatus !== 'None' + const status = + requestStatus === 'EnableRequested' + ? 'Enabling' + : requestStatus === 'DisableRequested' + ? 'Disabling' + : pipelineStatus === 'Stopped' + ? 'Disabled' + : 'Enabled' + return ( + <> + {isLoading && } + {isError && } + {isSuccess && ( +
+ {requestInFlight ? ( + + ) : ( +
+ )} + {status} +
+ )} + + ) +} + +export default PipelineStatus diff --git a/apps/studio/components/interfaces/Database/Replication/PublicationsComboBox.tsx b/apps/studio/components/interfaces/Database/Replication/PublicationsComboBox.tsx new file mode 100644 index 00000000000..7ed6599795b --- /dev/null +++ b/apps/studio/components/interfaces/Database/Replication/PublicationsComboBox.tsx @@ -0,0 +1,132 @@ +import { Check, ChevronsUpDown, Loader2, Plus } from 'lucide-react' +import { useState } from 'react' +import { + Button, + Command_Shadcn_, + CommandEmpty_Shadcn_, + CommandGroup_Shadcn_, + CommandInput_Shadcn_, + CommandItem_Shadcn_, + CommandList_Shadcn_, + CommandSeparator_Shadcn_, + Popover_Shadcn_, + PopoverContent_Shadcn_, + PopoverTrigger_Shadcn_, + ScrollArea, +} from 'ui' +import { ControllerRenderProps } from 'react-hook-form' + +interface PublicationsComboBoxProps { + publications: string[] + loading: boolean + onNewPublicationClick: () => void + field: ControllerRenderProps +} + +const PublicationsComboBox = ({ + publications, + loading, + onNewPublicationClick, + field, +}: PublicationsComboBoxProps) => { + const [dropdownOpen, setDropdownOpen] = useState(false) + const [selectedPublication, setSelectedPublication] = useState(field?.value || '') + const [searchTerm, setSearchTerm] = useState('') + + function handleSearchChange(value: string) { + setSearchTerm(value) + } + + function handlePublicationSelect(pub: string) { + setSelectedPublication(pub) + setDropdownOpen(false) + field.onChange(pub) + } + + return ( + { + setDropdownOpen(open) + if (!open && field?.onBlur) { + field.onBlur() + } + }} + > + + + + + + + + + {loading ? ( +
+ + Loading... +
+ ) : ( + 'No publications found' + )} +
+ + 7 ? 'h-[210px]' : ''}> + {publications.map((pub) => ( + { + handlePublicationSelect(pub) + }} + onClick={() => { + handlePublicationSelect(pub) + }} + > + {pub} + {selectedPublication === pub && ( + + )} + + ))} + + + + + + +

New publication

+
+
+
+
+
+
+ ) +} + +export default PublicationsComboBox diff --git a/apps/studio/components/interfaces/Database/Replication/RowMenu.tsx b/apps/studio/components/interfaces/Database/Replication/RowMenu.tsx new file mode 100644 index 00000000000..2852aea3577 --- /dev/null +++ b/apps/studio/components/interfaces/Database/Replication/RowMenu.tsx @@ -0,0 +1,72 @@ +import AlertError from 'components/ui/AlertError' +import { Edit, MoreVertical, Pause, Play, Trash } from 'lucide-react' +import { ResponseError } from 'types' +import { + Button, + DropdownMenu, + DropdownMenuContent, + DropdownMenuItem, + DropdownMenuSeparator, + DropdownMenuTrigger, +} from 'ui' +import ShimmeringLoader from 'ui-patterns/ShimmeringLoader' + +interface RowMenuProps { + pipelineStatus: string | undefined + error: ResponseError | null + isLoading: boolean + isError: boolean + onEnableClick: () => void + onDisableClick: () => void + onEditClick: () => void + onDeleteClick: () => void +} + +const RowMenu = ({ + pipelineStatus, + error, + isLoading, + isError, + onEnableClick, + onDisableClick, + onEditClick, + onDeleteClick, +}: RowMenuProps) => { + const pipelineEnabled = pipelineStatus === 'Stopped' ? false : true + return ( +
+ {isLoading && } + {isError && } + + + +
+ ) +} + +export default RowMenu diff --git a/apps/studio/components/layouts/DatabaseLayout/DatabaseLayout.tsx b/apps/studio/components/layouts/DatabaseLayout/DatabaseLayout.tsx index 82594b17652..232a284f306 100644 --- a/apps/studio/components/layouts/DatabaseLayout/DatabaseLayout.tsx +++ b/apps/studio/components/layouts/DatabaseLayout/DatabaseLayout.tsx @@ -9,6 +9,7 @@ import { useSelectedProject } from 'hooks/misc/useSelectedProject' import { withAuth } from 'hooks/misc/withAuth' import ProjectLayout from '../ProjectLayout/ProjectLayout' import { generateDatabaseMenu } from './DatabaseMenu.utils' +import { useFlag } from 'hooks/ui/useFlag' export interface DatabaseLayoutProps { title?: string @@ -29,6 +30,7 @@ const DatabaseProductMenu = () => { const pgNetExtensionExists = (data ?? []).find((ext) => ext.name === 'pg_net') !== undefined const pitrEnabled = addons?.selected_addons.find((addon) => addon.type === 'pitr') !== undefined const columnLevelPrivileges = useIsColumnLevelPrivilegesEnabled() + const enablePgReplicate = useFlag('enablePgReplicate') return ( <> @@ -38,6 +40,7 @@ const DatabaseProductMenu = () => { pgNetExtensionExists, pitrEnabled, columnLevelPrivileges, + enablePgReplicate, })} /> diff --git a/apps/studio/components/layouts/DatabaseLayout/DatabaseMenu.utils.tsx b/apps/studio/components/layouts/DatabaseLayout/DatabaseMenu.utils.tsx index de282e97098..c440c327ed6 100644 --- a/apps/studio/components/layouts/DatabaseLayout/DatabaseMenu.utils.tsx +++ b/apps/studio/components/layouts/DatabaseLayout/DatabaseMenu.utils.tsx @@ -9,10 +9,12 @@ export const generateDatabaseMenu = ( pgNetExtensionExists: boolean pitrEnabled: boolean columnLevelPrivileges: boolean + enablePgReplicate: boolean } ): ProductMenuGroup[] => { const ref = project?.ref ?? 'default' - const { pgNetExtensionExists, pitrEnabled, columnLevelPrivileges } = flags || {} + const { pgNetExtensionExists, pitrEnabled, columnLevelPrivileges, enablePgReplicate } = + flags || {} return [ { @@ -62,6 +64,16 @@ export const generateDatabaseMenu = ( url: `/project/${ref}/database/publications`, items: [], }, + ...(enablePgReplicate + ? [ + { + name: 'Replication', + key: 'replication', + url: `/project/${ref}/database/replication`, + items: [], + }, + ] + : []), ], }, { diff --git a/apps/studio/data/replication/create-pipeline-mutation.ts b/apps/studio/data/replication/create-pipeline-mutation.ts new file mode 100644 index 00000000000..98cf620977f --- /dev/null +++ b/apps/studio/data/replication/create-pipeline-mutation.ts @@ -0,0 +1,77 @@ +import { useMutation, UseMutationOptions, useQueryClient } from '@tanstack/react-query' +import { toast } from 'sonner' + +import type { ResponseError } from 'types' +import { replicationKeys } from './keys' +import { handleError, post } from 'data/fetchers' + +export type CreatePipelineParams = { + projectRef: string + sourceId: number + sinkId: number + publicationName: string + config: { config: { maxSize: number; maxFillSecs: number } } +} + +async function createPipeline( + { + projectRef, + sourceId, + sinkId, + publicationName, + config: { + config: { maxSize, maxFillSecs }, + }, + }: CreatePipelineParams, + signal?: AbortSignal +) { + if (!projectRef) throw new Error('projectRef is required') + + const { data, error } = await post('/platform/replication/{ref}/pipelines', { + params: { path: { ref: projectRef } }, + body: { + source_id: sourceId, + sink_id: sinkId, + publication_name: publicationName, + config: { config: { max_size: maxSize, max_fill_secs: maxFillSecs } }, + }, + signal, + }) + if (error) { + handleError(error) + } + + return data +} + +type CreatePipelineData = Awaited> + +export const useCreatePipelineMutation = ({ + onSuccess, + onError, + ...options +}: Omit< + UseMutationOptions, + 'mutationFn' +> = {}) => { + const queryClient = useQueryClient() + + return useMutation( + (vars) => createPipeline(vars), + { + async onSuccess(data, variables, context) { + const { projectRef } = variables + await queryClient.invalidateQueries(replicationKeys.pipelines(projectRef)) + await onSuccess?.(data, variables, context) + }, + async onError(data, variables, context) { + if (onError === undefined) { + toast.error(`Failed to create pipeline: ${data.message}`) + } else { + onError(data, variables, context) + } + }, + ...options, + } + ) +} diff --git a/apps/studio/data/replication/create-publication-mutation.ts b/apps/studio/data/replication/create-publication-mutation.ts new file mode 100644 index 00000000000..8b721eb6862 --- /dev/null +++ b/apps/studio/data/replication/create-publication-mutation.ts @@ -0,0 +1,66 @@ +import { useMutation, UseMutationOptions, useQueryClient } from '@tanstack/react-query' +import { toast } from 'sonner' + +import type { ResponseError } from 'types' +import { replicationKeys } from './keys' +import { handleError, post } from 'data/fetchers' + +export type CreatePublicationParams = { + projectRef: string + sourceId: number + name: string + tables: { schema: string; name: string }[] +} + +async function createPublication( + { projectRef, sourceId, name, tables }: CreatePublicationParams, + signal?: AbortSignal +) { + if (!projectRef) throw new Error('projectRef is required') + + const { data, error } = await post( + '/platform/replication/{ref}/sources/{source_id}/publications', + { + params: { path: { ref: projectRef, source_id: sourceId } }, + body: { name, tables }, + signal, + } + ) + if (error) { + handleError(error) + } + + return data +} + +type CreatePublicationData = Awaited> + +export const useCreatePublicationMutation = ({ + onSuccess, + onError, + ...options +}: Omit< + UseMutationOptions, + 'mutationFn' +> = {}) => { + const queryClient = useQueryClient() + + return useMutation( + (vars) => createPublication(vars), + { + async onSuccess(data, variables, context) { + const { projectRef, sourceId } = variables + await queryClient.invalidateQueries(replicationKeys.publications(projectRef, sourceId)) + await onSuccess?.(data, variables, context) + }, + async onError(data, variables, context) { + if (onError === undefined) { + toast.error(`Failed to create publication: ${data.message}`) + } else { + onError(data, variables, context) + } + }, + ...options, + } + ) +} diff --git a/apps/studio/data/replication/create-sink-mutation.ts b/apps/studio/data/replication/create-sink-mutation.ts new file mode 100644 index 00000000000..d71aac67304 --- /dev/null +++ b/apps/studio/data/replication/create-sink-mutation.ts @@ -0,0 +1,75 @@ +import { useMutation, UseMutationOptions, useQueryClient } from '@tanstack/react-query' +import { toast } from 'sonner' + +import type { ResponseError } from 'types' +import { replicationKeys } from './keys' +import { handleError, post } from 'data/fetchers' + +export type CreateSinkParams = { + projectRef: string + sinkName: string + projectId: string + datasetId: string + serviceAccountKey: string + maxStalenessMins: number +} + +async function createSink( + { + projectRef, + sinkName, + projectId, + datasetId, + serviceAccountKey, + maxStalenessMins, + }: CreateSinkParams, + signal?: AbortSignal +) { + if (!projectRef) throw new Error('projectRef is required') + + const { data, error } = await post('/platform/replication/{ref}/sinks', { + params: { path: { ref: projectRef } }, + body: { + project_id: projectId, + dataset_id: datasetId, + service_account_key: serviceAccountKey, + sink_name: sinkName, + max_staleness_mins: maxStalenessMins, + }, + signal, + }) + if (error) { + handleError(error) + } + + return data +} + +type CreateSinkData = Awaited> + +export const useCreateSinkMutation = ({ + onSuccess, + onError, + ...options +}: Omit< + UseMutationOptions, + 'mutationFn' +> = {}) => { + const queryClient = useQueryClient() + + return useMutation((vars) => createSink(vars), { + async onSuccess(data, variables, context) { + const { projectRef } = variables + await queryClient.invalidateQueries(replicationKeys.sinks(projectRef)) + await onSuccess?.(data, variables, context) + }, + async onError(data, variables, context) { + if (onError === undefined) { + toast.error(`Failed to create sink: ${data.message}`) + } else { + onError(data, variables, context) + } + }, + ...options, + }) +} diff --git a/apps/studio/data/replication/create-source-mutation.ts b/apps/studio/data/replication/create-source-mutation.ts new file mode 100644 index 00000000000..c68438bcf29 --- /dev/null +++ b/apps/studio/data/replication/create-source-mutation.ts @@ -0,0 +1,56 @@ +import { useMutation, UseMutationOptions, useQueryClient } from '@tanstack/react-query' +import { toast } from 'sonner' + +import type { ResponseError } from 'types' +import { replicationKeys } from './keys' +import { handleError, post } from 'data/fetchers' + +export type CreateSourceParams = { + projectRef: string +} + +async function createSource({ projectRef }: CreateSourceParams, signal?: AbortSignal) { + if (!projectRef) throw new Error('projectRef is required') + + const { data, error } = await post('/platform/replication/{ref}/sources', { + params: { path: { ref: projectRef } }, + signal, + }) + if (error) { + handleError(error) + } + + return data +} + +type CreateSourceData = Awaited> + +export const useCreateSourceMutation = ({ + onSuccess, + onError, + ...options +}: Omit< + UseMutationOptions, + 'mutationFn' +> = {}) => { + const queryClient = useQueryClient() + + return useMutation( + (vars) => createSource(vars), + { + async onSuccess(data, variables, context) { + const { projectRef } = variables + await queryClient.invalidateQueries(replicationKeys.sources(projectRef)) + await onSuccess?.(data, variables, context) + }, + async onError(data, variables, context) { + if (onError === undefined) { + toast.error(`Failed to create source: ${data.message}`) + } else { + onError(data, variables, context) + } + }, + ...options, + } + ) +} diff --git a/apps/studio/data/replication/delete-pipeline-mutation.ts b/apps/studio/data/replication/delete-pipeline-mutation.ts new file mode 100644 index 00000000000..75833d37757 --- /dev/null +++ b/apps/studio/data/replication/delete-pipeline-mutation.ts @@ -0,0 +1,60 @@ +import { useMutation, UseMutationOptions, useQueryClient } from '@tanstack/react-query' +import { toast } from 'sonner' + +import type { ResponseError } from 'types' +import { replicationKeys } from './keys' +import { handleError, del } from 'data/fetchers' + +export type DeletePipelineParams = { + projectRef: string + pipelineId: number +} + +async function deletePipeline( + { projectRef, pipelineId }: DeletePipelineParams, + signal?: AbortSignal +) { + if (!projectRef) throw new Error('projectRef is required') + + const { data, error } = await del('/platform/replication/{ref}/pipelines/{pipeline_id}', { + params: { path: { ref: projectRef, pipeline_id: pipelineId } }, + signal, + }) + if (error) { + handleError(error) + } + + return data +} + +type DeletePipelineData = Awaited> + +export const useDeletePipelineMutation = ({ + onSuccess, + onError, + ...options +}: Omit< + UseMutationOptions, + 'mutationFn' +> = {}) => { + const queryClient = useQueryClient() + + return useMutation( + (vars) => deletePipeline(vars), + { + async onSuccess(data, variables, context) { + const { projectRef, pipelineId } = variables + await queryClient.removeQueries(replicationKeys.pipelines(projectRef)) + await onSuccess?.(data, variables, context) + }, + async onError(data, variables, context) { + if (onError === undefined) { + toast.error(`Failed to delete pipeline: ${data.message}`) + } else { + onError(data, variables, context) + } + }, + ...options, + } + ) +} diff --git a/apps/studio/data/replication/delete-publication-mutation.ts b/apps/studio/data/replication/delete-publication-mutation.ts new file mode 100644 index 00000000000..5fac136e151 --- /dev/null +++ b/apps/studio/data/replication/delete-publication-mutation.ts @@ -0,0 +1,64 @@ +import { useMutation, UseMutationOptions, useQueryClient } from '@tanstack/react-query' +import { toast } from 'sonner' + +import type { ResponseError } from 'types' +import { replicationKeys } from './keys' +import { handleError, del } from 'data/fetchers' + +export type DeletePublicationParams = { + projectRef: string + sourceId: number + publicationName: string +} + +async function deletePublication( + { projectRef, sourceId, publicationName: publication_name }: DeletePublicationParams, + signal?: AbortSignal +) { + if (!projectRef) throw new Error('projectRef is required') + + const { data, error } = await del( + '/platform/replication/{ref}/sources/{source_id}/publications/{publication_name}', + { + params: { path: { ref: projectRef, source_id: sourceId, publication_name } }, + signal, + } + ) + if (error) { + handleError(error) + } + + return data +} + +type DeletePublicationData = Awaited> + +export const useDeletePublicationMutation = ({ + onSuccess, + onError, + ...options +}: Omit< + UseMutationOptions, + 'mutationFn' +> = {}) => { + const queryClient = useQueryClient() + + return useMutation( + (vars) => deletePublication(vars), + { + async onSuccess(data, variables, context) { + const { projectRef, sourceId } = variables + await queryClient.invalidateQueries(replicationKeys.publications(projectRef, sourceId)) + await onSuccess?.(data, variables, context) + }, + async onError(data, variables, context) { + if (onError === undefined) { + toast.error(`Failed to delete publication: ${data.message}`) + } else { + onError(data, variables, context) + } + }, + ...options, + } + ) +} diff --git a/apps/studio/data/replication/delete-sink-mutation.ts b/apps/studio/data/replication/delete-sink-mutation.ts new file mode 100644 index 00000000000..101d56255dd --- /dev/null +++ b/apps/studio/data/replication/delete-sink-mutation.ts @@ -0,0 +1,54 @@ +import { useMutation, UseMutationOptions, useQueryClient } from '@tanstack/react-query' +import { toast } from 'sonner' + +import type { ResponseError } from 'types' +import { replicationKeys } from './keys' +import { handleError, del } from 'data/fetchers' + +export type DeleteSinkParams = { + projectRef: string + sinkId: number +} + +async function deleteSink({ projectRef, sinkId }: DeleteSinkParams, signal?: AbortSignal) { + if (!projectRef) throw new Error('projectRef is required') + + const { data, error } = await del('/platform/replication/{ref}/sinks/{sink_id}', { + params: { path: { ref: projectRef, sink_id: sinkId } }, + signal, + }) + if (error) { + handleError(error) + } + + return data +} + +type DeleteSinkData = Awaited> + +export const useDeleteSinkMutation = ({ + onSuccess, + onError, + ...options +}: Omit< + UseMutationOptions, + 'mutationFn' +> = {}) => { + const queryClient = useQueryClient() + + return useMutation((vars) => deleteSink(vars), { + async onSuccess(data, variables, context) { + const { projectRef, sinkId } = variables + await queryClient.invalidateQueries(replicationKeys.sinks(projectRef)) + await onSuccess?.(data, variables, context) + }, + async onError(data, variables, context) { + if (onError === undefined) { + toast.error(`Failed to delete sink: ${data.message}`) + } else { + onError(data, variables, context) + } + }, + ...options, + }) +} diff --git a/apps/studio/data/replication/keys.ts b/apps/studio/data/replication/keys.ts new file mode 100644 index 00000000000..1dcdd61c0f6 --- /dev/null +++ b/apps/studio/data/replication/keys.ts @@ -0,0 +1,15 @@ +export const replicationKeys = { + sources: (projectRef: string | undefined) => ['projects', projectRef, 'sources'] as const, + sinks: (projectRef: string | undefined) => ['projects', projectRef, 'sinks'] as const, + sinkById: (projectRef: string | undefined, sinkId: number | undefined) => + ['projects', projectRef, 'sinks', sinkId] as const, + publications: (projectRef: string | undefined, source_id: number | undefined) => + ['projects', projectRef, 'sources', source_id, 'publications'] as const, + tables: (projectRef: string | undefined, source_id: number | undefined) => + ['projects', projectRef, 'sources', source_id, 'tables'] as const, + pipelines: (projectRef: string | undefined) => ['projects', projectRef, 'pipelines'] as const, + pipelineById: (projectRef: string | undefined, pipelineId: number | undefined) => + ['projects', projectRef, 'pipelines', pipelineId] as const, + pipelinesStatus: (projectRef: string | undefined, pipelineId: number | undefined) => + ['projects', projectRef, 'pipelines', pipelineId, 'status'] as const, +} diff --git a/apps/studio/data/replication/pipeline-by-id-query.ts b/apps/studio/data/replication/pipeline-by-id-query.ts new file mode 100644 index 00000000000..6482dc49c82 --- /dev/null +++ b/apps/studio/data/replication/pipeline-by-id-query.ts @@ -0,0 +1,42 @@ +import { UseQueryOptions, useQuery } from '@tanstack/react-query' + +import { get, handleError } from 'data/fetchers' +import { ResponseError } from 'types' +import { replicationKeys } from './keys' + +type ReplicationPipelineByIdParams = { projectRef?: string; pipelineId?: number } + +async function fetchReplicationPipelineById( + { projectRef, pipelineId }: ReplicationPipelineByIdParams, + signal?: AbortSignal +) { + if (!projectRef) throw new Error('projectRef is required') + if (!pipelineId) throw new Error('pipelineId is required') + const { data, error } = await get('/platform/replication/{ref}/pipelines/{pipeline_id}', { + params: { path: { ref: projectRef, pipeline_id: pipelineId } }, + signal, + }) + if (error) { + handleError(error) + } + + return data +} + +export type ReplicationPipelineByIdData = Awaited> + +export const useReplicationPipelineByIdQuery = ( + { projectRef, pipelineId }: ReplicationPipelineByIdParams, + { + enabled = true, + ...options + }: UseQueryOptions = {} +) => + useQuery( + replicationKeys.pipelineById(projectRef, pipelineId), + ({ signal }) => fetchReplicationPipelineById({ projectRef, pipelineId }, signal), + { + enabled: enabled && typeof projectRef !== 'undefined' && typeof pipelineId !== 'undefined', + ...options, + } + ) diff --git a/apps/studio/data/replication/pipeline-status-query.ts b/apps/studio/data/replication/pipeline-status-query.ts new file mode 100644 index 00000000000..66aaa6633e2 --- /dev/null +++ b/apps/studio/data/replication/pipeline-status-query.ts @@ -0,0 +1,45 @@ +import { UseQueryOptions, useQuery } from '@tanstack/react-query' + +import { get, handleError } from 'data/fetchers' +import { ResponseError } from 'types' +import { replicationKeys } from './keys' + +type ReplicationPipelinesStatusParams = { projectRef?: string; pipelineId?: number } + +async function fetchReplicationPipelineStatus( + { projectRef, pipelineId }: ReplicationPipelinesStatusParams, + signal?: AbortSignal +) { + if (!projectRef) throw new Error('projectRef is required') + if (!pipelineId) throw new Error('pipelineId is required') + + const { data, error } = await get('/platform/replication/{ref}/pipelines/{pipeline_id}/status', { + params: { path: { ref: projectRef, pipeline_id: pipelineId } }, + signal, + }) + if (error) { + handleError(error) + } + + return data +} + +export type ReplicationPipelineStatusData = Awaited< + ReturnType +> + +export const useReplicationPipelineStatusQuery = ( + { projectRef, pipelineId }: ReplicationPipelinesStatusParams, + { + enabled = true, + ...options + }: UseQueryOptions = {} +) => + useQuery( + replicationKeys.pipelinesStatus(projectRef, pipelineId), + ({ signal }) => fetchReplicationPipelineStatus({ projectRef, pipelineId }, signal), + { + enabled: enabled && typeof projectRef !== 'undefined' && typeof pipelineId !== 'undefined', + ...options, + } + ) diff --git a/apps/studio/data/replication/pipelines-query.ts b/apps/studio/data/replication/pipelines-query.ts new file mode 100644 index 00000000000..6e09e1cc409 --- /dev/null +++ b/apps/studio/data/replication/pipelines-query.ts @@ -0,0 +1,39 @@ +import { UseQueryOptions, useQuery } from '@tanstack/react-query' + +import { get, handleError } from 'data/fetchers' +import { ResponseError } from 'types' +import { replicationKeys } from './keys' + +type ReplicationPipelinesParams = { projectRef?: string } + +async function fetchReplicationPipelines( + { projectRef }: ReplicationPipelinesParams, + signal?: AbortSignal +) { + if (!projectRef) throw new Error('projectRef is required') + + const { data, error } = await get('/platform/replication/{ref}/pipelines', { + params: { path: { ref: projectRef } }, + signal, + }) + if (error) { + handleError(error) + } + + return data +} + +export type ReplicationPipelinesData = Awaited> + +export const useReplicationPipelinesQuery = ( + { projectRef }: ReplicationPipelinesParams, + { + enabled = true, + ...options + }: UseQueryOptions = {} +) => + useQuery( + replicationKeys.pipelines(projectRef), + ({ signal }) => fetchReplicationPipelines({ projectRef }, signal), + { enabled: enabled && typeof projectRef !== 'undefined', ...options } + ) diff --git a/apps/studio/data/replication/publications-query.ts b/apps/studio/data/replication/publications-query.ts new file mode 100644 index 00000000000..da56da7d824 --- /dev/null +++ b/apps/studio/data/replication/publications-query.ts @@ -0,0 +1,47 @@ +import { UseQueryOptions, useQuery } from '@tanstack/react-query' + +import { get, handleError } from 'data/fetchers' +import { ResponseError } from 'types' +import { replicationKeys } from './keys' + +type ReplicationPublicationsParams = { projectRef?: string; sourceId?: number } + +async function fetchReplicationPublications( + { projectRef, sourceId }: ReplicationPublicationsParams, + signal?: AbortSignal +) { + if (!projectRef) throw new Error('projectRef is required') + + if (!sourceId) throw new Error('sourceId is required') + + const { data, error } = await get( + '/platform/replication/{ref}/sources/{source_id}/publications', + { + params: { path: { ref: projectRef, source_id: sourceId } }, + signal, + } + ) + if (error) { + handleError(error) + } + + return data.publications.filter((pub) => pub.name !== 'supabase_realtime') +} + +export type ReplicationPublicationsData = Awaited> + +export const useReplicationPublicationsQuery = ( + { projectRef, sourceId }: ReplicationPublicationsParams, + { + enabled = true, + ...options + }: UseQueryOptions = {} +) => + useQuery( + replicationKeys.publications(projectRef, sourceId), + ({ signal }) => fetchReplicationPublications({ projectRef, sourceId }, signal), + { + enabled: enabled && typeof projectRef !== 'undefined' && typeof sourceId !== 'undefined', + ...options, + } + ) diff --git a/apps/studio/data/replication/sink-by-id-query.ts b/apps/studio/data/replication/sink-by-id-query.ts new file mode 100644 index 00000000000..6ad00e2d831 --- /dev/null +++ b/apps/studio/data/replication/sink-by-id-query.ts @@ -0,0 +1,42 @@ +import { UseQueryOptions, useQuery } from '@tanstack/react-query' + +import { get, handleError } from 'data/fetchers' +import { ResponseError } from 'types' +import { replicationKeys } from './keys' + +type ReplicationSinkByIdParams = { projectRef?: string; sinkId?: number } + +async function fetchReplicationSinkById( + { projectRef, sinkId }: ReplicationSinkByIdParams, + signal?: AbortSignal +) { + if (!projectRef) throw new Error('projectRef is required') + if (!sinkId) throw new Error('sinkId is required') + const { data, error } = await get('/platform/replication/{ref}/sinks/{sink_id}', { + params: { path: { ref: projectRef, sink_id: sinkId } }, + signal, + }) + if (error) { + handleError(error) + } + + return data +} + +export type ReplicationSinkByIdData = Awaited> + +export const useReplicationSinkByIdQuery = ( + { projectRef, sinkId }: ReplicationSinkByIdParams, + { + enabled = true, + ...options + }: UseQueryOptions = {} +) => + useQuery( + replicationKeys.sinkById(projectRef, sinkId), + ({ signal }) => fetchReplicationSinkById({ projectRef, sinkId }, signal), + { + enabled: enabled && typeof projectRef !== 'undefined' && typeof sinkId !== 'undefined', + ...options, + } + ) diff --git a/apps/studio/data/replication/sinks-query.ts b/apps/studio/data/replication/sinks-query.ts new file mode 100644 index 00000000000..de08ebb4f1f --- /dev/null +++ b/apps/studio/data/replication/sinks-query.ts @@ -0,0 +1,33 @@ +import { UseQueryOptions, useQuery } from '@tanstack/react-query' + +import { get, handleError } from 'data/fetchers' +import { ResponseError } from 'types' +import { replicationKeys } from './keys' + +type ReplicationSinksParams = { projectRef?: string } + +async function fetchReplicationSinks({ projectRef }: ReplicationSinksParams, signal?: AbortSignal) { + if (!projectRef) throw new Error('projectRef is required') + + const { data, error } = await get('/platform/replication/{ref}/sinks', { + params: { path: { ref: projectRef } }, + signal, + }) + if (error) { + handleError(error) + } + + return data +} + +export type ReplicationSinksData = Awaited> + +export const useReplicationSinksQuery = ( + { projectRef }: ReplicationSinksParams, + { enabled = true, ...options }: UseQueryOptions = {} +) => + useQuery( + replicationKeys.sinks(projectRef), + ({ signal }) => fetchReplicationSinks({ projectRef }, signal), + { enabled: enabled && typeof projectRef !== 'undefined', ...options } + ) diff --git a/apps/studio/data/replication/sources-query.ts b/apps/studio/data/replication/sources-query.ts new file mode 100644 index 00000000000..3009b432fa7 --- /dev/null +++ b/apps/studio/data/replication/sources-query.ts @@ -0,0 +1,36 @@ +import { UseQueryOptions, useQuery } from '@tanstack/react-query' + +import { get, handleError } from 'data/fetchers' +import { ResponseError } from 'types' +import { replicationKeys } from './keys' + +type ReplicationSourcesParams = { projectRef?: string } + +async function fetchReplicationSources( + { projectRef }: ReplicationSourcesParams, + signal?: AbortSignal +) { + if (!projectRef) throw new Error('projectRef is required') + + const { data, error } = await get('/platform/replication/{ref}/sources', { + params: { path: { ref: projectRef } }, + signal, + }) + if (error) { + handleError(error) + } + + return data +} + +export type ReplicationSourcesData = Awaited> + +export const useReplicationSourcesQuery = ( + { projectRef }: ReplicationSourcesParams, + { enabled = true, ...options }: UseQueryOptions = {} +) => + useQuery( + replicationKeys.sources(projectRef), + ({ signal }) => fetchReplicationSources({ projectRef }, signal), + { enabled: enabled && typeof projectRef !== 'undefined', ...options } + ) diff --git a/apps/studio/data/replication/start-pipeline-mutation.ts b/apps/studio/data/replication/start-pipeline-mutation.ts new file mode 100644 index 00000000000..7bcf93917c9 --- /dev/null +++ b/apps/studio/data/replication/start-pipeline-mutation.ts @@ -0,0 +1,60 @@ +import { useMutation, UseMutationOptions, useQueryClient } from '@tanstack/react-query' +import { toast } from 'sonner' + +import type { ResponseError } from 'types' +import { replicationKeys } from './keys' +import { handleError, post } from 'data/fetchers' + +export type StartPipelineParams = { + projectRef: string + pipelineId: number +} + +async function startPipeline( + { projectRef, pipelineId }: StartPipelineParams, + signal?: AbortSignal +) { + if (!projectRef) throw new Error('projectRef is required') + + const { data, error } = await post('/platform/replication/{ref}/pipelines/{pipeline_id}/start', { + params: { path: { ref: projectRef, pipeline_id: pipelineId } }, + signal, + }) + if (error) { + handleError(error) + } + + return data +} + +type StartPipelineData = Awaited> + +export const useStartPipelineMutation = ({ + onSuccess, + onError, + ...options +}: Omit< + UseMutationOptions, + 'mutationFn' +> = {}) => { + const queryClient = useQueryClient() + + return useMutation( + (vars) => startPipeline(vars), + { + async onSuccess(data, variables, context) { + const { projectRef, pipelineId } = variables + await queryClient.invalidateQueries(replicationKeys.pipelinesStatus(projectRef, pipelineId)) + await onSuccess?.(data, variables, context) + }, + async onError(data, variables, context) { + if (onError === undefined) { + toast.error(`Failed to start pipeline: ${data.message}`) + } else { + onError(data, variables, context) + } + }, + ...options, + } + ) +} diff --git a/apps/studio/data/replication/stop-pipeline-mutation.ts b/apps/studio/data/replication/stop-pipeline-mutation.ts new file mode 100644 index 00000000000..ab4c63bcbec --- /dev/null +++ b/apps/studio/data/replication/stop-pipeline-mutation.ts @@ -0,0 +1,57 @@ +import { useMutation, UseMutationOptions, useQueryClient } from '@tanstack/react-query' +import { toast } from 'sonner' + +import type { ResponseError } from 'types' +import { replicationKeys } from './keys' +import { handleError, post } from 'data/fetchers' + +export type StopPipelineParams = { + projectRef: string + pipelineId: number +} + +async function stopPipeline({ projectRef, pipelineId }: StopPipelineParams, signal?: AbortSignal) { + if (!projectRef) throw new Error('projectRef is required') + + const { data, error } = await post('/platform/replication/{ref}/pipelines/{pipeline_id}/stop', { + params: { path: { ref: projectRef, pipeline_id: pipelineId } }, + signal, + }) + if (error) { + handleError(error) + } + + return data +} + +type StartPipelineData = Awaited> + +export const useStopPipelineMutation = ({ + onSuccess, + onError, + ...options +}: Omit< + UseMutationOptions, + 'mutationFn' +> = {}) => { + const queryClient = useQueryClient() + + return useMutation( + (vars) => stopPipeline(vars), + { + async onSuccess(data, variables, context) { + const { projectRef, pipelineId } = variables + await queryClient.invalidateQueries(replicationKeys.pipelinesStatus(projectRef, pipelineId)) + await onSuccess?.(data, variables, context) + }, + async onError(data, variables, context) { + if (onError === undefined) { + toast.error(`Failed to stop pipeline: ${data.message}`) + } else { + onError(data, variables, context) + } + }, + ...options, + } + ) +} diff --git a/apps/studio/data/replication/tables-query.ts b/apps/studio/data/replication/tables-query.ts new file mode 100644 index 00000000000..42a7256f1d9 --- /dev/null +++ b/apps/studio/data/replication/tables-query.ts @@ -0,0 +1,40 @@ +import { UseQueryOptions, useQuery } from '@tanstack/react-query' + +import { get, handleError } from 'data/fetchers' +import { ResponseError } from 'types' +import { replicationKeys } from './keys' + +type ReplicationTablesParams = { projectRef?: string; sourceId?: number } + +async function fetchReplicationTables( + { projectRef, sourceId }: ReplicationTablesParams, + signal?: AbortSignal +) { + if (!projectRef) throw new Error('projectRef is required') + if (!sourceId) throw new Error('sourceId is required') + + const { data, error } = await get('/platform/replication/{ref}/sources/{source_id}/tables', { + params: { path: { ref: projectRef, source_id: sourceId } }, + signal, + }) + if (error) { + handleError(error) + } + + return data +} + +export type ReplicationTablesData = Awaited> + +export const useReplicationTablesQuery = ( + { projectRef, sourceId }: ReplicationTablesParams, + { enabled = true, ...options }: UseQueryOptions = {} +) => + useQuery( + replicationKeys.tables(projectRef, sourceId), + ({ signal }) => fetchReplicationTables({ projectRef, sourceId }, signal), + { + enabled: enabled && typeof projectRef !== 'undefined' && typeof sourceId !== 'undefined', + ...options, + } + ) diff --git a/apps/studio/data/replication/update-pipeline-mutation.ts b/apps/studio/data/replication/update-pipeline-mutation.ts new file mode 100644 index 00000000000..63b287a635c --- /dev/null +++ b/apps/studio/data/replication/update-pipeline-mutation.ts @@ -0,0 +1,79 @@ +import { useMutation, UseMutationOptions, useQueryClient } from '@tanstack/react-query' +import { toast } from 'sonner' + +import type { ResponseError } from 'types' +import { replicationKeys } from './keys' +import { handleError, post } from 'data/fetchers' + +export type UpdatePipelineParams = { + pipelineId: number + projectRef: string + sourceId: number + sinkId: number + publicationName: string + config: { config: { maxSize: number; maxFillSecs: number } } +} + +async function updatePipeline( + { + pipelineId, + projectRef, + sourceId, + sinkId, + publicationName, + config: { + config: { maxSize, maxFillSecs }, + }, + }: UpdatePipelineParams, + signal?: AbortSignal +) { + if (!projectRef) throw new Error('projectRef is required') + + const { data, error } = await post('/platform/replication/{ref}/pipelines/{pipeline_id}', { + params: { path: { ref: projectRef, pipeline_id: pipelineId } }, + body: { + source_id: sourceId, + sink_id: sinkId, + publication_name: publicationName, + config: { config: { max_size: maxSize, max_fill_secs: maxFillSecs } }, + }, + signal, + }) + if (error) { + handleError(error) + } + + return data +} + +type UpdatePipelineData = Awaited> + +export const useUpdatePipelineMutation = ({ + onSuccess, + onError, + ...options +}: Omit< + UseMutationOptions, + 'mutationFn' +> = {}) => { + const queryClient = useQueryClient() + + return useMutation( + (vars) => updatePipeline(vars), + { + async onSuccess(data, variables, context) { + const { projectRef } = variables + await queryClient.invalidateQueries(replicationKeys.pipelines(projectRef)) + await onSuccess?.(data, variables, context) + }, + async onError(data, variables, context) { + if (onError === undefined) { + toast.error(`Failed to update pipeline: ${data.message}`) + } else { + onError(data, variables, context) + } + }, + ...options, + } + ) +} diff --git a/apps/studio/data/replication/update-sink-mutation.ts b/apps/studio/data/replication/update-sink-mutation.ts new file mode 100644 index 00000000000..5f9029aef1e --- /dev/null +++ b/apps/studio/data/replication/update-sink-mutation.ts @@ -0,0 +1,77 @@ +import { useMutation, UseMutationOptions, useQueryClient } from '@tanstack/react-query' +import { toast } from 'sonner' + +import type { ResponseError } from 'types' +import { replicationKeys } from './keys' +import { handleError, post } from 'data/fetchers' + +export type UpdateSinkParams = { + sinkId: number + projectRef: string + sinkName: string + projectId: string + datasetId: string + serviceAccountKey: string + maxStalenessMins: number +} + +async function updateSink( + { + sinkId, + projectRef, + sinkName, + projectId, + datasetId, + serviceAccountKey, + maxStalenessMins, + }: UpdateSinkParams, + signal?: AbortSignal +) { + if (!projectRef) throw new Error('projectRef is required') + + const { data, error } = await post('/platform/replication/{ref}/sinks/{sink_id}', { + params: { path: { ref: projectRef, sink_id: sinkId } }, + body: { + project_id: projectId, + dataset_id: datasetId, + service_account_key: serviceAccountKey, + sink_name: sinkName, + max_staleness_mins: maxStalenessMins, + }, + signal, + }) + if (error) { + handleError(error) + } + + return data +} + +type UpdateSinkData = Awaited> + +export const useUpdateSinkMutation = ({ + onSuccess, + onError, + ...options +}: Omit< + UseMutationOptions, + 'mutationFn' +> = {}) => { + const queryClient = useQueryClient() + + return useMutation((vars) => updateSink(vars), { + async onSuccess(data, variables, context) { + const { projectRef } = variables + await queryClient.invalidateQueries(replicationKeys.sinks(projectRef)) + await onSuccess?.(data, variables, context) + }, + async onError(data, variables, context) { + if (onError === undefined) { + toast.error(`Failed to update sink: ${data.message}`) + } else { + onError(data, variables, context) + } + }, + ...options, + }) +} diff --git a/apps/studio/next.config.js b/apps/studio/next.config.js index adc90c7eb26..94104ae469c 100644 --- a/apps/studio/next.config.js +++ b/apps/studio/next.config.js @@ -222,11 +222,6 @@ const nextConfig = { destination: '/project/:ref/database/tables', permanent: true, }, - { - source: '/project/:ref/database/replication', - destination: '/project/:ref/database/publications', - permanent: true, - }, { source: '/project/:ref/database/graphiql', destination: '/project/:ref/api/graphiql', diff --git a/apps/studio/pages/project/[ref]/database/replication.tsx b/apps/studio/pages/project/[ref]/database/replication.tsx new file mode 100644 index 00000000000..bedb7735d5a --- /dev/null +++ b/apps/studio/pages/project/[ref]/database/replication.tsx @@ -0,0 +1,42 @@ +import type { NextPageWithLayout } from 'types' +import Destinations from 'components/interfaces/Database/Replication/Destinations' +import { ScaffoldContainer, ScaffoldSection } from 'components/layouts/Scaffold' +import DefaultLayout from 'components/layouts/DefaultLayout' +import { useFlag } from 'hooks/ui/useFlag' +import { PageLayout } from 'components/layouts/PageLayout/PageLayout' +import { Admonition } from 'ui-patterns' +import DatabaseLayout from 'components/layouts/DatabaseLayout/DatabaseLayout' + +const DatabaseReplicationPage: NextPageWithLayout = () => { + const enablePgReplicate = useFlag('enablePgReplicate') + + return ( + <> + {enablePgReplicate ? ( + + + + ) : ( + + + +

Replication is not yet available for your project

+
+
+
+ )} + + ) +} + +DatabaseReplicationPage.getLayout = (page) => ( + + + + {page} + + + +) + +export default DatabaseReplicationPage