feat: Queues v0.5 (#30418)

* Rename the regular/normal queues to basic.

* Move the partitioned type of queues to the end.

* Add metrics for the queues.

* Add actions for postponing and deleting a message.

* Add queue actions for purging and deleting a queue.

* Minor fixes for the empty stateof the queues page.

* Add a modal to send a message to the queue.

* Minor fixes.

* Correct some button texts.

* Refactor the metrics to switch to a rougher estimate method if it timeouts.

* Fix type errors.

* More minor fixes.
This commit is contained in:
Ivan Vasilov authored and GitHub committed 2024-11-14 13:02:29 +01:00
1 parent d6b32f05e9
commit 1157a20fbd
17 files changed
+828 -105

No files matched your search

@@ -2,10 +2,17 @@ import { Rows2, Rows3, Rows4 } from 'lucide-react'
export const QUEUE_TYPES = [
{
value: 'normal',
value: 'basic',
icon: <Rows4 strokeWidth={1} />,
label: 'Normal queue',
description: 'Create a normal queue.',
label: 'Basic queue',
description: 'Create a basic queue.',
},
{
value: 'unlogged',
icon: <Rows2 strokeWidth={1} />,
label: 'Unlogged queue',
description:
'Creates an unlogged queue which loses all data on database restart. Can be useful when write throughput is more important than durability.',
},
{
value: 'partitioned',
@@ -13,11 +20,4 @@ export const QUEUE_TYPES = [
label: 'Partitioned queue',
description: 'Create a partitioned queue which is optimized for large amount of messages',
},
{
value: 'unlogged',
icon: <Rows2 strokeWidth={1} />,
label: 'Unlogged queue',
description: 'Creates an unlogged queue which loses all data on database restart.',
},
] as const
@@ -1,12 +1,11 @@
import * as Tooltip from '@radix-ui/react-tooltip'
import { PermissionAction } from '@supabase/shared-types/out/constants'
import { ExternalLink } from 'lucide-react'
import Link from 'next/link'
import { useState } from 'react'
import EnableExtensionModal from 'components/interfaces/Database/Extensions/EnableExtensionModal'
import { useProjectContext } from 'components/layouts/ProjectLayout/ProjectContext'
import ProductEmptyState from 'components/to-be-cleaned/ProductEmptyState'
import { DocsButton } from 'components/ui/DocsButton'
import { useDatabaseExtensionsQuery } from 'data/database-extensions/database-extensions-query'
import { useCheckPermissions } from 'hooks/misc/useCheckPermissions'
import { Button } from 'ui'
@@ -65,15 +64,7 @@ export const QueuesDisabledState = () => {
</Tooltip.Portal>
)}
</Tooltip.Root>
<Button asChild type="text" icon={<ExternalLink />}>
<Link
href="https://supabase.com/docs/guides/database/extensions/pgmq"
target="_blank"
rel="noreferrer"
>
Documentation
</Link>
</Button>
<DocsButton href="https://supabase.com/docs/guides/database/extensions/pgmq" />
</div>
</ProductEmptyState>
</div>
@@ -69,10 +69,10 @@ export const QueuesListing = () => {
{queues.length === 0 ? (
<div
className={
'border rounded border-default px-20 py-16 flex flex-col items-center justify-center space-y-4'
'border rounded border-default px-20 py-16 flex flex-col items-center justify-center space-y-4 border-dashed'
}
>
<p className="text-sm text-foreground-light">No queues created yet</p>
<p className="text-sm text-foreground">No queues created yet</p>
<Button onClick={() => setCreateQueueSheetShown(true)}>Add a new queue</Button>
</div>
) : (
@@ -98,9 +98,12 @@ export const QueuesListing = () => {
<Table.th key="arguments" className="table-cell">
Type
</Table.th>
<Table.th key="created_at" className="table-cell">
<Table.th key="created_at" className="table-cell w-60">
Created at
</Table.th>
<Table.th key="queue_size" className="table-cell">
<div className="flex justify-center">Size</div>
</Table.th>
<Table.th key="buttons" className="table-cell"></Table.th>
</>
}
@@ -1,22 +1,76 @@
import dayjs from 'dayjs'
import { includes, sortBy } from 'lodash'
import { ChevronRight } from 'lucide-react'
import { ChevronRight, Loader2 } from 'lucide-react'
import { useRouter } from 'next/router'
import { useProjectContext } from 'components/layouts/ProjectLayout/ProjectContext'
import Table from 'components/to-be-cleaned/Table'
import { useQueuesMetricsQuery } from 'data/database-queues/database-queues-metrics-query'
import { PostgresQueue } from 'data/database-queues/database-queues-query'
import { DATETIME_FORMAT } from 'lib/constants'
import { useRouter } from 'next/router'
interface QueuesRowsProps {
queues: PostgresQueue[]
filterString: string
}
export const QueuesRows = ({ queues: fetchedQueues, filterString }: QueuesRowsProps) => {
const QueueRow = ({ queue }: { queue: PostgresQueue }) => {
const router = useRouter()
const { project: selectedProject } = useProjectContext()
const { data: metrics, isLoading } = useQueuesMetricsQuery(
{
queueName: queue.queue_name,
projectRef: selectedProject?.ref,
connectionString: selectedProject?.connectionString,
},
{
staleTime: 30 * 1000, // 60 seconds, talk with Oli whether this is ok to call every minute
}
)
const type = queue.is_partitioned ? 'Partitioned' : queue.is_unlogged ? 'Unlogged' : 'Basic'
return (
<Table.tr
key={queue.queue_name}
onClick={() => {
router.push(`/project/${selectedProject?.ref}/integrations/queues/${queue.queue_name}`)
}}
className="hover:"
>
<Table.td className="truncate">
<p title={queue.queue_name}>{queue.queue_name}</p>
</Table.td>
<Table.td className="table-cell overflow-auto">
<p title={type.toLocaleLowerCase()} className="truncate">
{type}
</p>
</Table.td>
<Table.td className="table-cell">
<p title={queue.created_at}>{dayjs(queue.created_at).format(DATETIME_FORMAT)}</p>
</Table.td>
<Table.td className="table-cell">
<div className="flex justify-center">
{isLoading ? (
<Loader2 className="animate-spin" size={16} />
) : (
<p>
{metrics?.queue_length} {metrics?.method === 'estimated' ? '(Approximate)' : null}
</p>
)}
</div>
</Table.td>
<Table.td>
<div className="flex items-center justify-end">
<ChevronRight size="18" />
</div>
</Table.td>
</Table.tr>
)
}
export const QueuesRows = ({ queues: fetchedQueues, filterString }: QueuesRowsProps) => {
const filteredQueues = fetchedQueues.filter((x) =>
includes(x.queue_name.toLowerCase(), filterString.toLowerCase())
)
@@ -25,7 +79,7 @@ export const QueuesRows = ({ queues: fetchedQueues, filterString }: QueuesRowsPr
if (queues.length === 0 && filterString.length > 0) {
return (
<Table.tr>
<Table.td colSpan={4}>
<Table.td colSpan={5}>
<p className="text-sm text-foreground">No results found</p>
<p className="text-sm text-foreground-light">
Your search for "{filterString}" did not return any results
@@ -37,34 +91,9 @@ export const QueuesRows = ({ queues: fetchedQueues, filterString }: QueuesRowsPr
return (
<>
{queues.map((q) => {
const type = q.is_partitioned ? 'Partitioned' : q.is_unlogged ? 'Unlogged' : 'Regular'
return (
<Table.tr
key={q.queue_name}
onClick={() => {
router.push(`/project/${selectedProject?.ref}/integrations/queues/${q.queue_name}`)
}}
className="hover:"
>
<Table.td className="truncate">
<p title={q.queue_name}>{q.queue_name}</p>
</Table.td>
<Table.td className="table-cell overflow-auto">
<p title={type.toLocaleLowerCase()} className="truncate">
{type}
</p>
</Table.td>
<Table.td className="table-cell">
<p title={q.created_at}>{dayjs(q.created_at).format(DATETIME_FORMAT)}</p>
</Table.td>
<Table.td className="flex items-center justify-end">
<ChevronRight />
</Table.td>
</Table.tr>
)
})}
{queues.map((q) => (
<QueueRow key={q.queue_name} queue={q} />
))}
</>
)
}
@@ -1,12 +1,15 @@
import { useEscapeKeydown } from '@radix-ui/react-use-escape-keydown'
import { isNil, noop } from 'lodash'
import { Archive, X } from 'lucide-react'
import { Archive, Clock12, Trash2, X } from 'lucide-react'
import { useState } from 'react'
import { useParams } from 'common'
import { MonacoEditor } from 'components/grid/components/common/MonacoEditor'
import { RowAction, RowData } from 'components/interfaces/Auth/Users/UserOverview'
import { useDatabaseQueueMessageArchiveMutation } from 'data/database-queues/database-queue-messages-archive-mutation'
import { useDatabaseQueueMessageDeleteMutation } from 'data/database-queues/database-queue-messages-delete-mutation'
import { PostgresQueueMessage } from 'data/database-queues/database-queue-messages-infinite-query'
import { useDatabaseQueueMessageReadMutation } from 'data/database-queues/database-queue-messages-read-mutation'
import dayjs from 'dayjs'
import { useSelectedProject } from 'hooks/misc/useSelectedProject'
import { prettifyJSON } from 'lib/helpers'
@@ -48,7 +51,25 @@ export const MessageDetailsPanel = ({
const { name: queueName } = useParams()
const project = useSelectedProject()
const { mutate, isLoading, isSuccess } = useDatabaseQueueMessageArchiveMutation()
useEscapeKeydown(() => setSelectedMessage(null))
const {
mutate: archiveMessage,
isLoading: isLoadingArchive,
isSuccess: isSuccessArchive,
} = useDatabaseQueueMessageArchiveMutation()
const {
mutate: readMessage,
isLoading: isLoadingRead,
isSuccess: isSuccessRead,
} = useDatabaseQueueMessageReadMutation()
const {
mutate: deleteMessage,
isLoading: isLoadingDelete,
isSuccess: isSuccessDelete,
} = useDatabaseQueueMessageDeleteMutation()
const initialValue = JSON.stringify(selectedMessage?.message)
const jsonString = prettifyJSON(!isNil(initialValue) ? tryFormatInitialValue(initialValue) : '')
@@ -100,7 +121,7 @@ export const MessageDetailsPanel = ({
<RowData property="Retries" value={`${selectedMessage.read_ct}`} />
<div>
<h3 className="text-foreground-light pt-1">Payload</h3>
<h3 className="text-foreground-light py-1">Payload</h3>
<MonacoEditor
key={selectedMessage.msg_id}
onChange={noop}
@@ -112,33 +133,88 @@ export const MessageDetailsPanel = ({
</div>
</div>
<Separator />
<div className="flex flex-col px-4 py-4">
<div className="flex flex-col px-4 py-4 -space-y-1">
{!selectedMessage.archived_at ? (
<RowAction
title="Archive message"
description="The message will be marked as archived and hidden from future reads by consumers"
button={{
icon: <Archive />,
text: 'Archive',
isLoading: isLoading,
onClick: () => {
mutate({
projectRef: project!.ref,
connectionString: project?.connectionString,
queryName: queueName!,
messageId: selectedMessage.msg_id,
})
},
}}
success={
isSuccess
? {
title: 'Archived',
description: 'The message is archived successfully.',
}
: undefined
}
/>
<>
<RowAction
title="Postpone message"
description="The message will be postponed and won't show up in reads for 60 seconds."
button={{
icon: <Clock12 />,
text: 'Postpone',
isLoading: isLoadingRead,
onClick: () => {
readMessage({
projectRef: project!.ref,
connectionString: project?.connectionString,
queryName: queueName!,
messageId: selectedMessage.msg_id,
duration: 60,
})
},
}}
success={
isSuccessRead
? {
title: 'Postponed',
description: 'The message was postponed for 60 seconds.',
}
: undefined
}
/>
<RowAction
title="Archive message"
description="The message will be marked as archived and hidden from future reads by consumers. You can still access the message later."
button={{
icon: <Archive />,
text: 'Archive',
isLoading: isLoadingArchive,
type: 'warning',
onClick: () => {
archiveMessage({
projectRef: project!.ref,
connectionString: project?.connectionString,
queryName: queueName!,
messageId: selectedMessage.msg_id,
})
},
}}
success={
isSuccessArchive
? {
title: 'Archived',
description: 'The message was archived successfully.',
}
: undefined
}
/>
<RowAction
title="Delete message"
description="The message cannot be recovered afterwards."
button={{
icon: <Trash2 />,
text: 'Delete',
type: 'danger',
isLoading: isLoadingDelete,
onClick: () => {
deleteMessage({
projectRef: project!.ref,
connectionString: project?.connectionString,
queueName: queueName!,
messageId: selectedMessage.msg_id,
})
},
}}
success={
isSuccessDelete
? {
title: 'Deleted',
description: 'The message was deleted successfully.',
}
: undefined
}
/>
</>
) : null}
</div>
</TabsContent_Shadcn_>
@@ -0,0 +1,65 @@
import { useRouter } from 'next/router'
import { toast } from 'sonner'
import { useProjectContext } from 'components/layouts/ProjectLayout/ProjectContext'
import { useDatabaseQueuePurgeMutation } from 'data/database-queues/database-queues-purge-mutation'
import TextConfirmModal from 'ui-patterns/Dialogs/TextConfirmModal'
interface PurgeQueueProps {
queueName: string
visible: boolean
onClose: () => void
}
const PurgeQueue = ({ queueName, visible, onClose }: PurgeQueueProps) => {
const { project } = useProjectContext()
const router = useRouter()
const { mutate: purgeDatabaseQueue, isLoading } = useDatabaseQueuePurgeMutation({
onSuccess: () => {
toast.success(`Successfully purged queue ${queueName}`)
router.push(`/project/${project?.ref}/integrations/queues`)
onClose()
},
})
async function handlePurge() {
if (!project) return console.error('Project is required')
purgeDatabaseQueue({
queueName: queueName,
projectRef: project.ref,
connectionString: project.connectionString,
})
}
if (!queueName) {
return null
}
return (
<TextConfirmModal
variant="warning"
visible={visible}
onCancel={() => onClose()}
onConfirm={handlePurge}
title="Purge this queue"
loading={isLoading}
confirmLabel={`Purge queue ${queueName}`}
confirmPlaceholder="Type in name of queue"
confirmString={queueName ?? 'Unknown'}
text={
<>
<span>This will purge the queue</span>{' '}
<span className="text-bold text-foreground">{queueName}</span>
</>
}
alert={{
title:
"This action will delete all messages from the queue. They can't be recovered afterwards.",
}}
/>
)
}
export default PurgeQueue
@@ -6,17 +6,18 @@ import { UIEvent, useMemo, useRef } from 'react'
import DataGrid, { Column, DataGridHandle, Row } from 'react-data-grid'
import { PostgresQueueMessage } from 'data/database-queues/database-queue-messages-infinite-query'
import { Badge, ResizableHandle, ResizablePanel, ResizablePanelGroup, cn } from 'ui'
import { Badge, Button, ResizableHandle, ResizablePanel, ResizablePanelGroup, cn } from 'ui'
import { GenericSkeletonLoader } from 'ui-patterns/ShimmeringLoader'
import { DATE_FORMAT, MessageDetailsPanel } from './MessageDetailsPanel'
interface QueueDataGridProps {
isLoading: boolean
messages: PostgresQueueMessage[]
showMessageModal: () => void
fetchNextPage: () => void
}
function isAtBottom({ currentTarget }: React.UIEvent<HTMLDivElement>): boolean {
function isAtBottom({ currentTarget }: UIEvent<HTMLDivElement>): boolean {
return currentTarget.scrollTop + 10 >= currentTarget.scrollHeight - currentTarget.clientHeight
}
@@ -123,6 +124,7 @@ const columns = messagesCols.map((col) => {
export const QueueMessagesDataGrid = ({
isLoading,
messages,
showMessageModal,
fetchNextPage,
}: QueueDataGridProps) => {
const gridRef = useRef<DataGridHandle>(null)
@@ -188,6 +190,9 @@ export const QueueMessagesDataGrid = ({
<p className="text-foreground-light">
The selected queue doesn't have any messages.
</p>
<Button className="mt-2" onClick={() => showMessageModal()}>
Add message
</Button>
</div>
</div>
),
@@ -0,0 +1,138 @@
import { zodResolver } from '@hookform/resolvers/zod'
import { SubmitHandler, useForm } from 'react-hook-form'
import z from 'zod'
import { useParams } from 'common'
import { useProjectContext } from 'components/layouts/ProjectLayout/ProjectContext'
import CodeEditor from 'components/ui/CodeEditor/CodeEditor'
import { useDatabaseQueueMessageSendMutation } from 'data/database-queues/database-queue-messages-send-mutation'
import { useEffect } from 'react'
import { toast } from 'sonner'
import { Form_Shadcn_, FormControl_Shadcn_, FormField_Shadcn_, Input, Modal } from 'ui'
import { FormItemLayout } from 'ui-patterns/form/FormItemLayout/FormItemLayout'
interface SendMessageModalProps {
visible: boolean
onClose: () => void
}
const FormSchema = z.object({
delay: z.coerce.number().int().gte(0).default(5),
payload: z.string().refine(
(val) => {
try {
JSON.parse(val)
} catch {
return false
}
},
{
message: 'The payload should be a JSON object',
}
),
})
export type SendMessageForm = z.infer<typeof FormSchema>
const FORM_ID = 'QUEUES_SEND_MESSAGE_FORM'
export const SendMessageModal = ({ visible, onClose }: SendMessageModalProps) => {
const { name: queueName } = useParams()
const { project } = useProjectContext()
const form = useForm<SendMessageForm>({
resolver: zodResolver(FormSchema),
defaultValues: {
delay: 1,
payload: '{}',
},
})
const { isLoading, mutate } = useDatabaseQueueMessageSendMutation({
onSuccess: () => {
toast.success(`Successfully added a message to the queue.`)
onClose()
},
})
const onSubmit: SubmitHandler<SendMessageForm> = (values) => {
mutate({
projectRef: project?.ref!,
connectionString: project?.connectionString,
queueName: queueName!,
payload: values.payload,
delay: values.delay,
})
}
useEffect(() => {
if (visible) {
form.reset({ delay: 1, payload: '{}' })
}
}, [visible])
return (
<Modal
size="medium"
alignFooter="right"
header="Add a message to the queue"
visible={visible}
loading={isLoading}
onCancel={onClose}
confirmText="Add"
onConfirm={() => {
const values = form.getValues()
onSubmit(values)
}}
>
<Modal.Content className="flex flex-col gap-y-4">
<Form_Shadcn_ {...form}>
<form
id={FORM_ID}
className="flex-grow overflow-auto gap-2 flex flex-col"
onSubmit={form.handleSubmit(onSubmit)}
>
<FormField_Shadcn_
control={form.control}
name="delay"
render={({ field: { ref, ...rest } }) => (
<FormItemLayout
label="Delay"
layout="vertical"
className="gap-1"
description="Time in seconds before the message becomes available for reading."
>
<FormControl_Shadcn_>
<Input
{...rest}
type="number"
placeholder="1"
actions={<p className="text-foreground-light pr-2">sec</p>}
/>
</FormControl_Shadcn_>
</FormItemLayout>
)}
/>
<FormField_Shadcn_
control={form.control}
name="payload"
render={({ field }) => (
<FormItemLayout label="Message payload" layout="vertical" className="gap-1">
<FormControl_Shadcn_>
<CodeEditor
id="message-payload"
language="json"
className="!mb-0 h-32 overflow-hidden rounded border"
onInputChange={(e: string | undefined) => field.onChange(e)}
options={{ wordWrap: 'off', contextmenu: false }}
value={field.value}
/>
</FormControl_Shadcn_>
</FormItemLayout>
)}
/>
</form>
</Form_Shadcn_>
</Modal.Content>
</Modal>
)
}
@@ -0,0 +1,68 @@
import { useMutation, UseMutationOptions, useQueryClient } from '@tanstack/react-query'
import { toast } from 'sonner'
import { executeSql } from 'data/sql/execute-sql-query'
import type { ResponseError } from 'types'
import { databaseQueuesKeys } from './keys'
export type DatabaseQueueMessageDeleteVariables = {
projectRef: string
connectionString?: string
queueName: string
messageId: number
}
export async function deleteDatabaseQueueMessage({
projectRef,
connectionString,
queueName,
messageId,
}: DatabaseQueueMessageDeleteVariables) {
const { result } = await executeSql({
projectRef,
connectionString,
sql: `SELECT * FROM pgmq.delete('${queueName}', ${messageId})`,
queryKey: databaseQueuesKeys.create(),
})
return result
}
type DatabaseQueueMessageDeleteData = Awaited<ReturnType<typeof deleteDatabaseQueueMessage>>
export const useDatabaseQueueMessageDeleteMutation = ({
onSuccess,
onError,
...options
}: Omit<
UseMutationOptions<
DatabaseQueueMessageDeleteData,
ResponseError,
DatabaseQueueMessageDeleteVariables
>,
'mutationFn'
> = {}) => {
const queryClient = useQueryClient()
return useMutation<
DatabaseQueueMessageDeleteData,
ResponseError,
DatabaseQueueMessageDeleteVariables
>((vars) => deleteDatabaseQueueMessage(vars), {
async onSuccess(data, variables, context) {
const { projectRef, queueName } = variables
await queryClient.invalidateQueries(
databaseQueuesKeys.getMessagesInfinite(projectRef, queueName)
)
await onSuccess?.(data, variables, context)
},
async onError(data, variables, context) {
if (onError === undefined) {
toast.error(`Failed to delete database queue message: ${data.message}`)
} else {
onError(data, variables, context)
}
},
...options,
})
}
@@ -0,0 +1,70 @@
import { useMutation, UseMutationOptions, useQueryClient } from '@tanstack/react-query'
import { toast } from 'sonner'
import { executeSql } from 'data/sql/execute-sql-query'
import type { ResponseError } from 'types'
import { databaseQueuesKeys } from './keys'
export type DatabaseQueueMessageReadVariables = {
projectRef: string
connectionString?: string
queryName: string
duration: number
messageId: number
}
export async function readDatabaseQueueMessage({
projectRef,
connectionString,
queryName,
messageId,
duration,
}: DatabaseQueueMessageReadVariables) {
const { result } = await executeSql({
projectRef,
connectionString,
sql: `select * from pgmq.set_vt('${queryName}', ${messageId}, ${duration})`,
queryKey: databaseQueuesKeys.create(),
})
return result
}
type DatabaseQueueMessageReadData = Awaited<ReturnType<typeof readDatabaseQueueMessage>>
export const useDatabaseQueueMessageReadMutation = ({
onSuccess,
onError,
...options
}: Omit<
UseMutationOptions<
DatabaseQueueMessageReadData,
ResponseError,
DatabaseQueueMessageReadVariables
>,
'mutationFn'
> = {}) => {
const queryClient = useQueryClient()
return useMutation<
DatabaseQueueMessageReadData,
ResponseError,
DatabaseQueueMessageReadVariables
>((vars) => readDatabaseQueueMessage(vars), {
async onSuccess(data, variables, context) {
const { projectRef, queryName } = variables
await queryClient.invalidateQueries(
databaseQueuesKeys.getMessagesInfinite(projectRef, queryName)
)
await onSuccess?.(data, variables, context)
},
async onError(data, variables, context) {
if (onError === undefined) {
toast.error(`Failed to postpone database queue message: ${data.message}`)
} else {
onError(data, variables, context)
}
},
...options,
})
}
@@ -0,0 +1,70 @@
import { useMutation, UseMutationOptions, useQueryClient } from '@tanstack/react-query'
import { toast } from 'sonner'
import { executeSql } from 'data/sql/execute-sql-query'
import type { ResponseError } from 'types'
import { databaseQueuesKeys } from './keys'
export type DatabaseQueueMessageSendVariables = {
projectRef: string
connectionString?: string
queueName: string
payload: string
delay: number
}
export async function sendDatabaseQueueMessage({
projectRef,
connectionString,
queueName,
payload,
delay,
}: DatabaseQueueMessageSendVariables) {
const { result } = await executeSql({
projectRef,
connectionString,
sql: `select * from pgmq.send( '${queueName}', '${payload}', ${delay})`,
queryKey: databaseQueuesKeys.create(),
})
return result
}
type DatabaseQueueMessageSendData = Awaited<ReturnType<typeof sendDatabaseQueueMessage>>
export const useDatabaseQueueMessageSendMutation = ({
onSuccess,
onError,
...options
}: Omit<
UseMutationOptions<
DatabaseQueueMessageSendData,
ResponseError,
DatabaseQueueMessageSendVariables
>,
'mutationFn'
> = {}) => {
const queryClient = useQueryClient()
return useMutation<
DatabaseQueueMessageSendData,
ResponseError,
DatabaseQueueMessageSendVariables
>((vars) => sendDatabaseQueueMessage(vars), {
async onSuccess(data, variables, context) {
const { projectRef, queueName } = variables
await queryClient.invalidateQueries(
databaseQueuesKeys.getMessagesInfinite(projectRef, queueName)
)
await onSuccess?.(data, variables, context)
},
async onError(data, variables, context) {
if (onError === undefined) {
toast.error(`Failed to send database queue message: ${data.message}`)
} else {
onError(data, variables, context)
}
},
...options,
})
}
@@ -20,7 +20,7 @@ export async function deleteDatabaseQueue({
projectRef,
connectionString,
sql: `select * from pgmq.drop_queue('${queueName}');`,
queryKey: databaseQueuesKeys.delete(),
queryKey: databaseQueuesKeys.delete(queueName),
})
return result
@@ -0,0 +1,90 @@
import { UseQueryOptions, useQuery } from '@tanstack/react-query'
import { handleError } from 'data/fetchers'
import { executeSql } from 'data/sql/execute-sql-query'
import { ResponseError } from 'types'
import { databaseQueuesKeys } from './keys'
export type DatabaseQueuesMetricsVariables = {
projectRef?: string
connectionString?: string
queueName: string
}
export type PostgresQueueMetric = {
queue_name: string
queue_length: number
method: 'estimated' | 'precise'
}
const preciseMetricsSqlQuery = (queueName: string) => `
set local statement_timeout = '1s';
SELECT
COUNT(*) AS row_count
FROM
"pgmq"."q_${queueName}";
`
const estimateMetricsSqlQuery = (queueName: string) => `
select
reltuples::bigint as estimated_rows
from
pg_class
where
relname = 'q_${queueName}'
and relnamespace = 'pgmq'::regnamespace;
`
export async function getDatabaseQueuesMetrics({
projectRef,
connectionString,
queueName,
}: DatabaseQueuesMetricsVariables) {
if (!projectRef) throw new Error('Project ref is required')
try {
const { result } = await executeSql({
projectRef,
connectionString,
sql: preciseMetricsSqlQuery(queueName),
})
return {
queue_name: queueName,
queue_length: result[0].row_count,
method: 'precise',
} as PostgresQueueMetric
} catch (error: any) {
// if the error is caused because the count timeouted, try to fetch an approximate count
if (error?.message === 'canceling statement due to statement timeout') {
const { result } = await executeSql({
projectRef,
connectionString,
sql: estimateMetricsSqlQuery(queueName),
})
return {
queue_name: queueName,
queue_length: result[0].estimated_rows,
method: 'estimated',
} as PostgresQueueMetric
}
return handleError(error)
}
}
export type DatabaseQueuesMetricsData = PostgresQueueMetric
export type DatabaseQueuesMetricsError = ResponseError
export const useQueuesMetricsQuery = <TData = DatabaseQueuesMetricsData>(
{ projectRef, connectionString, queueName }: DatabaseQueuesMetricsVariables,
{
enabled = true,
...options
}: UseQueryOptions<DatabaseQueuesMetricsData, DatabaseQueuesMetricsError, TData> = {}
) =>
useQuery<DatabaseQueuesMetricsData, DatabaseQueuesMetricsError, TData>(
databaseQueuesKeys.metrics(projectRef, queueName),
() => getDatabaseQueuesMetrics({ projectRef, connectionString, queueName }),
{
enabled: enabled && typeof projectRef !== 'undefined',
...options,
}
)
@@ -0,0 +1,61 @@
import { useMutation, UseMutationOptions, useQueryClient } from '@tanstack/react-query'
import { toast } from 'sonner'
import { executeSql } from 'data/sql/execute-sql-query'
import type { ResponseError } from 'types'
import { databaseQueuesKeys } from './keys'
export type DatabaseQueuePurgeVariables = {
projectRef: string
connectionString?: string
queueName: string
}
export async function purgeDatabaseQueue({
projectRef,
connectionString,
queueName,
}: DatabaseQueuePurgeVariables) {
const { result } = await executeSql({
projectRef,
connectionString,
sql: `select * from pgmq.purge_queue('${queueName}');`,
queryKey: databaseQueuesKeys.purge(queueName),
})
return result
}
type DatabaseQueuePurgeData = Awaited<ReturnType<typeof purgeDatabaseQueue>>
export const useDatabaseQueuePurgeMutation = ({
onSuccess,
onError,
...options
}: Omit<
UseMutationOptions<DatabaseQueuePurgeData, ResponseError, DatabaseQueuePurgeVariables>,
'mutationFn'
> = {}) => {
const queryClient = useQueryClient()
return useMutation<DatabaseQueuePurgeData, ResponseError, DatabaseQueuePurgeVariables>(
(vars) => purgeDatabaseQueue(vars),
{
async onSuccess(data, variables, context) {
const { projectRef, queueName } = variables
await queryClient.invalidateQueries(
databaseQueuesKeys.getMessagesInfinite(projectRef, queueName)
)
await onSuccess?.(data, variables, context)
},
async onError(data, variables, context) {
if (onError === undefined) {
toast.error(`Failed to purge database queue: ${data.message}`)
} else {
onError(data, variables, context)
}
},
...options,
}
)
}
+5 -1
View File
@@ -1,7 +1,11 @@
export const databaseQueuesKeys = {
create: () => ['queues', 'create'] as const,
delete: () => ['queues', 'delete'] as const,
delete: (name: string) => ['queues', name, 'delete'] as const,
purge: (name: string) => ['queues', name, 'purge'] as const,
getMessagesInfinite: (projectRef: string | undefined, queueName: string, options?: object) =>
['projects', projectRef, 'queues', queueName, options].filter(Boolean),
list: (projectRef: string | undefined) => ['projects', projectRef, 'queues'] as const,
// invalidating queues.list will also invalidate queues.metrics
metrics: (projectRef: string | undefined, queueName: string) =>
['projects', projectRef, 'queues', 'metrics', queueName] as const,
}
@@ -1,16 +1,29 @@
import { Paintbrush, Trash2 } from 'lucide-react'
import Link from 'next/link'
import { useMemo, useState } from 'react'
import { useParams } from 'common'
import DeleteQueue from 'components/interfaces/Integrations/Queues/SingleQueue/DeleteQueue'
import PurgeQueue from 'components/interfaces/Integrations/Queues/SingleQueue/PurgeQueue'
import { QUEUE_MESSAGE_TYPE } from 'components/interfaces/Integrations/Queues/SingleQueue/Queue.utils'
import { QueueMessagesDataGrid } from 'components/interfaces/Integrations/Queues/SingleQueue/QueueDataGrid'
import { QueueFilters } from 'components/interfaces/Integrations/Queues/SingleQueue/QueueFilters'
import { SendMessageModal } from 'components/interfaces/Integrations/Queues/SingleQueue/SendMessageModal'
import ProjectIntegrationsLayout from 'components/layouts/ProjectIntegrationsLayout/ProjectIntegrationsLayout'
import { useProjectContext } from 'components/layouts/ProjectLayout/ProjectContext'
import { FormHeader } from 'components/ui/Forms/FormHeader'
import { useQueueMessagesInfiniteQuery } from 'data/database-queues/database-queue-messages-infinite-query'
import type { NextPageWithLayout } from 'types'
import { Button, LoadingLine } from 'ui'
import {
Breadcrumb_Shadcn_,
BreadcrumbItem_Shadcn_,
BreadcrumbLink_Shadcn_,
BreadcrumbList_Shadcn_,
BreadcrumbPage_Shadcn_,
BreadcrumbSeparator_Shadcn_,
Button,
LoadingLine,
Separator,
} from 'ui'
const QueueMessagesPage: NextPageWithLayout = () => {
// TODO: Change this to the correct permissions
@@ -23,6 +36,8 @@ const QueueMessagesPage: NextPageWithLayout = () => {
const { name: queueName } = useParams()
const { project } = useProjectContext()
const [sendMessageModalShown, setSendMessageModalShown] = useState(false)
const [purgeQueueModalShown, setPurgeQueueModalShown] = useState(false)
const [deleteQueueModalShown, setDeleteQueueModalShown] = useState(false)
const [selectedTypes, setSelectedTypes] = useState<QUEUE_MESSAGE_TYPE[]>([])
@@ -38,34 +53,72 @@ const QueueMessagesPage: NextPageWithLayout = () => {
)
const messages = useMemo(() => data?.pages.flatMap((p) => p), [data?.pages])
if (isLoading && isError) {
if (isError) {
return null
}
return (
<div className="h-full flex flex-col">
<FormHeader
className="py-4 px-6 !mb-0"
title={`Queue ${queueName}`}
actions={
<Button type="danger" onClick={() => setDeleteQueueModalShown(true)}>
Delete the queue
<div className="flex items-center justify-between gap-x-4 py-4 px-6 mb-0">
<div className="space-y-1 flex-shrink">
<Breadcrumb_Shadcn_>
<BreadcrumbList_Shadcn_>
<BreadcrumbItem_Shadcn_>
<BreadcrumbLink_Shadcn_ asChild className="text-xl">
<Link href={`/project/${project?.ref}/integrations/queues`}>Queues</Link>
</BreadcrumbLink_Shadcn_>
</BreadcrumbItem_Shadcn_>
<BreadcrumbSeparator_Shadcn_ />
<BreadcrumbItem_Shadcn_>
<BreadcrumbPage_Shadcn_ className="text-xl">{queueName}</BreadcrumbPage_Shadcn_>
</BreadcrumbItem_Shadcn_>
</BreadcrumbList_Shadcn_>
</Breadcrumb_Shadcn_>
</div>
<div className="flex gap-x-2">
<Button
type="text"
onClick={() => setPurgeQueueModalShown(true)}
icon={<Paintbrush />}
title="Purge messages"
/>
<Button
type="text"
onClick={() => setDeleteQueueModalShown(true)}
icon={<Trash2 />}
title="Delete queue"
/>
<Separator orientation="vertical" className="h-[26px]" />
<Button type="primary" onClick={() => setSendMessageModalShown(true)}>
Add message
</Button>
}
/>
{/* <DocsButton href={docsUrl} />} */}
</div>
</div>
<QueueFilters selectedTypes={selectedTypes} setSelectedTypes={setSelectedTypes} />
<LoadingLine loading={isFetching} />
<QueueMessagesDataGrid
messages={messages || []}
isLoading={isLoading}
showMessageModal={() => setSendMessageModalShown(true)}
fetchNextPage={fetchNextPage}
/>
<SendMessageModal
visible={sendMessageModalShown}
onClose={() => setSendMessageModalShown(false)}
/>
<DeleteQueue
queueName={queueName!}
visible={deleteQueueModalShown}
onClose={() => setDeleteQueueModalShown(false)}
/>
<PurgeQueue
queueName={queueName!}
visible={purgeQueueModalShown}
onClose={() => setPurgeQueueModalShown(false)}
/>
</div>
)
}
@@ -1,6 +1,6 @@
import * as React from 'react'
import { ChevronRightIcon } from 'lucide-react'
import { Slot } from '@radix-ui/react-slot'
import { ChevronRightIcon } from 'lucide-react'
import * as React from 'react'
import { cn } from '../../../lib/utils'
const Breadcrumb = React.forwardRef<
@@ -64,7 +64,7 @@ const BreadcrumbPage = React.forwardRef<HTMLSpanElement, React.ComponentPropsWit
role="link"
aria-disabled="true"
aria-current="page"
className={cn('no-underline', className)}
className={cn('no-underline text-foreground', className)}
{...props}
/>
)
@@ -108,10 +108,10 @@ BreadcrumbEllipsis.displayName = 'BreadcrumbEllipsis'
export {
Breadcrumb,
BreadcrumbList,
BreadcrumbEllipsis,
BreadcrumbItem,
BreadcrumbLink,
BreadcrumbList,
BreadcrumbPage,
BreadcrumbSeparator,
BreadcrumbEllipsis,
}