From f2c4e1fbbaa9d1b70f2d14dceb37707d5413e297 Mon Sep 17 00:00:00 2001 From: Igor Savin Date: Thu, 1 Oct 2026 18:14:12 +0300 Subject: [PATCH 1/4] feat(sns): check subscription attributes before writing them assertSubscription now looks up an existing subscription and reads its attributes before doing anything. When they already match the config, no Subscribe or SetSubscriptionAttributes call is made. Otherwise only the differing attributes are written (or an error is thrown when updateAttributesIfExists is off). The locateOnly path gets the same read-first behaviour through setSubscriptionAttributes. Adds subscriptionConfig.manageOnlyFilterPolicy, which limits checking and writing to FilterPolicy and FilterPolicyScope. --- packages/sns/README.md | 12 +- packages/sns/lib/utils/snsInitter.ts | 14 +- packages/sns/lib/utils/snsSubscriber.spec.ts | 151 +++++++++++++++- packages/sns/lib/utils/snsSubscriber.ts | 176 +++++++++++++------ packages/sns/lib/utils/snsUtils.ts | 12 ++ 5 files changed, 293 insertions(+), 72 deletions(-) diff --git a/packages/sns/README.md b/packages/sns/README.md index 10c23b512..05f9d17b4 100644 --- a/packages/sns/README.md +++ b/packages/sns/README.md @@ -506,12 +506,16 @@ your application or by external tooling (e.g. Terraform): | 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, 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 +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. Set `subscriptionConfig.manageOnlyFilterPolicy` to manage only `FilterPolicy` and `FilterPolicyScope`: other +attributes such as `RawMessageDelivery` or `RedrivePolicy` are then neither checked nor written, and are left to +external tooling. When both a locator and a creation config are given for the same resource, the locator takes precedence and the creation config is ignored. With `subscriptionConfig: { locateOnly: true, Attributes }` the subscription is located, never created, and the given -`Attributes` (e.g. `FilterPolicy`) are applied to it on startup. This lets external tooling own the subscription while +`Attributes` (e.g. `FilterPolicy`) are applied to it on startup when they differ from the current ones. 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. @@ -824,6 +828,7 @@ SNS consumers use the same options as SQS consumers, plus SNS-specific subscript // SNS-Specific - Subscription Configuration subscriptionConfig: { updateAttributesIfExists: false, // Update subscription attributes if exists + manageOnlyFilterPolicy: false, // Only check and update FilterPolicy and FilterPolicyScope // Optional: Message filtering filterPolicy: { @@ -1407,6 +1412,7 @@ type SNSSubscriptionOptions = // Creation: the subscription is created, or updated if it already exists | { updateAttributesIfExists?: boolean + manageOnlyFilterPolicy?: boolean filterPolicy?: Record rawMessageDelivery?: boolean redrivePolicy?: { diff --git a/packages/sns/lib/utils/snsInitter.ts b/packages/sns/lib/utils/snsInitter.ts index c629626ee..3550b43d5 100644 --- a/packages/sns/lib/utils/snsInitter.ts +++ b/packages/sns/lib/utils/snsInitter.ts @@ -34,7 +34,7 @@ import { assertTopic, deleteSubscription, deleteTopic, - findSubscriptionByTopicAndQueue, + findConfirmedSubscriptionArn, getTopicAttributes, } from './snsUtils.ts' import { buildTopicArn } from './stsUtils.ts' @@ -231,18 +231,6 @@ 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 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 diff --git a/packages/sns/lib/utils/snsSubscriber.spec.ts b/packages/sns/lib/utils/snsSubscriber.spec.ts index 57560c202..7150bb7ea 100644 --- a/packages/sns/lib/utils/snsSubscriber.spec.ts +++ b/packages/sns/lib/utils/snsSubscriber.spec.ts @@ -1,13 +1,17 @@ -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 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 +85,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 +236,118 @@ 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('ignores attributes other than the filter policy when manageOnlyFilterPolicy is enabled', 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"]}`, RawMessageDelivery: 'false' }, + updateAttributesIfExists: false, + manageOnlyFilterPolicy: true, + }, + { 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('throws when attributes are different and update is disabled', async () => { const logger = new FakeLogger() const topicArn = await testAdmin.createTopic(TOPIC_NAME) @@ -270,4 +384,27 @@ 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', + }) + + 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 56df6d3f9..b06556c12 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,7 +18,7 @@ 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' /** @@ -32,6 +33,12 @@ export type SNSSubscriptionOptions = /** Creation */ | (Omit & { updateAttributesIfExists: boolean + /** + * When enabled, only `FilterPolicy` and `FilterPolicyScope` are managed. Other subscription + * attributes (e.g. `RawMessageDelivery`, `RedrivePolicy`) are neither checked nor written, + * and are expected to be set by external tooling such as Terraform. + */ + manageOnlyFilterPolicy?: boolean /** Should not be present when the subscription is created */ locateOnly?: never }) @@ -116,9 +123,10 @@ 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 attributes differ from the + * configured ones and `updateAttributesIfExists` is enabled. Only the differing attributes are + * written. */ export async function assertSubscription( snsClient: SNSClient, @@ -128,42 +136,65 @@ export async function assertSubscription( errorContext: { queueName?: string; topicName?: string }, logger?: CommonLogger, ): Promise { - const subscribeCommand = new SubscribeCommand({ - TopicArn: topicArn, - Endpoint: queueArn, - Protocol: 'sqs', - ReturnSubscriptionArn: true, - ...subscriptionConfiguration, - }) + const { updateAttributesIfExists, manageOnlyFilterPolicy, locateOnly, ...subscribeInput } = + subscriptionConfiguration + const attributes = resolveManagedAttributes(subscribeInput.Attributes, manageOnlyFilterPolicy) + 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, + attributes, + ) + 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, + Attributes: attributes, + 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, @@ -173,34 +204,28 @@ export async function assertSubscription( } } -async function tryToUpdateSubscription( +/** + * Applies the given attributes to an existing subscription. Current attributes are read first, and + * only the ones that differ are written. + */ +export async function setSubscriptionAttributes( snsClient: SNSClient, - topicArn: string, - queueArn: string, - subscriptionConfiguration: SNSSubscriptionCreationOptions, -) { - const subscription = await findSubscriptionByTopicAndQueue(snsClient, topicArn, queueArn) - if (!subscription?.SubscriptionArn || !subscriptionConfiguration.Attributes) { - return undefined - } - - await setSubscriptionAttributes( + subscriptionArn: string, + attributes: NonNullable, +): Promise { + const changedAttributes = await findChangedSubscriptionAttributes( snsClient, - subscription.SubscriptionArn, - subscriptionConfiguration.Attributes, + subscriptionArn, + attributes, ) - - return subscription + await writeSubscriptionAttributes(snsClient, subscriptionArn, changedAttributes) } -/** - * Applies the given attributes to an existing subscription, one by one as SNS only allows setting a single - * attribute per call. - */ -export async function setSubscriptionAttributes( +// SNS only allows setting a single attribute per call +async function writeSubscriptionAttributes( snsClient: SNSClient, subscriptionArn: string, - attributes: NonNullable, + attributes: Record, ): Promise { for (const [key, value] of Object.entries(attributes)) { await snsClient.send( @@ -212,3 +237,56 @@ export async function setSubscriptionAttributes( ) } } + +const FILTER_POLICY_ATTRIBUTES = ['FilterPolicy', 'FilterPolicyScope'] + +// Values SNS reports for attributes that were never set explicitly +const SUBSCRIPTION_ATTRIBUTE_DEFAULTS: Record = { + FilterPolicyScope: 'MessageAttributes', + RawMessageDelivery: 'false', +} + +function resolveManagedAttributes( + attributes: Record | undefined, + manageOnlyFilterPolicy: boolean | undefined, +): Record | undefined { + if (!attributes || !manageOnlyFilterPolicy) return attributes + + return Object.fromEntries( + Object.entries(attributes).filter(([name]) => FILTER_POLICY_ATTRIBUTES.includes(name)), + ) +} + +async function findChangedSubscriptionAttributes( + snsClient: SNSClient, + subscriptionArn: string, + attributes: Record | undefined, +): Promise> { + if (!attributes || Object.keys(attributes).length === 0) return {} + + const currentAttributesResult = await getSubscriptionAttributes(snsClient, subscriptionArn) + const currentAttributes = currentAttributesResult.result?.attributes ?? {} + + return Object.fromEntries( + Object.entries(attributes).filter( + ([name, value]) => + !isSameAttributeValue( + name, + value, + currentAttributes[name] ?? SUBSCRIPTION_ATTRIBUTE_DEFAULTS[name], + ), + ), + ) +} + +// 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 (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. * From 9a087cc676943406bd95e2592517f677d2195b23 Mon Sep 17 00:00:00 2001 From: Igor Savin Date: Thu, 1 Oct 2026 18:22:26 +0300 Subject: [PATCH 2/4] feat(sns): replace manageOnlyFilterPolicy with typed managedAttributes subscriptionConfig.managedAttributes lists the subscription attributes the application owns, defaulting to FilterPolicy, FilterPolicyScope, RawMessageDelivery and RedrivePolicy. A managed attribute missing from Attributes is now reset on the existing subscription, unmanaged ones are never read or written, and configuring an unmanaged one throws. The consumer DLQ RedrivePolicy write goes through the read-first setSubscriptionAttributes and RedrivePolicy is excluded from the managed set when reuseConsumerDeadLetterQueue is on, so the two never fight over it. --- packages/sns/README.md | 38 ++++- packages/sns/lib/index.ts | 2 + .../sns/lib/sns/AbstractSnsSqsConsumer.ts | 49 ++++-- packages/sns/lib/utils/snsInitter.ts | 9 +- packages/sns/lib/utils/snsSubscriber.spec.ts | 85 +++++++++- packages/sns/lib/utils/snsSubscriber.ts | 146 +++++++++++++----- 6 files changed, 263 insertions(+), 66 deletions(-) diff --git a/packages/sns/README.md b/packages/sns/README.md index 05f9d17b4..b138961e4 100644 --- a/packages/sns/README.md +++ b/packages/sns/README.md @@ -509,13 +509,30 @@ your application or by external tooling (e.g. Terraform): 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. Set `subscriptionConfig.manageOnlyFilterPolicy` to manage only `FilterPolicy` and `FilterPolicyScope`: other -attributes such as `RawMessageDelivery` or `RedrivePolicy` are then neither checked nor written, and are left to -external tooling. When both a locator and a creation config are given for the -same resource, the locator takes precedence and the creation config is ignored. +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. + +```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 the given -`Attributes` (e.g. `FilterPolicy`) are applied to it on startup when they differ from the current ones. This lets external tooling own the subscription while +managed attributes (e.g. `FilterPolicy`) are applied to it on startup when they differ from the current ones. 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. @@ -828,7 +845,7 @@ SNS consumers use the same options as SQS consumers, plus SNS-specific subscript // SNS-Specific - Subscription Configuration subscriptionConfig: { updateAttributesIfExists: false, // Update subscription attributes if exists - manageOnlyFilterPolicy: false, // Only check and update FilterPolicy and FilterPolicyScope + managedAttributes: ['FilterPolicy', 'FilterPolicyScope'], // Defaults to all supported attributes // Optional: Message filtering filterPolicy: { @@ -1412,7 +1429,7 @@ type SNSSubscriptionOptions = // Creation: the subscription is created, or updated if it already exists | { updateAttributesIfExists?: boolean - manageOnlyFilterPolicy?: boolean + managedAttributes?: SubscriptionManagedAttributeName[] filterPolicy?: Record rawMessageDelivery?: boolean redrivePolicy?: { @@ -1422,8 +1439,15 @@ type SNSSubscriptionOptions = // 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 ce91fd42a..784d01956 100644 --- a/packages/sns/lib/index.ts +++ b/packages/sns/lib/index.ts @@ -36,6 +36,8 @@ export { deserializeSNSMessage } from './utils/snsMessageDeserializer.ts' export { readSnsMessage } from './utils/snsMessageReader.ts' export { type SNSSubscriptionOptions, + SUBSCRIPTION_MANAGED_ATTRIBUTE_NAMES, + type SubscriptionManagedAttributeName, subscribeToTopic, } from './utils/snsSubscriber.ts' export { diff --git a/packages/sns/lib/sns/AbstractSnsSqsConsumer.ts b/packages/sns/lib/sns/AbstractSnsSqsConsumer.ts index 1718354d4..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 & { @@ -120,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({ @@ -262,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'], ) } @@ -328,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.ts b/packages/sns/lib/utils/snsInitter.ts index 3550b43d5..d31ad175c 100644 --- a/packages/sns/lib/utils/snsInitter.ts +++ b/packages/sns/lib/utils/snsInitter.ts @@ -277,8 +277,13 @@ async function resolveSubscriptionArn( options, )) - if (subscriptionConfig?.locateOnly && subscriptionConfig?.Attributes) { - await setSubscriptionAttributes(snsClient, subscriptionArn, subscriptionConfig.Attributes) + if (subscriptionConfig?.locateOnly) { + await setSubscriptionAttributes( + snsClient, + subscriptionArn, + subscriptionConfig.Attributes, + subscriptionConfig.managedAttributes, + ) } return subscriptionArn diff --git a/packages/sns/lib/utils/snsSubscriber.spec.ts b/packages/sns/lib/utils/snsSubscriber.spec.ts index 7150bb7ea..666d606e0 100644 --- a/packages/sns/lib/utils/snsSubscriber.spec.ts +++ b/packages/sns/lib/utils/snsSubscriber.spec.ts @@ -314,7 +314,7 @@ describe('snsSubscriber', () => { expect(subscriptionAttributes.result?.attributes?.FilterPolicy).toBe(`{"type":["add"]}`) }) - it('ignores attributes other than the filter policy when manageOnlyFilterPolicy is enabled', async () => { + 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( @@ -334,9 +334,9 @@ describe('snsSubscriber', () => { topicArn, queueArn, { - Attributes: { FilterPolicy: `{"type":["add"]}`, RawMessageDelivery: 'false' }, + Attributes: { FilterPolicy: `{"type":["add"]}` }, + managedAttributes: ['FilterPolicy', 'FilterPolicyScope'], updateAttributesIfExists: false, - manageOnlyFilterPolicy: true, }, { queueName: QUEUE_NAME, topicName: TOPIC_NAME }, ) @@ -348,6 +348,75 @@ describe('snsSubscriber', () => { 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', + RedrivePolicy: JSON.stringify({ deadLetterTargetArn: queueArn }), + }, + 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', + }) + expect(subscriptionAttributes.result?.attributes?.RedrivePolicy || undefined).toBeUndefined() + + // 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) + }) + + 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) @@ -395,10 +464,12 @@ describe('snsSubscriber', () => { }) const sendSpy = vi.spyOn(snsClient, 'send') - await setSubscriptionAttributes(snsClient, subscriptionArn, { - FilterPolicy: `{ "type": ["add"] }`, - RawMessageDelivery: 'true', - }) + await setSubscriptionAttributes( + snsClient, + subscriptionArn, + { FilterPolicy: `{ "type": ["add"] }`, RawMessageDelivery: 'true' }, + ['FilterPolicy', 'RawMessageDelivery'], + ) const writtenAttributes = sendSpy.mock.calls .map(([command]) => command) diff --git a/packages/sns/lib/utils/snsSubscriber.ts b/packages/sns/lib/utils/snsSubscriber.ts index b06556c12..8fb1d5002 100644 --- a/packages/sns/lib/utils/snsSubscriber.ts +++ b/packages/sns/lib/utils/snsSubscriber.ts @@ -21,34 +21,52 @@ import { import { assertTopic, findConfirmedSubscriptionArn, getSubscriptionAttributes } from './snsUtils.ts' import { buildTopicArn } from './stsUtils.ts' +/** 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 the given attributes (e.g. `FilterPolicy`) are applied to it. + * and its managed attributes (e.g. `FilterPolicy`) are applied to it. */ export type SNSSubscriptionOptions = /** Creation */ - | (Omit & { - updateAttributesIfExists: boolean - /** - * When enabled, only `FilterPolicy` and `FilterPolicyScope` are managed. Other subscription - * attributes (e.g. `RawMessageDelivery`, `RedrivePolicy`) are neither checked nor written, - * and are expected to be set by external tooling such as Terraform. - */ - manageOnlyFilterPolicy?: boolean - /** Should not be present when the subscription is created */ - locateOnly?: never - }) + | (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 @@ -124,7 +142,7 @@ export async function subscribeToTopic( /** * 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 attributes differ from the + * 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. */ @@ -136,9 +154,12 @@ export async function assertSubscription( errorContext: { queueName?: string; topicName?: string }, logger?: CommonLogger, ): Promise { - const { updateAttributesIfExists, manageOnlyFilterPolicy, locateOnly, ...subscribeInput } = + const { updateAttributesIfExists, managedAttributes, locateOnly, ...subscribeInput } = subscriptionConfiguration - const attributes = resolveManagedAttributes(subscribeInput.Attributes, manageOnlyFilterPolicy) + const resolvedManagedAttributes = resolveManagedAttributes( + subscribeInput.Attributes, + managedAttributes, + ) const resolvedLogger = logger ?? console const errMessagePrefix = `Error while creating subscription for queue "${errorContext.queueName}", topic "${errorContext.topicName}"` @@ -146,7 +167,8 @@ export async function assertSubscription( const changedAttributes = await findChangedSubscriptionAttributes( snsClient, subscriptionArn, - attributes, + subscribeInput.Attributes, + resolvedManagedAttributes, ) if (Object.keys(changedAttributes).length === 0) return subscriptionArn @@ -172,7 +194,6 @@ export async function assertSubscription( const subscriptionResult = await snsClient.send( new SubscribeCommand({ ...subscribeInput, - Attributes: attributes, TopicArn: topicArn, Endpoint: queueArn, Protocol: 'sqs', @@ -205,18 +226,20 @@ export async function assertSubscription( } /** - * Applies the given attributes to an existing subscription. Current attributes are read first, and - * only the ones that differ are written. + * 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, subscriptionArn: string, - attributes: NonNullable, + attributes: SubscribeCommandInput['Attributes'], + managedAttributes?: readonly SubscriptionManagedAttributeName[], ): Promise { const changedAttributes = await findChangedSubscriptionAttributes( snsClient, subscriptionArn, attributes, + resolveManagedAttributes(attributes, managedAttributes), ) await writeSubscriptionAttributes(snsClient, subscriptionArn, changedAttributes) } @@ -238,50 +261,91 @@ async function writeSubscriptionAttributes( } } -const FILTER_POLICY_ATTRIBUTES = ['FilterPolicy', 'FilterPolicyScope'] - -// Values SNS reports for attributes that were never set explicitly -const SUBSCRIPTION_ATTRIBUTE_DEFAULTS: Record = { +// 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: '', +} + +function isManagedAttributeName(name: string): name is SubscriptionManagedAttributeName { + return (SUBSCRIPTION_MANAGED_ATTRIBUTE_NAMES as readonly string[]).includes(name) } function resolveManagedAttributes( attributes: Record | undefined, - manageOnlyFilterPolicy: boolean | undefined, -): Record | undefined { - if (!attributes || !manageOnlyFilterPolicy) return attributes - - return Object.fromEntries( - Object.entries(attributes).filter(([name]) => FILTER_POLICY_ATTRIBUTES.includes(name)), + 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> { - if (!attributes || Object.keys(attributes).length === 0) return {} + 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 {} const currentAttributesResult = await getSubscriptionAttributes(snsClient, subscriptionArn) const currentAttributes = currentAttributesResult.result?.attributes ?? {} - return Object.fromEntries( - Object.entries(attributes).filter( - ([name, value]) => - !isSameAttributeValue( - name, - value, - currentAttributes[name] ?? SUBSCRIPTION_ATTRIBUTE_DEFAULTS[name], - ), + 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 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 { From 8aaaea29cc08bed61ebe5c758605131e388014f9 Mon Sep 17 00:00:00 2001 From: Igor Savin Date: Thu, 1 Oct 2026 18:30:35 +0300 Subject: [PATCH 3/4] test(sns): skip redrive policy removal on localstack LocalStack rejects any RedrivePolicy that is not a policy with a valid deadLetterTargetArn, so it cannot remove one. The reset test no longer sets RedrivePolicy, and its removal is covered by a separate test that runs on fauxqs only. --- packages/sns/lib/utils/snsSubscriber.spec.ts | 36 ++++++++++++++++++-- packages/sns/test/utils/fauxqsInstance.ts | 2 +- 2 files changed, 35 insertions(+), 3 deletions(-) diff --git a/packages/sns/lib/utils/snsSubscriber.spec.ts b/packages/sns/lib/utils/snsSubscriber.spec.ts index 666d606e0..3a0fcfc2d 100644 --- a/packages/sns/lib/utils/snsSubscriber.spec.ts +++ b/packages/sns/lib/utils/snsSubscriber.spec.ts @@ -8,6 +8,7 @@ import type { STSClient } from '@aws-sdk/client-sts' import type { AwilixContainer } from 'awilix' 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' @@ -360,7 +361,6 @@ describe('snsSubscriber', () => { FilterPolicy: `{"type":["add"]}`, FilterPolicyScope: 'MessageBody', RawMessageDelivery: 'true', - RedrivePolicy: JSON.stringify({ deadLetterTargetArn: queueArn }), }, updateAttributesIfExists: false, }, @@ -382,7 +382,6 @@ describe('snsSubscriber', () => { FilterPolicyScope: 'MessageAttributes', RawMessageDelivery: 'false', }) - expect(subscriptionAttributes.result?.attributes?.RedrivePolicy || undefined).toBeUndefined() // A subscription that is already reset is not written to again const sendSpy = vi.spyOn(snsClient, 'send') @@ -398,6 +397,39 @@ describe('snsSubscriber', () => { ).toBe(false) }) + // LocalStack rejects any RedrivePolicy value that is not a policy with a valid deadLetterTargetArn + 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 }, + ) + + await assertSubscription( + snsClient, + topicArn, + queueArn, + { updateAttributesIfExists: true }, + { queueName: QUEUE_NAME, topicName: TOPIC_NAME }, + new FakeLogger(), + ) + + 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) 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 From a5735a3fd322ee6acbef857d71919d1f022a662c Mon Sep 17 00:00:00 2001 From: Igor Savin Date: Thu, 1 Oct 2026 18:41:37 +0300 Subject: [PATCH 4/4] fix(sns): remove RedrivePolicy by omitting its value Real AWS rejects SetSubscriptionAttributes with an empty RedrivePolicy; the attribute is removed by leaving AttributeValue out, which is what the Terraform AWS provider does. FilterPolicy is reset with "{}" as in the SNS docs. --- packages/sns/lib/utils/snsSubscriber.spec.ts | 15 ++++++++++++++- packages/sns/lib/utils/snsSubscriber.ts | 5 +++-- 2 files changed, 17 insertions(+), 3 deletions(-) diff --git a/packages/sns/lib/utils/snsSubscriber.spec.ts b/packages/sns/lib/utils/snsSubscriber.spec.ts index 3a0fcfc2d..181319be8 100644 --- a/packages/sns/lib/utils/snsSubscriber.spec.ts +++ b/packages/sns/lib/utils/snsSubscriber.spec.ts @@ -397,7 +397,7 @@ describe('snsSubscriber', () => { ).toBe(false) }) - // LocalStack rejects any RedrivePolicy value that is not a policy with a valid deadLetterTargetArn + // 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 () => { @@ -413,6 +413,7 @@ describe('snsSubscriber', () => { }, { queueName: QUEUE_NAME, topicName: TOPIC_NAME }, ) + const sendSpy = vi.spyOn(snsClient, 'send') await assertSubscription( snsClient, @@ -423,6 +424,18 @@ describe('snsSubscriber', () => { 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, diff --git a/packages/sns/lib/utils/snsSubscriber.ts b/packages/sns/lib/utils/snsSubscriber.ts index 8fb1d5002..a2e86a29b 100644 --- a/packages/sns/lib/utils/snsSubscriber.ts +++ b/packages/sns/lib/utils/snsSubscriber.ts @@ -255,7 +255,8 @@ async function writeSubscriptionAttributes( new SetSubscriptionAttributesCommand({ SubscriptionArn: subscriptionArn, AttributeName: key, - AttributeValue: value, + // AWS rejects an empty RedrivePolicy, omitting the value is what removes it + AttributeValue: key === 'RedrivePolicy' && value === '' ? undefined : value, }), ) } @@ -263,7 +264,7 @@ async function writeSubscriptionAttributes( // Values that reset an attribute to the state of a subscription where it was never set const SUBSCRIPTION_ATTRIBUTE_UNSET_VALUES: Record = { - FilterPolicy: '', + FilterPolicy: '{}', FilterPolicyScope: 'MessageAttributes', RawMessageDelivery: 'false', RedrivePolicy: '',