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
6 changes: 2 additions & 4 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -474,7 +473,6 @@ const result = await initSnsSqs(
{
topicArn: '...',
queueUrl: '...',
subscriptionArn: '...',
startupResourcePolling: {
enabled: true,
timeoutMs: 5 * 60 * 1000,
Expand All @@ -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}`)
},
},
Expand Down
134 changes: 102 additions & 32 deletions packages/sns/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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',
},
}
```
Expand All @@ -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
Expand All @@ -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).
Expand All @@ -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: {
Expand All @@ -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,
Expand All @@ -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,
Expand Down Expand Up @@ -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,
Expand All @@ -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:
Expand All @@ -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, {
Expand All @@ -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()
```

Expand Down Expand Up @@ -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: {
Expand All @@ -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,
Expand Down Expand Up @@ -1382,14 +1437,29 @@ type SNSSQSConsumerDependencies = SNSDependencies & SQSDependencies & {
}

// Subscription options
type SNSSubscriptionOptions = {
updateAttributesIfExists?: boolean
filterPolicy?: Record<string, string[]>
rawMessageDelivery?: boolean
redrivePolicy?: {
deadLetterTargetArn: string
}
}
type SNSSubscriptionOptions =
// Creation: the subscription is created, or updated if it already exists
| {
updateAttributesIfExists?: boolean
managedAttributes?: SubscriptionManagedAttributeName[]
filterPolicy?: Record<string, string[]>
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<string, string>
}

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

### Utility Functions
Expand Down
7 changes: 6 additions & 1 deletion packages/sns/lib/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
Loading
Loading