import type { Either } from '@lokalise/node-core' import type { BarrierResult, PreHandlingOutputs, Prehandler } from '@message-queue-toolkit/core' import { MessageHandlerConfigBuilder } from '@message-queue-toolkit/core' import type { SQSConsumerDependencies, SQSConsumerOptions, } from '../../lib/sqs/AbstractSqsConsumer.ts' import { AbstractSqsConsumer } from '../../lib/sqs/AbstractSqsConsumer.ts' import type { PERMISSIONS_ADD_MESSAGE_TYPE, PERMISSIONS_REMOVE_MESSAGE_TYPE, } from './userConsumerSchemas.ts' import { PERMISSIONS_ADD_MESSAGE_SCHEMA, PERMISSIONS_REMOVE_MESSAGE_SCHEMA, } from './userConsumerSchemas.ts' export type SupportedMessages = PERMISSIONS_ADD_MESSAGE_TYPE | PERMISSIONS_REMOVE_MESSAGE_TYPE type SqsPermissionConsumerOptions = Pick< SQSConsumerOptions, | 'creationConfig' | 'locatorConfig' | 'logMessages' | 'deletionConfig' | 'deadLetterQueue' | 'consumerOverrides' | 'consumerPollingWaitTimeSeconds' | 'maxRetryDuration' | 'payloadStoreConfig' | 'messageDeduplicationConfig' | 'enableConsumerDeduplication' | 'codecs' | 'disableCodecAutoDetection' > & { addPreHandlerBarrier?: ( message: SupportedMessages, _executionContext: ExecutionContext, preHandlerOutput: PrehandlerOutput, ) => Promise> removeHandlerOverride?: ( _message: SupportedMessages, context: ExecutionContext, preHandlingOutputs: PreHandlingOutputs, ) => Promise> addHandlerOverride?: ( message: SupportedMessages, context: ExecutionContext, preHandlingOutputs: PreHandlingOutputs, ) => Promise> removePreHandlers?: Prehandler[] concurrentConsumersAmount?: number } type ExecutionContext = { incrementAmount: number } type PrehandlerOutput = { messageId: string } export class SqsPermissionConsumer extends AbstractSqsConsumer< SupportedMessages, ExecutionContext, PrehandlerOutput > { public addCounter = 0 public removeCounter = 0 public processedMessagesIds: Set = new Set() public static readonly QUEUE_NAME = 'user_permissions_multi' constructor( dependencies: SQSConsumerDependencies, options: SqsPermissionConsumerOptions = { creationConfig: { queue: { QueueName: SqsPermissionConsumer.QUEUE_NAME, }, }, }, ) { const defaultRemoveHandler = ( _message: SupportedMessages, context: ExecutionContext, _preHandlingOutputs: PreHandlingOutputs, ): Promise> => { this.removeCounter += context.incrementAmount return Promise.resolve({ result: 'success', }) } const defaultAddHandler = ( message: SupportedMessages, context: ExecutionContext, barrierOutput: PreHandlingOutputs, ): Promise> => { if (options.addPreHandlerBarrier && !barrierOutput) { return Promise.resolve({ error: 'retryLater' }) } this.addCounter += context.incrementAmount this.processedMessagesIds.add(message.id) return Promise.resolve({ result: 'success' }) } super( dependencies, { ...(options.locatorConfig ? { locatorConfig: options.locatorConfig } : { creationConfig: options.creationConfig ?? { queue: { QueueName: SqsPermissionConsumer.QUEUE_NAME }, }, }), logMessages: options.logMessages, deletionConfig: options.deletionConfig ?? { deleteIfExists: true, }, deadLetterQueue: options.deadLetterQueue, messageTypeResolver: { messageTypePath: 'messageType' }, handlerSpy: true, consumerOverrides: options.consumerOverrides ?? { terminateVisibilityTimeout: true, // this allows to retry failed messages immediately }, consumerPollingWaitTimeSeconds: options.consumerPollingWaitTimeSeconds ?? 0, concurrentConsumersAmount: options.concurrentConsumersAmount, maxRetryDuration: options.maxRetryDuration, payloadStoreConfig: options.payloadStoreConfig, messageDeduplicationConfig: options.messageDeduplicationConfig, enableConsumerDeduplication: options.enableConsumerDeduplication, codecs: options.codecs, disableCodecAutoDetection: options.disableCodecAutoDetection, messageDeduplicationIdField: 'deduplicationId', messageDeduplicationOptionsField: 'deduplicationOptions', handlers: new MessageHandlerConfigBuilder< SupportedMessages, ExecutionContext, PrehandlerOutput >() .addConfig( PERMISSIONS_ADD_MESSAGE_SCHEMA, options.addHandlerOverride ?? defaultAddHandler, { preHandlerBarrier: options.addPreHandlerBarrier, preHandlers: [ (message, _context, preHandlerOutput, next) => { preHandlerOutput.messageId = message.id next({ result: 'success', }) }, ], }, ) .addConfig( PERMISSIONS_REMOVE_MESSAGE_SCHEMA, options.removeHandlerOverride ?? defaultRemoveHandler, { preHandlers: options.removePreHandlers, }, ) .build(), }, { incrementAmount: 1, }, ) } public get queueProps() { return { name: this.queue.name, url: this.queue.url, arn: this.queue.arn, } } public get dlqUrl() { return this.deadLetterQueue?.url ?? '' } }