Saga
Wow's @StatelessSaga is a stateless event orchestrator: it receives domain or state events and generates commands for the next step. Each target aggregate still handles its command in a local transaction. A Saga does not create a cross-aggregate ACID transaction.
A Stateless Saga converts an occurred fact into 0..N follow-up commands sent in order.
When to Use a Saga
Use a Saga when a committed event must drive business behavior in other aggregates, such as sending an entry command to a target account after a transfer is prepared. Use an Event Processor for notifications, audit, or external integrations. Do not wrap an ordinary side effect that generates no command in a Saga.
A stateless Saga stores no workflow instance. Each event function must depend only on the current event and injectable context, then explicitly return the 0..N commands required for this invocation.
Stateless Saga Contract
committed event -> matching Saga function -> 0..N commands -> CommandGateway.sendA Saga shares event-function parsing, type/topic matching, domain/state-event dispatch, and the reactive filter chain with an ordinary Processor. The difference is that StatelessSagaFunction converts function results into commands, sends them in order, and stores the command stream on the exchange for the SAGA_HANDLED notification.
The source event is committed before the Saga runs. Neither a function failure nor a command-send failure can roll it back.
Define Saga Functions
@StatelessSaga is also a Spring component marker. Functions may use the onEvent / onStateEvent conventional names or explicit @OnEvent / @OnStateEvent annotations. Their first parameter may likewise be the event body, DomainEvent<T>, or DomainEventExchange<T>.
@StatelessSaga
class CartSaga {
@OnEvent
fun onOrderCreated(event: DomainEvent<OrderCreated>): CommandBuilder? {
if (!event.body.fromCart) return null
return RemoveCartItem(
productIds = event.body.items.map { it.productId }.toSet(),
).commandBuilder()
.aggregateId(event.ownerId)
}
}Use @OnStateEvent and inject state as a later parameter only when the aggregate's latest state is required. Return a CommandBuilder when the target aggregate ID, request ID, or another command field must be selected explicitly.
Generate 0..N Commands from an Event
A function can synchronously, asynchronously, or reactively return these results:
| Result | Send behavior |
|---|---|
null / Mono.empty() | Send 0 commands |
| Command body | Convert it to a CommandMessage and send 1 command |
CommandBuilder | Create and send 1 command with the builder's target and explicit fields |
CommandMessage<*> | Preserve the existing message and send 1 command |
Iterable<*>, Flux, Publisher, or Flow | Collect and send N commands in result order |
StatelessSagaFunction uses concatMap to create and send commands one at a time: the next command is created only after the previous CommandGateway.send completes. Keep result order stable. Reordering changes both the business flow and the default request IDs.
requestId and Context Propagation
For a command body or CommandBuilder, the framework fills only missing fields:
- the default
requestIdof a returned command body orCommandBuilderis${domainEvent.id}-${index}, starting at index0, as in 9.2 (a mixed 9.2/9.3 cluster derives the same ID). It does not name the Saga, so two Sagas that both return a body or builder for the same event, targeting the same aggregate, share request IDs and the second command is skipped as a duplicate request; give such commands an explicitrequestId(for example with the Saga's name) or return aCommandMessage(below); - an explicit
requestIdis preserved; - a command that names no aggregate (no aggregate ID in the body or on the builder) gets one derived from the event, the Saga function, the command's index and the target aggregate type, instead of a random one (since 9.3.0). Handling the same event again therefore sends a create to the same aggregate ID with the same request ID, which is rejected as a duplicate request instead of creating a second aggregate. The derived ID has the format of the target aggregate's ID generator, with the event's creation time as its timestamp; Wow always uses CosId, and deterministic Saga IDs apply to the time-based CosId generators (the default CosId and Snowflake); segment and custom generators, and generators whose string form depends on the current date or the time zone (a date prefix, a friendly ID), keep random IDs. The ID is set before any
CommandBuilderRewriterruns, so a rewriter sees it and may replace it. When the retried create reaches a store whose request-ID check comes after its version check, the processing node still reportsDuplicateRequestId(it finds the request ID on the existing stream), not a duplicate aggregate ID. Two different creates whose derived IDs collide are a realDuplicateAggregateId: the second create fails on the processing node and is logged there; it is not recorded by compensation. The derived ID keeps the generator's machine and sequence bits, so collisions need the same timestamp unit: negligible for the default CosId (36 hash bits per millisecond), but a Snowflake layout has only 22 (milliseconds) or 32 (seconds) bits per unit, so many Saga creates of one aggregate type from events in the same millisecond (or second) can collide; prefer the default CosId generator for aggregates that Sagas create at a high rate; - missing
tenantIdandspaceIdpropagate from the source event; - the source event becomes upstream and its message header is propagated.
A prebuilt CommandMessage keeps its message, aggregate ID and an explicit requestId while receiving source-event header propagation. A requestId it did not set (equal to its ID, the default) becomes ${domainEvent.id}-${index}-${producer hash} (since 9.3.0), where the hash is taken over the Saga function's context, processor and function name: two Sagas returning such a message for the same event get different request IDs, so the second is not skipped as a duplicate of the first. Its aggregate ID is not derived, because a generated ID cannot be told from a chosen one: return a command body or CommandBuilder to get a derived aggregate ID.
Replaying the same event with the same result order produces stable default request IDs that can cooperate with command-gateway idempotency checks. This does not make external side effects idempotent and does not deduplicate semantically repeated commands generated from different events.
During a rolling upgrade from 9.2, a retry handled by a 9.2 node still draws a random aggregate ID, so a create retried across a 9.2 and a 9.3 node can still create two aggregates, as every retry could on 9.2; two 9.3 attempts converge.
Immediate retries run the Saga function again and resend its returned commands. Newly created command messages keep the existing globally unique commandId generation rules, while default request IDs remain stable for idempotency checks. When an individual send returns DuplicateRequestIdException, the Saga skips that duplicate send and continues with subsequent commands; other errors still propagate. The Saga does not cache send progress.
Business Compensation
Business compensation in a Saga is an explicit domain action. For example, EntryFailed can generate UnlockAmount:
@OnEvent
fun onEntryFailed(event: EntryFailed): UnlockAmount =
UnlockAmount(event.sourceId, event.amount)This command describes how the business offsets an earlier effect. It does not delete committed events such as Prepared or EntryFailed, and it is not a database rollback.
Do not conflate business compensation with processing-failure recovery. A Saga decides which business command follows a failure fact. Compensation decides how a failed invocation of the same processing function is durably recorded and replayed. See Event Compensation for the authoritative recovery guide.
Wait Integration
The runtime produces SAGA_HANDLED after the Saga function completes and every generated CommandGateway.send completes. The signal carries the command-stream commandId values and confirms that dispatch finished, including skipped duplicate requests. It does not confirm successful processing by the target aggregates.
Wait for the matching SAGA_HANDLED when the caller needs only to know that the Saga sent its commands. If every downstream command must reach another stage, use CommandWait.chain(...) with the Saga function and the tail stage/function. Wait state is process-local and subject to caller timeouts and cancellation; Saga retries do not guarantee that an earlier wait remains active or can be resumed. See Completion Semantics for stages, function matching, and early-arriving signals.
Testing and Failure Boundaries
Use SagaSpec to verify the event-to-command mapping directly without starting a message broker:
class CartSagaSpec : SagaSpec<CartSaga>({
on {
whenEvent(orderCreatedFromCart, ownerId = ownerId) {
expectNoError()
expectCommandType(RemoveCartItem::class)
expectCommand<RemoveCartItem> {
aggregateId.id.assert().isEqualTo(ownerId)
}
}
}
on {
whenEvent(orderCreatedNotFromCart, ownerId = ownerId) {
expectNoCommand()
}
}
})Cover the normal command, zero-command path, and every business-compensation branch. Add assertions on the full CommandMessage when relying on default request IDs, context propagation, multi-command order, or prebuilt messages. Add integration coverage only when relying on real sends, chained waits, or failure recovery.
An error from the Saga function or a CommandGateway.send error other than DuplicateRequestIdException propagates through the reactive chain before immediate retry or enabled durable Compensation handles it. Even when failure follows partial command acceptance, neither the source event nor accepted commands are automatically undone. Command handling and business compensation must therefore remain safe under replay.
Completion signal: tests cover the 0..N mapping, command order, request ID/context propagation, and business-compensation branches, and the wait contract does not misstate SAGA_HANDLED as downstream command completion.