From 1157a20fbd6fa1ded377bed247d3035c58c43552 Mon Sep 17 00:00:00 2001
From: Ivan Vasilov
Date: Thu, 14 Nov 2024 13:02:29 +0100
Subject: [PATCH] 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.
---
.../Integrations/Queues/Queues.constants.tsx | 20 +--
.../Queues/QueuesDisabledState.tsx | 13 +-
.../Integrations/Queues/QueuesListing.tsx | 9 +-
.../Integrations/Queues/QueuesRows.tsx | 93 ++++++++----
.../SingleQueue/MessageDetailsPanel.tsx | 134 +++++++++++++----
.../Queues/SingleQueue/PurgeQueue.tsx | 65 +++++++++
.../Queues/SingleQueue/QueueDataGrid.tsx | 9 +-
.../Queues/SingleQueue/SendMessageModal.tsx | 138 ++++++++++++++++++
...database-queue-messages-delete-mutation.ts | 68 +++++++++
.../database-queue-messages-read-mutation.ts | 70 +++++++++
.../database-queue-messages-send-mutation.ts | 70 +++++++++
.../database-queues-delete-mutation.ts | 2 +-
.../database-queues-metrics-query.ts | 90 ++++++++++++
.../database-queues-purge-mutation.ts | 61 ++++++++
apps/studio/data/database-queues/keys.ts | 6 +-
.../[ref]/integrations/queues/[name].tsx | 75 ++++++++--
.../src/components/shadcn/ui/breadcrumb.tsx | 10 +-
17 files changed, 828 insertions(+), 105 deletions(-)
create mode 100644 apps/studio/components/interfaces/Integrations/Queues/SingleQueue/PurgeQueue.tsx
create mode 100644 apps/studio/components/interfaces/Integrations/Queues/SingleQueue/SendMessageModal.tsx
create mode 100644 apps/studio/data/database-queues/database-queue-messages-delete-mutation.ts
create mode 100644 apps/studio/data/database-queues/database-queue-messages-read-mutation.ts
create mode 100644 apps/studio/data/database-queues/database-queue-messages-send-mutation.ts
create mode 100644 apps/studio/data/database-queues/database-queues-metrics-query.ts
create mode 100644 apps/studio/data/database-queues/database-queues-purge-mutation.ts
diff --git a/apps/studio/components/interfaces/Integrations/Queues/Queues.constants.tsx b/apps/studio/components/interfaces/Integrations/Queues/Queues.constants.tsx
index 6ba001477a5..a80a56be71e 100644
--- a/apps/studio/components/interfaces/Integrations/Queues/Queues.constants.tsx
+++ b/apps/studio/components/interfaces/Integrations/Queues/Queues.constants.tsx
@@ -2,10 +2,17 @@ import { Rows2, Rows3, Rows4 } from 'lucide-react'
export const QUEUE_TYPES = [
{
- value: 'normal',
+ value: 'basic',
icon: ,
- label: 'Normal queue',
- description: 'Create a normal queue.',
+ label: 'Basic queue',
+ description: 'Create a basic queue.',
+ },
+ {
+ value: 'unlogged',
+ icon: ,
+ 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: ,
- label: 'Unlogged queue',
- description: 'Creates an unlogged queue which loses all data on database restart.',
- },
] as const
diff --git a/apps/studio/components/interfaces/Integrations/Queues/QueuesDisabledState.tsx b/apps/studio/components/interfaces/Integrations/Queues/QueuesDisabledState.tsx
index ad4331f6706..5de88c73820 100644
--- a/apps/studio/components/interfaces/Integrations/Queues/QueuesDisabledState.tsx
+++ b/apps/studio/components/interfaces/Integrations/Queues/QueuesDisabledState.tsx
@@ -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 = () => {
)}
- }>
-
- Documentation
-
-
+
diff --git a/apps/studio/components/interfaces/Integrations/Queues/QueuesListing.tsx b/apps/studio/components/interfaces/Integrations/Queues/QueuesListing.tsx
index f34ad985e9d..eba0a0779b3 100644
--- a/apps/studio/components/interfaces/Integrations/Queues/QueuesListing.tsx
+++ b/apps/studio/components/interfaces/Integrations/Queues/QueuesListing.tsx
@@ -69,10 +69,10 @@ export const QueuesListing = () => {
{queues.length === 0 ? (
-
No queues created yet
+
No queues created yet
) : (
@@ -98,9 +98,12 @@ export const QueuesListing = () => {
Type
-
+
Created at
+
+ Size
+
>
}
diff --git a/apps/studio/components/interfaces/Integrations/Queues/QueuesRows.tsx b/apps/studio/components/interfaces/Integrations/Queues/QueuesRows.tsx
index ec68ed69772..629b0347d9f 100644
--- a/apps/studio/components/interfaces/Integrations/Queues/QueuesRows.tsx
+++ b/apps/studio/components/interfaces/Integrations/Queues/QueuesRows.tsx
@@ -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 (
+ {
+ router.push(`/project/${selectedProject?.ref}/integrations/queues/${queue.queue_name}`)
+ }}
+ className="hover:"
+ >
+
+ {queue.queue_name}
+
+
+
+ {type}
+
+
+
+ {dayjs(queue.created_at).format(DATETIME_FORMAT)}
+
+
+
+ {isLoading ? (
+
+ ) : (
+
+ {metrics?.queue_length} {metrics?.method === 'estimated' ? '(Approximate)' : null}
+
+ )}
+
+
+
+
+
+
+
+
+ )
+}
+
+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 (
-
+
No results found
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 (
-
{
- router.push(`/project/${selectedProject?.ref}/integrations/queues/${q.queue_name}`)
- }}
- className="hover:"
- >
-
- {q.queue_name}
-
-
-
- {type}
-
-
-
- {dayjs(q.created_at).format(DATETIME_FORMAT)}
-
-
-
-
-
- )
- })}
+ {queues.map((q) => (
+
+ ))}
>
)
}
diff --git a/apps/studio/components/interfaces/Integrations/Queues/SingleQueue/MessageDetailsPanel.tsx b/apps/studio/components/interfaces/Integrations/Queues/SingleQueue/MessageDetailsPanel.tsx
index e7fed0b97e2..9de86466dec 100644
--- a/apps/studio/components/interfaces/Integrations/Queues/SingleQueue/MessageDetailsPanel.tsx
+++ b/apps/studio/components/interfaces/Integrations/Queues/SingleQueue/MessageDetailsPanel.tsx
@@ -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 = ({
-
Payload
+ Payload
-
+
{!selectedMessage.archived_at ? (
- ,
- 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
- }
- />
+ <>
+ ,
+ 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
+ }
+ />
+ ,
+ 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
+ }
+ />
+ ,
+ 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}
diff --git a/apps/studio/components/interfaces/Integrations/Queues/SingleQueue/PurgeQueue.tsx b/apps/studio/components/interfaces/Integrations/Queues/SingleQueue/PurgeQueue.tsx
new file mode 100644
index 00000000000..5743748e97a
--- /dev/null
+++ b/apps/studio/components/interfaces/Integrations/Queues/SingleQueue/PurgeQueue.tsx
@@ -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 (
+
onClose()}
+ onConfirm={handlePurge}
+ title="Purge this queue"
+ loading={isLoading}
+ confirmLabel={`Purge queue ${queueName}`}
+ confirmPlaceholder="Type in name of queue"
+ confirmString={queueName ?? 'Unknown'}
+ text={
+ <>
+ This will purge the queue{' '}
+ {queueName}
+ >
+ }
+ alert={{
+ title:
+ "This action will delete all messages from the queue. They can't be recovered afterwards.",
+ }}
+ />
+ )
+}
+
+export default PurgeQueue
diff --git a/apps/studio/components/interfaces/Integrations/Queues/SingleQueue/QueueDataGrid.tsx b/apps/studio/components/interfaces/Integrations/Queues/SingleQueue/QueueDataGrid.tsx
index e5e30783d41..c77211a1c7f 100644
--- a/apps/studio/components/interfaces/Integrations/Queues/SingleQueue/QueueDataGrid.tsx
+++ b/apps/studio/components/interfaces/Integrations/Queues/SingleQueue/QueueDataGrid.tsx
@@ -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): boolean {
+function isAtBottom({ currentTarget }: UIEvent): 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(null)
@@ -188,6 +190,9 @@ export const QueueMessagesDataGrid = ({
The selected queue doesn't have any messages.
+
),
diff --git a/apps/studio/components/interfaces/Integrations/Queues/SingleQueue/SendMessageModal.tsx b/apps/studio/components/interfaces/Integrations/Queues/SingleQueue/SendMessageModal.tsx
new file mode 100644
index 00000000000..26fb2c27c1b
--- /dev/null
+++ b/apps/studio/components/interfaces/Integrations/Queues/SingleQueue/SendMessageModal.tsx
@@ -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
+
+const FORM_ID = 'QUEUES_SEND_MESSAGE_FORM'
+
+export const SendMessageModal = ({ visible, onClose }: SendMessageModalProps) => {
+ const { name: queueName } = useParams()
+ const { project } = useProjectContext()
+ const form = useForm({
+ 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 = (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 (
+ {
+ const values = form.getValues()
+ onSubmit(values)
+ }}
+ >
+
+
+ }
+ />
+
+
+ )}
+ />
+ (
+
+
+ field.onChange(e)}
+ options={{ wordWrap: 'off', contextmenu: false }}
+ value={field.value}
+ />
+
+
+ )}
+ />
+
+
+
+
+ )
+}
diff --git a/apps/studio/data/database-queues/database-queue-messages-delete-mutation.ts b/apps/studio/data/database-queues/database-queue-messages-delete-mutation.ts
new file mode 100644
index 00000000000..65c0c124bf2
--- /dev/null
+++ b/apps/studio/data/database-queues/database-queue-messages-delete-mutation.ts
@@ -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>
+
+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,
+ })
+}
diff --git a/apps/studio/data/database-queues/database-queue-messages-read-mutation.ts b/apps/studio/data/database-queues/database-queue-messages-read-mutation.ts
new file mode 100644
index 00000000000..601675001cb
--- /dev/null
+++ b/apps/studio/data/database-queues/database-queue-messages-read-mutation.ts
@@ -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>
+
+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,
+ })
+}
diff --git a/apps/studio/data/database-queues/database-queue-messages-send-mutation.ts b/apps/studio/data/database-queues/database-queue-messages-send-mutation.ts
new file mode 100644
index 00000000000..e06600ff413
--- /dev/null
+++ b/apps/studio/data/database-queues/database-queue-messages-send-mutation.ts
@@ -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>
+
+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,
+ })
+}
diff --git a/apps/studio/data/database-queues/database-queues-delete-mutation.ts b/apps/studio/data/database-queues/database-queues-delete-mutation.ts
index 7cdd8b3df41..83c2c828359 100644
--- a/apps/studio/data/database-queues/database-queues-delete-mutation.ts
+++ b/apps/studio/data/database-queues/database-queues-delete-mutation.ts
@@ -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
diff --git a/apps/studio/data/database-queues/database-queues-metrics-query.ts b/apps/studio/data/database-queues/database-queues-metrics-query.ts
new file mode 100644
index 00000000000..a7617213b7c
--- /dev/null
+++ b/apps/studio/data/database-queues/database-queues-metrics-query.ts
@@ -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 = (
+ { projectRef, connectionString, queueName }: DatabaseQueuesMetricsVariables,
+ {
+ enabled = true,
+ ...options
+ }: UseQueryOptions = {}
+) =>
+ useQuery(
+ databaseQueuesKeys.metrics(projectRef, queueName),
+ () => getDatabaseQueuesMetrics({ projectRef, connectionString, queueName }),
+ {
+ enabled: enabled && typeof projectRef !== 'undefined',
+ ...options,
+ }
+ )
diff --git a/apps/studio/data/database-queues/database-queues-purge-mutation.ts b/apps/studio/data/database-queues/database-queues-purge-mutation.ts
new file mode 100644
index 00000000000..7a4419b64ae
--- /dev/null
+++ b/apps/studio/data/database-queues/database-queues-purge-mutation.ts
@@ -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>
+
+export const useDatabaseQueuePurgeMutation = ({
+ onSuccess,
+ onError,
+ ...options
+}: Omit<
+ UseMutationOptions,
+ 'mutationFn'
+> = {}) => {
+ const queryClient = useQueryClient()
+
+ return useMutation(
+ (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,
+ }
+ )
+}
diff --git a/apps/studio/data/database-queues/keys.ts b/apps/studio/data/database-queues/keys.ts
index c0c484c93aa..228c5caa694 100644
--- a/apps/studio/data/database-queues/keys.ts
+++ b/apps/studio/data/database-queues/keys.ts
@@ -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,
}
diff --git a/apps/studio/pages/project/[ref]/integrations/queues/[name].tsx b/apps/studio/pages/project/[ref]/integrations/queues/[name].tsx
index 0acaa3f36ca..1ed8a0dfd97 100644
--- a/apps/studio/pages/project/[ref]/integrations/queues/[name].tsx
+++ b/apps/studio/pages/project/[ref]/integrations/queues/[name].tsx
@@ -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([])
@@ -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 (
-
setDeleteQueueModalShown(true)}>
- Delete the queue
+
+
+
+
+
+
+ Queues
+
+
+
+
+ {queueName}
+
+
+
+
+
+
+
setSendMessageModalShown(true)}
fetchNextPage={fetchNextPage}
/>
+ setSendMessageModalShown(false)}
+ />
setDeleteQueueModalShown(false)}
/>
+ setPurgeQueueModalShown(false)}
+ />
)
}
diff --git a/packages/ui/src/components/shadcn/ui/breadcrumb.tsx b/packages/ui/src/components/shadcn/ui/breadcrumb.tsx
index 7bddf8fa745..e76be186190 100644
--- a/packages/ui/src/components/shadcn/ui/breadcrumb.tsx
+++ b/packages/ui/src/components/shadcn/ui/breadcrumb.tsx
@@ -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
)
@@ -108,10 +108,10 @@ BreadcrumbEllipsis.displayName = 'BreadcrumbEllipsis'
export {
Breadcrumb,
- BreadcrumbList,
+ BreadcrumbEllipsis,
BreadcrumbItem,
BreadcrumbLink,
+ BreadcrumbList,
BreadcrumbPage,
BreadcrumbSeparator,
- BreadcrumbEllipsis,
}