Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
38 changes: 34 additions & 4 deletions packages/sns/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -506,12 +506,33 @@ 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
same resource, the locator takes precedence and the creation config is ignored.
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.

```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. 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.

Expand Down Expand Up @@ -824,6 +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
managedAttributes: ['FilterPolicy', 'FilterPolicyScope'], // Defaults to all supported attributes

// Optional: Message filtering
filterPolicy: {
Expand Down Expand Up @@ -1407,6 +1429,7 @@ type SNSSubscriptionOptions =
// Creation: the subscription is created, or updated if it already exists
| {
updateAttributesIfExists?: boolean
managedAttributes?: SubscriptionManagedAttributeName[]
filterPolicy?: Record<string, string[]>
rawMessageDelivery?: boolean
redrivePolicy?: {
Expand All @@ -1416,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<string, string>
}

type SubscriptionManagedAttributeName =
| 'FilterPolicy'
| 'FilterPolicyScope'
| 'RawMessageDelivery'
| 'RedrivePolicy'
```

### Utility Functions
Expand Down
2 changes: 2 additions & 0 deletions packages/sns/lib/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
49 changes: 40 additions & 9 deletions packages/sns/lib/sns/AbstractSnsSqsConsumer.ts
Original file line number Diff line number Diff line change
@@ -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 {
Expand All @@ -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 & {
Expand Down Expand Up @@ -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({
Expand Down Expand Up @@ -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'],
)
}

Expand Down Expand Up @@ -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'),
}
}
23 changes: 8 additions & 15 deletions packages/sns/lib/utils/snsInitter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -34,7 +34,7 @@ import {
assertTopic,
deleteSubscription,
deleteTopic,
findSubscriptionByTopicAndQueue,
findConfirmedSubscriptionArn,
getTopicAttributes,
} from './snsUtils.ts'
import { buildTopicArn } from './stsUtils.ts'
Expand Down Expand Up @@ -231,18 +231,6 @@ async function resolveQueue(
)
}

async function findConfirmedSubscriptionArn(
snsClient: SNSClient,
topicArn: string,
queueArn: string,
): Promise<string | undefined> {
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
Expand Down Expand Up @@ -289,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
Expand Down
Loading
Loading