Skip to content

Example ​

This example models fulfilling a purchase from an online store. When a customer buys an item, the workflow ships it, then emails the customer a receipt, then completes.

ItemPurchased starts the workflow, which sends ShipItem. ItemShipped makes it send EmailReceipt, and ReceiptEmailed completes it.

The messages and state ​

ts
// The messages of the fulfilment workflow used in the workflow guides
import { Command, Event } from '@node-ts/bus-messages'

/**
 * A customer bought an item from the store
 */
export class ItemPurchased extends Event {
  static NAME = 'my-app/store/item-purchased'
  $name = ItemPurchased.NAME
  $version = 0

  constructor(
    readonly itemId: string,
    readonly customerId: string
  ) {
    super()
  }
}

/**
 * Ship an item to a customer
 */
export class ShipItem extends Command {
  static NAME = 'my-app/shipping/ship-item'
  $name = ShipItem.NAME
  $version = 0

  constructor(
    readonly itemId: string,
    readonly customerId: string
  ) {
    super()
  }
}

/**
 * An item was shipped to a customer
 */
export class ItemShipped extends Event {
  static NAME = 'my-app/shipping/item-shipped'
  $name = ItemShipped.NAME
  $version = 0

  constructor(
    readonly itemId: string,
    readonly shippedAt: Date
  ) {
    super()
  }
}

/**
 * Email a customer the receipt for an item
 */
export class EmailReceipt extends Command {
  static NAME = 'my-app/email/email-receipt'
  $name = EmailReceipt.NAME
  $version = 0

  constructor(
    readonly itemId: string,
    readonly customerId: string
  ) {
    super()
  }
}

/**
 * The receipt for an item was emailed to its customer
 */
export class ReceiptEmailed extends Event {
  static NAME = 'my-app/email/receipt-emailed'
  $name = ReceiptEmailed.NAME
  $version = 0

  constructor(readonly itemId: string) {
    super()
  }
}
ts
import { WorkflowState } from '@node-ts/bus-core'

export class FulfilmentWorkflowState extends WorkflowState {
  // Unique among all of your workflow states
  static NAME = 'my-app/store/fulfilment-workflow-state'
  $name = FulfilmentWorkflowState.NAME

  // The workflow's own fields
  itemId: string
  customerId: string
  status: 'shipping-item' | 'emailing-receipt' | 'complete'
  shippedAt?: Date
}

The workflow ​

ts
import { defineWorkflow } from '@node-ts/bus-core'
import {
  EmailReceipt,
  ItemPurchased,
  ItemShipped,
  ReceiptEmailed,
  ShipItem
} from '../messages'
import { FulfilmentWorkflowState } from './fulfilment-workflow-state'

export const fulfilmentWorkflow = defineWorkflow(FulfilmentWorkflowState)
  // A purchase starts a new workflow, which ships the item
  .startedBy(ItemPurchased, async ({ itemId, customerId }, _state, ctx) => {
    await ctx.send(new ShipItem(itemId, customerId))
    return { itemId, customerId, status: 'shipping-item' as const }
  })
  // ItemShipped is published by the ShipItem handler, so it carries this
  // workflow's id and is routed back to this instance
  .when(ItemShipped, async (event, { itemId, customerId }, ctx) => {
    await ctx.send(new EmailReceipt(itemId, customerId))
    return { status: 'emailing-receipt' as const, shippedAt: event.shippedAt }
  })
  .when(ReceiptEmailed, (_event, _state, ctx) =>
    ctx.complete({ status: 'complete' })
  )
ts
import { HandlerContext, Workflow, WorkflowMapper } from '@node-ts/bus-core'
import { MessageAttributes } from '@node-ts/bus-messages'
import {
  EmailReceipt,
  ItemPurchased,
  ItemShipped,
  ReceiptEmailed,
  ShipItem
} from '../messages'
import { FulfilmentWorkflowState } from './fulfilment-workflow-state'

export class FulfilmentWorkflow extends Workflow<FulfilmentWorkflowState> {
  configureWorkflow(
    mapper: WorkflowMapper<FulfilmentWorkflowState, FulfilmentWorkflow>
  ): void {
    mapper
      .withState(FulfilmentWorkflowState)
      .startedBy(ItemPurchased, 'shipItem')
      .when(ItemShipped, 'emailReceipt')
      .when(ReceiptEmailed, 'complete')
  }

  async shipItem(
    { itemId, customerId }: ItemPurchased,
    _state: FulfilmentWorkflowState,
    _attributes: MessageAttributes,
    ctx: HandlerContext
  ) {
    await ctx.send(new ShipItem(itemId, customerId))
    return { itemId, customerId, status: 'shipping-item' as const }
  }

  async emailReceipt(
    event: ItemShipped,
    { itemId, customerId }: FulfilmentWorkflowState,
    _attributes: MessageAttributes,
    ctx: HandlerContext
  ) {
    await ctx.send(new EmailReceipt(itemId, customerId))
    return { status: 'emailing-receipt' as const, shippedAt: event.shippedAt }
  }

  complete() {
    return this.completeWorkflow({ status: 'complete' })
  }
}

ItemShipped and ReceiptEmailed are published by the handlers of the commands the workflow sent, so they carry its id and use the default mapping.

The handlers and the bus ​

The commands are handled by ordinary handlers, which in a real system would usually run in other services:

ts
import { handlerFor } from '@node-ts/bus-core'
import {
  EmailReceipt,
  ItemShipped,
  ReceiptEmailed,
  ShipItem
} from '../messages'
import { shippingService } from '../services'

// The ShipItem command's sticky attributes, including the workflow id, are
// copied to the ItemShipped event it publishes
export const shipItemHandler = handlerFor(
  ShipItem,
  async (command, _attributes, ctx) => {
    await shippingService.ship(command.itemId, command.customerId)
    await ctx.publish(new ItemShipped(command.itemId, new Date()))
  }
)

export const emailReceiptHandler = handlerFor(
  EmailReceipt,
  async (command, _attributes, ctx) => {
    // ...send the email
    await ctx.publish(new ReceiptEmailed(command.itemId))
  }
)

Register the workflow, give the bus the message types of the messages and the state, and publish an ItemPurchased to start it:

ts
import { Bus } from '@node-ts/bus-core'
import {
  emailReceiptHandler,
  shipItemHandler
} from './handlers/fulfilment-handlers'
import { messageTypes } from './message-types.generated'
import { ItemPurchased } from './messages'
import { fulfilmentWorkflow } from './workflows/fulfilment-workflow'

const bus = Bus.configure()
  // Includes the workflow state, so its Dates are restored when it's read
  .withMessageTypes(messageTypes)
  .withWorkflow(fulfilmentWorkflow)
  // In a real system these would usually run in other services
  .withHandler(shipItemHandler, emailReceiptHandler)
  .build()

await bus.initialize()
await bus.start()

await bus.publish(new ItemPurchased('item-1', 'customer-1'))

See also ​

  • Persistence, to store the state in Postgres or MongoDB
  • State, for testing workflow handlers

Released under the MIT License.