import type { SqlFragment, QueryDef } from '@risingwave/wavelet'; export interface DiffRow { op: 'insert' | 'update_insert' | 'update_delete' | 'delete'; row: Record; rw_timestamp: string; } export interface ViewDiff { cursor: string; inserted: Record[]; updated: Record[]; deleted: Record[]; } export interface BootstrapResult { snapshotRows: Record[]; diffs: ViewDiff[]; lastCursor: string | null; } type DiffCallback = (queryName: string, diff: ViewDiff) => void; /** * Manages persistent subscription cursors against RisingWave. * * Each query gets its own dedicated pg connection and a persistent cursor. * Uses blocking FETCH (WITH timeout) so there is no polling interval - * diffs are dispatched as soon as RisingWave produces them. */ export declare class CursorManager { private connectionString; private queries; private client; private queryConnections; private cursorNames; private subscriptions; private running; constructor(connectionString: string, queries: Record); initialize(): Promise; /** * Start listening for diffs on all queries. * Each query runs its own async loop with blocking FETCH. * No polling interval - FETCH blocks until data arrives or timeout. */ startPolling(callback: DiffCallback): void; stopPolling(): void; private listenLoop; parseDiffs(rows: any[]): ViewDiff; parseDiffBatches(rows: any[]): ViewDiff[]; bootstrap(queryName: string): Promise; query(sql: string): Promise; execute(sql: string): Promise; close(): Promise; private stripSubscriptionMetadata; private normalizeCursor; } export {}; //# sourceMappingURL=cursor-manager.d.ts.map