From 7800eafb83c1b2e5e558c4f8f0a8cbce28e83f8c Mon Sep 17 00:00:00 2001 From: aaron-rab Date: Tue, 4 Aug 2026 18:25:29 -0400 Subject: [PATCH 1/3] added queue explorer page and related resolver/application service; made necessary additions to cellix service-queue-storage to support queue explorer functionality --- apps/ui-community/package.json | 5 +- .../src/hooks/use-staff-permissions.ts | 8 +- codegen.yml | 14 + knip.json | 1 + .../cellix/service-queue-storage/README.md | 44 +++ .../cellix/service-queue-storage/src/index.ts | 1 + .../service-queue-storage/src/interfaces.ts | 7 + .../internal-queue-storage-service.test.ts | 19 ++ .../src/internal-queue-storage-service.ts | 8 + .../src/queue-consumer.test.ts | 21 ++ .../src/queue-consumer.ts | 14 +- .../src/queue-producer.test.ts | 108 ++++++- .../src/queue-producer.ts | 18 +- .../src/register-queues.test.ts | 16 + .../src/register-queues.ts | 65 +++- .../src/contexts/tech-admin/index.ts | 13 + .../queue/get-queue-message-count.test.ts | 17 + .../queue/get-queue-message-count.ts | 26 ++ .../src/contexts/tech-admin/queue/index.ts | 28 ++ .../tech-admin/queue/list-queues.test.ts | 11 + .../contexts/tech-admin/queue/list-queues.ts | 8 + .../tech-admin/queue/peek-messages.ts | 18 ++ .../contexts/tech-admin/queue/queue-list.ts | 41 +++ .../tech-admin/queue/queue-operations.test.ts | 50 +++ .../tech-admin/queue/queue-operations.ts | 63 ++++ .../queue/queue-permissions.test.ts | 38 +++ .../tech-admin/queue/queue-permissions.ts | 27 ++ .../tech-admin/queue/send-message.test.ts | 21 ++ .../contexts/tech-admin/queue/send-message.ts | 20 ++ .../user/staff-role/apply-permissions.test.ts | 3 + .../user/staff-role/apply-permissions.ts | 4 + .../ocom/application-services/src/index.ts | 3 + .../src/models/role/staff-role.model.ts | 2 + .../staff-role-tech-admin-permissions.feature | 5 + .../staff-role/staff-role-defaults.test.ts | 2 + .../user/staff-role/staff-role-permissions.ts | 1 + .../staff-role-tech-admin-permissions.test.ts | 14 + .../staff-role-tech-admin-permissions.ts | 9 + .../contexts/user/staff-role/staff-role.ts | 2 + .../src/schema/types/staff-role.graphql | 2 + .../src/schema/types/tech-admin.graphql | 40 +++ .../schema/types/tech-admin.resolvers.test.ts | 41 +++ .../src/schema/types/tech-admin.resolvers.ts | 48 +++ .../staff-role/staff-role.domain-adapter.ts | 8 + packages/ocom/service-queue-storage/README.md | 2 + .../src/queue-storage.contract.ts | 2 +- .../ui-staff-route-tech-admin/package.json | 3 + .../queue-explorer.container.graphql | 27 ++ .../components/queue-explorer.container.tsx | 83 +++++ .../src/components/queue-explorer.tsx | 305 ++++++++++++++++++ .../ui-staff-route-tech-admin/src/index.tsx | 19 +- .../src/pages/queue-explorer.tsx | 19 ++ .../src/pages/tech-admin.tsx | 44 +++ .../staff-role-create.container.tsx | 1 + .../src/components/staff-role-create.tsx | 3 + .../staff-role-edit.container.graphql | 1 + .../components/staff-role-edit.container.tsx | 2 + .../ui-staff-shared/src/staff-route-shell.tsx | 2 + pnpm-lock.yaml | 270 +++++++--------- pnpm-workspace.yaml | 6 +- 60 files changed, 1528 insertions(+), 175 deletions(-) create mode 100644 packages/ocom/application-services/src/contexts/tech-admin/index.ts create mode 100644 packages/ocom/application-services/src/contexts/tech-admin/queue/get-queue-message-count.test.ts create mode 100644 packages/ocom/application-services/src/contexts/tech-admin/queue/get-queue-message-count.ts create mode 100644 packages/ocom/application-services/src/contexts/tech-admin/queue/index.ts create mode 100644 packages/ocom/application-services/src/contexts/tech-admin/queue/list-queues.test.ts create mode 100644 packages/ocom/application-services/src/contexts/tech-admin/queue/list-queues.ts create mode 100644 packages/ocom/application-services/src/contexts/tech-admin/queue/peek-messages.ts create mode 100644 packages/ocom/application-services/src/contexts/tech-admin/queue/queue-list.ts create mode 100644 packages/ocom/application-services/src/contexts/tech-admin/queue/queue-operations.test.ts create mode 100644 packages/ocom/application-services/src/contexts/tech-admin/queue/queue-operations.ts create mode 100644 packages/ocom/application-services/src/contexts/tech-admin/queue/queue-permissions.test.ts create mode 100644 packages/ocom/application-services/src/contexts/tech-admin/queue/queue-permissions.ts create mode 100644 packages/ocom/application-services/src/contexts/tech-admin/queue/send-message.test.ts create mode 100644 packages/ocom/application-services/src/contexts/tech-admin/queue/send-message.ts create mode 100644 packages/ocom/graphql/src/schema/types/tech-admin.graphql create mode 100644 packages/ocom/graphql/src/schema/types/tech-admin.resolvers.test.ts create mode 100644 packages/ocom/graphql/src/schema/types/tech-admin.resolvers.ts create mode 100644 packages/ocom/ui-staff-route-tech-admin/src/components/queue-explorer.container.graphql create mode 100644 packages/ocom/ui-staff-route-tech-admin/src/components/queue-explorer.container.tsx create mode 100644 packages/ocom/ui-staff-route-tech-admin/src/components/queue-explorer.tsx create mode 100644 packages/ocom/ui-staff-route-tech-admin/src/pages/queue-explorer.tsx create mode 100644 packages/ocom/ui-staff-route-tech-admin/src/pages/tech-admin.tsx diff --git a/apps/ui-community/package.json b/apps/ui-community/package.json index 1055de59f..bf6bf4741 100644 --- a/apps/ui-community/package.json +++ b/apps/ui-community/package.json @@ -52,6 +52,7 @@ "@types/react-dom": "^19.1.6", "@vitejs/plugin-react": "^6.0.1", "@vitest/coverage-istanbul": "catalog:", + "autoprefixer": "^10.5.4", "esbuild": "catalog:", "jsdom": "^26.1.0", "rollup-plugin-visualizer": "^6.0.5", @@ -60,7 +61,7 @@ "tailwindcss": "^3.4.17", "typescript": "catalog:", "vite": "catalog:", - "vitest": "catalog:", - "vite-plugin-node-polyfills": "catalog:" + "vite-plugin-node-polyfills": "catalog:", + "vitest": "catalog:" } } diff --git a/apps/ui-staff/src/hooks/use-staff-permissions.ts b/apps/ui-staff/src/hooks/use-staff-permissions.ts index 3d42e9ed0..6a96ccb51 100644 --- a/apps/ui-staff/src/hooks/use-staff-permissions.ts +++ b/apps/ui-staff/src/hooks/use-staff-permissions.ts @@ -34,6 +34,8 @@ const CURRENT_STAFF_USER_QUERY = gql` } techAdminPermissions { canManageTechAdmin + canViewQueues + canSendQueueMessages } } } @@ -49,6 +51,8 @@ interface StaffPermissions { canViewStaffUsers: boolean; canManageFinance: boolean; canManageTechAdmin: boolean; + canViewQueues: boolean; + canSendQueueMessages: boolean; canViewRoles: boolean; canAddRole: boolean; canEditRole: boolean; @@ -72,7 +76,7 @@ interface StaffUserQueryResult { userPermissions: { canManageUsers: boolean; canAssignStaffRoles: boolean; canViewStaffUsers: boolean }; staffRolePermissions: { canViewRoles: boolean; canAddRole: boolean; canEditRole: boolean; canRemoveRole: boolean }; financePermissions: { canManageFinance: boolean }; - techAdminPermissions: { canManageTechAdmin: boolean }; + techAdminPermissions: { canManageTechAdmin: boolean; canViewQueues: boolean; canSendQueueMessages: boolean }; }; }; }; @@ -104,6 +108,8 @@ export const useStaffPermissions = (): { canViewStaffUsers: rolePermissions.userPermissions.canViewStaffUsers || rolePermissions.userPermissions.canManageUsers || isTechAdmin, canManageFinance: rolePermissions.financePermissions.canManageFinance || isTechAdmin, canManageTechAdmin: isTechAdmin, + canViewQueues: rolePermissions.techAdminPermissions.canViewQueues, + canSendQueueMessages: rolePermissions.techAdminPermissions.canSendQueueMessages, canViewRoles: rolePermissions.staffRolePermissions.canViewRoles || rolePermissions.communityPermissions.canManageStaffRolesAndPermissions || isTechAdmin, canAddRole: rolePermissions.staffRolePermissions.canAddRole || rolePermissions.communityPermissions.canManageStaffRolesAndPermissions || isTechAdmin, canEditRole: rolePermissions.staffRolePermissions.canEditRole || rolePermissions.communityPermissions.canManageStaffRolesAndPermissions || isTechAdmin, diff --git a/codegen.yml b/codegen.yml index 09958a8fe..94a7e4d7f 100644 --- a/codegen.yml +++ b/codegen.yml @@ -169,6 +169,20 @@ generates: - typescript-operations - typed-document-node + './packages/ocom/ui-staff-route-tech-admin/src/generated.tsx': + documents: + - './packages/ocom/ui-staff-route-tech-admin/src/**/**.graphql' + config: + withHooks: true + withHOC: false + withComponent: false + useTypeImports: true + enumsAsTypes: true + plugins: + - typescript + - typescript-operations + - typed-document-node + # Cellix core base type defs (static array for rolldown bundling) './packages/cellix/graphql-core/src/schema/base-type-defs.generated.ts': plugins: diff --git a/knip.json b/knip.json index 04aeba72b..f10b164d1 100644 --- a/knip.json +++ b/knip.json @@ -136,6 +136,7 @@ "@cellix/graphql-codegen", "@graphql-typed-document-node/core", "@vitest/coverage-v8", + "autoprefixer", "ts-scope-trimmer-plugin", "chrome-devtools-mcp" ], diff --git a/packages/cellix/service-queue-storage/README.md b/packages/cellix/service-queue-storage/README.md index 5c1673462..e95883362 100644 --- a/packages/cellix/service-queue-storage/README.md +++ b/packages/cellix/service-queue-storage/README.md @@ -141,6 +141,28 @@ await service.sendMessageToOrderCreatedQueue({ }); ``` +### Send to a registered queue selected at runtime + +Operational workflows can send to the physical name of any queue registered as +either inbound or outbound. The service rejects unregistered names and validates +the payload against the selected queue's schema before enqueueing it. + +```ts +await service.sendMessageToRegisteredQueue('import-requests', { + importId: 'import-123', +}, { + visibilityTimeoutSeconds: 30, + loggingDirection: 'inbound', + loggingTags: { source: 'operations' }, + loggingMetadata: { reason: 'replay' }, +}); +``` + +This is intentionally narrower than the raw Azure transport. Prefer generated +`sendMessageTo...Queue` methods when the destination is known at compile time. +The operation accepts all `SendMessageOptions`; logging values from the selected +queue definition are defaults, and explicitly supplied options take precedence. + ### Receive from an inbound queue ```ts @@ -159,6 +181,27 @@ const message = await service.receiveFromImportRequestsQueue(queueItem, { const messages = await service.peekAtImportRequestsQueue(); ``` +### Peek at a poison queue + +Use the generated poison-queue method to inspect messages Azure Functions moved +after retry exhaustion. It is read-only and returns the same payload type as the +primary queue. + +```ts +const messages = await service.peekAtImportRequestsPoisonQueue(); +``` + +### Get an approximate queue message count + +Use the generated count methods when an operational view needs the number of +messages currently reported by Azure Queue Storage. Counts are approximate and +include messages that are not presently visible. + +```ts +const primaryCount = await service.getImportRequestsQueueMessageCount(); +const poisonCount = await service.getImportRequestsPoisonQueueMessageCount(); +``` + ## Queue Naming Each queue has: @@ -206,6 +249,7 @@ If you want to provision only a subset, pass `serviceDefaults.provisionQueues` t - `createRegisteredQueueService` - `QueueRegistryOperations` - `QueueRegistryService` +- `RegisteredQueueSender` - `QueueStorageConfig` - `QueueLoggingConfig` - `QueueTriggerMetadata` diff --git a/packages/cellix/service-queue-storage/src/index.ts b/packages/cellix/service-queue-storage/src/index.ts index 1e24249e1..1d4decc75 100644 --- a/packages/cellix/service-queue-storage/src/index.ts +++ b/packages/cellix/service-queue-storage/src/index.ts @@ -13,6 +13,7 @@ export type { QueueRegistryService, QueueServiceConstructorOptions, RegisteredQueueRegistry, + RegisteredQueueSender, RegisteredQueueService, } from './register-queues.ts'; export { createRegisteredQueueService, deriveProvisionQueues, registerQueues } from './register-queues.ts'; diff --git a/packages/cellix/service-queue-storage/src/interfaces.ts b/packages/cellix/service-queue-storage/src/interfaces.ts index 98a65034d..dd678d551 100644 --- a/packages/cellix/service-queue-storage/src/interfaces.ts +++ b/packages/cellix/service-queue-storage/src/interfaces.ts @@ -213,6 +213,13 @@ export interface IQueueStorageOperations { * @returns Decoded queue messages without altering visibility or dequeue state. */ peekMessages<_T = unknown>(queue: string, opts?: PeekMessagesOptions): Promise[]>; + /** + * Reads Azure Queue Storage's approximate visible and invisible message count. + * + * @param queue - Physical Azure Queue Storage queue name. + * @returns The approximate number of messages currently in the queue. + */ + getApproximateMessageCount(queue: string): Promise; } type QueueMessageSchema = Readonly>; diff --git a/packages/cellix/service-queue-storage/src/internal-queue-storage-service.test.ts b/packages/cellix/service-queue-storage/src/internal-queue-storage-service.test.ts index 8ff15448a..03119290c 100644 --- a/packages/cellix/service-queue-storage/src/internal-queue-storage-service.test.ts +++ b/packages/cellix/service-queue-storage/src/internal-queue-storage-service.test.ts @@ -97,6 +97,25 @@ describe('InternalQueueStorageService', () => { expect(sendMessage).toHaveBeenCalledWith(expect.any(String), { visibilityTimeout: 45 }); }); + it('returns Azure Queue Storage approximate message counts', async () => { + const getProperties = vi.fn(async () => ({ approximateMessagesCount: 12 })); + fromConnectionStringMock.mockImplementation((_conn: string) => ({ + getQueueClient: vi.fn((_q: string) => ({ + sendMessage: vi.fn(async (_m: string) => ({ messageId: 'mid' })), + createIfNotExists: vi.fn(async () => ({ succeeded: true })), + receiveMessages: vi.fn(async () => ({ receivedMessageItems: [] })), + peekMessages: vi.fn(async () => ({ peekedMessageItems: [] })), + deleteMessage: vi.fn(async () => ({})), + getProperties, + })), + })); + const svc = new InternalQueueStorageService({ connectionString: 'UseDevelopmentStorage=true' }); + await svc.startUp(); + + await expect(svc.getApproximateMessageCount('q')).resolves.toBe(12); + expect(getProperties).toHaveBeenCalledOnce(); + }); + it('createQueueIfNotExists does not throw for missing queue', async () => { const svc = new InternalQueueStorageService({ connectionString: 'UseDevelopmentStorage=true' }); await svc.startUp(); diff --git a/packages/cellix/service-queue-storage/src/internal-queue-storage-service.ts b/packages/cellix/service-queue-storage/src/internal-queue-storage-service.ts index b1e7bdb9c..8c8825afd 100644 --- a/packages/cellix/service-queue-storage/src/internal-queue-storage-service.ts +++ b/packages/cellix/service-queue-storage/src/internal-queue-storage-service.ts @@ -331,4 +331,12 @@ export class InternalQueueStorageService implements InternalQueueTransport { } return out; } + + /** + * Gets Azure Queue Storage's approximate number of messages in a queue. + */ + public async getApproximateMessageCount(queue: string): Promise { + const properties = await this.getQueueClient(queue).getProperties(); + return properties.approximateMessagesCount ?? 0; + } } diff --git a/packages/cellix/service-queue-storage/src/queue-consumer.test.ts b/packages/cellix/service-queue-storage/src/queue-consumer.test.ts index a0d9f0d7b..38297e2a3 100644 --- a/packages/cellix/service-queue-storage/src/queue-consumer.test.ts +++ b/packages/cellix/service-queue-storage/src/queue-consumer.test.ts @@ -135,4 +135,25 @@ describe('registerQueues', () => { }, ]); }); + + it('peeks at the inbound poison queue', async () => { + const registry = createInboundRegistry(); + const svc = new registry.Service({ connectionString: 'UseDevelopmentStorage=true' }); + peekedMessageItems = [ + { + messageId: 'poison-msg-1', + messageText: Buffer.from(JSON.stringify({ requestId: 'r1' })).toString('base64'), + dequeueCount: 5, + }, + ]; + await svc.startUp(); + + await expect(svc.peekAtImportRequestsPoisonQueue(8)).resolves.toEqual([ + { + id: 'poison-msg-1', + payload: { requestId: 'r1' }, + dequeueCount: 5, + }, + ]); + }); }); diff --git a/packages/cellix/service-queue-storage/src/queue-consumer.ts b/packages/cellix/service-queue-storage/src/queue-consumer.ts index 69b183389..3602fb61c 100644 --- a/packages/cellix/service-queue-storage/src/queue-consumer.ts +++ b/packages/cellix/service-queue-storage/src/queue-consumer.ts @@ -11,7 +11,8 @@ type Capitalize = S extends `${infer F}${infer R}` ? `${Upperc * Public consumer methods generated for an application's inbound queues. * * Each queue key becomes a strongly-typed `receiveFrom...Queue` method and a - * matching `peekAt...Queue` method on the registered service surface. + * matching `peekAt...Queue` and `get...QueueMessageCount` method on the + * registered service surface. * * @typeParam I - Inbound queue definition map passed to `registerQueues()`. * @@ -31,12 +32,18 @@ export type QueueConsumerContext = { [K in keyof I as `receiveFrom${Capitalize}Queue`]: (payload: unknown, metadata?: QueueTriggerMetadata) => Promise>>; } & { [K in keyof I as `peekAt${Capitalize}Queue`]: (maxMessages?: number) => Promise>[]>; +} & { + [K in keyof I as `peekAt${Capitalize}PoisonQueue`]: (maxMessages?: number) => Promise>[]>; +} & { + [K in keyof I as `get${Capitalize}QueueMessageCount`]: () => Promise; +} & { + [K in keyof I as `get${Capitalize}PoisonQueueMessageCount`]: () => Promise; }; type QueueMessage = { id: string; popReceipt?: string; payload: T; dequeueCount?: number }; export function createQueueConsumer( - service: Pick, + service: Pick, definitions: I, validators: Record, ): QueueConsumerContext { @@ -91,6 +98,9 @@ export function createQueueConsumer( }; context[`peekAt${cap}Queue`] = (maxMessages?: number) => service.peekMessages(def.queueName, { maxMessages: maxMessages ?? 32 }); + context[`peekAt${cap}PoisonQueue`] = (maxMessages?: number) => service.peekMessages(`${def.queueName}-poison`, { maxMessages: maxMessages ?? 32 }); + context[`get${cap}QueueMessageCount`] = () => service.getApproximateMessageCount(def.queueName); + context[`get${cap}PoisonQueueMessageCount`] = () => service.getApproximateMessageCount(`${def.queueName}-poison`); } return context as QueueConsumerContext; diff --git a/packages/cellix/service-queue-storage/src/queue-producer.test.ts b/packages/cellix/service-queue-storage/src/queue-producer.test.ts index e7d2fa0f7..7d403a9cc 100644 --- a/packages/cellix/service-queue-storage/src/queue-producer.test.ts +++ b/packages/cellix/service-queue-storage/src/queue-producer.test.ts @@ -4,23 +4,25 @@ import { describeFeature, loadFeature } from '@amiceli/vitest-cucumber'; import { describe, expect, it, vi } from 'vitest'; import { registerQueues } from './index.ts'; -type SentMessage = { queue: string; messageText: string }; +type SentMessage = { queue: string; messageText: string; options: { visibilityTimeout?: number } }; type MockPeekedMessage = { messageId: string; messageText: string; dequeueCount?: number }; let sentMessages: SentMessage[] = []; let peekedMessageItems: MockPeekedMessage[] = []; +let approximateMessagesCount = 0; vi.mock('@azure/storage-queue', () => ({ QueueServiceClient: { fromConnectionString: vi.fn(() => ({ getQueueClient: vi.fn((queue: string) => ({ - sendMessage: vi.fn((messageText: string) => { - sentMessages.push({ queue, messageText }); + sendMessage: vi.fn((messageText: string, options: { visibilityTimeout?: number }) => { + sentMessages.push({ queue, messageText, options }); return Promise.resolve({ messageId: 'mid' }); }), createIfNotExists: vi.fn(async () => ({ succeeded: true })), receiveMessages: vi.fn(async () => ({ receivedMessageItems: [] })), peekMessages: vi.fn(async () => ({ peekedMessageItems })), + getProperties: vi.fn(async () => ({ approximateMessagesCount })), deleteMessage: vi.fn(async () => ({})), })), })), @@ -62,6 +64,23 @@ function createPeekRegistry() { }); } +function createBidirectionalRegistry() { + return registerQueues({ + outbound: { + emailNotifications: { + queueName: 'email-notifications', + schema: { type: 'object', properties: { to: { type: 'string' } }, required: ['to'], additionalProperties: false }, + }, + }, + inbound: { + importRequests: { + queueName: 'import-requests', + schema: { type: 'object', properties: { importId: { type: 'string' } }, required: ['importId'], additionalProperties: false }, + }, + }, + }); +} + type OutboundRegistry = ReturnType; type PeekRegistry = ReturnType; type OutboundService = InstanceType; @@ -77,6 +96,7 @@ describe('registerQueues', () => { vi.clearAllMocks(); sentMessages = []; peekedMessageItems = []; + approximateMessagesCount = 0; }); Scenario('Successfully sending a valid message to an outbound queue', ({ Given, When, Then, And }) => { @@ -169,6 +189,52 @@ describe('registerQueues', () => { await expect(svc.sendMessageToEmailNotificationsQueue({ to: 'not-an-email', subject: 'hello' })).rejects.toThrow('Invalid payload for queue "email-notifications": /to must match format "email"'); }); + it('sends valid messages to registered outbound and inbound queues only', async () => { + const registry = createBidirectionalRegistry(); + const svc = new registry.Service({ connectionString: 'UseDevelopmentStorage=true' }); + await svc.startUp(); + + await svc.sendMessageToRegisteredQueue('email-notifications', { to: 'user@example.com' }); + await svc.sendMessageToRegisteredQueue('import-requests', { importId: 'import-1' }); + + expect(sentMessages.map(({ queue }) => queue)).toEqual(['email-notifications', 'import-requests']); + await expect(svc.sendMessageToRegisteredQueue('unknown-queue', { value: 'ignored' })).rejects.toThrow('Queue "unknown-queue" is not registered'); + await expect(svc.sendMessageToRegisteredQueue('import-requests', { invalid: true })).rejects.toThrow('Invalid payload for queue "import-requests"'); + }); + + it('passes all send options through a registered queue send', async () => { + const registry = createBidirectionalRegistry(); + const logMessage = vi.fn().mockResolvedValue(undefined); + const svc = new registry.Service({ + connectionString: 'UseDevelopmentStorage=true', + logger: { logMessage }, + logging: { enabled: true, container: '', await: true }, + }); + await svc.startUp(); + + await svc.sendMessageToRegisteredQueue( + 'import-requests', + { importId: 'registered-options-test' }, + { + visibilityTimeoutSeconds: 45, + loggingDirection: 'inbound', + loggingTags: { source: 'tech-admin' }, + loggingMetadata: { reason: 'replay' }, + }, + ); + + expect(sentMessages.find((message) => message.queue === 'import-requests' && JSON.parse(Buffer.from(message.messageText, 'base64').toString('utf-8')).importId === 'registered-options-test')?.options).toEqual({ + visibilityTimeout: 45, + }); + expect(logMessage).toHaveBeenCalledWith( + expect.objectContaining({ + direction: 'inbound', + metadata: { reason: 'replay' }, + tags: { source: 'tech-admin', queueName: 'import-requests' }, + }), + ); + }); + it('allows peeking invalid outbound payloads without throwing', async () => { const registry = createPeekRegistry(); const svc = new registry.Service({ connectionString: 'UseDevelopmentStorage=true' }); @@ -189,4 +255,40 @@ describe('registerQueues', () => { }, ]); }); + + it('peeks at the outbound poison queue', async () => { + const registry = createPeekRegistry(); + const svc = new registry.Service({ connectionString: 'UseDevelopmentStorage=true' }); + peekedMessageItems = [ + { + messageId: 'poison-msg-1', + messageText: Buffer.from(JSON.stringify({ to: 'user@example.com', subject: 'failed' })).toString('base64'), + dequeueCount: 5, + }, + ]; + await svc.startUp(); + + await expect(svc.peekAtEmailNotificationsPoisonQueue(8)).resolves.toEqual([ + { + id: 'poison-msg-1', + payload: { to: 'user@example.com', subject: 'failed' }, + dequeueCount: 5, + }, + ]); + }); + + it('gets approximate message counts for outbound primary and poison queues', async () => { + const registry = createPeekRegistry(); + const svc = new registry.Service({ connectionString: 'UseDevelopmentStorage=true' }); + approximateMessagesCount = 12; + await svc.startUp(); + + const queueStorage = svc as unknown as { + getEmailNotificationsQueueMessageCount: () => Promise; + getEmailNotificationsPoisonQueueMessageCount: () => Promise; + }; + + await expect(queueStorage.getEmailNotificationsQueueMessageCount()).resolves.toBe(12); + await expect(queueStorage.getEmailNotificationsPoisonQueueMessageCount()).resolves.toBe(12); + }); }); diff --git a/packages/cellix/service-queue-storage/src/queue-producer.ts b/packages/cellix/service-queue-storage/src/queue-producer.ts index a803306ca..382c676ef 100644 --- a/packages/cellix/service-queue-storage/src/queue-producer.ts +++ b/packages/cellix/service-queue-storage/src/queue-producer.ts @@ -9,7 +9,8 @@ type Capitalize = S extends `${infer F}${infer R}` ? `${Upperc * Public producer methods generated for an application's outbound queues. * * Each queue key becomes a strongly-typed `sendMessageTo...Queue` method and a - * matching `peekAt...Queue` method on the registered service surface. + * matching `peekAt...Queue` and `get...QueueMessageCount` method on the + * registered service surface. * * @typeParam O - Outbound queue definition map passed to `registerQueues()`. * @@ -29,9 +30,19 @@ export type QueueProducerContext = { [K in keyof O as `sendMessageTo${Capitalize}Queue`]: (payload: MessagePayload) => Promise; } & { [K in keyof O as `peekAt${Capitalize}Queue`]: (maxMessages?: number) => Promise>[]>; +} & { + [K in keyof O as `peekAt${Capitalize}PoisonQueue`]: (maxMessages?: number) => Promise>[]>; +} & { + [K in keyof O as `get${Capitalize}QueueMessageCount`]: () => Promise; +} & { + [K in keyof O as `get${Capitalize}PoisonQueueMessageCount`]: () => Promise; }; -export function createQueueProducer(service: Pick, definitions: O, validators: Record): QueueProducerContext { +export function createQueueProducer( + service: Pick, + definitions: O, + validators: Record, +): QueueProducerContext { const context = {} as Record; for (const [key, def] of Object.entries(definitions)) { @@ -54,6 +65,9 @@ export function createQueueProducer(service: Pick service.peekMessages(def.queueName, { maxMessages: maxMessages ?? 32 }); + context[`peekAt${cap}PoisonQueue`] = (maxMessages?: number) => service.peekMessages(`${def.queueName}-poison`, { maxMessages: maxMessages ?? 32 }); + context[`get${cap}QueueMessageCount`] = () => service.getApproximateMessageCount(def.queueName); + context[`get${cap}PoisonQueueMessageCount`] = () => service.getApproximateMessageCount(`${def.queueName}-poison`); } return context as unknown as QueueProducerContext; diff --git a/packages/cellix/service-queue-storage/src/register-queues.test.ts b/packages/cellix/service-queue-storage/src/register-queues.test.ts index ff0760c4c..f59925dc5 100644 --- a/packages/cellix/service-queue-storage/src/register-queues.test.ts +++ b/packages/cellix/service-queue-storage/src/register-queues.test.ts @@ -54,6 +54,22 @@ describe('registerQueues', () => { }); }); }); + + it('provides poison queue peek stubs for outbound and inbound queues', () => { + const registry = registerQueues({ outbound: { a: { queueName: 'q-a', schema: {} } }, inbound: { b: { queueName: 'q-b', schema: {} } } }); + + expect(registry.producer.peekAtAPoisonQueue).toBeDefined(); + expect(registry.consumer.peekAtBPoisonQueue).toBeDefined(); + }); + + it('provides queue message count stubs for outbound and inbound queues', () => { + const registry = registerQueues({ outbound: { a: { queueName: 'q-a', schema: {} } }, inbound: { b: { queueName: 'q-b', schema: {} } } }); + + expect(registry.producer.getAQueueMessageCount).toBeDefined(); + expect(registry.producer.getAPoisonQueueMessageCount).toBeDefined(); + expect(registry.consumer.getBQueueMessageCount).toBeDefined(); + expect(registry.consumer.getBPoisonQueueMessageCount).toBeDefined(); + }); }); describe('deriveProvisionQueues', () => { diff --git a/packages/cellix/service-queue-storage/src/register-queues.ts b/packages/cellix/service-queue-storage/src/register-queues.ts index b68b6c46d..ea846354f 100644 --- a/packages/cellix/service-queue-storage/src/register-queues.ts +++ b/packages/cellix/service-queue-storage/src/register-queues.ts @@ -1,10 +1,11 @@ import Ajv from 'ajv'; import addFormats from 'ajv-formats'; -import type { QueueMap, QueueStorageConfig } from './interfaces.ts'; +import type { QueueMap, QueueStorageConfig, SendMessageOptions } from './interfaces.ts'; import { InternalQueueStorageService, type QueueServiceLifecycle, type QueueServiceLogging } from './internal-queue-storage-service.ts'; +import { resolveLoggingFields } from './logging-fields.ts'; import { createQueueConsumer, type QueueConsumerContext } from './queue-consumer.ts'; import { createQueueProducer, type QueueProducerContext } from './queue-producer.ts'; -import type { QueuePayloadValidator } from './validation.ts'; +import { formatQueueValidationErrors, type QueuePayloadValidator } from './validation.ts'; type QueueServiceDefaults = Pick; @@ -58,7 +59,33 @@ if (typeof addFormatsAny === 'function') { * type ServiceQueueStorage = RegisteredQueueService; * ``` */ -export type RegisteredQueueService = QueueServiceLifecycle & QueueServiceLogging & QueueProducerContext & QueueConsumerContext; +export type RegisteredQueueService = QueueServiceLifecycle & QueueServiceLogging & RegisteredQueueSender & QueueProducerContext & QueueConsumerContext; + +/** + * Validated generic send capability for queues registered with a service. + * + * Unlike the low-level transport's `sendMessage`, this operation accepts only a + * physical queue name declared in the registry and validates the payload using + * that queue's JSON Schema. It is intended for controlled operational workflows + * that select a registered queue at runtime. + */ +export interface RegisteredQueueSender { + /** + * Sends a message to a registered inbound or outbound queue. + * + * @param queueName - Physical name of a queue registered with `registerQueues`. + * @param payload - Object payload validated against the registered queue schema. + * @param options - Optional Azure delivery and logging options. Registry-derived logging values are defaults and caller-provided values take precedence. + * @returns Resolves when Azure Queue Storage accepts the validated message. + * @throws Error when `queueName` is unregistered or `payload` fails schema validation. + * + * @example + * ```ts + * await service.sendMessageToRegisteredQueue('import-requests', { importId: 'import-123' }); + * ``` + */ + sendMessageToRegisteredQueue(queueName: string, payload: object, options?: SendMessageOptions): Promise; +} /** * Full registry shape returned by {@link registerQueues}. @@ -245,6 +272,9 @@ export function registerQueues(config: { const cap = capitalizeQueueKey(key); out[`sendMessageTo${cap}Queue`] = () => Promise.reject(new Error('Queue producer not bound to a registered queue service')); out[`peekAt${cap}Queue`] = (_maxMessages?: number) => Promise.resolve([]); + out[`peekAt${cap}PoisonQueue`] = (_maxMessages?: number) => Promise.resolve([]); + out[`get${cap}QueueMessageCount`] = () => Promise.resolve(0); + out[`get${cap}PoisonQueueMessageCount`] = () => Promise.resolve(0); } return out as QueueProducerContext; }; @@ -255,6 +285,9 @@ export function registerQueues(config: { const cap = capitalizeQueueKey(key); out[`receiveFrom${cap}Queue`] = () => Promise.resolve(undefined); out[`peekAt${cap}Queue`] = (_maxMessages?: number) => Promise.resolve([]); + out[`peekAt${cap}PoisonQueue`] = (_maxMessages?: number) => Promise.resolve([]); + out[`get${cap}QueueMessageCount`] = () => Promise.resolve(0); + out[`get${cap}PoisonQueueMessageCount`] = () => Promise.resolve(0); } return out as QueueConsumerContext; }; @@ -262,6 +295,13 @@ export function registerQueues(config: { const producer = makeProducerStub(config.outbound); const consumer = makeConsumerStub(config.inbound); const defaultProvisionQueues = config.serviceDefaults?.provisionQueues ?? deriveProvisionQueues(config.outbound, config.inbound); + const registeredQueues = new Map(); + for (const [key, definition] of Object.entries(config.outbound)) { + registeredQueues.set(definition.queueName, { definition, validator: outboundValidators[key] as QueuePayloadValidator }); + } + for (const [key, definition] of Object.entries(config.inbound)) { + registeredQueues.set(definition.queueName, { definition, validator: inboundValidators[key] as QueuePayloadValidator }); + } /** * Base class returned by `registerQueues`. Extends the internal queue transport with @@ -284,6 +324,25 @@ export function registerQueues(config: { Object.assign(this, createQueueProducer(this, config.outbound, outboundValidators)); Object.assign(this, createQueueConsumer(this, config.inbound, inboundValidators)); } + + public async sendMessageToRegisteredQueue(queueName: string, payload: object, options?: SendMessageOptions): Promise { + const registeredQueue = registeredQueues.get(queueName); + if (!registeredQueue) { + throw new Error(`Queue "${queueName}" is not registered`); + } + if (!registeredQueue.validator(payload)) { + throw new Error(`Invalid payload for queue "${queueName}": ${formatQueueValidationErrors(registeredQueue.validator.errors)}`); + } + + const loggingTags = resolveLoggingFields(registeredQueue.definition.loggingTags, payload); + const loggingMetadata = resolveLoggingFields(registeredQueue.definition.loggingMetadata, payload); + await this.sendMessage(queueName, payload, { + loggingDirection: 'outbound', + ...(loggingTags === undefined ? {} : { loggingTags }), + ...(loggingMetadata === undefined ? {} : { loggingMetadata }), + ...options, + }); + } } return { diff --git a/packages/ocom/application-services/src/contexts/tech-admin/index.ts b/packages/ocom/application-services/src/contexts/tech-admin/index.ts new file mode 100644 index 000000000..6d4db0482 --- /dev/null +++ b/packages/ocom/application-services/src/contexts/tech-admin/index.ts @@ -0,0 +1,13 @@ +import type { DataSources } from '@ocom/persistence'; +import type { QueueStorageOperations } from '@ocom/service-queue-storage'; +import { TechAdminQueue, type TechAdminQueueApplicationService } from './queue/index.ts'; + +export interface TechAdminContextApplicationService { + Queue: TechAdminQueueApplicationService; +} + +export const TechAdmin = (dataSources: DataSources, queueStorageService: QueueStorageOperations, staffUserExternalId: string | undefined): TechAdminContextApplicationService => { + return { + Queue: TechAdminQueue(dataSources, queueStorageService, staffUserExternalId), + }; +}; diff --git a/packages/ocom/application-services/src/contexts/tech-admin/queue/get-queue-message-count.test.ts b/packages/ocom/application-services/src/contexts/tech-admin/queue/get-queue-message-count.test.ts new file mode 100644 index 000000000..961270947 --- /dev/null +++ b/packages/ocom/application-services/src/contexts/tech-admin/queue/get-queue-message-count.test.ts @@ -0,0 +1,17 @@ +import type { QueueStorageOperations } from '@ocom/service-queue-storage'; +import { describe, expect, it, vi } from 'vitest'; +import { getQueueMessageCount } from './get-queue-message-count.ts'; + +describe('getQueueMessageCount', () => { + it('returns a queue-specific error when the physical queue is missing', async () => { + const checkPermission = vi.fn().mockResolvedValue(undefined); + const queueStorageService = { + getCommunityCreationQueueMessageCount: vi.fn().mockRejectedValue(new Error('The specified queue does not exist.')), + } as unknown as QueueStorageOperations; + + await expect(getQueueMessageCount(queueStorageService, checkPermission)({ queueName: 'community-creation' })).resolves.toEqual({ + errorMessage: 'The specified queue does not exist.', + }); + expect(checkPermission).toHaveBeenCalledOnce(); + }); +}); diff --git a/packages/ocom/application-services/src/contexts/tech-admin/queue/get-queue-message-count.ts b/packages/ocom/application-services/src/contexts/tech-admin/queue/get-queue-message-count.ts new file mode 100644 index 000000000..ad27a2891 --- /dev/null +++ b/packages/ocom/application-services/src/contexts/tech-admin/queue/get-queue-message-count.ts @@ -0,0 +1,26 @@ +import type { QueueStorageOperations } from '@ocom/service-queue-storage'; +import { techAdminQueueDefinitions } from './queue-list.ts'; + +export interface TechAdminQueueMessageCount { + value?: number; + errorMessage?: string; +} + +export interface GetQueueMessageCountCommand { + queueName: string; +} + +export function getQueueMessageCount(queueStorageService: QueueStorageOperations, checkPermission: () => Promise): (command: GetQueueMessageCountCommand) => Promise { + return async ({ queueName }): Promise => { + await checkPermission(); + + const queue = techAdminQueueDefinitions.find(({ name }) => name === queueName); + if (!queue) return { errorMessage: `Queue "${queueName}" is not registered` }; + + try { + return { value: await queue.getMessageCount(queueStorageService) }; + } catch (error) { + return { errorMessage: error instanceof Error ? error.message : String(error) }; + } + }; +} diff --git a/packages/ocom/application-services/src/contexts/tech-admin/queue/index.ts b/packages/ocom/application-services/src/contexts/tech-admin/queue/index.ts new file mode 100644 index 000000000..f1cad49ad --- /dev/null +++ b/packages/ocom/application-services/src/contexts/tech-admin/queue/index.ts @@ -0,0 +1,28 @@ +import type { DataSources } from '@ocom/persistence'; +import type { QueueStorageOperations } from '@ocom/service-queue-storage'; +import { type GetQueueMessageCountCommand, getQueueMessageCount, type TechAdminQueueMessageCount } from './get-queue-message-count.ts'; +import { listQueues } from './list-queues.ts'; +import { type PeekQueueMessagesCommand, peekMessages } from './peek-messages.ts'; +import type { TechAdminQueue as TechAdminQueueListItem, TechAdminQueueMessage } from './queue-list.ts'; +import { checkCanSendQueueMessages, checkCanViewQueues, checkPermissionOnce } from './queue-permissions.ts'; +import { type SendQueueMessageCommand, sendMessage } from './send-message.ts'; + +export interface TechAdminQueueApplicationService { + listQueues: () => Promise; + getMessageCount: (command: GetQueueMessageCountCommand) => Promise; + sendMessage: (command: SendQueueMessageCommand) => Promise; + peekMessages: (command: PeekQueueMessagesCommand) => Promise; +} + +export const TechAdminQueue = (dataSources: DataSources, queueStorageService: QueueStorageOperations, staffUserExternalId: string | undefined): TechAdminQueueApplicationService => { + const checkViewPermission = checkCanViewQueues(dataSources, staffUserExternalId); + const checkQueueViewPermissionOnce = checkPermissionOnce(checkViewPermission); + const checkSendPermission = checkCanSendQueueMessages(dataSources, staffUserExternalId); + + return { + listQueues: listQueues(checkQueueViewPermissionOnce), + getMessageCount: getQueueMessageCount(queueStorageService, checkQueueViewPermissionOnce), + sendMessage: sendMessage(queueStorageService, checkSendPermission), + peekMessages: peekMessages(queueStorageService, checkViewPermission), + }; +}; diff --git a/packages/ocom/application-services/src/contexts/tech-admin/queue/list-queues.test.ts b/packages/ocom/application-services/src/contexts/tech-admin/queue/list-queues.test.ts new file mode 100644 index 000000000..1ca5868da --- /dev/null +++ b/packages/ocom/application-services/src/contexts/tech-admin/queue/list-queues.test.ts @@ -0,0 +1,11 @@ +import { describe, expect, it, vi } from 'vitest'; +import { listQueues } from './list-queues.ts'; + +describe('listQueues', () => { + it('returns every centralized queue after authorization', async () => { + const checkPermission = vi.fn().mockResolvedValue(undefined); + + await expect(listQueues(checkPermission)()).resolves.toEqual([{ name: 'community-creation' }, { name: 'community-creation-poison' }, { name: 'end-user-update' }, { name: 'end-user-update-poison' }]); + expect(checkPermission).toHaveBeenCalledOnce(); + }); +}); diff --git a/packages/ocom/application-services/src/contexts/tech-admin/queue/list-queues.ts b/packages/ocom/application-services/src/contexts/tech-admin/queue/list-queues.ts new file mode 100644 index 000000000..213b5382a --- /dev/null +++ b/packages/ocom/application-services/src/contexts/tech-admin/queue/list-queues.ts @@ -0,0 +1,8 @@ +import { type TechAdminQueue, techAdminQueueDefinitions } from './queue-list.ts'; + +export function listQueues(checkPermission: () => Promise): () => Promise { + return async (): Promise => { + await checkPermission(); + return techAdminQueueDefinitions.map(({ name }) => ({ name })).sort((left, right) => left.name.localeCompare(right.name)); + }; +} diff --git a/packages/ocom/application-services/src/contexts/tech-admin/queue/peek-messages.ts b/packages/ocom/application-services/src/contexts/tech-admin/queue/peek-messages.ts new file mode 100644 index 000000000..bd8bb107a --- /dev/null +++ b/packages/ocom/application-services/src/contexts/tech-admin/queue/peek-messages.ts @@ -0,0 +1,18 @@ +import type { QueueStorageOperations } from '@ocom/service-queue-storage'; +import { type TechAdminQueueMessage, techAdminQueueDefinitions } from './queue-list.ts'; + +export interface PeekQueueMessagesCommand { + queueName: string; + maxMessages?: number; +} + +export function peekMessages(queueStorageService: QueueStorageOperations, checkPermission: () => Promise): (command: PeekQueueMessagesCommand) => Promise { + return async ({ queueName, maxMessages }): Promise => { + await checkPermission(); + + const queue = techAdminQueueDefinitions.find(({ name }) => name === queueName); + if (!queue) throw new Error(`Queue "${queueName}" is not registered`); + + return await queue.peek(queueStorageService, maxMessages); + }; +} diff --git a/packages/ocom/application-services/src/contexts/tech-admin/queue/queue-list.ts b/packages/ocom/application-services/src/contexts/tech-admin/queue/queue-list.ts new file mode 100644 index 000000000..c1424c006 --- /dev/null +++ b/packages/ocom/application-services/src/contexts/tech-admin/queue/queue-list.ts @@ -0,0 +1,41 @@ +import type { QueueStorageOperations } from '@ocom/service-queue-storage'; + +export interface TechAdminQueueMessage { + id: string; + popReceipt?: string; + payload: unknown; + dequeueCount?: number; +} + +interface TechAdminQueueDefinition { + name: string; + peek: (queueStorageService: QueueStorageOperations, maxMessages?: number) => Promise; + getMessageCount: (queueStorageService: QueueStorageOperations) => Promise; +} + +export interface TechAdminQueue { + name: string; +} + +export const techAdminQueueDefinitions: readonly TechAdminQueueDefinition[] = [ + { + name: 'community-creation', + peek: (queueStorageService, maxMessages) => queueStorageService.peekAtCommunityCreationQueue(maxMessages), + getMessageCount: (queueStorageService) => queueStorageService.getCommunityCreationQueueMessageCount(), + }, + { + name: 'community-creation-poison', + peek: (queueStorageService, maxMessages) => queueStorageService.peekAtCommunityCreationPoisonQueue(maxMessages), + getMessageCount: (queueStorageService) => queueStorageService.getCommunityCreationPoisonQueueMessageCount(), + }, + { + name: 'end-user-update', + peek: (queueStorageService, maxMessages) => queueStorageService.peekAtEndUserUpdateQueue(maxMessages), + getMessageCount: (queueStorageService) => queueStorageService.getEndUserUpdateQueueMessageCount(), + }, + { + name: 'end-user-update-poison', + peek: (queueStorageService, maxMessages) => queueStorageService.peekAtEndUserUpdatePoisonQueue(maxMessages), + getMessageCount: (queueStorageService) => queueStorageService.getEndUserUpdatePoisonQueueMessageCount(), + }, +]; diff --git a/packages/ocom/application-services/src/contexts/tech-admin/queue/queue-operations.test.ts b/packages/ocom/application-services/src/contexts/tech-admin/queue/queue-operations.test.ts new file mode 100644 index 000000000..21a9bd07b --- /dev/null +++ b/packages/ocom/application-services/src/contexts/tech-admin/queue/queue-operations.test.ts @@ -0,0 +1,50 @@ +import type { QueueStorageOperations } from '@ocom/service-queue-storage'; +import { describe, expect, it, vi } from 'vitest'; +import { ensureCanViewQueues, registeredQueueOperations } from './queue-operations.ts'; + +describe('registeredQueueOperations', () => { + it('includes peek-only poison queue operations for every registered primary queue', async () => { + const peekMessages = vi.fn().mockResolvedValue([]); + const queueStorageService = { + sendMessageToCommunityCreationQueue: vi.fn(), + peekAtCommunityCreationQueue: vi.fn(), + peekAtEndUserUpdateQueue: vi.fn(), + peekMessages, + } as unknown as QueueStorageOperations; + + const operations = registeredQueueOperations(queueStorageService); + + expect([...operations.keys()]).toEqual(['community-creation', 'end-user-update', 'community-creation-poison', 'end-user-update-poison']); + + await operations.get('community-creation-poison')?.peek(8); + + expect(peekMessages).toHaveBeenCalledWith('community-creation-poison', { maxMessages: 8 }); + expect(operations.get('community-creation-poison')?.send).toBeUndefined(); + }); +}); + +describe('ensureCanViewQueues', () => { + function makeDataSources(canViewQueues: boolean) { + return { + readonlyDataSource: { + User: { + StaffUser: { + StaffUserReadRepo: { + getByExternalId: vi.fn().mockResolvedValue({ + role: { permissions: { techAdminPermissions: { canViewQueues } } }, + }), + }, + }, + }, + }, + } as unknown as Parameters[0]; + } + + it('allows a staff user with canViewQueues', async () => { + await expect(ensureCanViewQueues(makeDataSources(true), 'staff-1')()).resolves.toBeUndefined(); + }); + + it('rejects a staff user without canViewQueues', async () => { + await expect(ensureCanViewQueues(makeDataSources(false), 'staff-1')()).rejects.toThrow('Unauthorized to view queues'); + }); +}); diff --git a/packages/ocom/application-services/src/contexts/tech-admin/queue/queue-operations.ts b/packages/ocom/application-services/src/contexts/tech-admin/queue/queue-operations.ts new file mode 100644 index 000000000..846770332 --- /dev/null +++ b/packages/ocom/application-services/src/contexts/tech-admin/queue/queue-operations.ts @@ -0,0 +1,63 @@ +import type { DataSources } from '@ocom/persistence'; +import type { QueueStorageOperations } from '@ocom/service-queue-storage'; + +interface TechAdminQueueMessage { + id: string; + popReceipt?: string; + payload: unknown; + dequeueCount?: number; +} + +interface QueueOperations { + send?: (payload: unknown) => Promise; + peek: (maxMessages?: number) => Promise; +} + +type QueueStorageWithRawPeek = QueueStorageOperations & { + peekMessages: (queueName: string, options?: { maxMessages?: number }) => Promise; +}; + +export function registeredQueueOperations(queueStorageService: QueueStorageOperations): Map { + const queueStorageWithRawPeek = queueStorageService as QueueStorageWithRawPeek; + const primaryQueueOperations: [string, QueueOperations][] = [ + [ + 'community-creation', + { + send: (payload) => queueStorageService.sendMessageToCommunityCreationQueue(payload as Parameters[0]), + peek: (maxMessages) => queueStorageService.peekAtCommunityCreationQueue(maxMessages), + }, + ], + [ + 'end-user-update', + { + peek: (maxMessages) => queueStorageService.peekAtEndUserUpdateQueue(maxMessages), + }, + ], + ]; + + return new Map([ + ...primaryQueueOperations, + ...primaryQueueOperations.map( + ([queueName]) => + [ + `${queueName}-poison`, + { + peek: (maxMessages?: number) => queueStorageWithRawPeek.peekMessages(`${queueName}-poison`, maxMessages === undefined ? undefined : { maxMessages }), + }, + ] as [string, QueueOperations], + ), + ]); +} + +export function ensureCanViewQueues(dataSources: DataSources, staffUserExternalId: string | undefined): () => Promise { + return async (): Promise => { + if (!staffUserExternalId) { + throw new Error('Unauthorized to view queues'); + } + + const staffUser = await dataSources.readonlyDataSource.User.StaffUser.StaffUserReadRepo.getByExternalId(staffUserExternalId); + if (!staffUser?.role?.permissions.techAdminPermissions.canViewQueues) { + throw new Error('Unauthorized to view queues'); + } + }; +} diff --git a/packages/ocom/application-services/src/contexts/tech-admin/queue/queue-permissions.test.ts b/packages/ocom/application-services/src/contexts/tech-admin/queue/queue-permissions.test.ts new file mode 100644 index 000000000..4efc9ef45 --- /dev/null +++ b/packages/ocom/application-services/src/contexts/tech-admin/queue/queue-permissions.test.ts @@ -0,0 +1,38 @@ +import type { DataSources } from '@ocom/persistence'; +import { describe, expect, it, vi } from 'vitest'; +import { checkCanSendQueueMessages, checkCanViewQueues, checkPermissionOnce } from './queue-permissions.ts'; + +function makeDataSources(canViewQueues: boolean, canSendQueueMessages: boolean): DataSources { + return { + readonlyDataSource: { + User: { + StaffUser: { + StaffUserReadRepo: { + getByExternalId: vi.fn().mockResolvedValue({ + role: { permissions: { techAdminPermissions: { canViewQueues, canSendQueueMessages } } }, + }), + }, + }, + }, + }, + } as unknown as DataSources; +} + +describe('queue permissions', () => { + it('runs a permission check once for concurrent calls', async () => { + const permissionCheck = vi.fn().mockResolvedValue(undefined); + const checkPermission = checkPermissionOnce(permissionCheck); + + await Promise.all([checkPermission(), checkPermission(), checkPermission(), checkPermission()]); + + expect(permissionCheck).toHaveBeenCalledOnce(); + }); + + it('allows a staff user with queue-view permission', async () => { + await expect(checkCanViewQueues(makeDataSources(true, false), 'staff-1')()).resolves.toBeUndefined(); + }); + + it('rejects a staff user without queue-send permission', async () => { + await expect(checkCanSendQueueMessages(makeDataSources(true, false), 'staff-1')()).rejects.toThrow('Unauthorized to send queue messages'); + }); +}); diff --git a/packages/ocom/application-services/src/contexts/tech-admin/queue/queue-permissions.ts b/packages/ocom/application-services/src/contexts/tech-admin/queue/queue-permissions.ts new file mode 100644 index 000000000..d904ce378 --- /dev/null +++ b/packages/ocom/application-services/src/contexts/tech-admin/queue/queue-permissions.ts @@ -0,0 +1,27 @@ +import type { DataSources } from '@ocom/persistence'; + +// only needed becasue passport for techadmin is not implemented +export function checkPermissionOnce(checkPermission: () => Promise): () => Promise { + let permissionCheck: Promise | undefined; + + return (): Promise => { + permissionCheck ??= checkPermission(); + return permissionCheck; + }; +} + +export function checkCanViewQueues(dataSources: DataSources, staffUserExternalId: string | undefined): () => Promise { + return async (): Promise => { + if (!staffUserExternalId) throw new Error('Unauthorized to view queues'); + const staffUser = await dataSources.readonlyDataSource.User.StaffUser.StaffUserReadRepo.getByExternalId(staffUserExternalId); + if (!staffUser?.role?.permissions.techAdminPermissions.canViewQueues) throw new Error('Unauthorized to view queues'); + }; +} + +export function checkCanSendQueueMessages(dataSources: DataSources, staffUserExternalId: string | undefined): () => Promise { + return async (): Promise => { + if (!staffUserExternalId) throw new Error('Unauthorized to send queue messages'); + const staffUser = await dataSources.readonlyDataSource.User.StaffUser.StaffUserReadRepo.getByExternalId(staffUserExternalId); + if (!staffUser?.role?.permissions.techAdminPermissions.canSendQueueMessages) throw new Error('Unauthorized to send queue messages'); + }; +} diff --git a/packages/ocom/application-services/src/contexts/tech-admin/queue/send-message.test.ts b/packages/ocom/application-services/src/contexts/tech-admin/queue/send-message.test.ts new file mode 100644 index 000000000..a06165418 --- /dev/null +++ b/packages/ocom/application-services/src/contexts/tech-admin/queue/send-message.test.ts @@ -0,0 +1,21 @@ +import { describe, expect, it, vi } from 'vitest'; +import { sendMessage } from './send-message.ts'; + +describe('sendMessage', () => { + it('requires an operator reason without adding it to the queue payload', async () => { + const sendMessageToRegisteredQueue = vi.fn().mockResolvedValue(undefined); + const sendQueueMessage = sendMessage({ sendMessageToRegisteredQueue }, vi.fn().mockResolvedValue(undefined)); + + await expect(sendQueueMessage({ queueName: 'community-creation', payload: { communityId: 'community-1' }, reason: 'Replaying a failed request' })).resolves.toBeUndefined(); + expect(sendMessageToRegisteredQueue).toHaveBeenCalledWith( + 'community-creation', + { communityId: 'community-1' }, + { + loggingTags: { source: 'TECH-ADMIN' }, + loggingMetadata: { reason: 'Replaying a failed request' }, + }, + ); + + await expect(sendQueueMessage({ queueName: 'community-creation', payload: {}, reason: ' ' })).rejects.toThrow('Reason is required'); + }); +}); diff --git a/packages/ocom/application-services/src/contexts/tech-admin/queue/send-message.ts b/packages/ocom/application-services/src/contexts/tech-admin/queue/send-message.ts new file mode 100644 index 000000000..a4960c95d --- /dev/null +++ b/packages/ocom/application-services/src/contexts/tech-admin/queue/send-message.ts @@ -0,0 +1,20 @@ +import type { QueueStorageOperations } from '@ocom/service-queue-storage'; + +export interface SendQueueMessageCommand { + queueName: string; + payload: unknown; + reason: string; +} + +export function sendMessage(queueStorageService: Pick, checkPermission: () => Promise): (command: SendQueueMessageCommand) => Promise { + return async ({ queueName, payload, reason }): Promise => { + await checkPermission(); + if (!reason.trim()) { + throw new Error('Reason is required'); + } + await queueStorageService.sendMessageToRegisteredQueue(queueName, payload as object, { + loggingTags: { source: 'TECH-ADMIN' }, + loggingMetadata: { reason }, + }); + }; +} diff --git a/packages/ocom/application-services/src/contexts/user/staff-role/apply-permissions.test.ts b/packages/ocom/application-services/src/contexts/user/staff-role/apply-permissions.test.ts index cae33193f..9e599e8b2 100644 --- a/packages/ocom/application-services/src/contexts/user/staff-role/apply-permissions.test.ts +++ b/packages/ocom/application-services/src/contexts/user/staff-role/apply-permissions.test.ts @@ -37,6 +37,7 @@ function makeStaffRole() { canViewDatabaseExplorer: false, canViewBlobExplorer: false, canViewQueueDashboard: false, + canViewQueues: false, canSendQueueMessages: false, }, }, @@ -163,6 +164,7 @@ describe('applyTechAdminPermissions', () => { canViewDatabaseExplorer: true, canViewBlobExplorer: true, canViewQueueDashboard: true, + canViewQueues: true, canSendQueueMessages: true, }); const tp = staffRole.permissions.techAdminPermissions; @@ -170,6 +172,7 @@ describe('applyTechAdminPermissions', () => { expect(tp.canViewDatabaseExplorer).toBe(true); expect(tp.canViewBlobExplorer).toBe(true); expect(tp.canViewQueueDashboard).toBe(true); + expect(tp.canViewQueues).toBe(true); expect(tp.canSendQueueMessages).toBe(true); }); }); diff --git a/packages/ocom/application-services/src/contexts/user/staff-role/apply-permissions.ts b/packages/ocom/application-services/src/contexts/user/staff-role/apply-permissions.ts index cc4d651e5..3397bb878 100644 --- a/packages/ocom/application-services/src/contexts/user/staff-role/apply-permissions.ts +++ b/packages/ocom/application-services/src/contexts/user/staff-role/apply-permissions.ts @@ -34,6 +34,7 @@ export interface StaffRoleCommandTechAdminPermissions { canViewDatabaseExplorer?: boolean; canViewBlobExplorer?: boolean; canViewQueueDashboard?: boolean; + canViewQueues?: boolean; canSendQueueMessages?: boolean; } @@ -131,6 +132,9 @@ export const applyTechAdminPermissions = (staffRole: Domain.Contexts.User.StaffR if (permissions.canViewQueueDashboard !== undefined) { techAdminPermissions.canViewQueueDashboard = permissions.canViewQueueDashboard; } + if (permissions.canViewQueues !== undefined) { + techAdminPermissions.canViewQueues = permissions.canViewQueues; + } if (permissions.canSendQueueMessages !== undefined) { techAdminPermissions.canSendQueueMessages = permissions.canSendQueueMessages; } diff --git a/packages/ocom/application-services/src/index.ts b/packages/ocom/application-services/src/index.ts index 9e5407886..041aef824 100644 --- a/packages/ocom/application-services/src/index.ts +++ b/packages/ocom/application-services/src/index.ts @@ -2,6 +2,7 @@ import type { ApiContextSpec } from '@ocom/context-spec'; import { Domain } from '@ocom/domain'; import { Community, type CommunityContextApplicationService } from './contexts/community/index.ts'; import { Service, type ServiceContextApplicationService } from './contexts/service/index.ts'; +import { TechAdmin, type TechAdminContextApplicationService } from './contexts/tech-admin/index.ts'; import { User, type UserContextApplicationService } from './contexts/user/index.ts'; export type { CommunityUpdateSettingsCommand } from './contexts/community/index.ts'; @@ -9,6 +10,7 @@ export type { CommunityUpdateSettingsCommand } from './contexts/community/index. export interface ApplicationServices { Community: CommunityContextApplicationService; Service: ServiceContextApplicationService; + TechAdmin: TechAdminContextApplicationService; User: UserContextApplicationService; get verifiedUser(): VerifiedUser | null; } @@ -73,6 +75,7 @@ export const buildApplicationServicesFactory = (context: ApiContextSpec): Applic return { Community: Community(dataSources, blobStorageService, queueStorageService), Service: Service(dataSources), + TechAdmin: TechAdmin(dataSources, queueStorageService, tokenValidationResult?.verifiedJwt.sub), User: User(dataSources), get verifiedUser(): VerifiedUser | null { return { ...tokenValidationResult, hints: hints }; diff --git a/packages/ocom/data-sources-mongoose-models/src/models/role/staff-role.model.ts b/packages/ocom/data-sources-mongoose-models/src/models/role/staff-role.model.ts index 6099c443f..30c059fbe 100644 --- a/packages/ocom/data-sources-mongoose-models/src/models/role/staff-role.model.ts +++ b/packages/ocom/data-sources-mongoose-models/src/models/role/staff-role.model.ts @@ -60,6 +60,7 @@ export interface StaffRoleTechAdminPermissions { canViewDatabaseExplorer: boolean; canViewBlobExplorer: boolean; canViewQueueDashboard: boolean; + canViewQueues: boolean; canSendQueueMessages: boolean; } @@ -155,6 +156,7 @@ const StaffRoleSchema = new Schema, StaffRole>( canViewDatabaseExplorer: { type: Boolean, required: true, default: false }, canViewBlobExplorer: { type: Boolean, required: true, default: false }, canViewQueueDashboard: { type: Boolean, required: true, default: false }, + canViewQueues: { type: Boolean, required: true, default: false }, canSendQueueMessages: { type: Boolean, required: true, default: false }, } as SchemaDefinition, userPermissions: { diff --git a/packages/ocom/domain/src/domain/contexts/user/staff-role/features/staff-role-tech-admin-permissions.feature b/packages/ocom/domain/src/domain/contexts/user/staff-role/features/staff-role-tech-admin-permissions.feature index 33816afec..ca4bfb3f6 100644 --- a/packages/ocom/domain/src/domain/contexts/user/staff-role/features/staff-role-tech-admin-permissions.feature +++ b/packages/ocom/domain/src/domain/contexts/user/staff-role/features/staff-role-tech-admin-permissions.feature @@ -49,6 +49,11 @@ Feature: StaffRoleTechAdminPermissions When I try to set canViewQueueDashboard to true Then a PermissionError should be thrown + Scenario: Changing canViewQueues with manage staff roles permission + Given a StaffRoleTechAdminPermissions entity with permission to manage staff roles + When I set canViewQueues to true + Then the property should be updated to true + Scenario: Changing canSendQueueMessages with manage staff roles permission Given a StaffRoleTechAdminPermissions entity with permission to manage staff roles When I set canSendQueueMessages to true diff --git a/packages/ocom/domain/src/domain/contexts/user/staff-role/staff-role-defaults.test.ts b/packages/ocom/domain/src/domain/contexts/user/staff-role/staff-role-defaults.test.ts index 6696bcf75..7ff1deb9f 100644 --- a/packages/ocom/domain/src/domain/contexts/user/staff-role/staff-role-defaults.test.ts +++ b/packages/ocom/domain/src/domain/contexts/user/staff-role/staff-role-defaults.test.ts @@ -114,6 +114,8 @@ test('applyDefaultSpec sets TechAdmin permissions correctly and marks default', expect(role.permissions.communityPermissions.canManageStaffRolesAndPermissions).toBe(true); expect(role.permissions.financePermissions.canManageFinance).toBe(true); expect(role.permissions.techAdminPermissions.canManageTechAdmin).toBe(true); + expect(role.permissions.techAdminPermissions.canViewQueues).toBe(true); + expect(role.permissions.techAdminPermissions.canSendQueueMessages).toBe(true); expect(role.permissions.userPermissions.canManageUsers).toBe(true); expect(role.isDefault).toBe(true); }); diff --git a/packages/ocom/domain/src/domain/contexts/user/staff-role/staff-role-permissions.ts b/packages/ocom/domain/src/domain/contexts/user/staff-role/staff-role-permissions.ts index 4f0f8677d..c1dbccfe7 100644 --- a/packages/ocom/domain/src/domain/contexts/user/staff-role/staff-role-permissions.ts +++ b/packages/ocom/domain/src/domain/contexts/user/staff-role/staff-role-permissions.ts @@ -116,6 +116,7 @@ export class StaffRolePermissions extends ValueObject canViewDatabaseExplorer: false, canViewBlobExplorer: false, canViewQueueDashboard: false, + canViewQueues: false, canSendQueueMessages: false, }, this.visa, diff --git a/packages/ocom/domain/src/domain/contexts/user/staff-role/staff-role-tech-admin-permissions.test.ts b/packages/ocom/domain/src/domain/contexts/user/staff-role/staff-role-tech-admin-permissions.test.ts index 07bcae739..5c2ac87a5 100644 --- a/packages/ocom/domain/src/domain/contexts/user/staff-role/staff-role-tech-admin-permissions.test.ts +++ b/packages/ocom/domain/src/domain/contexts/user/staff-role/staff-role-tech-admin-permissions.test.ts @@ -21,6 +21,7 @@ function makeProps(overrides = {}) { canViewDatabaseExplorer: false, canViewBlobExplorer: false, canViewQueueDashboard: false, + canViewQueues: false, canSendQueueMessages: false, ...overrides, }; @@ -179,6 +180,19 @@ test.for(feature, ({ Scenario, Background, BeforeEachScenario }) => { }); }); + Scenario('Changing canViewQueues with manage staff roles permission', ({ Given, When, Then }) => { + Given('a StaffRoleTechAdminPermissions entity with permission to manage staff roles', () => { + visa = makeVisa({ canManageStaffRolesAndPermissions: true, isSystemAccount: false }); + entity = new StaffRoleTechAdminPermissions(makeProps(), visa); + }); + When('I set canViewQueues to true', () => { + entity.canViewQueues = true; + }); + Then('the property should be updated to true', () => { + expect(entity.canViewQueues).toBe(true); + }); + }); + Scenario('Changing canSendQueueMessages with manage staff roles permission', ({ Given, When, Then }) => { Given('a StaffRoleTechAdminPermissions entity with permission to manage staff roles', () => { visa = makeVisa({ canManageStaffRolesAndPermissions: true, isSystemAccount: false }); diff --git a/packages/ocom/domain/src/domain/contexts/user/staff-role/staff-role-tech-admin-permissions.ts b/packages/ocom/domain/src/domain/contexts/user/staff-role/staff-role-tech-admin-permissions.ts index 9d225e6c7..9d2b79ce4 100644 --- a/packages/ocom/domain/src/domain/contexts/user/staff-role/staff-role-tech-admin-permissions.ts +++ b/packages/ocom/domain/src/domain/contexts/user/staff-role/staff-role-tech-admin-permissions.ts @@ -8,6 +8,7 @@ interface StaffRoleTechAdminPermissionsSpec { canViewDatabaseExplorer: boolean; canViewBlobExplorer: boolean; canViewQueueDashboard: boolean; + canViewQueues: boolean; canSendQueueMessages: boolean; } @@ -60,6 +61,14 @@ export class StaffRoleTechAdminPermissions extends ValueObject extends AggregateRoot { + it('returns registered queues without resolving their message counts', async () => { + const { context, listQueues, getMessageCount } = createContext(); + const resolver = techAdminResolvers.Query?.techAdminQueues as (parent: object, args: object, context: GraphContext, info: GraphQLResolveInfo) => Promise; + + await expect(resolver({}, {}, context, info)).resolves.toEqual([{ name: 'community-creation', messageCount: null }]); + expect(listQueues).toHaveBeenCalledOnce(); + expect(getMessageCount).not.toHaveBeenCalled(); + }); + + it('resolves a selected queue message count through the application service', async () => { + const { context, getMessageCount } = createContext(); + const resolver = techAdminResolvers.TechAdminQueue?.messageCount as (parent: { name: string }, args: object, context: GraphContext, info: GraphQLResolveInfo) => Promise; + + await expect(resolver({ name: 'community-creation' }, {}, context, info)).resolves.toEqual({ value: 3 }); + expect(getMessageCount).toHaveBeenCalledWith({ queueName: 'community-creation' }); + }); +}); diff --git a/packages/ocom/graphql/src/schema/types/tech-admin.resolvers.ts b/packages/ocom/graphql/src/schema/types/tech-admin.resolvers.ts new file mode 100644 index 000000000..f9614fd0a --- /dev/null +++ b/packages/ocom/graphql/src/schema/types/tech-admin.resolvers.ts @@ -0,0 +1,48 @@ +import type { GraphQLResolveInfo } from 'graphql'; +import type { Resolvers } from '../builder/generated.ts'; +import type { GraphContext } from '../context.ts'; + +const techAdmin: Resolvers = { + Query: { + techAdminQueues: async (_parent, _args, context: GraphContext, _info: GraphQLResolveInfo) => { + ensureAuthenticated(context); + const queues = await context.applicationServices.TechAdmin.Queue.listQueues(); + return queues.map((queue) => ({ ...queue, messageCount: null })); + }, + techAdminQueuePeek: async (_parent, args, context: GraphContext, _info: GraphQLResolveInfo) => { + ensureAuthenticated(context); + return await context.applicationServices.TechAdmin.Queue.peekMessages({ + queueName: args.input.queueName, + maxMessages: args.input.maxMessages ?? 32, + }); + }, + }, + TechAdminQueue: { + messageCount: (queue, _args, context: GraphContext, _info: GraphQLResolveInfo) => { + ensureAuthenticated(context); + return context.applicationServices.TechAdmin.Queue.getMessageCount({ queueName: queue.name }); + }, + }, + Mutation: { + techAdminQueueSend: async (_parent, args, context: GraphContext, _info: GraphQLResolveInfo) => { + if (!context.applicationServices.verifiedUser?.verifiedJwt) { + return { status: { success: false, errorMessage: 'Unauthorized' } }; + } + + try { + await context.applicationServices.TechAdmin.Queue.sendMessage(args.input); + return { status: { success: true } }; + } catch (error) { + return { status: { success: false, errorMessage: error instanceof Error ? error.message : String(error) } }; + } + }, + }, +}; + +export default techAdmin; + +function ensureAuthenticated(context: GraphContext): void { + if (!context.applicationServices.verifiedUser?.verifiedJwt) { + throw new Error('Unauthorized'); + } +} diff --git a/packages/ocom/persistence/src/datasources/domain/user/staff-role/staff-role.domain-adapter.ts b/packages/ocom/persistence/src/datasources/domain/user/staff-role/staff-role.domain-adapter.ts index 1a94a9907..a2bb359b6 100644 --- a/packages/ocom/persistence/src/datasources/domain/user/staff-role/staff-role.domain-adapter.ts +++ b/packages/ocom/persistence/src/datasources/domain/user/staff-role/staff-role.domain-adapter.ts @@ -143,6 +143,7 @@ export class StaffRolePermissionsAdapter implements Domain.Contexts.User.StaffRo canViewDatabaseExplorer: false, canViewBlobExplorer: false, canViewQueueDashboard: false, + canViewQueues: false, canSendQueueMessages: false, }; } @@ -427,6 +428,13 @@ export class StaffRoleTechAdminPermissionsAdapter implements Domain.Contexts.Use this.doc.canViewQueueDashboard = value; } + get canViewQueues(): boolean { + return this.ensureValue(this.doc.canViewQueues); + } + set canViewQueues(value: boolean) { + this.doc.canViewQueues = value; + } + get canSendQueueMessages(): boolean { return this.ensureValue(this.doc.canSendQueueMessages); } diff --git a/packages/ocom/service-queue-storage/README.md b/packages/ocom/service-queue-storage/README.md index a48e83bc3..61c950bee 100644 --- a/packages/ocom/service-queue-storage/README.md +++ b/packages/ocom/service-queue-storage/README.md @@ -24,6 +24,8 @@ Application queue registration for OCOM. This package is the consumer-facing pla `ApiContextSpec` should depend on `QueueStorageOperations`, not `ServiceQueueStorage`. The constructor is for bootstrap; the operations type is for application-service injection. +`QueueStorageOperations` also exposes `sendMessageToRegisteredQueue(queueName, payload)` for controlled operational flows that select a registered physical queue at runtime. It accepts both inbound and outbound queues, validates the payload against the selected schema, and rejects unregistered names. Prefer generated `sendMessageTo...Queue` methods for normal application behavior. + Example: ```ts diff --git a/packages/ocom/service-queue-storage/src/queue-storage.contract.ts b/packages/ocom/service-queue-storage/src/queue-storage.contract.ts index 3e3893278..a045d153d 100644 --- a/packages/ocom/service-queue-storage/src/queue-storage.contract.ts +++ b/packages/ocom/service-queue-storage/src/queue-storage.contract.ts @@ -7,7 +7,7 @@ import type { ServiceQueueStorage } from './registry.ts'; * This intentionally excludes lifecycle concerns such as `startUp`, * `shutDown`, and logging toggles. Application services depend only on the * strongly-typed queue methods generated from the registered queue - * definitions. + * definitions, plus the validated generic sender for registered queues. * * This stays aligned automatically as queues are added or removed from the * registered OCOM queue service. diff --git a/packages/ocom/ui-staff-route-tech-admin/package.json b/packages/ocom/ui-staff-route-tech-admin/package.json index a1622d8fb..0d43a041d 100644 --- a/packages/ocom/ui-staff-route-tech-admin/package.json +++ b/packages/ocom/ui-staff-route-tech-admin/package.json @@ -15,7 +15,10 @@ }, "dependencies": { "@ant-design/icons": "catalog:", + "@apollo/client": "^3.13.9", + "@graphql-typed-document-node/core": "^3.2.0", "@ocom/ui-staff-shared": "workspace:*", + "antd": "catalog:", "react": "catalog:", "react-dom": "catalog:", "react-router-dom": "catalog:" diff --git a/packages/ocom/ui-staff-route-tech-admin/src/components/queue-explorer.container.graphql b/packages/ocom/ui-staff-route-tech-admin/src/components/queue-explorer.container.graphql new file mode 100644 index 000000000..adca9c5ac --- /dev/null +++ b/packages/ocom/ui-staff-route-tech-admin/src/components/queue-explorer.container.graphql @@ -0,0 +1,27 @@ +query TechAdminQueueExplorerContainerQueues { + techAdminQueues { + name + messageCount { + value + errorMessage + } + } +} + +query TechAdminQueueExplorerContainerQueueMessages($input: TechAdminQueuePeekInput!) { + techAdminQueuePeek(input: $input) { + id + popReceipt + payload + dequeueCount + } +} + +mutation TechAdminQueueExplorerContainerSendQueueMessage($input: TechAdminQueueSendInput!) { + techAdminQueueSend(input: $input) { + status { + success + errorMessage + } + } +} diff --git a/packages/ocom/ui-staff-route-tech-admin/src/components/queue-explorer.container.tsx b/packages/ocom/ui-staff-route-tech-admin/src/components/queue-explorer.container.tsx new file mode 100644 index 000000000..951d37b1a --- /dev/null +++ b/packages/ocom/ui-staff-route-tech-admin/src/components/queue-explorer.container.tsx @@ -0,0 +1,83 @@ +import { useLazyQuery, useMutation, useQuery } from '@apollo/client'; +import { App } from 'antd'; +import type React from 'react'; +import { useState } from 'react'; +import { TechAdminQueueExplorerContainerQueueMessagesDocument, TechAdminQueueExplorerContainerQueuesDocument, TechAdminQueueExplorerContainerSendQueueMessageDocument } from '../generated.tsx'; +import { QueueExplorer, type QueueExplorerProps } from './queue-explorer.tsx'; + +interface QueueExplorerContainerProps { + canSendQueueMessages: boolean; +} + +export const QueueExplorerContainer: React.FC = ({ canSendQueueMessages }) => { + const { message } = App.useApp(); + const [selectedQueue, setSelectedQueue] = useState(); + const { + data: queueData, + loading: queueLoading, + error: queueError, + refetch: refetchQueues, + } = useQuery(TechAdminQueueExplorerContainerQueuesDocument, { + fetchPolicy: 'network-only', + }); + const [loadMessages, { data: messageData, loading: messageLoading, error: messageError }] = useLazyQuery(TechAdminQueueExplorerContainerQueueMessagesDocument, { + fetchPolicy: 'network-only', + }); + const [sendQueueMessage, { loading: sendLoading }] = useMutation(TechAdminQueueExplorerContainerSendQueueMessageDocument); + + const refreshMessages = () => { + if (!selectedQueue) return; + void loadMessages({ variables: { input: { queueName: selectedQueue, maxMessages: 32 } } }); + }; + + const handleSelectQueue = (queueName: string) => { + setSelectedQueue(queueName); + void loadMessages({ variables: { input: { queueName, maxMessages: 32 } } }); + }; + + const handleSendMessage = async (payloadText: string, reason: string): Promise => { + if (!selectedQueue) return false; + + let payload: unknown; + try { + payload = JSON.parse(payloadText) as unknown; + } catch { + message.error('Enter a valid JSON payload'); + return false; + } + + try { + const result = await sendQueueMessage({ variables: { input: { queueName: selectedQueue, payload, reason } } }); + const { status } = result.data?.techAdminQueueSend ?? {}; + if (!status?.success) { + message.error(status?.errorMessage ?? 'Unable to send the queue message'); + return false; + } + message.success('Queue message sent'); + refreshMessages(); + return true; + } catch { + message.error('Unable to send the queue message'); + return false; + } + }; + + const errorMessage = queueError?.message ?? messageError?.message; + const props: QueueExplorerProps = { + queues: queueData?.techAdminQueues ?? [], + messages: messageData?.techAdminQueuePeek ?? [], + selectedQueue, + queueLoading, + messageLoading, + sendLoading, + canSendQueueMessages, + errorMessage, + onRefreshQueues: () => void refetchQueues(), + onSelectQueue: handleSelectQueue, + onCloseMessages: () => setSelectedQueue(undefined), + onRefreshMessages: refreshMessages, + onSendMessage: handleSendMessage, + }; + + return ; +}; diff --git a/packages/ocom/ui-staff-route-tech-admin/src/components/queue-explorer.tsx b/packages/ocom/ui-staff-route-tech-admin/src/components/queue-explorer.tsx new file mode 100644 index 000000000..dcc7c8e95 --- /dev/null +++ b/packages/ocom/ui-staff-route-tech-admin/src/components/queue-explorer.tsx @@ -0,0 +1,305 @@ +import { PlusOutlined, ReloadOutlined } from '@ant-design/icons'; +import { Button, Empty, Input, Modal, Space, Spin, Table, type TableColumnsType, Typography } from 'antd'; +import type React from 'react'; +import { useState } from 'react'; + +const { TextArea } = Input; +const { Title, Text } = Typography; + +interface QueueExplorerMessage { + id: string; + popReceipt?: string | null; + payload: unknown; + dequeueCount?: number | null; +} + +interface QueueExplorerQueueRow { + key: string; + regularQueue?: string; + regularQueueCount?: number | null; + regularQueueError?: string | null; + poisonQueue?: string; + poisonQueueCount?: number | null; + poisonQueueError?: string | null; +} + +interface QueueExplorerQueue { + name: string; + messageCount?: { + value?: number | null; + errorMessage?: string | null; + } | null; +} + +export interface QueueExplorerProps { + queues: readonly QueueExplorerQueue[]; + messages: readonly QueueExplorerMessage[]; + selectedQueue?: string | undefined; + queueLoading: boolean; + messageLoading: boolean; + sendLoading: boolean; + canSendQueueMessages: boolean; + errorMessage?: string | undefined; + onRefreshQueues: () => void; + onSelectQueue: (queueName: string) => void; + onCloseMessages: () => void; + onRefreshMessages: () => void; + onSendMessage: (payload: string, reason: string) => Promise; +} + +const formatPayload = (payload: unknown): string => { + if (typeof payload === 'string') return payload; + return JSON.stringify(payload, null, 2); +}; + +const poisonQueueSuffix = '-poison'; + +const renderQueueCount = (messageCount: number | null | undefined, errorMessage: string | null | undefined): React.ReactNode => { + if (errorMessage) return {errorMessage}; + return messageCount ?? 'N/A'; +}; + +const toQueueRows = (queues: readonly QueueExplorerQueue[]): QueueExplorerQueueRow[] => { + const queuesByName = new Map(queues.map((queue) => [queue.name, queue])); + const regularQueues = queues.filter(({ name }) => !name.endsWith(poisonQueueSuffix)); + const orphanedPoisonQueues = queues.filter(({ name }) => name.endsWith(poisonQueueSuffix) && !queuesByName.has(name.slice(0, -poisonQueueSuffix.length))); + + return [ + ...regularQueues.map((regularQueue) => { + const poisonQueue = queuesByName.get(`${regularQueue.name}${poisonQueueSuffix}`); + + return { + key: regularQueue.name, + regularQueue: regularQueue.name, + ...(regularQueue.messageCount?.value === undefined ? {} : { regularQueueCount: regularQueue.messageCount.value }), + ...(regularQueue.messageCount?.errorMessage === undefined ? {} : { regularQueueError: regularQueue.messageCount.errorMessage }), + ...(poisonQueue === undefined + ? {} + : { + poisonQueue: poisonQueue.name, + ...(poisonQueue.messageCount?.value === undefined ? {} : { poisonQueueCount: poisonQueue.messageCount.value }), + ...(poisonQueue.messageCount?.errorMessage === undefined ? {} : { poisonQueueError: poisonQueue.messageCount.errorMessage }), + }), + }; + }), + ...orphanedPoisonQueues.map((poisonQueue) => ({ + key: poisonQueue.name, + poisonQueue: poisonQueue.name, + ...(poisonQueue.messageCount?.value === undefined ? {} : { poisonQueueCount: poisonQueue.messageCount.value }), + ...(poisonQueue.messageCount?.errorMessage === undefined ? {} : { poisonQueueError: poisonQueue.messageCount.errorMessage }), + })), + ]; +}; + +export const QueueExplorer: React.FC = ({ + queues, + messages, + selectedQueue, + queueLoading, + messageLoading, + sendLoading, + canSendQueueMessages, + errorMessage, + onRefreshQueues, + onSelectQueue, + onCloseMessages, + onRefreshMessages, + onSendMessage, +}) => { + const [isSendModalOpen, setIsSendModalOpen] = useState(false); + const [payload, setPayload] = useState(''); + const [reason, setReason] = useState(''); + const queueRows = toQueueRows(queues); + const isPoisonQueue = selectedQueue?.endsWith(poisonQueueSuffix) === true; + + const queueColumns: TableColumnsType = [ + { + title: 'Regular Queue', + dataIndex: 'regularQueue', + key: 'regularQueue', + render: (queueName: string | undefined) => + queueName && ( + + ), + sorter: (left, right) => (left.regularQueue ?? '').localeCompare(right.regularQueue ?? ''), + defaultSortOrder: 'ascend', + }, + { + title: 'Message Count', + dataIndex: 'regularQueueCount', + key: 'regularQueueCount', + render: (messageCount: number | null | undefined, record) => record.regularQueue && renderQueueCount(messageCount, record.regularQueueError), + sorter: (left, right) => (left.regularQueueCount ?? -1) - (right.regularQueueCount ?? -1), + }, + { + title: 'Poison Queue', + dataIndex: 'poisonQueue', + key: 'poisonQueue', + render: (queueName: string | undefined) => + queueName && ( + + ), + sorter: (left, right) => (left.poisonQueue ?? '').localeCompare(right.poisonQueue ?? ''), + }, + { + title: 'Message Count', + dataIndex: 'poisonQueueCount', + key: 'poisonQueueCount', + render: (messageCount: number | null | undefined, record) => record.poisonQueue && renderQueueCount(messageCount, record.poisonQueueError), + sorter: (left, right) => (left.poisonQueueCount ?? -1) - (right.poisonQueueCount ?? -1), + }, + ]; + + const messageColumns: TableColumnsType = [ + { + title: 'Message ID', + dataIndex: 'id', + key: 'id', + sorter: (left, right) => left.id.localeCompare(right.id), + }, + { + title: 'Payload', + dataIndex: 'payload', + key: 'payload', + render: (messagePayload: unknown) =>
{formatPayload(messagePayload)}
, + }, + { + title: 'Dequeue Count', + dataIndex: 'dequeueCount', + key: 'dequeueCount', + width: 140, + render: (dequeueCount: number | null | undefined) => dequeueCount ?? 'N/A', + sorter: (left, right) => (left.dequeueCount ?? 0) - (right.dequeueCount ?? 0), + }, + ]; + + const closeSendModal = () => { + setPayload(''); + setReason(''); + setIsSendModalOpen(false); + }; + + const handleSendMessage = async () => { + if (await onSendMessage(payload, reason)) { + closeSendModal(); + } + }; + + return ( + <> + +
+ Queue Explorer + Inspect registered queues and their most recent messages. +
+ + {errorMessage && {errorMessage}} + }} + /> + + + + + Results are limited to the most recent 32 messages. + + + {!isPoisonQueue && ( + + )} + + {messageLoading ? ( +
+ +
+ ) : ( +
}} + /> + )} + + + + void handleSendMessage()} + okText="Send" + okButtonProps={{ loading: sendLoading, disabled: payload.trim().length === 0 || reason.trim().length === 0 }} + destroyOnHidden + > +