feat: call the new api endpoints to handle partial failures in pg_replicate (#35418)

* feat: use the new tenant-source api to atomically create the tenant and source in pg_replicate

* feat: use the new sink-pipeline api to atomically create the sink and pipeline in pg_replicate

* fix: remove some unused imports

* feat: use the new sink-pipeline api to atomically update the sink and pipeline in pg_replicate

* feat: deleting the sink cascades to delete the pipeline

* chore: update api types

* fix: revert accidental update to types

* remove unused code

* add a comment explaining why deleting only sink is enough
This commit is contained in:
Raminder Singh authored and GitHub committed 2025-05-10 12:04:37 +05:30
1 parent 6d09404779
commit 520428680d
14 files changed
+510 -495

No files matched your search

@@ -1,12 +1,8 @@
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 { useCreateTenantSourceMutation } from 'data/replication/create-tenant-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 {
@@ -47,6 +43,8 @@ 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'
import { useCreateSinkPipelineMutation } from 'data/replication/create-sink-pipeline-mutation'
import { useUpdateSinkPipelineMutation } from 'data/replication/update-sink-pipeline-mutation'
interface DestinationPanelProps {
visible: boolean
@@ -68,13 +66,14 @@ const DestinationPanel = ({
}: 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: createTenantSource, isLoading: creatingTenantSource } =
useCreateTenantSourceMutation()
const { mutateAsync: createSinkPipeline, isLoading: creatingSinkPipeline } =
useCreateSinkPipelineMutation()
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 { mutateAsync: updateSinkPipeline, isLoading: updatingSinkPipeline } =
useUpdateSinkPipelineMutation()
const { data: publications, isLoading: loadingPublications } = useReplicationPublicationsQuery({
projectRef,
sourceId,
@@ -90,8 +89,8 @@ const DestinationPanel = ({
pipelineId: existingDestination?.pipelineId,
})
const isCreating = creatingSource || creatingSink || creatingPipeline || startingPipeline
const isUpdating = updatingSink || updatingPipeline || stoppingPipeline || startingPipeline
const isCreating = creatingTenantSource || creatingSinkPipeline || startingPipeline
const isUpdating = updatingSinkPipeline || stoppingPipeline || startingPipeline
const isSubmitting = isCreating || isUpdating
const editMode = !!existingDestination
@@ -144,26 +143,25 @@ const DestinationPanel = ({
return
}
// Update existing destination
await updateSink({
projectRef,
await updateSinkPipeline({
sinkId: existingDestination.sinkId,
pipelineId: existingDestination.pipelineId,
projectRef,
sinkName: data.name,
projectId: data.projectId,
datasetId: data.datasetId,
serviceAccountKey: data.serviceAccountKey,
maxStalenessMins: data.maxStalenessMins,
sinkConfig: {
bigQuery: {
projectId: data.projectId,
datasetId: data.datasetId,
serviceAccountKey: data.serviceAccountKey,
maxStalenessMins: data.maxStalenessMins,
},
},
pipelinConfig: {
config: { maxSize: data.maxSize, maxFillSecs: data.maxFillSecs },
},
publicationName: data.publicationName,
sourceId,
})
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 {
@@ -177,20 +175,22 @@ const DestinationPanel = ({
console.error('Source id is required')
return
}
const { id: sinkId } = await createSink({
const { pipeline_id: pipelineId } = await createSinkPipeline({
projectRef,
sinkName: data.name,
projectId: data.projectId,
datasetId: data.datasetId,
serviceAccountKey: data.serviceAccountKey,
maxStalenessMins: data.maxStalenessMins,
})
const { id: pipelineId } = await createPipeline({
projectRef,
sinkConfig: {
bigQuery: {
projectId: data.projectId,
datasetId: data.datasetId,
serviceAccountKey: data.serviceAccountKey,
maxStalenessMins: data.maxStalenessMins,
},
},
sourceId,
sinkId,
publicationName: data.publicationName,
config: { config: { maxSize: data.maxSize, maxFillSecs: data.maxFillSecs } },
pipelinConfig: {
config: { maxSize: data.maxSize, maxFillSecs: data.maxFillSecs },
},
})
if (data.enabled) {
await startPipeline({ projectRef, pipelineId })
@@ -204,7 +204,7 @@ const DestinationPanel = ({
}
const onEnableReplication = async () => {
if (!projectRef) return console.error('Project ref is required')
await createSource({ projectRef })
await createTenantSource({ projectRef })
}
const { enabled } = form.watch()
@@ -12,7 +12,6 @@ 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'
@@ -110,11 +109,6 @@ const DestinationRow = ({
setRefetchInterval(5000)
}
const { mutateAsync: deleteSink } = useDeleteSinkMutation({})
const { mutateAsync: deletePipeline } = useDeletePipelineMutation({
onSuccess: (_res: any) => {
toast.success('Successfully deleted destination')
},
})
const onDeleteClick = async () => {
if (!projectRef) {
@@ -128,7 +122,8 @@ const DestinationRow = ({
try {
await stopPipeline({ projectRef, pipelineId: pipeline.id })
await deletePipeline({ projectRef, pipelineId: pipeline.id })
// deleting the sink also deletes the pipeline because of cascade delete
// so we don't need to call deletePipeline explicitly
await deleteSink({ projectRef, sinkId })
} catch (error) {
toast.error('Failed to delete destination')
@@ -1,14 +1,12 @@
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'
@@ -16,14 +16,9 @@ import {
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'
@@ -1,77 +0,0 @@
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<ReturnType<typeof createPipeline>>
export const useCreatePipelineMutation = ({
onSuccess,
onError,
...options
}: Omit<
UseMutationOptions<CreatePipelineData, ResponseError, CreatePipelineParams>,
'mutationFn'
> = {}) => {
const queryClient = useQueryClient()
return useMutation<CreatePipelineData, ResponseError, CreatePipelineParams>(
(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,
}
)
}
@@ -1,75 +0,0 @@
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<ReturnType<typeof createSink>>
export const useCreateSinkMutation = ({
onSuccess,
onError,
...options
}: Omit<
UseMutationOptions<CreateSinkData, ResponseError, CreateSinkParams>,
'mutationFn'
> = {}) => {
const queryClient = useQueryClient()
return useMutation<CreateSinkData, ResponseError, CreateSinkParams>((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,
})
}
@@ -0,0 +1,104 @@
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 BigQuerySinkConfig = {
projectId: string
datasetId: string
serviceAccountKey: string
maxStalenessMins: number
}
export type CreateSinkPipelineParams = {
projectRef: string
sinkName: string
sinkConfig: {
bigQuery: BigQuerySinkConfig
}
sourceId: number
publicationName: string
pipelinConfig: { config: { maxSize: number; maxFillSecs: number } }
}
async function createSinkPipeline(
{
projectRef,
sinkName,
sinkConfig: {
bigQuery: { projectId, datasetId, serviceAccountKey, maxStalenessMins },
},
pipelinConfig: {
config: { maxSize, maxFillSecs },
},
publicationName,
sourceId,
}: CreateSinkPipelineParams,
signal?: AbortSignal
) {
if (!projectRef) throw new Error('projectRef is required')
const { data, error } = await post('/platform/replication/{ref}/sinks-pipelines', {
params: { path: { ref: projectRef } },
body: {
sink_name: sinkName,
sink_config: {
big_query: {
project_id: projectId,
dataset_id: datasetId,
service_account_key: serviceAccountKey,
max_staleness_mins: maxStalenessMins,
},
},
pipeline_config: {
config: {
max_size: maxSize,
max_fill_secs: maxFillSecs,
},
},
publication_name: publicationName,
source_id: sourceId,
},
signal,
})
if (error) {
handleError(error)
}
return data
}
type CreateSinkPipelineData = Awaited<ReturnType<typeof createSinkPipeline>>
export const useCreateSinkPipelineMutation = ({
onSuccess,
onError,
...options
}: Omit<
UseMutationOptions<CreateSinkPipelineData, ResponseError, CreateSinkPipelineParams>,
'mutationFn'
> = {}) => {
const queryClient = useQueryClient()
return useMutation<CreateSinkPipelineData, ResponseError, CreateSinkPipelineParams>(
(vars) => createSinkPipeline(vars),
{
async onSuccess(data, variables, context) {
const { projectRef } = variables
await queryClient.invalidateQueries(replicationKeys.sinks(projectRef))
await queryClient.invalidateQueries(replicationKeys.pipelines(projectRef))
await onSuccess?.(data, variables, context)
},
async onError(data, variables, context) {
if (onError === undefined) {
toast.error(`Failed to create sink or pipeline: ${data.message}`)
} else {
onError(data, variables, context)
}
},
...options,
}
)
}
@@ -5,14 +5,14 @@ import type { ResponseError } from 'types'
import { replicationKeys } from './keys'
import { handleError, post } from 'data/fetchers'
export type CreateSourceParams = {
export type CreateTenantSourceParams = {
projectRef: string
}
async function createSource({ projectRef }: CreateSourceParams, signal?: AbortSignal) {
async function createTenantSource({ projectRef }: CreateTenantSourceParams, signal?: AbortSignal) {
if (!projectRef) throw new Error('projectRef is required')
const { data, error } = await post('/platform/replication/{ref}/sources', {
const { data, error } = await post('/platform/replication/{ref}/tenants-sources', {
params: { path: { ref: projectRef } },
signal,
})
@@ -23,20 +23,20 @@ async function createSource({ projectRef }: CreateSourceParams, signal?: AbortSi
return data
}
type CreateSourceData = Awaited<ReturnType<typeof createSource>>
type CreateTenantSourceData = Awaited<ReturnType<typeof createTenantSource>>
export const useCreateSourceMutation = ({
export const useCreateTenantSourceMutation = ({
onSuccess,
onError,
...options
}: Omit<
UseMutationOptions<CreateSourceData, ResponseError, CreateSourceParams>,
UseMutationOptions<CreateTenantSourceData, ResponseError, CreateTenantSourceParams>,
'mutationFn'
> = {}) => {
const queryClient = useQueryClient()
return useMutation<CreateSourceData, ResponseError, CreateSourceParams>(
(vars) => createSource(vars),
return useMutation<CreateTenantSourceData, ResponseError, CreateTenantSourceParams>(
(vars) => createTenantSource(vars),
{
async onSuccess(data, variables, context) {
const { projectRef } = variables
@@ -45,7 +45,7 @@ export const useCreateSourceMutation = ({
},
async onError(data, variables, context) {
if (onError === undefined) {
toast.error(`Failed to create source: ${data.message}`)
toast.error(`Failed to create tenant or source: ${data.message}`)
} else {
onError(data, variables, context)
}
@@ -1,60 +0,0 @@
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<ReturnType<typeof deletePipeline>>
export const useDeletePipelineMutation = ({
onSuccess,
onError,
...options
}: Omit<
UseMutationOptions<DeletePipelineData, ResponseError, DeletePipelineParams>,
'mutationFn'
> = {}) => {
const queryClient = useQueryClient()
return useMutation<DeletePipelineData, ResponseError, DeletePipelineParams>(
(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,
}
)
}
@@ -1,64 +0,0 @@
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<ReturnType<typeof deletePublication>>
export const useDeletePublicationMutation = ({
onSuccess,
onError,
...options
}: Omit<
UseMutationOptions<DeletePublicationData, ResponseError, DeletePublicationParams>,
'mutationFn'
> = {}) => {
const queryClient = useQueryClient()
return useMutation<DeletePublicationData, ResponseError, DeletePublicationParams>(
(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,
}
)
}
@@ -1,79 +0,0 @@
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<ReturnType<typeof updatePipeline>>
export const useUpdatePipelineMutation = ({
onSuccess,
onError,
...options
}: Omit<
UseMutationOptions<UpdatePipelineData, ResponseError, UpdatePipelineParams>,
'mutationFn'
> = {}) => {
const queryClient = useQueryClient()
return useMutation<UpdatePipelineData, ResponseError, UpdatePipelineParams>(
(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,
}
)
}
@@ -1,77 +0,0 @@
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<ReturnType<typeof updateSink>>
export const useUpdateSinkMutation = ({
onSuccess,
onError,
...options
}: Omit<
UseMutationOptions<UpdateSinkData, ResponseError, UpdateSinkParams>,
'mutationFn'
> = {}) => {
const queryClient = useQueryClient()
return useMutation<UpdateSinkData, ResponseError, UpdateSinkParams>((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,
})
}
@@ -0,0 +1,111 @@
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 BigQuerySinkConfig = {
projectId: string
datasetId: string
serviceAccountKey: string
maxStalenessMins: number
}
export type UpdateSinkPipelineParams = {
sinkId: number
pipelineId: number
projectRef: string
sinkName: string
sinkConfig: {
bigQuery: BigQuerySinkConfig
}
sourceId: number
publicationName: string
pipelinConfig: { config: { maxSize: number; maxFillSecs: number } }
}
async function updateSinkPipeline(
{
sinkId,
pipelineId,
projectRef,
sinkName,
sinkConfig: {
bigQuery: { projectId, datasetId, serviceAccountKey, maxStalenessMins },
},
pipelinConfig: {
config: { maxSize, maxFillSecs },
},
publicationName,
sourceId,
}: UpdateSinkPipelineParams,
signal?: AbortSignal
) {
if (!projectRef) throw new Error('projectRef is required')
const { data, error } = await post(
'/platform/replication/{ref}/sinks-pipelines/{sink_id}/{pipeline_id}',
{
params: { path: { ref: projectRef, sink_id: sinkId, pipeline_id: pipelineId } },
body: {
sink_name: sinkName,
sink_config: {
big_query: {
project_id: projectId,
dataset_id: datasetId,
service_account_key: serviceAccountKey,
max_staleness_mins: maxStalenessMins,
},
},
pipeline_config: {
config: {
max_size: maxSize,
max_fill_secs: maxFillSecs,
},
},
publication_name: publicationName,
source_id: sourceId,
},
signal,
}
)
if (error) {
handleError(error)
}
return data
}
type UpdateSinkPipelineData = Awaited<ReturnType<typeof updateSinkPipeline>>
export const useUpdateSinkPipelineMutation = ({
onSuccess,
onError,
...options
}: Omit<
UseMutationOptions<UpdateSinkPipelineData, ResponseError, UpdateSinkPipelineParams>,
'mutationFn'
> = {}) => {
const queryClient = useQueryClient()
return useMutation<UpdateSinkPipelineData, ResponseError, UpdateSinkPipelineParams>(
(vars) => updateSinkPipeline(vars),
{
async onSuccess(data, variables, context) {
const { projectRef } = variables
await queryClient.invalidateQueries(replicationKeys.sinks(projectRef))
await queryClient.invalidateQueries(replicationKeys.pipelines(projectRef))
await onSuccess?.(data, variables, context)
},
async onError(data, variables, context) {
if (onError === undefined) {
toast.error(`Failed to update sink or pipeline: ${data.message}`)
} else {
onError(data, variables, context)
}
},
...options,
}
)
}
+244
View File
@@ -3281,6 +3281,40 @@ export interface paths {
patch?: never
trace?: never
}
'/platform/replication/{ref}/sinks-pipelines': {
parameters: {
query?: never
header?: never
path?: never
cookie?: never
}
get?: never
put?: never
/** Creates a replication sink and pipeline. */
post: operations['ReplicationSinksPipelinesController_createSinkPipeline']
delete?: never
options?: never
head?: never
patch?: never
trace?: never
}
'/platform/replication/{ref}/sinks-pipelines/{sink_id}/{pipeline_id}': {
parameters: {
query?: never
header?: never
path?: never
cookie?: never
}
get?: never
put?: never
/** Updates a replication sink and pipeline */
post: operations['ReplicationSinksPipelinesController_updateSinkPipeline']
delete?: never
options?: never
head?: never
patch?: never
trace?: never
}
'/platform/replication/{ref}/sinks/{sink_id}': {
parameters: {
query?: never
@@ -3370,6 +3404,23 @@ export interface paths {
patch?: never
trace?: never
}
'/platform/replication/{ref}/tenants-sources': {
parameters: {
query?: never
header?: never
path?: never
cookie?: never
}
get?: never
put?: never
/** Creates a replication tenant and source. */
post: operations['ReplicationTenantsSourcesController_createTenantSource']
delete?: never
options?: never
head?: never
patch?: never
trace?: never
}
'/platform/reset-password': {
parameters: {
query?: never
@@ -4605,10 +4656,44 @@ export interface components {
/** @description Sink name */
sink_name: string
}
CreateReplicationSinkPipelineBody: {
/** @description Pipeline config */
pipeline_config: {
config: {
max_fill_secs: number
max_size: number
}
}
/** @description Publication name */
publication_name: string
/** @description Sink config */
sink_config: {
big_query: {
/** @description BigQuery dataset id */
dataset_id: string
/** @description Max staleness in minutes */
max_staleness_mins: number
/** @description BigQuery project id */
project_id: string
/** @description BigQuery service account key */
service_account_key: string
}
}
/** @description Sink name */
sink_name: string
/** @description Source id */
source_id: number
}
CreateSchemaBody: {
name: string
owner: string
}
CreateSinkPipelineResponse: {
/** @description Pipeline id */
pipeline_id: number
/** @description Sink id */
sink_id: number
}
CreateSinkResponse: {
id: number
}
@@ -4640,6 +4725,12 @@ export interface components {
type: string
value: string
}
CreateTenantSourceResponse: {
/** @description Source id */
source_id: number
/** @description Tenant id */
tenant_id: string
}
CreateTriggerBody: {
/** @enum {string} */
activation: 'AFTER' | 'BEFORE'
@@ -6934,9 +7025,13 @@ export interface components {
ReplicationSinkResponse: {
config: {
big_query: {
/** @description BigQuery dataset id */
dataset_id: string
/** @description Max staleness in minutes */
max_staleness_mins: number
/** @description BigQuery project id */
project_id: string
/** @description BigQuery service account key */
service_account_key: string
}
}
@@ -6948,9 +7043,13 @@ export interface components {
sinks: {
config: {
big_query: {
/** @description BigQuery dataset id */
dataset_id: string
/** @description Max staleness in minutes */
max_staleness_mins: number
/** @description BigQuery project id */
project_id: string
/** @description BigQuery service account key */
service_account_key: string
}
}
@@ -7858,6 +7957,34 @@ export interface components {
/** @description Sink name */
sink_name: string
}
UpdateReplicationSinkPipelineBody: {
/** @description Pipeline config */
pipeline_config: {
config: {
max_fill_secs: number
max_size: number
}
}
/** @description Publication name */
publication_name: string
/** @description Sink config */
sink_config: {
big_query: {
/** @description BigQuery dataset id */
dataset_id: string
/** @description Max staleness in minutes */
max_staleness_mins: number
/** @description BigQuery project id */
project_id: string
/** @description BigQuery service account key */
service_account_key: string
}
}
/** @description Sink name */
sink_name: string
/** @description Source id */
source_id: number
}
UpdateSchemaBody: {
name?: string
owner?: string
@@ -17475,6 +17602,87 @@ export interface operations {
}
}
}
ReplicationSinksPipelinesController_createSinkPipeline: {
parameters: {
query?: never
header?: never
path: {
/** @description Project ref */
ref: string
}
cookie?: never
}
requestBody: {
content: {
'application/json': components['schemas']['CreateReplicationSinkPipelineBody']
}
}
responses: {
/** @description Returns the created replication sink and pipeline ids. */
201: {
headers: {
[name: string]: unknown
}
content: {
'application/json': components['schemas']['CreateSinkPipelineResponse']
}
}
403: {
headers: {
[name: string]: unknown
}
content?: never
}
/** @description Returned when the API fails to create the replication sink or pipeline. */
500: {
headers: {
[name: string]: unknown
}
content?: never
}
}
}
ReplicationSinksPipelinesController_updateSinkPipeline: {
parameters: {
query?: never
header?: never
path: {
/** @description Pipeline id */
pipeline_id: number
/** @description Project reference */
ref: string
/** @description Sink id */
sink_id: number
}
cookie?: never
}
requestBody: {
content: {
'application/json': components['schemas']['UpdateReplicationSinkPipelineBody']
}
}
responses: {
201: {
headers: {
[name: string]: unknown
}
content?: never
}
403: {
headers: {
[name: string]: unknown
}
content?: never
}
/** @description Returned when the API fails to update the replication sink or pipeline */
500: {
headers: {
[name: string]: unknown
}
content?: never
}
}
}
ReplicationSinksController_getSink: {
parameters: {
query?: never
@@ -17818,6 +18026,42 @@ export interface operations {
}
}
}
ReplicationTenantsSourcesController_createTenantSource: {
parameters: {
query?: never
header?: never
path: {
/** @description Project ref */
ref: string
}
cookie?: never
}
requestBody?: never
responses: {
/** @description Returns the created replication tenant and source ids. */
201: {
headers: {
[name: string]: unknown
}
content: {
'application/json': components['schemas']['CreateTenantSourceResponse']
}
}
403: {
headers: {
[name: string]: unknown
}
content?: never
}
/** @description Returned when the API fails to create the replication tenant or source. */
500: {
headers: {
[name: string]: unknown
}
content?: never
}
}
}
ResetPasswordController_resetPassword: {
parameters: {
query?: never