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 bybuild()of each bus that uses it, with that bus' logger factory.initialize()anddispose(), both optional, connect to and disconnect from the database.initializeWorkflow(State, mappings)is called for each workflow state atinitialize(), 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 whosemapping.mapsTofield matches the valuemapping.lookupreturns for the message.saveWorkflowState(state)inserts a new state, when its$versionis 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:
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:
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:
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
- Custom transports
PersistenceandworkflowStateRoundTripTestsin the API reference