diff --git a/services/process/src/main.ts b/services/process/src/main.ts index 991c6390be..7db1d3b659 100644 --- a/services/process/src/main.ts +++ b/services/process/src/main.ts @@ -41,6 +41,7 @@ import process, { State, Step, Transition, + Trigger, UserResult } from '@hcengineering/process' import serverProcess, { @@ -73,7 +74,7 @@ export async function messageHandler (record: ProcessMessage, ws: WorkspaceUuid, ctx.info('Processing event', { event: record.event, ws, record }) if (record.execution !== undefined) { const execution = await control.client.findOne(process.class.Execution, { _id: record.execution }) - if (execution !== undefined) { + if (execution !== undefined && isActiveExecution(execution, record.event)) { await processExecution(control, record, execution) } } else if (record.card !== undefined) { @@ -89,6 +90,12 @@ export async function messageHandler (record: ProcessMessage, ws: WorkspaceUuid, } } +function isActiveExecution (execution: Execution, trigger: Ref): boolean { + if (execution.status !== ExecutionStatus.Active) return false + if (execution.error != null && trigger !== process.trigger.OnExecutionContinue) return false + return true +} + async function processBroadcast (control: ProcessControl, record: ProcessMessage): Promise { const transitions = control.client.getModel().findAllSync(process.class.Transition, { trigger: record.event @@ -100,7 +107,9 @@ async function processBroadcast (control: ProcessControl, record: ProcessMessage currentState: transition.from }) for (const execution of executions) { - await execute(execution, transition, control) + if (isActiveExecution(execution, record.event)) { + await execute(execution, transition, control) + } } } } @@ -122,7 +131,9 @@ async function processCardExecutions (control: ProcessControl, record: ProcessMe currentState: { $in: Array.from(states) } }) for (const execution of executions) { - await processExecution(control, record, execution) + if (isActiveExecution(execution, record.event)) { + await processExecution(control, record, execution) + } } } @@ -309,115 +320,123 @@ async function execute (execution: Execution, transition: Transition, control: P } } -async function executeTransition (execution: Execution, transition: Transition, control: ProcessControl): Promise { - let deep = control.cache.get(execution._id + 'transition') ?? 0 - deep++ - control.cache.set(execution._id + 'transition', deep) - if (deep > 100) { - const error = parseError( - processError(process.error.TooDeepTransitionRecursion, undefined, undefined, true), - transition._id - ) - await control.client.update(execution, { error: [error] }) - return - } - const trigger = control.client.getModel().findObject(transition.trigger) - if (trigger === undefined) return - const rollback: Tx[] = [] - const triggerImpl = control.client.getHierarchy().as(trigger, serverProcess.mixin.TriggerImpl) - const disableRollback = triggerImpl?.preventRollback ?? false - const triggerRollback = await getTriggerRollback(triggerImpl, control) - if (triggerRollback !== undefined) { - rollback.push(triggerRollback) - } - const client = control.client - const res: Tx[] = [] - const _process = client.getModel().findObject(execution.process) - if (_process === undefined) return - - const state = client.getModel().findObject(transition.to) - if (state === undefined) return - const isDone = - client.getModel().findAllSync(process.class.Transition, { - from: transition.to, - process: transition.process - }).length === 0 - const errors: ExecutionError[] = [] - if (execution.currentState !== null) { - rollback.push( - client.txFactory.createTxUpdateDoc(execution._class, execution.space, execution._id, { - currentState: execution.currentState, - status: execution.status - }) - ) - } else { - rollback.push(client.txFactory.createTxRemoveDoc(execution._class, execution.space, execution._id)) - } - if (!control.cache.has(execution.card)) { - const card = await control.client.findOne(cardPlugin.class.Card, { _id: execution.card }) - if (card !== undefined) { - control.cache.set(execution.card, card) +async function executeTransition ( + execution: Execution, + _transition: Transition, + control: ProcessControl +): Promise { + let transition: Transition | undefined = _transition + while (transition !== undefined) { + let deep = control.cache.get(execution._id + 'transition') ?? 0 + deep++ + control.cache.set(execution._id + 'transition', deep) + if (deep > 100) { + const error = parseError( + processError(process.error.TooDeepTransitionRecursion, undefined, undefined, true), + transition._id + ) + await control.client.update(execution, { error: [error] }) + return } - } - for (const action of transition.actions) { - const actionResult = await executeAction(action, transition._id, execution, control) - if (isError(actionResult)) { - errors.push(actionResult) + const trigger = control.client.getModel().findObject(transition.trigger) + if (trigger === undefined) return + const rollback: Tx[] = [] + const triggerImpl = control.client.getHierarchy().as(trigger, serverProcess.mixin.TriggerImpl) + const disableRollback = triggerImpl?.preventRollback ?? false + const triggerRollback = await getTriggerRollback(triggerImpl, control) + if (triggerRollback !== undefined) { + rollback.push(triggerRollback) + } + const client = control.client + const res: Tx[] = [] + const _process = client.getModel().findObject(execution.process) + if (_process === undefined) return + + const state = client.getModel().findObject(transition.to) + if (state === undefined) return + const isDone = + client.getModel().findAllSync(process.class.Transition, { + from: transition.to, + process: transition.process + }).length === 0 + const errors: ExecutionError[] = [] + if (execution.currentState !== null) { + rollback.push( + client.txFactory.createTxUpdateDoc(execution._class, execution.space, execution._id, { + currentState: execution.currentState, + status: execution.status + }) + ) } else { - if (actionResult.rollback !== undefined) { - rollback.push(...actionResult.rollback) + rollback.push(client.txFactory.createTxRemoveDoc(execution._class, execution.space, execution._id)) + } + if (!control.cache.has(execution.card)) { + const card = await control.client.findOne(cardPlugin.class.Card, { _id: execution.card }) + if (card !== undefined) { + control.cache.set(execution.card, card) } - res.push(...actionResult.txes) - for (const tx of actionResult.txes) { - const updateTx = tx as TxUpdateDoc - if (updateTx._class === core.class.TxUpdateDoc && updateTx.objectId === execution.card) { - const updatedCard = TxProcessor.updateDoc2Doc( - control.client.getHierarchy().clone(control.cache.get(execution.card)), - updateTx - ) - control.cache.set(execution.card, updatedCard) + } + for (const action of transition.actions) { + const actionResult = await executeAction(action, transition._id, execution, control) + if (isError(actionResult)) { + errors.push(actionResult) + } else { + if (actionResult.rollback !== undefined) { + rollback.push(...actionResult.rollback) + } + res.push(...actionResult.txes) + for (const tx of actionResult.txes) { + const updateTx = tx as TxUpdateDoc + if (updateTx._class === core.class.TxUpdateDoc && updateTx.objectId === execution.card) { + const updatedCard = TxProcessor.updateDoc2Doc( + control.client.getHierarchy().clone(control.cache.get(execution.card)), + updateTx + ) + control.cache.set(execution.card, updatedCard) + } } } } - } - if (!disableRollback) { - execution.rollback.push(rollback) - } - const executionUpdate = getDiffUpdate(execution, { - currentState: state._id, - status: isDone ? ExecutionStatus.Done : ExecutionStatus.Active, - error: null - }) - executionUpdate.context = execution.context - executionUpdate.rollback = execution.rollback.length > 30 ? execution.rollback.slice(-30) : execution.rollback - res.push(client.txFactory.createTxUpdateDoc(execution._class, execution.space, execution._id, executionUpdate)) - res.push( - client.txFactory.createTxCreateDoc(process.class.ExecutionLog, execution.space, { - execution: execution._id, - process: execution.process, - card: execution.card, - transition: transition._id, - action: transition.from === null ? ExecutionLogAction.Started : ExecutionLogAction.Transition + if (!disableRollback) { + execution.rollback.push(rollback) + } + const executionUpdate = getDiffUpdate(execution, { + currentState: state._id, + status: isDone ? ExecutionStatus.Done : ExecutionStatus.Active, + error: null }) - ) - if (errors.length === 0) { - for (const tx of res) { - const timeout = setTimeout(() => { - control.ctx.warn('TX HANG', tx) - }, 30000) - await client.tx(tx) - clearTimeout(timeout) + executionUpdate.context = execution.context + executionUpdate.rollback = execution.rollback.length > 30 ? execution.rollback.slice(-30) : execution.rollback + res.push(client.txFactory.createTxUpdateDoc(execution._class, execution.space, execution._id, executionUpdate)) + res.push( + client.txFactory.createTxCreateDoc(process.class.ExecutionLog, execution.space, { + execution: execution._id, + process: execution.process, + card: execution.card, + transition: transition._id, + action: transition.from === null ? ExecutionLogAction.Started : ExecutionLogAction.Transition + }) + ) + if (errors.length === 0) { + for (const tx of res) { + const timeout = setTimeout(() => { + control.ctx.warn('TX HANG', tx) + }, 30000) + await client.tx(tx) + clearTimeout(timeout) + } + if (isDone && execution.parentId !== undefined) { + await checkParent(execution, control) + } + TxProcessor.applyUpdate(execution, executionUpdate) + transition = await checkNext(control, execution) + if (transition === undefined) { + await setNextTimers(control, execution) + } + } else { + await client.update(execution, { error: errors }) + break } - if (isDone && execution.parentId !== undefined) { - await checkParent(execution, control) - } - TxProcessor.applyUpdate(execution, executionUpdate) - const next = await checkNext(control, execution) - if (!next) { - await setNextTimers(control, execution) - } - } else { - await client.update(execution, { error: errors }) } } @@ -623,7 +642,7 @@ async function checkParent (execution: Execution, control: ProcessControl): Prom } } -async function checkNext (control: ProcessControl, execution: Execution): Promise { +async function checkNext (control: ProcessControl, execution: Execution): Promise { const autoTriggers = control.client.getModel().findAllSync(process.class.Trigger, { auto: true }) const transitions = control.client.getModel().findAllSync( process.class.Transition, @@ -641,13 +660,9 @@ async function checkNext (control: ProcessControl, execution: Execution): Promis const transition = await pickTransition(control, execution, transitions, { card: doc }) - if (transition !== undefined) { - await executeTransition(execution, transition, control) - return true - } + return transition } } - return false } export async function pickTransition (