feat(replication): Significantly improve the replication UI behavior (#38237)

This commit is contained in:
Riccardo Busetti authored and GitHub committed 2025-08-29 12:20:51 +02:00
1 parent a448877f24
commit 77356bf946
20 files changed
+1234 -845

No files matched your search

@@ -2,40 +2,32 @@ import TextConfirmModal from 'ui-patterns/Dialogs/TextConfirmModal'
interface DeleteDestinationProps {
visible: boolean
setVisible: (value: boolean) => void
onDelete: () => void
isLoading: boolean
name: string
setVisible: (value: boolean) => void
onDelete: () => void
}
const DeleteDestination = ({
export const DeleteDestination = ({
visible,
setVisible,
onDelete,
isLoading,
name,
setVisible,
onDelete,
}: DeleteDestinationProps) => {
return (
<>
<TextConfirmModal
variant={'warning'}
visible={visible}
onCancel={() => setVisible(!visible)}
onConfirm={onDelete}
title="Delete this destination"
loading={isLoading}
confirmLabel={`Delete destination`}
confirmPlaceholder="Type in name of destination"
confirmString={name ?? 'Unknown'}
text={
<>
<span>This will delete the destination</span>{' '}
</>
}
alert={{ title: 'You cannot recover this destination once deleted.' }}
/>
</>
<TextConfirmModal
variant="destructive"
visible={visible}
loading={isLoading}
title="Delete this destination"
confirmLabel={isLoading ? 'Deleting…' : `Delete destination`}
confirmPlaceholder="Type in name of destination"
confirmString={name ?? 'Unknown'}
text={`This will delete the destination "${name}"`}
alert={{ title: 'You cannot recover this destination once deleted.' }}
onCancel={() => setVisible(!visible)}
onConfirm={onDelete}
/>
)
}
export default DeleteDestination
@@ -1,10 +1,21 @@
import { zodResolver } from '@hookform/resolvers/zod'
import { useParams } from 'common'
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 { useEffect, useMemo, useState } from 'react'
import { useForm } from 'react-hook-form'
import { toast } from 'sonner'
import * as z from 'zod'
import { useParams } from 'common'
import { useCreateDestinationPipelineMutation } from 'data/replication/create-destination-pipeline-mutation'
import { useCreateTenantSourceMutation } from 'data/replication/create-tenant-source-mutation'
import { useReplicationDestinationByIdQuery } from 'data/replication/destination-by-id-query'
import { useReplicationPipelineByIdQuery } from 'data/replication/pipeline-by-id-query'
import { useReplicationPublicationsQuery } from 'data/replication/publications-query'
import { useStartPipelineMutation } from 'data/replication/start-pipeline-mutation'
import { useUpdateDestinationPipelineMutation } from 'data/replication/update-destination-pipeline-mutation'
import {
PipelineStatusRequestStatus,
usePipelineRequestStatus,
} from 'state/replication-pipeline-request-status'
import {
Accordion_Shadcn_,
AccordionContent_Shadcn_,
@@ -23,6 +34,7 @@ import {
SelectGroup_Shadcn_,
SelectItem_Shadcn_,
SelectTrigger_Shadcn_,
Separator,
Sheet,
SheetContent,
SheetDescription,
@@ -30,21 +42,28 @@ import {
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 { useReplicationDestinationByIdQuery } from 'data/replication/destination-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 { useCreateDestinationPipelineMutation } from 'data/replication/create-destination-pipeline-mutation'
import { useUpdateDestinationPipelineMutation } from 'data/replication/update-destination-pipeline-mutation'
import NewPublicationPanel from './NewPublicationPanel'
import PublicationsComboBox from './PublicationsComboBox'
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().optional(),
maxFillMs: z.number().min(1, 'Max Fill milliseconds should be greater than 0').int().optional(),
maxStalenessMins: z.number().nonnegative().optional(),
})
interface DestinationPanelProps {
visible: boolean
@@ -58,22 +77,33 @@ interface DestinationPanelProps {
}
}
const DestinationPanel = ({
export const DestinationPanel = ({
visible,
sourceId,
onClose,
existingDestination,
}: DestinationPanelProps) => {
const { ref: projectRef } = useParams()
const { setRequestStatus } = usePipelineRequestStatus()
const editMode = !!existingDestination
const [publicationPanelVisible, setPublicationPanelVisible] = useState(false)
const { mutateAsync: createTenantSource, isLoading: creatingTenantSource } =
useCreateTenantSourceMutation()
const { mutateAsync: createDestinationPipeline, isLoading: creatingDestinationPipeline } =
useCreateDestinationPipelineMutation()
const { mutateAsync: startPipeline, isLoading: startingPipeline } = useStartPipelineMutation()
const { mutateAsync: stopPipeline, isLoading: stoppingPipeline } = useStopPipelineMutation()
useCreateDestinationPipelineMutation({
onSuccess: () => form.reset(defaultValues),
})
const { mutateAsync: updateDestinationPipeline, isLoading: updatingDestinationPipeline } =
useUpdateDestinationPipelineMutation()
useUpdateDestinationPipelineMutation({
onSuccess: () => form.reset(defaultValues),
})
const { mutateAsync: startPipeline, isLoading: startingPipeline } = useStartPipelineMutation()
const { data: publications, isLoading: loadingPublications } = useReplicationPublicationsQuery({
projectRef,
sourceId,
@@ -89,26 +119,6 @@ const DestinationPanel = ({
pipelineId: existingDestination?.pipelineId,
})
const isCreating = creatingTenantSource || creatingDestinationPipeline || startingPipeline
const isUpdating = updatingDestinationPipeline || 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(),
maxFillMs: z.number().min(1, 'Max Fill milliseconds should be greater than 0').int(),
maxStalenessMins: z.number().nonnegative(),
enabled: z.boolean(),
})
const defaultValues = useMemo(
() => ({
type: TypeEnum.enum.BigQuery,
@@ -118,403 +128,406 @@ const DestinationPanel = ({
// For now, the password will always be set as empty for security reasons.
serviceAccountKey: destinationData?.config?.big_query?.service_account_key ?? '',
publicationName: pipelineData?.config.publication_name ?? '',
maxSize: pipelineData?.config?.batch?.max_size ?? 1000,
maxFillMs: pipelineData?.config?.batch?.max_fill_ms ?? 10,
maxStalenessMins: destinationData?.config?.big_query?.max_staleness_mins ?? 5,
enabled: existingDestination?.enabled ?? true,
maxSize: pipelineData?.config?.batch?.max_size,
maxFillMs: pipelineData?.config?.batch?.max_fill_ms,
maxStalenessMins: destinationData?.config?.big_query?.max_staleness_mins,
}),
[destinationData, pipelineData, existingDestination]
[destinationData, pipelineData]
)
const form = useForm<z.infer<typeof FormSchema>>({
mode: 'onBlur',
reValidateMode: 'onBlur',
resolver: zodResolver(FormSchema),
defaultValues,
})
const isSaving = creatingDestinationPipeline || updatingDestinationPipeline || startingPipeline
const onSubmit = async (data: z.infer<typeof FormSchema>) => {
if (!projectRef) return console.error('Project ref is required')
if (!sourceId) return console.error('Source id is required')
try {
if (editMode && existingDestination) {
if (!sourceId) {
console.error('Source id is required')
return
if (!existingDestination.pipelineId) return console.error('Pipeline id is required')
const bigQueryConfig: any = {
projectId: data.projectId,
datasetId: data.datasetId,
serviceAccountKey: data.serviceAccountKey,
}
if (!existingDestination.pipelineId) {
console.error('Pipeline id is required')
return
if (!!data.maxStalenessMins) {
bigQueryConfig.maxStalenessMins = data.maxStalenessMins
}
// Update existing destination
const batchConfig: any = {}
if (!!data.maxSize) batchConfig.maxSize = data.maxSize
if (!!data.maxFillMs) batchConfig.maxFillMs = data.maxFillMs
const hasBothBatchFields = Object.keys(batchConfig).length === 2
await updateDestinationPipeline({
destinationId: existingDestination.destinationId,
pipelineId: existingDestination.pipelineId,
projectRef,
destinationName: data.name,
destinationConfig: {
bigQuery: {
projectId: data.projectId,
datasetId: data.datasetId,
serviceAccountKey: data.serviceAccountKey,
maxStalenessMins: data.maxStalenessMins,
},
},
destinationConfig: { bigQuery: bigQueryConfig },
pipelineConfig: {
publicationName: data.publicationName,
batch: { maxSize: data.maxSize, maxFillMs: data.maxFillMs },
...(hasBothBatchFields ? { batch: batchConfig } : {}),
},
sourceId,
})
if (data.enabled) {
await startPipeline({ projectRef, pipelineId: existingDestination.pipelineId })
// Set request status only right before starting, then fire and close
const snapshot = existingDestination.enabled ? 'started' : 'stopped'
if (existingDestination.enabled) {
setRequestStatus(
existingDestination.pipelineId,
PipelineStatusRequestStatus.RestartRequested,
snapshot
)
toast.success('Settings applied. Restarting the pipeline...')
} else {
await stopPipeline({ projectRef, pipelineId: existingDestination.pipelineId })
setRequestStatus(
existingDestination.pipelineId,
PipelineStatusRequestStatus.StartRequested,
snapshot
)
toast.success('Settings applied. Starting the pipeline...')
}
startPipeline({ projectRef, pipelineId: existingDestination.pipelineId })
onClose()
} else {
const bigQueryConfig: any = {
projectId: data.projectId,
datasetId: data.datasetId,
serviceAccountKey: data.serviceAccountKey,
}
if (!!data.maxStalenessMins) {
bigQueryConfig.maxStalenessMins = data.maxStalenessMins
}
toast.success('Successfully updated destination')
} else {
// Create new destination
if (!sourceId) {
console.error('Source id is required')
return
}
const batchConfig: any = {}
if (!!data.maxSize) batchConfig.maxSize = data.maxSize
if (!!data.maxFillMs) batchConfig.maxFillMs = data.maxFillMs
const hasBothBatchFields = Object.keys(batchConfig).length === 2
const { pipeline_id: pipelineId } = await createDestinationPipeline({
projectRef,
destinationName: data.name,
destinationConfig: {
bigQuery: {
projectId: data.projectId,
datasetId: data.datasetId,
serviceAccountKey: data.serviceAccountKey,
maxStalenessMins: data.maxStalenessMins,
},
},
destinationConfig: { bigQuery: bigQueryConfig },
sourceId,
pipelineConfig: {
publicationName: data.publicationName,
batch: { maxSize: data.maxSize, maxFillMs: data.maxFillMs },
...(hasBothBatchFields ? { batch: batchConfig } : {}),
},
})
if (data.enabled) {
await startPipeline({ projectRef, pipelineId })
}
toast.success('Successfully created destination')
// Set request status only right before starting, then fire and close
setRequestStatus(pipelineId, PipelineStatusRequestStatus.StartRequested, undefined)
toast.success('Destination created. Starting the pipeline...')
startPipeline({ projectRef, pipelineId })
onClose()
}
onClose()
} catch (error) {
toast.error(`Failed to ${editMode ? 'update' : 'create'} destination`)
const action = editMode ? 'apply and run' : 'create and start'
toast.error(`Failed to ${action} destination`)
}
}
const onEnableReplication = async () => {
if (!projectRef) return console.error('Project ref is required')
await createTenantSource({ projectRef })
}
const { enabled } = form.watch()
useEffect(() => {
if (editMode && destinationData && pipelineData) {
form.reset(defaultValues)
}
}, [destinationData, pipelineData, editMode, defaultValues, form])
return (
return sourceId ? (
<>
{sourceId ? (
<>
<Sheet open={visible} onOpenChange={onClose}>
<SheetContent showClose={false} size="default">
<div className="flex flex-col h-full" tabIndex={-1}>
<SheetHeader className="flex justify-between items-center">
<div>
<SheetTitle>{editMode ? 'Edit Destination' : 'New Destination'}</SheetTitle>
<SheetDescription>
{editMode ? null : 'Send data to a new destination'}
</SheetDescription>
</div>
<div className="flex items-center gap-2">
<Switch
checked={enabled}
onCheckedChange={(checked) => {
form.setValue('enabled', checked)
}}
<Sheet open={visible} onOpenChange={onClose}>
<SheetContent showClose={false} size="default">
<div className="flex flex-col h-full" tabIndex={-1}>
<SheetHeader>
<SheetTitle>{editMode ? 'Edit destination' : 'Create a new destination'}</SheetTitle>
<SheetDescription>
{editMode ? null : 'Send data to a new destination'}
</SheetDescription>
</SheetHeader>
<SheetSection className="flex-grow overflow-auto px-0 pb-0">
<Form_Shadcn_ {...form}>
<form id={formId} onSubmit={form.handleSubmit(onSubmit)}>
<div className="px-5 pb-4">
<FormField_Shadcn_
control={form.control}
name="name"
render={({ field }) => (
<FormItemLayout
className="mb-8"
label="Name"
layout="vertical"
description="A name you will use to identify this destination"
>
<FormControl_Shadcn_>
<Input_Shadcn_ {...field} placeholder="Name" />
</FormControl_Shadcn_>
</FormItemLayout>
)}
/>
<h3 className="mb-4">What data to send</h3>
<FormField_Shadcn_
control={form.control}
name="publicationName"
render={({ field }) => (
<FormItemLayout
className="mb-4"
label="Publication"
layout="vertical"
description="A publication is a collection of tables that you want to replicate "
>
<FormControl_Shadcn_>
<PublicationsComboBox
publications={publications?.map((pub) => pub.name) || []}
loading={loadingPublications}
field={field}
onNewPublicationClick={() => setPublicationPanelVisible(true)}
/>
</FormControl_Shadcn_>
</FormItemLayout>
)}
/>
<h3 className="mb-4 mt-8">Where to send that data</h3>
<FormField_Shadcn_
name="type"
control={form.control}
render={({ field }) => (
<FormItemLayout
className="mb-4"
label="Type"
layout="vertical"
description="The type of destination to send the data to"
>
<FormControl_Shadcn_>
<Select_Shadcn_ value={field.value}>
<SelectTrigger_Shadcn_>{field.value}</SelectTrigger_Shadcn_>
<SelectContent_Shadcn_>
<SelectGroup_Shadcn_>
<SelectItem_Shadcn_ value="BigQuery">BigQuery</SelectItem_Shadcn_>
</SelectGroup_Shadcn_>
</SelectContent_Shadcn_>
</Select_Shadcn_>
</FormControl_Shadcn_>
</FormItemLayout>
)}
/>
<FormField_Shadcn_
control={form.control}
name="projectId"
render={({ field }) => (
<FormItemLayout
className="mb-4"
label="Project ID"
layout="vertical"
description="Which BigQuery project to send data to"
>
<FormControl_Shadcn_>
<Input_Shadcn_ {...field} placeholder="Project ID" />
</FormControl_Shadcn_>
</FormItemLayout>
)}
/>
<FormField_Shadcn_
control={form.control}
name="datasetId"
render={({ field }) => (
<FormItemLayout
className="mb-4"
label="Project's Dataset ID"
layout="vertical"
>
<FormControl_Shadcn_>
<Input_Shadcn_ {...field} placeholder="Dataset ID" />
</FormControl_Shadcn_>
</FormItemLayout>
)}
/>
<FormField_Shadcn_
control={form.control}
name="serviceAccountKey"
render={({ field }) => (
<FormItemLayout
label="Service Account Key"
layout="vertical"
description="The service account key for BigQuery"
>
<FormControl_Shadcn_>
<TextArea_Shadcn_
{...field}
rows={4}
maxLength={5000}
placeholder="Service account key"
/>
</FormControl_Shadcn_>
</FormItemLayout>
)}
/>
<Label className="text-sm mx-2">Enable</Label>
</div>
</SheetHeader>
<SheetSection className="flex-grow overflow-auto">
<Form_Shadcn_ {...form}>
<form id={formId} onSubmit={form.handleSubmit(onSubmit)}>
<FormField_Shadcn_
control={form.control}
name="name"
render={({ field }) => (
<FormItemLayout
className="mb-8"
label="Name"
layout="vertical"
description="A name you will use to identify this destination"
>
<FormControl_Shadcn_>
<Input_Shadcn_ {...field} placeholder="Name" />
</FormControl_Shadcn_>
</FormItemLayout>
)}
/>
<h3 className="mb-4">What data to send</h3>
<FormField_Shadcn_
control={form.control}
name="publicationName"
render={({ field }) => (
<FormItemLayout
className="mb-4"
label="Publication"
layout="vertical"
description="A publication is a collection of tables that you want to replicate "
>
<FormControl_Shadcn_>
<PublicationsComboBox
publications={publications?.map((pub) => pub.name) || []}
loading={loadingPublications}
field={field}
onNewPublicationClick={() => setPublicationPanelVisible(true)}
/>
</FormControl_Shadcn_>
</FormItemLayout>
)}
/>
<h3 className="mb-4 mt-8">Where to send that data</h3>
<Separator />
<FormField_Shadcn_
name="type"
control={form.control}
render={({ field }) => (
<FormItemLayout
className="mb-4"
label="Type"
layout="vertical"
description="The type of destination to send the data to"
>
<FormControl_Shadcn_>
<Select_Shadcn_ value={field.value}>
<SelectTrigger_Shadcn_>{field.value}</SelectTrigger_Shadcn_>
<SelectContent_Shadcn_>
<SelectGroup_Shadcn_>
<SelectItem_Shadcn_ value="BigQuery">
BigQuery
</SelectItem_Shadcn_>
</SelectGroup_Shadcn_>
</SelectContent_Shadcn_>
</Select_Shadcn_>
</FormControl_Shadcn_>
</FormItemLayout>
)}
></FormField_Shadcn_>
<FormField_Shadcn_
control={form.control}
name="projectId"
render={({ field }) => (
<FormItemLayout
className="mb-4"
label="Project Id"
layout="vertical"
description="Which BigQuery project to send data to"
>
<FormControl_Shadcn_>
<Input_Shadcn_ {...field} placeholder="Project id" />
</FormControl_Shadcn_>
</FormItemLayout>
)}
/>
<FormField_Shadcn_
control={form.control}
name="datasetId"
render={({ field }) => (
<FormItemLayout
className="mb-4"
label="Project's Dataset Id"
layout="vertical"
>
<FormControl_Shadcn_>
<Input_Shadcn_ {...field} placeholder="Dataset id" />
</FormControl_Shadcn_>
</FormItemLayout>
)}
/>
<FormField_Shadcn_
control={form.control}
name="serviceAccountKey"
render={({ field }) => (
<FormItemLayout
label="Service Account Key"
layout="vertical"
description="The service account key for BigQuery"
>
<FormControl_Shadcn_>
<TextArea_Shadcn_
{...field}
rows={4}
maxLength={5000}
placeholder="Service account key"
/>
</FormControl_Shadcn_>
</FormItemLayout>
)}
/>
<Accordion_Shadcn_ type="single" collapsible>
<AccordionItem_Shadcn_ value="item-1" className="border-none">
<AccordionTrigger_Shadcn_ className="font-normal gap-2 justify-start mb-0 mt-8">
Advanced Settings
</AccordionTrigger_Shadcn_>
<AccordionContent_Shadcn_ asChild className="!pb-0">
<FormField_Shadcn_
control={form.control}
name="maxSize"
render={({ field }) => (
<FormItemLayout
className="mb-4"
label="Max Size"
layout="vertical"
description="The maximum size of the data to send"
>
<FormControl_Shadcn_>
<Input_Shadcn_
{...field}
type="number"
{...form.register('maxSize', {
valueAsNumber: true, // Ensure the value is handled as a number
})}
placeholder="Max size"
/>
</FormControl_Shadcn_>
</FormItemLayout>
)}
/>
<FormField_Shadcn_
control={form.control}
name="maxFillMs"
render={({ field }) => (
<FormItemLayout
className="mb-4"
label="Max Fill Seconds"
layout="vertical"
description="The maximum amount of time to fill the data"
>
<FormControl_Shadcn_>
<Input_Shadcn_
{...field}
type="number"
{...form.register('maxFillMs', {
valueAsNumber: true, // Ensure the value is handled as a number
})}
placeholder="Max fill seconds"
/>
</FormControl_Shadcn_>
</FormItemLayout>
)}
/>
<FormField_Shadcn_
control={form.control}
name="maxStalenessMins"
render={({ field }) => (
<FormItemLayout
className="mb-4"
label="Max Staleness"
layout="vertical"
description="Maximum staleness time allowed"
>
<FormControl_Shadcn_>
<Input_Shadcn_
{...field}
type="number"
{...form.register('maxStalenessMins', {
valueAsNumber: true, // Ensure the value is handled as a number
})}
placeholder="Max staleness in minutes"
/>
</FormControl_Shadcn_>
</FormItemLayout>
)}
/>
</AccordionContent_Shadcn_>
</AccordionItem_Shadcn_>
</Accordion_Shadcn_>
<div className="hidden">
<FormField_Shadcn_
control={form.control}
name="enabled"
render={({ field }) => (
<FormItemLayout className="mb-4" layout="vertical" label="Enabled">
<FormControl_Shadcn_>
<Switch checked={field.value} onCheckedChange={field.onChange} />
</FormControl_Shadcn_>
</FormItemLayout>
)}
/>
</div>
</form>
</Form_Shadcn_>
</SheetSection>
<SheetFooter>
<Button disabled={isSubmitting} type="default" onClick={onClose}>
Cancel
</Button>
<Button
disabled={isSubmitting}
loading={isSubmitting}
form={formId}
htmlType="submit"
>
{editMode ? 'Update' : 'Create'}
</Button>
</SheetFooter>
</div>
</SheetContent>
</Sheet>
<NewPublicationPanel
visible={publicationPanelVisible}
sourceId={sourceId}
onClose={() => setPublicationPanelVisible(false)}
/>
</>
) : (
<>
<Sheet open={visible} onOpenChange={onClose}>
<SheetContent showClose={false} size="default">
<div className="flex flex-col h-full" tabIndex={-1}>
<SheetHeader>
<SheetTitle>New Destination</SheetTitle>
</SheetHeader>
<SheetSection className="flex-grow overflow-auto">
<Alert_Shadcn_>
<WarningIcon />
<AlertTitle_Shadcn_>
{/* Pricing to be decided yet */}
Enabling replication will cost additional $xx.xx
</AlertTitle_Shadcn_>
<AlertDescription_Shadcn_>
<span></span>
<div className="flex items-center gap-x-2 mt-3">
<Button type="default" onClick={onEnableReplication}>
Enable replication
</Button>
</div>
</AlertDescription_Shadcn_>
</Alert_Shadcn_>
</SheetSection>
<SheetFooter>
<Button disabled={isSubmitting} type="default" onClick={onClose}>
Cancel
</Button>
</SheetFooter>
</div>
</SheetContent>
</Sheet>
</>
)}
<div className="px-5">
<Accordion_Shadcn_ type="single" collapsible>
<AccordionItem_Shadcn_ value="item-1" className="border-none">
<AccordionTrigger_Shadcn_ className="font-normal gap-2 justify-between text-sm">
Advanced Settings
</AccordionTrigger_Shadcn_>
<AccordionContent_Shadcn_ asChild className="!pb-0">
<FormField_Shadcn_
control={form.control}
name="maxSize"
render={({ field }) => (
<FormItemLayout
className="mb-4"
label="Max size"
layout="vertical"
description="The maximum size of the data to send. Leave empty to use default value."
>
<FormControl_Shadcn_>
<Input_Shadcn_
{...field}
type="number"
value={field.value ?? ''}
onChange={(e) => {
const val = e.target.value
field.onChange(val === '' ? undefined : Number(val))
}}
placeholder="Leave empty for default"
/>
</FormControl_Shadcn_>
</FormItemLayout>
)}
/>
<FormField_Shadcn_
control={form.control}
name="maxFillMs"
render={({ field }) => (
<FormItemLayout
className="mb-4"
label="Max fill milliseconds"
layout="vertical"
description="The maximum amount of time to fill the data in milliseconds. Leave empty to use default value."
>
<FormControl_Shadcn_>
<Input_Shadcn_
{...field}
type="number"
value={field.value ?? ''}
onChange={(e) => {
const val = e.target.value
field.onChange(val === '' ? undefined : Number(val))
}}
placeholder="Leave empty for default"
/>
</FormControl_Shadcn_>
</FormItemLayout>
)}
/>
<FormField_Shadcn_
control={form.control}
name="maxStalenessMins"
render={({ field }) => (
<FormItemLayout
className="mb-4"
label="Max staleness minutes"
layout="vertical"
description="Maximum staleness time allowed in minutes. Leave empty to use default value."
>
<FormControl_Shadcn_>
<Input_Shadcn_
{...field}
type="number"
value={field.value ?? ''}
onChange={(e) => {
const val = e.target.value
field.onChange(val === '' ? undefined : Number(val))
}}
placeholder="Leave empty for default"
/>
</FormControl_Shadcn_>
</FormItemLayout>
)}
/>
</AccordionContent_Shadcn_>
</AccordionItem_Shadcn_>
</Accordion_Shadcn_>
</div>
</form>
</Form_Shadcn_>
</SheetSection>
<SheetFooter>
<Button disabled={isSaving} type="default" onClick={onClose}>
Cancel
</Button>
<Button loading={isSaving} form={formId} htmlType="submit">
{editMode
? existingDestination?.enabled
? 'Apply and restart'
: 'Apply and start'
: 'Create and start'}
</Button>
</SheetFooter>
</div>
</SheetContent>
</Sheet>
<NewPublicationPanel
visible={publicationPanelVisible}
sourceId={sourceId}
onClose={() => setPublicationPanelVisible(false)}
/>
</>
) : (
<Sheet open={visible} onOpenChange={onClose}>
<SheetContent showClose={false} size="default">
<div className="flex flex-col h-full" tabIndex={-1}>
<SheetHeader>
<SheetTitle>Create a new destination</SheetTitle>
</SheetHeader>
<SheetSection className="flex-grow overflow-auto">
<Alert_Shadcn_>
<WarningIcon />
<AlertTitle_Shadcn_>
{/* Pricing to be decided yet */}
Enabling replication will cost additional $xx.xx
</AlertTitle_Shadcn_>
<AlertDescription_Shadcn_>
<span></span>
<div className="flex items-center gap-x-2 mt-3">
<Button
type="default"
loading={creatingTenantSource}
onClick={onEnableReplication}
>
Enable replication
</Button>
</div>
</AlertDescription_Shadcn_>
</Alert_Shadcn_>
</SheetSection>
<SheetFooter>
<Button disabled={creatingTenantSource} type="default" onClick={onClose}>
Cancel
</Button>
</SheetFooter>
</div>
</SheetContent>
</Sheet>
)
}
export default DestinationPanel
@@ -6,18 +6,20 @@ import { useParams } from 'common'
import Table from 'components/to-be-cleaned/Table'
import AlertError from 'components/ui/AlertError'
import { useDeleteDestinationPipelineMutation } from 'data/replication/delete-destination-pipeline-mutation'
import { useReplicationPipelineReplicationStatusQuery } from 'data/replication/pipeline-replication-status-query'
import { useReplicationPipelineStatusQuery } from 'data/replication/pipeline-status-query'
import { Pipeline } from 'data/replication/pipelines-query'
import { useStopPipelineMutation } from 'data/replication/stop-pipeline-mutation'
import { AlertCircle } from 'lucide-react'
import {
PipelineStatusRequestStatus,
usePipelineRequestStatus,
} from 'state/replication-pipeline-request-status'
import { ResponseError } from 'types'
import { Button } from 'ui'
import { Button, Tooltip, TooltipContent, TooltipTrigger } from 'ui'
import ShimmeringLoader from 'ui-patterns/ShimmeringLoader'
import DeleteDestination from './DeleteDestination'
import DestinationPanel from './DestinationPanel'
import { DeleteDestination } from './DeleteDestination'
import { DestinationPanel } from './DestinationPanel'
import { getStatusName, PIPELINE_ERROR_MESSAGES } from './Pipeline.utils'
import { PipelineStatus, PipelineStatusName } from './PipelineStatus'
import { STATUS_REFRESH_FREQUENCY_MS } from './Replication.constants'
@@ -48,6 +50,7 @@ export const DestinationRow = ({
}: DestinationRowProps) => {
const { ref: projectRef } = useParams()
const [showDeleteDestinationForm, setShowDeleteDestinationForm] = useState(false)
const [isDeleting, setIsDeleting] = useState(false)
const [showEditDestinationPanel, setShowEditDestinationPanel] = useState(false)
const {
@@ -74,6 +77,15 @@ export const DestinationRow = ({
const pipelineStatus = pipelineStatusData?.status
const statusName = getStatusName(pipelineStatus)
// Fetch table-level replication status to surface errors in list view
const { data: replicationStatusData } = useReplicationPipelineReplicationStatusQuery(
{ projectRef, pipelineId: pipeline?.id },
{ refetchInterval: STATUS_REFRESH_FREQUENCY_MS }
)
const tableStatuses = replicationStatusData?.table_statuses ?? []
const errorCount = tableStatuses.filter((t) => t.state?.name === 'error').length
const hasTableErrors = errorCount > 0
const onDeleteClick = async () => {
if (!projectRef) {
return console.error('Project ref is required')
@@ -83,14 +95,20 @@ export const DestinationRow = ({
}
try {
setIsDeleting(true)
await stopPipeline({ projectRef, pipelineId: pipeline.id })
await deleteDestinationPipeline({
projectRef,
destinationId: destinationId,
pipelineId: pipeline.id,
})
// Close dialog after successful deletion
setShowDeleteDestinationForm(false)
toast.success(`Deleted destination "${destinationName}"`)
} catch (error) {
toast.error(PIPELINE_ERROR_MESSAGES.DELETE_DESTINATION)
} finally {
setIsDeleting(false)
}
}
@@ -132,9 +150,21 @@ export const DestinationRow = ({
</Table.td>
<Table.td>
<div className="flex items-center justify-end gap-x-2">
<Button asChild type="default">
<Button asChild type="default" className="relative">
<Link href={`/project/${projectRef}/database/replication/${pipeline?.id}`}>
View status
<span className="inline-flex items-center gap-2">
<span>View status</span>
{hasTableErrors && (
<Tooltip>
<TooltipTrigger>
<AlertCircle size={14} />
</TooltipTrigger>
<TooltipContent side="bottom">
{errorCount} table{errorCount === 1 ? '' : 's'} have replication errors
</TooltipContent>
</Tooltip>
)}
</span>
</Link>
</Button>
<RowMenu
@@ -154,7 +184,7 @@ export const DestinationRow = ({
visible={showDeleteDestinationForm}
setVisible={setShowDeleteDestinationForm}
onDelete={onDeleteClick}
isLoading={isPipelineStatusLoading}
isLoading={isDeleting}
name={destinationName}
/>
<DestinationPanel
@@ -9,7 +9,7 @@ import { useReplicationPipelinesQuery } from 'data/replication/pipelines-query'
import { useReplicationSourcesQuery } from 'data/replication/sources-query'
import { Button, cn, Input_Shadcn_ } from 'ui'
import { GenericSkeletonLoader } from 'ui-patterns'
import NewDestinationPanel from './DestinationPanel'
import { DestinationPanel } from './DestinationPanel'
import { DestinationRow } from './DestinationRow'
import { PIPELINE_ERROR_MESSAGES } from './Pipeline.utils'
@@ -159,7 +159,7 @@ export const Destinations = () => {
</div>
)}
<NewDestinationPanel
<DestinationPanel
visible={showNewDestinationPanel}
sourceId={sourceId}
onClose={() => setShowNewDestinationPanel(false)}
@@ -22,16 +22,25 @@ export const getStatusName = (
return undefined
}
export const PIPELINE_ENABLE_ALLOWED_FROM = ['stopped'] as const
export const PIPELINE_DISABLE_ALLOWED_FROM = ['started', 'failed'] as const
export const PIPELINE_ACTIONABLE_STATES = ['failed', 'started', 'stopped'] as const
const PIPELINE_STATE_MESSAGES = {
enabling: {
title: 'Pipeline enabling',
message: 'Starting the pipeline. Table replication will resume once enabled.',
badge: 'Enabling',
title: 'Starting pipeline',
message: 'Starting the pipeline. Table replication will resume once running.',
badge: 'Starting',
},
disabling: {
title: 'Pipeline disabling',
message: 'Stopping the pipeline. Table replication will be paused once disabled.',
badge: 'Disabling',
title: 'Stopping pipeline',
message: 'Stopping the pipeline. Table replication will be paused once stopped.',
badge: 'Stopping',
},
restarting: {
title: 'Restarting pipeline',
message: 'Applying settings and restarting the pipeline.',
badge: 'Restarting',
},
failed: {
title: 'Pipeline failed',
@@ -40,7 +49,7 @@ const PIPELINE_STATE_MESSAGES = {
},
stopped: {
title: 'Pipeline stopped',
message: 'Replication is paused. Enable the pipeline to resume data synchronization.',
message: 'Replication is paused. Start the pipeline to resume data synchronization.',
badge: 'Stopped',
},
starting: {
@@ -60,8 +69,8 @@ const PIPELINE_STATE_MESSAGES = {
},
notRunning: {
title: 'Pipeline not running',
message: 'Replication is not active. Enable the pipeline to start data synchronization.',
badge: 'Disabled',
message: 'Replication is not active. Start the pipeline to begin data synchronization.',
badge: 'Stopped',
},
} as const
@@ -69,23 +78,25 @@ export const getPipelineStateMessages = (
requestStatus: PipelineStatusRequestStatus | undefined,
statusName: string | undefined
) => {
// Always prioritize request status (enabling/disabling) over pipeline status
if (requestStatus === PipelineStatusRequestStatus.EnableRequested) {
// Reflect optimistic request intent immediately after click
if (requestStatus === PipelineStatusRequestStatus.RestartRequested) {
return PIPELINE_STATE_MESSAGES.restarting
}
if (requestStatus === PipelineStatusRequestStatus.StartRequested) {
return PIPELINE_STATE_MESSAGES.enabling
}
if (requestStatus === PipelineStatusRequestStatus.DisableRequested) {
if (requestStatus === PipelineStatusRequestStatus.StopRequested) {
return PIPELINE_STATE_MESSAGES.disabling
}
// Only check pipeline status if no request is in progress
// Fall back to steady states
switch (statusName) {
case 'starting':
return PIPELINE_STATE_MESSAGES.starting
case 'failed':
return PIPELINE_STATE_MESSAGES.failed
case 'stopped':
return PIPELINE_STATE_MESSAGES.stopped
case 'starting':
return PIPELINE_STATE_MESSAGES.starting
case 'started':
return PIPELINE_STATE_MESSAGES.running
case 'unknown':
@@ -45,18 +45,26 @@ export const PipelineStatus = ({
// Get consistent tooltip message using the same logic as other components
const stateMessages = getPipelineStateMessages(requestStatus, statusName)
if (requestStatus === PipelineStatusRequestStatus.EnableRequested) {
// Show optimistic request state while backend still reports steady states
if (requestStatus === PipelineStatusRequestStatus.RestartRequested) {
return {
label: 'Enabling...',
dot: <Loader2 className="animate-spin w-3 h-3 text-brand-600" />,
color: 'text-brand-600',
label: 'Restarting',
dot: <Loader2 className="animate-spin w-3 h-3 text-warning-600" />,
color: 'text-warning-600',
tooltip: stateMessages.message,
}
}
if (requestStatus === PipelineStatusRequestStatus.DisableRequested) {
if (requestStatus === PipelineStatusRequestStatus.StartRequested) {
return {
label: 'Disabling...',
label: 'Starting',
dot: <Loader2 className="animate-spin w-3 h-3 text-warning-600" />,
color: 'text-warning-600',
tooltip: stateMessages.message,
}
}
if (requestStatus === PipelineStatusRequestStatus.StopRequested) {
return {
label: 'Stopping',
dot: <Loader2 className="animate-spin w-3 h-3 text-warning-600" />,
color: 'text-warning-600',
tooltip: stateMessages.message,
@@ -1,7 +1,3 @@
// @ts-nocheck [Joshen] Temporarily silencing the TS checks here but please eventually remove
// it's cause the API types are conflicting a bit - API types have probably been updated for the UI here
// but the UI hasn't been updated yet to fit the new API types
import { Activity, ChevronLeft, ExternalLink, Search, X } from 'lucide-react'
import Link from 'next/link'
import { useEffect, useState } from 'react'
@@ -24,7 +20,13 @@ import { Badge, Button, cn } from 'ui'
import { GenericSkeletonLoader } from 'ui-patterns'
import { Input } from 'ui-patterns/DataInputs/Input'
import { ErroredTableDetails } from './ErroredTableDetails'
import { getStatusName, PIPELINE_ERROR_MESSAGES } from './Pipeline.utils'
import {
PIPELINE_ACTIONABLE_STATES,
PIPELINE_DISABLE_ALLOWED_FROM,
PIPELINE_ENABLE_ALLOWED_FROM,
PIPELINE_ERROR_MESSAGES,
getStatusName,
} from './Pipeline.utils'
import { PipelineStatus } from './PipelineStatus'
import { STATUS_REFRESH_FREQUENCY_MS } from './Replication.constants'
import { TableState } from './ReplicationPipelineStatus.types'
@@ -97,8 +99,9 @@ export const ReplicationPipelineStatus = () => {
const isPipelineRunning = statusName === 'started'
const hasTableData = tableStatuses.length > 0
const isEnablingDisabling =
requestStatus === PipelineStatusRequestStatus.EnableRequested ||
requestStatus === PipelineStatusRequestStatus.DisableRequested
requestStatus === PipelineStatusRequestStatus.StartRequested ||
requestStatus === PipelineStatusRequestStatus.StopRequested ||
requestStatus === PipelineStatusRequestStatus.RestartRequested
const showDisabledState = !isPipelineRunning || isEnablingDisabling
const onTogglePipeline = async () => {
@@ -110,12 +113,12 @@ export const ReplicationPipelineStatus = () => {
}
try {
if (statusName === 'stopped') {
if (PIPELINE_ENABLE_ALLOWED_FROM.includes(statusName as any)) {
setRequestStatus(pipeline.id, PipelineStatusRequestStatus.StartRequested, statusName)
await startPipeline({ projectRef, pipelineId: pipeline.id })
setRequestStatus(pipeline.id, PipelineStatusRequestStatus.EnableRequested)
} else if (statusName === 'started') {
} else if (PIPELINE_DISABLE_ALLOWED_FROM.includes(statusName as any)) {
setRequestStatus(pipeline.id, PipelineStatusRequestStatus.StopRequested, statusName)
await stopPipeline({ projectRef, pipelineId: pipeline.id })
setRequestStatus(pipeline.id, PipelineStatusRequestStatus.DisableRequested)
}
} catch (error) {
toast.error(PIPELINE_ERROR_MESSAGES.ENABLE_DESTINATION)
@@ -181,9 +184,9 @@ export const ReplicationPipelineStatus = () => {
type={statusName === 'stopped' ? 'primary' : 'default'}
onClick={() => onTogglePipeline()}
loading={isPipelineError || isStartingPipeline || isStoppingPipeline}
disabled={!['failed', 'started', 'stopped', 'stopping'].includes(statusName ?? '')}
disabled={!PIPELINE_ACTIONABLE_STATES.includes((statusName ?? '') as any)}
>
{statusName === 'stopped' ? 'Enable' : 'Disable'} pipeline
{statusName === 'stopped' ? 'Start' : 'Stop'} pipeline
</Button>
</div>
</div>
@@ -286,7 +289,7 @@ export const ReplicationPipelineStatus = () => {
Status unavailable while pipeline is {config.badge.toLowerCase()}
</p>
) : (
<div className="space-y-3">
<div className="space-y-1">
<div className="text-sm text-foreground">
{statusConfig.description}
</div>
@@ -56,9 +56,10 @@ export const getDisabledStateConfig = ({
const { title, message, badge } = getPipelineStateMessages(requestStatus, statusName)
// Get icon and colors based on current state
const isEnabling = requestStatus === PipelineStatusRequestStatus.EnableRequested
const isDisabling = requestStatus === PipelineStatusRequestStatus.DisableRequested
const isTransitioning = isEnabling || isDisabling
const isEnabling = requestStatus === PipelineStatusRequestStatus.StartRequested
const isDisabling = requestStatus === PipelineStatusRequestStatus.StopRequested
const isRestarting = requestStatus === PipelineStatusRequestStatus.RestartRequested
const isTransitioning = isEnabling || isDisabling || isRestarting
const icon = isTransitioning ? (
<Loader2 className="w-6 h-6 animate-spin" />
@@ -72,37 +73,38 @@ export const getDisabledStateConfig = ({
<Activity className="w-6 h-6" />
)
const colors = isEnabling
? {
bg: 'bg-brand-50',
text: 'text-brand-900',
subtext: 'text-brand-700',
iconBg: 'bg-brand-600',
icon: 'text-white dark:text-black',
}
: isDisabling || statusName === 'starting' || statusName === 'unknown'
const colors =
isEnabling || isRestarting
? {
bg: 'bg-warning-50',
text: 'text-warning-900',
subtext: 'text-warning-700',
iconBg: 'bg-warning-600',
bg: 'bg-brand-50',
text: 'text-brand-900',
subtext: 'text-brand-700',
iconBg: 'bg-brand-600',
icon: 'text-white dark:text-black',
}
: statusName === 'failed'
: isDisabling || statusName === 'starting' || statusName === 'unknown'
? {
bg: 'bg-destructive-50',
text: 'text-destructive-900',
subtext: 'text-destructive-700',
iconBg: 'bg-destructive-600',
icon: 'text-white dark:text-black',
}
: {
bg: 'bg-surface-100',
text: 'text-foreground',
subtext: 'text-foreground-light',
iconBg: 'bg-foreground-lighter',
bg: 'bg-warning-50',
text: 'text-warning-900',
subtext: 'text-warning-700',
iconBg: 'bg-warning-600',
icon: 'text-white dark:text-black',
}
: statusName === 'failed'
? {
bg: 'bg-destructive-50',
text: 'text-destructive-900',
subtext: 'text-destructive-700',
iconBg: 'bg-destructive-600',
icon: 'text-white dark:text-black',
}
: {
bg: 'bg-surface-100',
text: 'text-foreground',
subtext: 'text-foreground-light',
iconBg: 'bg-foreground-lighter',
icon: 'text-white dark:text-black',
}
return { title, message, badge, icon, colors }
}
@@ -40,7 +40,7 @@ export const RetryOptionsDropdown = ({ tableId, tableName }: RetryOptionsDropdow
const { mutate: rollbackTable, isLoading: isRollingBack } = useRollbackTableMutation({
onSuccess: (_, vars) => {
const { projectRef, pipelineId } = vars
toast.success(`Table "${tableName}" rolled back successfully`)
toast.success(`Table "${tableName}" rolled back successfully and pipeline is being restarted`)
startPipeline({ projectRef, pipelineId })
},
onError: (error, vars) => {
@@ -20,7 +20,12 @@ import {
DropdownMenuTrigger,
} from 'ui'
import ShimmeringLoader from 'ui-patterns/ShimmeringLoader'
import { PIPELINE_ERROR_MESSAGES } from './Pipeline.utils'
import {
PIPELINE_DISABLE_ALLOWED_FROM,
PIPELINE_ENABLE_ALLOWED_FROM,
PIPELINE_ERROR_MESSAGES,
getStatusName,
} from './Pipeline.utils'
import { PipelineStatusName } from './PipelineStatus'
interface RowMenuProps {
@@ -44,13 +49,6 @@ export const RowMenu = ({
}: RowMenuProps) => {
const { ref: projectRef } = useParams()
const getStatusName = (status: any) => {
if (status && typeof status === 'object' && 'name' in status) {
return status.name
}
return status
}
const statusName = getStatusName(pipelineStatus)
const pipelineEnabled = statusName !== PipelineStatusName.STOPPED
@@ -67,10 +65,13 @@ export const RowMenu = ({
}
try {
// Only show 'enabling' when transitioning from allowed states
if (PIPELINE_ENABLE_ALLOWED_FROM.includes(statusName as any)) {
setGlobalRequestStatus(pipeline.id, PipelineStatusRequestStatus.StartRequested, statusName)
}
await startPipeline({ projectRef, pipelineId: pipeline.id })
toast(`Enabling pipeline ${pipeline.destination_name}`)
setGlobalRequestStatus(pipeline.id, PipelineStatusRequestStatus.EnableRequested)
} catch (error) {
setGlobalRequestStatus(pipeline.id, PipelineStatusRequestStatus.None)
toast.error(PIPELINE_ERROR_MESSAGES.ENABLE_DESTINATION)
}
}
@@ -86,10 +87,13 @@ export const RowMenu = ({
}
try {
// Only show 'disabling' when transitioning from allowed states
if (PIPELINE_DISABLE_ALLOWED_FROM.includes(statusName as any)) {
setGlobalRequestStatus(pipeline.id, PipelineStatusRequestStatus.StopRequested, statusName)
}
await stopPipeline({ projectRef, pipelineId: pipeline.id })
toast(`Disabling pipeline ${pipeline.destination_name}`)
setGlobalRequestStatus(pipeline.id, PipelineStatusRequestStatus.DisableRequested)
} catch (error) {
setGlobalRequestStatus(pipeline.id, PipelineStatusRequestStatus.None)
toast.error(PIPELINE_ERROR_MESSAGES.DISABLE_DESTINATION)
}
}
@@ -109,12 +113,12 @@ export const RowMenu = ({
{pipelineEnabled ? (
<DropdownMenuItem className="space-x-2" onClick={onDisablePipeline}>
<Pause size={14} />
<p>Disable pipeline</p>
<p>Stop pipeline</p>
</DropdownMenuItem>
) : (
<DropdownMenuItem className="space-x-2" onClick={onEnablePipeline}>
<Play size={14} />
<p>Enable pipeline</p>
<p>Start pipeline</p>
</DropdownMenuItem>
)}
<DropdownMenuSeparator />
@@ -6,7 +6,7 @@ import { analyticsKeys } from './keys'
export type FunctionsReqStatsVariables = {
projectRef?: string
functionId?: string
interval?: operations['FunctionRequestLogsController_getStatus']['parameters']['query']['interval']
interval?: operations['FunctionsLogsController_getRequestStats']['parameters']['query']['interval']
}
export type FunctionsReqStatsResponse = any
@@ -6,7 +6,7 @@ import { analyticsKeys } from './keys'
export type FunctionsResourceUsageVariables = {
projectRef?: string
functionId?: string
interval?: operations['FunctionResourceLogsController_getStatus']['parameters']['query']['interval']
interval?: operations['FunctionsLogsController_getRequestStats']['parameters']['query']['interval']
}
export type FunctionsResourceUsageResponse = any
@@ -15,7 +15,7 @@ export async function deleteContents(
headers: { Version: '2' },
params: {
path: { ref: projectRef },
query: { ids },
query: { ids: ids.join(',') },
},
signal,
})
@@ -12,7 +12,7 @@ export type LintRuleDeleteVariables = {
export async function deleteLintRule({ projectRef, ids }: LintRuleDeleteVariables) {
const { data, error } = await del('/platform/projects/{ref}/notifications/advisor/exceptions', {
params: { path: { ref: projectRef }, query: { ids } },
params: { path: { ref: projectRef }, query: { ids: ids.join(',') } },
})
if (error) handleError(error)
@@ -1,15 +1,15 @@
import { useMutation, UseMutationOptions, useQueryClient } from '@tanstack/react-query'
import { toast } from 'sonner'
import { handleError, post } from 'data/fetchers'
import type { ResponseError } from 'types'
import { replicationKeys } from './keys'
import { handleError, post } from 'data/fetchers'
export type BigQueryDestinationConfig = {
projectId: string
datasetId: string
serviceAccountKey: string
maxStalenessMins: number
maxStalenessMins?: number
}
export type CreateDestinationPipelineParams = {
@@ -21,7 +21,7 @@ export type CreateDestinationPipelineParams = {
sourceId: number
pipelineConfig: {
publicationName: string
batch: {
batch?: {
maxSize: number
maxFillMs: number
}
@@ -35,10 +35,7 @@ async function createDestinationPipeline(
destinationConfig: {
bigQuery: { projectId, datasetId, serviceAccountKey, maxStalenessMins },
},
pipelineConfig: {
publicationName,
batch: { maxSize, maxFillMs },
},
pipelineConfig: { publicationName, batch },
sourceId,
}: CreateDestinationPipelineParams,
signal?: AbortSignal
@@ -48,23 +45,27 @@ async function createDestinationPipeline(
const { data, error } = await post('/platform/replication/{ref}/destinations-pipelines', {
params: { path: { ref: projectRef } },
body: {
source_id: sourceId,
destination_name: destinationName,
destination_config: {
big_query: {
project_id: projectId,
dataset_id: datasetId,
service_account_key: serviceAccountKey,
max_staleness_mins: maxStalenessMins,
...(maxStalenessMins != null && { max_staleness_mins: maxStalenessMins }),
},
},
pipeline_config: {
publication_name: publicationName,
batch: {
max_size: maxSize,
max_fill_ms: maxFillMs,
},
...(batch
? {
batch: {
max_size: batch.maxSize,
max_fill_ms: batch.maxFillMs,
},
}
: {}),
},
source_id: sourceId,
},
signal,
})
@@ -92,8 +93,12 @@ export const useCreateDestinationPipelineMutation = ({
{
async onSuccess(data, variables, context) {
const { projectRef } = variables
await queryClient.invalidateQueries(replicationKeys.destinations(projectRef))
await queryClient.invalidateQueries(replicationKeys.pipelines(projectRef))
await Promise.all([
queryClient.invalidateQueries(replicationKeys.destinations(projectRef)),
queryClient.invalidateQueries(replicationKeys.pipelines(projectRef)),
])
await onSuccess?.(data, variables, context)
},
async onError(data, variables, context) {
@@ -1,9 +1,9 @@
import { useMutation, UseMutationOptions, useQueryClient } from '@tanstack/react-query'
import { toast } from 'sonner'
import { del, handleError } from 'data/fetchers'
import type { ResponseError } from 'types'
import { replicationKeys } from './keys'
import { handleError, del } from 'data/fetchers'
export type DeleteDestinationPipelineParams = {
projectRef: string
@@ -47,8 +47,19 @@ export const useDeleteDestinationPipelineMutation = ({
(vars) => deleteDestinationPipeline(vars),
{
async onSuccess(data, variables, context) {
const { projectRef } = variables
await queryClient.invalidateQueries(replicationKeys.destinations(projectRef))
const { projectRef, destinationId, pipelineId } = variables
await Promise.all([
queryClient.invalidateQueries(replicationKeys.destinations(projectRef)),
queryClient.invalidateQueries(replicationKeys.pipelines(projectRef)),
queryClient.invalidateQueries(replicationKeys.pipelineById(projectRef, pipelineId)),
queryClient.invalidateQueries(replicationKeys.pipelinesStatus(projectRef, pipelineId)),
queryClient.invalidateQueries(
replicationKeys.pipelinesReplicationStatus(projectRef, pipelineId)
),
queryClient.invalidateQueries(replicationKeys.destinationById(projectRef, destinationId)),
])
await onSuccess?.(data, variables, context)
},
async onError(data, variables, context) {
@@ -1,15 +1,15 @@
import { useMutation, UseMutationOptions, useQueryClient } from '@tanstack/react-query'
import { toast } from 'sonner'
import { handleError, post } from 'data/fetchers'
import type { ResponseError } from 'types'
import { replicationKeys } from './keys'
import { handleError, post } from 'data/fetchers'
export type BigQueryDestinationConfig = {
projectId: string
datasetId: string
serviceAccountKey: string
maxStalenessMins: number
maxStalenessMins?: number
}
export type UpdateDestinationPipelineParams = {
@@ -23,7 +23,7 @@ export type UpdateDestinationPipelineParams = {
sourceId: number
pipelineConfig: {
publicationName: string
batch: {
batch?: {
maxSize: number
maxFillMs: number
}
@@ -39,10 +39,7 @@ async function updateDestinationPipeline(
destinationConfig: {
bigQuery: { projectId, datasetId, serviceAccountKey, maxStalenessMins },
},
pipelineConfig: {
publicationName,
batch: { maxSize, maxFillMs },
},
pipelineConfig: { publicationName, batch },
sourceId,
}: UpdateDestinationPipelineParams,
signal?: AbortSignal
@@ -60,15 +57,17 @@ async function updateDestinationPipeline(
project_id: projectId,
dataset_id: datasetId,
service_account_key: serviceAccountKey,
max_staleness_mins: maxStalenessMins,
...(maxStalenessMins != null && { max_staleness_mins: maxStalenessMins }),
},
},
pipeline_config: {
publication_name: publicationName,
batch: {
max_size: maxSize,
max_fill_ms: maxFillMs,
},
...(batch && {
batch: {
max_size: batch.maxSize,
max_fill_ms: batch.maxFillMs,
},
}),
},
source_id: sourceId,
},
@@ -99,8 +98,12 @@ export const useUpdateDestinationPipelineMutation = ({
{
async onSuccess(data, variables, context) {
const { projectRef } = variables
await queryClient.invalidateQueries(replicationKeys.destinations(projectRef))
await queryClient.invalidateQueries(replicationKeys.pipelines(projectRef))
await Promise.all([
queryClient.invalidateQueries(replicationKeys.destinations(projectRef)),
queryClient.invalidateQueries(replicationKeys.pipelines(projectRef)),
])
await onSuccess?.(data, variables, context)
},
async onError(data, variables, context) {
@@ -1,14 +1,28 @@
import { createContext, useContext, useState, ReactNode, useCallback } from 'react'
import {
createContext,
useContext,
useState,
ReactNode,
useCallback,
useRef,
useEffect,
} from 'react'
export enum PipelineStatusRequestStatus {
None = 'None',
EnableRequested = 'EnableRequested',
DisableRequested = 'DisableRequested',
StartRequested = 'StartRequested',
StopRequested = 'StopRequested',
RestartRequested = 'RestartRequested',
}
interface PipelineRequestStatusContextType {
requestStatus: Record<number, PipelineStatusRequestStatus>
setRequestStatus: (pipelineId: number, status: PipelineStatusRequestStatus) => void
pipelineStatusSnapshot: Record<number, string | undefined>
setRequestStatus: (
pipelineId: number,
status: PipelineStatusRequestStatus,
snapshotStatus?: string
) => void
getRequestStatus: (pipelineId: number) => PipelineStatusRequestStatus
updatePipelineStatus: (pipelineId: number, backendStatus: string | undefined) => void
}
@@ -25,12 +39,58 @@ export const PipelineRequestStatusProvider = ({ children }: PipelineRequestStatu
const [requestStatus, setRequestStatusState] = useState<
Record<number, PipelineStatusRequestStatus>
>({})
const [pipelineStatusSnapshot, setPipelineStatusSnapshot] = useState<
Record<number, string | undefined>
>({})
const timeoutsRef = useRef<Record<number, number>>({})
const REQUEST_TIMEOUT_MS = 10_000
const setRequestStatus = (pipelineId: number, status: PipelineStatusRequestStatus) => {
const setRequestStatus = (
pipelineId: number,
status: PipelineStatusRequestStatus,
snapshotStatus?: string
) => {
setRequestStatusState((prev) => ({
...prev,
[pipelineId]: status,
}))
setPipelineStatusSnapshot((prev) => {
if (status === PipelineStatusRequestStatus.None) {
const { [pipelineId]: _omit, ...rest } = prev
return rest
}
// Only set snapshot when provided to avoid undefined entries
if (snapshotStatus !== undefined) {
return { ...prev, [pipelineId]: snapshotStatus }
}
return prev
})
// Clear existing timeout for this pipeline
const existing = timeoutsRef.current[pipelineId]
if (existing !== undefined) {
clearTimeout(existing)
delete timeoutsRef.current[pipelineId]
}
// Start auto-reset timer for non-None states
if (status !== PipelineStatusRequestStatus.None) {
const id = window.setTimeout(() => {
// If still pending, clear to None to show backend state
setRequestStatusState((prev) => {
if (prev[pipelineId] && prev[pipelineId] !== PipelineStatusRequestStatus.None) {
return { ...prev, [pipelineId]: PipelineStatusRequestStatus.None }
}
return prev
})
setPipelineStatusSnapshot((prev) => {
const { [pipelineId]: _omit, ...rest } = prev
return rest
})
delete timeoutsRef.current[pipelineId]
}, REQUEST_TIMEOUT_MS)
timeoutsRef.current[pipelineId] = id
}
}
const getRequestStatus = (pipelineId: number): PipelineStatusRequestStatus => {
@@ -38,25 +98,32 @@ export const PipelineRequestStatusProvider = ({ children }: PipelineRequestStatu
}
const updatePipelineStatus = useCallback(
(pipelineId: number, backendStatus: string | undefined) => {
(pipelineId: number, newStatus: string | undefined) => {
const currentRequestStatus = requestStatus[pipelineId] || PipelineStatusRequestStatus.None
if (currentRequestStatus === PipelineStatusRequestStatus.None) return
if (
(currentRequestStatus === PipelineStatusRequestStatus.EnableRequested &&
(backendStatus === 'started' || backendStatus === 'failed')) ||
(currentRequestStatus === PipelineStatusRequestStatus.DisableRequested &&
(backendStatus === 'stopped' || backendStatus === 'failed'))
) {
// Only remove when backend status differs from snapshot
const snapshotStatus = pipelineStatusSnapshot[pipelineId]
if (newStatus !== snapshotStatus) {
setRequestStatus(pipelineId, PipelineStatusRequestStatus.None)
}
},
[requestStatus, setRequestStatus]
[requestStatus, pipelineStatusSnapshot]
)
// Cleanup all timers on unmount
useEffect(() => {
return () => {
Object.values(timeoutsRef.current).forEach((id) => clearTimeout(id))
timeoutsRef.current = {}
}
}, [])
return (
<PipelineRequestStatusContext.Provider
value={{
requestStatus,
pipelineStatusSnapshot,
setRequestStatus,
getRequestStatus,
updatePipelineStatus,
+175 -5
View File
@@ -337,6 +337,23 @@ export interface paths {
patch?: never
trace?: never
}
'/v1/projects/{ref}/analytics/endpoints/functions.combined-stats': {
parameters: {
query?: never
header?: never
path?: never
cookie?: never
}
/** Gets a project's function combined statistics */
get: operations['v1-get-project-function-combined-stats']
put?: never
post?: never
delete?: never
options?: never
head?: never
patch?: never
trace?: never
}
'/v1/projects/{ref}/analytics/endpoints/logs.all': {
parameters: {
query?: never
@@ -946,7 +963,11 @@ export interface paths {
* @description Modifies the roles that can be assumed and for how long
*/
put: operations['v1-update-jit-access']
post?: never
/**
* Authorize user-id to role mappings for JIT access
* @description Authorizes the request to assume a role in the project database
*/
post: operations['v1-authorize-jit-access']
delete?: never
options?: never
head?: never
@@ -1828,6 +1849,8 @@ export interface components {
mfa_totp_verify_enabled: boolean | null
mfa_web_authn_enroll_enabled: boolean | null
mfa_web_authn_verify_enabled: boolean | null
nimbus_oauth_client_id: string | null
nimbus_oauth_client_secret: string | null
password_hibp_enabled: boolean | null
password_min_length: number | null
/** @enum {string|null} */
@@ -1894,6 +1917,10 @@ export interface components {
smtp_user: string | null
uri_allow_list: string | null
}
AuthorizeJitAccessBody: {
rhost: string
role: string
}
BranchActionBody: {
migration_version?: string
}
@@ -1959,6 +1986,7 @@ export interface components {
| 'FUNCTIONS_FAILED'
/** Format: date-time */
updated_at: string
with_data: boolean
}
BranchUpdateResponse: {
/** @enum {string} */
@@ -2336,16 +2364,48 @@ export interface components {
/** Format: uuid */
user_id: string
user_roles: {
expires_at?: string
allowed_networks?: {
allowed_cidrs?: {
cidr: string
}[]
allowed_cidrs_v6?: {
cidr: string
}[]
}
expires_at?: number
role: string
}[]
}
JitAuthorizeAccessResponse: {
/** Format: uuid */
user_id: string
user_role: {
allowed_networks?: {
allowed_cidrs?: {
cidr: string
}[]
allowed_cidrs_v6?: {
cidr: string
}[]
}
expires_at?: number
role: string
}
}
JitListAccessResponse: {
items: {
/** Format: uuid */
user_id: string
user_roles: {
expires_at?: string
allowed_networks?: {
allowed_cidrs?: {
cidr: string
}[]
allowed_cidrs_v6?: {
cidr: string
}[]
}
expires_at?: number
role: string
}[]
}[]
@@ -2981,6 +3041,8 @@ export interface components {
mfa_totp_verify_enabled?: boolean | null
mfa_web_authn_enroll_enabled?: boolean | null
mfa_web_authn_verify_enabled?: boolean | null
nimbus_oauth_client_id?: string | null
nimbus_oauth_client_secret?: string | null
password_hibp_enabled?: boolean | null
password_min_length?: number | null
/** @enum {string|null} */
@@ -3107,7 +3169,15 @@ export interface components {
}
UpdateJitAccessBody: {
roles: {
expires_at?: string
allowed_networks?: {
allowed_cidrs?: {
cidr: string
}[]
allowed_cidrs_v6?: {
cidr: string
}[]
}
expires_at?: number
role: string
}[]
/** Format: uuid */
@@ -4283,6 +4353,44 @@ export interface operations {
}
}
}
'v1-get-project-function-combined-stats': {
parameters: {
query: {
function_id: string
interval: '15min' | '1hr' | '3hr' | '1day'
}
header?: never
path: {
/** @description Project ref */
ref: string
}
cookie?: never
}
requestBody?: never
responses: {
200: {
headers: {
[name: string]: unknown
}
content: {
'application/json': components['schemas']['AnalyticsResponse']
}
}
403: {
headers: {
[name: string]: unknown
}
content?: never
}
/** @description Failed to get project's function combined statistics */
500: {
headers: {
[name: string]: unknown
}
content?: never
}
}
}
'v1-get-project-logs': {
parameters: {
query?: {
@@ -4692,7 +4800,30 @@ export interface operations {
query?: never
header?: never
path: {
addon_variant: unknown
addon_variant:
| (
| 'ci_micro'
| 'ci_small'
| 'ci_medium'
| 'ci_large'
| 'ci_xlarge'
| 'ci_2xlarge'
| 'ci_4xlarge'
| 'ci_8xlarge'
| 'ci_12xlarge'
| 'ci_16xlarge'
| 'ci_24xlarge'
| 'ci_24xlarge_optimized_cpu'
| 'ci_24xlarge_optimized_memory'
| 'ci_24xlarge_high_memory'
| 'ci_48xlarge'
| 'ci_48xlarge_optimized_cpu'
| 'ci_48xlarge_optimized_memory'
| 'ci_48xlarge_high_memory'
)
| 'cd_default'
| ('pitr_7' | 'pitr_14' | 'pitr_28')
| 'ipv4_default'
/** @description Project ref */
ref: string
}
@@ -6229,6 +6360,45 @@ export interface operations {
}
}
}
'v1-authorize-jit-access': {
parameters: {
query?: never
header?: never
path: {
/** @description Project ref */
ref: string
}
cookie?: never
}
requestBody: {
content: {
'application/json': components['schemas']['AuthorizeJitAccessBody']
}
}
responses: {
200: {
headers: {
[name: string]: unknown
}
content: {
'application/json': components['schemas']['JitAuthorizeAccessResponse']
}
}
403: {
headers: {
[name: string]: unknown
}
content?: never
}
/** @description Failed to authorize database jit access */
500: {
headers: {
[name: string]: unknown
}
content?: never
}
}
}
'v1-delete-jit-access': {
parameters: {
query?: never
+365 -295
View File
File diff suppressed because it is too large. Load diff