/** * Event Emitter for CDC (Change Data Capture) Subscriptions * * A framework-agnostic, type-safe event emitter designed for PostgreSQL * change data capture streams. This module provides: * * **Features:** * - Fully typed event handling with TypeScript generics * - One-time (`once`) and persistent (`on`) listeners * - Memory-efficient listener management with configurable limits * - Error isolation between handlers (configurable) * - Async handler support with `emitAsync` for sequential execution * - Event piping between emitter instances * - Promise-based `waitFor` for async/await patterns * - Node.js EventEmitter-compatible API (`addListener`, `removeListener`, etc.) * * **Event Flow:** * ``` * CDC Transport -> CDCEventEmitter -> Application Handlers * | * +-> Error isolation (optional) * +-> Max listener enforcement * +-> Once listener cleanup * +-> newListener/removeListener events * ``` * * **Memory Safety:** * The emitter enforces a maximum listener count per event type (default: 100) * to prevent memory leaks from unbounded listener registration. When the * limit is exceeded, a `CDCError` with code `TOO_MANY_SUBSCRIPTIONS` is thrown. * * @example Basic usage * ```typescript * import { createEventEmitter } from './event-emitter' * * interface User { * id: number * name: string * } * * const emitter = createEventEmitter() * * // Subscribe to inserts * emitter.on('insert', (user) => { * console.log('New user:', user.name) * }) * * // One-time connection handler * emitter.once('connected', () => { * console.log('CDC stream connected') * }) * * // Wait for an event with timeout * const user = await emitter.waitFor('insert', 5000) * ``` * * @module cdc/event-emitter */ import type { CDCChangeEvent, Row } from '../types' import { CDCError, CDCErrorCode } from './errors' import { CDCLogger, defaultLogger } from './logger' /** * Supported CDC event types. * * **Data Events:** * - `change` - Any data change (INSERT, UPDATE, DELETE, TRUNCATE) * - `insert` - New row inserted * - `update` - Existing row modified * - `delete` - Row deleted * - `truncate` - Table truncated * * **Connection Events:** * - `connected` - Successfully connected to CDC stream * - `disconnected` - Disconnected from CDC stream * - `reconnecting` - Attempting to reconnect after failure * * **Status Events:** * - `error` - Error occurred during CDC operations * - `stateChange` - Internal state machine transition * - `healthChange` - Connection health status changed * * **Meta Events (Node.js EventEmitter compatible):** * - `newListener` - Emitted when a new listener is added * - `removeListener` - Emitted when a listener is removed * * @example Listening to specific event types * ```typescript * emitter.on('insert', (row) => handleInsert(row)) * emitter.on('update', ({ newRow, oldRow }) => handleUpdate(newRow, oldRow)) * emitter.on('delete', (row) => handleDelete(row)) * emitter.on('error', (error) => handleError(error)) * ``` * * @example Listening to meta events * ```typescript * emitter.on('newListener', ({ event, listener }) => { * console.log(`New listener added for ${event}`) * }) * ``` */ export type CDCEventType = | 'change' | 'insert' | 'update' | 'delete' | 'truncate' | 'error' | 'connected' | 'disconnected' | 'reconnecting' | 'stateChange' | 'healthChange' | 'newListener' | 'removeListener' /** * Type-safe payload definitions for each CDC event type. * * This interface maps event types to their corresponding payload structures, * ensuring type safety when emitting and receiving events. The generic * parameter `T` allows customization of row types for data events. * * @typeParam T - The row type for data events (default: `Row`) * * @example Type-safe event handling * ```typescript * interface User { * id: number * name: string * email: string * } * * const emitter = createEventEmitter() * * // TypeScript knows `row` is of type `User` * emitter.on('insert', (row) => { * console.log(row.email) // Fully typed * }) * * // TypeScript knows this is { newRow: User; oldRow?: User } * emitter.on('update', ({ newRow, oldRow }) => { * console.log(newRow.name, oldRow?.name) * }) * ``` */ export interface CDCEventPayloads { /** Full change event with operation type and row data */ change: CDCChangeEvent & { newRow?: T; oldRow?: T } /** The newly inserted row */ insert: T /** Update event with new and optionally old row data */ update: { newRow: T; oldRow?: T } /** The deleted row (requires REPLICA IDENTITY FULL) */ delete: T /** Truncate event with table and schema info */ truncate: { table: string; schema: string } /** Error that occurred during CDC operations */ error: CDCError | Error /** Connection established (void payload) */ connected: void /** Connection closed with optional reason */ disconnected: { reason?: string } /** Reconnection attempt in progress */ reconnecting: { attempt: number; lastLsn?: string } /** Internal state machine transition */ stateChange: { from: string; to: string; event: string } /** Health status change */ healthChange: { status: string; score: number } /** Emitted when a new listener is added (Node.js EventEmitter compatible) */ newListener: { event: CDCEventType; listener: CDCEventListener } /** Emitted when a listener is removed (Node.js EventEmitter compatible) */ removeListener: { event: CDCEventType; listener: CDCEventListener } } /** * Event listener callback function type. * * Listeners can be synchronous or asynchronous. When using `emit()`, * async listeners run without waiting (fire-and-forget). Use `emitAsync()` * to wait for all handlers to complete. * * @typeParam T - The payload type for this listener * * @example Sync vs async listeners * ```typescript * // Synchronous handler * emitter.on('insert', (row) => { * cache.set(row.id, row) * }) * * // Asynchronous handler * emitter.on('insert', async (row) => { * await database.replicate(row) * }) * ``` */ export type CDCEventListener = (payload: T) => void | Promise /** * Internal listener registration entry. * * Tracks the handler function and whether it should be removed * after the first invocation (one-time listener). * * @internal */ interface ListenerEntry { /** The callback function to invoke */ handler: CDCEventListener /** If true, remove after first invocation */ once: boolean } /** * Configuration options for the CDC event emitter. * * @example Custom configuration * ```typescript * const emitter = createEventEmitter({ * maxListeners: 50, // Stricter limit for memory-constrained env * catchErrors: true, // Don't let one handler break others * logger: customLogger, // Custom logging implementation * }) * ``` */ export interface EventEmitterConfig { /** * Maximum number of listeners allowed per event type. * * This limit helps prevent memory leaks from unbounded listener * registration. When exceeded, `on()` and `once()` throw a * `CDCError` with code `TOO_MANY_SUBSCRIPTIONS`. * * Can be changed at runtime using `setMaxListeners()`. * * @default 100 */ maxListeners?: number /** * Whether to catch and log handler errors instead of propagating. * * When `true` (default), errors thrown by handlers are caught, * logged, and execution continues with remaining handlers. * * When `false`, errors propagate immediately, potentially * preventing subsequent handlers from executing. * * @default true */ catchErrors?: boolean /** * Logger instance for error and warning output. * * Uses the default CDC logger if not specified. * * @default defaultLogger */ logger?: CDCLogger } /** * CDC Event Emitter - A type-safe event emitter for change data capture. * * Provides a robust event system for handling PostgreSQL CDC events with: * - Full TypeScript type safety for event payloads * - Memory leak prevention via max listener limits * - Error isolation between handlers * - Async handler support * - Event piping and composition * - Node.js EventEmitter-compatible API * * **Listener Patterns:** * * 1. **Persistent Listeners** (`on`/`addListener`) - Called every time the event fires * 2. **One-time Listeners** (`once`) - Called once, then auto-removed * 3. **Prepend Listeners** (`prependListener`/`prependOnceListener`) - Added at the beginning * 4. **Promise-based** (`waitFor`) - Resolves on first event occurrence * * **Error Handling:** * * By default (`catchErrors: true`), errors in handlers are caught and logged, * allowing other handlers to continue executing. This prevents one faulty * handler from breaking the entire event chain. * * **Thread Safety:** * * The emitter creates a copy of the listener array before invoking handlers, * ensuring safe modification of listeners during emission (e.g., a handler * that removes itself). * * **Meta Events:** * * Like Node.js EventEmitter, this emitter fires `newListener` and `removeListener` * events when listeners are added or removed, allowing monitoring of listener * registration. * * @typeParam T - The row type for data events (default: `Row`) * * @example Complete usage example * ```typescript * interface User { * id: number * name: string * email: string * } * * // Create typed emitter * const emitter = new CDCEventEmitter({ * maxListeners: 50, * catchErrors: true, * }) * * // Persistent listener * emitter.on('insert', (user) => { * console.log('New user:', user.name) * }) * * // One-time listener * emitter.once('connected', () => { * console.log('Connected!') * }) * * // Error handling * emitter.on('error', (error) => { * if (error instanceof CDCError) { * console.error(error.code, error.message) * } * }) * * // Wait for specific event with timeout * try { * const user = await emitter.waitFor('insert', 5000) * console.log('Got user:', user) * } catch (e) { * console.log('Timeout waiting for insert') * } * * // Monitor listener changes * emitter.on('newListener', ({ event }) => { * console.log(`Listener added for: ${event}`) * }) * * // Cleanup * emitter.destroy() * ``` */ export class CDCEventEmitter { /** Map of event types to their registered listeners */ private listenersMap = new Map[]>() /** Maximum listeners per event (mutable via setMaxListeners) */ private _maxListeners: number /** Whether to catch and log handler errors */ private readonly catchErrors: boolean /** Logger instance for error and warning output */ private readonly logger: CDCLogger /** * Create a new CDC event emitter. * * @param config - Optional configuration options * * @example Create with defaults * ```typescript * const emitter = new CDCEventEmitter() * ``` * * @example Create with custom config * ```typescript * const emitter = new CDCEventEmitter({ * maxListeners: 50, * catchErrors: false, * logger: customLogger, * }) * ``` */ constructor(config: EventEmitterConfig = {}) { this._maxListeners = config.maxListeners ?? 100 this.catchErrors = config.catchErrors ?? true this.logger = config.logger ?? defaultLogger } /** * Static method to get listener count for an emitter. * * Provides Node.js EventEmitter-compatible static API for checking * listener counts without having a reference to the event type. * * @typeParam R - The row type of the emitter * @param emitter - The emitter to check * @param event - The event type to count listeners for * @returns The number of registered listeners * * @example Static listener count * ```typescript * const count = CDCEventEmitter.listenerCount(emitter, 'insert') * console.log(`${count} listeners for insert`) * ``` */ static listenerCount( emitter: CDCEventEmitter, event: CDCEventType ): number { return emitter.listenerCount(event) } /** * Add a persistent event listener. * * The listener will be called every time the specified event is emitted * until it is explicitly removed with `off()` or `removeAllListeners()`. * * Emits a `newListener` event before adding the listener. * * @typeParam E - The event type (inferred from first argument) * @param event - The event type to listen for * @param handler - The callback function to invoke * @returns `this` for method chaining * @throws {CDCError} When max listeners limit is exceeded (code: `TOO_MANY_SUBSCRIPTIONS`) * * @example Basic usage * ```typescript * emitter.on('insert', (row) => { * console.log('Inserted:', row) * }) * ``` * * @example Chaining * ```typescript * emitter * .on('connected', () => console.log('Connected')) * .on('disconnected', () => console.log('Disconnected')) * .on('error', (e) => console.error(e)) * ``` */ on( event: E, handler: CDCEventListener[E]> ): this { return this.addListenerInternal(event, handler, false, false) } /** * Alias for `on()`. Adds a persistent event listener. * * This method is provided for Node.js EventEmitter API compatibility. * * @typeParam E - The event type (inferred from first argument) * @param event - The event type to listen for * @param handler - The callback function to invoke * @returns `this` for method chaining * @throws {CDCError} When max listeners limit is exceeded * * @example Node.js style * ```typescript * emitter.addListener('insert', handler) * ``` */ addListener( event: E, handler: CDCEventListener[E]> ): this { return this.on(event, handler) } /** * Add a one-time event listener. * * The listener will be called only once when the event is first emitted, * then automatically removed. Useful for initialization or single-event * workflows. * * Emits a `newListener` event before adding the listener. * * @typeParam E - The event type (inferred from first argument) * @param event - The event type to listen for * @param handler - The callback function to invoke * @returns `this` for method chaining * @throws {CDCError} When max listeners limit is exceeded (code: `TOO_MANY_SUBSCRIPTIONS`) * * @example Wait for connection * ```typescript * emitter.once('connected', () => { * console.log('First connection established') * // This won't fire on reconnections * }) * ``` * * @example Capture first insert * ```typescript * emitter.once('insert', (row) => { * console.log('First row inserted:', row.id) * }) * ``` */ once( event: E, handler: CDCEventListener[E]> ): this { return this.addListenerInternal(event, handler, true, false) } /** * Add a listener at the beginning of the listeners array. * * Similar to `on()`, but the listener is added to the front of the * listener array, ensuring it will be called before listeners that * were added with `on()` or `addListener()`. * * Emits a `newListener` event before adding the listener. * * @typeParam E - The event type (inferred from first argument) * @param event - The event type to listen for * @param handler - The callback function to invoke * @returns `this` for method chaining * @throws {CDCError} When max listeners limit is exceeded * * @example Priority handler * ```typescript * // This handler will be called first * emitter.prependListener('insert', (row) => { * console.log('Priority handler:', row.id) * }) * * // This handler will be called second * emitter.on('insert', (row) => { * console.log('Normal handler:', row.id) * }) * ``` */ prependListener( event: E, handler: CDCEventListener[E]> ): this { return this.addListenerInternal(event, handler, false, true) } /** * Add a one-time listener at the beginning of the listeners array. * * Combines `once()` and `prependListener()` behavior: the listener * is called only once and is placed at the front of the listener array. * * Emits a `newListener` event before adding the listener. * * @typeParam E - The event type (inferred from first argument) * @param event - The event type to listen for * @param handler - The callback function to invoke * @returns `this` for method chaining * @throws {CDCError} When max listeners limit is exceeded * * @example Priority one-time handler * ```typescript * emitter.prependOnceListener('connected', () => { * console.log('Called first, only once') * }) * ``` */ prependOnceListener( event: E, handler: CDCEventListener[E]> ): this { return this.addListenerInternal(event, handler, true, true) } /** * Remove a specific event listener. * * Removes the listener by reference equality. The same function reference * that was passed to `on()` or `once()` must be provided. * * Emits a `removeListener` event after removing the listener. * * @typeParam E - The event type (inferred from first argument) * @param event - The event type to remove the listener from * @param handler - The exact handler function to remove * @returns `this` for method chaining * * @example Remove a listener * ```typescript * const handler = (row: User) => console.log(row) * * emitter.on('insert', handler) * // ... later ... * emitter.off('insert', handler) * ``` * * @example Safe removal (no-op if not found) * ```typescript * // This is safe even if the handler was never registered * emitter.off('insert', unknownHandler) * ``` */ off( event: E, handler: CDCEventListener[E]> ): this { const eventListeners = this.listenersMap.get(event) if (!eventListeners) return this const index = eventListeners.findIndex((l) => l.handler === handler) if (index !== -1) { eventListeners.splice(index, 1) // Emit removeListener event (but not for removeListener itself to avoid infinite loop) if (event !== 'newListener' && event !== 'removeListener') { this.emitInternal('removeListener', event, handler) } } return this } /** * Alias for `off()`. Removes a specific event listener. * * This method is provided for Node.js EventEmitter API compatibility. * * @typeParam E - The event type (inferred from first argument) * @param event - The event type to remove the listener from * @param handler - The exact handler function to remove * @returns `this` for method chaining * * @example Node.js style * ```typescript * emitter.removeListener('insert', handler) * ``` */ removeListener( event: E, handler: CDCEventListener[E]> ): this { return this.off(event, handler) } /** * Remove all listeners for a specific event or all events. * * When called with an event type, removes all listeners for that event. * When called without arguments, removes all listeners for all events. * * Note: Does not emit `removeListener` events for removed listeners. * * @param event - Optional event type to clear listeners for * @returns `this` for method chaining * * @example Remove all listeners for one event * ```typescript * emitter.removeAllListeners('insert') * ``` * * @example Remove all listeners for all events * ```typescript * emitter.removeAllListeners() * ``` */ removeAllListeners(event?: CDCEventType): this { if (event) { this.listenersMap.delete(event) } else { this.listenersMap.clear() } return this } /** * Emit an event to all registered listeners. * * Calls all listeners synchronously in registration order. Async listeners * are invoked but not awaited (fire-and-forget). Use `emitAsync()` if you * need to wait for async handlers. * * **Error Handling:** * - If `catchErrors: true` (default), errors are caught and logged * - If `catchErrors: false`, errors propagate and stop subsequent handlers * * **Listener Removal:** * One-time listeners (`once`) are removed before handlers are called, * ensuring they only execute once even if the handler throws. * * @typeParam E - The event type (inferred from first argument) * @param event - The event type to emit * @param payload - The event payload (type-checked against event type) * @returns `true` if there were listeners, `false` otherwise * * @example Emit an insert event * ```typescript * const hadListeners = emitter.emit('insert', { id: 1, name: 'Alice' }) * if (!hadListeners) { * console.log('No listeners registered for insert') * } * ``` * * @example Emit connection event * ```typescript * emitter.emit('connected', undefined) // void payload * ``` */ emit(event: E, payload: CDCEventPayloads[E]): boolean { const eventListeners = this.listenersMap.get(event) if (!eventListeners || eventListeners.length === 0) { return false } // Create a copy to handle removal during iteration const toCall = [...eventListeners] // Remove once listeners before calling this.listenersMap.set( event, eventListeners.filter((l) => !l.once) ) for (const { handler } of toCall) { try { const result = handler(payload) // Handle async handlers if (result instanceof Promise) { result.catch((error) => { if (this.catchErrors) { this.logger.error('Event handler error (async)', { data: { event }, error: error instanceof Error ? error : new Error(String(error)), }) } }) } } catch (error) { if (this.catchErrors) { this.logger.error('Event handler error', { data: { event }, error: error instanceof Error ? error : new Error(String(error)), }) } else { throw error } } } return true } /** * Emit an event and wait for all handlers to complete. * * Unlike `emit()`, this method awaits all handlers (including async ones) * in sequence before returning. Useful when you need to ensure all * handlers have completed before proceeding. * * **Error Handling:** * - If `catchErrors: true`, errors are collected but execution continues * - If `catchErrors: false`, first error is thrown immediately * * @typeParam E - The event type (inferred from first argument) * @param event - The event type to emit * @param payload - The event payload (type-checked against event type) * @returns `true` if all handlers succeeded, `false` if any errors occurred * * @example Wait for all handlers * ```typescript * const success = await emitter.emitAsync('insert', newRow) * if (!success) { * console.log('Some handlers failed') * } * ``` * * @example Use with catchErrors: false * ```typescript * try { * await emitter.emitAsync('insert', newRow) * } catch (e) { * console.error('Handler threw:', e) * } * ``` */ async emitAsync( event: E, payload: CDCEventPayloads[E] ): Promise { const eventListeners = this.listenersMap.get(event) if (!eventListeners || eventListeners.length === 0) { return false } const toCall = [...eventListeners] // Remove once listeners before calling this.listenersMap.set( event, eventListeners.filter((l) => !l.once) ) const errors: Error[] = [] for (const { handler } of toCall) { try { await handler(payload) } catch (error) { const err = error instanceof Error ? error : new Error(String(error)) if (this.catchErrors) { this.logger.error('Event handler error', { data: { event }, error: err, }) errors.push(err) } else { throw error } } } return errors.length === 0 } /** * Get the number of listeners for a specific event. * * @param event - The event type to count listeners for * @returns The number of registered listeners * * @example Check listener count * ```typescript * if (emitter.listenerCount('insert') === 0) { * console.warn('No insert handlers registered') * } * ``` */ listenerCount(event: CDCEventType): number { return this.listenersMap.get(event)?.length ?? 0 } /** * Get all event types that have registered listeners. * * @returns Array of event type names with at least one listener * * @example List active events * ```typescript * const events = emitter.eventNames() * console.log('Listening to:', events.join(', ')) * ``` */ eventNames(): CDCEventType[] { return Array.from(this.listenersMap.keys()) } /** * Check if an event has any registered listeners. * * @param event - The event type to check * @returns `true` if the event has at least one listener * * @example Conditional emission * ```typescript * if (emitter.hasListeners('insert')) { * emitter.emit('insert', newRow) * } * ``` */ hasListeners(event: CDCEventType): boolean { return this.listenerCount(event) > 0 } /** * Get a copy of the listeners array for an event. * * Returns the unwrapped handler functions, not the internal * `ListenerEntry` objects. Useful for inspecting registered handlers. * * @typeParam E - The event type (inferred from first argument) * @param event - The event type to get listeners for * @returns Array of handler functions (copy, not the internal array) * * @example Inspect listeners * ```typescript * const handlers = emitter.listeners('insert') * console.log(`${handlers.length} insert handlers registered`) * ``` */ listeners( event: E ): CDCEventListener[E]>[] { const eventListeners = this.listenersMap.get(event) if (!eventListeners) return [] return eventListeners.map((l) => l.handler as CDCEventListener[E]>) } /** * Get a copy of the raw listener entries for an event. * * Unlike `listeners()`, this returns the internal `ListenerEntry` objects * which include the `once` flag. Useful for debugging listener state. * * @typeParam E - The event type (inferred from first argument) * @param event - The event type to get listeners for * @returns Array of `ListenerEntry` objects (copy, not the internal array) * * @example Inspect listener entries * ```typescript * const entries = emitter.rawListeners('insert') * const onceCount = entries.filter(e => e.once).length * console.log(`${onceCount} one-time listeners`) * ``` */ rawListeners( event: E ): ListenerEntry[E]>[] { const eventListeners = this.listenersMap.get(event) if (!eventListeners) return [] return [...eventListeners] as ListenerEntry[E]>[] } /** * Get the current maximum listeners per event. * * @returns The maximum number of listeners allowed per event * * @example Check max listeners * ```typescript * console.log(`Max listeners: ${emitter.getMaxListeners()}`) * ``` */ getMaxListeners(): number { return this._maxListeners } /** * Set the maximum listeners per event. * * This can be used to increase or decrease the limit at runtime. * Setting to `Infinity` (or `0`) will disable the limit warning. * * @param n - The new maximum listeners limit * @returns `this` for method chaining * * @example Increase limit for high-throughput scenarios * ```typescript * emitter.setMaxListeners(200) * ``` * * @example Disable limit (use with caution) * ```typescript * emitter.setMaxListeners(Infinity) * ``` */ setMaxListeners(n: number): this { this._maxListeners = n return this } /** * Internal method to add a listener with configurable options. * * Handles the common logic for all listener addition methods including: * - Creating the listener array if needed * - Checking max listener limit * - Logging warnings when limit is exceeded * - Emitting `newListener` event * - Prepending vs appending to the array * * @internal * @throws {CDCError} When max listeners limit is exceeded */ private addListenerInternal( event: E, handler: CDCEventListener[E]>, once: boolean, prepend: boolean ): this { let eventListeners = this.listenersMap.get(event) if (!eventListeners) { eventListeners = [] this.listenersMap.set(event, eventListeners) } if (eventListeners.length >= this._maxListeners) { this.logger.warn('Max listeners exceeded', { data: { event, count: eventListeners.length, max: this._maxListeners }, }) throw new CDCError( `Max listeners (${this._maxListeners}) exceeded for event "${event}"`, { code: CDCErrorCode.TOO_MANY_SUBSCRIPTIONS, context: { event, count: eventListeners.length }, } ) } const entry = { handler: handler as CDCEventListener, once } if (prepend) { eventListeners.unshift(entry) } else { eventListeners.push(entry) } // Emit newListener event (but not for newListener itself to avoid infinite loop) if (event !== 'newListener' && event !== 'removeListener') { this.emitInternal('newListener', event, handler) } return this } /** * Internal emit for newListener/removeListener meta events. * * This is a separate method to handle the special case of meta events * that need to pass the event name and handler as a payload object. * * @internal */ private emitInternal( type: 'newListener' | 'removeListener', event: E, handler: CDCEventListener[E]> ): void { const eventListeners = this.listenersMap.get(type) if (!eventListeners || eventListeners.length === 0) return // Build the payload object expected by newListener/removeListener handlers const payload = { event, listener: handler as CDCEventListener } for (const { handler: listener } of eventListeners) { try { listener(payload) } catch (error) { if (this.catchErrors) { this.logger.error(`Event handler error (${type})`, { data: { event: type }, error: error instanceof Error ? error : new Error(String(error)), }) } else { throw error } } } } /** * Create a promise that resolves when an event is emitted. * * Registers a one-time listener and returns a Promise that resolves * with the event payload. Optionally supports a timeout. * * **Timeout Behavior:** * When a timeout is specified and reached, the promise rejects with * a `CDCError` (code: `CONNECTION_TIMEOUT`) and the listener is removed. * * @typeParam E - The event type (inferred from first argument) * @param event - The event type to wait for * @param timeoutMs - Optional timeout in milliseconds * @returns Promise that resolves with the event payload * @throws {CDCError} When timeout is reached (code: `CONNECTION_TIMEOUT`) * * @example Wait for connection * ```typescript * await emitter.waitFor('connected') * console.log('Connected!') * ``` * * @example Wait with timeout * ```typescript * try { * const user = await emitter.waitFor('insert', 5000) * console.log('Got user:', user) * } catch (e) { * if (e instanceof CDCError && e.code === CDCErrorCode.CONNECTION_TIMEOUT) { * console.log('Timeout waiting for insert') * } * } * ``` */ waitFor( event: E, timeoutMs?: number ): Promise[E]> { return new Promise((resolve, reject) => { let timeoutId: ReturnType | undefined const handler: CDCEventListener[E]> = (payload) => { if (timeoutId) { clearTimeout(timeoutId) } resolve(payload) } this.once(event, handler) if (timeoutMs) { timeoutId = setTimeout(() => { this.off(event, handler) reject( new CDCError(`Timeout waiting for event "${event}"`, { code: CDCErrorCode.CONNECTION_TIMEOUT, context: { event, timeoutMs }, }) ) }, timeoutMs) } }) } /** * Pipe events from another emitter to this one. * * Subscribes to events on the source emitter and re-emits them on this * emitter. Useful for composing emitters or forwarding events between * layers. * * By default, pipes all standard CDC event types (excluding meta events * `newListener` and `removeListener`). Optionally specify a subset of * events to pipe. * * @param source - The source emitter to pipe events from * @param events - Optional array of event types to pipe (default: all data/connection events) * @returns Unsubscribe function to stop piping * * @example Pipe all events * ```typescript * const unsubscribe = childEmitter.pipe(parentEmitter) * * // Later, stop piping * unsubscribe() * ``` * * @example Pipe specific events * ```typescript * const unsubscribe = logEmitter.pipe(mainEmitter, ['error', 'connected']) * ``` */ pipe(source: CDCEventEmitter, events?: CDCEventType[]): () => void { const eventsToPipe = events ?? [ 'change', 'insert', 'update', 'delete', 'truncate', 'error', 'connected', 'disconnected', 'reconnecting', 'stateChange', 'healthChange', ] const handlers = new Map>() for (const event of eventsToPipe) { const handler = (payload: unknown) => { this.emit(event, payload as CDCEventPayloads[typeof event]) } handlers.set(event, handler) source.on(event, handler as CDCEventListener[typeof event]>) } // Return unsubscribe function return () => { for (const [event, handler] of handlers) { source.off(event, handler as CDCEventListener[typeof event]>) } } } /** * Remove all listeners and clean up resources. * * Call this when the emitter is no longer needed to release memory. * The emitter can still be used after calling `destroy()` - new * listeners can be added. * * @example Cleanup on shutdown * ```typescript * process.on('SIGTERM', () => { * emitter.destroy() * }) * ``` */ destroy(): void { this.listenersMap.clear() } } /** * Factory function to create a new CDC event emitter. * * Convenience function that creates and returns a new `CDCEventEmitter` * instance. Equivalent to `new CDCEventEmitter(config)`. * * @typeParam T - The row type for data events (default: `Row`) * @param config - Optional configuration options * @returns A new `CDCEventEmitter` instance * * @example Create with defaults * ```typescript * const emitter = createEventEmitter() * ``` * * @example Create with config * ```typescript * const emitter = createEventEmitter({ * maxListeners: 50, * catchErrors: true, * }) * ``` */ export function createEventEmitter( config?: EventEmitterConfig ): CDCEventEmitter { return new CDCEventEmitter(config) }