Skip to content

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.

ts
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))
  })
ts
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:

  1. An ItemPurchased event starts a new workflow, with a new $workflowId in its state.
  2. The start handler sends ShipItem, with the $workflowId in its sticky attributes.
  3. The ShipItem handler, maybe in another service, ships the item and publishes ItemShipped. That event gets the sticky attributes of the command, including the $workflowId.
  4. ItemShipped arrives, and its $workflowId finds the workflow's state.
  5. The ItemShipped handler 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:

ts
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))
    }
  )
ts
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:

ts
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))
    }
  )

See also ​

Released under the MIT License.