fix(core): Confirm messages immediately when no destination is listening (#27334)

Co-authored-by: Claude Haiku 4.5 <noreply@anthropic.com>
This commit is contained in:
Guillaume Jacquart
2026-03-23 14:32:43 +00:00
committed by GitHub
co-authored by Claude Haiku 4.5
parent 050aef73db
commit d2da928429
4 changed files with 82 additions and 16 deletions
@@ -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);
}
@@ -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<Publisher>();
@@ -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;
}
@@ -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,
});