From d2da9284298d539613292e6dbd2852f4f14d658d Mon Sep 17 00:00:00 2001 From: Guillaume Jacquart Date: Mon, 23 Mar 2026 15:32:43 +0100 Subject: [PATCH] fix(core): Confirm messages immediately when no destination is listening (#27334) Co-authored-by: Claude Haiku 4.5 --- .../message-event-bus/message-event-bus.ts | 29 +++++---- .../log-streaming-destination.service.test.ts | 2 +- .../log-streaming-destination.service.ts | 5 +- .../cli/test/integration/eventbus.ee.test.ts | 62 ++++++++++++++++++- 4 files changed, 82 insertions(+), 16 deletions(-) diff --git a/packages/cli/src/eventbus/message-event-bus/message-event-bus.ts b/packages/cli/src/eventbus/message-event-bus/message-event-bus.ts index f838e402f52..8eac9c2dfe3 100644 --- a/packages/cli/src/eventbus/message-event-bus/message-event-bus.ts +++ b/packages/cli/src/eventbus/message-event-bus/message-event-bus.ts @@ -206,7 +206,7 @@ export class MessageEventBus extends EventEmitter { this.logger.debug(`Found unsent event messages: ${unsentMessages.length}`); for (const unsentMsg of unsentMessages) { this.logger.debug(`Retrying: ${unsentMsg.id} ${unsentMsg.__type}`); - await this.emitMessage(unsentMsg); + this.emitMessageWithCallback('message', unsentMsg); } } } @@ -223,26 +223,31 @@ export class MessageEventBus extends EventEmitter { msgs = [msgs]; } for (const msg of msgs) { + // 1. Write the message to the log file this.logWriter?.putMessage(msg); - await this.emitMessage(msg); + + // 2. Emit for internal metrics (e.g. Prometheus) + this.emit('metrics.eventBus.event', msg); + + // 3. Emit for external destinations (e.g. log-streaming to syslog/webhook/sentry). + // Returns true if at least one listener is registered, false otherwise. + const hasDestinationListener = this.emitMessageWithCallback('message', msg); + + // 4. If no external listener is registered (e.g. log-streaming module + // not licensed), confirm it immediately to prevent unbounded log accumulation. + if (!hasDestinationListener) { + this.confirmMessageDelivered(msg, { id: '0', name: 'eventBus' }); + } } } - confirmSent(msg: EventMessageTypes, source?: EventMessageConfirmSource) { + confirmMessageDelivered(msg: EventMessageTypes, source?: EventMessageConfirmSource) { this.logWriter?.confirmMessageSent(msg.id, source); } - private async emitMessage(msg: EventMessageTypes) { - this.emit('metrics.eventBus.event', msg); - - // generic emit for external modules to capture events - // this is for internal use ONLY and not for use with custom destinations! - this.emitMessageWithCallback('message', msg); - } - private emitMessageWithCallback(eventName: string, msg: EventMessageTypes): boolean { const confirmCallback = (message: EventMessageTypes, src: EventMessageConfirmSource) => - this.confirmSent(message, src); + this.confirmMessageDelivered(message, src); return this.emit(eventName, msg, confirmCallback); } diff --git a/packages/cli/src/modules/log-streaming.ee/__tests__/log-streaming-destination.service.test.ts b/packages/cli/src/modules/log-streaming.ee/__tests__/log-streaming-destination.service.test.ts index ad8ed886539..ad505f0b40b 100644 --- a/packages/cli/src/modules/log-streaming.ee/__tests__/log-streaming-destination.service.test.ts +++ b/packages/cli/src/modules/log-streaming.ee/__tests__/log-streaming-destination.service.test.ts @@ -20,7 +20,7 @@ describe('LogStreamingDestinationService', () => { const eventBus = { on: jest.fn(), removeListener: jest.fn(), - confirmSent: jest.fn(), + confirmMessageDelivered: jest.fn(), } as unknown as MessageEventBus; const publisher = mock(); diff --git a/packages/cli/src/modules/log-streaming.ee/log-streaming-destination.service.ts b/packages/cli/src/modules/log-streaming.ee/log-streaming-destination.service.ts index d05d586caea..546a9ea7dea 100644 --- a/packages/cli/src/modules/log-streaming.ee/log-streaming-destination.service.ts +++ b/packages/cli/src/modules/log-streaming.ee/log-streaming-destination.service.ts @@ -158,7 +158,7 @@ export class LogStreamingDestinationService { ) { // If there are no destinations that should receive this message, mark it as sent immediately if (!this.shouldSendMsg(msg)) { - this.eventBus.confirmSent(msg, { id: '0', name: 'eventBus' }); + confirmCallback(msg, { id: '0', name: 'eventBus' }); return; } @@ -200,7 +200,8 @@ export class LogStreamingDestinationService { if (destination.length > 0) { const sendResult = await this.destinations[destinationId].receiveFromEventBus({ msg, - confirmCallback: () => this.eventBus.confirmSent(msg, { id: '0', name: 'eventBus' }), + confirmCallback: () => + this.eventBus.confirmMessageDelivered(msg, { id: '0', name: 'eventBus' }), }); return sendResult; } diff --git a/packages/cli/test/integration/eventbus.ee.test.ts b/packages/cli/test/integration/eventbus.ee.test.ts index 7f2123a3b7f..0edcfe08e5e 100644 --- a/packages/cli/test/integration/eventbus.ee.test.ts +++ b/packages/cli/test/integration/eventbus.ee.test.ts @@ -113,6 +113,66 @@ test('should have a running logwriter process', () => { expect(thread).toBeDefined(); }); +describe('message confirmation', () => { + afterEach(async () => { + // Restore the log-streaming destination listener for subsequent tests + destinationService['isListening'] = false; + await destinationService.initialize(); + }); + + test('should confirm messages immediately when no listener is registered', async () => { + // Simulate an unlicensed instance: remove all message listeners + eventBus.removeAllListeners('message'); + + const testMessage = new EventMessageGeneric({ + eventName: 'n8n.test.message' as EventNamesTypes, + id: uuid(), + }); + + await eventBus.send(testMessage); + await new Promise((resolve) => { + eventBus.logWriter.worker?.on( + 'message', + async function handler(msg: { command: string; data: any }) { + if (msg.command === 'confirmMessageSent') { + await confirmIdSent(testMessage.id); + eventBus.logWriter.worker?.removeListener('message', handler); + resolve(true); + } + }, + ); + }); + }); + + test('should delegate confirmation to listener when one is registered', async () => { + const testMessage = new EventMessageGeneric({ + eventName: 'n8n.test.message' as EventNamesTypes, + id: uuid(), + }); + + await eventBus.send(testMessage); + // The first worker message should be appendMessageToLog, not confirmMessageSent. + // This proves the event bus delegated to the handler instead of auto-confirming. + await new Promise((resolve) => { + const workerMessages: string[] = []; + eventBus.logWriter.worker?.on( + 'message', + function handler(msg: { command: string; data: unknown }) { + workerMessages.push(msg.command); + if ( + workerMessages.includes('appendMessageToLog') && + workerMessages.includes('confirmMessageSent') + ) { + expect(workerMessages[0]).toBe('appendMessageToLog'); + eventBus.logWriter.worker?.removeListener('message', handler); + resolve(true); + } + }, + ); + }); + }); +}); + test('should have logwriter log messages', async () => { const testMessage = new EventMessageGeneric({ eventName: 'n8n.test.message' as EventNamesTypes, @@ -284,7 +344,7 @@ test('should send message to sentry ', async () => { const mockedSentryCaptureMessage = jest.spyOn(sentryDestination.sentryClient!, 'captureMessage'); mockedSentryCaptureMessage.mockImplementation((_m, _level, _hint, _scope) => { - eventBus.confirmSent(testMessage, { + eventBus.confirmMessageDelivered(testMessage, { id: sentryDestination.id, name: sentryDestination.label, });