import { NonRecoverablePipelineError } from './errors'; import { Terminal } from './spi/types'; import type { TransitionRecord, TransitionResolver } from './spi/types'; import type { HandlerContext, OnAfterHandler, OnBeforeHandler, OnErrorHandler, StatefulPipelineEntity, StateRepository } from './types'; import { createLogger } from './logger'; class Pipeline, S, C extends HandlerContext> { private readonly onError: OnErrorHandler; private readonly onBefore: OnBeforeHandler; private readonly onAfter: OnAfterHandler; private readonly logger = createLogger(Pipeline.name); constructor( private readonly transitionResolver: TransitionResolver, private readonly stateRepository: StateRepository, private readonly failedState?: S, onError?: OnErrorHandler, onBefore?: OnBeforeHandler, onAfter?: OnAfterHandler ) { this.onError = onError || this.getDefaultErrorHandler(); this.onBefore = onBefore || ((entity: T) => Promise.resolve(entity)); this.onAfter = onAfter || (() => Promise.resolve()); } async handle(entity: T, ctx: C): Promise { try { this.logger.debug('Going to handler entity. State: [%s]', entity.state); let modifiedEntity = await this.onBefore(entity, ctx); const mappingOrTerminal = this.transitionResolver.resolveTransitionFrom(entity, ctx); if (mappingOrTerminal !== Terminal) { const { targetState, handler } = mappingOrTerminal as TransitionRecord; modifiedEntity = await handler.handle(modifiedEntity, ctx); modifiedEntity.state = targetState; modifiedEntity = await this.stateRepository.update(modifiedEntity, ctx); await this.onAfter(modifiedEntity, ctx); } return modifiedEntity; } catch (e) { this.logger.error(e); return this.onError(e, entity, ctx); } } private readonly getDefaultErrorHandler = () => { return (error: Error, entity: T, ctx: C) => { if (error instanceof NonRecoverablePipelineError && this.failedState) { entity.state = this.failedState; return this.stateRepository.update(entity, ctx); } else { throw error; } }; }; } export { Pipeline };