Skip to content
Open
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
5 changes: 5 additions & 0 deletions .changeset/container-context-delivery.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
'@node-ts/bus-core': minor
---

`ContainerAdapter.get(type, context)` is also given the delivery being handled, as `context.transportMessage`, so a container can scope what it resolves to one delivery of a message (#353). Every class handler and workflow that handles a delivery is given the same `transportMessage`, and each retry or send of the message a new one. Key per-message scopes on it rather than on `context.message`, which `InMemoryQueue` hands out again on a retry and which can be sent more than once. The context's type is exported as `ContainerContext`. A custom transport or `Receiver` must return a new `TransportMessage` for each delivery, as every transport in this repo does.
2 changes: 1 addition & 1 deletion .changeset/nestjs-module.md
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,7 @@
Add `@node-ts/bus-nestjs`, a NestJS module for the bus (#268), for NestJS 11 and 12.

- `BusModule.forRoot({ configure })` and `BusModule.forRootAsync({ imports, inject, useFactory })` register a bus, configured by a factory that's given a `BusConfiguration` (logging through Nest's `Logger`, with no interrupt signals) and returns it with the transport, persistence and message types. The bus is injected as `BusInstance`; named buses with `@InjectBus(name)` and `getBusToken(name)`.
- Class handlers and workflows are providers of your modules, registered with the bus by `@BusHandler()` and `@BusWorkflow()` or by listing them in `BusModule.forFeature()`, and resolved from Nest's container for each message by `nestContainer(moduleRef)`. Function handlers and `defineWorkflow()` workflows are registered with `BusModule.forFeatureAsync()`, given the providers they close over. Request-scoped providers get a request scope per received message, with the message and its attributes as Nest's `REQUEST` (`BusRequest`).
- Class handlers and workflows are providers of your modules, registered with the bus by `@BusHandler()` and `@BusWorkflow()` or by listing them in `BusModule.forFeature()`, and resolved from Nest's container for each message by `nestContainer(moduleRef)`. Function handlers and `defineWorkflow()` workflows are registered with `BusModule.forFeatureAsync()`, given the providers they close over. Request-scoped providers get a request scope per delivery of a message, so a retry gets new ones (#353), with the message and its attributes as Nest's `REQUEST` (`BusRequest`).
- The bus is built when the application initializes, and initialized and started when it bootstraps (or by the application, with `lifecycle: 'manual'`); if that fails, it's disposed. Since `BusModule` is global, Nest stops it after the `onModuleDestroy` of every module that isn't global, and disposes it after their `onApplicationShutdown`, so release what handlers use in `beforeApplicationShutdown` or `onApplicationShutdown`. Startup fails with `BusClassNotProvided`, `BusNotRegistered`, `BusAlreadyRegistered`, `BusFeatureNotStatic`, `WorkflowResolvedWithoutMessage` or `BusCoreVersionNotSupported`, naming the fix.
- `createBusForProvisioning(AppModule)` returns the application's bus, built but not initialized, for `bus provision`.

Expand Down
2 changes: 1 addition & 1 deletion CLAUDE.md
Original file line number Diff line number Diff line change
Expand Up @@ -72,7 +72,7 @@ pnpm exec dotenv -e test.env -- jest packages/bus-sqs/src/sqs-transport.spec.ts

### Handlers

- Function handlers via `handlerFor(MessageClass | definition, fn)`; class handlers implement `Handler` and are resolved from the `ContainerAdapter` (`withContainer`) when there is one, or constructed with `new` otherwise, like class workflows (`build()` throws `ContainerNotRegistered` for a constructor with arguments). `handlerFor<TMessage, TAttributes>` types the attributes; the third generic keeps the handler's own type so it can be called directly in tests.
- Function handlers via `handlerFor(MessageClass | definition, fn)`; class handlers implement `Handler` and are resolved from the `ContainerAdapter` (`withContainer`) when there is one, or constructed with `new` otherwise, like class workflows (`build()` throws `ContainerNotRegistered` for a constructor with arguments). `ContainerAdapter.get(type, context)` gets a `ContainerContext` with the message, its attributes and the `transportMessage`, the key for per-message scopes: the handlers and workflows of one delivery get the same one (the registry passes `getReceived()`, not its workflow copy), and a retry or a re-send a new one, since transports return a new `TransportMessage` per read while `InMemoryQueue` reuses the message object (#353). `handlerFor<TMessage, TAttributes>` types the attributes; the third generic keeps the handler's own type so it can be called directly in tests.
- Handlers are called with `(message, attributes, ctx)` and class workflow handlers with `(message, state, attributes, ctx)`. `ctx` is a `HandlerContext` (`handler/handler-context.ts`: `send`, `publish`, `reply`, `failMessage`, `returnMessage`, and the optional `transaction` with `withOutbox()`, which the fakes take as an override) that `BusInstance.createHandlerContext` builds per dispatch; its methods delegate to the bus, so sends use the handler's outbox and the bus' handling context. Extend this interface rather than adding new ways to reach the bus, and update its test fakes in the same PR: `handlerContext()`, `workflowContext()` and `testWorkflow()` in `testing/` (exported for consumers) build every `HandlerContext` and `WorkflowContext` member, record the calls, and take overrides for optional members.
- `DefaultHandlerRegistry` maps message `$name` to handlers; custom handlers (`withCustomHandler`) take a resolver for external/system messages that don't follow the `Message` shape, optionally with a topic identifier the transport subscribes to.

Expand Down
10 changes: 9 additions & 1 deletion docs/guide/dependency-injection.md
Original file line number Diff line number Diff line change
Expand Up @@ -25,10 +25,18 @@ To construct them with their dependencies, pass an adapter to your IoC container

<<< @/snippets/dependency-injection.ts#container

`get` is also given the message being handled and its attributes, for containers that resolve differently per message, and may return a promise.
`get` may return a promise.

In a NestJS application, [`@node-ts/bus-nestjs`](/guide/nestjs) registers handlers and workflows as providers and resolves them from Nest's container.

## A scope per message

`get` is also given a [`ContainerContext`](/api/bus-core/interfaces/ContainerContext): the message being handled, its attributes, and its delivery, the `transportMessage` the bus read from the transport. Every class handler and workflow that handles one delivery is given the same `transportMessage`, and each retry of the message gets a new one, as does each send of the same message. To give a message's handlers dependencies of their own, such as a database transaction or a tenant's repository, key a child container on it:

<<< @/snippets/dependency-injection.ts#scope-per-delivery

Key it on the `transportMessage`, not the `message`. A transport may hand out the same message object again when it retries it, and the same object can be sent more than once, so a scope keyed on the message would give a retry the dependencies of the attempt that failed. A `WeakMap` drops each scope once its delivery has been handled.

## See also

- [Commands](/guide/messages/commands), for class handlers
Expand Down
2 changes: 1 addition & 1 deletion docs/guide/nestjs.md
Original file line number Diff line number Diff line change
Expand Up @@ -100,7 +100,7 @@ It's still stopped and disposed with the application. To stop it before anything

## Request-scoped providers

The bus resolves class handlers and workflows for each message. Request-scoped providers, and providers that depend on them, get a request scope for each received message, shared by every handler and workflow that handles it. Nest's `REQUEST` is a `BusRequest`, with the message and its attributes:
The bus resolves class handlers and workflows for each message. Request-scoped providers, and providers that depend on them, get a request scope for each delivery of a message, shared by every handler and workflow that handles it. A retry of the message is a new delivery, so it gets new request-scoped providers rather than those of the attempt that failed, and so does each send of the same message object. Nest's `REQUEST` is a `BusRequest`, with the message and its attributes:

<<< @/snippets/nestjs.ts#request-scope

Expand Down
33 changes: 32 additions & 1 deletion docs/snippets/dependency-injection.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,10 @@
import { Bus, ClassConstructor, Handler, handlerFor } from '@node-ts/bus-core'
import {
Bus,
ClassConstructor,
ContainerContext,
Handler,
handlerFor
} from '@node-ts/bus-core'
import { messageTypes } from './message-types.generated'
import { ChargeCreditCard } from './messages'

Expand Down Expand Up @@ -33,6 +39,7 @@ export class ChargeCreditCardHandler implements Handler<ChargeCreditCard> {
*/
interface Container {
get<T>(type: ClassConstructor<T>): T
createChild(): Container
}
declare const container: Container
declare const gateway: PaymentGateway
Expand All @@ -47,6 +54,30 @@ const bus = Bus.configure()
.build()
// #endregion container

// #region scope-per-delivery
// One child container for each delivery, shared by the handlers and workflows that handle it
const scopes = new WeakMap<object, Container>()

export const scopedBus = Bus.configure()
.withMessageTypes(messageTypes)
.withHandler(ChargeCreditCardHandler)
.withContainer({
get: <T>(type: ClassConstructor<T>, context?: ContainerContext) => {
const delivery = context?.transportMessage
if (!delivery) {
return container.get<T>(type)
}
let scope = scopes.get(delivery)
if (!scope) {
scope = container.createChild()
scopes.set(delivery, scope)
}
return scope.get<T>(type)
}
})
.build()
// #endregion scope-per-delivery

Bus.configure()
.withMessageTypes(messageTypes)
.withHandler(chargeCreditCardHandler(gateway))
Expand Down
1 change: 1 addition & 0 deletions docs/transports/custom.md
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,7 @@ A few rules the bus relies on:

- `readNextMessage()` returns `undefined` when there's nothing to read, rather than throwing.
- `readNextMessage()` sets `failedAttempts` to how many times handling the message has failed before, `0` on its first delivery. Count it from the broker's delivery count, or a header the transport writes when it returns a message.
- `readNextMessage()` returns a new `TransportMessage` object for each delivery, including each retry of a message it returned, even when it keeps the message itself in memory. The bus freezes it while it's handled, and [container adapters](/guide/dependency-injection#a-scope-per-message) scope what they resolve to it.
- `returnMessage(message, delay)` makes the message available again after `delay` milliseconds, with one more failed attempt. The bus' [recoverability policy](/guide/recoverability) decides the delay and when a message is out of attempts, so the transport never dead-letters a returned message itself.
- `fail(message, failure)` moves a message to the dead letter queue and removes it from the service queue, keeping its attributes and headers. Write `failure` on the copy as a `bus-failure` header with `toFailureHeader()`, and reserve that name. The bus calls `fail()` instead of `deleteMessage()`, not before it.
- The bus settles each message once: it calls exactly one of `deleteMessage()`, `returnMessage()` or `fail()`.
Expand Down
107 changes: 107 additions & 0 deletions packages/bus-core/src/container/container-adapter.spec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,11 +5,16 @@ import { ClassHandlerNotResolved, ContainerNotRegistered } from '../error'
import { Handler, HandlerDispatchRejected } from '../handler'
import { Logger } from '../logger'
import { BusMiddleware } from '../middleware'
import { retry } from '../recoverability'
import { Bus, BusInstance } from '../service-bus'
import { TestEvent, TestEvent2, testMessageTypes } from '../test'
import { TestEventClassHandler } from '../test/test-event-class-handler'
import { MessageLogger } from '../test/test-event-handler'
import { InMemoryQueue, TransportMessage } from '../transport'
import { ClassConstructor, sleep } from '../util'
import { Workflow, WorkflowMapper } from '../workflow'
import { TestWorkflowState } from '../workflow/test'
import { ContainerContext } from './container-context'

// Lets a test see messages handled by UnregisteredClassHandler when the bus constructs it with new
let handled: ((message: TestEvent2) => void) | undefined
Expand Down Expand Up @@ -328,6 +333,108 @@ describe('ContainerAdapter', () => {
})
})

describe('when class handlers and workflows are resolved for a message that is retried and sent again', () => {
const queue = new InMemoryQueue()
const resolutions: { type: string; context: ContainerContext }[] = []
const sentEvent = new TestEvent2()
let failNextHandling = true

class FailsOnceHandler implements Handler<TestEvent2> {
get messageType() {
return TestEvent2
}

async handle(): Promise<void> {
if (failNextHandling) {
failNextHandling = false
throw new Error('Failed on the first attempt')
}
}
}

class StartedByEventWorkflow extends Workflow<TestWorkflowState> {
configureWorkflow(
mapper: WorkflowMapper<TestWorkflowState, StartedByEventWorkflow>
): void {
mapper.withState(TestWorkflowState).startedBy(TestEvent2, 'start')
}

async start(): Promise<Partial<TestWorkflowState>> {
return {}
}
}

/**
* The deliveries the container was given, in the order it first saw them, each with the classes it resolved
*/
const deliveries = () => {
const byDelivery = new Map<TransportMessage<unknown>, string[]>()
for (const { type, context } of resolutions) {
const delivery = context.transportMessage!
byDelivery.set(delivery, [...(byDelivery.get(delivery) ?? []), type])
}
return [...byDelivery]
}

beforeAll(async () => {
bus = Bus.configure()
.withTransport(queue)
.withMessageTypes(testMessageTypes)
.withLogger(() => Mock.ofType<Logger>().object)
.withRecoverability(() => retry(0))
.withContainer({
get<T>(type: ClassConstructor<T>, context?: ContainerContext) {
// Without a message, the bus is reading a class workflow's configureWorkflow()
if (context?.message) {
resolutions.push({ type: type.name, context })
}
return new type()
}
})
.withHandler(FailsOnceHandler)
.withWorkflow(StartedByEventWorkflow)
.build()
await bus.initialize()
await bus.start()

await bus.publish(sentEvent)
await queue.idle()
await bus.publish(sentEvent)
await queue.idle()
})

afterAll(async () => bus.dispose())

it('should give the class handler and workflow of each delivery the same transport message', () => {
expect(deliveries().map(([, types]) => types.sort())).toEqual([
['FailsOnceHandler', 'StartedByEventWorkflow'],
['FailsOnceHandler', 'StartedByEventWorkflow'],
['FailsOnceHandler', 'StartedByEventWorkflow']
])
})

it('should give the retry and the second send transport messages of their own', () => {
expect(deliveries().map(([delivery]) => delivery.failedAttempts)).toEqual(
[0, 1, 0]
)
})

it('should give every delivery the same message object, so only the transport message tells them apart', () => {
const messages = resolutions.map(({ context }) => context.message)
expect(new Set(messages).size).toEqual(1)
expect(messages[0]).toBe(sentEvent)
})

it('should give the transport message with the message and attributes it carries', () => {
for (const { context } of resolutions) {
expect(context.transportMessage!.domainMessage).toBe(context.message)
expect(context.transportMessage!.attributes).toBe(
context.messageAttributes
)
}
})
})

describe('when no adapter is installed', () => {
describe('and no class handlers are registered', () => {
it('should initialize without errors', async () => {
Expand Down
12 changes: 6 additions & 6 deletions packages/bus-core/src/container/container-adapter.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
import { Message, MessageAttributes } from '@node-ts/bus-messages'
import { ClassConstructor } from '../util'
import { ContainerContext } from './container-context'

/**
* An adapter so that resolvers can use a local DI/IoC container
Expand All @@ -9,11 +9,11 @@ export interface ContainerAdapter {
/**
* Fetch a class instance from the container
* @param type Type of the class to fetch an instance for
* @param context Optional context to pass to the container. This is used in order to allow different resolving containers to be used for different messages.
* @param context The message being handled, its attributes and its delivery (`transportMessage`), so a container
* can resolve differently for each message. Every class handler and workflow that handles one delivery is given
* the same `transportMessage`, and a retry of the message a new one.
* @returns the instance, or a promise of it
* @example get(MessageHandler)
*/
get<T>(
type: ClassConstructor<T>,
context?: { message?: Message; messageAttributes?: MessageAttributes }
): T | Promise<T>
get<T>(type: ClassConstructor<T>, context?: ContainerContext): T | Promise<T>
}
27 changes: 27 additions & 0 deletions packages/bus-core/src/container/container-context.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,27 @@
import { Message, MessageAttributes } from '@node-ts/bus-messages'
import { TransportMessage } from '../transport'

/**
* What the bus passes to `ContainerAdapter.get` when it resolves a class handler or workflow to handle a message,
* so a container can resolve differently for each message, such as in a request scope of its own.
*/
export interface ContainerContext {
/**
* The message being handled
*/
message?: Message

/**
* The attributes of the message being handled
*/
messageAttributes?: MessageAttributes

/**
* The delivery of the message being handled, as it was read from the transport or given to a `Receiver`. It's a
* new object for each delivery, including each retry of the same message, and the same object for every class
* handler and workflow that handles that delivery. Key a scope per message on it rather than on `message`, which a
* transport may hand out again on a retry, or which may be sent more than once, so a retry gets a fresh scope
* instead of the state of the attempt that failed.
*/
transportMessage?: TransportMessage<unknown>
}
1 change: 1 addition & 0 deletions packages/bus-core/src/container/index.ts
Original file line number Diff line number Diff line change
@@ -1 +1,2 @@
export * from './container-adapter'
export * from './container-context'
3 changes: 2 additions & 1 deletion packages/bus-core/src/receiver/receiver.ts
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,8 @@ export interface Receiver<
> {
/**
* Invoked when a message is received by the application and needs to be converted into a transport message
* so that it can be passed to the dispatcher and send to handlers.
* so that it can be passed to the dispatcher and send to handlers. Return new transport messages each time, as
* `Transport.readNextMessage()` does.
*
* @param receivedMessage The message received by the app
* @param messageSerializer The configured serializer, which can be used to deserialize the incoming message
Expand Down
Loading
Loading