import { useParams } from 'common' import { Activity, ArrowUpCircle, Ban, ChevronDown, ChevronLeft, Info, Pause, Play, RotateCcw, Search, WifiOff, X, } from 'lucide-react' import Link from 'next/link' import { parseAsString, useQueryState } from 'nuqs' import { useEffect, useMemo, useState } from 'react' import { toast } from 'sonner' import { Button, Card, CardContent, DropdownMenu, DropdownMenuContent, DropdownMenuTrigger, Table, TableBody, TableHead, TableHeader, TableRow, } from 'ui' import { GenericSkeletonLoader } from 'ui-patterns' import { Input } from 'ui-patterns/DataInputs/Input' import { BatchRestartDialog } from '../BatchRestartDialog' import { ErrorDetailsDialog } from '../ErrorDetailsDialog' import { getPipelineDisplayState, getStatusName, PIPELINE_ACTIONABLE_STATES, } from '../Pipeline.utils' import { PipelineStatus } from '../PipelineStatus' import { PipelineStatusName, STATUS_REFRESH_FREQUENCY_MS } from '../Replication.constants' import { RestartTableDialog } from '../RestartTableDialog' import { UpdateVersionModal } from '../UpdateVersionModal' import { SlotLagMetrics } from './ReplicationPipelineStatus.types' import { getDisabledStateConfig } from './ReplicationPipelineStatus.utils' import { SlotLagMetricsInline, SlotLagMetricsList } from './SlotLagMetrics' import { TableReplicationRow } from './TableReplicationRow' import { AlertError } from '@/components/ui/AlertError' import { DropdownMenuItemTooltip } from '@/components/ui/DropdownMenuItemTooltip' import { useReplicationPipelineByIdQuery } from '@/data/replication/pipeline-by-id-query' import { useReplicationPipelineReplicationStatusQuery } from '@/data/replication/pipeline-replication-status-query' import { useReplicationPipelineStatusQuery } from '@/data/replication/pipeline-status-query' import { useReplicationPipelineVersionQuery } from '@/data/replication/pipeline-version-query' import { useRestartPipelineHelper } from '@/data/replication/restart-pipeline-helper' import { useStartPipelineMutation } from '@/data/replication/start-pipeline-mutation' import { useStopPipelineMutation } from '@/data/replication/stop-pipeline-mutation' import { PipelineStatusRequestStatus, usePipelineRequestStatus, } from '@/state/replication-pipeline-request-status' import { type ResponseError } from '@/types' /** * Component for displaying replication pipeline status and table replication details. * Supports both legacy 'error' state and new 'errored' state with retry policies. */ export const ReplicationPipelineStatus = () => { const { ref: projectRef, pipelineId: _pipelineId } = useParams() const [searchString, setSearchString] = useQueryState('search', parseAsString.withDefault('')) const [showUpdateVersionModal, setShowUpdateVersionModal] = useState(false) const [showErrorDialog, setShowErrorDialog] = useState(false) const [selectedTableError, setSelectedTableError] = useState<{ tableName: string reason: string solution?: string } | null>(null) const [showRestartDialog, setShowRestartDialog] = useState(false) const [selectedTableForRestart, setSelectedTableForRestart] = useState<{ tableId: number tableName: string } | null>(null) const [showBatchRestartDialog, setShowBatchRestartDialog] = useState(false) const [batchRestartMode, setBatchRestartMode] = useState<'all' | 'errored' | null>(null) const [restartingTableIds, setRestartingTableIds] = useState>(new Set()) const pipelineId = Number(_pipelineId) const { getRequestStatus, updatePipelineStatus, setRequestStatus } = usePipelineRequestStatus() const requestStatus = getRequestStatus(pipelineId) const { data: pipeline, error: pipelineError, isPending: isPipelineLoading, isError: isPipelineError, } = useReplicationPipelineByIdQuery({ projectRef, pipelineId, }) const { data: pipelineStatusData, error: pipelineStatusError, isLoading: isPipelineStatusLoading, isError: isPipelineStatusError, isSuccess: isPipelineStatusSuccess, } = useReplicationPipelineStatusQuery( { projectRef, pipelineId }, { enabled: !!pipelineId, refetchInterval: STATUS_REFRESH_FREQUENCY_MS, } ) const { data: replicationStatusData, isPending: isStatusLoading, isError: isStatusError, } = useReplicationPipelineReplicationStatusQuery( { projectRef, pipelineId }, { enabled: !!pipelineId, refetchInterval: STATUS_REFRESH_FREQUENCY_MS, } ) const { data: versionData } = useReplicationPipelineVersionQuery({ projectRef, pipelineId: pipeline?.id, }) const hasUpdate = Boolean(versionData?.new_version) const { mutateAsync: startPipeline, isPending: isStartingPipeline } = useStartPipelineMutation() const { mutateAsync: stopPipeline, isPending: isStoppingPipeline } = useStopPipelineMutation() const { restartPipeline } = useRestartPipelineHelper() const destinationName = pipeline?.destination_name const statusName = getStatusName(pipelineStatusData?.status) const displayState = getPipelineDisplayState(requestStatus, statusName) const config = getDisabledStateConfig({ requestStatus, statusName }) // Sort tables by name for consistent ordering (memoized) const tableStatuses = useMemo( () => (replicationStatusData?.table_statuses || []).sort((a, b) => a.table_name.localeCompare(b.table_name) ), [replicationStatusData?.table_statuses] ) const applyLagMetrics = replicationStatusData?.apply_lag // Filter tables based on search (memoized) const filteredTableStatuses = useMemo( () => searchString.length === 0 ? tableStatuses : tableStatuses.filter((table) => table.table_name.toLowerCase().includes(searchString.toLowerCase()) ), [tableStatuses, searchString] ) const tablesWithLag = useMemo( () => tableStatuses.filter((table) => Boolean(table.table_sync_lag)), [tableStatuses] ) const erroredTables = useMemo( () => tableStatuses.filter( (table) => table.state.name === 'error' && 'retry_policy' in table.state && table.state.retry_policy?.policy === 'manual_retry' ), [tableStatuses] ) const hasErroredTables = erroredTables.length > 0 const isAnyRestartInProgress = restartingTableIds.size > 0 const hasTableData = tableStatuses.length > 0 const isPipelineActionable = statusName === PipelineStatusName.STARTED || statusName === PipelineStatusName.STOPPED || statusName === PipelineStatusName.FAILED const isEnablingDisabling = requestStatus === PipelineStatusRequestStatus.StartRequested || requestStatus === PipelineStatusRequestStatus.StopRequested || requestStatus === PipelineStatusRequestStatus.RestartRequested const isPipelineBusy = isEnablingDisabling || isAnyRestartInProgress const showDisabledState = isPipelineBusy || !isPipelineActionable const lastKnownStateMessage = statusName === PipelineStatusName.STOPPED ? 'Showing the last known table state before the pipeline was stopped.' : statusName === PipelineStatusName.FAILED ? 'Showing the last reported table state before the pipeline failed.' : null const refreshIntervalLabel = STATUS_REFRESH_FREQUENCY_MS >= 1000 ? `${Math.round(STATUS_REFRESH_FREQUENCY_MS / 1000)}s` : `${STATUS_REFRESH_FREQUENCY_MS}ms` const logsUrl = `/project/${projectRef}/logs/replication-logs${ pipelineId ? `?f=${encodeURIComponent(JSON.stringify({ pipeline_id: pipelineId }))}` : '' }` const label = isEnablingDisabling ? displayState.label : statusName === PipelineStatusName.STOPPED ? 'Start' : statusName === PipelineStatusName.STARTED ? 'Stop' : statusName === PipelineStatusName.FAILED ? 'Restart' : displayState.label const icon = statusName === PipelineStatusName.STOPPED ? ( ) : statusName === PipelineStatusName.STARTED ? ( ) : statusName === PipelineStatusName.FAILED ? ( ) : ( ) const onPrimaryAction = async () => { if (!projectRef) return console.error('Project ref is required') if (!pipeline) return toast.error('No pipeline found') const action = statusName === PipelineStatusName.STOPPED ? 'start' : statusName === PipelineStatusName.STARTED ? 'stop' : 'restart' try { if (statusName === PipelineStatusName.STOPPED) { setRequestStatus(pipeline.id, PipelineStatusRequestStatus.StartRequested, statusName) await startPipeline({ projectRef, pipelineId: pipeline.id }) } else if (statusName === PipelineStatusName.STARTED) { setRequestStatus(pipeline.id, PipelineStatusRequestStatus.StopRequested, statusName) await stopPipeline({ projectRef, pipelineId: pipeline.id }) } else if (statusName === PipelineStatusName.FAILED) { setRequestStatus(pipeline.id, PipelineStatusRequestStatus.RestartRequested, statusName) await restartPipeline({ projectRef, pipelineId: pipeline.id }) } } catch (error) { setRequestStatus(pipeline.id, PipelineStatusRequestStatus.None) toast.error(`Failed to ${action} pipeline: ${(error as ResponseError).message}`) } } useEffect(() => { updatePipelineStatus(pipelineId, statusName) }, [pipelineId, statusName, updatePipelineStatus]) return ( <>

{destinationName || 'Pipeline'}

{hasUpdate && ( )}
{isPipelineError && ( )} {isStatusError && (
Live updates paused Retrying automatically
)} {(isPipelineLoading || isStatusLoading) && (
)} {applyLagMetrics && (

Replication lag

Snapshot of how far this pipeline is trailing behind right now.

Updates every {refreshIntervalLabel}

{isStatusError && (

Unable to refresh data. Showing the last values we received.

)} {tablesWithLag.length > 0 && ( <>
During initial sync, tables can copy and stream independently before reconciling with the overall pipeline.
    {tablesWithLag.map((table) => (
  • ))}
)}
)} {!isPipelineLoading && !isStatusLoading && hasTableData && (
} size="tiny" className="text-xs w-52" placeholder="Search for tables" value={searchString} disabled={isPipelineError} onChange={(e) => setSearchString(e.target.value)} actions={ searchString.length > 0 && [ setSearchString('')} />, ] } />
{lastKnownStateMessage !== null && !showDisabledState && (
{lastKnownStateMessage}
)} Table Status Details {filteredTableStatuses.map((table) => { const isRestarting = restartingTableIds.has(table.table_id) const isErrorState = table.state.name === 'error' const errorReason = isErrorState && 'reason' in table.state ? table.state.reason : undefined const errorSolution = isErrorState && 'solution' in table.state ? table.state.solution : undefined return ( { setSelectedTableForRestart({ tableId: table.table_id, tableName: table.table_name, }) setShowRestartDialog(true) }} onSelectShowError={ isErrorState && errorReason ? () => { setSelectedTableError({ tableName: table.table_name, reason: errorReason, solution: errorSolution, }) setShowErrorDialog(true) } : () => {} } /> ) })}
)} {!isPipelineLoading && !isStatusLoading && tableStatuses.length === 0 && (

{showDisabledState ? config.title : statusName === PipelineStatusName.STOPPED ? 'Pipeline stopped' : statusName === PipelineStatusName.FAILED ? 'Pipeline failed' : 'No table data yet'}

{showDisabledState ? config.message : statusName === PipelineStatusName.STOPPED ? 'Start the pipeline to begin replication.' : statusName === PipelineStatusName.FAILED ? 'The pipeline encountered an error. Restart it or reset your tables to recover.' : 'Table status will appear here once replication begins.'}

{statusName !== PipelineStatusName.STOPPED && (

Data refreshes every {refreshIntervalLabel}

)}
)}
setShowUpdateVersionModal(false)} confirmLabel={ statusName === PipelineStatusName.STARTED || statusName === PipelineStatusName.FAILED ? 'Update and restart' : 'Update version' } /> {/* Restart Table Confirmation Dialog */} {selectedTableForRestart && ( { setRestartingTableIds((prev) => new Set(prev).add(selectedTableForRestart.tableId)) }} onRestartComplete={() => { setRestartingTableIds((prev) => { const next = new Set(prev) next.delete(selectedTableForRestart.tableId) return next }) }} /> )} {/* Error Details Dialog */} {selectedTableError && ( )} {/* Batch Restart Dialog */} {batchRestartMode && ( { setRestartingTableIds((prev) => new Set([...prev, ...tableIds])) }} onRestartComplete={(tableIds) => { setRestartingTableIds((prev) => { const next = new Set(prev) tableIds.forEach((id) => next.delete(id)) return next }) }} /> )} ) }