diff --git a/README.md b/README.md index 5215b2322..6f440adfd 100644 --- a/README.md +++ b/README.md @@ -110,7 +110,7 @@ Multi-schema consumers support multiple message types via handler configs. They * `queueName`; (for SNS publishers this is a misnomer which actually refers to a topic name) * `locatorConfig` - configuration for resolving existing queue and/or topic. Should not be specified together with the `creationConfig`. * `creationConfig` - configuration for queue and/or topic to create, if one does not exist. Should not be specified together with the `locatorConfig`. - * `subscriptionConfig` - SNS SQS consumer only - configuration for SNS -> SQS subscription to create, if one doesn't exist. + * `subscriptionConfig` - SNS SQS consumer only - configuration for SNS -> SQS subscription to create, if one doesn't exist. With `locateOnly: true`, the subscription is only located (never created) and its `Attributes` (e.g. `FilterPolicy`) are applied to it, which requires the topic and the queue to be located too (see the [SNS README](packages/sns/README.md#resource-resolution)). * `policyConfig` - SQS only - configuration for queue access policies (see [SQS Policy Configuration](#sqs-policy-configuration) for more information); * `deletionConfig` - automatic cleanup of resources; * `consumerOverrides` – available only for SQS consumers; @@ -337,7 +337,6 @@ const consumer = new MySnsSqsConsumer(dependencies, { locatorConfig: { topicArn: 'arn:aws:sns:...', queueUrl: 'https://sqs...', - subscriptionArn: '...', // Enable eventual consistency mode startupResourcePolling: { enabled: true, // Enable polling for resource availability @@ -474,7 +473,6 @@ const result = await initSnsSqs( { topicArn: '...', queueUrl: '...', - subscriptionArn: '...', startupResourcePolling: { enabled: true, timeoutMs: 5 * 60 * 1000, @@ -485,7 +483,7 @@ const result = await initSnsSqs( undefined, { onResourcesReady: ({ topicArn, queueUrl }) => { - // Called only when BOTH topic and queue are available + // Called only when the topic, the queue and the subscription are all available console.log(`Resources ready: topic=${topicArn}, queue=${queueUrl}`) }, }, diff --git a/packages/sns/README.md b/packages/sns/README.md index d175c1566..e3f9acb15 100644 --- a/packages/sns/README.md +++ b/packages/sns/README.md @@ -488,9 +488,9 @@ When using `locatorConfig`, you connect to an existing topic without creating it // or // queueName: 'my-queue', - // Optional: Existing subscription ARN. When omitted and no `subscriptionConfig` is given, the subscription of - // the queue to the topic is looked up instead (see Resource Resolution) - subscriptionArn: 'arn:aws:sns:us-east-1:123456789012:my-topic:uuid', + // Deprecated: the subscription is located from the topic and the queue (see Resource Resolution), so its ARN is + // no longer needed. It will be removed in the next major version + // subscriptionArn: 'arn:aws:sns:us-east-1:123456789012:my-topic:uuid', }, } ``` @@ -500,16 +500,49 @@ When using `locatorConfig`, you connect to an existing topic without creating it Consumers resolve the topic, the queue and the subscription independently, so each of them can either be managed by your application or by external tooling (e.g. Terraform): -| Resource | Located when | Created when | -|--------------|------------------------------------------------|----------------------------------------| -| Topic | `locatorConfig.topicArn` or `topicName` is set | otherwise, from `creationConfig.topic` | -| Queue | `locatorConfig.queueUrl` or `queueName` is set | otherwise, from `creationConfig.queue` | -| Subscription | `subscriptionConfig` is not set | `subscriptionConfig` is set | +| Resource | Located when | Created when | +|--------------|------------------------------------------------|-------------------------------------------| +| Topic | `locatorConfig.topicArn` or `topicName` is set | otherwise, from `creationConfig.topic` | +| Queue | `locatorConfig.queueUrl` or `queueName` is set | otherwise, from `creationConfig.queue` | +| Subscription | no `subscriptionConfig`, or `locateOnly: true` | `subscriptionConfig` without `locateOnly` | + +A subscription that already exists is reused. Its attributes are read first and compared with the configured ones +(JSON policies structurally), so a subscription that is already up to date receives no writes at all. Differing +attributes are written one by one when `subscriptionConfig.updateAttributesIfExists` is enabled, and an error is thrown +otherwise. When both a locator and a creation config are given for the same resource, the locator takes precedence and +the creation config is ignored. + +`subscriptionConfig.managedAttributes` lists the subscription attributes the application owns. It defaults to all of +`FilterPolicy`, `FilterPolicyScope`, `RawMessageDelivery` and `RedrivePolicy`. A managed attribute that is missing from +`Attributes` is reset on the existing subscription: `FilterPolicy` and `RedrivePolicy` are removed, `FilterPolicyScope` +goes back to `MessageAttributes` and `RawMessageDelivery` to `false`. Attributes left out of `managedAttributes` are +neither checked nor written, so external tooling can own them, and setting one of them in `Attributes` is a +configuration error. -A subscription that already exists is reused, and its attributes are updated when they differ and -`subscriptionConfig.updateAttributesIfExists` is enabled. When both a locator and a creation config are given for the -same resource, the locator takes precedence and the creation config is ignored. When `locatorConfig.subscriptionArn` -is set, every resource is located and both `creationConfig` and `subscriptionConfig` are ignored. +```typescript +// Filter policy owned by the application, everything else by Terraform +subscriptionConfig: { + locateOnly: true, + managedAttributes: ['FilterPolicy', 'FilterPolicyScope'], + Attributes: { FilterPolicy: JSON.stringify({ type: ['entity.created'] }) }, +} +``` + +With `subscriptionDeadLetterQueue.reuseConsumerDeadLetterQueue`, `RedrivePolicy` is set from the consumer DLQ and is +excluded from the managed attributes. + +With `subscriptionConfig: { locateOnly: true, Attributes }` the subscription is located, never created, and its +managed attributes are applied to it on startup when they differ from the current ones. Managed attributes missing +from `Attributes` are reset as well, so set `managedAttributes` to the ones the application owns (usually +`['FilterPolicy', 'FilterPolicyScope']`), otherwise attributes set by the external tooling, such as `RedrivePolicy`, +are removed. This lets external tooling own the subscription while the application keeps its filter policy in sync +with the consumer handlers. A locate-only subscription requires both the topic and the queue to be located, and none +of the located resources are deleted when `deletionConfig` is set. + +> **Deprecated**: `locatorConfig.subscriptionArn` is no longer needed, as the subscription is located from the topic +> and the queue. When it is set, every resource is located and `creationConfig` and `subscriptionConfig` are ignored, +> except for a locate-only `subscriptionConfig`, whose attributes are still applied. It will be removed in the next +> major version. ```typescript // Application manages everything @@ -535,14 +568,26 @@ is set, every resource is located and both `creationConfig` and `subscriptionCon { locatorConfig: { topicName: 'my-topic', queueName: 'my-queue' }, } + +// Everything managed externally, the application only manages the subscription filter policy +{ + locatorConfig: { topicName: 'my-topic', queueName: 'my-queue' }, + subscriptionConfig: { + locateOnly: true, + managedAttributes: ['FilterPolicy', 'FilterPolicyScope'], + Attributes: { FilterPolicy: JSON.stringify({ messageType: ['user.created'] }) }, + }, +} ``` Some things to keep in mind when resources are managed externally: - **Queue policy**: when the queue is located, the application does not set its policy. The external tooling must allow the topic to send messages to the queue, otherwise the subscription exists but messages are never delivered. -- **Filter policy**: when the subscription is located, its filter policy is not derived from the consumer handlers. - The external tooling must keep it in sync with the message types the consumer handles. +- **Filter policy**: when the subscription is located without `subscriptionConfig`, its filter policy is not derived + from the consumer handlers, so the external tooling must keep it in sync with the message types the consumer + handles. Use `subscriptionConfig.locateOnly` with `Attributes.FilterPolicy` to keep managing it from the application + (requires `sns:SetSubscriptionAttributes`, and the external tooling should ignore changes to it). - **Permissions**: locating a subscription requires `sns:ListSubscriptionsByTopic`. Subscriptions pending confirmation (e.g. cross-account) are not considered until they are confirmed. - A missing located resource makes `init()` fail, unless startup resource polling is enabled (see below). @@ -556,7 +601,6 @@ When your SNS topic or SQS queue may not exist at startup (e.g., created by anot locatorConfig: { topicName: 'my-topic', queueUrl: 'https://sqs.us-east-1.amazonaws.com/123456789012/my-queue', - subscriptionArn: 'arn:aws:sns:us-east-1:123456789012:my-topic:uuid', // Enable startup resource polling startupResourcePolling: { @@ -582,7 +626,6 @@ const consumer = new MyConsumer(deps, { locatorConfig: { topicName: 'my-topic', queueUrl: 'https://sqs...', - subscriptionArn: 'arn:aws:sns:...', startupResourcePolling: { enabled: true, pollingIntervalMs: 5000, @@ -606,7 +649,6 @@ const consumer = new MyConsumer(deps, { locatorConfig: { topicName: 'my-topic', queueUrl: 'https://sqs...', - subscriptionArn: 'arn:aws:sns:...', startupResourcePolling: { enabled: true, pollingIntervalMs: 5000, @@ -680,13 +722,13 @@ if (result.resourcesReady) { #### Subscription Creation Mode -When you want to **create** a subscription (no existing `subscriptionArn`), startup resource polling will wait for the topic to exist before attempting to subscribe: +When you want to **create** a subscription (`subscriptionConfig` without `locateOnly`), startup resource polling will +wait for the topic to exist before attempting to subscribe: ```typescript const consumer = new MyConsumer(deps, { locatorConfig: { topicName: 'my-topic', // Topic created by another service - // No subscriptionArn - we'll create the subscription startupResourcePolling: { enabled: true, pollingIntervalMs: 5000, @@ -696,6 +738,7 @@ const consumer = new MyConsumer(deps, { creationConfig: { queue: { QueueName: 'my-consumer-queue' }, // Queue will be created }, + subscriptionConfig: { updateAttributesIfExists: true }, // Subscription will be created }) // This will: @@ -708,8 +751,8 @@ await consumer.start() #### Subscription Locate Mode -When the subscription is managed externally (no `subscriptionConfig`), startup resource polling will also wait for the -subscription of the queue to the topic to exist: +When the subscription is managed externally (no `subscriptionConfig`, or `subscriptionConfig.locateOnly`), startup +resource polling will also wait for the subscription of the queue to the topic to exist: ```typescript const consumer = new MyConsumer(deps, { @@ -722,13 +765,19 @@ const consumer = new MyConsumer(deps, { timeoutMs: 60000, }, }, - // No subscriptionConfig - the subscription is located, never created + // No subscriptionConfig - the subscription is located, never created. Alternatively, keep managing its filter policy: + // subscriptionConfig: { + // locateOnly: true, + // managedAttributes: ['FilterPolicy', 'FilterPolicyScope'], + // Attributes: { FilterPolicy: '...' }, + // }, }) // This will: // 1. Poll until the topic and the queue exist // 2. Poll until the queue is subscribed to the topic -// 3. Start consuming +// 3. Apply the subscription attributes, if a locate-only subscriptionConfig is given +// 4. Start consuming await consumer.start() ``` @@ -799,12 +848,12 @@ SNS consumers use the same options as SQS consumers, plus SNS-specific subscript locatorConfig: { topicArn: 'arn:aws:sns:...', queueUrl: 'https://sqs...', - subscriptionArn: 'arn:aws:sns:...', }, // SNS-Specific - Subscription Configuration subscriptionConfig: { updateAttributesIfExists: false, // Update subscription attributes if exists + managedAttributes: ['FilterPolicy', 'FilterPolicyScope'], // Defaults to all supported attributes // Optional: Message filtering filterPolicy: { @@ -819,6 +868,12 @@ SNS consumers use the same options as SQS consumers, plus SNS-specific subscript deadLetterTargetArn: 'arn:aws:sqs:us-east-1:123456789012:my-dlq', }, }, + // or, for a subscription managed externally (requires topic and queue locators) + // subscriptionConfig: { + // locateOnly: true, + // managedAttributes: ['FilterPolicy', 'FilterPolicyScope'], + // Attributes: { FilterPolicy: '...' }, + // }, // Optional - FIFO Configuration fifoQueue: false, @@ -1382,14 +1437,29 @@ type SNSSQSConsumerDependencies = SNSDependencies & SQSDependencies & { } // Subscription options -type SNSSubscriptionOptions = { - updateAttributesIfExists?: boolean - filterPolicy?: Record - rawMessageDelivery?: boolean - redrivePolicy?: { - deadLetterTargetArn: string - } -} +type SNSSubscriptionOptions = + // Creation: the subscription is created, or updated if it already exists + | { + updateAttributesIfExists?: boolean + managedAttributes?: SubscriptionManagedAttributeName[] + filterPolicy?: Record + rawMessageDelivery?: boolean + redrivePolicy?: { + deadLetterTargetArn: string + } + } + // Locate only: the subscription is located, never created, and the given attributes are applied to it + | { + locateOnly: true + managedAttributes?: SubscriptionManagedAttributeName[] + Attributes?: Record + } + +type SubscriptionManagedAttributeName = + | 'FilterPolicy' + | 'FilterPolicyScope' + | 'RawMessageDelivery' + | 'RedrivePolicy' ``` ### Utility Functions diff --git a/packages/sns/lib/index.ts b/packages/sns/lib/index.ts index dce81d925..784d01956 100644 --- a/packages/sns/lib/index.ts +++ b/packages/sns/lib/index.ts @@ -34,7 +34,12 @@ export { } from './utils/snsInitter.ts' export { deserializeSNSMessage } from './utils/snsMessageDeserializer.ts' export { readSnsMessage } from './utils/snsMessageReader.ts' -export { type SNSSubscriptionOptions, subscribeToTopic } from './utils/snsSubscriber.ts' +export { + type SNSSubscriptionOptions, + SUBSCRIPTION_MANAGED_ATTRIBUTE_NAMES, + type SubscriptionManagedAttributeName, + subscribeToTopic, +} from './utils/snsSubscriber.ts' export { assertTopic, calculateOutgoingMessageSize, diff --git a/packages/sns/lib/sns/AbstractSnsSqsConsumer.ts b/packages/sns/lib/sns/AbstractSnsSqsConsumer.ts index 3a75a5f4b..c1710f0a4 100644 --- a/packages/sns/lib/sns/AbstractSnsSqsConsumer.ts +++ b/packages/sns/lib/sns/AbstractSnsSqsConsumer.ts @@ -1,5 +1,4 @@ import type { SNSClient } from '@aws-sdk/client-sns' -import { SetSubscriptionAttributesCommand } from '@aws-sdk/client-sns' import type { STSClient } from '@aws-sdk/client-sts' import { type Either, InternalError } from '@lokalise/node-core' import type { @@ -17,7 +16,11 @@ import type { import { AbstractSqsConsumer, deleteSqs } from '@message-queue-toolkit/sqs' import { deleteSnsSqs, initSnsSqs } from '../utils/snsInitter.ts' import { readSnsMessage } from '../utils/snsMessageReader.ts' -import type { SNSSubscriptionOptions } from '../utils/snsSubscriber.ts' +import { + type SNSSubscriptionOptions, + SUBSCRIPTION_MANAGED_ATTRIBUTE_NAMES, + setSubscriptionAttributes, +} from '../utils/snsSubscriber.ts' import type { SNSCreationConfig, SNSOptions, SNSTopicLocatorType } from './AbstractSnsService.ts' export type SNSSQSConsumerDependencies = SQSConsumerDependencies & { @@ -28,6 +31,11 @@ export type SNSSQSCreationConfig = Omit & SNS export type SNSSQSQueueLocatorType = Partial & SNSTopicLocatorType & { + /** + * @deprecated The subscription is located from the topic and queue locators, so its ARN is no longer needed. + * To avoid creating it, omit `subscriptionConfig` or use `subscriptionConfig.locateOnly`. It will be removed in + * the next major version. + */ subscriptionArn?: string } @@ -115,9 +123,11 @@ export abstract class AbstractSnsSqsConsumer< ) { super(dependencies, { ...options }, executionContext) - this.subscriptionConfig = options.subscriptionConfig this.reuseConsumerDeadLetterQueueForSubscription = !!options.subscriptionDeadLetterQueue?.reuseConsumerDeadLetterQueue + this.subscriptionConfig = this.reuseConsumerDeadLetterQueueForSubscription + ? withoutManagedRedrivePolicy(options.subscriptionConfig) + : options.subscriptionConfig if (this.reuseConsumerDeadLetterQueueForSubscription && !options.deadLetterQueue) { throw new InternalError({ @@ -132,7 +142,12 @@ export abstract class AbstractSnsSqsConsumer< } override async init(): Promise { - if (this.deletionConfig && this.creationConfig && this.subscriptionConfig) { + if ( + this.deletionConfig && + this.creationConfig && + this.subscriptionConfig && + !this.subscriptionConfig.locateOnly + ) { await deleteSnsSqs( this.sqsClient, this.snsClient, @@ -144,7 +159,7 @@ export abstract class AbstractSnsSqsConsumer< undefined, this.locatorConfig, ) - } else if (this.deletionConfig && this.creationConfig) { + } else if (this.deletionConfig && this.creationConfig && !this.subscriptionConfig?.locateOnly) { await deleteSqs(this.sqsClient, this.deletionConfig, this.creationConfig) } @@ -252,12 +267,11 @@ export abstract class AbstractSnsSqsConsumer< const dlq = this.deadLetterQueue if (!dlq) return - await this.snsClient.send( - new SetSubscriptionAttributesCommand({ - SubscriptionArn: this.subscription.subscriptionArn, - AttributeName: 'RedrivePolicy', - AttributeValue: JSON.stringify({ deadLetterTargetArn: dlq.arn }), - }), + await setSubscriptionAttributes( + this.snsClient, + this.subscription.subscriptionArn, + { RedrivePolicy: JSON.stringify({ deadLetterTargetArn: dlq.arn }) }, + ['RedrivePolicy'], ) } @@ -318,3 +332,30 @@ export abstract class AbstractSnsSqsConsumer< return this._messageSchemaContainer.resolveSchema(messagePayload) } } + +/** + * With `subscriptionDeadLetterQueue.reuseConsumerDeadLetterQueue`, the redrive policy is set once the + * DLQ is resolved, so the subscription config must not manage (and reset) it on its own. + */ +function withoutManagedRedrivePolicy( + subscriptionConfig: SNSSubscriptionOptions | undefined, +): SNSSubscriptionOptions | undefined { + if (!subscriptionConfig) return undefined + if ( + subscriptionConfig.managedAttributes?.includes('RedrivePolicy') || + subscriptionConfig.Attributes?.RedrivePolicy + ) { + throw new InternalError({ + errorCode: 'invalid_subscription_dlq_configuration', + message: + 'subscriptionDeadLetterQueue.reuseConsumerDeadLetterQueue sets the subscription RedrivePolicy, so subscriptionConfig must not manage it', + }) + } + + return { + ...subscriptionConfig, + managedAttributes: ( + subscriptionConfig.managedAttributes ?? SUBSCRIPTION_MANAGED_ATTRIBUTE_NAMES + ).filter((name) => name !== 'RedrivePolicy'), + } +} diff --git a/packages/sns/lib/utils/snsInitter.spec.ts b/packages/sns/lib/utils/snsInitter.spec.ts index a6bc9778b..0ecce14a1 100644 --- a/packages/sns/lib/utils/snsInitter.spec.ts +++ b/packages/sns/lib/utils/snsInitter.spec.ts @@ -1,8 +1,10 @@ import { setTimeout } from 'node:timers/promises' import { CreateTopicCommand, + GetSubscriptionAttributesCommand, GetTopicAttributesCommand, ListSubscriptionsByTopicCommand, + SetSubscriptionAttributesCommand, type SNSClient, SubscribeCommand, } from '@aws-sdk/client-sns' @@ -182,6 +184,140 @@ describe('snsInitter', () => { }) }) + describe('locate-only subscription', () => { + const filterPolicy = JSON.stringify({ type: ['entity.created'] }) + const getFilterPolicy = async (subscriptionArn: string) => { + const { Attributes } = await snsClient.send( + new GetSubscriptionAttributesCommand({ SubscriptionArn: subscriptionArn }), + ) + return Attributes?.FilterPolicy + } + + it('locates existing subscription and applies its attributes without subscribing', async () => { + const topicArn = await testAdmin.createTopic(topicName) + const { queueArn } = await testAdmin.createQueue(queueName) + const subscriptionArn = await testAdmin.createSubscription(topicArn, queueArn) + const snsSpy = vi.spyOn(snsClient, 'send') + + const result = await initSnsSqs( + sqsClient, + snsClient, + stsClient, + { topicName, queueName }, + undefined, + { locateOnly: true, Attributes: { FilterPolicy: filterPolicy } }, + ) + + expect(result).toEqual({ topicArn, queueUrl, queueArn, queueName, subscriptionArn }) + expect(snsSpy).not.toHaveBeenCalledWith(expect.any(SubscribeCommand)) + expect(await getFilterPolicy(subscriptionArn)).toBe(filterPolicy) + }) + + it('does not update the subscription when no attributes are provided', async () => { + const topicArn = await testAdmin.createTopic(topicName) + const { queueArn } = await testAdmin.createQueue(queueName) + const subscriptionArn = await testAdmin.createSubscription(topicArn, queueArn) + const snsSpy = vi.spyOn(snsClient, 'send') + + const result = await initSnsSqs( + sqsClient, + snsClient, + stsClient, + { topicName, queueName }, + undefined, + { locateOnly: true }, + ) + + expect(result?.subscriptionArn).toBe(subscriptionArn) + expect(snsSpy).not.toHaveBeenCalledWith(expect.any(SubscribeCommand)) + expect(snsSpy).not.toHaveBeenCalledWith(expect.any(SetSubscriptionAttributesCommand)) + }) + + it('applies attributes to the subscription when its ARN is provided', async () => { + const topicArn = await testAdmin.createTopic(topicName) + const { queueArn } = await testAdmin.createQueue(queueName) + const subscriptionArn = await testAdmin.createSubscription(topicArn, queueArn) + const snsSpy = vi.spyOn(snsClient, 'send') + + const result = await initSnsSqs( + sqsClient, + snsClient, + stsClient, + { topicName, queueName, subscriptionArn }, + undefined, + { locateOnly: true, Attributes: { FilterPolicy: filterPolicy } }, + ) + + expect(result?.subscriptionArn).toBe(subscriptionArn) + // The subscription is not looked up, as its ARN is provided + expect(snsSpy).not.toHaveBeenCalledWith(expect.any(ListSubscriptionsByTopicCommand)) + expect(snsSpy).not.toHaveBeenCalledWith(expect.any(SubscribeCommand)) + expect(await getFilterPolicy(subscriptionArn)).toBe(filterPolicy) + }) + + it('throws without creating the subscription when it does not exist', async () => { + const topicArn = await testAdmin.createTopic(topicName) + const { queueArn } = await testAdmin.createQueue(queueName) + + await expect( + initSnsSqs(sqsClient, snsClient, stsClient, { topicName, queueName }, undefined, { + locateOnly: true, + Attributes: { FilterPolicy: filterPolicy }, + }), + ).rejects.toThrow(/Subscription of queue .* to topic .* does not exist/) + expect(await findSubscriptionByTopicAndQueue(snsClient, topicArn, queueArn)).toBeUndefined() + }) + + it('throws when the located topic does not exist', async () => { + await testAdmin.createQueue(queueName) + + await expect( + initSnsSqs(sqsClient, snsClient, stsClient, { topicName, queueName }, undefined, { + locateOnly: true, + }), + ).rejects.toThrow(/Topic with topicArn .* does not exist/) + }) + + it('applies attributes once the subscription becomes available in non-blocking mode', async () => { + const topicArn = await testAdmin.createTopic(topicName) + const { queueArn } = await testAdmin.createQueue(queueName) + const onResourcesReady = vi.fn() + + const result = await initSnsSqs( + sqsClient, + snsClient, + stsClient, + { + topicName, + queueName, + startupResourcePolling: { + enabled: true, + pollingIntervalMs: 50, + timeoutMs: 5000, + nonBlocking: true, + }, + }, + undefined, + { locateOnly: true, Attributes: { FilterPolicy: filterPolicy } }, + { onResourcesReady }, + ) + + expect(result).toBeUndefined() + + const subscriptionArn = await testAdmin.createSubscription(topicArn, queueArn) + + await waitAndRetry(() => onResourcesReady.mock.calls.length > 0, 50, 40) + expect(onResourcesReady).toHaveBeenCalledWith({ + topicArn, + queueUrl, + queueArn, + queueName, + subscriptionArn, + }) + expect(await getFilterPolicy(subscriptionArn)).toBe(filterPolicy) + }) + }) + describe('config validation', () => { it('throws when neither topic creation config nor topic locator is provided', async () => { await expect( @@ -229,6 +365,58 @@ describe('snsInitter', () => { 'If creationConfig.queue is specified, subscriptionConfig is mandatory, as the subscription of a queue being created cannot be located', ) }) + + it('throws when queue creation config is provided with a locate-only subscription config', async () => { + await expect( + initSnsSqs( + sqsClient, + snsClient, + stsClient, + { topicName }, + { queue: { QueueName: queueName } }, + { locateOnly: true }, + ), + ).rejects.toThrow( + 'If subscriptionConfig.locateOnly is specified, both the topic and the queue must be located', + ) + }) + + it('throws when topic creation config is provided with a locate-only subscription config', async () => { + await expect( + initSnsSqs( + sqsClient, + snsClient, + stsClient, + { queueName }, + { topic: { Name: topicName }, queue: { QueueName: queueName } }, + { locateOnly: true }, + ), + ).rejects.toThrow( + 'If subscriptionConfig.locateOnly is specified, both the topic and the queue must be located', + ) + }) + + it('ignores creation config when topic and queue are located with a locate-only subscription config', async () => { + const topicArn = await testAdmin.createTopic(topicName) + const { queueArn } = await testAdmin.createQueue(queueName) + const subscriptionArn = await testAdmin.createSubscription(topicArn, queueArn) + const snsSpy = vi.spyOn(snsClient, 'send') + const sqsSpy = vi.spyOn(sqsClient, 'send') + + const result = await initSnsSqs( + sqsClient, + snsClient, + stsClient, + { topicName, queueName }, + { topic: { Name: topicName }, queue: { QueueName: queueName } }, + { locateOnly: true }, + ) + + expect(result).toEqual({ topicArn, queueUrl, queueArn, queueName, subscriptionArn }) + expect(snsSpy).not.toHaveBeenCalledWith(expect.any(CreateTopicCommand)) + expect(snsSpy).not.toHaveBeenCalledWith(expect.any(SubscribeCommand)) + expect(sqsSpy).not.toHaveBeenCalledWith(expect.any(CreateQueueCommand)) + }) }) describe('non-blocking startup resource polling', () => { diff --git a/packages/sns/lib/utils/snsInitter.ts b/packages/sns/lib/utils/snsInitter.ts index 56d8d99ab..d31ad175c 100644 --- a/packages/sns/lib/utils/snsInitter.ts +++ b/packages/sns/lib/utils/snsInitter.ts @@ -23,13 +23,18 @@ import { import type { SNSCreationConfig, SNSTopicLocatorType } from '../sns/AbstractSnsService.ts' import type { SNSSQSQueueLocatorType } from '../sns/AbstractSnsSqsConsumer.ts' import { isCreateTopicCommand } from '../types/TopicTypes.ts' -import type { SNSSubscriptionOptions } from './snsSubscriber.ts' -import { assertSubscription, subscribeToTopic } from './snsSubscriber.ts' +import { + assertSubscription, + type SNSSubscriptionCreationOptions, + type SNSSubscriptionOptions, + setSubscriptionAttributes, + subscribeToTopic, +} from './snsSubscriber.ts' import { assertTopic, deleteSubscription, deleteTopic, - findSubscriptionByTopicAndQueue, + findConfirmedSubscriptionArn, getTopicAttributes, } from './snsUtils.ts' import { buildTopicArn } from './stsUtils.ts' @@ -226,21 +231,12 @@ async function resolveQueue( ) } -async function findConfirmedSubscriptionArn( - snsClient: SNSClient, - topicArn: string, - queueArn: string, -): Promise { - const subscription = await findSubscriptionByTopicAndQueue(snsClient, topicArn, queueArn) - const subscriptionArn = subscription?.SubscriptionArn - - // Unconfirmed subscriptions are listed with a status placeholder instead of an ARN - return subscriptionArn?.startsWith('arn:') ? subscriptionArn : undefined -} - /** - * Subscription is used as is when its ARN is given, created (or updated) when subscriptionConfig is - * given, and otherwise located by looking up the queue's subscription on the topic. + * Subscription is created (or updated) when a creation subscriptionConfig is given. Otherwise, it is located: + * used as is when its ARN is given, or looked up on the topic, and the attributes of a locate-only + * subscriptionConfig are applied to it. + * + * Note: a creation subscriptionConfig is never given along with the subscription ARN, as initSnsSqs discards it. */ async function resolveSubscriptionArn( snsClient: SNSClient, @@ -250,9 +246,7 @@ async function resolveSubscriptionArn( subscriptionConfig: SNSSubscriptionOptions | undefined, options: ResourceResolutionOptions, ): Promise { - if (locatorConfig?.subscriptionArn) return locatorConfig.subscriptionArn - - if (subscriptionConfig) { + if (subscriptionConfig && !subscriptionConfig.locateOnly) { const subscriptionArn = await assertSubscription( snsClient, topicArn, @@ -266,20 +260,33 @@ async function resolveSubscriptionArn( return subscriptionArn } - return await waitForOrCheckResource( - `SNS subscription of SQS queue ${queue.queueArn} to topic ${topicArn}`, - `Subscription of queue ${queue.queueArn} to topic ${topicArn} does not exist.`, - async () => { - const subscriptionArn = await findConfirmedSubscriptionArn( - snsClient, - topicArn, - queue.queueArn, - ) - if (!subscriptionArn) return { isAvailable: false } - return { isAvailable: true, result: subscriptionArn } - }, - options, - ) + const subscriptionArn = + locatorConfig?.subscriptionArn ?? + (await waitForOrCheckResource( + `SNS subscription of SQS queue ${queue.queueArn} to topic ${topicArn}`, + `Subscription of queue ${queue.queueArn} to topic ${topicArn} does not exist.`, + async () => { + const subscriptionArn = await findConfirmedSubscriptionArn( + snsClient, + topicArn, + queue.queueArn, + ) + if (!subscriptionArn) return { isAvailable: false } + return { isAvailable: true, result: subscriptionArn } + }, + options, + )) + + if (subscriptionConfig?.locateOnly) { + await setSubscriptionAttributes( + snsClient, + subscriptionArn, + subscriptionConfig.Attributes, + subscriptionConfig.managedAttributes, + ) + } + + return subscriptionArn } async function resolveSnsSqsResources( @@ -329,6 +336,14 @@ function validateInitSnsSqsConfig( } // Queue creation config is ignored when the queue is located const isQueueLocated = !!(locatorConfig?.queueUrl || locatorConfig?.queueName) + const isTopicLocated = !!(locatorConfig?.topicArn || locatorConfig?.topicName) + + // A subscription managed externally can only belong to a topic and a queue managed externally too + if (subscriptionConfig?.locateOnly && (!isTopicLocated || !isQueueLocated)) { + throw new Error( + 'If subscriptionConfig.locateOnly is specified, both the topic and the queue must be located', + ) + } if (!isQueueLocated && creationConfig?.queue) { if (!creationConfig.queue.QueueName) { throw new Error( @@ -351,10 +366,12 @@ function validateInitSnsSqsConfig( * - Queue: located when `locatorConfig.queueUrl` or `locatorConfig.queueName` is given, otherwise * created from `creationConfig.queue`. * - Subscription: created (or updated) when `subscriptionConfig` is given, otherwise located by - * looking up the queue's subscription on the topic. + * looking up the queue's subscription on the topic. With `subscriptionConfig.locateOnly`, it is + * always located and its `Attributes` (e.g. `FilterPolicy`) are applied to it. * * When `locatorConfig.subscriptionArn` is given, every resource is located and `creationConfig` and - * `subscriptionConfig` are ignored. + * `subscriptionConfig` are ignored, except for a locate-only `subscriptionConfig`, whose attributes + * are still applied. * * Located resources are waited for when `locatorConfig.startupResourcePolling` is enabled. In * non-blocking mode, `undefined` is returned if they are not immediately available, and @@ -371,7 +388,10 @@ export async function initSnsSqs( ): Promise { const isFullyLocated = !!locatorConfig?.subscriptionArn const resolvedCreationConfig = isFullyLocated ? undefined : creationConfig - const resolvedSubscriptionConfig = isFullyLocated ? undefined : subscriptionConfig + // A locate-only subscriptionConfig doesn't create anything, so it is kept to apply its attributes + const resolvedSubscriptionConfig = + isFullyLocated && !subscriptionConfig?.locateOnly ? undefined : subscriptionConfig + const createsSubscription = !!resolvedSubscriptionConfig && !resolvedSubscriptionConfig.locateOnly validateInitSnsSqsConfig(locatorConfig, resolvedCreationConfig, resolvedSubscriptionConfig) @@ -388,7 +408,7 @@ export async function initSnsSqs( const startupResourcePolling = locatorConfig?.startupResourcePolling if (!isStartupResourcePollingEnabled(startupResourcePolling)) { - return await resolve({ checkLocatedTopic: !resolvedSubscriptionConfig }) + return await resolve({ checkLocatedTopic: !createsSubscription }) } if (startupResourcePolling.nonBlocking !== true) { return await resolve({ polling: startupResourcePolling, checkLocatedTopic: true }) @@ -429,7 +449,7 @@ export async function deleteSnsSqs( deletionConfig: DeletionConfig, queueConfiguration: CreateQueueCommandInput, topicConfiguration: CreateTopicCommandInput | undefined, - subscriptionConfiguration: SNSSubscriptionOptions, + subscriptionConfiguration: SNSSubscriptionCreationOptions, extraParams?: ExtraParams, topicLocator?: SNSTopicLocatorType, ) { diff --git a/packages/sns/lib/utils/snsSubscriber.spec.ts b/packages/sns/lib/utils/snsSubscriber.spec.ts index 57560c202..181319be8 100644 --- a/packages/sns/lib/utils/snsSubscriber.spec.ts +++ b/packages/sns/lib/utils/snsSubscriber.spec.ts @@ -1,13 +1,18 @@ -import type { SNSClient } from '@aws-sdk/client-sns' +import { + SetSubscriptionAttributesCommand, + type SNSClient, + SubscribeCommand, +} from '@aws-sdk/client-sns' import type { SQSClient } from '@aws-sdk/client-sqs' import type { STSClient } from '@aws-sdk/client-sts' import type { AwilixContainer } from 'awilix' -import { afterEach, beforeEach, describe, expect, it } from 'vitest' +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' import { FakeLogger } from '../../test/fakes/FakeLogger.ts' +import { isLocalstack } from '../../test/utils/fauxqsInstance.ts' import type { TestAwsResourceAdmin } from '../../test/utils/testAdmin.ts' import type { Dependencies } from '../../test/utils/testContext.ts' import { registerDependencies } from '../../test/utils/testContext.ts' -import { assertSubscription, subscribeToTopic } from './snsSubscriber.ts' +import { assertSubscription, setSubscriptionAttributes, subscribeToTopic } from './snsSubscriber.ts' import { findSubscriptionByTopicAndQueue, getSubscriptionAttributes } from './snsUtils.ts' const TOPIC_NAME = 'topic' @@ -81,13 +86,11 @@ describe('snsSubscriber', () => { logger, }, ), - ).rejects.toThrow( - /Invalid parameter: Attributes Reason: Subscription already exists with different attributes/, - ) + ).rejects.toThrow(/Subscription already exists with different attributes/) expect(logger.loggedErrors).toHaveLength(1) expect(logger.loggedErrors[0]).toBe( - 'Error while creating subscription for queue "queue", topic "topic": Invalid parameter: Attributes Reason: Subscription already exists with different attributes', + 'Error while creating subscription for queue "queue", topic "topic": Subscription already exists with different attributes: FilterPolicy, FilterPolicyScope', ) }) @@ -234,6 +237,231 @@ describe('snsSubscriber', () => { }) }) + it('does not write to an existing subscription when its attributes are unchanged', async () => { + const topicArn = await testAdmin.createTopic(TOPIC_NAME) + const { queueArn } = await testAdmin.createQueue(QUEUE_NAME) + const subscriptionConfig = { + Attributes: { + FilterPolicy: `{"type":["add","remove"]}`, + FilterPolicyScope: 'MessageBody', + RawMessageDelivery: 'true', + }, + updateAttributesIfExists: true, + } + const existingSubscriptionArn = await assertSubscription( + snsClient, + topicArn, + queueArn, + subscriptionConfig, + { queueName: QUEUE_NAME, topicName: TOPIC_NAME }, + ) + const sendSpy = vi.spyOn(snsClient, 'send') + + const subscriptionArn = await assertSubscription( + snsClient, + topicArn, + queueArn, + { + ...subscriptionConfig, + Attributes: { + ...subscriptionConfig.Attributes, + FilterPolicy: `{ "type": [ "add", "remove" ] }`, + }, + }, + { queueName: QUEUE_NAME, topicName: TOPIC_NAME }, + ) + + expect(subscriptionArn).toBe(existingSubscriptionArn) + const sentCommands = sendSpy.mock.calls.map(([command]) => command) + expect(sentCommands.some((command) => command instanceof SubscribeCommand)).toBe(false) + expect( + sentCommands.some((command) => command instanceof SetSubscriptionAttributesCommand), + ).toBe(false) + }) + + it('writes only the attributes that changed', async () => { + const topicArn = await testAdmin.createTopic(TOPIC_NAME) + const { queueArn } = await testAdmin.createQueue(QUEUE_NAME) + const subscriptionArn = await assertSubscription( + snsClient, + topicArn, + queueArn, + { + Attributes: { FilterPolicy: `{"type":["remove"]}`, RawMessageDelivery: 'true' }, + updateAttributesIfExists: false, + }, + { queueName: QUEUE_NAME, topicName: TOPIC_NAME }, + ) + const sendSpy = vi.spyOn(snsClient, 'send') + + await assertSubscription( + snsClient, + topicArn, + queueArn, + { + Attributes: { FilterPolicy: `{"type":["add"]}`, RawMessageDelivery: 'true' }, + updateAttributesIfExists: true, + }, + { queueName: QUEUE_NAME, topicName: TOPIC_NAME }, + new FakeLogger(), + ) + + const writtenAttributes = sendSpy.mock.calls + .map(([command]) => command) + .filter((command) => command instanceof SetSubscriptionAttributesCommand) + .map((command) => command.input.AttributeName) + expect(writtenAttributes).toEqual(['FilterPolicy']) + const subscriptionAttributes = await getSubscriptionAttributes(snsClient, subscriptionArn!) + expect(subscriptionAttributes.result?.attributes?.FilterPolicy).toBe(`{"type":["add"]}`) + }) + + it('neither checks nor writes attributes that are not managed', async () => { + const topicArn = await testAdmin.createTopic(TOPIC_NAME) + const { queueArn } = await testAdmin.createQueue(QUEUE_NAME) + const subscriptionArn = await assertSubscription( + snsClient, + topicArn, + queueArn, + { + Attributes: { FilterPolicy: `{"type":["add"]}`, RawMessageDelivery: 'true' }, + updateAttributesIfExists: false, + }, + { queueName: QUEUE_NAME, topicName: TOPIC_NAME }, + ) + const sendSpy = vi.spyOn(snsClient, 'send') + + await assertSubscription( + snsClient, + topicArn, + queueArn, + { + Attributes: { FilterPolicy: `{"type":["add"]}` }, + managedAttributes: ['FilterPolicy', 'FilterPolicyScope'], + updateAttributesIfExists: false, + }, + { queueName: QUEUE_NAME, topicName: TOPIC_NAME }, + ) + + expect( + sendSpy.mock.calls.some(([command]) => command instanceof SetSubscriptionAttributesCommand), + ).toBe(false) + const subscriptionAttributes = await getSubscriptionAttributes(snsClient, subscriptionArn!) + expect(subscriptionAttributes.result?.attributes?.RawMessageDelivery).toBe('true') + }) + + it('resets managed attributes that were removed from the config', async () => { + const topicArn = await testAdmin.createTopic(TOPIC_NAME) + const { queueArn } = await testAdmin.createQueue(QUEUE_NAME) + const subscriptionArn = await assertSubscription( + snsClient, + topicArn, + queueArn, + { + Attributes: { + FilterPolicy: `{"type":["add"]}`, + FilterPolicyScope: 'MessageBody', + RawMessageDelivery: 'true', + }, + updateAttributesIfExists: false, + }, + { queueName: QUEUE_NAME, topicName: TOPIC_NAME }, + ) + + await assertSubscription( + snsClient, + topicArn, + queueArn, + { Attributes: { FilterPolicy: `{"type":["add"]}` }, updateAttributesIfExists: true }, + { queueName: QUEUE_NAME, topicName: TOPIC_NAME }, + new FakeLogger(), + ) + + const subscriptionAttributes = await getSubscriptionAttributes(snsClient, subscriptionArn!) + expect(subscriptionAttributes.result?.attributes).toMatchObject({ + FilterPolicy: `{"type":["add"]}`, + FilterPolicyScope: 'MessageAttributes', + RawMessageDelivery: 'false', + }) + + // A subscription that is already reset is not written to again + const sendSpy = vi.spyOn(snsClient, 'send') + await assertSubscription( + snsClient, + topicArn, + queueArn, + { Attributes: { FilterPolicy: `{"type":["add"]}` }, updateAttributesIfExists: true }, + { queueName: QUEUE_NAME, topicName: TOPIC_NAME }, + ) + expect( + sendSpy.mock.calls.some(([command]) => command instanceof SetSubscriptionAttributesCommand), + ).toBe(false) + }) + + // LocalStack cannot remove a RedrivePolicy: it rejects both an empty and an omitted value + it.skipIf(isLocalstack)( + 'removes a managed redrive policy missing from the config', + async () => { + const topicArn = await testAdmin.createTopic(TOPIC_NAME) + const { queueArn } = await testAdmin.createQueue(QUEUE_NAME) + const subscriptionArn = await assertSubscription( + snsClient, + topicArn, + queueArn, + { + Attributes: { RedrivePolicy: JSON.stringify({ deadLetterTargetArn: queueArn }) }, + updateAttributesIfExists: false, + }, + { queueName: QUEUE_NAME, topicName: TOPIC_NAME }, + ) + const sendSpy = vi.spyOn(snsClient, 'send') + + await assertSubscription( + snsClient, + topicArn, + queueArn, + { updateAttributesIfExists: true }, + { queueName: QUEUE_NAME, topicName: TOPIC_NAME }, + new FakeLogger(), + ) + + const redrivePolicyWrite = sendSpy.mock.calls + .map(([command]) => command) + .find( + (command) => + command instanceof SetSubscriptionAttributesCommand && + command.input.AttributeName === 'RedrivePolicy', + ) + expect(redrivePolicyWrite?.input).toEqual({ + SubscriptionArn: subscriptionArn, + AttributeName: 'RedrivePolicy', + AttributeValue: undefined, + }) + const subscriptionAttributes = await getSubscriptionAttributes(snsClient, subscriptionArn!) + expect( + subscriptionAttributes.result?.attributes?.RedrivePolicy || undefined, + ).toBeUndefined() + }, + ) + + it('throws when a configured attribute is not managed', async () => { + const topicArn = await testAdmin.createTopic(TOPIC_NAME) + const { queueArn } = await testAdmin.createQueue(QUEUE_NAME) + + await expect( + assertSubscription( + snsClient, + topicArn, + queueArn, + { + Attributes: { FilterPolicy: `{"type":["add"]}`, RawMessageDelivery: 'true' }, + managedAttributes: ['FilterPolicy'], + updateAttributesIfExists: true, + }, + { queueName: QUEUE_NAME, topicName: TOPIC_NAME }, + ), + ).rejects.toThrow(/RawMessageDelivery are set, but not listed in managedAttributes/) + }) + it('throws when attributes are different and update is disabled', async () => { const logger = new FakeLogger() const topicArn = await testAdmin.createTopic(TOPIC_NAME) @@ -270,4 +498,29 @@ describe('snsSubscriber', () => { ) }) }) + + describe('setSubscriptionAttributes', () => { + it('writes only the attributes that differ from the current ones', async () => { + const topicArn = await testAdmin.createTopic(TOPIC_NAME) + const { queueArn } = await testAdmin.createQueue(QUEUE_NAME) + const subscriptionArn = await testAdmin.createSubscription(topicArn, queueArn) + await setSubscriptionAttributes(snsClient, subscriptionArn, { + FilterPolicy: `{"type":["add"]}`, + }) + const sendSpy = vi.spyOn(snsClient, 'send') + + await setSubscriptionAttributes( + snsClient, + subscriptionArn, + { FilterPolicy: `{ "type": ["add"] }`, RawMessageDelivery: 'true' }, + ['FilterPolicy', 'RawMessageDelivery'], + ) + + const writtenAttributes = sendSpy.mock.calls + .map(([command]) => command) + .filter((command) => command instanceof SetSubscriptionAttributesCommand) + .map((command) => command.input.AttributeName) + expect(writtenAttributes).toEqual(['RawMessageDelivery']) + }) + }) }) diff --git a/packages/sns/lib/utils/snsSubscriber.ts b/packages/sns/lib/utils/snsSubscriber.ts index 025e963cb..a2e86a29b 100644 --- a/packages/sns/lib/utils/snsSubscriber.ts +++ b/packages/sns/lib/utils/snsSubscriber.ts @@ -1,3 +1,4 @@ +import { isDeepStrictEqual } from 'node:util' import type { SNSClient, SubscribeCommandInput } from '@aws-sdk/client-sns' import { SetSubscriptionAttributesCommand, SubscribeCommand } from '@aws-sdk/client-sns' import type { CreateQueueCommandInput, SQSClient } from '@aws-sdk/client-sqs' @@ -17,13 +18,58 @@ import { isSNSTopicLocatorType, type TopicResolutionOptions, } from '../types/TopicTypes.ts' -import { assertTopic, findSubscriptionByTopicAndQueue } from './snsUtils.ts' +import { assertTopic, findConfirmedSubscriptionArn, getSubscriptionAttributes } from './snsUtils.ts' import { buildTopicArn } from './stsUtils.ts' -export type SNSSubscriptionOptions = Omit< - SubscribeCommandInput, - 'TopicArn' | 'Endpoint' | 'Protocol' | 'ReturnSubscriptionArn' -> & { updateAttributesIfExists: boolean } +/** Subscription attributes that can be managed for an SQS subscription */ +export const SUBSCRIPTION_MANAGED_ATTRIBUTE_NAMES = [ + 'FilterPolicy', + 'FilterPolicyScope', + 'RawMessageDelivery', + 'RedrivePolicy', +] as const + +export type SubscriptionManagedAttributeName = (typeof SUBSCRIPTION_MANAGED_ATTRIBUTE_NAMES)[number] + +type SubscriptionAttributeManagementOptions = { + /** + * Attributes whose value is owned by this config. Defaults to all of + * {@link SUBSCRIPTION_MANAGED_ATTRIBUTE_NAMES}. + * + * A managed attribute that is missing from `Attributes` is reset on the existing subscription + * (`FilterPolicy` and `RedrivePolicy` are removed, the others go back to their defaults). + * Attributes that are not managed are neither checked nor written, so external tooling such as + * Terraform can own them, and must not be present in `Attributes`. + */ + managedAttributes?: readonly SubscriptionManagedAttributeName[] +} + +/** + * Options for the subscription of a queue to a topic. + * + * There are two modes for this config: + * 1. Creation (default): the subscription is created, or updated if it already exists. + * 2. Locate only: for subscriptions managed externally. The subscription is only located, never created, + * and its managed attributes (e.g. `FilterPolicy`) are applied to it. + */ +export type SNSSubscriptionOptions = + /** Creation */ + | (Omit & + SubscriptionAttributeManagementOptions & { + updateAttributesIfExists: boolean + /** Should not be present when the subscription is created */ + locateOnly?: never + }) + /** Locate only */ + | (SubscriptionAttributeManagementOptions & { + /** Marks the subscription as located only */ + locateOnly: true + /** Attributes applied to the located subscription */ + Attributes?: SubscribeCommandInput['Attributes'] + }) + +/** Subscription options that create the subscription, required when subscribing */ +export type SNSSubscriptionCreationOptions = Exclude async function resolveTopicArnToSubscribeTo( snsClient: SNSClient, @@ -61,7 +107,7 @@ export async function subscribeToTopic( stsClient: STSClient, queueConfiguration: CreateQueueCommandInput, topicConfiguration: TopicResolutionOptions, - subscriptionConfiguration: SNSSubscriptionOptions, + subscriptionConfiguration: SNSSubscriptionCreationOptions, extraParams?: ExtraSNSCreationParams & ExtraSQSCreationParams & ExtraParams, ) { const topicArn = await resolveTopicArnToSubscribeTo( @@ -95,54 +141,81 @@ export async function subscribeToTopic( } /** - * Subscribes an existing SQS queue to an existing SNS topic. Subscribing is idempotent, so an - * existing subscription is returned as is, or has its attributes updated when they differ and - * `updateAttributesIfExists` is enabled. + * Subscribes an existing SQS queue to an existing SNS topic. An existing subscription is checked + * with read-only calls first and is only written to when its managed attributes differ from the + * configured ones and `updateAttributesIfExists` is enabled. Only the differing attributes are + * written. */ export async function assertSubscription( snsClient: SNSClient, topicArn: string, queueArn: string, - subscriptionConfiguration: SNSSubscriptionOptions, + subscriptionConfiguration: SNSSubscriptionCreationOptions, errorContext: { queueName?: string; topicName?: string }, logger?: CommonLogger, ): Promise { - const subscribeCommand = new SubscribeCommand({ - TopicArn: topicArn, - Endpoint: queueArn, - Protocol: 'sqs', - ReturnSubscriptionArn: true, - ...subscriptionConfiguration, - }) + const { updateAttributesIfExists, managedAttributes, locateOnly, ...subscribeInput } = + subscriptionConfiguration + const resolvedManagedAttributes = resolveManagedAttributes( + subscribeInput.Attributes, + managedAttributes, + ) + const resolvedLogger = logger ?? console + const errMessagePrefix = `Error while creating subscription for queue "${errorContext.queueName}", topic "${errorContext.topicName}"` + + const reconcileExistingSubscription = async (subscriptionArn: string) => { + const changedAttributes = await findChangedSubscriptionAttributes( + snsClient, + subscriptionArn, + subscribeInput.Attributes, + resolvedManagedAttributes, + ) + if (Object.keys(changedAttributes).length === 0) return subscriptionArn + + const errMessage = `${errMessagePrefix}: Subscription already exists with different attributes: ${Object.keys(changedAttributes).join(', ')}` + if (!updateAttributesIfExists) { + resolvedLogger.error(errMessage) + throw new InternalError({ + errorCode: 'sns_subscription_creation_failed', + message: errMessage, + details: { queueName: errorContext.queueName, topicArn }, + }) + } + + resolvedLogger.warn(`${errMessage}. Updating subscription`) + await writeSubscriptionAttributes(snsClient, subscriptionArn, changedAttributes) + return subscriptionArn + } + + const existingSubscriptionArn = await findConfirmedSubscriptionArn(snsClient, topicArn, queueArn) + if (existingSubscriptionArn) return reconcileExistingSubscription(existingSubscriptionArn) try { - const subscriptionResult = await snsClient.send(subscribeCommand) + const subscriptionResult = await snsClient.send( + new SubscribeCommand({ + ...subscribeInput, + TopicArn: topicArn, + Endpoint: queueArn, + Protocol: 'sqs', + ReturnSubscriptionArn: true, + }), + ) return subscriptionResult.SubscriptionArn } catch (err) { if (!isError(err)) throw err - const resolvedLogger = logger ?? console - const errMessage = `Error while creating subscription for queue "${errorContext.queueName}", topic "${errorContext.topicName}": ${err.message}` - - if ( - subscriptionConfiguration.updateAttributesIfExists && - err.message.includes('Subscription already exists with different attributes') - ) { - resolvedLogger.warn(`${errMessage}. Trying to update subscription`) - - const result = await tryToUpdateSubscription( + // Another instance may have created the subscription after the lookup above + if (err.message.includes('Subscription already exists with different attributes')) { + const concurrentSubscriptionArn = await findConfirmedSubscriptionArn( snsClient, topicArn, queueArn, - subscriptionConfiguration, ) - if (result) return result.SubscriptionArn - - resolvedLogger.error('Failed to update subscription') - } else { - resolvedLogger.error(errMessage) + if (concurrentSubscriptionArn) return reconcileExistingSubscription(concurrentSubscriptionArn) } + const errMessage = `${errMessagePrefix}: ${err.message}` + resolvedLogger.error(errMessage) throw new InternalError({ errorCode: 'sns_subscription_creation_failed', message: errMessage, @@ -152,30 +225,133 @@ export async function assertSubscription( } } -async function tryToUpdateSubscription( +/** + * Applies the given attributes to an existing subscription, resetting managed attributes that are + * missing from them. Current attributes are read first, and only the ones that differ are written. + */ +export async function setSubscriptionAttributes( snsClient: SNSClient, - topicArn: string, - queueArn: string, - subscriptionConfiguration: SNSSubscriptionOptions, -) { - const subscription = await findSubscriptionByTopicAndQueue(snsClient, topicArn, queueArn) - if (!subscription || !subscriptionConfiguration.Attributes) { - return undefined + subscriptionArn: string, + attributes: SubscribeCommandInput['Attributes'], + managedAttributes?: readonly SubscriptionManagedAttributeName[], +): Promise { + const changedAttributes = await findChangedSubscriptionAttributes( + snsClient, + subscriptionArn, + attributes, + resolveManagedAttributes(attributes, managedAttributes), + ) + await writeSubscriptionAttributes(snsClient, subscriptionArn, changedAttributes) +} + +// SNS only allows setting a single attribute per call +async function writeSubscriptionAttributes( + snsClient: SNSClient, + subscriptionArn: string, + attributes: Record, +): Promise { + for (const [key, value] of Object.entries(attributes)) { + await snsClient.send( + new SetSubscriptionAttributesCommand({ + SubscriptionArn: subscriptionArn, + AttributeName: key, + // AWS rejects an empty RedrivePolicy, omitting the value is what removes it + AttributeValue: key === 'RedrivePolicy' && value === '' ? undefined : value, + }), + ) } +} + +// Values that reset an attribute to the state of a subscription where it was never set +const SUBSCRIPTION_ATTRIBUTE_UNSET_VALUES: Record = { + FilterPolicy: '{}', + FilterPolicyScope: 'MessageAttributes', + RawMessageDelivery: 'false', + RedrivePolicy: '', +} - const setSubscriptionAttributesCommands = Object.entries( - subscriptionConfiguration.Attributes, - ).map(([key, value]) => { - return new SetSubscriptionAttributesCommand({ - SubscriptionArn: subscription.SubscriptionArn, - AttributeName: key, - AttributeValue: value, +function isManagedAttributeName(name: string): name is SubscriptionManagedAttributeName { + return (SUBSCRIPTION_MANAGED_ATTRIBUTE_NAMES as readonly string[]).includes(name) +} + +function resolveManagedAttributes( + attributes: Record | undefined, + managedAttributes: readonly SubscriptionManagedAttributeName[] | undefined, +): readonly SubscriptionManagedAttributeName[] { + const resolvedManagedAttributes = managedAttributes ?? SUBSCRIPTION_MANAGED_ATTRIBUTE_NAMES + const unmanagedAttributes = Object.keys(attributes ?? {}).filter( + (name) => isManagedAttributeName(name) && !resolvedManagedAttributes.includes(name), + ) + if (unmanagedAttributes.length > 0) { + throw new InternalError({ + errorCode: 'invalid_subscription_configuration', + message: `Subscription attributes ${unmanagedAttributes.join(', ')} are set, but not listed in managedAttributes`, }) - }) + } + + return resolvedManagedAttributes +} + +/** + * Compares the current attributes of the subscription with the expected ones. Every managed + * attribute is compared, falling back to its unset value when it is not configured. Configured + * attributes outside {@link SUBSCRIPTION_MANAGED_ATTRIBUTE_NAMES} are compared, but never reset. + */ +async function findChangedSubscriptionAttributes( + snsClient: SNSClient, + subscriptionArn: string, + attributes: Record | undefined, + managedAttributes: readonly SubscriptionManagedAttributeName[], +): Promise> { + const expectedAttributes: Record = {} + // Iterating in the declared order writes FilterPolicy before FilterPolicyScope + for (const name of SUBSCRIPTION_MANAGED_ATTRIBUTE_NAMES) { + if (managedAttributes.includes(name)) { + expectedAttributes[name] = attributes?.[name] ?? SUBSCRIPTION_ATTRIBUTE_UNSET_VALUES[name] + } + } + for (const [name, value] of Object.entries(attributes ?? {})) { + if (!isManagedAttributeName(name)) expectedAttributes[name] = value + } + if (Object.keys(expectedAttributes).length === 0) return {} - for (const command of setSubscriptionAttributesCommands) { - await snsClient.send(command) + const currentAttributesResult = await getSubscriptionAttributes(snsClient, subscriptionArn) + const currentAttributes = currentAttributesResult.result?.attributes ?? {} + + const changedAttributes = Object.fromEntries( + Object.entries(expectedAttributes).filter( + ([name, value]) => !isSameAttributeValue(name, value, currentAttributes[name]), + ), + ) + + // Filter policy scope means nothing without a filter policy, so it is not reset on its own + const resultingFilterPolicy = expectedAttributes.FilterPolicy ?? currentAttributes.FilterPolicy + if ( + isUnsetAttributeValue('FilterPolicy', resultingFilterPolicy) && + isUnsetAttributeValue('FilterPolicyScope', changedAttributes.FilterPolicyScope) + ) { + delete changedAttributes.FilterPolicyScope } - return subscription + return changedAttributes +} + +function isUnsetAttributeValue(name: string, value: string | undefined): boolean { + if (value === undefined || value === '') return true + if (name === 'FilterPolicy' && value === '{}') return true + + return isManagedAttributeName(name) && SUBSCRIPTION_ATTRIBUTE_UNSET_VALUES[name] === value +} + +// JSON attributes are compared structurally, as SNS does not preserve their formatting +function isSameAttributeValue(name: string, expected: string, current: string | undefined) { + if (expected === current) return true + if (isUnsetAttributeValue(name, expected) && isUnsetAttributeValue(name, current)) return true + if (current === undefined || !name.endsWith('Policy')) return false + + try { + return isDeepStrictEqual(JSON.parse(expected), JSON.parse(current)) + } catch { + return false + } } diff --git a/packages/sns/lib/utils/snsUtils.ts b/packages/sns/lib/utils/snsUtils.ts index 5d7c8b827..0960bd644 100644 --- a/packages/sns/lib/utils/snsUtils.ts +++ b/packages/sns/lib/utils/snsUtils.ts @@ -190,6 +190,18 @@ export async function findSubscriptionByTopicAndQueue( return undefined } +export async function findConfirmedSubscriptionArn( + snsClient: SNSClient, + topicArn: string, + queueArn: string, +): Promise { + const subscription = await findSubscriptionByTopicAndQueue(snsClient, topicArn, queueArn) + const subscriptionArn = subscription?.SubscriptionArn + + // Unconfirmed subscriptions are listed with a status placeholder instead of an ARN + return subscriptionArn?.startsWith('arn:') ? subscriptionArn : undefined +} + /** * Calculates the size of an outgoing SNS message. * diff --git a/packages/sns/test/consumers/SnsSqsPermissionConsumer.spec.ts b/packages/sns/test/consumers/SnsSqsPermissionConsumer.spec.ts index ca12d2eae..399873867 100644 --- a/packages/sns/test/consumers/SnsSqsPermissionConsumer.spec.ts +++ b/packages/sns/test/consumers/SnsSqsPermissionConsumer.spec.ts @@ -1,6 +1,12 @@ import { setTimeout } from 'node:timers/promises' -import { ListTagsForResourceCommand, type SNSClient, SubscribeCommand } from '@aws-sdk/client-sns' -import { ListQueueTagsCommand, type SQSClient } from '@aws-sdk/client-sqs' +import { + DeleteTopicCommand, + ListTagsForResourceCommand, + type SNSClient, + SubscribeCommand, + UnsubscribeCommand, +} from '@aws-sdk/client-sns' +import { DeleteQueueCommand, ListQueueTagsCommand, type SQSClient } from '@aws-sdk/client-sqs' import type { STSClient } from '@aws-sdk/client-sts' import { waitAndRetry } from '@lokalise/node-core' import { getQueueAttributes } from '@message-queue-toolkit/sqs' @@ -161,6 +167,29 @@ describe('SnsSqsPermissionConsumer', () => { expect(snsSpy).not.toHaveBeenCalledWith(expect.any(SubscribeCommand)) }) + it('does not delete nor create any resource with a locate-only subscription config', async () => { + const topicArn = await testAdmin.createTopic(topicName) + const { queueArn } = await testAdmin.createQueue(queueName) + const subscriptionArn = await testAdmin.createSubscription(topicArn, queueArn) + const snsSpy = vi.spyOn(snsClient, 'send') + const sqsSpy = vi.spyOn(sqsClient, 'send') + + const newConsumer = new SnsSqsPermissionConsumer(diContainer.cradle, { + locatorConfig: { topicName, queueName }, + creationConfig: { queue: { QueueName: queueName }, topic: { Name: topicName } }, + deletionConfig: { deleteIfExists: true }, + subscriptionConfig: { locateOnly: true }, + }) + + await newConsumer.init() + expect(newConsumer.subscriptionProps.subscriptionArn).toBe(subscriptionArn) + expect(snsSpy).not.toHaveBeenCalledWith(expect.any(SubscribeCommand)) + expect(snsSpy).not.toHaveBeenCalledWith(expect.any(UnsubscribeCommand)) + expect(snsSpy).not.toHaveBeenCalledWith(expect.any(DeleteTopicCommand)) + expect(sqsSpy).not.toHaveBeenCalledWith(expect.any(DeleteQueueCommand)) + expect(await findSubscriptionByTopicAndQueue(snsClient, topicArn, queueArn)).toBeDefined() + }) + describe('tags update', () => { const getQueueTags = (queueUrl: string) => sqsClient.send(new ListQueueTagsCommand({ QueueUrl: queueUrl })) diff --git a/packages/sns/test/utils/fauxqsInstance.ts b/packages/sns/test/utils/fauxqsInstance.ts index c2dafa5f0..4a6366098 100644 --- a/packages/sns/test/utils/fauxqsInstance.ts +++ b/packages/sns/test/utils/fauxqsInstance.ts @@ -1,7 +1,7 @@ /** biome-ignore-all lint/suspicious/noConsole: test **/ import { type FauxqsServer, startFauxqs } from 'fauxqs' -const isLocalstack = process.env.QUEUE_BACKEND === 'localstack' +export const isLocalstack = process.env.QUEUE_BACKEND === 'localstack' const LOCALSTACK_PORT = 4566 const LOCALSTACK_HOST = 'localstack' const FAUXQS_PORT = 4567 diff --git a/packages/sqs/lib/sqs/AbstractSqsConsumer.ts b/packages/sqs/lib/sqs/AbstractSqsConsumer.ts index c2799aad0..516850a11 100644 --- a/packages/sqs/lib/sqs/AbstractSqsConsumer.ts +++ b/packages/sqs/lib/sqs/AbstractSqsConsumer.ts @@ -78,7 +78,7 @@ type SQSDeadLetterQueueOptions = | { /** * A located DLQ can omit the redrive policy, in which case the redrive policy of the source queue is left - * untouched (e.g. when both queues are managed externally) + * untouched */ locatorConfig: SQSQueueLocatorType creationConfig?: never diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index b45e1e61e..8f720f010 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -370,7 +370,7 @@ importers: version: 5.11.1(supports-color@7.2.0) redis-semaphore: specifier: ^5.7.0 - version: 5.7.0(ioredis@5.11.1(supports-color@7.2.0))(supports-color@7.2.0) + version: 5.8.0(ioredis@5.11.1(supports-color@7.2.0))(supports-color@7.2.0) devDependencies: '@biomejs/biome': specifier: ^2.5.14 @@ -631,74 +631,38 @@ packages: resolution: {integrity: sha512-bgLbxtF/GBSTDJtl+WweClwoQkNd9TxmoGikHCs+c16rvvCs0sFOsmEZ0AT5DLIxGvRZOO/NT9d3FBQNsIs9cg==} engines: {node: '>=20.0.0'} - '@aws-sdk/core@3.978.0': - resolution: {integrity: sha512-2yX9LUmxPklVjSGTb8dfnWRJSiFQ3TeH2nn7G1mdKHTfnabzF0+gfrS8rYfLWmZrQ8A3mEcxMJjRc51dL5KWaA==} - engines: {node: '>=20.0.0'} - '@aws-sdk/core@3.978.1': resolution: {integrity: sha512-LbY9aGsEiznDWmUc30Nwv3aIX/+dbwTx8KfS0yOC3NPYMO+O91e6jkT1azf34FwjOndq8/Q+RcVVZz5xnerwdg==} engines: {node: '>=20.0.0'} - '@aws-sdk/credential-provider-env@3.972.71': - resolution: {integrity: sha512-JN+JHruYZw3GUZB8YGAlDk4wTDPOEAEEdEzj5nS0xodWR4smzHsN7PnK2j6IeOsDIj2aqua5DSbhXl9Gtf90FQ==} - engines: {node: '>=20.0.0'} - '@aws-sdk/credential-provider-env@3.972.72': resolution: {integrity: sha512-xTKO/FWJPozTIXbozVnVGoNBhaGba8TBcx+KyUjRVeOlXE+dUc7GTR1cLvu0uTdIdmemzaFbqqCshXeZA1fZew==} engines: {node: '>=20.0.0'} - '@aws-sdk/credential-provider-http@3.972.73': - resolution: {integrity: sha512-uyYYnJOnlis8uQzaYGPd7N1JoioCoNpXgnkXYixsWJXHXgXyYi8WXJSDfofxJeWfQIGWLe2Nwyq60Uc7MZdVOg==} - engines: {node: '>=20.0.0'} - '@aws-sdk/credential-provider-http@3.972.74': resolution: {integrity: sha512-u91E/hT8f4d1xy0Jl7VG4nVKJ3lxbrZkoBTeSVoJdWBiSEUMwMS/9+e0H/aJVQV//Lt5wuzP+E69v4aRSsNTmw==} engines: {node: '>=20.0.0'} - '@aws-sdk/credential-provider-ini@3.973.16': - resolution: {integrity: sha512-i++ly+0Uxa+u3ebSSyr0S/3CFhFJDxCXT3+Zj+mW2bXenEx5bKGCdTIKFu39SgXBNhWDjex/8cXUx9MUTMCrTw==} - engines: {node: '>=20.0.0'} - '@aws-sdk/credential-provider-ini@3.973.17': resolution: {integrity: sha512-ged4KXdBkvIC81bLvNHHuQKdKak/VXhQTR1NWYTTqW0474nlmsxy9O/vlgTIohDDWH3xpBdtVMZRyjb+DnocDA==} engines: {node: '>=20.0.0'} - '@aws-sdk/credential-provider-login@3.972.78': - resolution: {integrity: sha512-eUtswnXu0+Ii9ieRK+0L7aPFV3Z/dnW2VntJzjBP9xs8s+8p5nBNuymIXtXwZ+5r5+XJP3e32nMkuZ/r0HozEA==} - engines: {node: '>=20.0.0'} - '@aws-sdk/credential-provider-login@3.972.79': resolution: {integrity: sha512-L+Z85anONJd8MaiuraO4wRxATCdEejBZ3K3eymzWI5JPXa9sOS9CkIm72PBKqXKX+Z9p9NGMX5AIMXm0LEflgw==} engines: {node: '>=20.0.0'} - '@aws-sdk/credential-provider-node@3.972.83': - resolution: {integrity: sha512-jdso7ejzfRnatxMUZK4S/U6KbaDPCvfIV4XL+IQAPFDBt5rj5Fq595euqlK8Le4lNCMFR9oUpt+1l0aMgaayOQ==} - engines: {node: '>=20.0.0'} - '@aws-sdk/credential-provider-node@3.972.84': resolution: {integrity: sha512-oHt854odINVwzwsh+c5x69j0ajm4DbqqqVJ+O1ECsCIZeMDAbzFpXItaqP7UZstJj/ATdTk/KFSH0LaNAgV+kA==} engines: {node: '>=20.0.0'} - '@aws-sdk/credential-provider-process@3.972.71': - resolution: {integrity: sha512-lYmXJa4gvq4xN1lrT5NiP5vIYYKcGWAdj8y+8o6dlcateB5eF3Dn8DtmjjHKfMBrTPAMr2pebIiX/UOj8c1/UA==} - engines: {node: '>=20.0.0'} - '@aws-sdk/credential-provider-process@3.972.72': resolution: {integrity: sha512-rLIp2xbMjX/k9/od7APpqq1ZgXXnV0pOL1Th3ZsL8Wu0TRtBsDTVS8iPqcfRFcHakFxPvR04OSTv2ka2qOb/2A==} engines: {node: '>=20.0.0'} - '@aws-sdk/credential-provider-sso@3.973.15': - resolution: {integrity: sha512-6Jhcf4v0pSFdjk1EW2kvzuEBKD+UZ2uNcHUIglKKLndD20YhvkL2kdmDOV5/j4mYuWWwe/a1FQ1aomU86/Cg5Q==} - engines: {node: '>=20.0.0'} - '@aws-sdk/credential-provider-sso@3.973.16': resolution: {integrity: sha512-IGihaJfFZYacJJr/odqILCoK7W/mvrZ7cuK7ECn3sAu4vLC6u0V8bS7mCGbdugJ8Aum2tnvqmx0F2MRFp2rn9g==} engines: {node: '>=20.0.0'} - '@aws-sdk/credential-provider-web-identity@3.972.77': - resolution: {integrity: sha512-uylIQSUWpfLuH2LovxEEfwzJGM/SabLOfLMg6YXu/E8jJEKUdpdILCVCQCdFvHyu/7dLJOHPMfrSwduxO56NkQ==} - engines: {node: '>=20.0.0'} - '@aws-sdk/credential-provider-web-identity@3.972.78': resolution: {integrity: sha512-/y9WvNtlcPBGLR0qc1a+9J/xtYZfVczvLUOuXaVWylzttH7ewsxwHtjmiJSolNrVSDorIxHGHMU61CbonRkmwA==} engines: {node: '>=20.0.0'} @@ -711,10 +675,6 @@ packages: resolution: {integrity: sha512-D/O6iHAqlm0b430vIGewanSsr5RzaWbSDFlHYdKLWlZb4kaMKBmb44ReNHAFEKj5tNRY8pysw0qfciosB/VcJg==} engines: {node: '>=20.0.0'} - '@aws-sdk/nested-clients@3.997.45': - resolution: {integrity: sha512-mooq9Q+jLa18VoM7HouczmslZU60iiB0aKc/Ztnq/luIL1ud0z4DnYprLR/ZO1gp331S9tJctM1HZr7u6YKBXQ==} - engines: {node: '>=20.0.0'} - '@aws-sdk/nested-clients@3.997.46': resolution: {integrity: sha512-oRxtBcka/JGHGs9l9p9IVajGoTP8vTPmoAzdHGy4Qcy9P5vPnDf6nhIeM/COQNY9k/OahImTRaLkHftoXvfcmQ==} engines: {node: '>=20.0.0'} @@ -723,26 +683,14 @@ packages: resolution: {integrity: sha512-Zk08macMvQTHzQJCLJVkOlviVoqwYMrpXv4lmLN7b7sAbiMoOK7Go0NYdR5UeF+MW8LIbRmwrNy9u/5VvX1U5g==} engines: {node: '>=20.0.0'} - '@aws-sdk/token-providers@3.1129.0': - resolution: {integrity: sha512-Sbl3rpzQdsG4ZK2zh0JWUYyZPKKorJlVOddA2T0DVbKJFrsW8J6wgnslxxUH04+WaBMr4A1HzJZvZX0xUvkniA==} - engines: {node: '>=20.0.0'} - '@aws-sdk/token-providers@3.1138.0': resolution: {integrity: sha512-GpyAr0DD63YOEmYFM6Df+gJuIgC92MMTiBK4FTKfxii5MJ9ge20epR7LyroulscYlG89J+ZB2ivFDPjvfQhzdw==} engines: {node: '>=20.0.0'} - '@aws-sdk/types@3.974.5': - resolution: {integrity: sha512-LkwLL2BLbC6wNNm4JaH9mbEqBMdOZCct6VAYqhdN4U1xrWM+fUJQEfbHwQgDypapOWTRtlk25akb5afM0P8CIQ==} - engines: {node: '>=20.0.0'} - '@aws-sdk/types@3.974.6': resolution: {integrity: sha512-v/clNZzZnDxGyvpHMOGpJKVXFAExJzUNAAjaWGdcx8QAcXLGwTaOkw33p5SHAi0YAioK32xB3hWwOekRVfmfKg==} engines: {node: '>=20.0.0'} - '@aws-sdk/xml-builder@3.972.40': - resolution: {integrity: sha512-wlFmCIGUlwF4zx/kncw+bmxTQh1HeSJq4mYV/V5cZUSJadDP3kXvGW8Rn21cimj/7y9ju+47oYWXi97vF7czaA==} - engines: {node: '>=20.0.0'} - '@aws-sdk/xml-builder@3.972.41': resolution: {integrity: sha512-ctjVSyCMegrWfXlx6VqzSBFI6UqmQ5ZlnfMhdLIiWmhoH8UAQxSCP5N3OpG7X3k4LnS7ou74C4mt20+bfTW2aQ==} engines: {node: '>=20.0.0'} @@ -1258,10 +1206,6 @@ packages: '@rolldown/pluginutils@1.0.1': resolution: {integrity: sha512-2j9bGt5Jh8hj+vPtgzPtl72j0yRxHAyumoo6TNfAjsLB04UtpSvPbPcDcBMxz7n+9CYB0c1GxQFxYRg2jimqGw==} - '@smithy/core@3.33.3': - resolution: {integrity: sha512-CsOeKq/9kA3y6VJHt+/+VTCtBaxJ4OTFpgrjIUhPpDIKxBci1k2bJaQASF2h/ELWrulGp+t97DZ0mevfAD8idg==} - engines: {node: '>=18.0.0'} - '@smithy/core@3.35.0': resolution: {integrity: sha512-zRMhfkByhT2snNdr1si24vJitU6Cr9ix2MikUfWmkAgp4jrNP0GcKSP5YvwQ+TlI8AZXER5QOGJn3JsVtSD9/A==} engines: {node: '>=18.0.0'} @@ -1270,19 +1214,10 @@ packages: resolution: {integrity: sha512-A9uSdn72ozbRUSit0eib0TW7nXuNPlaeM0zcGkJ+nE6tFcSDbnmtwoxbTCFBukVQcszDAyvsd7+rTduPTXpygg==} engines: {node: '>=18.0.0'} - '@smithy/fetch-http-handler@5.7.2': - resolution: {integrity: sha512-nZyWTmSpJEXl6VtWVMBJve/7x12DZu6sIX1z1a+ZMaHlQQRs9Zpu6NbTe/gmxYXVRpkjxyDYpZ5gx2IM6f/Wkw==} - engines: {node: '>=18.0.0'} - '@smithy/fetch-http-handler@5.8.0': resolution: {integrity: sha512-ycSJu3tFAQ4v04CBB0agqFMVsSQ1iG3yw+SpgxRqKfaURpQD4CZ8Wn0zPMmSnOuTpTh65Vz+EA0rMrw089wvkA==} engines: {node: '>=18.0.0'} - '@smithy/node-http-handler@4.12.0': - resolution: {integrity: sha512-0mq1pHadfyXCYCqm2cNpbjNIT+fbaUpNxewZb/YNr2L0IrEVMOb8gM/Fl4K6XvHCW3uSNDFwPl/+iKm0bx9jYg==} - engines: {node: '>=18.0.0'} - deprecated: Depreacted due to a memory leak issue https://github.com/smithy-lang/smithy-typescript/issues/2258 - '@smithy/node-http-handler@4.12.1': resolution: {integrity: sha512-ThMkboGeONWXAelq9FvGsuJC4rOi+qyC4/zhUF58xYpxUg5sQKx2VXZYJmtNjr4dSuBJ1HeJXETQILCz3wOHvw==} engines: {node: '>=18.0.0'} @@ -1291,10 +1226,6 @@ packages: resolution: {integrity: sha512-7ImGm+FkHRLcBaRttIAMZ6bzJZWb2cJGoYjq46F2UjycujWzrL9GEN9h4w7eQyXJYnltrUhxbbieBAIRrdqpow==} engines: {node: '>=18.0.0'} - '@smithy/types@4.17.2': - resolution: {integrity: sha512-FOKpVZob9MPTn2znRzGrnsMHv7BOsKVw3XiP/cOyYLDVZ9qKp4nifIiSCuUU/fIj5Vu0UOAxCFr+qRAtG0NUkA==} - engines: {node: '>=18.0.0'} - '@smithy/types@4.19.0': resolution: {integrity: sha512-r7jh49VJxGerfAcTQA6gXcKc+98zOp/tqRwzYjgOE+iSQsP6cEU1hq2QzbuipmP68QtYdY9wKEhiCQZIzHgZ4Q==} engines: {node: '>=18.0.0'} @@ -2244,15 +2175,6 @@ packages: resolution: {integrity: sha512-DJnGAeenTdpMEH6uAJRK/uiyEIH9WVsUmoLwzudwGJUwZPp80PDBWPHXSAGNPwNvIXAbe7MSUB1zQFugFml66A==} engines: {node: '>=4'} - redis-semaphore@5.7.0: - resolution: {integrity: sha512-2xjl8eYVVTKX8O8C/H70iWEGx8BSOOzI5P1tQgXEO8A8dR2/BPgU7QctCfcXOHsQ7/IIkebv1ulEdo2iC2SwvA==} - engines: {node: '>= 14.17.0'} - peerDependencies: - ioredis: ^4.1.0 || ^5 - peerDependenciesMeta: - ioredis: - optional: true - redis-semaphore@5.8.0: resolution: {integrity: sha512-I8Y6MexOHOOO1iHQfAgRELH5RExgRX6igevq0qLtLqKHLhycXwoJCaYZE1unLx4fTJXBACkcDkfNoLCsSB9pqg==} engines: {node: '>= 14.17.0'} @@ -2656,22 +2578,10 @@ snapshots: tslib: 2.8.1 '@aws-sdk/client-sqs@3.1134.0': - dependencies: - '@aws-sdk/core': 3.978.0 - '@aws-sdk/credential-provider-node': 3.972.83 - '@aws-sdk/middleware-sdk-sqs': 3.972.42 - '@aws-sdk/types': 3.974.5 - '@smithy/core': 3.33.3 - '@smithy/fetch-http-handler': 5.7.2 - '@smithy/node-http-handler': 4.12.0 - '@smithy/types': 4.17.2 - tslib: 2.8.1 - - '@aws-sdk/client-sts@3.1138.0': dependencies: '@aws-sdk/core': 3.978.1 '@aws-sdk/credential-provider-node': 3.972.84 - '@aws-sdk/signature-v4-multi-region': 3.996.47 + '@aws-sdk/middleware-sdk-sqs': 3.972.42 '@aws-sdk/types': 3.974.6 '@smithy/core': 3.35.0 '@smithy/fetch-http-handler': 5.8.0 @@ -2679,15 +2589,16 @@ snapshots: '@smithy/types': 4.19.0 tslib: 2.8.1 - '@aws-sdk/core@3.978.0': + '@aws-sdk/client-sts@3.1138.0': dependencies: + '@aws-sdk/core': 3.978.1 + '@aws-sdk/credential-provider-node': 3.972.84 + '@aws-sdk/signature-v4-multi-region': 3.996.47 '@aws-sdk/types': 3.974.6 - '@aws-sdk/xml-builder': 3.972.40 - '@aws/lambda-invoke-store': 0.3.0 '@smithy/core': 3.35.0 - '@smithy/signature-v4': 5.7.3 + '@smithy/fetch-http-handler': 5.8.0 + '@smithy/node-http-handler': 4.12.1 '@smithy/types': 4.19.0 - bowser: 2.14.1 tslib: 2.8.1 '@aws-sdk/core@3.978.1': @@ -2701,14 +2612,6 @@ snapshots: bowser: 2.14.1 tslib: 2.8.1 - '@aws-sdk/credential-provider-env@3.972.71': - dependencies: - '@aws-sdk/core': 3.978.1 - '@aws-sdk/types': 3.974.6 - '@smithy/core': 3.35.0 - '@smithy/types': 4.19.0 - tslib: 2.8.1 - '@aws-sdk/credential-provider-env@3.972.72': dependencies: '@aws-sdk/core': 3.978.1 @@ -2717,16 +2620,6 @@ snapshots: '@smithy/types': 4.19.0 tslib: 2.8.1 - '@aws-sdk/credential-provider-http@3.972.73': - dependencies: - '@aws-sdk/core': 3.978.1 - '@aws-sdk/types': 3.974.6 - '@smithy/core': 3.35.0 - '@smithy/fetch-http-handler': 5.8.0 - '@smithy/node-http-handler': 4.12.1 - '@smithy/types': 4.19.0 - tslib: 2.8.1 - '@aws-sdk/credential-provider-http@3.972.74': dependencies: '@aws-sdk/core': 3.978.1 @@ -2737,22 +2630,6 @@ snapshots: '@smithy/types': 4.19.0 tslib: 2.8.1 - '@aws-sdk/credential-provider-ini@3.973.16': - dependencies: - '@aws-sdk/core': 3.978.1 - '@aws-sdk/credential-provider-env': 3.972.71 - '@aws-sdk/credential-provider-http': 3.972.73 - '@aws-sdk/credential-provider-login': 3.972.78 - '@aws-sdk/credential-provider-process': 3.972.71 - '@aws-sdk/credential-provider-sso': 3.973.15 - '@aws-sdk/credential-provider-web-identity': 3.972.77 - '@aws-sdk/nested-clients': 3.997.45 - '@aws-sdk/types': 3.974.6 - '@smithy/core': 3.35.0 - '@smithy/credential-provider-imds': 4.5.2 - '@smithy/types': 4.19.0 - tslib: 2.8.1 - '@aws-sdk/credential-provider-ini@3.973.17': dependencies: '@aws-sdk/core': 3.978.1 @@ -2769,15 +2646,6 @@ snapshots: '@smithy/types': 4.19.0 tslib: 2.8.1 - '@aws-sdk/credential-provider-login@3.972.78': - dependencies: - '@aws-sdk/core': 3.978.1 - '@aws-sdk/nested-clients': 3.997.46 - '@aws-sdk/types': 3.974.6 - '@smithy/core': 3.35.0 - '@smithy/types': 4.19.0 - tslib: 2.8.1 - '@aws-sdk/credential-provider-login@3.972.79': dependencies: '@aws-sdk/core': 3.978.1 @@ -2787,20 +2655,6 @@ snapshots: '@smithy/types': 4.19.0 tslib: 2.8.1 - '@aws-sdk/credential-provider-node@3.972.83': - dependencies: - '@aws-sdk/credential-provider-env': 3.972.71 - '@aws-sdk/credential-provider-http': 3.972.73 - '@aws-sdk/credential-provider-ini': 3.973.16 - '@aws-sdk/credential-provider-process': 3.972.71 - '@aws-sdk/credential-provider-sso': 3.973.15 - '@aws-sdk/credential-provider-web-identity': 3.972.77 - '@aws-sdk/types': 3.974.6 - '@smithy/core': 3.35.0 - '@smithy/credential-provider-imds': 4.5.2 - '@smithy/types': 4.19.0 - tslib: 2.8.1 - '@aws-sdk/credential-provider-node@3.972.84': dependencies: '@aws-sdk/credential-provider-env': 3.972.72 @@ -2815,14 +2669,6 @@ snapshots: '@smithy/types': 4.19.0 tslib: 2.8.1 - '@aws-sdk/credential-provider-process@3.972.71': - dependencies: - '@aws-sdk/core': 3.978.1 - '@aws-sdk/types': 3.974.6 - '@smithy/core': 3.35.0 - '@smithy/types': 4.19.0 - tslib: 2.8.1 - '@aws-sdk/credential-provider-process@3.972.72': dependencies: '@aws-sdk/core': 3.978.1 @@ -2831,16 +2677,6 @@ snapshots: '@smithy/types': 4.19.0 tslib: 2.8.1 - '@aws-sdk/credential-provider-sso@3.973.15': - dependencies: - '@aws-sdk/core': 3.978.1 - '@aws-sdk/nested-clients': 3.997.45 - '@aws-sdk/token-providers': 3.1129.0 - '@aws-sdk/types': 3.974.6 - '@smithy/core': 3.35.0 - '@smithy/types': 4.19.0 - tslib: 2.8.1 - '@aws-sdk/credential-provider-sso@3.973.16': dependencies: '@aws-sdk/core': 3.978.1 @@ -2851,15 +2687,6 @@ snapshots: '@smithy/types': 4.19.0 tslib: 2.8.1 - '@aws-sdk/credential-provider-web-identity@3.972.77': - dependencies: - '@aws-sdk/core': 3.978.1 - '@aws-sdk/nested-clients': 3.997.45 - '@aws-sdk/types': 3.974.6 - '@smithy/core': 3.35.0 - '@smithy/types': 4.19.0 - tslib: 2.8.1 - '@aws-sdk/credential-provider-web-identity@3.972.78': dependencies: '@aws-sdk/core': 3.978.1 @@ -2885,17 +2712,6 @@ snapshots: '@smithy/types': 4.19.0 tslib: 2.8.1 - '@aws-sdk/nested-clients@3.997.45': - dependencies: - '@aws-sdk/core': 3.978.1 - '@aws-sdk/signature-v4-multi-region': 3.996.47 - '@aws-sdk/types': 3.974.6 - '@smithy/core': 3.35.0 - '@smithy/fetch-http-handler': 5.8.0 - '@smithy/node-http-handler': 4.12.1 - '@smithy/types': 4.19.0 - tslib: 2.8.1 - '@aws-sdk/nested-clients@3.997.46': dependencies: '@aws-sdk/core': 3.978.1 @@ -2914,15 +2730,6 @@ snapshots: '@smithy/types': 4.19.0 tslib: 2.8.1 - '@aws-sdk/token-providers@3.1129.0': - dependencies: - '@aws-sdk/core': 3.978.1 - '@aws-sdk/nested-clients': 3.997.46 - '@aws-sdk/types': 3.974.6 - '@smithy/core': 3.35.0 - '@smithy/types': 4.19.0 - tslib: 2.8.1 - '@aws-sdk/token-providers@3.1138.0': dependencies: '@aws-sdk/core': 3.978.1 @@ -2932,21 +2739,11 @@ snapshots: '@smithy/types': 4.19.0 tslib: 2.8.1 - '@aws-sdk/types@3.974.5': - dependencies: - '@smithy/types': 4.19.0 - tslib: 2.8.1 - '@aws-sdk/types@3.974.6': dependencies: '@smithy/types': 4.19.0 tslib: 2.8.1 - '@aws-sdk/xml-builder@3.972.40': - dependencies: - '@smithy/types': 4.19.0 - tslib: 2.8.1 - '@aws-sdk/xml-builder@3.972.41': dependencies: '@smithy/types': 4.19.0 @@ -3175,7 +2972,7 @@ snapshots: '@lokalise/node-core': 14.9.1(zod@4.6.5) '@message-queue-toolkit/core': link:packages/core ioredis: 5.11.1(supports-color@7.2.0) - redis-semaphore: 5.7.0(ioredis@5.11.1(supports-color@7.2.0))(supports-color@7.2.0) + redis-semaphore: 5.8.0(ioredis@5.11.1(supports-color@7.2.0))(supports-color@7.2.0) transitivePeerDependencies: - supports-color - zod @@ -3393,11 +3190,6 @@ snapshots: '@rolldown/pluginutils@1.0.1': {} - '@smithy/core@3.33.3': - dependencies: - '@smithy/types': 4.19.0 - tslib: 2.8.1 - '@smithy/core@3.35.0': dependencies: '@smithy/types': 4.19.0 @@ -3409,24 +3201,12 @@ snapshots: '@smithy/types': 4.19.0 tslib: 2.8.1 - '@smithy/fetch-http-handler@5.7.2': - dependencies: - '@smithy/core': 3.35.0 - '@smithy/types': 4.19.0 - tslib: 2.8.1 - '@smithy/fetch-http-handler@5.8.0': dependencies: '@smithy/core': 3.35.0 '@smithy/types': 4.19.0 tslib: 2.8.1 - '@smithy/node-http-handler@4.12.0': - dependencies: - '@smithy/core': 3.35.0 - '@smithy/types': 4.19.0 - tslib: 2.8.1 - '@smithy/node-http-handler@4.12.1': dependencies: '@smithy/core': 3.35.0 @@ -3439,10 +3219,6 @@ snapshots: '@smithy/types': 4.19.0 tslib: 2.8.1 - '@smithy/types@4.17.2': - dependencies: - tslib: 2.8.1 - '@smithy/types@4.19.0': dependencies: tslib: 2.8.1 @@ -4346,7 +4122,7 @@ snapshots: dependencies: redis-errors: 1.2.0 - redis-semaphore@5.7.0(ioredis@5.11.1(supports-color@7.2.0))(supports-color@7.2.0): + redis-semaphore@5.8.0(ioredis@5.11.1(supports-color@7.2.0))(supports-color@7.2.0): dependencies: debug: 4.4.3(supports-color@7.2.0) optionalDependencies: