import { useQueries, useQueryClient } from '@tanstack/react-query' import { useParams } from 'common' import { MoreVertical, Plus, Search, Workflow, X } from 'lucide-react' import Link from 'next/link' import { parseAsStringEnum, useQueryState } from 'nuqs' import { useEffect, useMemo, useRef, useState } from 'react' import { Button, Card, CardContent, DropdownMenu, DropdownMenuContent, DropdownMenuItem, DropdownMenuSeparator, DropdownMenuTrigger, Table, TableBody, TableHead, TableHeader, TableHeadSort, TableRow, } from 'ui' import { Input } from 'ui-patterns/DataInputs/Input' import { EmptyStatePresentational } from 'ui-patterns/EmptyStatePresentational' import { GenericTableLoader } from 'ui-patterns/ShimmeringLoader' import { DestinationPanel } from './DestinationPanel/DestinationPanel' import { DestinationType } from './DestinationPanel/DestinationPanel.types' import { DestinationRow } from './DestinationRow' import { DisablePipelinesDialog } from './DisablePipelinesDialog' import { EnablePipelinesModal } from './EnablePipelinesCallout' import { getStatusName } from './Pipeline.utils' import { PipelineStatusName } from './Replication.constants' import { useIsETLBigQueryPrivateAlpha, useIsETLClickHousePrivateAlpha, useIsETLDucklakePrivateAlpha, useIsETLIcebergPrivateAlpha, useIsETLSnowflakePrivateAlpha, } from './useIsETLPrivateAlpha' import { useRedirectLegacyReadReplicaDestination } from './useRedirectLegacyReadReplicaDestination' import { AlertError } from '@/components/ui/AlertError' import { Shortcut } from '@/components/ui/Shortcut' import { TableRowNoResults } from '@/components/ui/TableRowNoResults' import { useReplicationDestinationsQuery } from '@/data/replication/destinations-query' import { replicationKeys } from '@/data/replication/keys' import { replicationPipelineStatusQueryOptions, type ReplicationPipelineStatusData, } from '@/data/replication/pipeline-status-query' import { fetchReplicationPipelineVersion } from '@/data/replication/pipeline-version-query' import { useReplicationPipelinesQuery } from '@/data/replication/pipelines-query' import { useReplicationSourcesQuery } from '@/data/replication/sources-query' import { checkLocalETLNotSetUp } from '@/data/replication/utils' import { useSelectedOrganizationQuery } from '@/hooks/misc/useSelectedOrganization' import { onSearchInputEscape } from '@/lib/keyboard' import { SHORTCUT_IDS } from '@/state/shortcuts/registry' import { useShortcut } from '@/state/shortcuts/useShortcut' type DestinationSortColumn = 'name' | 'status' type DestinationSort = `${DestinationSortColumn}:${'asc' | 'desc'}` // Worst first, so sorting ascending by status surfaces the pipelines that need attention. const STATUS_SORT_ORDER: PipelineStatusName[] = [ PipelineStatusName.FAILED, PipelineStatusName.STOPPED, PipelineStatusName.STOPPING, PipelineStatusName.STARTING, PipelineStatusName.STARTED, PipelineStatusName.UNKNOWN, ] // Keyed by pipeline id from the responses themselves, so this never closes over component state. const combinePipelineStatuses = ( results: { data?: ReplicationPipelineStatusData }[] ): Map => new Map( results .map((result) => result.data) .filter((data): data is ReplicationPipelineStatusData => data !== undefined) .map((data) => [data.pipeline_id, getStatusName(data.status)]) ) const compareStatusNames = ( a: PipelineStatusName | undefined, b: PipelineStatusName | undefined, direction: 'asc' | 'desc' ) => { if (a === undefined) return b === undefined ? 0 : 1 if (b === undefined) return -1 const comparison = STATUS_SORT_ORDER.indexOf(a) - STATUS_SORT_ORDER.indexOf(b) return direction === 'asc' ? comparison : -comparison } export const Destinations = () => { const queryClient = useQueryClient() const { ref: projectRef } = useParams() const { data: organization } = useSelectedOrganizationQuery() useRedirectLegacyReadReplicaDestination() const etlEnableBigQuery = useIsETLBigQueryPrivateAlpha() const etlEnableIceberg = useIsETLIcebergPrivateAlpha() const etlEnableDucklake = useIsETLDucklakePrivateAlpha() const etlEnableSnowflake = useIsETLSnowflakePrivateAlpha() const etlEnableClickHouse = useIsETLClickHousePrivateAlpha() const newDestinationDefaultType: DestinationType | null = etlEnableBigQuery ? 'BigQuery' : etlEnableIceberg ? 'Analytics Bucket' : etlEnableDucklake ? 'DuckLake' : etlEnableSnowflake ? 'Snowflake' : etlEnableClickHouse ? 'ClickHouse' : null const prefetchedRef = useRef(false) const searchInputRef = useRef(null) const [filterString, setFilterString] = useState('') const [showEnablePipelinesDialog, setShowEnablePipelinesDialog] = useState(false) const [showDisablePipelinesDialog, setShowDisablePipelinesDialog] = useState(false) const [, setDestinationType] = useQueryState( 'destinationType', parseAsStringEnum([ 'BigQuery', 'Analytics Bucket', 'DuckLake', 'Snowflake', 'ClickHouse', ]).withOptions({ history: 'push', clearOnDefault: true, }) ) const { data: destinationsData, error: destinationsError, isPending: isDestinationsLoading, isError: isDestinationsError, isSuccess: isDestinationsSuccess, } = useReplicationDestinationsQuery({ projectRef, }) const destinations = useMemo( () => destinationsData?.destinations ?? [], [destinationsData?.destinations] ) const hasDestinations = isDestinationsSuccess && destinationsData?.destinations.length > 0 const filteredDestinations = useMemo( () => filterString.length === 0 ? destinations : destinations.filter((destination) => destination.name.toLowerCase().includes(filterString.toLowerCase()) ), [destinations, filterString] ) const { data: pipelinesData, isSuccess: isPipelinesSuccess } = useReplicationPipelinesQuery({ projectRef, }) const pipelines = useMemo(() => pipelinesData?.pipelines ?? [], [pipelinesData?.pipelines]) // Sorting by status needs every pipeline's status up here, not just inside each row. These share // the rows' query keys, so each status is still only fetched once. const statusByPipelineId = useQueries({ queries: pipelines.map((pipeline) => replicationPipelineStatusQueryOptions({ projectRef, pipelineId: pipeline.id }) ), combine: combinePipelineStatuses, }) const getDestinationStatus = (destinationId: number) => { const pipeline = pipelines.find((p) => p.destination_id === destinationId) return pipeline === undefined ? undefined : statusByPipelineId.get(pipeline.id) } const [sort, setSort] = useState('name:asc') const [sortColumn, sortDirection] = sort.split(':') as [DestinationSortColumn, 'asc' | 'desc'] const getAriaSort = (column: DestinationSortColumn) => { if (sortColumn !== column) return 'none' return sortDirection === 'asc' ? 'ascending' : 'descending' } const handleSortChange = (column: DestinationSortColumn) => { if (sortColumn !== column) return setSort(`${column}:asc`) setSort(`${column}:${sortDirection === 'asc' ? 'desc' : 'asc'}`) } // Not memoized: the status map is rebuilt whenever a pipeline status refetches, so a useMemo // here would never hit. Sorting a handful of destinations per render costs nothing. const sortedDestinations = [...filteredDestinations].sort((a, b) => { if (sortColumn === 'status') { const nameComparison = a.name.localeCompare(b.name) return ( compareStatusNames(getDestinationStatus(a.id), getDestinationStatus(b.id), sortDirection) || (sortDirection === 'asc' ? nameComparison : -nameComparison) ) } const comparison = a.name.localeCompare(b.name) return sortDirection === 'asc' ? comparison : -comparison }) const { data: sourcesData, isSuccess: isSourcesSuccess } = useReplicationSourcesQuery({ projectRef, }) const externalReplicationSource = useMemo( () => sourcesData?.sources.find((source) => source.name === projectRef), [projectRef, sourcesData?.sources] ) const replicationNotEnabled = isSourcesSuccess && !externalReplicationSource const canDisablePipelines = isSourcesSuccess && isDestinationsSuccess && isPipelinesSuccess && !!externalReplicationSource && destinations.length === 0 && pipelines.length === 0 const isLocalETLNotSetUp = checkLocalETLNotSetUp(destinationsError) const hasErrorsFetchingData = !isLocalETLNotSetUp && isDestinationsError const openDestinationPanel = () => { if (!newDestinationDefaultType) return setDestinationType(newDestinationDefaultType) } useShortcut( SHORTCUT_IDS.LIST_PAGE_FOCUS_SEARCH, () => { searchInputRef.current?.focus() searchInputRef.current?.select() }, { label: 'Search pipelines' } ) useShortcut(SHORTCUT_IDS.LIST_PAGE_RESET_FILTERS, () => setFilterString('')) useEffect(() => { if ( projectRef && !prefetchedRef.current && pipelinesData?.pipelines && pipelinesData.pipelines.length > 0 && isPipelinesSuccess ) { prefetchedRef.current = true pipelinesData.pipelines.forEach((p) => { if (!p?.id) return queryClient.prefetchQuery({ queryKey: replicationKeys.pipelinesVersion(projectRef, p.id), queryFn: ({ signal }) => fetchReplicationPipelineVersion({ projectRef, pipelineId: p.id }, signal), staleTime: Infinity, }) }) } }, [projectRef, pipelinesData?.pipelines, isPipelinesSuccess, queryClient]) return (
} value={filterString} className="w-full lg:w-52" onChange={(e) => setFilterString(e.target.value)} onKeyDown={onSearchInputEscape(filterString, setFilterString)} actions={ filterString.length > 0 && (
{/* Mounted whether or not it has anything to say, so the update is announced */}

{isDestinationsLoading ? 'Loading pipelines' : ''}

{hasErrorsFetchingData && ( )} {isDestinationsLoading && ( )} {!isDestinationsLoading && hasDestinations && ( Name Status Lag Publication {sortedDestinations.map((destination) => ( ))} {!isDestinationsLoading && filteredDestinations.length === 0 && hasDestinations && }
)} {!isDestinationsLoading && !hasDestinations && !hasErrorsFetchingData && ( )}
) }