|
| 1 | +import { ExecutionContext } from 'ava'; |
| 2 | +import * as workflow from '@temporalio/workflow'; |
| 3 | +import { HandlerUnfinishedPolicy } from '@temporalio/common'; |
| 4 | +import { LogEntry } from '@temporalio/worker'; |
| 5 | +import { Context, helpers, makeTestFunction } from './helpers-integration'; |
| 6 | + |
| 7 | +const recordedLogs: { [workflowId: string]: LogEntry[] } = {}; |
| 8 | +const test = makeTestFunction({ |
| 9 | + workflowsPath: __filename, |
| 10 | + recordedLogs, |
| 11 | +}); |
| 12 | + |
| 13 | +export const unfinishedHandlersUpdate = workflow.defineUpdate<void>('unfinished-handlers-update'); |
| 14 | +export const unfinishedHandlersUpdate_ABANDON = workflow.defineUpdate<void>('unfinished-handlers-update-ABANDON'); |
| 15 | +export const unfinishedHandlersUpdate_WARN_AND_ABANDON = workflow.defineUpdate<void>( |
| 16 | + 'unfinished-handlers-update-WARN_AND_ABANDON' |
| 17 | +); |
| 18 | +export const unfinishedHandlersSignal = workflow.defineSignal('unfinished-handlers-signal'); |
| 19 | +export const unfinishedHandlersSignal_ABANDON = workflow.defineSignal('unfinished-handlers-signal-ABANDON'); |
| 20 | +export const unfinishedHandlersSignal_WARN_AND_ABANDON = workflow.defineSignal( |
| 21 | + 'unfinished-handlers-signal-WARN_AND_ABANDON' |
| 22 | +); |
| 23 | + |
| 24 | +/** |
| 25 | + * A workflow for testing `workflow.allHandlersFinished()` and control of |
| 26 | + * warnings by HandlerUnfinishedPolicy. |
| 27 | + */ |
| 28 | +export async function unfinishedHandlersWorkflow(waitAllHandlersFinished: boolean): Promise<boolean> { |
| 29 | + let startedHandler = false; |
| 30 | + let handlerMayReturn = false; |
| 31 | + let handlerFinished = false; |
| 32 | + |
| 33 | + const doUpdateOrSignal = async (): Promise<void> => { |
| 34 | + startedHandler = true; |
| 35 | + await workflow.condition(() => handlerMayReturn); |
| 36 | + handlerFinished = true; |
| 37 | + }; |
| 38 | + |
| 39 | + workflow.setHandler(unfinishedHandlersUpdate, doUpdateOrSignal); |
| 40 | + workflow.setHandler(unfinishedHandlersUpdate_ABANDON, doUpdateOrSignal, { |
| 41 | + unfinishedPolicy: HandlerUnfinishedPolicy.ABANDON, |
| 42 | + }); |
| 43 | + workflow.setHandler(unfinishedHandlersUpdate_WARN_AND_ABANDON, doUpdateOrSignal, { |
| 44 | + unfinishedPolicy: HandlerUnfinishedPolicy.WARN_AND_ABANDON, |
| 45 | + }); |
| 46 | + workflow.setHandler(unfinishedHandlersSignal, doUpdateOrSignal); |
| 47 | + workflow.setHandler(unfinishedHandlersSignal_ABANDON, doUpdateOrSignal, { |
| 48 | + unfinishedPolicy: HandlerUnfinishedPolicy.ABANDON, |
| 49 | + }); |
| 50 | + workflow.setHandler(unfinishedHandlersSignal_WARN_AND_ABANDON, doUpdateOrSignal, { |
| 51 | + unfinishedPolicy: HandlerUnfinishedPolicy.WARN_AND_ABANDON, |
| 52 | + }); |
| 53 | + workflow.setDefaultSignalHandler(doUpdateOrSignal); |
| 54 | + |
| 55 | + await workflow.condition(() => startedHandler); |
| 56 | + if (waitAllHandlersFinished) { |
| 57 | + handlerMayReturn = true; |
| 58 | + await workflow.condition(() => workflow.allHandlersFinished()); |
| 59 | + } |
| 60 | + return handlerFinished; |
| 61 | +} |
| 62 | + |
| 63 | +test('unfinished update handler', async (t) => { |
| 64 | + await new UnfinishedHandlersTest(t, 'update').testWaitAllHandlersFinishedAndUnfinishedHandlersWarning(); |
| 65 | +}); |
| 66 | + |
| 67 | +test('unfinished signal handler', async (t) => { |
| 68 | + await new UnfinishedHandlersTest(t, 'signal').testWaitAllHandlersFinishedAndUnfinishedHandlersWarning(); |
| 69 | +}); |
| 70 | + |
| 71 | +class UnfinishedHandlersTest { |
| 72 | + constructor( |
| 73 | + private readonly t: ExecutionContext<Context>, |
| 74 | + private readonly handlerType: 'update' | 'signal' |
| 75 | + ) {} |
| 76 | + |
| 77 | + async testWaitAllHandlersFinishedAndUnfinishedHandlersWarning() { |
| 78 | + // The unfinished handler warning is issued by default, |
| 79 | + let [handlerFinished, warning] = await this.getWorkflowResultAndWarning(false); |
| 80 | + this.t.false(handlerFinished); |
| 81 | + this.t.true(warning); |
| 82 | + |
| 83 | + // and when the workflow sets the unfinished_policy to WARN_AND_ABANDON, |
| 84 | + [handlerFinished, warning] = await this.getWorkflowResultAndWarning( |
| 85 | + false, |
| 86 | + HandlerUnfinishedPolicy.WARN_AND_ABANDON |
| 87 | + ); |
| 88 | + this.t.false(handlerFinished); |
| 89 | + this.t.true(warning); |
| 90 | + |
| 91 | + // and when a default (aka dynamic) handler is used |
| 92 | + if (this.handlerType == 'signal') { |
| 93 | + [handlerFinished, warning] = await this.getWorkflowResultAndWarning(false, undefined, true); |
| 94 | + this.t.false(handlerFinished); |
| 95 | + this.t.true(warning); |
| 96 | + } else { |
| 97 | + // default handlers not supported yet for update |
| 98 | + // https://github.com/temporalio/sdk-typescript/issues/1460 |
| 99 | + } |
| 100 | + |
| 101 | + // but not when the workflow waits for handlers to complete, |
| 102 | + [handlerFinished, warning] = await this.getWorkflowResultAndWarning(true); |
| 103 | + this.t.true(handlerFinished); |
| 104 | + this.t.false(warning); |
| 105 | + if (false) { |
| 106 | + // TODO: make default handlers honor HandlerUnfinishedPolicy |
| 107 | + [handlerFinished, warning] = await this.getWorkflowResultAndWarning(true, undefined, true); |
| 108 | + this.t.true(handlerFinished); |
| 109 | + this.t.false(warning); |
| 110 | + } |
| 111 | + |
| 112 | + // nor when the silence-warnings policy is set on the handler. |
| 113 | + [handlerFinished, warning] = await this.getWorkflowResultAndWarning(false, HandlerUnfinishedPolicy.ABANDON); |
| 114 | + this.t.false(handlerFinished); |
| 115 | + this.t.false(warning); |
| 116 | + } |
| 117 | + |
| 118 | + /** |
| 119 | + * Run workflow and send signal/update. Return two booleans: |
| 120 | + * - did the handler complete? (i.e. the workflow return value) |
| 121 | + * - was an unfinished handler warning emitted? |
| 122 | + */ |
| 123 | + async getWorkflowResultAndWarning( |
| 124 | + waitAllHandlersFinished: boolean, |
| 125 | + unfinishedPolicy?: HandlerUnfinishedPolicy, |
| 126 | + useDefaultHandler?: boolean |
| 127 | + ): Promise<[boolean, boolean]> { |
| 128 | + const { createWorker, startWorkflow } = helpers(this.t); |
| 129 | + const worker = await createWorker(); |
| 130 | + return await worker.runUntil(async () => { |
| 131 | + const handle = await startWorkflow(unfinishedHandlersWorkflow, { args: [waitAllHandlersFinished] }); |
| 132 | + let messageType: string; |
| 133 | + if (useDefaultHandler) { |
| 134 | + messageType = '__no_registered_handler__'; |
| 135 | + this.t.falsy(unfinishedPolicy); // default handlers do not support setting the unfinished policy |
| 136 | + } else { |
| 137 | + messageType = `unfinished-handlers-${this.handlerType}`; |
| 138 | + if (unfinishedPolicy) { |
| 139 | + messageType += '-' + HandlerUnfinishedPolicy[unfinishedPolicy]; |
| 140 | + } |
| 141 | + } |
| 142 | + switch (this.handlerType) { |
| 143 | + case 'signal': |
| 144 | + await handle.signal(messageType); |
| 145 | + break; |
| 146 | + case 'update': |
| 147 | + const executeUpdate = handle.executeUpdate(messageType, { updateId: 'my-update-id' }); |
| 148 | + if (!waitAllHandlersFinished) { |
| 149 | + const err: workflow.WorkflowNotFoundError = (await this.t.throwsAsync(executeUpdate, { |
| 150 | + instanceOf: workflow.WorkflowNotFoundError, |
| 151 | + })) as workflow.WorkflowNotFoundError; |
| 152 | + this.t.is(err.message, 'workflow execution already completed'); |
| 153 | + } else { |
| 154 | + await executeUpdate; |
| 155 | + } |
| 156 | + break; |
| 157 | + } |
| 158 | + const handlerFinished = await handle.result(); |
| 159 | + const unfinishedHandlerWarningEmitted = |
| 160 | + recordedLogs[handle.workflowId] && |
| 161 | + recordedLogs[handle.workflowId].findIndex((e) => this.isUnfinishedHandlerWarning(e)) >= 0; |
| 162 | + return [handlerFinished, unfinishedHandlerWarningEmitted]; |
| 163 | + }); |
| 164 | + } |
| 165 | + |
| 166 | + isUnfinishedHandlerWarning(logEntry: LogEntry): boolean { |
| 167 | + return ( |
| 168 | + logEntry.level === 'WARN' && |
| 169 | + new RegExp(`^Workflow finished while an? ${this.handlerType} handler was still running\.`).test(logEntry.message) |
| 170 | + ); |
| 171 | + } |
| 172 | +} |
0 commit comments