diff --git a/apps/studio/components/interfaces/Database/Replication/DestinationPanel/DestinationPanel.tsx b/apps/studio/components/interfaces/Database/Replication/DestinationPanel/DestinationPanel.tsx index de5192ff12c..44a6cde8bd5 100644 --- a/apps/studio/components/interfaces/Database/Replication/DestinationPanel/DestinationPanel.tsx +++ b/apps/studio/components/interfaces/Database/Replication/DestinationPanel/DestinationPanel.tsx @@ -1,3 +1,4 @@ +import { useParams } from 'common' import { ArrowUpRight } from 'lucide-react' import Link from 'next/link' import { parseAsInteger, parseAsStringEnum, useQueryState } from 'nuqs' @@ -23,6 +24,8 @@ import { DestinationForm } from './DestinationForm' import { DestinationType } from './DestinationPanel.types' import { DestinationTypeSelection } from './DestinationTypeSelection' import { ReadReplicaForm } from './ReadReplicaForm' +import { WarehouseDestinationForm } from './WarehouseDestinationForm' +import { useWarehouseProjectState } from '@/components/interfaces/Database/Warehouse/warehouseDemoStore' import { DocsButton } from '@/components/ui/DocsButton' import { useCheckEntitlements } from '@/hooks/misc/useCheckEntitlements' import { DOCS_URL } from '@/lib/constants' @@ -32,6 +35,8 @@ interface DestinationPanelProps { } export const DestinationPanel = ({ onSuccessCreateReadReplica }: DestinationPanelProps) => { + const { ref: projectRef } = useParams() + const warehouseState = useWarehouseProjectState(projectRef) const enablePgReplicate = useIsETLPrivateAlpha() const { hasAccess: hasETLReplicationAccess } = useCheckEntitlements('replication.etl') @@ -39,6 +44,7 @@ export const DestinationPanel = ({ onSuccessCreateReadReplica }: DestinationPane 'destinationType', parseAsStringEnum([ 'Read Replica', + 'Warehouse', 'BigQuery', 'Analytics Bucket', 'DuckLake', @@ -94,74 +100,85 @@ export const DestinationPanel = ({ onSuccessCreateReadReplica }: DestinationPane } }, [edit, invalidExistingDestination, setEdit]) + useEffect(() => { + if (urlDestinationType === 'Warehouse' && warehouseState.enabled) { + setDestinationType(null) + } + }, [urlDestinationType, warehouseState.enabled, setDestinationType]) + return ( - <> - - -
- - {editMode ? 'Edit destination' : 'Add destination'} - - {editMode - ? 'Update the configuration for this destination.' - : 'A destination can be a read replica or an external destination that receives replicated data in near real time.'} - - + { + if (!open) onClose() + }} + > + +
+ + {editMode ? 'Edit destination' : 'Add destination'} + + {editMode + ? 'Update the configuration for this destination.' + : 'A destination can be a read replica or an external destination that receives replicated data in near real time.'} + + - + - + - {destinationType === 'Read Replica' ? ( - onSuccessCreateReadReplica?.()} /> - ) : !enablePgReplicate ? ( - -
-
-

Request Pipelines access

-

- Pipelines is in alpha and being - rolled out gradually. Request access below to join the waitlist. Read replicas - are available now. -

-
-
- - -
+ {destinationType === 'Read Replica' ? ( + onSuccessCreateReadReplica?.()} /> + ) : destinationType === 'Warehouse' ? ( + + ) : !enablePgReplicate ? ( + +
+
+

Request Pipelines access

+

+ Pipelines is in alpha and being rolled + out gradually. Request access below to join the waitlist. Read replicas are + available now. +

- - ) : replicationNotEnabled ? ( - - - - ) : ( - + + +
+
+
+ ) : replicationNotEnabled ? ( + + - )} -
-
-
- + + ) : ( + + )} +
+
+
) } diff --git a/apps/studio/components/interfaces/Database/Replication/DestinationPanel/DestinationPanel.types.ts b/apps/studio/components/interfaces/Database/Replication/DestinationPanel/DestinationPanel.types.ts index 0d14e2ac3d4..6dc7a64b056 100644 --- a/apps/studio/components/interfaces/Database/Replication/DestinationPanel/DestinationPanel.types.ts +++ b/apps/studio/components/interfaces/Database/Replication/DestinationPanel/DestinationPanel.types.ts @@ -1,5 +1,6 @@ export type DestinationType = | 'Read Replica' + | 'Warehouse' | 'BigQuery' | 'Analytics Bucket' | 'DuckLake' diff --git a/apps/studio/components/interfaces/Database/Replication/DestinationPanel/DestinationTypeSelection.test.tsx b/apps/studio/components/interfaces/Database/Replication/DestinationPanel/DestinationTypeSelection.test.tsx index f5d2d442785..290b483f76a 100644 --- a/apps/studio/components/interfaces/Database/Replication/DestinationPanel/DestinationTypeSelection.test.tsx +++ b/apps/studio/components/interfaces/Database/Replication/DestinationPanel/DestinationTypeSelection.test.tsx @@ -32,6 +32,14 @@ vi.mock('@/hooks/misc/useIsFeatureEnabled', () => ({ useIsFeatureEnabled: () => ({ infrastructureReadReplicas: true }), })) +vi.mock('@/hooks/misc/useIsWarehouseEnabled', () => ({ + useIsWarehouseEnabled: () => false, +})) + +vi.mock('@/components/interfaces/Database/Warehouse/warehouseDemoStore', () => ({ + useWarehouseProjectState: () => ({ enabled: false }), +})) + // Background queries from useDestinationInformation (sources + pipelines fire // even in create mode). Prevent retries so unmatched handlers fail fast. vi.mock('@/data/replication/utils', () => ({ @@ -125,7 +133,7 @@ describe('DestinationTypeSelection', () => { fireEvent.click(await screen.findByRole('combobox')) fireEvent.click(await screen.findByText('BigQuery')) - expect(await screen.findByText(/This destination type is in alpha/)).toBeInTheDocument() + expect(await screen.findByText(/In alpha and may change/)).toBeInTheDocument() }) test('disables the selector in edit mode so the destination type cannot be changed', async () => { diff --git a/apps/studio/components/interfaces/Database/Replication/DestinationPanel/DestinationTypeSelection.tsx b/apps/studio/components/interfaces/Database/Replication/DestinationPanel/DestinationTypeSelection.tsx index c528ff2e08e..fa3d959d142 100644 --- a/apps/studio/components/interfaces/Database/Replication/DestinationPanel/DestinationTypeSelection.tsx +++ b/apps/studio/components/interfaces/Database/Replication/DestinationPanel/DestinationTypeSelection.tsx @@ -1,6 +1,8 @@ +import { useParams } from 'common' import { AnalyticsBucket, BigQuery, Database } from 'icons' -import { Snowflake } from 'lucide-react' +import { Snowflake, Warehouse } from 'lucide-react' import { parseAsInteger, parseAsStringEnum, useQueryState } from 'nuqs' +import type { ElementType } from 'react' import { Badge, Select, @@ -21,14 +23,16 @@ import { useIsETLSnowflakePrivateAlpha, } from '../useIsETLPrivateAlpha' import { DestinationType } from './DestinationPanel.types' +import { useWarehouseProjectState } from '@/components/interfaces/Database/Warehouse/warehouseDemoStore' import { InlineLink } from '@/components/ui/InlineLink' import { useIsFeatureEnabled } from '@/hooks/misc/useIsFeatureEnabled' +import { useIsWarehouseEnabled } from '@/hooks/misc/useIsWarehouseEnabled' interface DestinationTypeOption { value: DestinationType label: string description: string - icon: typeof Database + icon: ElementType<{ size?: number; className?: string }> isAlpha: boolean enabled: boolean } @@ -38,17 +42,34 @@ interface DestinationTypeGroup { options: DestinationTypeOption[] } +const LUCIDE_DESTINATION_TYPES = new Set(['Warehouse', 'Snowflake']) + +function DestinationOptionIcon({ option }: { option: DestinationTypeOption }) { + const Icon = option.icon + const className = 'shrink-0 text-foreground-light' + + if (LUCIDE_DESTINATION_TYPES.has(option.value)) { + return + } + + return +} + export const DestinationTypeSelection = () => { + const { ref: projectRef } = useParams() const etlEnableBigQuery = useIsETLBigQueryPrivateAlpha() const etlEnableIceberg = useIsETLIcebergPrivateAlpha() const etlEnableDucklake = useIsETLDucklakePrivateAlpha() const etlEnableSnowflake = useIsETLSnowflakePrivateAlpha() + const isWarehouseFeatureEnabled = useIsWarehouseEnabled() + const warehouseState = useWarehouseProjectState(projectRef) const { infrastructureReadReplicas } = useIsFeatureEnabled(['infrastructure:read_replicas']) const [urlDestinationType, setDestinationType] = useQueryState( 'destinationType', parseAsStringEnum([ 'Read Replica', + 'Warehouse', 'BigQuery', 'Analytics Bucket', 'DuckLake', @@ -91,6 +112,17 @@ export const DestinationTypeSelection = () => { { label: 'Pipelines', options: [ + { + value: 'Warehouse', + label: 'Warehouse', + description: 'Replicate to a managed Warehouse endpoint for analytics', + icon: Warehouse, + isAlpha: true, + enabled: isOptionVisible( + 'Warehouse', + isWarehouseFeatureEnabled && !warehouseState.enabled + ), + }, { value: 'Analytics Bucket', label: 'Analytics Bucket', @@ -146,8 +178,7 @@ export const DestinationTypeSelection = () => { description={ selectedOption?.isAlpha && ( - This destination type is in alpha and may be unstable or introduce breaking changes - while we iterate based on customer feedback.{' '} + In alpha and may change.{' '} Leave feedback @@ -163,7 +194,7 @@ export const DestinationTypeSelection = () => { {selectedOption ? (
- +
{selectedOption.label} {selectedOption.isAlpha && Alpha} @@ -181,7 +212,7 @@ export const DestinationTypeSelection = () => { {group.options.map((option) => (
- +
{option.label} diff --git a/apps/studio/components/interfaces/Database/Replication/DestinationPanel/WarehouseDestinationForm.tsx b/apps/studio/components/interfaces/Database/Replication/DestinationPanel/WarehouseDestinationForm.tsx new file mode 100644 index 00000000000..4c848098231 --- /dev/null +++ b/apps/studio/components/interfaces/Database/Replication/DestinationPanel/WarehouseDestinationForm.tsx @@ -0,0 +1,89 @@ +import { useParams } from 'common' +import { useState } from 'react' +import { toast } from 'sonner' +import { Button, Checkbox, Label, SheetFooter, SheetSection } from 'ui' +import { FormItemLayout } from 'ui-patterns/form/FormItemLayout/FormItemLayout' + +import { enableWarehouseProject } from '@/components/interfaces/Database/Warehouse/warehouseDemoStore' +import { useDefaultWarehouseSchemas } from '@/components/interfaces/Database/Warehouse/WarehouseSchemaScope' + +interface WarehouseDestinationFormProps { + onClose: () => void +} + +export function WarehouseDestinationForm({ onClose }: WarehouseDestinationFormProps) { + const { ref: projectRef } = useParams() + const defaultSchemas = useDefaultWarehouseSchemas() + + const [selectedSchemas, setSelectedSchemas] = useState(() => defaultSchemas) + const [isEnabling, setIsEnabling] = useState(false) + + const selectedSet = new Set(selectedSchemas) + + const onToggleSchema = (schema: string, checked: boolean) => { + setSelectedSchemas((prev) => (checked ? [...prev, schema] : prev.filter((s) => s !== schema))) + } + + const onEnable = async () => { + if (!projectRef || selectedSchemas.length === 0) return + setIsEnabling(true) + enableWarehouseProject(projectRef, selectedSchemas) + setIsEnabling(false) + toast.success('Warehouse replication started') + onClose() + } + + return ( + <> + +
+

+ Replicate your Postgres data to a separate Warehouse endpoint. Tables keep the same + schema and table names (for example, public.events). +

+
+

Replicated schemas

+

+ Choose which schemas to replicate. You can change this later. +

+
+ {defaultSchemas.map((schema) => ( + + onToggleSchema(schema, checked === true)} + /> + + + ))} +
+ +
+
+
+ + + + + + ) +} diff --git a/apps/studio/components/interfaces/Database/Replication/Destinations.tsx b/apps/studio/components/interfaces/Database/Replication/Destinations.tsx index 81c51548ba4..1482b7dd4ec 100644 --- a/apps/studio/components/interfaces/Database/Replication/Destinations.tsx +++ b/apps/studio/components/interfaces/Database/Replication/Destinations.tsx @@ -38,6 +38,7 @@ import { } from './useIsETLPrivateAlpha' import { useWarehouseProjectState } from '@/components/interfaces/Database/Warehouse/warehouseDemoStore' import { WarehouseDestinationRow } from '@/components/interfaces/Database/Warehouse/WarehouseDestinationRow' +import { WarehouseManageSheet } from '@/components/interfaces/Database/Warehouse/WarehouseManageSheet' import { AlertError } from '@/components/ui/AlertError' import { DocsButton } from '@/components/ui/DocsButton' import { Shortcut } from '@/components/ui/Shortcut' @@ -67,17 +68,22 @@ export const Destinations = () => { const hasManagedWarehouse = isWarehouseFeatureEnabled && warehouseState.enabled const { infrastructureReadReplicas } = useIsFeatureEnabled(['infrastructure:read_replicas']) - const newDestinationDefaultType = infrastructureReadReplicas - ? 'Read Replica' - : etlEnableBigQuery - ? 'BigQuery' - : etlEnableIceberg - ? 'Analytics Bucket' - : etlEnableDucklake - ? 'DuckLake' - : etlEnableSnowflake - ? 'Snowflake' - : null + const newDestinationDefaultType: DestinationType | null = + isWarehouseFeatureEnabled && !warehouseState.enabled + ? 'Warehouse' + : infrastructureReadReplicas + ? 'Read Replica' + : etlEnableBigQuery + ? 'BigQuery' + : etlEnableIceberg + ? 'Analytics Bucket' + : etlEnableDucklake + ? 'DuckLake' + : etlEnableSnowflake + ? 'Snowflake' + : isWarehouseFeatureEnabled + ? 'Warehouse' + : null const prefetchedRef = useRef(false) const searchInputRef = useRef(null) @@ -89,6 +95,7 @@ export const Destinations = () => { 'destinationType', parseAsStringEnum([ 'Read Replica', + 'Warehouse', 'BigQuery', 'Analytics Bucket', 'DuckLake', @@ -361,7 +368,7 @@ export const Destinations = () => {
setStatusRefetchInterval(5000)} /> +
diff --git a/apps/studio/components/interfaces/Database/Warehouse/WarehouseProjectCard.tsx b/apps/studio/components/interfaces/Database/Warehouse/WarehouseProjectCard.tsx deleted file mode 100644 index 1a6f865810f..00000000000 --- a/apps/studio/components/interfaces/Database/Warehouse/WarehouseProjectCard.tsx +++ /dev/null @@ -1,84 +0,0 @@ -import { useParams } from 'common' -import { parseAsBoolean, useQueryState } from 'nuqs' -import { useState } from 'react' -import { Button, Card, CardContent } from 'ui' - -import { useWarehouseProjectState } from './warehouseDemoStore' -import { WarehouseDisableModal } from './WarehouseDisableModal' -import { WarehouseEnableModal } from './WarehouseEnableModal' -import { WarehouseObservabilityPanel } from './WarehouseObservabilityPanel' -import { WarehouseSchemaScope } from './WarehouseSchemaScope' -import { WarehouseSyncChip } from './WarehouseSyncChip' -import { useIsWarehouseEnabled } from '@/hooks/misc/useIsWarehouseEnabled' - -export function WarehouseProjectCard() { - const { ref: projectRef } = useParams() - const isWarehouseFeatureEnabled = useIsWarehouseEnabled() - const warehouseState = useWarehouseProjectState(projectRef) - - const [showEnableModal, setShowEnableModal] = useState(false) - const [showDisableModal, setShowDisableModal] = useState(false) - const [, setShowConnect] = useQueryState('showConnect', parseAsBoolean.withDefault(false)) - const [, setConnectTab] = useQueryState('connectTab') - - if (!isWarehouseFeatureEnabled) return null - - return ( - <> - - -
-
-
-

Warehouse

- {warehouseState.enabled && ( - - )} -
-

- {warehouseState.enabled - ? 'Analytical replica of your Postgres data on a separate Warehouse endpoint. Query the same schema and table names against the Warehouse host.' - : 'Enable an analytical replica of your Postgres data on a separate Warehouse endpoint for analytics workloads.'} -

-
-
- {warehouseState.enabled ? ( - <> - - - - ) : ( - - )} -
-
- - {warehouseState.enabled && ( - <> - - - - )} -
-
- - - - - ) -} diff --git a/apps/studio/components/interfaces/Database/Warehouse/WarehouseReplicationPipelineStatus.tsx b/apps/studio/components/interfaces/Database/Warehouse/WarehouseReplicationPipelineStatus.tsx new file mode 100644 index 00000000000..4b9eb25455d --- /dev/null +++ b/apps/studio/components/interfaces/Database/Warehouse/WarehouseReplicationPipelineStatus.tsx @@ -0,0 +1,132 @@ +import { useParams } from 'common' +import { ChevronLeft } from 'lucide-react' +import Link from 'next/link' +import { useMemo } from 'react' +import { Badge, Button, Table, TableBody, TableCell, TableHead, TableHeader, TableRow } from 'ui' +import { GenericSkeletonLoader } from 'ui-patterns/ShimmeringLoader' + +import { isTableReplicated, useWarehouseProjectState } from './warehouseDemoStore' +import { WarehouseObservabilityPanel } from './WarehouseObservabilityPanel' +import { formatReplicationLagSeconds } from './warehouseReplication.utils' +import { WarehouseSyncChip } from './WarehouseSyncChip' +import { useTablesQuery } from '@/data/tables/tables-query' + +export function WarehouseReplicationPipelineStatus() { + const { ref: projectRef } = useParams() + const state = useWarehouseProjectState(projectRef) + + const { data: tables, isPending: isTablesLoading } = useTablesQuery( + { projectRef }, + { enabled: !!projectRef && state.enabled } + ) + + const replicatedTables = useMemo(() => { + if (!tables) return [] + return tables + .filter((table) => isTableReplicated(projectRef, table.schema, table.name)) + .sort((a, b) => `${a.schema}.${a.name}`.localeCompare(`${b.schema}.${b.name}`)) + }, [tables, projectRef]) + + if (!state.enabled || !state.pipelineId) return null + + const destinationName = `DuckLake (Pipeline ID: ${state.pipelineId})` + const isPipelineStopped = + state.pipelineStatus === 'stopped' && state.replicationPhase === 'streaming' + const tableStatusLabel = + state.replicationPhase === 'failed' + ? 'Error' + : state.replicationPhase === 'backfilling' + ? 'Backfilling' + : isPipelineStopped + ? 'Paused' + : state.replicationPhase === 'streaming' + ? 'Replicating' + : 'Pending' + + return ( +
+
+
+ +
+

{destinationName}

+ +
+
+
+ + + +
+
+

Table replication

+

+ Tables in replicated schemas use the same schema and table names on the Warehouse + endpoint. +

+
+ + {isTablesLoading ? ( +
+ +
+ ) : replicatedTables.length === 0 ? ( +

No replicated tables found.

+ ) : ( + + + + Table + Status + Replication lag + + + + {replicatedTables.map((table) => ( + + + {table.schema}.{table.name} + + + + {tableStatusLabel} + + + + {isPipelineStopped || state.replicationPhase === 'provisioning' + ? '—' + : formatReplicationLagSeconds(state.lagSeconds)} + + + ))} + +
+ )} +
+
+ ) +} + +export function isWarehouseMockPipelineId( + projectRef: string | undefined, + pipelineId: string | undefined, + warehouseState: ReturnType +): boolean { + if (!projectRef || !pipelineId || !warehouseState.enabled) return false + return warehouseState.pipelineId === pipelineId +} diff --git a/apps/studio/components/interfaces/Database/Warehouse/WarehouseSchemaScope.tsx b/apps/studio/components/interfaces/Database/Warehouse/WarehouseSchemaScope.tsx index 6ab1c96294f..c9437455242 100644 --- a/apps/studio/components/interfaces/Database/Warehouse/WarehouseSchemaScope.tsx +++ b/apps/studio/components/interfaces/Database/Warehouse/WarehouseSchemaScope.tsx @@ -10,6 +10,7 @@ import { useWarehouseProjectState, } from './warehouseDemoStore' import { useSchemasQuery } from '@/data/database/schemas-query' +import { EMPTY_ARR } from '@/lib/void' interface WarehouseSchemaScopeProps { disabled?: boolean @@ -17,12 +18,12 @@ interface WarehouseSchemaScopeProps { export function WarehouseSchemaScope({ disabled = false }: WarehouseSchemaScopeProps) { const { ref: projectRef } = useParams() - const { data: schemas = [] } = useSchemasQuery({ projectRef }) + const { data: schemas } = useSchemasQuery({ projectRef }) const warehouseState = useWarehouseProjectState(projectRef) const replicableSchemas = useMemo( () => - schemas + (schemas ?? EMPTY_ARR) .map((s) => s.name) .filter(isReplicableSchema) .sort(), @@ -80,6 +81,9 @@ export function WarehouseSchemaScope({ disabled = false }: WarehouseSchemaScopeP export function useDefaultWarehouseSchemas(): string[] { const { ref: projectRef } = useParams() - const { data: schemas = [] } = useSchemasQuery({ projectRef }) - return useMemo(() => getDefaultIncludedSchemas(schemas.map((s) => s.name)), [schemas]) + const { data: schemas } = useSchemasQuery({ projectRef }) + return useMemo( + () => getDefaultIncludedSchemas((schemas ?? EMPTY_ARR).map((s) => s.name)), + [schemas] + ) } diff --git a/apps/studio/components/interfaces/Database/Warehouse/WarehouseSyncChip.tsx b/apps/studio/components/interfaces/Database/Warehouse/WarehouseSyncChip.tsx index 6df6d8c4389..0696db9ced2 100644 --- a/apps/studio/components/interfaces/Database/Warehouse/WarehouseSyncChip.tsx +++ b/apps/studio/components/interfaces/Database/Warehouse/WarehouseSyncChip.tsx @@ -1,23 +1,26 @@ import { Badge } from 'ui' -import type { ReplicationPhase } from './warehouseDemoStore' +import type { ReplicationPhase, WarehousePipelineStatus } from './warehouseDemoStore' import { getWarehouseDestinationStatusLabel } from './warehouseReplication.utils' interface WarehouseSyncChipProps { phase: ReplicationPhase + pipelineStatus?: WarehousePipelineStatus } -export function WarehouseSyncChip({ phase }: WarehouseSyncChipProps) { - const label = getWarehouseDestinationStatusLabel(phase) +export function WarehouseSyncChip({ phase, pipelineStatus = 'running' }: WarehouseSyncChipProps) { + const label = getWarehouseDestinationStatusLabel(phase, pipelineStatus) const variant = - phase === 'failed' - ? 'destructive' - : phase === 'backfilling' || phase === 'provisioning' - ? 'warning' - : phase === 'streaming' - ? 'success' - : 'default' + pipelineStatus === 'stopped' && phase === 'streaming' + ? 'default' + : phase === 'failed' + ? 'destructive' + : phase === 'backfilling' || phase === 'provisioning' + ? 'warning' + : phase === 'streaming' + ? 'success' + : 'default' return {label} } diff --git a/apps/studio/components/interfaces/Database/Warehouse/managedWarehouse.resources.ts b/apps/studio/components/interfaces/Database/Warehouse/managedWarehouse.resources.ts index 1cb1fd73615..8b5427bbb14 100644 --- a/apps/studio/components/interfaces/Database/Warehouse/managedWarehouse.resources.ts +++ b/apps/studio/components/interfaces/Database/Warehouse/managedWarehouse.resources.ts @@ -5,4 +5,4 @@ export function isSupabaseManagedWarehousePublicationName(name: string): boolean } export const MANAGED_WAREHOUSE_PUBLICATION_TOOLTIP = - 'Managed by Supabase for Warehouse replication. Configure schemas on the Warehouse card above.' + 'Managed by Supabase for Warehouse replication. Use Manage to update replicated schemas.' diff --git a/apps/studio/components/interfaces/Database/Warehouse/warehouseDemoStore.ts b/apps/studio/components/interfaces/Database/Warehouse/warehouseDemoStore.ts index 37ccce76595..5babf58f634 100644 --- a/apps/studio/components/interfaces/Database/Warehouse/warehouseDemoStore.ts +++ b/apps/studio/components/interfaces/Database/Warehouse/warehouseDemoStore.ts @@ -6,9 +6,12 @@ import { INTERNAL_SCHEMAS } from '@/hooks/useProtectedSchemas' export type ReplicationPhase = 'idle' | 'provisioning' | 'backfilling' | 'streaming' | 'failed' +export type WarehousePipelineStatus = 'running' | 'stopped' + export interface WarehouseProjectState { enabled: boolean replicationPhase: ReplicationPhase + pipelineStatus: WarehousePipelineStatus lagSeconds: number | null includedSchemas: string[] pipelineId: string | null @@ -21,6 +24,7 @@ export interface WarehouseProjectState { const DEFAULT_STATE: WarehouseProjectState = { enabled: false, replicationPhase: 'idle', + pipelineStatus: 'stopped', lagSeconds: null, includedSchemas: [], pipelineId: null, @@ -71,7 +75,14 @@ export function getDefaultIncludedSchemas(availableSchemas: string[]): string[] } function getProjectState(projectRef: string): WarehouseProjectState { - return warehouseDemoStore.projects[projectRef] ?? { ...DEFAULT_STATE } + const raw = warehouseDemoStore.projects[projectRef] + if (!raw) return { ...DEFAULT_STATE } + return { + ...DEFAULT_STATE, + ...raw, + pipelineStatus: + raw.pipelineStatus ?? (raw.replicationPhase === 'streaming' ? 'running' : 'stopped'), + } } function setProjectState(projectRef: string, next: WarehouseProjectState): void { @@ -81,7 +92,14 @@ function setProjectState(projectRef: string, next: WarehouseProjectState): void export function useWarehouseProjectState(projectRef: string | undefined) { const snap = useSnapshot(warehouseDemoStore) if (!projectRef) return { ...DEFAULT_STATE } - return snap.projects[projectRef] ?? { ...DEFAULT_STATE } + const raw = snap.projects[projectRef] + if (!raw) return { ...DEFAULT_STATE } + return { + ...DEFAULT_STATE, + ...raw, + pipelineStatus: + raw.pipelineStatus ?? (raw.replicationPhase === 'streaming' ? 'running' : 'stopped'), + } } export function isWarehouseProjectEnabled(projectRef: string | undefined): boolean { @@ -129,6 +147,7 @@ function simulatePhaseProgression(projectRef: string): void { setProjectState(projectRef, { ...backfilling, replicationPhase: 'streaming', + pipelineStatus: 'running', lagSeconds: 2, }) }, 6000) @@ -140,6 +159,7 @@ export function enableWarehouseProject(projectRef: string, includedSchemas: stri setProjectState(projectRef, { enabled: true, replicationPhase: 'provisioning', + pipelineStatus: 'running', lagSeconds: null, includedSchemas, pipelineId, @@ -156,6 +176,45 @@ export function disableWarehouseProject(projectRef: string): void { setProjectState(projectRef, { ...DEFAULT_STATE }) } +export function stopWarehousePipeline(projectRef: string): void { + const current = getProjectState(projectRef) + if (!current.enabled || current.replicationPhase !== 'streaming') return + setProjectState(projectRef, { ...current, pipelineStatus: 'stopped' }) +} + +export function startWarehousePipeline(projectRef: string): void { + const current = getProjectState(projectRef) + if (!current.enabled || current.pipelineStatus !== 'stopped') return + setProjectState(projectRef, { + ...current, + pipelineStatus: 'running', + lagSeconds: current.lagSeconds ?? 2, + }) +} + +export function restartWarehousePipeline(projectRef: string): void { + const current = getProjectState(projectRef) + if (!current.enabled) return + clearPhaseTimer() + setProjectState(projectRef, { + ...current, + pipelineStatus: 'running', + replicationPhase: 'backfilling', + lagSeconds: 30, + }) + + phaseTimer = setTimeout(() => { + const state = getProjectState(projectRef) + if (!state.enabled || state.replicationPhase !== 'backfilling') return + setProjectState(projectRef, { + ...state, + replicationPhase: 'streaming', + pipelineStatus: 'running', + lagSeconds: 2, + }) + }, 4000) +} + export function updateWarehouseIncludedSchemas( projectRef: string, includedSchemas: string[] @@ -178,6 +237,7 @@ export function updateWarehouseIncludedSchemas( setProjectState(projectRef, { ...state, replicationPhase: 'streaming', + pipelineStatus: 'running', lagSeconds: 2, }) }, 4000) diff --git a/apps/studio/components/interfaces/Database/Warehouse/warehouseReplication.utils.ts b/apps/studio/components/interfaces/Database/Warehouse/warehouseReplication.utils.ts index d661110d1c7..4bf32bb2db7 100644 --- a/apps/studio/components/interfaces/Database/Warehouse/warehouseReplication.utils.ts +++ b/apps/studio/components/interfaces/Database/Warehouse/warehouseReplication.utils.ts @@ -1,4 +1,8 @@ -import type { ReplicationPhase, WarehouseProjectState } from './warehouseDemoStore' +import type { + ReplicationPhase, + WarehousePipelineStatus, + WarehouseProjectState, +} from './warehouseDemoStore' export function formatReplicationPhase(phase: ReplicationPhase): string { switch (phase) { @@ -25,7 +29,11 @@ export function formatReplicationLagSeconds(lagSeconds: number | null): string { return `${(lagSeconds / 3600).toFixed(1)}h` } -export function getWarehouseDestinationStatusLabel(phase: ReplicationPhase): string { +export function getWarehouseDestinationStatusLabel( + phase: ReplicationPhase, + pipelineStatus: WarehousePipelineStatus = 'running' +): string { + if (pipelineStatus === 'stopped' && phase === 'streaming') return 'Stopped' switch (phase) { case 'provisioning': return 'Starting' @@ -58,5 +66,10 @@ export function buildWarehouseConnectionString({ } export function isWarehouseReplicationHealthy(state: WarehouseProjectState): boolean { - return state.enabled && state.replicationPhase === 'streaming' && (state.lagSeconds ?? 0) < 120 + return ( + state.enabled && + state.replicationPhase === 'streaming' && + state.pipelineStatus === 'running' && + (state.lagSeconds ?? 0) < 120 + ) } diff --git a/apps/studio/pages/project/[ref]/database/replication/[pipelineId].tsx b/apps/studio/pages/project/[ref]/database/replication/[pipelineId].tsx index e1caff23d17..1e4625f801d 100644 --- a/apps/studio/pages/project/[ref]/database/replication/[pipelineId].tsx +++ b/apps/studio/pages/project/[ref]/database/replication/[pipelineId].tsx @@ -4,18 +4,29 @@ import { useContext, useEffect } from 'react' import { ReplicationPipelineStatus } from '@/components/interfaces/Database/Replication/ReplicationPipelineStatus/ReplicationPipelineStatus' import { useIsETLPrivateAlpha } from '@/components/interfaces/Database/Replication/useIsETLPrivateAlpha' +import { useWarehouseProjectState } from '@/components/interfaces/Database/Warehouse/warehouseDemoStore' +import { + isWarehouseMockPipelineId, + WarehouseReplicationPipelineStatus, +} from '@/components/interfaces/Database/Warehouse/WarehouseReplicationPipelineStatus' import DatabaseLayout from '@/components/layouts/DatabaseLayout/DatabaseLayout' import { DefaultLayout } from '@/components/layouts/DefaultLayout' import { ScaffoldContainer, ScaffoldSection } from '@/components/layouts/Scaffold' import { FormHeader } from '@/components/ui/Forms/FormHeader' +import { useIsWarehouseEnabled } from '@/hooks/misc/useIsWarehouseEnabled' import { PipelineRequestStatusProvider } from '@/state/replication-pipeline-request-status' import type { NextPageWithLayout } from '@/types' const DatabaseReplicationPage: NextPageWithLayout = () => { const router = useRouter() const { ref: projectRef } = useParams() + const pipelineId = router.query.pipelineId as string | undefined const { hasLoaded } = useContext(FeatureFlagContext) const enablePgReplicate = useIsETLPrivateAlpha() + const isWarehouseFeatureEnabled = useIsWarehouseEnabled() + const warehouseState = useWarehouseProjectState(projectRef) + const isWarehouseMockPipeline = + isWarehouseFeatureEnabled && isWarehouseMockPipelineId(projectRef, pipelineId, warehouseState) useEffect(() => { if (hasLoaded && !enablePgReplicate) { @@ -31,7 +42,11 @@ const DatabaseReplicationPage: NextPageWithLayout = () => {
- + {isWarehouseMockPipeline ? ( + + ) : ( + + )}
diff --git a/apps/studio/pages/project/[ref]/database/replication/index.tsx b/apps/studio/pages/project/[ref]/database/replication/index.tsx index 6a8dd329afa..75db3f5abc1 100644 --- a/apps/studio/pages/project/[ref]/database/replication/index.tsx +++ b/apps/studio/pages/project/[ref]/database/replication/index.tsx @@ -2,7 +2,6 @@ import { GenericSkeletonLoader } from 'ui-patterns/ShimmeringLoader' import { Destinations } from '@/components/interfaces/Database/Replication/Destinations' import { ReplicationDiagram } from '@/components/interfaces/Database/Replication/ReplicationDiagram' -import { WarehouseProjectCard } from '@/components/interfaces/Database/Warehouse/WarehouseProjectCard' import DatabaseLayout from '@/components/layouts/DatabaseLayout/DatabaseLayout' import { DefaultLayout } from '@/components/layouts/DefaultLayout' import { ScaffoldContainer, ScaffoldSection } from '@/components/layouts/Scaffold' @@ -43,8 +42,7 @@ const DatabaseReplicationPage: NextPageWithLayout = () => {

Replication

- Deploy Read Replicas across multiple regions, or use Pipelines to replicate database - changes to analytics destinations. + Replicate to another region, an analytics destination, or Warehouse.

@@ -56,11 +54,6 @@ const DatabaseReplicationPage: NextPageWithLayout = () => { ) : ( <> - - - - -