/** * CDC Client Module for postgres.do * * This module handles: * - CDC transport creation and management * - Subscription lifecycle (subscribe/unsubscribe) * - Event handling * - Transport selection (WebSocket vs SSE) */ import type { Row, PostgresConfig, CDCTransport, CDCClientSubscription, CDCClientSubscribeOptions, CDCClientSubscriptionInfo, } from './types' import { createCdcWsTransport } from './transport/cdc-ws' import { createCdcSseTransport } from './transport/cdc-sse' import { CDCSubscription, validateTableName, generateSubscriptionId, type SubscribeOptions, } from './subscription' // ============================================================================ // CDC State // ============================================================================ /** * State container for CDC functionality. * This is created per-client to manage CDC transport and subscriptions. */ export interface CDCState { /** The CDC transport (WebSocket or SSE) */ transport: CDCTransport | null /** Map of subscription ID to subscription instance */ subscriptions: Map } /** * Create a new CDC state container. */ export function createCDCState(): CDCState { return { transport: null, subscriptions: new Map(), } } // ============================================================================ // CDC Transport Management // ============================================================================ /** * Get or create the CDC transport. * * This lazily creates the CDC transport on first use. The transport type * (WebSocket or SSE) is determined by the client configuration. * * @param state - The CDC state container * @param baseUrl - The base URL for the postgres.do API * @param config - The client configuration * @returns The CDC transport */ export async function getOrCreateCdcTransport( state: CDCState, baseUrl: string, config: PostgresConfig ): Promise { if (state.transport) { return state.transport } // Determine the CDC endpoint URL from baseUrl const cdcUrl = baseUrl.replace('https://', config.transport === 'ws' ? 'wss://' : 'https://') + '/cdc' // Create transport based on config if (config.transport === 'ws') { // Build config object only with defined values const wsConfig: Parameters[0] = { url: cdcUrl, } if (config.apiKey !== undefined) wsConfig.apiKey = config.apiKey if (config.WebSocket !== undefined) wsConfig.WebSocket = config.WebSocket if (config.connectTimeout !== undefined) wsConfig.connectTimeout = config.connectTimeout state.transport = createCdcWsTransport(wsConfig) } else { // HTTP transport uses SSE for CDC // Build config object only with defined values const sseConfig: Parameters[0] = { url: cdcUrl, } if (config.apiKey !== undefined) sseConfig.apiKey = config.apiKey if (config.fetch !== undefined) sseConfig.fetch = config.fetch state.transport = createCdcSseTransport(sseConfig) } // Set up event handlers state.transport.onEvent((subscriptionId: string, event) => { const subscription = state.subscriptions.get(subscriptionId) if (subscription) { subscription.handleEvent(event) } }) state.transport.onError((subscriptionId: string, error) => { const subscription = state.subscriptions.get(subscriptionId) if (subscription) { subscription.handleError(error) } }) state.transport.onClose(() => { // Mark all subscriptions as inactive state.subscriptions.forEach(subscription => { subscription.setInactive() }) state.subscriptions.clear() state.transport = null }) // Connect await state.transport.connect() return state.transport } // ============================================================================ // Subscription Management // ============================================================================ /** * Subscribe to changes on a table. * * Creates a new CDC subscription for the specified table. The subscription * can be used with async iteration or callback handlers. * * @param state - The CDC state container * @param baseUrl - The base URL for the postgres.do API * @param config - The client configuration * @param tableName - Table name (optionally with schema: "schema.table") * @param options - Subscription options including event filters and callbacks * @returns A subscription that can be iterated over or used with callbacks * * @example * ```typescript * // Using async iteration * const sub = await subscribe(state, baseUrl, config, 'users') * for await (const event of sub) { * console.log(event.operation, event.newRow) * } * * // Using callbacks * const sub = await subscribe(state, baseUrl, config, 'users', { * onInsert: (row) => console.log('New user:', row), * onUpdate: (newRow, oldRow) => console.log('Updated:', newRow), * onDelete: (oldRow) => console.log('Deleted:', oldRow) * }) * ``` */ export async function subscribe( state: CDCState, baseUrl: string, config: PostgresConfig, tableName: string, options?: CDCClientSubscribeOptions ): Promise> { // Validate table name to prevent SQL injection const { table, schema } = validateTableName(tableName) // Get or create CDC transport const transport = await getOrCreateCdcTransport(state, baseUrl, config) // Generate subscription ID const subscriptionId = generateSubscriptionId() // Create unsubscribe callback const unsubscribeCallback = async (): Promise => { state.subscriptions.delete(subscriptionId) await transport.unsubscribe(subscriptionId) // If no more subscriptions, disconnect the transport if (state.subscriptions.size === 0 && state.transport) { await state.transport.disconnect() state.transport = null } } // Create subscription instance const subscription = new CDCSubscription( subscriptionId, table, schema, transport, unsubscribeCallback, options as SubscribeOptions ) // Register subscription state.subscriptions.set(subscriptionId, subscription as CDCSubscription) // Subscribe via transport - build options object only with defined values const subscriptionOptions: Parameters[3] = {} if (options?.events !== undefined) subscriptionOptions.events = options.events if (options?.filter !== undefined) subscriptionOptions.filter = options.filter if (options?.resumeFrom !== undefined) subscriptionOptions.resumeFrom = options.resumeFrom if (options?.includeOldRow !== undefined) subscriptionOptions.includeOldRow = options.includeOldRow if (options?.trackChangedColumns !== undefined) subscriptionOptions.trackChangedColumns = options.trackChangedColumns if (options?.batchSize !== undefined) subscriptionOptions.batchSize = options.batchSize if (options?.heartbeatInterval !== undefined) subscriptionOptions.heartbeatInterval = options.heartbeatInterval await transport.subscribe(subscriptionId, table, schema, subscriptionOptions) return subscription } /** * Unsubscribe from a subscription by ID. * * @param state - The CDC state container * @param subscriptionId - The subscription ID to unsubscribe from */ export async function unsubscribe( state: CDCState, subscriptionId: string ): Promise { const subscription = state.subscriptions.get(subscriptionId) if (subscription) { await subscription.unsubscribe() } // If subscription not found, silently succeed (idempotent behavior) } /** * Get all active subscriptions. * * @param state - The CDC state container * @returns Array of subscription info objects */ export function getSubscriptions(state: CDCState): CDCClientSubscriptionInfo[] { return Array.from(state.subscriptions.values()).map(sub => ({ id: sub.id, table: sub.table, schema: sub.schema, isActive: sub.isActive, })) }