From ecf4f8f9d968e9766e744d988b92a74328d456e8 Mon Sep 17 00:00:00 2001 From: Andrew den Hertog Date: Fri, 9 Oct 2026 15:37:36 +1000 Subject: [PATCH] pass the delivery to ContainerAdapter.get, so request scopes aren't shared across retries Co-Authored-By: Claude Opus 5.5 (1M context) --- .changeset/container-context-delivery.md | 5 + .changeset/nestjs-module.md | 2 +- CLAUDE.md | 2 +- docs/guide/dependency-injection.md | 10 +- docs/guide/nestjs.md | 2 +- docs/snippets/dependency-injection.ts | 33 +++++- docs/transports/custom.md | 1 + .../src/container/container-adapter.spec.ts | 107 ++++++++++++++++++ .../src/container/container-adapter.ts | 12 +- .../src/container/container-context.ts | 27 +++++ packages/bus-core/src/container/index.ts | 1 + packages/bus-core/src/receiver/receiver.ts | 3 +- .../bus-core/src/service-bus/bus-instance.ts | 11 +- packages/bus-core/src/transport/transport.ts | 3 + .../workflow/registry/workflow-registry.ts | 5 +- packages/bus-nestjs/CLAUDE.md | 2 +- .../bus-nestjs/src/bus-module.integration.ts | 88 +++++++++++++- packages/bus-nestjs/src/bus-request.ts | 5 +- .../bus-nestjs/src/nest-container.spec.ts | 62 +++++++--- packages/bus-nestjs/src/nest-container.ts | 55 ++++----- .../bus-nestjs/src/test/charge-attempt.ts | 88 ++++++++++++++ packages/bus-nestjs/src/test/index.ts | 1 + 22 files changed, 463 insertions(+), 62 deletions(-) create mode 100644 .changeset/container-context-delivery.md create mode 100644 packages/bus-core/src/container/container-context.ts create mode 100644 packages/bus-nestjs/src/test/charge-attempt.ts diff --git a/.changeset/container-context-delivery.md b/.changeset/container-context-delivery.md new file mode 100644 index 00000000..06f60e41 --- /dev/null +++ b/.changeset/container-context-delivery.md @@ -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. diff --git a/.changeset/nestjs-module.md b/.changeset/nestjs-module.md index 93b979fd..62b706f5 100644 --- a/.changeset/nestjs-module.md +++ b/.changeset/nestjs-module.md @@ -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`. diff --git a/CLAUDE.md b/CLAUDE.md index 2bdf0a1a..d97bcb41 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -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` 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` 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. diff --git a/docs/guide/dependency-injection.md b/docs/guide/dependency-injection.md index 04ec2dac..964f3f01 100644 --- a/docs/guide/dependency-injection.md +++ b/docs/guide/dependency-injection.md @@ -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 diff --git a/docs/guide/nestjs.md b/docs/guide/nestjs.md index d7679a8a..f0e0fdf6 100644 --- a/docs/guide/nestjs.md +++ b/docs/guide/nestjs.md @@ -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 diff --git a/docs/snippets/dependency-injection.ts b/docs/snippets/dependency-injection.ts index 2a548ee3..dd322ca7 100644 --- a/docs/snippets/dependency-injection.ts +++ b/docs/snippets/dependency-injection.ts @@ -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' @@ -33,6 +39,7 @@ export class ChargeCreditCardHandler implements Handler { */ interface Container { get(type: ClassConstructor): T + createChild(): Container } declare const container: Container declare const gateway: PaymentGateway @@ -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() + +export const scopedBus = Bus.configure() + .withMessageTypes(messageTypes) + .withHandler(ChargeCreditCardHandler) + .withContainer({ + get: (type: ClassConstructor, context?: ContainerContext) => { + const delivery = context?.transportMessage + if (!delivery) { + return container.get(type) + } + let scope = scopes.get(delivery) + if (!scope) { + scope = container.createChild() + scopes.set(delivery, scope) + } + return scope.get(type) + } + }) + .build() +// #endregion scope-per-delivery + Bus.configure() .withMessageTypes(messageTypes) .withHandler(chargeCreditCardHandler(gateway)) diff --git a/docs/transports/custom.md b/docs/transports/custom.md index dd204c25..4e44e003 100644 --- a/docs/transports/custom.md +++ b/docs/transports/custom.md @@ -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()`. diff --git a/packages/bus-core/src/container/container-adapter.spec.ts b/packages/bus-core/src/container/container-adapter.spec.ts index 1d7cfe18..04bb5dfc 100644 --- a/packages/bus-core/src/container/container-adapter.spec.ts +++ b/packages/bus-core/src/container/container-adapter.spec.ts @@ -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 @@ -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 { + get messageType() { + return TestEvent2 + } + + async handle(): Promise { + if (failNextHandling) { + failNextHandling = false + throw new Error('Failed on the first attempt') + } + } + } + + class StartedByEventWorkflow extends Workflow { + configureWorkflow( + mapper: WorkflowMapper + ): void { + mapper.withState(TestWorkflowState).startedBy(TestEvent2, 'start') + } + + async start(): Promise> { + 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, 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().object) + .withRecoverability(() => retry(0)) + .withContainer({ + get(type: ClassConstructor, 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 () => { diff --git a/packages/bus-core/src/container/container-adapter.ts b/packages/bus-core/src/container/container-adapter.ts index baee1f4d..c3cb7845 100644 --- a/packages/bus-core/src/container/container-adapter.ts +++ b/packages/bus-core/src/container/container-adapter.ts @@ -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 @@ -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( - type: ClassConstructor, - context?: { message?: Message; messageAttributes?: MessageAttributes } - ): T | Promise + get(type: ClassConstructor, context?: ContainerContext): T | Promise } diff --git a/packages/bus-core/src/container/container-context.ts b/packages/bus-core/src/container/container-context.ts new file mode 100644 index 00000000..52ac60e2 --- /dev/null +++ b/packages/bus-core/src/container/container-context.ts @@ -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 +} diff --git a/packages/bus-core/src/container/index.ts b/packages/bus-core/src/container/index.ts index c756bb04..80d06d3e 100644 --- a/packages/bus-core/src/container/index.ts +++ b/packages/bus-core/src/container/index.ts @@ -1 +1,2 @@ export * from './container-adapter' +export * from './container-context' diff --git a/packages/bus-core/src/receiver/receiver.ts b/packages/bus-core/src/receiver/receiver.ts index a44e3ab0..1836ad18 100644 --- a/packages/bus-core/src/receiver/receiver.ts +++ b/packages/bus-core/src/receiver/receiver.ts @@ -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 diff --git a/packages/bus-core/src/service-bus/bus-instance.ts b/packages/bus-core/src/service-bus/bus-instance.ts index 1e3441a4..271300fe 100644 --- a/packages/bus-core/src/service-bus/bus-instance.ts +++ b/packages/bus-core/src/service-bus/bus-instance.ts @@ -2072,7 +2072,7 @@ export class BusInstance implements BusSender { { name: handlerNameOf(handler) }, async () => this.middlewarePipeline.runHandler(invocationContext, async () => - this.invokeHandler(message, attributes, handler, context) + this.invokeHandler(transportMessage, handler, context) ) ) } @@ -2559,14 +2559,16 @@ export class BusInstance implements BusSender { /** * Resolves a class handler from the container, or constructs it, and calls it. A function handler is called as * it is. + * @param transportMessage the delivery being handled, which the container is given so it can scope what it + * resolves to that delivery * @throws ClassHandlerNotResolved if a class handler can't be resolved or constructed */ private async invokeHandler( - message: Message, - attributes: MessageAttributes, + transportMessage: TransportMessage, handler: HandlerDefinition, context: HandlerContext ): Promise { + const { domainMessage: message, attributes } = transportMessage if (!isClassHandler(handler)) { const fnHandler = handler as FunctionHandler await fnHandler(message, attributes, context) @@ -2591,7 +2593,8 @@ export class BusInstance implements BusSender { try { handlerInstance = await container.get(classHandler, { message, - messageAttributes: attributes + messageAttributes: attributes, + transportMessage }) } catch (e) { throw new ClassHandlerNotResolved( diff --git a/packages/bus-core/src/transport/transport.ts b/packages/bus-core/src/transport/transport.ts index d8868f50..856b0e9b 100644 --- a/packages/bus-core/src/transport/transport.ts +++ b/packages/bus-core/src/transport/transport.ts @@ -191,6 +191,9 @@ export interface Transport { * Fetch the next message from the underlying queue. If there are no messages, then `undefined` * should be returned. * + * Return a new `TransportMessage` object for each delivery, including each retry of a message that was returned to + * the queue. The bus freezes it while it's handled, and container adapters scope what they resolve to it. + * * @returns The message construct from the underlying transport, that includes both the raw message envelope * plus the contents or body that contains the `@node-ts/bus-messages` message. */ diff --git a/packages/bus-core/src/workflow/registry/workflow-registry.ts b/packages/bus-core/src/workflow/registry/workflow-registry.ts index d63c360b..314511e9 100644 --- a/packages/bus-core/src/workflow/registry/workflow-registry.ts +++ b/packages/bus-core/src/workflow/registry/workflow-registry.ts @@ -315,9 +315,12 @@ export class WorkflowRegistry { async (message, workflowState, attributes, context) => { let workflow: Workflow if (container) { + // The received message rather than the workflow's copy of the handling context, so the class handlers and + // workflows that handle one delivery are given the same one const workflowFromContainer = container.get(WorkflowCtor, { message, - messageAttributes: attributes + messageAttributes: attributes, + transportMessage: this.messageHandlingContext.getReceived() }) if (workflowFromContainer instanceof Promise) { workflow = await workflowFromContainer diff --git a/packages/bus-nestjs/CLAUDE.md b/packages/bus-nestjs/CLAUDE.md index 604cd2d0..1df77c7d 100644 --- a/packages/bus-nestjs/CLAUDE.md +++ b/packages/bus-nestjs/CLAUDE.md @@ -11,7 +11,7 @@ A NestJS module that runs a bus in a Nest application. Read the root `CLAUDE.md` - `forFeature`/`forFeatureAsync` only register: they never make classes providers (a dynamic module only sees its own `imports`, so a class provided there couldn't get its module's providers). Users keep classes in their own module's `providers`, and `BusClassNotProvided` reports one that isn't. `forFeatureAsync`'s `imports`/`inject` are only for its factory. - `BusLifecycleHost.onModuleInit` finds registrations with `DiscoveryService`: `BusFeatureRegistration` instances (from `forFeature`'s value provider or `forFeatureAsync`'s factory, each under a `Symbol` whose description is `BUS_FEATURE_TOKEN_PREFIX` + bus name, so one Nest didn't create, because it injects a request-scoped provider, throws `BusFeatureNotStatic`), and decorated providers (`classOf()` reads a factory or value provider's class from its instance, as Nest's discovery does). Every host checks for duplicate bus names (`BusAlreadyRegistered`) and registrations for unknown buses (`BusNotRegistered`), then registers its own, checks each class handler and workflow with `moduleRef.introspect` (`BusClassNotProvided`), builds, and throws `BusCoreVersionNotSupported` if the bus has no boolean `canStart` (an older bus-core peer). `forRootAsync` warns when the factory returns a configuration other than the seed it was given. - Lifecycle: `onApplicationBootstrap` initializes and, when `bus.canStart` (bus-core), starts, unless `lifecycle: 'manual'`, and disposes the bus (once) before rethrowing if either fails; `onModuleDestroy` stops a started bus; `onApplicationShutdown` disposes. `BusModule` is global, and Nest runs each shutdown hook of global modules after that hook of every non-global module, so the bus keeps handling messages through other modules' `onModuleDestroy`, stops before any `beforeApplicationShutdown`, and is disposed after other modules' `onApplicationShutdown` (pinned by the shutdown-order test; the same in Nest 11.0.0). Other `@Global()` modules run their hooks in import order relative to `BusModule`: one imported before `forRoot()` bootstraps before the bus is initialized, and one imported after it runs `onModuleDestroy` while the bus is still handling messages. Docs tell users to release resources in `beforeApplicationShutdown`/`onApplicationShutdown`, or use `lifecycle: 'manual'` and `bus.stop()` before `app.close()`. In an application from `createBusForProvisioning()` (which adds a global module providing `BUS_PROVISIONING`), hosts only build. The app is never closed there: `bus provision` disposes the bus and `bus.mjs` exits. -- `nestContainer(moduleRef)`: `introspect` gives the scope (`REQUEST` also for providers that depend on request-scoped ones). Singletons use `moduleRef.get(type, { strict: false })`; others `moduleRef.resolve` with a `ContextId` per message object (a `WeakMap`), registering `{ message, attributes }` (`BusRequest`) as `REQUEST` once. No message (bus-core resolving a class workflow to read `configureWorkflow()`) gets a fresh context with no request, and a failure then is rethrown as `WorkflowResolvedWithoutMessage`. The request scope is keyed by the message object, so the in-memory queue's retries of the same object share it (#353). Reading `configureWorkflow()` without resolving the workflow is #354. +- `nestContainer(moduleRef)`: `introspect` gives the scope (`REQUEST` also for providers that depend on request-scoped ones). Singletons use `moduleRef.get(type, { strict: false })`; others `moduleRef.resolve` with a `ContextId` per delivery (a `WeakMap` keyed by the context's `transportMessage`, never the message, which the in-memory queue reuses on retries and re-sends, #353), registering `{ message, attributes }` (`BusRequest`) as `REQUEST` once. No `transportMessage` gets a fresh context each call. No message (bus-core resolving a class workflow to read `configureWorkflow()`) gets a fresh context with no request, and a failure then is rethrown as `WorkflowResolvedWithoutMessage`. Reading `configureWorkflow()` without resolving the workflow is #354. ## Tests diff --git a/packages/bus-nestjs/src/bus-module.integration.ts b/packages/bus-nestjs/src/bus-module.integration.ts index ad46603f..dc2b6b55 100644 --- a/packages/bus-nestjs/src/bus-module.integration.ts +++ b/packages/bus-nestjs/src/bus-module.integration.ts @@ -19,7 +19,8 @@ import { InMemoryQueue, Logger, deadLetter, - handlerFor + handlerFor, + retry } from '@node-ts/bus-core' import { EventEmitter } from 'node:events' import { Mock } from 'typemoq' @@ -36,8 +37,12 @@ import { getBusToken } from './get-bus-token' import { InjectBus } from './inject-bus' import { BillingHandler, + ChargeAttempt, + ChargeAttemptHandler, + ChargeAttemptWorkflow, ChargeCreditCard, ChargeCreditCardHandler, + DECLINED_ONCE, FirstScopedHandler, FulfilmentWorkflow, MessageScope, @@ -198,6 +203,87 @@ describe('BusModule', () => { }) }) + describe('when request-scoped providers handle a message that is retried and sent again', () => { + const queue = new ProvisionedQueue() + const declinedOnce = new ChargeCreditCard(DECLINED_ONCE, 10) + const sentTwice = new ChargeCreditCard('sent-twice', 10) + let app: TestingModule + let recorder: Recorder + + /** + * What the handler and workflow recorded for one order, in the order they handled its deliveries + */ + const attemptsFor = (command: ChargeCreditCard) => { + const of = (name: string) => + recorder + .by(name) + .filter(({ message }) => message === command) + .map(({ detail }) => detail as ChargeAttempt) + return { + handler: of(ChargeAttemptHandler.name), + workflow: of(ChargeAttemptWorkflow.name) + } + } + + beforeAll(async () => { + app = await startApp({ + imports: [ + recorderModule(), + BusModule.forRoot({ + configure: configuration => + configuration + .withTransport(queue) + .withMessageTypes(messageTypes) + .withLogger(() => Mock.ofType().object) + .withRecoverability(() => retry(0)) + }), + BusModule.forFeature({ + handlers: [ChargeAttemptHandler], + workflows: [ChargeAttemptWorkflow] + }) + ], + providers: [ChargeAttempt, ChargeAttemptHandler, ChargeAttemptWorkflow] + }) + recorder = app.get(Recorder) + const bus = app.get(BusInstance) + + await bus.send(declinedOnce) + await queue.idle() + await bus.send(sentTwice) + await bus.send(sentTwice) + await queue.idle() + }) + + afterAll(async () => app.close()) + + it('should give the retry a request scope of its own, without the state of the attempt that failed', () => { + const { handler } = attemptsFor(declinedOnce) + expect(handler).toHaveLength(2) + const [failed, retried] = handler + expect(retried).not.toBe(failed) + expect(retried.charges).toEqual([DECLINED_ONCE]) + }) + + it('should give each send of the same message a request scope of its own', () => { + const { handler } = attemptsFor(sentTwice) + expect(handler).toHaveLength(2) + const [first, second] = handler + expect(second).not.toBe(first) + expect(second.charges).toEqual(['sent-twice']) + }) + + it('should give the handler and workflow of each delivery the same request scope, with the message as the request', () => { + for (const command of [declinedOnce, sentTwice]) { + const { handler, workflow } = attemptsFor(command) + expect(workflow).toHaveLength(2) + workflow.forEach((attempt, delivery) => { + expect(attempt).toBe(handler[delivery]) + expect(attempt.request.message).toBe(command) + }) + } + }) + }) + describe('when the application shuts down while a message is being handled', () => { const events: string[] = [] let bus: BusInstance diff --git a/packages/bus-nestjs/src/bus-request.ts b/packages/bus-nestjs/src/bus-request.ts index 5464fc1b..a88044d5 100644 --- a/packages/bus-nestjs/src/bus-request.ts +++ b/packages/bus-nestjs/src/bus-request.ts @@ -2,8 +2,9 @@ import { MessageAttributes } from '@node-ts/bus-messages' /** * What a request-scoped provider gets when it injects Nest's `REQUEST` while the bus resolves it to handle a - * message: the message and its attributes. Each received message has its own request scope, shared by every class - * handler and workflow that handles it. + * message: the message and its attributes. Each delivery of a message has its own request scope, shared by every + * class handler and workflow that handles it, so a retry gets new request-scoped providers rather than those of the + * attempt that failed. * * A request-scoped class workflow is also resolved once without a message, when the bus reads its * `configureWorkflow()`, so `REQUEST` is `undefined` then. diff --git a/packages/bus-nestjs/src/nest-container.spec.ts b/packages/bus-nestjs/src/nest-container.spec.ts index 6854c475..27ab554c 100644 --- a/packages/bus-nestjs/src/nest-container.spec.ts +++ b/packages/bus-nestjs/src/nest-container.spec.ts @@ -1,7 +1,7 @@ import { Scope } from '@nestjs/common' import { ContextId, ModuleRef } from '@nestjs/core' -import { ContainerAdapter } from '@node-ts/bus-core' -import { MessageAttributes } from '@node-ts/bus-messages' +import { ContainerAdapter, TransportMessage } from '@node-ts/bus-core' +import { Message, MessageAttributes } from '@node-ts/bus-messages' import { Mock } from 'typemoq' import { BusRequest } from './bus-request' import { BusClassNotProvided, WorkflowResolvedWithoutMessage } from './error' @@ -52,6 +52,20 @@ class FakeModuleRef { } } +/** + * A new delivery of a message, as a transport hands out for each read, including a retry of the same message + */ +const deliveryOf = ( + message: Message, + attributes: MessageAttributes +): TransportMessage => ({ + id: undefined, + domainMessage: message, + attributes, + raw: {}, + failedAttempts: 0 +}) + describe('nestContainer', () => { describe('when a singleton is resolved', () => { const moduleRef = new FakeModuleRef() @@ -71,32 +85,50 @@ describe('nestContainer', () => { }) }) - describe('when a request-scoped provider is resolved for the handlers of messages', () => { + describe('when a request-scoped provider is resolved for the handlers of deliveries', () => { const moduleRef = new FakeModuleRef() const message = new OrderPlaced('order-1') const attributes = Mock.ofType().object + const delivery = deliveryOf(message, attributes) + const retry = deliveryOf(message, attributes) beforeAll(async () => { const sut: ContainerAdapter = nestContainer( moduleRef as unknown as ModuleRef ) - await sut.get(ScopedHandler, { message, messageAttributes: attributes }) - await sut.get(ScopedHandler, { message, messageAttributes: attributes }) - await sut.get(ScopedHandler, { message: new OrderPlaced('order-2') }) + for (const transportMessage of [delivery, delivery, retry]) { + await sut.get(ScopedHandler, { + message, + messageAttributes: attributes, + transportMessage + }) + } + await sut.get(ScopedHandler, { message }) + await sut.get(ScopedHandler, { message }) await sut.get(ScopedHandler) }) - it('should resolve the handlers of one message in one request scope', () => { - const [first, second, third, fourth] = moduleRef.resolvedIn - expect(first).toBe(second) - expect(third).not.toBe(first) - expect(fourth).not.toBe(first) - expect(fourth).not.toBe(third) + it('should resolve the handlers of one delivery in one request scope', () => { + const [first, second] = moduleRef.resolvedIn + expect(second).toBe(first) + }) + + it('should resolve a retry of the same message object in a request scope of its own', () => { + const [first, , retried] = moduleRef.resolvedIn + expect(retried).not.toBe(first) + }) + + it('should resolve in a new request scope each time it is given no delivery', () => { + expect(new Set(moduleRef.resolvedIn).size).toEqual(5) }) - it('should register each message and its attributes as the request once', () => { - expect(moduleRef.requests).toHaveLength(2) - expect(moduleRef.requests[0]).toEqual({ message, attributes }) + it('should register the message and its attributes as the request once for each scope with a message', () => { + expect(moduleRef.requests).toEqual([ + { message, attributes }, + { message, attributes }, + { message, attributes: undefined }, + { message, attributes: undefined } + ]) }) }) diff --git a/packages/bus-nestjs/src/nest-container.ts b/packages/bus-nestjs/src/nest-container.ts index cca47e79..f8c329e6 100644 --- a/packages/bus-nestjs/src/nest-container.ts +++ b/packages/bus-nestjs/src/nest-container.ts @@ -1,7 +1,10 @@ import { Scope } from '@nestjs/common' import { ContextId, ContextIdFactory, ModuleRef } from '@nestjs/core' -import { ClassConstructor, ContainerAdapter } from '@node-ts/bus-core' -import { Message, MessageAttributes } from '@node-ts/bus-messages' +import { + ClassConstructor, + ContainerAdapter, + ContainerContext +} from '@node-ts/bus-core' import { BusRequest } from './bus-request' import { BusClassNotProvided, WorkflowResolvedWithoutMessage } from './error' @@ -10,9 +13,9 @@ import { BusClassNotProvided, WorkflowResolvedWithoutMessage } from './error' * configures each bus with one; use it directly to build a bus inside a Nest application without `BusModule`. * * Singleton providers are resolved with `moduleRef.get()`. Request-scoped and transient providers, and providers - * that depend on them, are resolved with `moduleRef.resolve()` in a request scope of their own for each received - * message, shared by every handler and workflow that handles it, with the message and its attributes registered - * as Nest's `REQUEST` (see `BusRequest`). + * that depend on them, are resolved with `moduleRef.resolve()` in a request scope of their own for each delivery of + * a message, shared by every handler and workflow that handles it, with the message and its attributes registered + * as Nest's `REQUEST` (see `BusRequest`). A retry of the message is a new delivery, so it gets a new scope. * @param moduleRef a `ModuleRef` from the application, which finds providers in any of its modules * @returns the adapter, for `withContainer()` * @throws BusClassNotProvided from `get`, when the class isn't a provider in the application @@ -22,23 +25,26 @@ import { BusClassNotProvided, WorkflowResolvedWithoutMessage } from './error' * Bus.configure().withContainer(nestContainer(app.get(ModuleRef))) */ export const nestContainer = (moduleRef: ModuleRef): ContainerAdapter => { - // Keyed by the message object, so the handlers and workflows of one message share a scope. Received messages - // are dropped once handled, which drops their scopes with them. + // Keyed by the delivery (the transport message), which the bus gives every handler and workflow of one delivery, + // and a retry a new one. Not by the message, which a transport may hand out again on a retry, or which may be sent + // more than once, so the retry would get the request-scoped state of the attempt that failed. Deliveries are + // dropped once handled, which drops their scopes with them. const scopes = new WeakMap() - const scopeFor = ( - message: Message | undefined, - attributes: MessageAttributes | undefined - ): ContextId => { - if (typeof message !== 'object' || message === null) { - return ContextIdFactory.create() + const scopeFor = (context: ContainerContext | undefined): ContextId => { + const delivery = context?.transportMessage + const existing = delivery ? scopes.get(delivery) : undefined + if (existing) { + return existing } - let contextId = scopes.get(message) - if (!contextId) { - contextId = ContextIdFactory.create() - scopes.set(message, contextId) + const contextId = ContextIdFactory.create() + if (delivery) { + scopes.set(delivery, contextId) + } + const message = context?.message + if (typeof message === 'object' && message !== null) { moduleRef.registerRequestByContextId( - { message, attributes }, + { message, attributes: context?.messageAttributes }, contextId ) } @@ -48,7 +54,7 @@ export const nestContainer = (moduleRef: ModuleRef): ContainerAdapter => { return { get( type: ClassConstructor, - context?: { message?: Message; messageAttributes?: MessageAttributes } + context?: ContainerContext ): T | Promise { let scope: Scope | undefined try { @@ -59,13 +65,10 @@ export const nestContainer = (moduleRef: ModuleRef): ContainerAdapter => { if (scope === Scope.DEFAULT) { return moduleRef.get(type, { strict: false }) } - const message = context?.message - const resolved = moduleRef.resolve( - type, - scopeFor(message, context?.messageAttributes), - { strict: false } - ) - if (message !== undefined) { + const resolved = moduleRef.resolve(type, scopeFor(context), { + strict: false + }) + if (context?.message !== undefined) { return resolved } // Only class workflows are resolved without a message, when the bus reads their configureWorkflow() diff --git a/packages/bus-nestjs/src/test/charge-attempt.ts b/packages/bus-nestjs/src/test/charge-attempt.ts new file mode 100644 index 00000000..1386db68 --- /dev/null +++ b/packages/bus-nestjs/src/test/charge-attempt.ts @@ -0,0 +1,88 @@ +import { Inject, Injectable, Scope } from '@nestjs/common' +import { REQUEST } from '@nestjs/core' +import { Handler, Workflow, WorkflowMapper } from '@node-ts/bus-core' +import { BusRequest } from '../bus-request' +import { ChargeCreditCard } from './charge-credit-card' +import { FulfilmentState } from './fulfilment-state' +import { Recorder } from './recorder' + +// The tests compile without experimentalDecorators, so decorators are applied by calling them + +/** + * An order whose charge fails the first time it's handled, so it's retried + */ +export const DECLINED_ONCE = 'declined-once' + +/** + * A request-scoped provider that builds up state while a message is handled, which must not reach a retry + */ +export class ChargeAttempt { + readonly charges: string[] = [] + + constructor(readonly request: BusRequest) {} +} +Injectable({ scope: Scope.REQUEST })(ChargeAttempt) +Inject(REQUEST)(ChargeAttempt, undefined, 0) + +/** + * A class handler that records a charge in the request-scoped `ChargeAttempt`, and fails the first time it charges + * `DECLINED_ONCE` + */ +export class ChargeAttemptHandler implements Handler { + constructor( + private readonly recorder: Recorder, + private readonly attempt: ChargeAttempt + ) {} + + get messageType() { + return ChargeCreditCard + } + + async handle(command: ChargeCreditCard): Promise { + this.attempt.charges.push(command.orderId) + this.recorder.record({ + by: ChargeAttemptHandler.name, + message: command, + detail: this.attempt + }) + const declinedAttempts = this.recorder + .by(ChargeAttemptHandler.name) + .filter(({ message }) => message === command) + if (command.orderId === DECLINED_ONCE && declinedAttempts.length === 1) { + throw new Error('Card declined') + } + } +} +Injectable()(ChargeAttemptHandler) +Inject(Recorder)(ChargeAttemptHandler, undefined, 0) +Inject(ChargeAttempt)(ChargeAttemptHandler, undefined, 1) + +/** + * A class workflow started by the same command, which depends on the request-scoped `ChargeAttempt` too + */ +export class ChargeAttemptWorkflow extends Workflow { + constructor( + private readonly recorder: Recorder, + private readonly attempt: ChargeAttempt + ) { + super() + } + + configureWorkflow( + mapper: WorkflowMapper + ): void { + mapper.withState(FulfilmentState).startedBy(ChargeCreditCard, 'start') + } + + async start(command: ChargeCreditCard): Promise> { + this.recorder.record({ + by: ChargeAttemptWorkflow.name, + message: command, + detail: this.attempt + }) + return { orderId: command.orderId } + } +} +Injectable()(ChargeAttemptWorkflow) +Inject(Recorder)(ChargeAttemptWorkflow, undefined, 0) +Inject(ChargeAttempt)(ChargeAttemptWorkflow, undefined, 1) diff --git a/packages/bus-nestjs/src/test/index.ts b/packages/bus-nestjs/src/test/index.ts index 440e60af..e2938f6a 100644 --- a/packages/bus-nestjs/src/test/index.ts +++ b/packages/bus-nestjs/src/test/index.ts @@ -1,3 +1,4 @@ +export * from './charge-attempt' export * from './charge-credit-card' export * from './fulfilment-state' export * from './handlers'