From 4dd369991fa87e68814c5963ae45359aa4ff8989 Mon Sep 17 00:00:00 2001 From: Jaakko Husso Date: Thu, 9 Jul 2026 10:44:20 +0300 Subject: [PATCH] feat(core): Use a shared per-thread event sequence for Instance AI multi-main (no-changelog) (#33558) Co-authored-by: Claude Fable 5 --- .../instance-ai/docs/streaming-protocol.md | 13 +- .../harness/in-memory-event-bus.ts | 4 +- .../src/event-bus/event-bus.interface.ts | 5 +- .../__tests__/instance-ai.controller.test.ts | 2 +- .../__tests__/in-process-event-bus.test.ts | 286 +++++++++++++++-- .../event-bus/in-process-event-bus.ts | 291 +++++++++++++++--- .../instance-ai/instance-ai.controller.ts | 5 +- .../src/scaling/pubsub/pubsub.event-map.ts | 7 +- .../instanceAi.threadRuntime.test.ts | 63 ++++ .../ai/instanceAi/instanceAi.threadRuntime.ts | 39 ++- 10 files changed, 626 insertions(+), 89 deletions(-) diff --git a/packages/@n8n/instance-ai/docs/streaming-protocol.md b/packages/@n8n/instance-ai/docs/streaming-protocol.md index 8df9f5cb1d0..f2c948702f5 100644 --- a/packages/@n8n/instance-ai/docs/streaming-protocol.md +++ b/packages/@n8n/instance-ai/docs/streaming-protocol.md @@ -52,7 +52,9 @@ data: {"type":"text-delta","runId":"run_abc","agentId":"agent-001","payload":{"t ``` Event IDs are monotonically increasing integers per thread channel and unique -within that thread. +within that thread. In multi-main deployments they come from a shared +per-thread sequence (Redis), so ids — and the replay cursor built from them — +are valid against any main. ## Event Schema @@ -407,8 +409,13 @@ The backend must accept the cursor from both `Last-Event-ID` header and `?lastEventId` query parameter. If neither is present, replay starts from event ID 0 (full history). -IDs are monotonically increasing integers per thread. Replay does not -require dedup. +IDs are monotonically increasing integers per thread, assigned from a shared +per-thread sequence in multi-main. Assignment order is monotonic, but delivery +order is not guaranteed to be: concurrent producers on different mains (e.g. a +background task while the orchestrator runs elsewhere) can interleave, so a +connection may occasionally deliver a lower id after a higher one. The +frontend therefore tracks its reconnect cursor as the max id seen and drops +already-seen ids on replay overlap. ## Abort Support diff --git a/packages/@n8n/instance-ai/evaluations/harness/in-memory-event-bus.ts b/packages/@n8n/instance-ai/evaluations/harness/in-memory-event-bus.ts index 46f8bfb4f46..b9dd8a59c92 100644 --- a/packages/@n8n/instance-ai/evaluations/harness/in-memory-event-bus.ts +++ b/packages/@n8n/instance-ai/evaluations/harness/in-memory-event-bus.ts @@ -41,8 +41,8 @@ export function createInMemoryEventBus(): InstanceAiEventBus { .map((event) => event.event) .filter((event) => 'runId' in event && runIdSet.has(event.runId)); }, - getNextEventId(threadId) { - return (storeByThread.get(threadId) ?? []).length + 1; + async getNextEventId(threadId) { + return await Promise.resolve((storeByThread.get(threadId) ?? []).length + 1); }, }; } diff --git a/packages/@n8n/instance-ai/src/event-bus/event-bus.interface.ts b/packages/@n8n/instance-ai/src/event-bus/event-bus.interface.ts index adc1625beb5..e312cf031cb 100644 --- a/packages/@n8n/instance-ai/src/event-bus/event-bus.interface.ts +++ b/packages/@n8n/instance-ai/src/event-bus/event-bus.interface.ts @@ -45,7 +45,8 @@ export interface InstanceAiEventBus { /** * Get the next event ID that will be assigned for a thread. - * Useful for the SSE endpoint to know whether there are events to replay. + * Used to seed the frontend's SSE replay cursor after message hydration. + * Async because multi-main implementations read a shared sequence. */ - getNextEventId(threadId: string): number; + getNextEventId(threadId: string): Promise; } diff --git a/packages/cli/src/modules/instance-ai/__tests__/instance-ai.controller.test.ts b/packages/cli/src/modules/instance-ai/__tests__/instance-ai.controller.test.ts index 2535d0bca3e..177dd947ff2 100644 --- a/packages/cli/src/modules/instance-ai/__tests__/instance-ai.controller.test.ts +++ b/packages/cli/src/modules/instance-ai/__tests__/instance-ai.controller.test.ts @@ -1121,7 +1121,7 @@ describe('InstanceAiController', () => { it('should return rich messages with nextEventId', async () => { const richResult = mock>(); memoryService.getRichMessages.mockResolvedValue(richResult); - eventBus.getNextEventId.mockReturnValue(42); + eventBus.getNextEventId.mockResolvedValue(42); const query = mock({ limit: 50, page: 0, diff --git a/packages/cli/src/modules/instance-ai/event-bus/__tests__/in-process-event-bus.test.ts b/packages/cli/src/modules/instance-ai/event-bus/__tests__/in-process-event-bus.test.ts index 24d264875a2..06db9ecdaea 100644 --- a/packages/cli/src/modules/instance-ai/event-bus/__tests__/in-process-event-bus.test.ts +++ b/packages/cli/src/modules/instance-ai/event-bus/__tests__/in-process-event-bus.test.ts @@ -1,5 +1,6 @@ import type { Logger } from '@n8n/backend-common'; import type { InstanceAiEvent } from '@n8n/api-types'; +import type { GlobalConfig } from '@n8n/config'; import { mock } from 'vitest-mock-extended'; import type { InstanceSettings } from 'n8n-core'; @@ -16,21 +17,81 @@ function makeEvent(type: string, runId: string): InstanceAiEvent { }; } +/** Flush the per-thread drain: each batch awaits one (mock) Redis round trip. */ +async function flushDrain() { + await new Promise((resolve) => setImmediate(resolve)); +} + describe('InProcessEventBus', () => { let bus: InProcessEventBus; let publisher: ReturnType>; let instanceSettings: { isMultiMain: boolean }; + /** Shared fake Redis sequence — one Map plays the role of the Redis server, + * so two bus instances built in one test behave like two mains. */ + let seqByKey: Map; + let redisFailure: Error | null; + let incrbyCalls: Array<{ key: string; count: number }>; + let deletedKeys: string[]; + + const redisClient = { + multi: (): unknown => { + let incrArgs: { key: string; count: number } | null = null; + const chain = { + incrby(key: string, count: number) { + incrArgs = { key, count }; + return chain; + }, + expire() { + return chain; + }, + async exec() { + if (redisFailure) throw redisFailure; + const { key, count } = incrArgs!; + incrbyCalls.push({ key, count }); + const value = (seqByKey.get(key) ?? 0) + count; + seqByKey.set(key, value); + return [ + [null, value], + [null, 1], + ]; + }, + }; + return chain; + }, + async get(key: string) { + if (redisFailure) throw redisFailure; + const value = seqByKey.get(key); + return value === undefined ? null : String(value); + }, + async del(key: string) { + deletedKeys.push(key); + seqByKey.delete(key); + return 1; + }, + }; + function buildBus() { const logger = mock(); logger.scoped.mockReturnValue(logger); publisher = mock(); publisher.publishCommand.mockResolvedValue(undefined); - return new InProcessEventBus(logger, instanceSettings as InstanceSettings, publisher); + publisher.getClient.mockReturnValue(redisClient as never); + const globalConfig = mock({ redis: { prefix: 'n8n' } }); + return new InProcessEventBus( + logger, + instanceSettings as InstanceSettings, + publisher, + globalConfig, + ); } beforeEach(() => { instanceSettings = { isMultiMain: false }; + seqByKey = new Map(); + redisFailure = null; + incrbyCalls = []; + deletedKeys = []; bus = buildBus(); }); @@ -38,8 +99,8 @@ describe('InProcessEventBus', () => { bus.clear(); }); - describe('publish', () => { - it('should assign monotonically increasing IDs per thread', () => { + describe('publish (single-main)', () => { + it('should assign monotonically increasing IDs per thread in the same tick', () => { bus.publish('thread-1', makeEvent('a', 'run_1')); bus.publish('thread-1', makeEvent('b', 'run_1')); bus.publish('thread-1', makeEvent('c', 'run_1')); @@ -66,6 +127,81 @@ describe('InProcessEventBus', () => { expect(events2).toHaveLength(1); expect(events2[0].id).toBe(1); }); + + it('should not touch Redis', () => { + bus.publish('thread-1', makeEvent('a', 'run_1')); + expect(seqByKey.size).toBe(0); + }); + }); + + describe('publish (multi-main, shared sequence)', () => { + beforeEach(() => { + instanceSettings = { isMultiMain: true }; + bus = buildBus(); + }); + + it('assigns ids from the shared Redis sequence in publish order', async () => { + bus.publish('thread-1', makeEvent('a', 'run_1')); + bus.publish('thread-1', makeEvent('b', 'run_1')); + bus.publish('thread-1', makeEvent('c', 'run_1')); + await flushDrain(); + + const events = bus.getEventsAfter('thread-1', 0); + expect(events.map((e) => e.id)).toEqual([1, 2, 3]); + expect(events.map((e) => e.event.payload)).toEqual([ + { text: 'a-run_1' }, + { text: 'b-run_1' }, + { text: 'c-run_1' }, + ]); + }); + + it('sequences events queued during a Redis round trip as one INCRBY batch', async () => { + bus.publish('thread-1', makeEvent('a', 'run_1')); // drains alone + bus.publish('thread-1', makeEvent('b', 'run_1')); // queued during the round trip + bus.publish('thread-1', makeEvent('c', 'run_1')); // queued during the round trip + await flushDrain(); + + expect(incrbyCalls.map((c) => c.count)).toEqual([1, 2]); + expect(bus.getEventsAfter('thread-1', 0).map((e) => e.id)).toEqual([1, 2, 3]); + }); + + it('continues the sequence started by another main', async () => { + const otherMain = buildBus(); + otherMain.publish('thread-1', makeEvent('a', 'run_1')); + otherMain.publish('thread-1', makeEvent('b', 'run_1')); + await flushDrain(); + + bus.publish('thread-1', makeEvent('c', 'run_1')); + await flushDrain(); + + expect(bus.getEventsAfter('thread-1', 0).map((e) => e.id)).toEqual([3]); + }); + + it('falls back to local ids above the high-water mark when Redis fails', async () => { + bus.publish('thread-1', makeEvent('a', 'run_1')); + bus.publish('thread-1', makeEvent('b', 'run_1')); + await flushDrain(); + + redisFailure = new Error('connection lost'); + bus.publish('thread-1', makeEvent('c', 'run_1')); + await flushDrain(); + + expect(bus.getEventsAfter('thread-1', 0).map((e) => e.id)).toEqual([1, 2, 3]); + }); + + it('keeps fallback ids above ids observed from relayed events', async () => { + redisFailure = new Error('connection lost'); + // A sibling produced up to id 7 — observed via relay without a subscriber. + bus.handleRelayInstanceAiEvent({ + threadId: 'thread-1', + storedEvent: { id: 7, event: makeEvent('x', 'run_1') }, + }); + + bus.publish('thread-1', makeEvent('a', 'run_1')); + await flushDrain(); + + expect(bus.getEventsAfter('thread-1', 0).map((e) => e.id)).toEqual([8]); + }); }); describe('subscribe', () => { @@ -140,15 +276,38 @@ describe('InProcessEventBus', () => { }); describe('getNextEventId', () => { - it('should return 1 for a new thread', () => { - expect(bus.getNextEventId('thread-1')).toBe(1); + it('should return 1 for a new thread', async () => { + await expect(bus.getNextEventId('thread-1')).resolves.toBe(1); }); - it('should return the next sequential ID after publishing', () => { + it('should return the next sequential ID after publishing', async () => { bus.publish('thread-1', makeEvent('a', 'run_1')); bus.publish('thread-1', makeEvent('b', 'run_1')); - expect(bus.getNextEventId('thread-1')).toBe(3); + await expect(bus.getNextEventId('thread-1')).resolves.toBe(3); + }); + + it('reads the shared sequence in multi-main, so any main returns the same cursor', async () => { + instanceSettings = { isMultiMain: true }; + bus = buildBus(); + const otherMain = buildBus(); + + otherMain.publish('thread-1', makeEvent('a', 'run_1')); + otherMain.publish('thread-1', makeEvent('b', 'run_1')); + await flushDrain(); + + // This main never buffered the thread, but agrees on the next id. + await expect(bus.getNextEventId('thread-1')).resolves.toBe(3); + }); + + it('falls back to the local high-water mark when Redis fails', async () => { + instanceSettings = { isMultiMain: true }; + bus = buildBus(); + bus.publish('thread-1', makeEvent('a', 'run_1')); + await flushDrain(); + + redisFailure = new Error('connection lost'); + await expect(bus.getNextEventId('thread-1')).resolves.toBe(2); }); }); @@ -186,10 +345,22 @@ describe('InProcessEventBus', () => { expect(events[0]).toHaveProperty('runId'); expect(events[0]).toHaveProperty('agentId'); }); + + it('includes events still awaiting a sequence number (same-main read-your-writes)', () => { + instanceSettings = { isMultiMain: true }; + bus = buildBus(); + + bus.publish('thread-1', makeEvent('a', 'run_1')); + + // No drain flush: the event has no id yet, but same-main callers + // (terminal outcomes, tracing, snapshots) must still see it. + expect(bus.getEventsForRun('thread-1', 'run_1')).toHaveLength(1); + expect(bus.getEventsAfter('thread-1', 0)).toHaveLength(0); + }); }); describe('clear', () => { - it('should remove all stored events and listeners', () => { + it('should remove all stored events and listeners', async () => { const received: Array<{ id: number; event: InstanceAiEvent }> = []; bus.subscribe('thread-1', (stored) => received.push(stored)); @@ -200,7 +371,7 @@ describe('InProcessEventBus', () => { // Events cleared expect(bus.getEventsAfter('thread-1', 0)).toEqual([]); - expect(bus.getNextEventId('thread-1')).toBe(1); + await expect(bus.getNextEventId('thread-1')).resolves.toBe(1); // Listener removed — new publish should not reach old handler bus.publish('thread-1', makeEvent('b', 'run_1')); @@ -208,39 +379,65 @@ describe('InProcessEventBus', () => { }); }); + describe('clearThread', () => { + it('deletes the shared sequence key in multi-main', async () => { + instanceSettings = { isMultiMain: true }; + bus = buildBus(); + bus.publish('thread-1', makeEvent('a', 'run_1')); + await flushDrain(); + expect(seqByKey.size).toBe(1); + + bus.clearThread('thread-1'); + await flushDrain(); + + expect(deletedKeys).toEqual(['n8n:instance-ai:event-seq:thread-1']); + expect(bus.getEventsAfter('thread-1', 0)).toEqual([]); + }); + + it('does not touch Redis in single-main', async () => { + bus.publish('thread-1', makeEvent('a', 'run_1')); + bus.clearThread('thread-1'); + await flushDrain(); + + expect(deletedKeys).toEqual([]); + }); + }); + describe('cross-main relay', () => { it('does not relay when single-main', () => { bus.publish('thread-1', makeEvent('a', 'run_1')); expect(publisher.publishCommand).not.toHaveBeenCalled(); }); - it('relays each event via pubsub when multi-main', () => { - instanceSettings.isMultiMain = true; + it('relays each event with its producer-assigned id when multi-main', async () => { + instanceSettings = { isMultiMain: true }; bus = buildBus(); const event = makeEvent('a', 'run_1'); bus.publish('thread-1', event); + await flushDrain(); expect(publisher.publishCommand).toHaveBeenCalledWith({ command: 'relay-instance-ai-event', - payload: { threadId: 'thread-1', event }, + payload: { threadId: 'thread-1', storedEvent: { id: 1, event } }, }); }); - it('still delivers locally even when relaying', () => { - instanceSettings.isMultiMain = true; + it('still delivers locally even when relaying', async () => { + instanceSettings = { isMultiMain: true }; bus = buildBus(); const received: number[] = []; bus.subscribe('thread-1', (e) => received.push(e.id)); bus.publish('thread-1', makeEvent('a', 'run_1')); + await flushDrain(); expect(received).toEqual([1]); expect(publisher.publishCommand).toHaveBeenCalledTimes(1); }); - it('skips relay for oversized events but still delivers locally', () => { - instanceSettings.isMultiMain = true; + it('skips relay for oversized events but still delivers locally', async () => { + instanceSettings = { isMultiMain: true }; bus = buildBus(); const received: number[] = []; bus.subscribe('thread-1', (e) => received.push(e.id)); @@ -248,42 +445,67 @@ describe('InProcessEventBus', () => { (huge.payload as { text: string }).text = 'x'.repeat(6 * 1024 * 1024); bus.publish('thread-1', huge); + await flushDrain(); // Relay skipped (would bloat pubsub), but the local SSE client still got it - // synchronously via the emit (even though the 2 MB store cap then evicts it). + // via the emit (even though the 2 MB store cap then evicts it). expect(publisher.publishCommand).not.toHaveBeenCalled(); expect(received).toEqual([1]); }); - - it('publishLocalOnly never relays', () => { - instanceSettings.isMultiMain = true; - bus = buildBus(); - - bus.publishLocalOnly('thread-1', makeEvent('a', 'run_1')); - - expect(publisher.publishCommand).not.toHaveBeenCalled(); - expect(bus.getEventsAfter('thread-1', 0)).toHaveLength(1); - }); }); describe('handleRelayInstanceAiEvent', () => { - it('re-emits a relayed event when this main holds an SSE subscriber', () => { + it('stores and re-emits a relayed event under its producer-assigned id', () => { const received: number[] = []; bus.subscribe('thread-1', (e) => received.push(e.id)); - bus.handleRelayInstanceAiEvent({ threadId: 'thread-1', event: makeEvent('a', 'run_1') }); + bus.handleRelayInstanceAiEvent({ + threadId: 'thread-1', + storedEvent: { id: 42, event: makeEvent('a', 'run_1') }, + }); - expect(received).toEqual([1]); - // Re-emit must not re-relay (loop guard): publishLocalOnly path. + expect(received).toEqual([42]); + expect(bus.getEventsAfter('thread-1', 0).map((e) => e.id)).toEqual([42]); + // Re-emit must not re-relay (loop guard). expect(publisher.publishCommand).not.toHaveBeenCalled(); }); it('ignores a relayed event when this main has no subscriber for the thread', () => { - bus.handleRelayInstanceAiEvent({ threadId: 'thread-1', event: makeEvent('a', 'run_1') }); + bus.handleRelayInstanceAiEvent({ + threadId: 'thread-1', + storedEvent: { id: 1, event: makeEvent('a', 'run_1') }, + }); // Nothing stored, since the thread has no local consumer here. expect(bus.getEventsAfter('thread-1', 0)).toHaveLength(0); }); + + it('keeps the store sorted when a concurrent producer relays a lower id', () => { + bus.subscribe('thread-1', () => {}); + + bus.handleRelayInstanceAiEvent({ + threadId: 'thread-1', + storedEvent: { id: 5, event: makeEvent('later', 'run_1') }, + }); + bus.handleRelayInstanceAiEvent({ + threadId: 'thread-1', + storedEvent: { id: 3, event: makeEvent('earlier', 'run_2') }, + }); + + expect(bus.getEventsAfter('thread-1', 0).map((e) => e.id)).toEqual([3, 5]); + }); + + it('drops a duplicate id instead of storing or emitting it twice', () => { + const received: number[] = []; + bus.subscribe('thread-1', (e) => received.push(e.id)); + const storedEvent = { id: 5, event: makeEvent('a', 'run_1') }; + + bus.handleRelayInstanceAiEvent({ threadId: 'thread-1', storedEvent }); + bus.handleRelayInstanceAiEvent({ threadId: 'thread-1', storedEvent }); + + expect(received).toEqual([5]); + expect(bus.getEventsAfter('thread-1', 0)).toHaveLength(1); + }); }); describe('hasSubscribers', () => { diff --git a/packages/cli/src/modules/instance-ai/event-bus/in-process-event-bus.ts b/packages/cli/src/modules/instance-ai/event-bus/in-process-event-bus.ts index 80447aadf3d..858d87736a6 100644 --- a/packages/cli/src/modules/instance-ai/event-bus/in-process-event-bus.ts +++ b/packages/cli/src/modules/instance-ai/event-bus/in-process-event-bus.ts @@ -1,5 +1,6 @@ import { Logger } from '@n8n/backend-common'; import type { InstanceAiEvent } from '@n8n/api-types'; +import { GlobalConfig } from '@n8n/config'; import { OnPubSubEvent } from '@n8n/decorators'; import { Service } from '@n8n/di'; import type { InstanceAiEventBus, StoredEvent } from '@n8n/instance-ai'; @@ -12,6 +13,15 @@ import { Publisher } from '@/scaling/pubsub/publisher.service'; const MAX_EVENTS_PER_THREAD = 500; const MAX_BYTES_PER_THREAD = 2 * 1024 * 1024; // 2 MB +/** + * How long an idle thread's shared sequence key lives in Redis (refreshed on + * every assignment). Generous on purpose: if it ever expires and the sequence + * restarts at 1, clients holding stale high cursors get an empty replay + * (recovered via run-sync / hydration), and fresh page loads re-seed their + * cursor from `GET /messages` anyway. + */ +const SEQ_KEY_TTL_SECONDS = 14 * 24 * 60 * 60; + @Service() export class InProcessEventBus implements InstanceAiEventBus { private readonly emitter = new EventEmitter(); @@ -21,50 +31,182 @@ export class InProcessEventBus implements InstanceAiEventBus { /** Approximate serialized size per thread for eviction. */ private readonly sizeBytes = new Map(); - /** Monotonic counter per thread — never resets even after eviction. */ - private readonly nextId = new Map(); + /** + * Highest event id this main has assigned or observed per thread. The id + * source in single-main, and the fallback when Redis is unavailable in + * multi-main (kept bumped from relayed events so fallback ids stay above + * what siblings have already used). + */ + private readonly lastLocalId = new Map(); + + /** + * Events awaiting a sequence number (multi-main only). `publish()` stays + * synchronous by enqueueing here; a single per-thread drain assigns ids. + */ + private readonly pendingByThread = new Map(); + + /** + * The batch currently being sequenced (multi-main only): taken off the + * pending queue but not yet in the store. Kept visible so run-scoped reads + * see events across the Redis round trip. + */ + private readonly inFlightByThread = new Map(); + + private readonly drainingThreads = new Set(); + + private readonly seqKeyPrefix: string; constructor( private readonly logger: Logger, private readonly instanceSettings: InstanceSettings, private readonly publisher: Publisher, + globalConfig: GlobalConfig, ) { this.logger = this.logger.scoped('instance-ai'); + this.seqKeyPrefix = `${globalConfig.redis.prefix}:instance-ai:event-seq:`; // Avoid warnings when many SSE clients connect (each adds a listener per thread) this.emitter.setMaxListeners(0); } /** - * Publish an event for a thread: store it, deliver it to local SSE - * subscribers, and — in multi-main — relay it to sibling mains so the main - * holding the client's SSE connection (which may not be this one) delivers it. + * Publish an event for a thread. + * + * Single-main: assign the next local id and deliver in the same tick. + * + * Multi-main: enqueue and drain asynchronously — event ids come from a + * shared per-thread Redis sequence, so every main agrees on them and the + * frontend's replay cursor is valid against any main. The queue preserves + * publish order; each sequenced event is stored, delivered to local SSE + * subscribers, and relayed to sibling mains with its id. */ publish(threadId: string, event: InstanceAiEvent): void { - // Serialize once: reused for the store's size accounting and the relay guard. - const sizeBytes = Buffer.byteLength(JSON.stringify(event), 'utf8'); - this.storeAndEmit(threadId, event, sizeBytes); - this.relayToSiblings(threadId, event, sizeBytes); + if (!this.instanceSettings.isMultiMain) { + const id = (this.lastLocalId.get(threadId) ?? 0) + 1; + this.lastLocalId.set(threadId, id); + this.storeAndEmit(threadId, { id, event }); + return; + } + + const pending = this.pendingByThread.get(threadId); + if (pending) { + pending.push(event); + } else { + this.pendingByThread.set(threadId, [event]); + } + void this.drainQueue(threadId); } /** - * Store + deliver locally WITHOUT relaying. Used by the pubsub handler when a - * relayed event arrives from another main — re-relaying would loop. The local - * `nextId` stamps the SSE id, so this main is the id authority for the - * connection it serves. + * Assign sequence ids to queued events and dispatch them, preserving + * publish order. Only one drain runs per thread; events queued while a + * Redis round trip is in flight are picked up by the next loop iteration + * and sequenced as one batch (single sequence round trip). */ - publishLocalOnly(threadId: string, event: InstanceAiEvent): void { - this.storeAndEmit(threadId, event, Buffer.byteLength(JSON.stringify(event), 'utf8')); + private async drainQueue(threadId: string): Promise { + if (this.drainingThreads.has(threadId)) return; + this.drainingThreads.add(threadId); + try { + let batch = this.takePending(threadId); + while (batch.length > 0) { + this.inFlightByThread.set(threadId, batch); + // Never throws — falls back to local ids on Redis failure. + const firstId = await this.assignSequenceBlock(threadId, batch.length); + for (let i = 0; i < batch.length; i++) { + const stored: StoredEvent = { id: firstId + i, event: batch[i] }; + // Serialize once: reused for the store's size accounting and the relay guard. + const sizeBytes = Buffer.byteLength(JSON.stringify(batch[i]), 'utf8'); + this.storeAndEmit(threadId, stored, sizeBytes); + this.relayToSiblings(threadId, stored, sizeBytes); + } + this.inFlightByThread.delete(threadId); + batch = this.takePending(threadId); + } + } finally { + this.inFlightByThread.delete(threadId); + this.drainingThreads.delete(threadId); + } } - private storeAndEmit(threadId: string, event: InstanceAiEvent, eventSizeBytes: number): void { + private takePending(threadId: string): InstanceAiEvent[] { + const pending = this.pendingByThread.get(threadId); + if (!pending) return []; + this.pendingByThread.delete(threadId); + return pending; + } + + /** + * Reserve a contiguous block of `count` ids from the shared per-thread + * sequence (atomic INCRBY). On Redis failure, continue monotonically from the + * local high-water mark — ids stay usable for this main's connections, at the + * cost of possible overlap with siblings until Redis recovers. + * + * Accepted degradation: after a Redis outage the shared counter can briefly + * sit below this main's local high-water mark (the fallback advanced local ids + * that never reached Redis), so INCRBY on recovery may re-issue an id already + * in this main's store — `insertById` then drops it as a duplicate, i.e. a few + * events can be lost from the live stream during recovery. Not worth an atomic + * conditional-max (Lua/WATCH) here: it only bites during a Redis incident, and + * the persisted run snapshot reconciles the tree via `run-sync` on reconnect. + * (A single-main→multi-main flip mid-thread would collide the same way, but a + * thread only starts producing events once the license has settled isMultiMain + * at boot, so that path isn't reached in practice.) + */ + private async assignSequenceBlock(threadId: string, count: number): Promise { + try { + const key = this.seqKey(threadId); + const results = await this.getRedisClient() + .multi() + .incrby(key, count) + .expire(key, SEQ_KEY_TTL_SECONDS) + .exec(); + const [incrError, incrResult] = results?.[0] ?? [new Error('empty transaction result'), null]; + if (incrError) throw incrError; + const endId = Number(incrResult); + if (!Number.isFinite(endId)) { + throw new Error(`non-numeric INCRBY result: ${String(incrResult)}`); + } + this.bumpLocalHighWaterMark(threadId, endId); + return endId - count + 1; + } catch (error) { + this.logger.error( + 'Failed to assign Instance AI event sequence from Redis, falling back to local ids', + { threadId, error }, + ); + const firstId = (this.lastLocalId.get(threadId) ?? 0) + 1; + this.lastLocalId.set(threadId, firstId + count - 1); + return firstId; + } + } + + /** + * The shared sequence lives on the pubsub publisher's Redis client. Only ever + * reached in multi-main, which implies queue mode — where the publisher's + * client is initialized. Reusing it avoids a second persistent connection per + * main. Publishing never puts a client in subscriber mode, so running + * sequence commands on it is safe. + */ + private getRedisClient() { + return this.publisher.getClient(); + } + + private seqKey(threadId: string): string { + return `${this.seqKeyPrefix}${threadId}`; + } + + private bumpLocalHighWaterMark(threadId: string, id: number): void { + if (id > (this.lastLocalId.get(threadId) ?? 0)) { + this.lastLocalId.set(threadId, id); + } + } + + private storeAndEmit(threadId: string, stored: StoredEvent, eventSizeBytes?: number): void { + const size = eventSizeBytes ?? Buffer.byteLength(JSON.stringify(stored.event), 'utf8'); const events = this.getOrCreateStore(threadId); - const id = (this.nextId.get(threadId) ?? 0) + 1; - this.nextId.set(threadId, id); - const stored: StoredEvent = { id, event }; + // Duplicate id (e.g. an event relayed twice): already stored and emitted. + if (!this.insertById(events, stored)) return; - events.push(stored); - this.sizeBytes.set(threadId, (this.sizeBytes.get(threadId) ?? 0) + eventSizeBytes); + this.sizeBytes.set(threadId, (this.sizeBytes.get(threadId) ?? 0) + size); // Evict oldest events if count or size exceeds caps this.evictIfNeeded(threadId, events); @@ -72,19 +214,40 @@ export class InProcessEventBus implements InstanceAiEventBus { this.emitter.emit(threadId, stored); } - private relayToSiblings(threadId: string, event: InstanceAiEvent, sizeBytes: number): void { + /** + * Insert keeping the store sorted by id. Local events always append, but a + * relayed event from a concurrent producer on another main (e.g. a + * background task while the orchestrator runs elsewhere) can arrive with a + * lower id than the latest stored one. Returns false for a duplicate id. + */ + private insertById(events: StoredEvent[], stored: StoredEvent): boolean { + if (events.length === 0 || events[events.length - 1].id < stored.id) { + events.push(stored); + return true; + } + let i = events.length - 1; + while (i >= 0 && events[i].id > stored.id) i--; + if (i >= 0 && events[i].id === stored.id) return false; + events.splice(i + 1, 0, stored); + return true; + } + + private relayToSiblings(threadId: string, stored: StoredEvent, sizeBytes: number): void { if (!this.instanceSettings.isMultiMain) return; if (sizeBytes > MAX_PUBSUB_PAYLOAD_BYTES) { this.logger.warn( - `Skipping cross-main relay of "${event.type}" event (${sizeBytes} bytes exceeds ${MAX_PUBSUB_PAYLOAD_BYTES})`, - { threadId, runId: event.runId }, + `Skipping cross-main relay of "${stored.event.type}" event (${sizeBytes} bytes exceeds ${MAX_PUBSUB_PAYLOAD_BYTES})`, + { threadId, runId: stored.event.runId }, ); return; } void this.publisher - .publishCommand({ command: 'relay-instance-ai-event', payload: { threadId, event } }) + .publishCommand({ + command: 'relay-instance-ai-event', + payload: { threadId, storedEvent: stored }, + }) .catch((error: unknown) => this.logger.error('Failed to relay Instance AI event to sibling mains', { threadId, @@ -93,16 +256,19 @@ export class InProcessEventBus implements InstanceAiEventBus { ); } - /** A relayed event from another main: re-emit to this main's SSE clients only - * if it actually holds a subscription for the thread (avoids every main - * buffering every thread). */ + /** A relayed event from another main, carrying its producer-assigned id + * from the shared sequence. Stored/re-emitted only if this main holds a + * subscription for the thread (avoids every main buffering every thread). */ @OnPubSubEvent('relay-instance-ai-event', { instanceType: 'main' }) handleRelayInstanceAiEvent({ threadId, - event, - }: { threadId: string; event: InstanceAiEvent }): void { + storedEvent, + }: { threadId: string; storedEvent: StoredEvent }): void { + // Track the shared-sequence high-water mark even without subscribers, so + // a Redis-outage fallback keeps assigning ids above what siblings used. + this.bumpLocalHighWaterMark(threadId, storedEvent.id); if (!this.hasSubscribers(threadId)) return; - this.publishLocalOnly(threadId, event); + this.storeAndEmit(threadId, storedEvent); } subscribe(threadId: string, handler: (storedEvent: StoredEvent) => void): () => void { @@ -115,6 +281,11 @@ export class InProcessEventBus implements InstanceAiEventBus { return this.emitter.listenerCount(threadId) > 0; } + /** + * Events still awaiting a sequence number are intentionally excluded: they + * have no id yet, and once sequenced they reach subscribers live — the SSE + * bootstrap subscribes before calling this, so nothing is missed. + */ getEventsAfter(threadId: string, afterId: number): StoredEvent[] { const events = this.store.get(threadId); if (!events) return []; @@ -122,35 +293,71 @@ export class InProcessEventBus implements InstanceAiEventBus { } getEventsForRun(threadId: string, runId: string): InstanceAiEvent[] { - const events = this.store.get(threadId); - if (!events) return []; - return events.filter((e) => e.event.runId === runId).map((e) => e.event); + return this.getEventsForRuns(threadId, [runId]); } getEventsForRuns(threadId: string, runIds: string[]): InstanceAiEvent[] { - const events = this.store.get(threadId); - if (!events || runIds.length === 0) return []; + if (runIds.length === 0) return []; const runIdSet = new Set(runIds); - return events.filter((e) => runIdSet.has(e.event.runId)).map((e) => e.event); + const stored = (this.store.get(threadId) ?? []) + .filter((e) => runIdSet.has(e.event.runId)) + .map((e) => e.event); + // Include events still awaiting a sequence number (both the batch being + // sequenced and the queue behind it) so same-main callers (terminal + // outcomes, tracing, snapshots) read their own writes. A run's events are + // produced on one main, so unsequenced ones are always newest. + const unsequenced = [ + ...(this.inFlightByThread.get(threadId) ?? []), + ...(this.pendingByThread.get(threadId) ?? []), + ].filter((e) => runIdSet.has(e.runId)); + return [...stored, ...unsequenced]; } - getNextEventId(threadId: string): number { - return (this.nextId.get(threadId) ?? 0) + 1; + async getNextEventId(threadId: string): Promise { + if (this.instanceSettings.isMultiMain) { + try { + const value = await this.getRedisClient().get(this.seqKey(threadId)); + if (value !== null) return Number(value) + 1; + } catch (error) { + this.logger.warn( + 'Failed to read Instance AI event sequence from Redis, falling back to local high-water mark', + { threadId, error }, + ); + } + } + return (this.lastLocalId.get(threadId) ?? 0) + 1; } /** Clear stored events for a specific thread (e.g. on thread expiration). */ clearThread(threadId: string): void { this.store.delete(threadId); this.sizeBytes.delete(threadId); - this.nextId.delete(threadId); + this.lastLocalId.delete(threadId); + this.pendingByThread.delete(threadId); + this.inFlightByThread.delete(threadId); this.emitter.removeAllListeners(threadId); + if (this.instanceSettings.isMultiMain) { + // Every main clears on thread deletion (task-control broadcast), so the + // shared key DEL is idempotent across mains. + void this.getRedisClient() + .del(this.seqKey(threadId)) + .catch((error: unknown) => + this.logger.warn('Failed to delete Instance AI event sequence key', { + threadId, + error, + }), + ); + } } - /** Clear all stored events. Used during module shutdown. */ + /** Clear all stored events. Used during module shutdown. Leaves the shared + * Redis sequence keys untouched — sibling mains still rely on them. */ clear(): void { this.store.clear(); this.sizeBytes.clear(); - this.nextId.clear(); + this.lastLocalId.clear(); + this.pendingByThread.clear(); + this.inFlightByThread.clear(); this.emitter.removeAllListeners(); } diff --git a/packages/cli/src/modules/instance-ai/instance-ai.controller.ts b/packages/cli/src/modules/instance-ai/instance-ai.controller.ts index 73ee8036185..a2eadab0b70 100644 --- a/packages/cli/src/modules/instance-ai/instance-ai.controller.ts +++ b/packages/cli/src/modules/instance-ai/instance-ai.controller.ts @@ -681,8 +681,9 @@ export class InstanceAiController { }); // Include the next SSE event ID so the frontend can skip past events - // already covered by these historical messages (prevents duplicates) - const nextEventId = this.eventBus.getNextEventId(threadId); + // already covered by these historical messages (prevents duplicates). + // Read from the shared sequence, so the cursor is valid against any main. + const nextEventId = await this.eventBus.getNextEventId(threadId); return { ...result, nextEventId }; } diff --git a/packages/cli/src/scaling/pubsub/pubsub.event-map.ts b/packages/cli/src/scaling/pubsub/pubsub.event-map.ts index 5c3611160bf..74353fa6d77 100644 --- a/packages/cli/src/scaling/pubsub/pubsub.event-map.ts +++ b/packages/cli/src/scaling/pubsub/pubsub.event-map.ts @@ -161,7 +161,12 @@ export type PubSubCommandMap = { */ 'relay-instance-ai-event': { threadId: string; - event: InstanceAiEvent; + /** + * Producer-assigned stored event. The id comes from the shared per-thread + * sequence, so every main stores and serves identical event ids and the + * frontend's replay cursor is valid against any main. + */ + storedEvent: { id: number; event: InstanceAiEvent }; }; /** diff --git a/packages/frontend/editor-ui/src/features/ai/instanceAi/__tests__/instanceAi.threadRuntime.test.ts b/packages/frontend/editor-ui/src/features/ai/instanceAi/__tests__/instanceAi.threadRuntime.test.ts index 04e997416f4..7fb53bcf8ea 100644 --- a/packages/frontend/editor-ui/src/features/ai/instanceAi/__tests__/instanceAi.threadRuntime.test.ts +++ b/packages/frontend/editor-ui/src/features/ai/instanceAi/__tests__/instanceAi.threadRuntime.test.ts @@ -395,6 +395,69 @@ describe('createThreadRuntime - SSE and hydration', () => { expect(registry.getRuntime(threadId)?.lastEventId).toBe(43); }); + test('an event replayed with an already-seen id is dropped', () => { + const threadId = activeThreadId; + const event = { + type: 'text-delta', + runId: 'run-1', + agentId: 'agent-root', + payload: { text: 'hello' }, + }; + + capturedOnMessage!(makeSSEEvent(validRunStartEvent('run-1', 'agent-root'), '1')); + capturedOnMessage!(makeSSEEvent(event, '2')); + // e.g. an auto-reconnect replaying an id that arrived just before the disconnect + capturedOnMessage!(makeSSEEvent(event, '2')); + + expect(registry.getRuntime(threadId)?.debugEvents).toHaveLength(2); + }); + + test('the reconnect cursor keeps the max seen id when producers interleave out of order', () => { + const threadId = activeThreadId; + + capturedOnMessage!(makeSSEEvent(validRunStartEvent('run-1', 'agent-root'), '43')); + // A concurrent producer on another main can relay a lower id afterwards. + capturedOnMessage!( + makeSSEEvent( + { + type: 'text-delta', + runId: 'run-1', + agentId: 'agent-root', + payload: { text: 'hello' }, + }, + '42', + ), + ); + + // The out-of-order event is still applied, but the cursor never regresses. + expect(registry.getRuntime(threadId)?.debugEvents).toHaveLength(2); + expect(registry.getRuntime(threadId)?.lastEventId).toBe(43); + }); + + test('a backend sequence reset (id 1 re-issued) drops stale dedup state and renders the fresh run', () => { + const threadId = activeThreadId; + const event = (text: string) => ({ + type: 'text-delta' as const, + runId: 'run-1', + agentId: 'agent-root', + payload: { text }, + }); + + // First run before the backend restarts. + capturedOnMessage!(makeSSEEvent(validRunStartEvent('run-1', 'agent-root'), '1')); + capturedOnMessage!(makeSSEEvent(event('a'), '2')); + expect(registry.getRuntime(threadId)?.lastEventId).toBe(2); + + // Backend restarts and re-issues ids from 1. Without reset detection these + // would be dropped as already-seen; instead the fresh sequence renders and + // the cursor snaps back down. + capturedOnMessage!(makeSSEEvent(validRunStartEvent('run-2', 'agent-root'), '1')); + capturedOnMessage!(makeSSEEvent(event('b'), '2')); + + expect(registry.getRuntime(threadId)?.debugEvents).toHaveLength(4); + expect(registry.getRuntime(threadId)?.lastEventId).toBe(2); + }); + test('deleting the last active thread clears stale routing state before the replacement thread starts', async () => { const deletedThreadId = activeThreadId; const previousEventSource = capturedInstance; diff --git a/packages/frontend/editor-ui/src/features/ai/instanceAi/instanceAi.threadRuntime.ts b/packages/frontend/editor-ui/src/features/ai/instanceAi/instanceAi.threadRuntime.ts index 6e32d0bd67c..c5274125a5d 100644 --- a/packages/frontend/editor-ui/src/features/ai/instanceAi/instanceAi.threadRuntime.ts +++ b/packages/frontend/editor-ui/src/features/ai/instanceAi/instanceAi.threadRuntime.ts @@ -75,6 +75,8 @@ export interface PendingConfirmationItem { export type HistoricalHydrationStatus = 'applied' | 'stale' | 'skipped'; const MAX_DEBUG_EVENTS = 1000; +/** Mirrors the backend's per-thread event buffer cap (MAX_EVENTS_PER_THREAD × 2). */ +const MAX_SEEN_EVENT_IDS = 1000; /** Silence window after which an active run with no stream traffic counts as stalled. */ const GENERATION_STALL_TIMEOUT_MS = 60_000; @@ -297,6 +299,10 @@ export function createThreadRuntime( const hydrationStatus = ref<'idle' | 'hydrating' | 'ready'>('idle'); const sseState = ref('disconnected'); const lastEventId = ref(undefined); + // Event ids already applied on this thread — guards against replay overlap, + // e.g. an auto-reconnect replaying an id that already arrived just before + // the disconnect. Not reactive: only consulted inside onSSEMessage. + const seenEventIds = new Set(); const amendContext = ref<{ agentId: string; role: string } | null>(null); const activePlanEdit = ref(null); const updatingPlanRequestIds = reactive(new Set()); @@ -606,9 +612,31 @@ export function createThreadRuntime( // --- SSE lifecycle --- function onSSEMessage(sseEvent: MessageEvent): void { - // Track last event ID for this thread (for reconnection) - if (sseEvent.lastEventId) { - lastEventId.value = Number(sseEvent.lastEventId); + // Event ids come from a shared per-thread sequence, so they are valid + // across mains — but concurrent producers on different mains can arrive + // interleaved out of order, so the reconnect cursor keeps the max seen + // rather than the latest, and duplicates are dropped by id. + const eventId = sseEvent.lastEventId ? Number(sseEvent.lastEventId) : undefined; + if (eventId !== undefined && Number.isFinite(eventId)) { + // A backend sequence reset (single-main restart, or seq-key TTL expiry) + // re-issues ids from 1. Seeing id 1 again while it's still in the dedup + // set means the sequence restarted, so drop the stale cursor + dedup + // state and render the fresh sequence instead of dropping it as + // duplicates. Precise, not a heuristic: id 1 is issued once per sequence + // lifetime, and a legit replay from cursor 0 can't reach here because + // the cursor and seenEventIds only ever reset together (in resetState). + if (eventId === 1 && seenEventIds.has(1)) { + seenEventIds.clear(); + lastEventId.value = undefined; + } + if (seenEventIds.has(eventId)) return; + seenEventIds.add(eventId); + if (seenEventIds.size > MAX_SEEN_EVENT_IDS) { + // Sets iterate in insertion order — evict the oldest (≈ lowest) id. + const oldest: number | undefined = seenEventIds.values().next().value; + if (oldest !== undefined) seenEventIds.delete(oldest); + } + lastEventId.value = Math.max(lastEventId.value ?? 0, eventId); } try { const parsed = instanceAiEventSchema.safeParse(JSON.parse(String(sseEvent.data))); @@ -827,6 +855,7 @@ export function createThreadRuntime( runStateByGroupId.clear(); groupIdByRunId.clear(); lastEventId.value = undefined; + seenEventIds.clear(); disarmGenerationStallWatchdog(); } @@ -866,8 +895,10 @@ export function createThreadRuntime( } // Set SSE cursor to skip past events already covered by historical messages. // This prevents duplicate messages when SSE replays in-memory events. + // Never move the cursor backwards: SSE may have advanced it while this + // request was in flight. if (result.nextEventId !== null && result.nextEventId !== undefined) { - lastEventId.value = result.nextEventId - 1; + lastEventId.value = Math.max(lastEventId.value ?? 0, result.nextEventId - 1); } if (result.projectId) projectId.value = result.projectId; return 'applied';