Improve granularity of replication status (#40329)

* Improve granularity of replication status

* Minor copy improvements

* Update toast message copy

* Decouple pipeline status from table replication status
This commit is contained in:
Joshen Lim authored and GitHub committed 2025-11-13 15:13:03 +08:00
1 parent 7100f39dde
commit d8678921c0
8 files changed
+215 -81

No files matched your search

@@ -20,7 +20,7 @@ export const EnableReplicationModal = () => {
const { ref: projectRef } = useParams()
const [open, setOpen] = useState(false)
const { mutateAsync: createTenantSource, isLoading: creatingTenantSource } =
const { mutate: createTenantSource, isLoading: creatingTenantSource } =
useCreateTenantSourceMutation({
onSuccess: () => {
toast.success('Replication has been successfully enabled!')
@@ -33,7 +33,7 @@ export const EnableReplicationModal = () => {
const onEnableReplication = async () => {
if (!projectRef) return console.error('Project ref is required')
await createTenantSource({ projectRef })
createTenantSource({ projectRef })
}
return (
@@ -0,0 +1,13 @@
import { snakeCase } from 'lodash'
import { ReplicationPublication } from 'data/etl/publications-query'
export const inferPostgresTableFromNamespaceTable = ({
publication,
tableName,
}: {
publication?: ReplicationPublication
tableName: string
}) => {
return publication?.tables.find((x) => tableName === snakeCase(`${x.schema}.${x.name}_changelog`))
}
@@ -1,7 +1,7 @@
import { snakeCase, uniq } from 'lodash'
import { MoreVertical, Pause, Play, Trash } from 'lucide-react'
import { Loader2, MoreVertical, Pause, Play, Trash } from 'lucide-react'
import Link from 'next/link'
import { useState } from 'react'
import { useMemo, useState } from 'react'
import { toast } from 'sonner'
import { useParams } from 'common'
@@ -10,7 +10,9 @@ import {
formatWrapperTables,
} from 'components/interfaces/Integrations/Wrappers/Wrappers.utils'
import { getDecryptedParameters } from 'components/interfaces/Storage/ImportForeignSchemaDialog.utils'
import { DotPing } from 'components/ui/DotPing'
import { DropdownMenuItemTooltip } from 'components/ui/DropdownMenuItemTooltip'
import { useReplicationPipelineStatusQuery } from 'data/etl/pipeline-status-query'
import { useUpdatePublicationMutation } from 'data/etl/publication-update-mutation'
import { useStartPipelineMutation } from 'data/etl/start-pipeline-mutation'
import { useReplicationTablesQuery } from 'data/etl/tables-query'
@@ -27,14 +29,17 @@ import {
DropdownMenuTrigger,
TableCell,
TableRow,
Tooltip,
TooltipContent,
TooltipTrigger,
} from 'ui'
import ConfirmationModal from 'ui-patterns/Dialogs/ConfirmationModal'
import { ConfirmationModal } from 'ui-patterns/Dialogs/ConfirmationModal'
import { getAnalyticsBucketFDWServerName } from '../AnalyticsBucketDetails.utils'
import { useAnalyticsBucketAssociatedEntities } from '../useAnalyticsBucketAssociatedEntities'
import { useAnalyticsBucketWrapperInstance } from '../useAnalyticsBucketWrapperInstance'
import { inferPostgresTableFromNamespaceTable } from './NamespaceWithTables.utils'
interface TableRowComponentProps {
index: number
table: { id: number; name: string; isConnected: boolean }
schema: string
namespace: string
@@ -43,7 +48,6 @@ interface TableRowComponentProps {
}
export const TableRowComponent = ({
index,
table,
schema,
namespace,
@@ -63,6 +67,8 @@ export const TableRowComponent = ({
projectRef,
bucketId,
})
const { data } = useReplicationPipelineStatusQuery({ projectRef, pipelineId: pipeline?.id })
const pipelineStatus = data?.status.name
const { data: tables } = useReplicationTablesQuery({ projectRef, sourceId })
const { data: wrapperInstance, meta: wrapperMeta } = useAnalyticsBucketWrapperInstance({
@@ -74,9 +80,29 @@ export const TableRowComponent = ({
const { mutateAsync: updatePublication } = useUpdatePublicationMutation()
const { mutateAsync: startPipeline } = useStartPipelineMutation()
const isReplicating = !!publication?.tables.find(
(x) => table.name === snakeCase(`${x.schema}.${x.name}_changelog`)
)
const inferredPostgresTable = inferPostgresTableFromNamespaceTable({
publication,
tableName: table.name,
})
const isTableUnderReplicationPublication = !!inferredPostgresTable
const hasReplication = !!pipeline && !!publication
const isPipelineRunning = pipelineStatus === 'started'
const isReplicating = isTableUnderReplicationPublication && isPipelineRunning
// [Joshen] Considers both the replication pipeline status + if the table is in the replication publication
const replicationStatusLabel = useMemo(() => {
if (isLoading) return 'Checking'
if (hasReplication) {
if (!isPipelineRunning) {
return '-'
} else if (isTableUnderReplicationPublication) {
return 'Running'
} else {
return 'Disabled'
}
}
}, [hasReplication, isLoading, isPipelineRunning, isTableUnderReplicationPublication])
const onConfirmStopReplication = async () => {
if (!projectRef) return console.error('Project ref is required')
@@ -100,9 +126,9 @@ export const TableRowComponent = ({
})
await startPipeline({ projectRef, pipelineId: pipeline.id })
setShowStopReplicationModal(false)
toast.success('Successfully stopped replication for table! Pipeline is being restarted.')
toast.success('Successfully disabled replication for table! Pipeline is being restarted.')
} catch (error: any) {
toast.error(`Failed to stop replication for table: ${error.message}`)
toast.error(`Failed to disable replication for table: ${error.message}`)
} finally {
setIsUpdatingReplication(false)
}
@@ -132,9 +158,9 @@ export const TableRowComponent = ({
})
await startPipeline({ projectRef, pipelineId: pipeline.id })
setShowStartReplicationModal(false)
toast.success('Successfully stopped replication for table! Pipeline is being restarted.')
toast.success('Successfully enabled replication for table! Pipeline is being restarted.')
} catch (error: any) {
toast.error(`Failed to stop replication for table: ${error.message}`)
toast.error(`Failed to enable replication for table: ${error.message}`)
} finally {
setIsUpdatingReplication(false)
}
@@ -209,38 +235,35 @@ export const TableRowComponent = ({
<>
<TableRow>
<TableCell className="min-w-[120px]">{table.name}</TableCell>
{!!publication && (
{!!hasReplication && (
<TableCell colSpan={table.isConnected ? 1 : 2} className="min-w-[150px]">
<div className="flex flex-row items-center text-foregroung-lighter">
<div className="relative mr-2 align-middle w-3 h-3">
<span
className={`absolute inset-0 rounded-full ${
isReplicating
? isLoading
? 'bg-brand/20 animate-ping'
: 'bg-brand/20 animate-ping'
: isLoading
? 'bg-selection/20 animate-ping'
: 'hidden'
}`}
style={{
animationDelay: `${1 + index * 0.15}s`,
animationDuration: '2s',
}}
/>
<span
className={`absolute top-1/2 left-1/2 transform -translate-x-1/2 -translate-y-1/2 inline-block w-2 h-2 rounded-full ${
isReplicating ? 'bg-brand' : 'bg-selection'
}`}
/>
</div>
<span className="text-foreground-lighter">
{isLoading && !isReplicating
? '-'
: isReplicating
? 'Replicating'
: 'Not replicating'}
</span>
<div className="flex items-center">
<Tooltip>
<TooltipTrigger asChild>
<div className="flex items-center gap-x-2">
{isLoading ? (
<Loader2 size={12} className="animate-spin text-foreground-lighter" />
) : isPipelineRunning ? (
<DotPing
animate={isReplicating}
variant={isReplicating ? 'primary' : 'default'}
/>
) : null}
<span className="text-foreground-lighter capitalize">
{replicationStatusLabel}
</span>
</div>
</TooltipTrigger>
{isPipelineRunning && (
<TooltipContent side="bottom">
{isReplicating
? `Table data is currently replicating${!!inferredPostgresTable ? ` from ${inferredPostgresTable.schema}.${inferredPostgresTable.name}` : ''}`
: !isTableUnderReplicationPublication
? 'Replication is disabled for this table'
: undefined}
</TooltipContent>
)}
</Tooltip>
</div>
</TableCell>
)}
@@ -270,13 +293,13 @@ export const TableRowComponent = ({
{!!publication && (
<>
{isReplicating ? (
{isTableUnderReplicationPublication ? (
<DropdownMenuItem
className="flex items-center gap-x-2"
onClick={() => setShowStopReplicationModal(true)}
>
<Pause size={12} className="text-foreground-lighter" />
<p>Stop replication</p>
<p>Disable replication</p>
</DropdownMenuItem>
) : (
<DropdownMenuItem
@@ -284,7 +307,7 @@ export const TableRowComponent = ({
onClick={() => setShowStartReplicationModal(true)}
>
<Play size={12} className="text-foreground-lighter" />
<p>Start replication</p>
<p>Enable replication</p>
</DropdownMenuItem>
)}
</>
@@ -317,14 +340,14 @@ export const TableRowComponent = ({
variant="warning"
visible={showStopReplicationModal}
loading={isUpdatingReplication}
title="Confirm to stop replication for table"
confirmLabel="Stop replication"
title="Confirm to disable replication for table"
confirmLabel="Disable replication"
onCancel={() => setShowStopReplicationModal(false)}
onConfirm={() => onConfirmStopReplication()}
>
<p className="text-sm text-foreground-light">
Data within the "{table.name}" table will stop replicating. However do note that,
restarting replication on the table will clear and re-sync all data in it. Are you sure?
re-enabling replication on this table will clear and re-sync all data in it. Are you sure?
</p>
</ConfirmationModal>
@@ -333,13 +356,13 @@ export const TableRowComponent = ({
variant="warning"
visible={showStartReplicationModal}
loading={isUpdatingReplication}
title="Confirm to start replication for table"
confirmLabel="Start replication"
title="Confirm to enable replication for table"
confirmLabel="Enable replication"
onCancel={() => setShowStartReplicationModal(false)}
onConfirm={() => onConfirmStartReplication()}
>
<p className="text-sm text-foreground-light">
Restarting replication on the "{table.name}" table will clear and re-sync all data in it.
Re-enabling replication on the "{table.name}" table will clear and re-sync all data in it.
Are you sure?
</p>
</ConfirmationModal>
@@ -234,7 +234,7 @@ export const NamespaceWithTables = ({
</TableCell>
</TableRow>
) : (
allTables.map((table, index) => (
allTables.map((table) => (
<TableRowComponent
key={table.name}
table={table}
@@ -242,7 +242,6 @@ export const NamespaceWithTables = ({
token={token}
schema={displaySchema}
isLoading={isImportingForeignSchema || isLoadingNamespaceTables}
index={index}
/>
))
)}
@@ -25,6 +25,8 @@ import {
DatabaseExtension,
useDatabaseExtensionsQuery,
} from 'data/database-extensions/database-extensions-query'
import { useReplicationPipelineStatusQuery } from 'data/etl/pipeline-status-query'
import { useStartPipelineMutation } from 'data/etl/start-pipeline-mutation'
import { AnalyticsBucket } from 'data/storage/analytics-buckets-query'
import { useIcebergNamespacesQuery } from 'data/storage/iceberg-namespaces-query'
import { useIcebergWrapperCreateMutation } from 'data/storage/iceberg-wrapper-create-mutation'
@@ -64,6 +66,8 @@ export const AnalyticBucketDetails = () => {
const [pollIntervalNamespaces, setPollIntervalNamespaces] = useState(0)
const [pollIntervalNamespaceTables, setPollIntervalNamespaceTables] = useState(0)
const { mutateAsync: startPipeline, isLoading: isStartingPipeline } = useStartPipelineMutation()
/** The wrapper instance is the wrapper that is installed for this Analytics bucket. */
const { data: wrapperInstance, isLoading } = useAnalyticsBucketWrapperInstance({
bucketId: bucket?.id,
@@ -72,6 +76,18 @@ export const AnalyticBucketDetails = () => {
projectRef,
bucketId: bucket?.id,
})
const { data } = useReplicationPipelineStatusQuery(
{ projectRef, pipelineId: pipeline?.id },
{
refetchInterval: (data) => {
if (data?.status.name !== 'started') return 2000
else return false
},
}
)
const pipelineStatus = data?.status.name
const isPipelineRunning = pipelineStatus === 'started'
const isPipelineStopped = ['failed', 'stopped'].includes(pipelineStatus ?? '')
const wrapperValues = convertKVStringArrayToJson(wrapperInstance?.server_options ?? [])
const integration = INTEGRATIONS.find((i) => i.id === 'iceberg_wrapper' && i.type === 'wrapper')
@@ -193,7 +209,7 @@ export const AnalyticBucketDetails = () => {
</ScaffoldSectionDescription>
</div>
<div className="flex items-center gap-x-2">
{!!pipeline && (
{!!pipeline && isPipelineRunning && (
<Button asChild type="default">
<Link
href={`/project/${projectRef}/database/etl/${pipeline.replicator_id}`}
@@ -255,25 +271,74 @@ export const AnalyticBucketDetails = () => {
)}
</>
) : (
<div className="flex flex-col gap-y-10">
{namespaces.map(({ namespace, schema, tables }) => (
<NamespaceWithTables
key={namespace}
bucketName={bucket?.id}
namespace={namespace}
sourceType="direct"
schema={schema}
tables={tables as any}
token={token!}
wrapperInstance={wrapperInstance}
wrapperValues={wrapperValues}
wrapperMeta={wrapperMeta}
tablesToPoll={tablesToPoll}
pollIntervalNamespaceTables={pollIntervalNamespaceTables}
setPollIntervalNamespaceTables={setPollIntervalNamespaceTables}
/>
))}
</div>
<>
{!!pipeline && !isPipelineRunning && (
<Admonition
type="note"
layout="horizontal"
className="[&>div]:pl-[2.5rem] [&>div]:-translate-y-[3px]"
childProps={{ title: { className: 'block capitalize-sentence' } }}
showIcon={isPipelineStopped}
title={
isPipelineStopped
? `Replication on the bucket has ${pipelineStatus}`
: `${pipelineStatus} replication on the bucket...`
}
description={
isPipelineStopped
? 'Data changes from Postgres tables is currently not streaming to their corresponding analytics bucket table'
: 'Data changes from Postgres tables will resume streaming once pipeline has started'
}
actions={
<div className="flex items-center gap-x-2">
<Button asChild type="default">
<Link
href={`/project/${projectRef}/database/etl/${pipeline.replicator_id}`}
>
View replication
</Link>
</Button>
{isPipelineStopped && (
<Button
type="default"
loading={isStartingPipeline}
onClick={async () => {
if (projectRef) {
await startPipeline({ projectRef, pipelineId: pipeline.id })
}
}}
>
Restart
</Button>
)}
</div>
}
>
{!isPipelineStopped && (
<Loader2 size={18} className="absolute top-1.5 left-[3px] animate-spin" />
)}
</Admonition>
)}
<div className="flex flex-col gap-y-10">
{namespaces.map(({ namespace, schema, tables }) => (
<NamespaceWithTables
key={namespace}
bucketName={bucket?.id}
namespace={namespace}
sourceType="direct"
schema={schema}
tables={tables as any}
token={token!}
wrapperInstance={wrapperInstance}
wrapperValues={wrapperValues}
wrapperMeta={wrapperMeta}
tablesToPoll={tablesToPoll}
pollIntervalNamespaceTables={pollIntervalNamespaceTables}
setPollIntervalNamespaceTables={setPollIntervalNamespaceTables}
/>
))}
</div>
</>
)}
</ScaffoldSection>
+34
View File
@@ -0,0 +1,34 @@
import { cn } from 'ui'
interface DotPingProps {
animate?: boolean
variant?: 'primary' | 'default' | 'warning'
}
export const DotPing = ({ animate = true, variant = 'primary' }: DotPingProps) => {
return (
<div className="relative align-middle w-2.5 h-2.5">
<span
className={cn(
'absolute inset-0 rounded-full',
animate && 'animate-ping',
variant === 'primary' && 'bg-brand/20',
variant === 'default' && 'bg-selection/20',
variant === 'warning' && 'bg-warning/20'
)}
style={{
animationDelay: '1s',
animationDuration: '1.5s',
}}
/>
<span
className={cn(
'absolute top-1/2 left-1/2 transform -translate-x-1/2 -translate-y-1/2 inline-block w-2 h-2 rounded-full',
variant === 'primary' && 'bg-brand',
variant === 'default' && 'bg-selection',
variant === 'warning' && 'bg-warning'
)}
/>
</div>
)
}
@@ -34,7 +34,7 @@ export interface ConfirmationModalProps {
}
}
const ConfirmationModal = forwardRef<
export const ConfirmationModal = forwardRef<
React.ElementRef<typeof DialogContent>,
React.ComponentPropsWithoutRef<typeof Dialog> & ConfirmationModalProps
>(
+4 -4
View File
@@ -1,5 +1,5 @@
import { cva } from 'class-variance-authority'
import { forwardRef, ReactNode } from 'react'
import { ComponentProps, forwardRef, ReactNode } from 'react'
import { Alert_Shadcn_, AlertDescription_Shadcn_, AlertTitle_Shadcn_, cn } from 'ui'
export interface AdmonitionProps {
@@ -14,11 +14,11 @@ export interface AdmonitionProps {
| 'warning'
label?: string
title?: string
description?: string | React.ReactNode
description?: string | ReactNode
showIcon?: boolean
childProps?: {
title?: React.ComponentProps<typeof AlertTitle_Shadcn_>
description?: React.ComponentProps<typeof AlertDescription_Shadcn_>
title?: ComponentProps<typeof AlertTitle_Shadcn_>
description?: ComponentProps<typeof AlertDescription_Shadcn_>
}
layout?: 'horizontal' | 'vertical'
actions?: ReactNode