From 568e5a24bf8f4e73d0b134dbac1631535bba10a7 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Iv=C3=A1n=20Ovejero?= Date: Mon, 4 May 2026 17:08:59 +0200 Subject: [PATCH] fix(core): Isolate expressions on chat resumption and test webhook deactivation (#29703) --- .../__tests__/chat-execution-manager.test.ts | 84 +++++++++++++++++++ .../cli/src/chat/chat-execution-manager.ts | 7 +- .../webhooks/__tests__/test-webhooks.test.ts | 76 +++++++++++++++++ packages/cli/src/webhooks/test-webhooks.ts | 16 +++- 4 files changed, 180 insertions(+), 3 deletions(-) diff --git a/packages/cli/src/chat/__tests__/chat-execution-manager.test.ts b/packages/cli/src/chat/__tests__/chat-execution-manager.test.ts index e2c394e5de9..62638fab5e9 100644 --- a/packages/cli/src/chat/__tests__/chat-execution-manager.test.ts +++ b/packages/cli/src/chat/__tests__/chat-execution-manager.test.ts @@ -434,6 +434,10 @@ describe('ChatExecutionManager', () => { const workflow = { getNode: jest.fn().mockReturnValue(node), nodeTypes: { getByNameAndVersion: jest.fn().mockReturnValue(nodeType) }, + expression: { + acquireIsolate: jest.fn().mockResolvedValue(undefined), + releaseIsolate: jest.fn().mockResolvedValue(undefined), + }, }; jest.spyOn(chatExecutionManager as any, 'getWorkflow').mockReturnValue(workflow); jest.spyOn(WorkflowExecuteAdditionalData, 'getBase').mockResolvedValue({} as any); @@ -465,6 +469,10 @@ describe('ChatExecutionManager', () => { const workflow = { getNode: jest.fn().mockReturnValue(node), nodeTypes: { getByNameAndVersion: jest.fn().mockReturnValue(nodeType) }, + expression: { + acquireIsolate: jest.fn().mockResolvedValue(undefined), + releaseIsolate: jest.fn().mockResolvedValue(undefined), + }, }; jest.spyOn(chatExecutionManager as any, 'getWorkflow').mockReturnValue(workflow); jest.spyOn(WorkflowExecuteAdditionalData, 'getBase').mockResolvedValue({} as any); @@ -474,6 +482,78 @@ describe('ChatExecutionManager', () => { expect(result).toEqual([[{ json: message }]]); }); + describe('expression isolate lifecycle', () => { + const message: ChatMessage = { + sessionId: '123', + action: 'sendMessage', + chatInput: 'input', + files: [], + }; + + function makeExecution() { + return { + id: '1', + workflowData: { id: 'workflowId' }, + data: { + resultData: { lastNodeExecuted: 'nodeId' }, + executionData: { nodeExecutionStack: [{ data: { main: [[{}]] } }] }, + }, + mode: 'manual', + } as any; + } + + function makeWorkflow(nodeType: { onMessage?: jest.Mock }) { + const expression = { + acquireIsolate: jest.fn().mockResolvedValue(undefined), + releaseIsolate: jest.fn().mockResolvedValue(undefined), + }; + const workflow = { + getNode: jest.fn().mockReturnValue({ type: 'testType', typeVersion: 1 }), + nodeTypes: { getByNameAndVersion: jest.fn().mockReturnValue(nodeType) }, + expression, + }; + jest.spyOn(chatExecutionManager as any, 'getWorkflow').mockReturnValue(workflow); + jest.spyOn(WorkflowExecuteAdditionalData, 'getBase').mockResolvedValue({} as any); + return { workflow, expression }; + } + + it('should acquire and release isolate around onMessage', async () => { + const onMessage = jest.fn().mockResolvedValue([[{ json: message }]]); + const { expression } = makeWorkflow({ onMessage }); + + await (chatExecutionManager as any).runNode(makeExecution(), message); + + expect(expression.acquireIsolate).toHaveBeenCalledTimes(1); + expect(expression.releaseIsolate).toHaveBeenCalledTimes(1); + const [acquireOrder] = expression.acquireIsolate.mock.invocationCallOrder; + const [onMessageOrder] = onMessage.mock.invocationCallOrder; + const [releaseOrder] = expression.releaseIsolate.mock.invocationCallOrder; + expect(acquireOrder).toBeLessThan(onMessageOrder); + expect(onMessageOrder).toBeLessThan(releaseOrder); + }); + + it('should release isolate when onMessage throws', async () => { + const onMessage = jest.fn().mockRejectedValue(new Error('boom')); + const { expression } = makeWorkflow({ onMessage }); + + await expect( + (chatExecutionManager as any).runNode(makeExecution(), message), + ).rejects.toThrow('boom'); + + expect(expression.acquireIsolate).toHaveBeenCalledTimes(1); + expect(expression.releaseIsolate).toHaveBeenCalledTimes(1); + }); + + it('should not acquire isolate when node type has no onMessage', async () => { + const { expression } = makeWorkflow({}); + + await (chatExecutionManager as any).runNode(makeExecution(), message); + + expect(expression.acquireIsolate).not.toHaveBeenCalled(); + expect(expression.releaseIsolate).not.toHaveBeenCalled(); + }); + }); + describe('when lastNodeExecuted is TOOL_EXECUTOR_NODE_NAME', () => { const message: ChatMessage = { sessionId: '123', action: 'sendMessage', chatInput: 'input' }; const toolNode = { @@ -550,6 +630,10 @@ describe('ChatExecutionManager', () => { nodeTypes: { getByNameAndVersion: jest.fn().mockReturnValue({ onMessage }), }, + expression: { + acquireIsolate: jest.fn().mockResolvedValue(undefined), + releaseIsolate: jest.fn().mockResolvedValue(undefined), + }, }; jest.spyOn(chatExecutionManager as any, 'getWorkflow').mockReturnValue(workflow); jest.spyOn(WorkflowExecuteAdditionalData, 'getBase').mockResolvedValue({} as any); diff --git a/packages/cli/src/chat/chat-execution-manager.ts b/packages/cli/src/chat/chat-execution-manager.ts index 3b114eecee5..cd7779e252d 100644 --- a/packages/cli/src/chat/chat-execution-manager.ts +++ b/packages/cli/src/chat/chat-execution-manager.ts @@ -132,7 +132,12 @@ export class ChatExecutionManager { } if (nodeType.onMessage) { - return await nodeType.onMessage(context, nodeExecutionData); + await workflow.expression.acquireIsolate(); + try { + return await nodeType.onMessage(context, nodeExecutionData); + } finally { + await workflow.expression.releaseIsolate(); + } } return [[nodeExecutionData]]; diff --git a/packages/cli/src/webhooks/__tests__/test-webhooks.test.ts b/packages/cli/src/webhooks/__tests__/test-webhooks.test.ts index 021ae8a341f..78f4a6b7c72 100644 --- a/packages/cli/src/webhooks/__tests__/test-webhooks.test.ts +++ b/packages/cli/src/webhooks/__tests__/test-webhooks.test.ts @@ -440,6 +440,82 @@ describe('TestWebhooks', () => { }); }); + describe('cancelWebhook()', () => { + const flushMicrotasks = async () => + await new Promise((resolve) => jest.requireActual('timers').setImmediate(resolve)); + + test('acquires and releases isolate around deactivateWebhooks', async () => { + const expression = mock(); + const workflow = mock({ id: workflowEntity.id, expression }); + + jest.spyOn(testWebhooks, 'toWorkflow').mockReturnValue(workflow); + registrations.getAllKeys.mockResolvedValue(['key1']); + registrations.get.mockResolvedValue({ + version: 1, + workflowEntity, + webhook, + } as TestWebhookRegistration); + const deactivateSpy = jest + .spyOn(testWebhooks, 'deactivateWebhooks') + .mockResolvedValue(undefined); + + await testWebhooks.cancelWebhook(workflowEntity.id); + await flushMicrotasks(); + + expect(expression.acquireIsolate).toHaveBeenCalledTimes(1); + expect(deactivateSpy).toHaveBeenCalledWith(workflow); + expect(expression.releaseIsolate).toHaveBeenCalledTimes(1); + const [acquireOrder] = (expression.acquireIsolate as jest.Mock).mock.invocationCallOrder; + const [deactivateOrder] = deactivateSpy.mock.invocationCallOrder; + const [releaseOrder] = (expression.releaseIsolate as jest.Mock).mock.invocationCallOrder; + expect(acquireOrder).toBeLessThan(deactivateOrder); + expect(deactivateOrder).toBeLessThan(releaseOrder); + }); + }); + + describe('handleClearTestWebhooks()', () => { + test('acquires and releases isolate around deactivateWebhooks', async () => { + const expression = mock(); + const workflow = mock({ id: workflowEntity.id, expression }); + + jest.spyOn(testWebhooks, 'toWorkflow').mockReturnValue(workflow); + ((testWebhooks as any).push.hasPushRef as jest.Mock).mockReturnValue(true); + const deactivateSpy = jest + .spyOn(testWebhooks, 'deactivateWebhooks') + .mockResolvedValue(undefined); + + await testWebhooks.handleClearTestWebhooks({ + webhookKey: 'key1', + workflowEntity, + pushRef: 'push-ref', + }); + + expect(expression.acquireIsolate).toHaveBeenCalledTimes(1); + expect(deactivateSpy).toHaveBeenCalledWith(workflow); + expect(expression.releaseIsolate).toHaveBeenCalledTimes(1); + }); + + test('releases isolate when deactivateWebhooks throws', async () => { + const expression = mock(); + const workflow = mock({ id: workflowEntity.id, expression }); + + jest.spyOn(testWebhooks, 'toWorkflow').mockReturnValue(workflow); + ((testWebhooks as any).push.hasPushRef as jest.Mock).mockReturnValue(true); + jest.spyOn(testWebhooks, 'deactivateWebhooks').mockRejectedValue(new Error('boom')); + + await expect( + testWebhooks.handleClearTestWebhooks({ + webhookKey: 'key1', + workflowEntity, + pushRef: 'push-ref', + }), + ).rejects.toThrow('boom'); + + expect(expression.acquireIsolate).toHaveBeenCalledTimes(1); + expect(expression.releaseIsolate).toHaveBeenCalledTimes(1); + }); + }); + describe('getWebhookMethods()', () => { beforeEach(() => { registrations.toKey.mockImplementation( diff --git a/packages/cli/src/webhooks/test-webhooks.ts b/packages/cli/src/webhooks/test-webhooks.ts index 0bdcf80534d..80159ddf1ed 100644 --- a/packages/cli/src/webhooks/test-webhooks.ts +++ b/packages/cli/src/webhooks/test-webhooks.ts @@ -225,7 +225,12 @@ export class TestWebhooks implements IWebhookManager { const workflow = this.toWorkflow(workflowEntity); - await this.deactivateWebhooks(workflow); + await workflow.expression.acquireIsolate(); + try { + await this.deactivateWebhooks(workflow); + } finally { + await workflow.expression.releaseIsolate(); + } } clearTimeout(key: string) { @@ -476,7 +481,14 @@ export class TestWebhooks implements IWebhookManager { if (!foundWebhook) { // As it removes all webhooks of the workflow execute only once - void this.deactivateWebhooks(workflow); + void (async () => { + await workflow.expression.acquireIsolate(); + try { + await this.deactivateWebhooks(workflow); + } finally { + await workflow.expression.releaseIsolate(); + } + })(); } foundWebhook = true;