Files
supabase/apps/studio/data/database-queues/database-queue-messages-infinite-query.ts
oniani1 3d101e2415 fix(studio): paginate queue messages on a unique cursor (#47016)
Closes #47015

## What kind of change does this PR introduce?

Bug fix.

## What is the current behavior?

The queue message list paginates with `WHERE enqueued_at > <last>` and
`ORDER BY enqueued_at`. `enqueued_at` is not unique: pgmq defaults it to
`now()`, so every message sent in one `send_batch` shares a timestamp.
When a group of same-timestamp messages straddles a page boundary, the
strict cursor skips the rest of that group, so those messages never
appear in the grid even though they are still in the queue. With 40
messages from one batch, only 30 render.

## What is the new behavior?

Pagination uses a composite `(enqueued_at, msg_id)` keyset cursor and
orders by the same pair. `msg_id` is unique within each queue/archive
table and breaks the tie, so no rows are dropped between pages. After
the change, all 40 messages render. This mirrors the cron-runs query,
which already paginates on a unique key.

## Additional context

Added a test in
`apps/studio/data/database-queues/database-queue-messages-infinite-query.test.ts`
asserting next pages use the composite cursor and order by `enqueued_at,
msg_id`.


<!-- This is an auto-generated comment: release notes by coderabbit.ai
-->

## Summary by CodeRabbit

* **Bug Fixes**
* Improved database queue message pagination to reliably retrieve all
messages, including those with identical enqueued timestamps, preventing
potential message skipping during pagination.

* **Tests**
  * Added test coverage for database queue message pagination behavior.

<!-- end of auto-generated comment: release notes by coderabbit.ai -->

Co-authored-by: Ali Waseem <waseema393@gmail.com>
2026-06-18 07:47:33 -06:00

144 lines
5.0 KiB
TypeScript

import { ident, literal, safeSql, type SafeSqlFragment } from '@supabase/pg-meta/src/pg-format'
import { InfiniteData, useInfiniteQuery } from '@tanstack/react-query'
import dayjs from 'dayjs'
import { last } from 'lodash'
import { databaseQueuesKeys } from './keys'
import {
isQueueNameValid,
pgmqArchiveTable,
pgmqQueueTable,
} from '@/components/interfaces/Integrations/Queues/Queues.utils'
import { QUEUE_MESSAGE_TYPE } from '@/components/interfaces/Integrations/Queues/SingleQueue/Queue.utils'
import { executeSql } from '@/data/sql/execute-sql-mutation'
import { DATE_FORMAT } from '@/lib/constants'
import type { ResponseError, UseCustomInfiniteQueryOptions } from '@/types'
export type DatabaseQueueVariables = {
projectRef?: string
connectionString?: string | null
queueName: string
status: QUEUE_MESSAGE_TYPE[]
}
export type PostgresQueueMessage = {
msg_id: number
read_ct: number
enqueued_at: string
archived_at: string
vt: Date
message: Record<string, never>
}
export const QUEUE_MESSAGES_PAGE_SIZE = 30
export type QueueMessagesPageParam = { enqueuedAt: string; msgId: number }
export async function getDatabaseQueue({
projectRef,
connectionString,
queueName,
after,
status,
}: DatabaseQueueVariables & { after: QueueMessagesPageParam | undefined }) {
if (!projectRef) throw new Error('Project ref is required')
if (!isQueueNameValid(queueName)) {
throw new Error(
'Invalid queue name: must contain only alphanumeric characters, underscores, and hyphens'
)
}
if (status.length === 0) {
return []
}
// handles when scheduled and available are deselected
const queueTable = safeSql`${ident('pgmq')}.${ident(pgmqQueueTable(queueName))}`
const archivedTable = safeSql`${ident('pgmq')}.${ident(pgmqArchiveTable(queueName))}`
const nowLiteral = literal(dayjs(new Date()).format(DATE_FORMAT))
let queueQuery: SafeSqlFragment | null = null
if (status.includes('available') && status.includes('scheduled')) {
queueQuery = safeSql`SELECT msg_id, enqueued_at, read_ct, vt, message, NULL as archived_at FROM ${queueTable}`
} else if (status.includes('available') && !status.includes('scheduled')) {
queueQuery = safeSql`SELECT msg_id, enqueued_at, read_ct, vt, message, NULL as archived_at FROM ${queueTable} WHERE vt < ${nowLiteral}`
} else if (!status.includes('available') && status.includes('scheduled')) {
queueQuery = safeSql`SELECT msg_id, enqueued_at, read_ct, vt, message, NULL as archived_at FROM ${queueTable} WHERE vt > ${nowLiteral}`
}
const archivedQuery = status.includes('archived')
? safeSql`SELECT msg_id, enqueued_at, read_ct, vt, message, archived_at FROM ${archivedTable}`
: null
const unionParts = [queueQuery, archivedQuery].filter(
(part): part is SafeSqlFragment => part !== null
)
const unionFragment = unionParts.reduce(
(acc, part, index) => (index === 0 ? part : safeSql`${acc} UNION ALL ${part}`),
safeSql``
)
// Keyset pagination on a composite (enqueued_at, msg_id) cursor. enqueued_at is
// not unique: pgmq defaults it to now(), so every message from one send_batch
// shares a timestamp. A plain `enqueued_at > last` cursor skips the rows that
// share the last page's timestamp, dropping them from the list. msg_id is unique
// within each queue/archive table and breaks the tie.
const whereClause = after
? safeSql` WHERE (enqueued_at, msg_id) > (${literal(after.enqueuedAt)}, ${literal(after.msgId)})`
: safeSql``
const sql = safeSql`SELECT
*
FROM
(
${unionFragment}
) AS combined${whereClause} order by enqueued_at, msg_id LIMIT ${literal(QUEUE_MESSAGES_PAGE_SIZE)}`
const { result } = await executeSql({
projectRef,
connectionString,
sql,
})
return result as DatabaseQueueData
}
export type DatabaseQueueData = PostgresQueueMessage[]
export type DatabaseQueueError = ResponseError
export const useQueueMessagesInfiniteQuery = <TData = DatabaseQueueData>(
{ projectRef, connectionString, queueName, status }: DatabaseQueueVariables,
{
enabled = true,
...options
}: UseCustomInfiniteQueryOptions<
DatabaseQueueData,
DatabaseQueueError,
InfiniteData<TData>,
readonly unknown[],
QueueMessagesPageParam | undefined
> = {}
) =>
useInfiniteQuery({
queryKey: databaseQueuesKeys.getMessagesInfinite(projectRef, queueName, { status }),
queryFn: ({ pageParam }) => {
return getDatabaseQueue({
projectRef,
connectionString,
queueName,
after: pageParam,
status,
})
},
staleTime: 0,
enabled: enabled && typeof projectRef !== 'undefined',
initialPageParam: undefined as QueueMessagesPageParam | undefined,
getNextPageParam(lastPage) {
const hasNextPage = lastPage.length >= QUEUE_MESSAGES_PAGE_SIZE
if (!hasNextPage) return undefined
const lastRow = last(lastPage)
if (!lastRow) return undefined
return { enqueuedAt: lastRow.enqueued_at, msgId: lastRow.msg_id }
},
...options,
})