DefaultCommandHandler

class DefaultCommandHandler(serviceProvider: ServiceProvider, aggregateProcessorFactory: AggregateProcessorFactory, domainEventBus: DomainEventBus?, stateEventBus: StateEventBus?, commandWaitNotifier: CommandWaitNotifier?, instrumentations: List<CommandInstrumentation> = emptyList(), requestIdChecker: RequestIdChecker? = null, errorHandler: ErrorHandler<ServerCommandExchange<*>> = LogResumeErrorHandler()) : CommandHandler

The command side's pipeline, in a fixed order (V5: there is no command filter chain):

  1. Each CommandInstrumentation wraps everything below; the first one is the outermost.

  2. When there is a commandWaitNotifier, the PROCESSED wait signal is reported once everything below completed or failed; a failed signal carries the error.

  3. When there is a requestIdChecker, the command's request ID is checked on this node, the one that processes it, before the aggregate runs: a request ID this aggregate already committed fails the command with DuplicateRequestIdException without running its handler. The event store's unique request ID still guards the append.

  4. The aggregate processes the command through an AggregateProcessorFactory processor. The transport message is acknowledged whatever the outcome; a failure skips the publication.

  5. The domain event stream the command committed is sent on domainEventBus; a failure propagates.

  6. When the state applied that stream (its version is the stream's), the state event is sent on stateEventBus; a failure is logged and resumed.

A failure that reaches the end is recorded on the exchange and given to errorHandler. A null bus or notifier skips its step; the Spring wiring passes all of them.

Constructors

Link copied to clipboard
constructor(serviceProvider: ServiceProvider, aggregateProcessorFactory: AggregateProcessorFactory, domainEventBus: DomainEventBus?, stateEventBus: StateEventBus?, commandWaitNotifier: CommandWaitNotifier?, instrumentations: List<CommandInstrumentation> = emptyList(), requestIdChecker: RequestIdChecker? = null, errorHandler: ErrorHandler<ServerCommandExchange<*>> = LogResumeErrorHandler())

Functions

Link copied to clipboard
open override fun handle(exchange: ServerCommandExchange<*>, aggregateMetadata: AggregateMetadata<*, *>): Mono<Void>

Handles exchange, a command for an aggregate described by aggregateMetadata. The returned Mono completes once the command has been handled, successfully or not; it does not fail for a command that failed.