/** * Change Streams Implementation * * Real-time change notifications for document collections * Compatible with MongoDB change streams API */ import { Document, ChangeEvent, ChangeStream, ChangeStreamOptions, AggregationPipeline, IChangeStreamManager } from './interface.js'; interface StreamConnection { id: string; collection: string; pipeline?: AggregationPipeline; options: ChangeStreamOptions; listeners: { change?: (event: ChangeEvent) => void; error?: (error: Error) => void; close?: () => void; }; isActive: boolean; createdAt: Date; lastEventId?: string; } /** * Change Stream implementation using Server-Sent Events */ export declare class DocumentChangeStream implements ChangeStream { private connection; private changeStreamManager; private closed; constructor(connection: StreamConnection, manager: ChangeStreamManager); on(event: 'change', listener: (change: ChangeEvent) => void): this; on(event: 'error', listener: (error: Error) => void): this; on(event: 'close', listener: () => void): this; close(): Promise; isClosed(): boolean; _emitChange(event: ChangeEvent): void; _emitError(error: Error): void; } /** * Change Stream Manager * Manages multiple change streams and event distribution */ export declare class ChangeStreamManager implements IChangeStreamManager { private streams; private collectionStreams; private eventQueue; private durableObject?; private connections; constructor(options?: { durableObject?: any; maxEventHistory?: number; eventTTL?: number; }); /** * Create a new change stream */ createStream(collection: string, pipeline?: AggregationPipeline, options?: ChangeStreamOptions): Promise>; /** * Close a change stream */ closeStream(streamId: string): Promise; /** * Emit change event to all relevant streams */ emitChange(collection: string, change: ChangeEvent): Promise; /** * List active streams */ listActiveStreams(): Promise; /** * Get stream statistics */ getStreamStats(): Promise<{ totalStreams: number; activeStreams: number; streamsByCollection: Record; }>; /** * Handle WebSocket connection for real-time streaming */ handleWebSocketConnection(request: Request, streamId: string): Promise; /** * Handle Server-Sent Events connection */ handleSSEConnection(request: Request, streamId: string): Promise; private generateStreamId; private generateEventId; private storeEvent; private sendMissedEvents; private matchesPipeline; private matchesFilter; private matchesDocumentFilter; private processChangeEvent; private handleWebSocketMessage; private cleanupInactiveStreams; /** * Get connection for stream (for external use) */ getStreamConnection(streamId: string): any; /** * Send event to specific stream connection */ sendToConnection(streamId: string, event: ChangeEvent): Promise; } export {}; //# sourceMappingURL=change-streams.d.ts.map