feat(core): Use a shared per-thread event sequence for Instance AI multi-main (no-changelog) (#33558)

Co-authored-by: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
Jaakko Husso
2026-07-09 07:44:20 +00:00
committed by GitHub
co-authored by Claude Fable 5
parent 52b72dd5a8
commit 4dd369991f
10 changed files with 626 additions and 89 deletions
@@ -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
@@ -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);
},
};
}
@@ -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<number>;
}
@@ -1121,7 +1121,7 @@ describe('InstanceAiController', () => {
it('should return rich messages with nextEventId', async () => {
const richResult = mock<Omit<InstanceAiRichMessagesResponse, 'nextEventId'>>();
memoryService.getRichMessages.mockResolvedValue(richResult);
eventBus.getNextEventId.mockReturnValue(42);
eventBus.getNextEventId.mockResolvedValue(42);
const query = mock<InstanceAiThreadMessagesQuery>({
limit: 50,
page: 0,
@@ -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<typeof mock<Publisher>>;
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<string, number>;
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>();
logger.scoped.mockReturnValue(logger);
publisher = mock<Publisher>();
publisher.publishCommand.mockResolvedValue(undefined);
return new InProcessEventBus(logger, instanceSettings as InstanceSettings, publisher);
publisher.getClient.mockReturnValue(redisClient as never);
const globalConfig = mock<GlobalConfig>({ 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', () => {
@@ -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<string, number>();
/** Monotonic counter per thread — never resets even after eviction. */
private readonly nextId = new Map<string, number>();
/**
* 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<string, number>();
/**
* 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<string, InstanceAiEvent[]>();
/**
* 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<string, InstanceAiEvent[]>();
private readonly drainingThreads = new Set<string>();
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<void> {
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<number> {
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<number> {
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();
}
@@ -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 };
}
@@ -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 };
};
/**
@@ -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;
@@ -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<InstanceAiSSEConnectionState>('disconnected');
const lastEventId = ref<number | undefined>(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<number>();
const amendContext = ref<{ agentId: string; role: string } | null>(null);
const activePlanEdit = ref<PlanEditContext | null>(null);
const updatingPlanRequestIds = reactive(new Set<string>());
@@ -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';