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.
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