Handling
Once a workflow has started, it usually waits for another message before taking the next step of its process. This page covers handling those messages with when, and the ways to find which instance of the workflow a message is for.
Default mapping
When a workflow starts, it's given a $workflowId that's saved in its state and never changes. Every message sent or published from a workflow handler carries it in its sticky attributes, so it also passes on to the messages sent while handling those. A when handler with no mapping finds the workflow instance by that id.
export const fulfilmentWorkflow = defineWorkflow(FulfilmentWorkflowState)
.startedBy(ItemPurchased, async ({ itemId, customerId }, _state, ctx) => {
// ShipItem carries this workflow's id in its sticky attributes
await ctx.send(new ShipItem(itemId, customerId))
return { itemId, customerId }
})
// When the item is shipped, email the customer their receipt
.when(ItemShipped, async (_event, { itemId, customerId }, ctx) => {
await ctx.send(new EmailReceipt(itemId, customerId))
})export class FulfilmentWorkflow extends Workflow<FulfilmentWorkflowState> {
configureWorkflow(
mapper: WorkflowMapper<FulfilmentWorkflowState, FulfilmentWorkflow>
): void {
mapper
.withState(FulfilmentWorkflowState)
.startedBy(ItemPurchased, 'shipItem')
// When the item is shipped, email the customer their receipt
.when(ItemShipped, 'emailReceipt')
}
async shipItem(
{ itemId, customerId }: ItemPurchased,
_state: FulfilmentWorkflowState,
_attributes: MessageAttributes,
ctx: HandlerContext
) {
await ctx.send(new ShipItem(itemId, customerId))
return { itemId, customerId }
}
async emailReceipt(
_event: ItemShipped,
{ itemId, customerId }: FulfilmentWorkflowState,
_attributes: MessageAttributes,
ctx: HandlerContext
) {
await ctx.send(new EmailReceipt(itemId, customerId))
}
}What happens is:
- An
ItemPurchasedevent starts a new workflow, with a new$workflowIdin its state. - The start handler sends
ShipItem, with the$workflowIdin its sticky attributes. - The
ShipItemhandler, maybe in another service, ships the item and publishesItemShipped. That event gets the sticky attributes of the command, including the$workflowId. ItemShippedarrives, and its$workflowIdfinds the workflow's state.- The
ItemShippedhandler of that workflow instance runs.
The default mapping suits messages that are replies to commands the workflow sent.
Mapping by message fields
A message can also be mapped to a workflow instance by matching one of its fields to a field of the state. Give when a lookup that gets the value from the message, and the state field it mapsTo:
export const fulfilmentByItemWorkflow = defineWorkflow(FulfilmentWorkflowState)
.startedBy(ItemPurchased, async ({ itemId, customerId }, _state, ctx) => {
await ctx.send(new ShipItem(itemId, customerId))
// Save the itemId, so later messages can find this workflow by it
return { itemId, customerId }
})
.when(
ItemShipped,
{
// When an ItemShipped event is received, get its itemId...
lookup: event => event.itemId,
// ...and find the workflows whose state has the same itemId
mapsTo: 'itemId'
},
async (_event, { itemId, customerId }, ctx) => {
await ctx.send(new EmailReceipt(itemId, customerId))
}
)mapper
.withState(FulfilmentWorkflowState)
.startedBy(ItemPurchased, 'shipItem')
.when(ItemShipped, 'emailReceipt', {
lookup: event => event.itemId,
mapsTo: 'itemId'
})When ItemShipped arrives, lookup gets its itemId, and the bus finds the running workflows whose state has the same itemId. mapsTo must be a field of the state, which the compiler checks.
Mapping by fields suits messages that the workflow didn't cause, which don't carry its id.
Mapping by message attributes
lookup also gets the message's attributes, so a workflow can be found by an attribute instead:
export const fulfilmentByAttributeWorkflow = defineWorkflow(
FulfilmentWorkflowState
)
.startedBy(ItemPurchased, async ({ itemId, customerId }, _state, ctx) => {
await ctx.send(new ShipItem(itemId, customerId), {
attributes: { itemId }
})
return { itemId, customerId }
})
.when(
ItemShipped,
{
// Get the itemId from the event's attributes...
lookup: (_event, { attributes }) => attributes.itemId,
// ...and find the workflows whose state has the same itemId
mapsTo: 'itemId'
},
async (_event, { itemId, customerId }, ctx) => {
await ctx.send(new EmailReceipt(itemId, customerId))
}
)