fix(engine): Handle Redis failure modes in engine v2 execution responses (#39486)

This commit is contained in:
Tomi Turtiainen
2026-09-26 18:44:11 +00:00
committed by GitHub
parent f21d644e70
commit 7cc0450ae4
5 changed files with 368 additions and 29 deletions
@@ -8,6 +8,7 @@ import {
} from '../response-channel/redis-execution-response-receiver';
import {
RedisExecutionResponseSender,
SHUTDOWN_DRAIN_TIMEOUT_MS,
type RedisResponsePublisher,
} from '../response-channel/redis-execution-response-sender';
@@ -58,6 +59,51 @@ describe('Redis execution response sender', () => {
expect(publisher.disconnect).toHaveBeenCalledTimes(1);
expect(publisher.publish).not.toHaveBeenCalled();
});
it('waits for an in-flight publish before it disconnects', async () => {
const publisher = mock<RedisResponsePublisher>();
let finishPublish: (() => void) | undefined;
publisher.publish.mockReturnValue(
new Promise((resolve) => {
finishPublish = () => resolve(1);
}),
);
const sender = new RedisExecutionResponseSender(publisher, getChannelName, mockLogger());
sender.send(ended());
const stop = sender.stop();
await Promise.resolve();
expect(publisher.disconnect).not.toHaveBeenCalled();
finishPublish?.();
await stop;
expect(publisher.disconnect).toHaveBeenCalledTimes(1);
});
it('disconnects after the drain timeout when a publish never settles', async () => {
vi.useFakeTimers();
try {
const publisher = mock<RedisResponsePublisher>();
publisher.publish.mockReturnValue(new Promise(() => {}));
const logger = mockLogger();
const sender = new RedisExecutionResponseSender(publisher, getChannelName, logger);
sender.send(ended());
const stop = sender.stop();
await vi.advanceTimersByTimeAsync(SHUTDOWN_DRAIN_TIMEOUT_MS - 1);
expect(publisher.disconnect).not.toHaveBeenCalled();
await vi.advanceTimersByTimeAsync(1);
await stop;
expect(publisher.disconnect).toHaveBeenCalledTimes(1);
expect(logger.warn).toHaveBeenCalledWith(
'Execution responses were still being published at shutdown',
{ count: 1 },
);
} finally {
vi.useRealTimers();
}
});
});
describe('Redis execution response receiver', () => {
@@ -67,6 +113,16 @@ describe('Redis execution response receiver', () => {
return registration![1];
};
const connectionHandler = (
subscriber: MockProxy<RedisResponseSubscriber>,
event: 'close' | 'ready',
): (() => void) => {
// Mock calls are typed by the last overload, which is the `message` one.
const registration = subscriber.on.mock.calls.find(([name]) => (name as string) === event);
expect(registration).toBeDefined();
return registration![1] as () => void;
};
it('refuses to receive before it starts', async () => {
const subscriber = mock<RedisResponseSubscriber>();
const receiver = new RedisExecutionResponseReceiver(subscriber, getChannelName, mockLogger());
@@ -173,7 +229,6 @@ describe('Redis execution response receiver', () => {
it('removes handlers and disconnects during shutdown', async () => {
const subscriber = mock<RedisResponseSubscriber>();
subscriber.subscribe.mockResolvedValue(1);
subscriber.unsubscribe.mockResolvedValue(1);
const receiver = new RedisExecutionResponseReceiver(subscriber, getChannelName, mockLogger());
const handler = vi.fn();
await receiver.start();
@@ -185,11 +240,82 @@ describe('Redis execution response receiver', () => {
await receiver.stop();
expect(handler).not.toHaveBeenCalled();
expect(subscriber.unsubscribe).toHaveBeenCalledExactlyOnceWith(executionChannel);
// Closing the connection ends its subscriptions, so shutdown sends no UNSUBSCRIBE.
expect(subscriber.unsubscribe).not.toHaveBeenCalled();
expect(subscriber.off).toHaveBeenCalledWith('message', onMessage);
expect(subscriber.off).toHaveBeenCalledWith('close', connectionHandler(subscriber, 'close'));
expect(subscriber.off).toHaveBeenCalledWith('ready', connectionHandler(subscriber, 'ready'));
expect(subscriber.disconnect).toHaveBeenCalledTimes(1);
});
it('stops without waiting for a subscription that is still in flight', async () => {
const subscriber = mock<RedisResponseSubscriber>();
subscriber.subscribe.mockReturnValue(new Promise(() => {}));
const receiver = new RedisExecutionResponseReceiver(subscriber, getChannelName, mockLogger());
await receiver.start();
void receiver.receive('exec-1', vi.fn());
await receiver.stop();
expect(subscriber.unsubscribe).not.toHaveBeenCalled();
expect(subscriber.disconnect).toHaveBeenCalledTimes(1);
});
it('subscribes again to active executions after Redis recovers', async () => {
const subscriber = mock<RedisResponseSubscriber>();
subscriber.subscribe.mockResolvedValue(1);
subscriber.unsubscribe.mockResolvedValue(1);
const receiver = new RedisExecutionResponseReceiver(subscriber, getChannelName, mockLogger());
await receiver.start();
await receiver.receive('exec-1', vi.fn());
await receiver.receive('exec-2', vi.fn());
const leave = await receiver.receive('exec-3', vi.fn());
leave();
subscriber.subscribe.mockClear();
connectionHandler(subscriber, 'close')();
connectionHandler(subscriber, 'ready')();
await vi.waitFor(() =>
expect(subscriber.subscribe).toHaveBeenCalledExactlyOnceWith(
executionChannel,
getChannelName('exec-2'),
),
);
});
it('does not subscribe again when the connection becomes ready for the first time', async () => {
const subscriber = mock<RedisResponseSubscriber>();
subscriber.subscribe.mockResolvedValue(1);
const receiver = new RedisExecutionResponseReceiver(subscriber, getChannelName, mockLogger());
await receiver.start();
await receiver.receive('exec-1', vi.fn());
subscriber.subscribe.mockClear();
connectionHandler(subscriber, 'ready')();
expect(subscriber.subscribe).not.toHaveBeenCalled();
});
it('logs a failure to subscribe again and retries on the next recovery', async () => {
const subscriber = mock<RedisResponseSubscriber>();
subscriber.subscribe.mockResolvedValue(1);
const logger = mockLogger();
const receiver = new RedisExecutionResponseReceiver(subscriber, getChannelName, logger);
await receiver.start();
await receiver.receive('exec-1', vi.fn());
subscriber.subscribe.mockClear();
subscriber.subscribe.mockRejectedValueOnce(new Error('Redis is unavailable'));
const onReady = connectionHandler(subscriber, 'ready');
connectionHandler(subscriber, 'close')();
onReady();
await vi.waitFor(() => expect(logger.error).toHaveBeenCalled());
onReady();
await vi.waitFor(() => expect(subscriber.subscribe).toHaveBeenCalledTimes(2));
});
it('removes the execution after its subscription fails', async () => {
const subscriber = mock<RedisResponseSubscriber>();
subscriber.subscribe.mockRejectedValue(new Error('Redis is unavailable'));
@@ -11,17 +11,19 @@ import type { RedisExecutionResponseChannelNameGenerator } from './redis-executi
type RedisMessageHandler = (channel: string, message: string) => void;
type RedisConnectionHandler = () => void;
type ExecutionSubscription = {
executionId: string;
handler: (response: ExecutionResponse) => void;
/** Resolves when the channel is ready. Shutdown uses it to wait for in-flight subscriptions. */
ready: Promise<void>;
};
export interface RedisResponseSubscriber {
subscribe(channel: string): Promise<unknown>;
subscribe(...channels: string[]): Promise<unknown>;
unsubscribe(channel: string): Promise<unknown>;
on(event: 'close' | 'ready', handler: RedisConnectionHandler): unknown;
on(event: 'message', handler: RedisMessageHandler): unknown;
off(event: 'close' | 'ready', handler: RedisConnectionHandler): unknown;
off(event: 'message', handler: RedisMessageHandler): unknown;
disconnect(): void;
}
@@ -33,6 +35,8 @@ export class RedisExecutionResponseReceiver implements ExecutionResponseReceiver
private stopped = false;
private lostConnection = false;
constructor(
private readonly subscriber: RedisResponseSubscriber,
private readonly getChannelName: RedisExecutionResponseChannelNameGenerator,
@@ -43,6 +47,8 @@ export class RedisExecutionResponseReceiver implements ExecutionResponseReceiver
if (this.started || this.stopped) return;
this.subscriber.on('message', this.handleMessage);
this.subscriber.on('close', this.handleClose);
this.subscriber.on('ready', this.handleReady);
this.started = true;
}
@@ -62,16 +68,14 @@ export class RedisExecutionResponseReceiver implements ExecutionResponseReceiver
);
}
const subscription: ExecutionSubscription = {
executionId,
handler,
ready: this.subscriber.subscribe(channel).then(() => {}),
};
const subscription: ExecutionSubscription = { executionId, handler };
this.subscriptionsByChannel.set(channel, subscription);
try {
await subscription.ready;
await this.subscriber.subscribe(channel);
} catch (error) {
this.subscriptionsByChannel.delete(channel);
if (this.subscriptionsByChannel.get(channel) === subscription) {
this.subscriptionsByChannel.delete(channel);
}
throw error;
}
@@ -101,18 +105,42 @@ export class RedisExecutionResponseReceiver implements ExecutionResponseReceiver
if (this.stopped) return;
this.stopped = true;
const subscriptions = [...this.subscriptionsByChannel.entries()];
this.subscriptionsByChannel.clear();
this.subscriber.off('message', this.handleMessage);
this.subscriber.off('close', this.handleClose);
this.subscriber.off('ready', this.handleReady);
// Closing the connection ends every Redis subscription, so there is nothing to
// unsubscribe. Waiting for Redis here could hold shutdown open while it is down.
this.subscriber.disconnect();
}
private readonly handleClose: RedisConnectionHandler = () => {
this.lostConnection = true;
};
/**
* ioredis replays SUBSCRIBE on reconnect without awaiting it, so a failed
* replay leaves the connection ready with zero subscriptions. Re-issue our
* own on every recovery. SUBSCRIBE is idempotent.
*/
private readonly handleReady: RedisConnectionHandler = () => {
if (!this.lostConnection) return;
this.lostConnection = false;
void this.resubscribe();
};
private async resubscribe(): Promise<void> {
const channels = Array.from(this.subscriptionsByChannel.keys());
if (channels.length === 0) return;
try {
await Promise.allSettled(
subscriptions.map(async ([, subscription]) => await subscription.ready),
);
await Promise.all(
subscriptions.map(async ([channel]) => await this.subscriber.unsubscribe(channel)),
);
} finally {
this.subscriber.disconnect();
await this.subscriber.subscribe(...channels);
} catch (error) {
this.lostConnection = true;
this.logger.error('Failed to resubscribe to execution responses after Redis reconnect', {
error,
});
}
}
@@ -11,6 +11,12 @@ import { UnexpectedError } from 'n8n-workflow';
import { serializeExecutionResponse } from './execution-response-frame';
import type { RedisExecutionResponseChannelNameGenerator } from './redis-execution-response-channel';
/**
* How long (in ms) shutdown waits for in-flight publishes. Redis queues commands
* while it is down, so an unbounded wait could hold shutdown open.
*/
export const SHUTDOWN_DRAIN_TIMEOUT_MS = 5_000;
export interface RedisResponsePublisher {
publish(channel: string, message: string): Promise<unknown>;
disconnect(): void;
@@ -20,6 +26,8 @@ export interface RedisResponsePublisher {
export class RedisExecutionResponseSender implements ExecutionResponseSender {
private stopped = false;
private readonly inFlightPublishes = new Set<Promise<unknown>>();
constructor(
private readonly publisher: RedisResponsePublisher,
private readonly getChannelName: RedisExecutionResponseChannelNameGenerator,
@@ -33,14 +41,16 @@ export class RedisExecutionResponseSender implements ExecutionResponseSender {
const frameResult = serializeExecutionResponse(response, this.logger);
void this.publisher
const publishTask: Promise<unknown> = this.publisher
.publish(this.getChannelName(response.executionId), frameResult.frame)
.catch((error: unknown) => {
this.logger.error('Failed to publish an execution response', {
executionId: response.executionId,
error,
});
});
})
.finally(() => this.inFlightPublishes.delete(publishTask));
this.inFlightPublishes.add(publishTask);
return frameResult.ok ? createResultOk(undefined) : createResultError(frameResult.error);
}
@@ -55,6 +65,28 @@ export class RedisExecutionResponseSender implements ExecutionResponseSender {
if (this.stopped) return;
this.stopped = true;
await this.drainInFlightPublishes();
this.publisher.disconnect();
}
/** Gives in-flight publishes a bounded time to reach Redis before it disconnects. */
private async drainInFlightPublishes(): Promise<void> {
if (this.inFlightPublishes.size === 0) return;
let timer: NodeJS.Timeout | undefined;
const timeout = new Promise<void>((resolve) => {
timer = setTimeout(resolve, SHUTDOWN_DRAIN_TIMEOUT_MS);
});
try {
await Promise.race([Promise.allSettled(this.inFlightPublishes), timeout]);
} finally {
clearTimeout(timer);
}
if (this.inFlightPublishes.size > 0) {
this.logger.warn('Execution responses were still being published at shutdown', {
count: this.inFlightPublishes.size,
});
}
}
}
@@ -1,6 +1,7 @@
import type { Logger } from '@n8n/backend-common';
import type { EngineConfig } from '@n8n/config';
import type { ExecutionResponse } from '@n8n/engine';
import { createDeferredPromise, type IDeferredPromise } from '@n8n/utils/promise/deferred-promise';
import { ENCODED_BUFFER_KEY } from 'n8n-core';
import { mock } from 'vitest-mock-extended';
@@ -9,6 +10,7 @@ import type { ExecutionResponseReceiver } from '@/modules/engine-v2/response-cha
import {
EngineV2WebhookResponder,
MAX_PENDING_WEBHOOKS,
SUBSCRIBE_TIMEOUT_MS,
} from '@/services/engine-v2-webhook-responder.service';
const TIMEOUT_MS = 50_000;
@@ -64,12 +66,17 @@ describe('EngineV2WebhookResponder', () => {
let responder: EngineV2WebhookResponder;
beforeEach(() => {
vi.useFakeTimers();
const fake = fakeReceiver();
deliver = fake.deliver;
responder = newResponder();
responder.useReceiver(fake.receiver);
});
afterEach(() => {
vi.useRealTimers();
});
it('listens under the id the run is started with', async () => {
const executionId = createExecutionIdV2();
@@ -250,6 +257,7 @@ describe('EngineV2WebhookResponder', () => {
impatient.useReceiver(fakeReceiver().receiver);
const pending = await impatient.waitForResponse(createExecutionIdV2());
await vi.advanceTimersByTimeAsync(1);
await expect(pending.settled).resolves.toEqual({
status: 'timeout',
});
@@ -275,3 +283,95 @@ describe('EngineV2WebhookResponder', () => {
);
});
});
describe('EngineV2WebhookResponder while the receiver subscribes', () => {
/** A receiver whose subscriptions complete only when the test says so. */
function slowReceiver() {
const subscriptions: Array<IDeferredPromise<() => void>> = [];
return {
receiver: {
receive: async () => {
const subscription = createDeferredPromise<() => void>();
subscriptions.push(subscription);
return await subscription.promise;
},
stop: async () => {},
} satisfies ExecutionResponseReceiver,
subscriptions,
};
}
beforeEach(() => {
vi.useFakeTimers();
});
afterEach(() => {
vi.useRealTimers();
});
it('counts runs that are still subscribing against the limit', async () => {
const { receiver } = slowReceiver();
const responder = newResponder();
responder.useReceiver(receiver);
for (let i = 0; i < MAX_PENDING_WEBHOOKS; i++) {
void responder.waitForResponse(createExecutionIdV2());
}
await expect(responder.waitForResponse(createExecutionIdV2())).rejects.toThrow(
'Try again later',
);
});
it('frees the slot and the timer when the subscription fails', async () => {
const { receiver, subscriptions } = slowReceiver();
const responder = newResponder();
responder.useReceiver(receiver);
const executionId = createExecutionIdV2();
const failed = responder.waitForResponse(executionId);
await vi.waitFor(() => expect(subscriptions).toHaveLength(1));
subscriptions[0].reject(new Error('Redis is unavailable'));
await expect(failed).rejects.toThrow('Redis is unavailable');
expect(vi.getTimerCount()).toBe(0);
const retried = responder.waitForResponse(executionId);
await vi.waitFor(() => expect(subscriptions).toHaveLength(2));
subscriptions[1].resolve(() => {});
await expect(retried).resolves.toBeDefined();
});
it('gives up at the subscribe timeout and drops a late subscription', async () => {
const { receiver, subscriptions } = slowReceiver();
// The response timeout is shorter, but it must not end the subscribe wait.
const responder = newResponder(1_000);
responder.useReceiver(receiver);
const executionId = createExecutionIdV2();
let settled = false;
const waiting = responder.waitForResponse(executionId).finally(() => {
settled = true;
});
const rejection = expect(waiting).rejects.toThrow(`within ${SUBSCRIBE_TIMEOUT_MS / 1000}s`);
await vi.advanceTimersByTimeAsync(SUBSCRIBE_TIMEOUT_MS - 1);
expect(settled).toBe(false);
await vi.advanceTimersByTimeAsync(1);
await rejection;
expect(vi.getTimerCount()).toBe(0);
// The slot is free before the late subscription completes.
const retried = responder.waitForResponse(executionId);
await vi.waitFor(() => expect(subscriptions).toHaveLength(2));
const unsubscribe = vi.fn();
subscriptions[0].resolve(unsubscribe);
await vi.waitFor(() => expect(unsubscribe).toHaveBeenCalledTimes(1));
subscriptions[1].resolve(() => {});
await expect(retried).resolves.toBeDefined();
});
});
@@ -24,6 +24,9 @@ type PendingWebhook = {
*/
export const MAX_PENDING_WEBHOOKS = 5000;
/** How long to wait for the transport to start listening for a run. */
export const SUBSCRIBE_TIMEOUT_MS = 30_000;
/**
* A service for hooking up responses from the data plane with webhook callers
* expecting that response.
@@ -60,7 +63,8 @@ export class EngineV2WebhookResponder {
* `StartExecution`.
* @throws {UnexpectedError} If the execution response receiver is not set, or
* if the service already waits for this execution.
* @throws {OperationalError} If the service is at capacity.
* @throws {OperationalError} If the service is at capacity, or if it cannot
* listen for the response within `SUBSCRIBE_TIMEOUT_MS`.
*/
async waitForResponse(
executionId: ExecutionIdV2,
@@ -79,7 +83,7 @@ export class EngineV2WebhookResponder {
// A second entry would take over the first one's slot and release.
if (this.pendingWebhooks.has(executionId)) {
throw new UnexpectedError('Engine 2.0 already waits for a response for this execution', {
throw new UnexpectedError('Engine v2 already waits for a response for this execution', {
extra: { executionId },
});
}
@@ -90,14 +94,63 @@ export class EngineV2WebhookResponder {
timeoutMs: this.engineConfig.webhookResponseTimeout,
onRelease: (id) => this.release(id),
});
const unsubscribe = await receiver.receive(executionId, (received) =>
this.handle(received, response),
);
this.pendingWebhooks.set(executionId, { response, unsubscribe });
// Hold the slot before the subscription is ready, so requests that arrive
// meanwhile still count against the limit.
const pending: PendingWebhook = { response, unsubscribe: () => {} };
this.pendingWebhooks.set(executionId, pending);
pending.unsubscribe = await this.subscribe(receiver, response);
return response;
}
/**
* A transport can wait for its broker, for example Redis while it reconnects.
* A separate timer bounds that wait, so the request cannot stay open while
* the broker is down. The run is not started when the wait runs out.
*
* If the subscription fails or times out, this releases the slot.
*/
private async subscribe(
receiver: ExecutionResponseReceiver,
response: PendingWebhookResponse,
): Promise<UnsubscribeExecutionResponse> {
const { executionId } = response;
const subscription = receiver.receive(executionId, (received) =>
this.handle(received, response),
);
let timeoutTimer: NodeJS.Timeout | undefined;
const timedOut = new Promise<'timed-out'>((resolve) => {
timeoutTimer = setTimeout(() => resolve('timed-out'), SUBSCRIBE_TIMEOUT_MS).unref();
});
let result: UnsubscribeExecutionResponse | 'timed-out';
try {
result = await Promise.race([subscription, timedOut]);
} catch (error) {
response.release();
throw error;
} finally {
clearTimeout(timeoutTimer);
}
if (result === 'timed-out') {
response.release();
// The subscription can still complete. It must not outlive the request.
void subscription.then(
(unsubscribe) => unsubscribe(),
() => {},
);
throw new OperationalError(
`Engine v2 could not listen for the execution response within ${SUBSCRIBE_TIMEOUT_MS / 1000}s.`,
{ extra: { executionId } },
);
}
return result;
}
private handle(received: ExecutionResponse, response: PendingWebhookResponse): void {
try {
this.route(received, response);