mirror of
https://github.com/n8n-io/n8n.git
synced 2026-09-24 23:22:38 +08:00
fix(core): Bound Instance AI checkpoint growth and enforce confirmation expiry (no-changelog) (#33456)
Co-authored-by: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 4.8
parent
29aa414b1d
commit
41681795b9
@@ -148,6 +148,10 @@ export class InstanceAiConfig {
|
||||
@Env('N8N_INSTANCE_AI_SNAPSHOT_RETENTION')
|
||||
snapshotRetention: number = 24 * Time.hours.toMilliseconds;
|
||||
|
||||
/** Retention period in milliseconds for expired checkpoint tombstones before they are hard-deleted. Must exceed snapshotRetention. 0 = never hard-delete. */
|
||||
@Env('N8N_INSTANCE_AI_CHECKPOINT_GC_RETENTION')
|
||||
checkpointGcRetention: number = 7 * Time.days.toMilliseconds;
|
||||
|
||||
/** Timeout in milliseconds for HITL confirmation requests. 0 = no timeout. */
|
||||
@Env('N8N_INSTANCE_AI_CONFIRMATION_TIMEOUT')
|
||||
confirmationTimeout: number = 24 * Time.hours.toMilliseconds;
|
||||
|
||||
@@ -340,6 +340,7 @@ describe('GlobalConfig', () => {
|
||||
threadTtlDays: 90,
|
||||
pruneInterval: 3_600_000,
|
||||
snapshotRetention: 86_400_000,
|
||||
checkpointGcRetention: 604_800_000,
|
||||
confirmationTimeout: 86_400_000,
|
||||
outputRedactionEnabled: true,
|
||||
outputRedactionSecrets: true,
|
||||
|
||||
@@ -456,12 +456,14 @@ type CheckpointPruneServiceInternals = {
|
||||
scheduleCheckpointPrune: MockedFunction<(delayMs?: number) => void>;
|
||||
checkpointStore: {
|
||||
markExpiredOlderThan: MockedFunction<(olderThan: Date) => Promise<number>>;
|
||||
hardDeleteExpiredOlderThan: MockedFunction<(olderThan: Date) => Promise<number>>;
|
||||
};
|
||||
checkpointPruneTimer?: NodeJS.Timeout;
|
||||
checkpointPruningStopped: boolean;
|
||||
instanceAiConfig: {
|
||||
pruneInterval: number;
|
||||
snapshotRetention: number;
|
||||
checkpointGcRetention: number;
|
||||
};
|
||||
logger: { info: Mock; debug: Mock; warn: Mock };
|
||||
};
|
||||
@@ -477,11 +479,13 @@ function createCheckpointPruneService(): CheckpointPruneServiceInternals {
|
||||
service.pruneExpiredThreads = vi.fn(async () => undefined);
|
||||
service.checkpointStore = {
|
||||
markExpiredOlderThan: vi.fn(async (_olderThan: Date) => 0),
|
||||
hardDeleteExpiredOlderThan: vi.fn(async (_olderThan: Date) => 0),
|
||||
};
|
||||
service.checkpointPruningStopped = true;
|
||||
service.instanceAiConfig = {
|
||||
pruneInterval: 60 * 60 * 1000,
|
||||
snapshotRetention: 7 * 24 * 60 * 60 * 1000,
|
||||
snapshotRetention: 24 * 60 * 60 * 1000,
|
||||
checkpointGcRetention: 7 * 24 * 60 * 60 * 1000,
|
||||
};
|
||||
service.logger = {
|
||||
info: vi.fn(),
|
||||
@@ -1269,7 +1273,12 @@ describe('InstanceAiService — scheduled pruning', () => {
|
||||
|
||||
await service.runScheduledPrune(now);
|
||||
|
||||
// snapshotRetention = 24h → tombstone anything untouched since 05-12
|
||||
expect(service.checkpointStore.markExpiredOlderThan).toHaveBeenCalledWith(
|
||||
new Date('2026-05-12T12:00:00.000Z'),
|
||||
);
|
||||
// checkpointGcRetention = 7d → hard-delete tombstones expired before 05-06
|
||||
expect(service.checkpointStore.hardDeleteExpiredOlderThan).toHaveBeenCalledWith(
|
||||
new Date('2026-05-06T12:00:00.000Z'),
|
||||
);
|
||||
expect(service.suspendedThreads.pruneStalePendingConfirmations).toHaveBeenCalledWith(now);
|
||||
@@ -1277,6 +1286,31 @@ describe('InstanceAiService — scheduled pruning', () => {
|
||||
expect(service.scheduleCheckpointPrune).toHaveBeenCalledWith();
|
||||
});
|
||||
|
||||
it('skips hard-deleting tombstones when the GC retention is disabled', async () => {
|
||||
const service = createCheckpointPruneService();
|
||||
service.instanceAiConfig.checkpointGcRetention = 0;
|
||||
|
||||
await service.runScheduledPrune(new Date('2026-05-13T12:00:00.000Z').getTime());
|
||||
|
||||
expect(service.checkpointStore.hardDeleteExpiredOlderThan).not.toHaveBeenCalled();
|
||||
// The rest of the cycle still runs.
|
||||
expect(service.checkpointStore.markExpiredOlderThan).toHaveBeenCalled();
|
||||
expect(service.scheduleCheckpointPrune).toHaveBeenCalledWith();
|
||||
});
|
||||
|
||||
it('continues the prune cycle when hard-deleting tombstones fails', async () => {
|
||||
const service = createCheckpointPruneService();
|
||||
service.checkpointStore.hardDeleteExpiredOlderThan.mockRejectedValueOnce(new Error('db down'));
|
||||
|
||||
await service.runScheduledPrune(new Date('2026-05-13T12:00:00.000Z').getTime());
|
||||
|
||||
// A GC failure is swallowed and never forces the short retry cadence.
|
||||
expect(service.suspendedThreads.pruneStalePendingConfirmations).toHaveBeenCalled();
|
||||
expect(service.pruneExpiredThreads).toHaveBeenCalled();
|
||||
expect(service.scheduleCheckpointPrune).toHaveBeenCalledWith();
|
||||
expect(service.logger.warn).toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it('starts checkpoint pruning when configured', () => {
|
||||
const service = createCheckpointPruneService();
|
||||
|
||||
@@ -1446,6 +1480,9 @@ type ResolveConfirmationServiceInternals = {
|
||||
suspendedThreads: {
|
||||
dropPendingConfirmation: Mock;
|
||||
};
|
||||
pendingConfirmationRepo: {
|
||||
isPastExpiry: Mock<(...args: [string, string, Date]) => Promise<boolean>>;
|
||||
};
|
||||
logger: { debug: Mock; warn: Mock; error: Mock; info: Mock };
|
||||
};
|
||||
|
||||
@@ -1462,6 +1499,9 @@ function createResolveConfirmationService(): ResolveConfirmationServiceInternals
|
||||
findSuspendedByRequestId: vi.fn(),
|
||||
rejectPendingConfirmation: vi.fn(),
|
||||
};
|
||||
service.pendingConfirmationRepo = {
|
||||
isPastExpiry: vi.fn(async () => false),
|
||||
};
|
||||
service.resumeSuspendedRun = vi.fn(async () => null);
|
||||
service.suspendedRunRestorer = {
|
||||
resolveOrphanedConfirmation: vi.fn(async () => null),
|
||||
@@ -1674,6 +1714,45 @@ describe('InstanceAiService — resolveConfirmation', () => {
|
||||
expect(service.suspendedThreads.dropPendingConfirmation).toHaveBeenCalledWith('req-1');
|
||||
});
|
||||
|
||||
it('refuses a click on a confirmation whose row is already past its expiry', async () => {
|
||||
const service = createResolveConfirmationService();
|
||||
service.revalidateActiveUser.mockResolvedValue({ id: 'user-1' } as unknown as User);
|
||||
service.pendingConfirmationRepo.isPastExpiry.mockResolvedValue(true);
|
||||
// A still-present in-memory entry must not be resolved once the row expired.
|
||||
service.runState.getPendingConfirmation.mockReturnValue({
|
||||
userId: 'user-1',
|
||||
threadId: 'thread-1',
|
||||
});
|
||||
|
||||
await expect(service.resolveConfirmation('user-1', 'req-1', approval)).rejects.toThrow(
|
||||
/expired/i,
|
||||
);
|
||||
|
||||
expect(service.pendingConfirmationRepo.isPastExpiry).toHaveBeenCalledWith(
|
||||
'req-1',
|
||||
'user-1',
|
||||
expect.any(Date),
|
||||
);
|
||||
expect(service.runState.resolvePendingConfirmation).not.toHaveBeenCalled();
|
||||
expect(service.suspendedRunRestorer.resolveOrphanedConfirmation).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it('resolves normally when the row is not past its expiry', async () => {
|
||||
const service = createResolveConfirmationService();
|
||||
service.revalidateActiveUser.mockResolvedValue({ id: 'user-1' } as unknown as User);
|
||||
service.pendingConfirmationRepo.isPastExpiry.mockResolvedValue(false);
|
||||
service.runState.getPendingConfirmation.mockReturnValue({
|
||||
userId: 'user-1',
|
||||
threadId: 'thread-1',
|
||||
});
|
||||
service.runState.resolvePendingConfirmation.mockReturnValue(true);
|
||||
service.runState.getActiveRunId.mockReturnValue('run-1');
|
||||
|
||||
const result = await service.resolveConfirmation('user-1', 'req-1', approval);
|
||||
|
||||
expect(result).toEqual({ ok: true, runId: 'run-1' });
|
||||
});
|
||||
|
||||
it('delegates to the orphan-restoration path when no live run resumes', async () => {
|
||||
// The detailed orphan claim/rebuild/finalize scenarios live in
|
||||
// suspended-run-restorer.service.test.ts; here we only assert the
|
||||
|
||||
@@ -237,6 +237,9 @@ function isTelemetryConfigurableAgent(
|
||||
const INSTANCE_AI_CHECKPOINT_PRUNE_RETRY_MS = 30 * 1000;
|
||||
const WORKFLOW_SETUP_ROUTING_CLAIM_TTL_MS = 15 * 60 * 1000;
|
||||
|
||||
const CONFIRMATION_EXPIRED_MESSAGE =
|
||||
'This confirmation has expired. Send a new message to continue.';
|
||||
|
||||
/**
|
||||
* Upper bound on how long `shutdown()` will wait for in-flight executeRun /
|
||||
* processResumedStream promises to drain after their abortControllers fire.
|
||||
@@ -1437,10 +1440,11 @@ export class InstanceAiService {
|
||||
|
||||
/**
|
||||
* One tick of the recurring leader prune cycle: expire stale checkpoints,
|
||||
* drop expired pending confirmations, and delete expired conversation
|
||||
* threads, then schedule the next run. A checkpoint failure reschedules
|
||||
* with a shorter retry delay; the confirmation and thread steps swallow
|
||||
* their own errors so they never disrupt the cycle.
|
||||
* hard-delete tombstones past the GC horizon, drop expired pending
|
||||
* confirmations, and delete expired conversation threads, then schedule the
|
||||
* next run. A checkpoint failure reschedules with a shorter retry delay; the
|
||||
* confirmation and thread steps swallow their own errors so they never
|
||||
* disrupt the cycle.
|
||||
*/
|
||||
private async runScheduledPrune(now = Date.now()): Promise<void> {
|
||||
const olderThan = new Date(now - this.instanceAiConfig.snapshotRetention);
|
||||
@@ -1452,6 +1456,7 @@ export class InstanceAiService {
|
||||
} else {
|
||||
this.logger.debug('No stale Instance AI checkpoints to expire');
|
||||
}
|
||||
await this.hardDeleteExpiredCheckpoints(now);
|
||||
await this.suspendedThreads.pruneStalePendingConfirmations(now);
|
||||
await this.pruneExpiredThreads();
|
||||
this.scheduleCheckpointPrune();
|
||||
@@ -1463,6 +1468,31 @@ export class InstanceAiService {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Hard-delete expired checkpoint tombstones (`markExpiredOlderThan` releases
|
||||
* the blob but keeps the row so a stale resume gets a clear "expired" error;
|
||||
* this drops the row entirely once it's past the GC horizon so tombstones
|
||||
* don't grow unbounded). Has its own try/catch so a GC failure never forces
|
||||
* the checkpoint prune into its short retry cadence. No-op when
|
||||
* `checkpointGcRetention` is 0.
|
||||
*/
|
||||
private async hardDeleteExpiredCheckpoints(now: number): Promise<void> {
|
||||
const retention = this.instanceAiConfig.checkpointGcRetention;
|
||||
if (retention <= 0) return;
|
||||
|
||||
try {
|
||||
const olderThan = new Date(now - retention);
|
||||
const count = await this.checkpointStore.hardDeleteExpiredOlderThan(olderThan);
|
||||
if (count > 0) {
|
||||
this.logger.info('Hard-deleted expired Instance AI checkpoint tombstones', { count });
|
||||
}
|
||||
} catch (error: unknown) {
|
||||
this.logger.warn('Failed to hard-delete expired Instance AI checkpoint tombstones', {
|
||||
error: getErrorMessage(error),
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Delete conversation threads older than the configured TTL as part of the
|
||||
* recurring leader prune. Has its own try/catch so a failure here never
|
||||
@@ -3828,6 +3858,19 @@ export class InstanceAiService {
|
||||
return null;
|
||||
}
|
||||
|
||||
// Close the same-process window: a click on a pre-rendered card can still
|
||||
// reach the in-memory fast path after the row's `expiresAt` but before the
|
||||
// prune/liveness sweep drops it. The persisted row is the source of truth
|
||||
// for expiry, so consult it first and refuse a stale click; the sweep still
|
||||
// owns releasing the suspended run. Scoped by `freshUser.id` so another
|
||||
// user's request ID falls through to the existing not-found handling
|
||||
// rather than leaking an "expired" signal. One extra SELECT per click —
|
||||
// the click path isn't hot.
|
||||
if (await this.pendingConfirmationRepo.isPastExpiry(requestId, freshUser.id, new Date())) {
|
||||
this.logger.debug('Rejecting expired confirmation', { requestId });
|
||||
throw new UserError(CONFIRMATION_EXPIRED_MESSAGE);
|
||||
}
|
||||
|
||||
const pending = this.runState.getPendingConfirmation(requestId);
|
||||
if (
|
||||
pending &&
|
||||
|
||||
+39
@@ -119,3 +119,42 @@ describe('InstanceAiPendingConfirmationRepository.claim', () => {
|
||||
expect(where).toHaveProperty('expiresAt');
|
||||
});
|
||||
});
|
||||
|
||||
describe('InstanceAiPendingConfirmationRepository.isPastExpiry', () => {
|
||||
function buildRepoWithCount(countResult: number) {
|
||||
const repo = Object.create(
|
||||
InstanceAiPendingConfirmationRepository.prototype,
|
||||
) as InstanceAiPendingConfirmationRepository;
|
||||
const countMock = vi.fn(async (_opts?: { where: Record<string, unknown> }) => countResult);
|
||||
Object.defineProperty(repo, 'count', { value: countMock, configurable: true });
|
||||
return { repo, countMock };
|
||||
}
|
||||
|
||||
it('is true when the user owns a row with expiresAt in the past', async () => {
|
||||
const now = new Date('2026-05-13T12:00:00.000Z');
|
||||
const { repo, countMock } = buildRepoWithCount(1);
|
||||
|
||||
await expect(repo.isPastExpiry('req-1', 'user-1', now)).resolves.toBe(true);
|
||||
// Scoped by userId, and LessThan(now) excludes null expiresAt (timeout
|
||||
// disabled) at the SQL layer.
|
||||
const where = countMock.mock.calls[0][0]!.where;
|
||||
expect(where).toMatchObject({ requestId: 'req-1', userId: 'user-1' });
|
||||
expect(where).toHaveProperty('expiresAt');
|
||||
});
|
||||
|
||||
it('is false when no row is past its expiry', async () => {
|
||||
const { repo } = buildRepoWithCount(0);
|
||||
|
||||
await expect(repo.isPastExpiry('req-1', 'user-1', new Date())).resolves.toBe(false);
|
||||
});
|
||||
|
||||
it('does not treat another user’s expired request as expired', async () => {
|
||||
// The userId scope means the count query matches nothing for a request
|
||||
// owned by someone else, so the caller falls through to its existing
|
||||
// not-found/not-authorized handling instead of leaking an expired signal.
|
||||
const { repo, countMock } = buildRepoWithCount(0);
|
||||
|
||||
await expect(repo.isPastExpiry('req-1', 'attacker-user', new Date())).resolves.toBe(false);
|
||||
expect(countMock.mock.calls[0][0]!.where).toMatchObject({ userId: 'attacker-user' });
|
||||
});
|
||||
});
|
||||
|
||||
+13
@@ -65,6 +65,19 @@ export class InstanceAiPendingConfirmationRepository extends Repository<Instance
|
||||
return result.affected ?? 0;
|
||||
}
|
||||
|
||||
/**
|
||||
* True when this user's confirmation row exists and its `expiresAt` is
|
||||
* already in the past. Scoped by `userId` — mirroring `claim` and the
|
||||
* in-memory path — so another user's request ID is treated as not-expired
|
||||
* here and left to the caller's existing not-found/not-authorized handling.
|
||||
* Rows with a null `expiresAt` (timeout disabled) — and a missing row — are
|
||||
* not expired.
|
||||
*/
|
||||
async isPastExpiry(requestId: string, userId: string, now: Date): Promise<boolean> {
|
||||
const count = await this.count({ where: { requestId, userId, expiresAt: LessThan(now) } });
|
||||
return count > 0;
|
||||
}
|
||||
|
||||
async findByThreadId(threadId: string): Promise<InstanceAiPendingConfirmation[]> {
|
||||
return await this.find({ where: { threadId } });
|
||||
}
|
||||
|
||||
+10
@@ -1,4 +1,5 @@
|
||||
import type { SerializableAgentState } from '@n8n/instance-ai';
|
||||
import { LessThan } from '@n8n/typeorm';
|
||||
import { UserError } from 'n8n-workflow';
|
||||
import type { Mock } from 'vitest';
|
||||
import { mock } from 'vitest-mock-extended';
|
||||
@@ -161,6 +162,15 @@ describe('TypeORMAgentCheckpointStore', () => {
|
||||
expect(spy).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
|
||||
it('hard-deletes tombstones expired before the GC horizon', async () => {
|
||||
checkpointRepo.delete.mockResolvedValueOnce({ affected: 3, raw: {} });
|
||||
|
||||
const olderThan = new Date('2026-05-01T00:00:00.000Z');
|
||||
await expect(store.hardDeleteExpiredOlderThan(olderThan)).resolves.toBe(3);
|
||||
|
||||
expect(checkpointRepo.delete).toHaveBeenCalledWith({ expiredAt: LessThan(olderThan) });
|
||||
});
|
||||
|
||||
describe('findSuspendedSubAgentResumeInfo', () => {
|
||||
it('picks the suspended tool call when parallel non-suspended ones are present', async () => {
|
||||
const state = makeState({
|
||||
|
||||
Reference in New Issue
Block a user