import { Logger } from '@n8n/backend-common'; import type { InstanceAiEvent } from '@n8n/api-types'; import { GlobalConfig } from '@n8n/config'; import type { InstanceAiEventBus, StoredEvent } from '@n8n/instance-ai'; import { InstanceSettings } from 'n8n-core'; import { Publisher } from '../../../scaling/pubsub/publisher.service'; import { DurableEventLog } from './durable-event-log'; export declare class InProcessEventBus implements InstanceAiEventBus { private readonly logger; private readonly instanceSettings; private readonly publisher; private readonly eventLog; private readonly emitter; private readonly store; private readonly sizeBytes; private readonly lastLocalId; private readonly pendingByThread; private readonly inFlightByThread; private readonly drainingThreads; private readonly seqKeyPrefix; private readonly durableLogEnabled; constructor(logger: Logger, instanceSettings: InstanceSettings, publisher: Publisher, eventLog: DurableEventLog, globalConfig: GlobalConfig); publish(threadId: string, event: InstanceAiEvent): void; private onDrained; private cacheSequencedEvent; private drainQueue; private takePending; private assignSequenceBlock; private getRedisClient; private seqKey; private bumpLocalHighWaterMark; private storeAndEmit; private insertById; private relayToSiblings; handleRelayInstanceAiEvent({ threadId, storedEvent, }: { threadId: string; storedEvent: StoredEvent; }): void; subscribe(threadId: string, handler: (storedEvent: StoredEvent) => void): () => void; hasSubscribers(threadId: string): boolean; getEventsAfter(threadId: string, afterId: number): StoredEvent[]; getEventsForRun(threadId: string, runId: string): InstanceAiEvent[]; getEventsForRuns(threadId: string, runIds: string[]): InstanceAiEvent[]; getNextEventId(threadId: string): Promise; clearThread(threadId: string): void; clear(): void; private evictIfNeeded; private getOrCreateStore; }