Skip to content

Custom persistence ​

To store workflow state in a database that doesn't have a persistence adapter yet, implement the Persistence interface from @node-ts/bus-core, and run the @node-ts/bus-test round trip suite against it. This page walks through both.

Implementing Persistence ​

A persistence stores and finds workflow state:

  • prepare(coreDependencies) is called by build() of each bus that uses it, with that bus' logger factory.
  • initialize() and dispose(), both optional, connect to and disconnect from the database.
  • initializeWorkflow(State, mappings) is called for each workflow state at initialize(), to create somewhere to store it and indexes for the fields its messages are looked up by.
  • getWorkflowState(State, mapping, message, attributes, includeCompleted) finds the running workflows whose mapping.mapsTo field matches the value mapping.lookup returns for the message.
  • saveWorkflowState(state) inserts a new state, when its $version is 0, or updates one.

The bus passes state to the persistence as plain JSON values, and restores its classes itself, so return the state as it was stored. Saving must use optimistic concurrency: only update a state if its version is still the one that was read, and throw if it isn't, so the message is retried with the latest state.

This skeleton adapts an imaginary document database:

ts
import {
  ClassConstructor,
  CoreDependencies,
  Logger,
  MessageWorkflowMapping,
  Persistence,
  WorkflowState,
  WorkflowStateVersionConflict,
  WorkflowStatus
} from '@node-ts/bus-core'
import { Message, MessageAttributes } from '@node-ts/bus-messages'
import { DocumentStore } from './document-store'

export class MyPersistence implements Persistence {
  private logger: Logger

  constructor(private readonly store: DocumentStore) {}

  // Called by each bus that uses the persistence
  prepare(coreDependencies: CoreDependencies): void {
    this.logger = coreDependencies.loggerFactory('my-org:my-persistence')
  }

  async initialize(): Promise<void> {
    await this.store.connect()
  }

  async dispose(): Promise<void> {
    await this.store.close()
  }

  // Create somewhere to store each workflow state, indexed by the fields
  // its messages are looked up by
  async initializeWorkflow<TWorkflowState extends WorkflowState>(
    workflowStateConstructor: ClassConstructor<TWorkflowState>,
    messageWorkflowMappings: MessageWorkflowMapping<Message, WorkflowState>[]
  ): Promise<void> {
    const indexedFields = [
      ...new Set(messageWorkflowMappings.map(mapping => mapping.mapsTo))
    ]
    await this.store.createCollection(
      collectionName(workflowStateConstructor),
      indexedFields
    )
  }

  async getWorkflowState<
    TWorkflowState extends WorkflowState,
    TMessage extends Message
  >(
    workflowStateConstructor: ClassConstructor<TWorkflowState>,
    messageMap: MessageWorkflowMapping<TMessage, TWorkflowState>,
    message: TMessage,
    attributes: MessageAttributes,
    includeCompleted = false
  ): Promise<TWorkflowState[]> {
    const lookupValue = messageMap.lookup(message, attributes)
    if (lookupValue === undefined) {
      return []
    }
    const documents = await this.store.find(
      collectionName(workflowStateConstructor),
      {
        [`data.${messageMap.mapsTo}`]: lookupValue,
        ...(!includeCompleted && { 'data.$status': WorkflowStatus.Running })
      }
    )
    // Return the state as it was stored. The bus restores its classes.
    return documents.map(document => document.data as TWorkflowState)
  }

  async saveWorkflowState<TWorkflowState extends WorkflowState>(
    workflowState: TWorkflowState
  ): Promise<void> {
    const collection = collectionNameOf(workflowState)
    const document = {
      id: workflowState.$workflowId,
      version: workflowState.$version + 1,
      data: { ...workflowState, $version: workflowState.$version + 1 }
    }
    if (workflowState.$version === 0) {
      await this.store.insert(collection, document)
      return
    }
    // Optimistic concurrency: only save over the version that was read
    const saved = await this.store.replaceIfVersion(
      collection,
      workflowState.$version,
      document
    )
    if (!saved) {
      this.logger.debug('Workflow state was changed by another handler', {
        workflowId: workflowState.$workflowId
      })
      // Throwing returns the message to the queue, so it's retried with the latest state
      throw new WorkflowStateVersionConflict(
        workflowState.$name,
        workflowState.$workflowId,
        workflowState.$version,
        undefined
      )
    }
  }
}

const collectionNameOf = (workflowState: WorkflowState) =>
  workflowState.$name.replace(/[^a-z0-9]/gi, '_')

const collectionName = (
  workflowStateConstructor: ClassConstructor<WorkflowState>
) => collectionNameOf(new workflowStateConstructor())

Pass the persistence to the bus configuration:

ts
Bus.configure()
  .withMessageTypes(messageTypes)
  .withPersistence(new MyPersistence(documentStore))

Testing with the round trip suite ​

workflowStateRoundTripTests() from @node-ts/bus-test starts a workflow whose state has Dates and nested class instances, and checks that the next handler reads it back with its types restored. Run it from your persistence's integration test, with a persistence instance of its own:

ts
import { workflowStateRoundTripTests } from '@node-ts/bus-test'
import { documentStore } from './document-store'
import { MyPersistence } from './my-persistence'

jest.setTimeout(30_000)

describe('MyPersistence', () => {
  // The suite disposes the persistence when it's done, so give it its own
  workflowStateRoundTripTests(new MyPersistence(documentStore))
})

For complete examples, see the Postgres and MongoDB adapters' tests.

Contributing a persistence

To contribute your persistence to @node-ts/bus, add it as packages/bus-<database> in the repository and open a pull request. CONTRIBUTING.md covers the conventions it follows.

See also ​

Released under the MIT License.