fix(pipelines): Make pipeline actions and status updates reliable (#50085)

## Summary

Make pipeline actions and status feedback reliable while requests are
running or fail. Let the backend coordinate table resets and restarts,
keep stopped pipelines stopped after resets or settings changes, and
refresh the UI from confirmed backend state.

## Pipeline actions and recovery

- Reset one table, all errored tables, or all tables through the
rollback endpoint without separate frontend stop/start requests. Explain
which destination data is deleted, which rows are copied again, initial
sync charges, and the skip-initial-sync setting.
- Keep pending feedback until the action and a fresh status read finish,
including across navigation and polling errors. Prevent overlapping
actions and disable start/stop controls when status is unavailable or
transitioning.
- Close the creation form once the pipeline is created. If its initial
start fails, users can retry Start on the existing pipeline without
creating a duplicate.
- Wait for confirmed shutdown before deletion; a shutdown error or
timeout leaves deletion retryable. Keep failed version updates open and
avoid reporting success.
- Clarify recovery guidance and pending labels, suppress duplicate error
toasts, and hide stale table errors during transitions.

## Status updates and shared UI

- Poll pipeline status and table metrics one second after each response,
share in-flight reads, pause dashboard polling in background tabs, and
respect rate-limit backoff. The shutdown waiter continues in the
background.
- Refresh metadata after mutations even when an older read is in flight,
while preserving shared polling requests. Refresh affected data after
failures that may follow a committed reset or settings change.
- Move pending request state into the shared, project-keyed
`DatabaseLayout` so the list, detail page, and diagram stay consistent.
The surrounding database-page changes update named imports in both
Next.js and TanStack routes.
- Simplify action, status, and form rendering; announce status changes
to assistive technology; and sort table statuses without mutating cached
data.

---------

Co-authored-by: Joshen Lim <joshenlimek@gmail.com>
Co-authored-by: Danny White <3104761+dnywh@users.noreply.github.com>
This commit is contained in:
authored and GitHub committed 2026-09-18 11:32:48 +08:00
1 parent b1d2efd99e
commit edec85d1ca
67 files changed
+2164 -972

No files matched your search

@@ -2,6 +2,7 @@ import { useMutation, useQueryClient } from '@tanstack/react-query'
import type { components } from 'api-types'
import { toast } from 'sonner'
import { invalidateReplicationPipelineQueries } from './invalidate-pipeline-queries'
import { replicationKeys } from './keys'
import type {
BigQueryDestinationConfig,
@@ -216,7 +217,7 @@ export const useCreateDestinationPipelineMutation = ({
await Promise.all([
queryClient.invalidateQueries({ queryKey: replicationKeys.destinations(projectRef) }),
queryClient.invalidateQueries({ queryKey: replicationKeys.pipelines(projectRef) }),
invalidateReplicationPipelineQueries(queryClient, projectRef),
])
await onSuccess?.(data, variables, context)
@@ -1,6 +1,7 @@
import { useMutation, useQueryClient } from '@tanstack/react-query'
import { toast } from 'sonner'
import { invalidateReplicationPipelineQueries } from './invalidate-pipeline-queries'
import { replicationKeys } from './keys'
import { del, handleError } from '@/data/fetchers'
import type { ResponseError, UseCustomMutationOptions } from '@/types'
@@ -51,23 +52,10 @@ export const useDeleteDestinationPipelineMutation = ({
{
mutationFn: (vars) => deleteDestinationPipeline(vars),
async onSuccess(data, variables, context) {
const { projectRef, destinationId, pipelineId } = variables
const { projectRef } = variables
await Promise.all([
queryClient.invalidateQueries({ queryKey: replicationKeys.destinations(projectRef) }),
queryClient.invalidateQueries({ queryKey: replicationKeys.pipelines(projectRef) }),
queryClient.invalidateQueries({
queryKey: replicationKeys.pipelineById(projectRef, pipelineId),
}),
queryClient.invalidateQueries({
queryKey: replicationKeys.pipelinesStatus(projectRef, pipelineId),
}),
queryClient.invalidateQueries({
queryKey: replicationKeys.pipelinesReplicationStatus(projectRef, pipelineId),
}),
queryClient.invalidateQueries({
queryKey: replicationKeys.destinationById(projectRef, destinationId),
}),
invalidateReplicationPipelineQueries(queryClient, projectRef),
])
await onSuccess?.(data, variables, context)
@@ -0,0 +1,41 @@
import { QueryClient, QueryObserver } from '@tanstack/react-query'
import { describe, expect, test, vi } from 'vitest'
import { invalidateReplicationPipelineQueries } from './invalidate-pipeline-queries'
import { replicationKeys } from './keys'
describe('pipeline cache invalidation', () => {
test.each([
{ key: replicationKeys.pipelines('default'), shouldReplaceRead: true },
{ key: replicationKeys.pipelineById('default', 1), shouldReplaceRead: true },
{ key: replicationKeys.pipelinesVersion('default', 1), shouldReplaceRead: true },
{ key: replicationKeys.pipelinesStatus('default', 1), shouldReplaceRead: false },
{ key: replicationKeys.pipelinesReplicationStatus('default', 1), shouldReplaceRead: false },
])('refreshes metadata while sharing polls: $key', async ({ key, shouldReplaceRead }) => {
const queryClient = new QueryClient()
queryClient.setQueryData(key, 'cached')
let completeOldRead!: (value: string) => void
const oldRead = new Promise<string>((resolve) => {
completeOldRead = resolve
})
const aborted = vi.fn()
const queryFn = vi.fn(({ signal }: { signal: AbortSignal }) => {
signal.addEventListener('abort', aborted)
return queryFn.mock.calls.length === 1 ? oldRead : Promise.resolve('saved')
})
const observer = new QueryObserver(queryClient, { queryKey: key, queryFn })
const unsubscribe = observer.subscribe(() => {})
try {
const refresh = invalidateReplicationPipelineQueries(queryClient, 'default')
completeOldRead('before mutation')
await refresh
expect(queryClient.getQueryData(key)).toBe(shouldReplaceRead ? 'saved' : 'before mutation')
expect(queryFn).toHaveBeenCalledTimes(shouldReplaceRead ? 2 : 1)
expect(aborted).toHaveBeenCalledTimes(shouldReplaceRead ? 1 : 0)
} finally {
unsubscribe()
queryClient.clear()
}
})
})
@@ -0,0 +1,21 @@
import type { QueryClient, QueryKey } from '@tanstack/react-query'
import { replicationKeys } from './keys'
const isPollingQuery = ({ queryKey }: { queryKey: QueryKey }) =>
queryKey.at(-1) === 'status' || queryKey.at(-1) === 'replication-status'
export const invalidateReplicationPipelineQueries = (
queryClient: QueryClient,
projectRef: string | undefined
) => {
const queryKey = replicationKeys.pipelines(projectRef)
return Promise.all([
queryClient.invalidateQueries({ queryKey, predicate: (query) => !isPollingQuery(query) }),
// Polls will refresh again after the current read; metadata needs a post-mutation read now.
queryClient.invalidateQueries(
{ queryKey, predicate: isPollingQuery },
{ cancelRefetch: false }
),
])
}
@@ -0,0 +1,261 @@
import { focusManager, QueryClient, QueryObserver } from '@tanstack/react-query'
import { act, waitFor } from '@testing-library/react'
import type { components } from 'api-types'
import { HttpResponse } from 'msw'
import { afterEach, describe, expect, test, vi } from 'vitest'
import { replicationKeys } from './keys'
import { useReplicationPipelineReplicationStatusQuery } from './pipeline-replication-status-query'
import {
replicationPipelineStatusQueryOptions,
useReplicationPipelineStatusQuery,
waitForPipelineStopped,
} from './pipeline-status-query'
import { customRenderHook } from '@/tests/lib/custom-render'
import { addAPIMock, type APIErrorBody } from '@/tests/lib/msw'
const variables = { projectRef: 'default', pipelineId: 1 }
const statusKey = replicationKeys.pipelinesStatus('default', 1)
type StatusResponse = components['schemas']['PipelineStatusResponse_Output']
type MetricsResponse = components['schemas']['PipelineReplicationStatusResponse_Output']
const stopped: StatusResponse = { pipeline_id: 1, status: { name: 'stopped' } }
const stopping: StatusResponse = { pipeline_id: 1, status: { name: 'stopping' } }
const deferred = <T,>() => {
let resolve!: (value: T) => void
const promise = new Promise<T>((complete) => {
resolve = complete
})
return { promise, resolve }
}
afterEach(() => {
vi.useRealTimers()
focusManager.setFocused(true)
})
describe('pipeline polling', () => {
test.each([
{ endpoint: 'status', retryAfter: '30', delay: 30_000 },
{ endpoint: 'replication-status', retryAfter: '60', delay: 60_000 },
{ endpoint: 'status', retryAfter: undefined, delay: 30_000 },
] as const)(
'$endpoint respects rate-limit backoff ($retryAfter) and resumes normal polling after recovery',
async ({ endpoint, retryAfter, delay }) => {
let requests = 0
addAPIMock({
method: 'get',
path: `/platform/replication/:ref/pipelines/:pipeline_id/${endpoint}`,
response: () => {
requests += 1
if (requests === 1) {
return HttpResponse.json<APIErrorBody>(
{ message: 'Rate limited' },
{ status: 429, headers: retryAfter ? { 'Retry-After': retryAfter } : undefined }
)
}
return endpoint === 'status'
? HttpResponse.json<StatusResponse>(stopped)
: HttpResponse.json<MetricsResponse>({ pipeline_id: 1, table_statuses: [] })
},
})
const useResource =
endpoint === 'status'
? useReplicationPipelineStatusQuery
: useReplicationPipelineReplicationStatusQuery
vi.useFakeTimers()
const { result, unmount } = customRenderHook(() => useResource(variables))
await act(async () => {
await vi.advanceTimersByTimeAsync(0)
})
expect(result.current.isError).toBe(true)
await act(async () => {
await vi.advanceTimersByTimeAsync(delay - 1)
})
expect(requests).toBe(1)
await act(async () => {
await vi.advanceTimersByTimeAsync(1)
})
expect(requests).toBe(2)
await act(async () => {
await vi.advanceTimersByTimeAsync(5_000)
})
expect(result.current.isSuccess).toBe(true)
expect(requests).toBe(3)
unmount()
}
)
test.each(['status', 'replication-status'] as const)(
'shares slow %s requests and polls five seconds after completion',
async (endpoint) => {
const response = deferred<void>()
const requests = vi.fn()
const aborted = vi.fn()
addAPIMock({
method: 'get',
path: `/platform/replication/:ref/pipelines/:pipeline_id/${endpoint}`,
response: async ({ request }) => {
requests()
request.signal.addEventListener('abort', aborted)
await response.promise
return endpoint === 'status'
? HttpResponse.json<StatusResponse>(stopped)
: HttpResponse.json<MetricsResponse>({
pipeline_id: 1,
apply_lag: null,
table_statuses: [],
})
},
})
const queryClient = new QueryClient()
// Exercise a background refresh with cached data, where invalidation can otherwise
// cancel and replace a request that is already on the server.
queryClient.setQueryData(
endpoint === 'status'
? statusKey
: replicationKeys.pipelinesReplicationStatus('default', 1),
endpoint === 'status' ? stopped : { pipeline_id: 1, apply_lag: null, table_statuses: [] }
)
const useResource =
endpoint === 'status'
? useReplicationPipelineStatusQuery
: useReplicationPipelineReplicationStatusQuery
const { result, unmount } = customRenderHook(
() => ({ first: useResource(variables), second: useResource(variables) }),
{ queryClient }
)
await waitFor(() => expect(requests).toHaveBeenCalledTimes(1))
vi.useFakeTimers()
await act(async () => {
await vi.advanceTimersByTimeAsync(5_000)
})
expect(requests).toHaveBeenCalledTimes(1)
let refresh!: Promise<void>
act(() => {
refresh = queryClient.invalidateQueries(
{
queryKey:
endpoint === 'status'
? statusKey
: replicationKeys.pipelinesReplicationStatus('default', 1),
},
{ cancelRefetch: false }
)
})
await act(async () => {
response.resolve()
await refresh
})
expect(aborted).not.toHaveBeenCalled()
expect(requests).toHaveBeenCalledTimes(1)
await act(async () => {
await vi.advanceTimersByTimeAsync(4_999)
})
expect(result.current.first.isSuccess).toBe(true)
expect(result.current.second.isSuccess).toBe(true)
expect(requests).toHaveBeenCalledTimes(1)
await act(async () => {
await vi.advanceTimersByTimeAsync(1)
})
expect(requests).toHaveBeenCalledTimes(2)
unmount()
queryClient.clear()
}
)
test('pauses dashboard polling while unfocused and refreshes on return', async () => {
const requests = vi.fn()
addAPIMock({
method: 'get',
path: '/platform/replication/:ref/pipelines/:pipeline_id/status',
response: () => {
requests()
return HttpResponse.json<StatusResponse>(stopped)
},
})
const { result } = customRenderHook(() => useReplicationPipelineStatusQuery(variables))
await waitFor(() => expect(result.current.isSuccess).toBe(true))
vi.useFakeTimers()
focusManager.setFocused(false)
await act(async () => {
await vi.advanceTimersByTimeAsync(5_000)
})
expect(requests).toHaveBeenCalledTimes(1)
await act(async () => {
focusManager.setFocused(true)
await vi.advanceTimersByTimeAsync(0)
})
expect(requests).toHaveBeenCalledTimes(2)
})
})
describe('waiting for pipeline shutdown', () => {
test('checks fresh status even when stopped is cached, sharing a dashboard request', async () => {
const response = deferred<StatusResponse>()
const requests = vi.fn()
addAPIMock({
method: 'get',
path: '/platform/replication/:ref/pipelines/:pipeline_id/status',
response: async () => {
requests()
return HttpResponse.json<StatusResponse>(await response.promise)
},
})
const queryClient = new QueryClient()
queryClient.setQueryData(statusKey, stopped)
const observer = new QueryObserver(
queryClient,
replicationPipelineStatusQueryOptions(variables)
)
const unsubscribe = observer.subscribe(() => {})
await waitFor(() => expect(requests).toHaveBeenCalledTimes(1))
const complete = vi.fn()
const shutdown = waitForPipelineStopped(queryClient, variables).then(complete)
expect(complete).not.toHaveBeenCalled()
response.resolve(stopping)
await waitFor(() => expect(queryClient.getQueryState(statusKey)?.fetchStatus).toBe('idle'))
expect(complete).not.toHaveBeenCalled()
expect(requests).toHaveBeenCalledTimes(1)
addAPIMock({
method: 'get',
path: '/platform/replication/:ref/pipelines/:pipeline_id/status',
response: () => HttpResponse.json<StatusResponse>(stopped),
})
await queryClient.invalidateQueries({ queryKey: statusKey }, { cancelRefetch: false })
await shutdown
expect(complete).toHaveBeenCalledOnce()
unsubscribe()
queryClient.clear()
})
test('rejects on timeout instead of proceeding with deletion', async () => {
addAPIMock({
method: 'get',
path: '/platform/replication/:ref/pipelines/:pipeline_id/status',
response: () => HttpResponse.json<StatusResponse>(stopping),
})
const queryClient = new QueryClient()
vi.useFakeTimers()
const shutdown = waitForPipelineStopped(queryClient, variables)
const rejection = expect(shutdown).rejects.toThrow('Pipeline is still stopping')
await vi.advanceTimersByTimeAsync(30_000)
await rejection
expect(queryClient.getQueryCache().find({ queryKey: statusKey })?.getObserversCount()).toBe(0)
queryClient.clear()
})
test('rejects if shutdown cannot be verified', async () => {
addAPIMock({
method: 'get',
path: '/platform/replication/:ref/pipelines/:pipeline_id/status',
response: () =>
HttpResponse.json<APIErrorBody>({ message: 'Status unavailable' }, { status: 503 }),
})
const queryClient = new QueryClient()
await expect(waitForPipelineStopped(queryClient, variables)).rejects.toMatchObject({
message: 'Status unavailable',
})
expect(queryClient.getQueryCache().find({ queryKey: statusKey })?.getObserversCount()).toBe(0)
queryClient.clear()
})
})
@@ -1,17 +1,23 @@
import { useQuery } from '@tanstack/react-query'
import { queryOptions, useQuery } from '@tanstack/react-query'
import { components } from 'api-types'
import { replicationKeys } from './keys'
import { replicationPollingOptions } from './polling'
import { get, handleError } from '@/data/fetchers'
import type { ResponseError, UseCustomQueryOptions } from '@/types'
type ReplicationPipelineReplicationStatusParams = { projectRef?: string; pipelineId?: number }
export type ReplicationPipelineReplicationStatusVariables = {
projectRef?: string
pipelineId?: number
}
export type ReplicationPipelineReplicationStatusError = ResponseError
export type ReplicationPipelineTableStatus =
components['schemas']['PipelineReplicationStatusResponse_Output']['table_statuses'][number]
async function fetchReplicationPipelineReplicationStatus(
{ projectRef, pipelineId }: ReplicationPipelineReplicationStatusParams,
{ projectRef, pipelineId }: ReplicationPipelineReplicationStatusVariables,
signal?: AbortSignal
) {
if (!projectRef) throw new Error('projectRef is required')
@@ -35,19 +41,43 @@ export type ReplicationPipelineReplicationStatusData = Awaited<
ReturnType<typeof fetchReplicationPipelineReplicationStatus>
>
export const useReplicationPipelineReplicationStatusQuery = <
export const replicationPipelineReplicationStatusQueryOptions = <
TData = ReplicationPipelineReplicationStatusData,
>(
{ projectRef, pipelineId }: ReplicationPipelineReplicationStatusParams,
{
enabled = true,
...options
}: UseCustomQueryOptions<ReplicationPipelineReplicationStatusData, ResponseError, TData> = {}
) =>
useQuery<ReplicationPipelineReplicationStatusData, ResponseError, TData>({
>({
projectRef,
pipelineId,
}: ReplicationPipelineReplicationStatusVariables) =>
queryOptions<
ReplicationPipelineReplicationStatusData,
ReplicationPipelineReplicationStatusError,
TData
>({
queryKey: replicationKeys.pipelinesReplicationStatus(projectRef, pipelineId),
queryFn: ({ signal }) =>
fetchReplicationPipelineReplicationStatus({ projectRef, pipelineId }, signal),
enabled: enabled && typeof projectRef !== 'undefined' && typeof pipelineId !== 'undefined',
...replicationPollingOptions,
enabled: typeof projectRef !== 'undefined' && typeof pipelineId !== 'undefined',
})
export const useReplicationPipelineReplicationStatusQuery = <
TData = ReplicationPipelineReplicationStatusData,
>(
variables: ReplicationPipelineReplicationStatusVariables,
options: UseCustomQueryOptions<
ReplicationPipelineReplicationStatusData,
ReplicationPipelineReplicationStatusError,
TData
> = {}
) =>
useQuery<
ReplicationPipelineReplicationStatusData,
ReplicationPipelineReplicationStatusError,
TData
>({
...replicationPipelineReplicationStatusQueryOptions<TData>(variables),
...options,
enabled:
options.enabled !== false &&
typeof variables.projectRef !== 'undefined' &&
typeof variables.pipelineId !== 'undefined',
})
@@ -1,11 +1,14 @@
import { queryOptions, useQuery } from '@tanstack/react-query'
import { QueryObserver, queryOptions, useQuery, type QueryClient } from '@tanstack/react-query'
import { components } from 'api-types'
import { replicationKeys } from './keys'
import { replicationPollingOptions } from './polling'
import { get, handleError } from '@/data/fetchers'
import type { ResponseError, UseCustomQueryOptions } from '@/types'
type ReplicationPipelinesStatusParams = { projectRef?: string; pipelineId?: number }
const PIPELINE_STOP_TIMEOUT_MS = 30_000
export type ReplicationPipelineStatusResponse =
components['schemas']['PipelineStatusResponse_Output']
export type ReplicationPipelineStatus = ReplicationPipelineStatusResponse['status']['name']
@@ -36,26 +39,55 @@ export type ReplicationPipelineStatusData = Awaited<
* Shared definition so callers that need many pipeline statuses at once (`useQueries`) hit the
* same cache entries as the per-pipeline hook below, rather than fetching each status twice.
*/
export const replicationPipelineStatusQueryOptions = ({
export const replicationPipelineStatusQueryOptions = <TData = ReplicationPipelineStatusData>({
projectRef,
pipelineId,
}: ReplicationPipelinesStatusParams) =>
queryOptions<ReplicationPipelineStatusData, ResponseError>({
queryOptions<ReplicationPipelineStatusData, ResponseError, TData>({
queryKey: replicationKeys.pipelinesStatus(projectRef, pipelineId),
queryFn: ({ signal }) => fetchReplicationPipelineStatus({ projectRef, pipelineId }, signal),
...replicationPollingOptions,
enabled: typeof projectRef !== 'undefined' && typeof pipelineId !== 'undefined',
})
/** Shares the dashboard's status request while waiting for shutdown before deletion. */
export function waitForPipelineStopped(
queryClient: QueryClient,
variables: ReplicationPipelinesStatusParams
): Promise<void> {
return new Promise((resolve, reject) => {
const observer = new QueryObserver(queryClient, {
...replicationPipelineStatusQueryOptions(variables),
staleTime: 0,
refetchIntervalInBackground: true,
})
const timer = setTimeout(() => {
observer.destroy()
reject(
new Error('Pipeline is still stopping. Wait for it to stop, then try deleting it again.')
)
}, PIPELINE_STOP_TIMEOUT_MS)
observer.subscribe((result) => {
if (result.fetchStatus !== 'idle') return
if (result.isError || result.data?.status.name === 'stopped') {
clearTimeout(timer)
observer.destroy()
if (result.isError) reject(result.error)
else resolve()
}
})
})
}
export const useReplicationPipelineStatusQuery = <TData = ReplicationPipelineStatusData>(
{ projectRef, pipelineId }: ReplicationPipelinesStatusParams,
{
enabled = true,
...options
}: UseCustomQueryOptions<ReplicationPipelineStatusData, ResponseError, TData> = {}
variables: ReplicationPipelinesStatusParams,
options: UseCustomQueryOptions<ReplicationPipelineStatusData, ResponseError, TData> = {}
) =>
useQuery<ReplicationPipelineStatusData, ResponseError, TData>({
queryKey: replicationKeys.pipelinesStatus(projectRef, pipelineId),
queryFn: ({ signal }) => fetchReplicationPipelineStatus({ projectRef, pipelineId }, signal),
enabled: enabled && typeof projectRef !== 'undefined' && typeof pipelineId !== 'undefined',
...replicationPipelineStatusQueryOptions<TData>(variables),
...options,
enabled:
options.enabled !== false &&
typeof variables.projectRef !== 'undefined' &&
typeof variables.pipelineId !== 'undefined',
})
+20
View File
@@ -0,0 +1,20 @@
import type { FetchStatus } from '@tanstack/react-query'
import type { ResponseError } from '@/types'
// Restart the interval after a response. Slow endpoints never accumulate overlapping polls.
export const replicationPollingOptions = {
refetchInterval: ({
state,
}: {
state: { fetchStatus: FetchStatus; error: ResponseError | null }
}) => {
if (state.fetchStatus === 'fetching') return false
const retryAfter = state.error?.retryAfter
if (retryAfter && retryAfter > 0) return Math.max(1_000, retryAfter * 1_000)
if (state.error?.code && state.error.code >= 400) return 30_000
return 5_000
},
refetchIntervalInBackground: false,
retry: false,
} as const
@@ -46,21 +46,11 @@ export const useRestartPipelineMutation = ({
async onSuccess(data, variables, context) {
const { projectRef, pipelineId } = variables
await queryClient.invalidateQueries({
queryKey: replicationKeys.pipelinesStatus(projectRef, pipelineId),
})
// [Joshen] We're manually updating the query client here as the pipeline status is async
// So setting it so starting while letting long poll update the actual status thereafter
queryClient.setQueriesData(
await queryClient.invalidateQueries(
{
queryKey: replicationKeys.pipelinesStatus(projectRef, pipelineId),
exact: true,
},
(prev) => {
if (!prev) return prev
return { ...prev, status: { name: 'starting' } }
}
{ cancelRefetch: false }
)
await onSuccess?.(data, variables, context)
@@ -68,6 +58,13 @@ export const useRestartPipelineMutation = ({
// No default error toast here: callers already show one from their try/catch around
// mutateAsync, so a default here would double up. onError is only for opt-in callers.
async onError(data, variables, context) {
await queryClient.invalidateQueries(
{
queryKey: replicationKeys.pipelinesStatus(variables.projectRef, variables.pipelineId),
},
{ cancelRefetch: false }
)
await onError?.(data, variables, context)
},
...options,
@@ -1,93 +1,39 @@
import { useMutation, useQueryClient } from '@tanstack/react-query'
import type { components } from 'api-types'
import { toast } from 'sonner'
import { replicationKeys } from './keys'
import { startPipeline } from './start-pipeline-mutation'
import { stopPipeline } from './stop-pipeline-mutation'
import { PipelineStatusName } from '@/components/interfaces/Database/Replication/Replication.constants'
import { handleError, post } from '@/data/fetchers'
import type { ResponseError, UseCustomMutationOptions } from '@/types'
export type RollbackType = 'individual' | 'full'
export type RollbackTablesTarget =
| { type: 'single_table'; table_id: number }
| { type: 'all_tables' }
| { type: 'all_errored_tables' }
export type RollbackTablesTarget = components['schemas']['RollbackTablesBody']['target']
type RollbackTablesParams = {
projectRef: string
pipelineId: number
target: RollbackTablesTarget
rollbackType: RollbackType
pipelineStatusName?: PipelineStatusName
}
type RolledBackTable = {
table_id: number
new_state: {
name: string
[key: string]: any
}
}
type RollbackTablesResponse = {
pipeline_id: number
tables: RolledBackTable[]
}
type RollbackTablesResponse = components['schemas']['RollbackTablesResponse_Output']
async function rollbackTables(
{ projectRef, pipelineId, target, rollbackType, pipelineStatusName }: RollbackTablesParams,
{ projectRef, pipelineId, target }: RollbackTablesParams,
signal?: AbortSignal
): Promise<RollbackTablesResponse> {
if (!projectRef) throw new Error('Project reference is required')
if (!pipelineId) throw new Error('Pipeline ID is required')
if (!rollbackType) throw new Error('Rollback type is required')
const { data, error } = await post(
'/platform/replication/{ref}/pipelines/{pipeline_id}/rollback-tables',
{
params: { path: { ref: projectRef, pipeline_id: pipelineId } },
body: { target, rollback_type: rollbackType },
// Production OpenAPI still includes the retired rollback_type field.
body: { target } as components['schemas']['RollbackTablesBody'],
signal,
}
)
if (error) handleError(error)
// Logic for starting the pipeline back up after a successful rollback
if (pipelineStatusName) {
const shouldStartPipelineAfterRollback = [
PipelineStatusName.STOPPED,
PipelineStatusName.STARTED,
PipelineStatusName.FAILED,
].includes(pipelineStatusName)
try {
if (pipelineStatusName === PipelineStatusName.STOPPED) {
await startPipeline({ projectRef, pipelineId })
} else if (
pipelineStatusName === PipelineStatusName.STARTED ||
pipelineStatusName === PipelineStatusName.FAILED
) {
await stopPipeline({ projectRef, pipelineId })
await startPipeline({ projectRef, pipelineId })
} else {
// [Joshen] This error sounds misleading as though the rollback failed?
throw new Error(
`Cannot apply rollback while pipeline status is ${
pipelineStatusName || 'unknown'
}. Retry once the pipeline status is started, failed, or stopped.`
)
}
} catch (error) {
if (shouldStartPipelineAfterRollback) {
throw new Error('RESTART_FAILED', { cause: error })
} else {
throw new Error('RESTART_SKIPPED')
}
}
}
return data
}
@@ -108,30 +54,43 @@ export const useRollbackTablesMutation = ({
async onSuccess(data, variables, context) {
const { projectRef, pipelineId } = variables
await Promise.all([
queryClient.invalidateQueries({
queryKey: replicationKeys.pipelinesStatus(projectRef, pipelineId),
}),
queryClient.invalidateQueries({
queryKey: replicationKeys.pipelinesReplicationStatus(projectRef, pipelineId),
}),
queryClient.invalidateQueries(
{
queryKey: replicationKeys.pipelinesStatus(projectRef, pipelineId),
},
{ cancelRefetch: false }
),
queryClient.invalidateQueries(
{
queryKey: replicationKeys.pipelinesReplicationStatus(projectRef, pipelineId),
},
{ cancelRefetch: false }
),
])
await onSuccess?.(data, variables, context)
},
async onError(data, variables, context) {
// A reset can commit before runtime recreation fails. Refresh both views after errors.
await Promise.all([
queryClient.invalidateQueries(
{
queryKey: replicationKeys.pipelinesStatus(variables.projectRef, variables.pipelineId),
},
{ cancelRefetch: false }
),
queryClient.invalidateQueries(
{
queryKey: replicationKeys.pipelinesReplicationStatus(
variables.projectRef,
variables.pipelineId
),
},
{ cancelRefetch: false }
),
])
if (onError === undefined) {
if (data.message === 'RESTART_FAILED') {
const cause = (data as Error).cause
const causeMessage = cause instanceof Error ? cause.message : undefined
toast.error(
`Rollback completed, but failed to start the pipeline${causeMessage ? `: ${causeMessage}` : ''}`
)
} else if (data.message === 'RESTART_SKIPPED') {
toast(
'Rollback completed, but the pipeline state changed before it could be resumed. Refresh the page and try again.'
)
} else {
toast.error(`Failed to rollback tables: ${data.message}`)
}
toast.error(`Failed to restart table replication: ${data.message}`)
} else {
onError(data, variables, context)
}
@@ -44,26 +44,23 @@ export const useStartPipelineMutation = ({
async onSuccess(data, variables, context) {
const { projectRef, pipelineId } = variables
await queryClient.invalidateQueries({
queryKey: replicationKeys.pipelinesStatus(projectRef, pipelineId),
})
// [Joshen] We're manually updating the query client here as the pipeline status is async
// So setting it so starting while letting long poll update the actual status thereafter
queryClient.setQueriesData(
await queryClient.invalidateQueries(
{
queryKey: replicationKeys.pipelinesStatus(projectRef, pipelineId),
exact: true,
},
(prev) => {
if (!prev) return prev
return { ...prev, status: { name: 'starting' } }
}
{ cancelRefetch: false }
)
await onSuccess?.(data, variables, context)
},
async onError(data, variables, context) {
await queryClient.invalidateQueries(
{
queryKey: replicationKeys.pipelinesStatus(variables.projectRef, variables.pipelineId),
},
{ cancelRefetch: false }
)
if (onError === undefined) {
toast.error(`Failed to start pipeline: ${data.message}`)
} else {
@@ -2,12 +2,14 @@ import { useMutation, useQueryClient } from '@tanstack/react-query'
import { toast } from 'sonner'
import { replicationKeys } from './keys'
import { waitForPipelineStopped } from './pipeline-status-query'
import { handleError, post } from '@/data/fetchers'
import type { ResponseError, UseCustomMutationOptions } from '@/types'
export type StopPipelineParams = {
projectRef: string
pipelineId: number
waitUntilStopped?: boolean
}
export async function stopPipeline(
@@ -27,28 +29,42 @@ export async function stopPipeline(
return data
}
type StartPipelineData = Awaited<ReturnType<typeof stopPipeline>>
type StopPipelineData = Awaited<ReturnType<typeof stopPipeline>>
export const useStopPipelineMutation = ({
onSuccess,
onError,
...options
}: Omit<
UseCustomMutationOptions<StartPipelineData, ResponseError, StopPipelineParams>,
UseCustomMutationOptions<StopPipelineData, ResponseError, StopPipelineParams>,
'mutationFn'
> = {}) => {
const queryClient = useQueryClient()
return useMutation<StartPipelineData, ResponseError, StopPipelineParams>({
mutationFn: (vars) => stopPipeline(vars),
return useMutation<StopPipelineData, ResponseError, StopPipelineParams>({
mutationFn: async (variables) => {
const data = await stopPipeline(variables)
if (variables.waitUntilStopped) await waitForPipelineStopped(queryClient, variables)
return data
},
async onSuccess(data, variables, context) {
const { projectRef, pipelineId } = variables
await queryClient.invalidateQueries({
queryKey: replicationKeys.pipelinesStatus(projectRef, pipelineId),
})
await queryClient.invalidateQueries(
{
queryKey: replicationKeys.pipelinesStatus(projectRef, pipelineId),
},
{ cancelRefetch: false }
)
await onSuccess?.(data, variables, context)
},
async onError(data, variables, context) {
await queryClient.invalidateQueries(
{
queryKey: replicationKeys.pipelinesStatus(variables.projectRef, variables.pipelineId),
},
{ cancelRefetch: false }
)
if (onError === undefined) {
toast.error(`Failed to stop pipeline: ${data.message}`)
} else {
@@ -3,6 +3,7 @@ import { components } from 'api-types'
import { toast } from 'sonner'
import { optionalSecret } from './destination-secret-utils'
import { invalidateReplicationPipelineQueries } from './invalidate-pipeline-queries'
import { replicationKeys } from './keys'
import type {
BigQueryDestinationConfig,
@@ -228,24 +229,23 @@ export const useUpdateDestinationPipelineMutation = ({
{
mutationFn: (vars) => updateDestinationPipeline(vars),
async onSuccess(data, variables, context) {
const { projectRef, destinationId, pipelineId } = variables
const { projectRef } = variables
// These prefixes include list, editor, pipeline status, and table metrics caches.
await Promise.all([
// Invalidate lists
queryClient.invalidateQueries({ queryKey: replicationKeys.destinations(projectRef) }),
queryClient.invalidateQueries({ queryKey: replicationKeys.pipelines(projectRef) }),
// Invalidate item-level caches used by the editor panel
queryClient.invalidateQueries({
queryKey: replicationKeys.destinationById(projectRef, destinationId),
}),
queryClient.invalidateQueries({
queryKey: replicationKeys.pipelineById(projectRef, pipelineId),
}),
invalidateReplicationPipelineQueries(queryClient, projectRef),
])
await onSuccess?.(data, variables, context)
},
async onError(data, variables, context) {
// Settings may commit before runtime recreation fails.
await Promise.all([
queryClient.invalidateQueries({
queryKey: replicationKeys.destinations(variables.projectRef),
}),
invalidateReplicationPipelineQueries(queryClient, variables.projectRef),
])
if (onError === undefined) {
toast.error(`Failed to update destination or pipeline: ${data.message}`)
} else {
@@ -9,6 +9,7 @@ type UpdatePipelineVersionParams = {
projectRef: string
pipelineId: number
versionId: number
skipStatusInvalidation?: boolean
}
async function updatePipelineVersion(
@@ -47,11 +48,20 @@ export const useUpdatePipelineVersionMutation = ({
return useMutation<UpdatePipelineVersionData, ResponseError, UpdatePipelineVersionParams>({
mutationFn: (vars) => updatePipelineVersion(vars),
async onSuccess(data, variables, context) {
const { projectRef, pipelineId } = variables
// Ensure the version dot updates promptly
await queryClient.invalidateQueries({
queryKey: replicationKeys.pipelinesVersion(projectRef, pipelineId),
})
const { projectRef, pipelineId, skipStatusInvalidation = true } = variables
await Promise.all([
queryClient.invalidateQueries({
queryKey: replicationKeys.pipelinesVersion(projectRef, pipelineId),
}),
...(skipStatusInvalidation
? []
: [
queryClient.invalidateQueries(
{ queryKey: replicationKeys.pipelinesStatus(projectRef, pipelineId) },
{ cancelRefetch: false }
),
]),
])
await onSuccess?.(data, variables, context)
},
async onError(error, variables, context) {