import { EngineConfig, StreamBoundaryData, StreamChunk, StreamConfig, StreamEventHandler, StreamMetrics, StreamPriority, StreamState } from './types'; export { StreamState, StreamPriority }; /** * Core streaming engine with priority-based scheduling and lifecycle management. * * @description * The StreamingEngine is the heart of the Dynamic HTML Streaming system. * It manages multiple stream boundaries, coordinating their execution * based on priority levels and system resource availability. * * Key features: * - Priority-based scheduling with preemption support * - Stream lifecycle management (start, pause, resume, abort) * - Backpressure handling with configurable strategies * - Automatic retry with exponential backoff * - Comprehensive event system for monitoring * - Thread-safe operations for concurrent access * * @example * ```typescript * const engine = new StreamingEngine({ debug: true }); * * // Register boundaries * engine.registerBoundary('nav', { priority: StreamPriority.Critical }); * engine.registerBoundary('hero', { priority: StreamPriority.High }); * engine.registerBoundary('sidebar', { priority: StreamPriority.Low }); * * // Subscribe to events * const unsubscribe = engine.subscribe((event) => { * console.log(`${event.type}: ${event.boundaryId}`); * }); * * // Start streaming * engine.startAll(); * ``` */ export declare class StreamingEngine { private readonly config; private readonly boundaries; private readonly priorityQueue; private readonly buffer; private readonly eventHandlers; private readonly metrics; private activeStreams; private isProcessing; private processingTimer; private readonly chunkLatencies; private readonly MAX_LATENCY_SAMPLES; /** * Creates a new StreamingEngine instance. * @param config - Engine configuration options */ constructor(config?: Partial); /** * Checks if streaming is supported in the current environment. */ static isStreamingSupported(): boolean; /** * Registers a new stream boundary with the engine. * * @param id - Unique identifier for the boundary * @param config - Stream configuration * @throws {Error} If boundary with same ID already exists * * @example * ```typescript * engine.registerBoundary('main-content', { * priority: StreamPriority.High, * placeholder: , * onStreamComplete: () => console.log('Content ready!'), * }); * ``` */ registerBoundary(id: string, config: StreamConfig): void; /** * Unregisters a stream boundary, aborting any active streams. * * @param id - Boundary identifier to unregister * @returns Whether the boundary was found and unregistered */ unregisterBoundary(id: string): boolean; /** * Gets the current state of a boundary. * @param id - Boundary identifier * @returns Current state or null if not found */ getState(id: string): StreamState | null; /** * Gets the boundary data for inspection. * @param id - Boundary identifier * @returns Boundary data or undefined if not found */ getBoundary(id: string): Readonly | undefined; /** * Returns all registered boundary IDs. */ getBoundaryIds(): string[]; /** * Starts streaming for a specific boundary. * * @param id - Boundary identifier * @throws {Error} If boundary not found or in invalid state * * @example * ```typescript * engine.start('hero'); * ``` */ start(id: string): void; /** * Starts streaming for all registered boundaries. */ startAll(): void; /** * Pauses streaming for a specific boundary. * * @param id - Boundary identifier * @throws {Error} If boundary not found or not streaming */ pause(id: string): void; /** * Resumes a paused stream. * * @param id - Boundary identifier * @throws {Error} If boundary not found or not paused */ resume(id: string): void; /** * Aborts streaming for a specific boundary. * * @param id - Boundary identifier * @param reason - Optional reason for abortion */ abort(id: string, reason?: string): void; /** * Aborts all active streams. * @param reason - Optional reason for abortion */ abortAll(reason?: string): void; /** * Resets a boundary to idle state for restart. * @param id - Boundary identifier */ reset(id: string): void; /** * Processes an incoming chunk for a boundary. * * @param chunk - The chunk to process * @returns Whether the chunk was accepted * * @example * ```typescript * const chunk: StreamChunk = { * id: crypto.randomUUID(), * sequence: 0, * data: '
Content
', * timestamp: Date.now(), * size: 18, * isFinal: false, * boundaryId: 'main', * }; * * engine.processChunk(chunk); * ``` */ processChunk(chunk: StreamChunk): Promise; /** * Delivers buffered chunks to the client. * * @param boundaryId - Boundary to deliver chunks for * @returns Array of delivered chunks */ deliverChunks(boundaryId: string): StreamChunk[]; /** * Subscribes to stream events. * * @param handler - Event handler function * @returns Unsubscribe function * * @example * ```typescript * const unsubscribe = engine.subscribe((event) => { * switch (event.type) { * case StreamEventType.Start: * console.log(`Stream ${event.boundaryId} started`); * break; * case StreamEventType.Complete: * console.log(`Stream ${event.boundaryId} completed`); * break; * } * }); * * // Later... * unsubscribe(); * ``` */ subscribe(handler: StreamEventHandler): () => void; /** * Returns current streaming metrics. */ getMetrics(): Readonly; /** * Returns the current engine configuration. */ getConfig(): Readonly; /** * Disposes of the engine, cleaning up all resources. */ dispose(): void; private queueBoundary; private processQueue; private processEntry; private startStreaming; private completeStream; private handleBackpressure; private handleBufferFull; private waitForBufferDrain; private handleChunkError; private retryChunk; private handleStreamError; private handleTimeout; private normalizeError; private applyTransformer; private transitionState; private canTransitionTo; private isActiveState; private isTerminalState; private createInitialMetrics; private updateMetrics; private recordChunkLatency; private recordBoundaryMetrics; private emitEvent; private getPriorityValue; private normalizePriority; private log; } /** * Creates a new StreamingEngine instance with optional configuration. * * @param config - Partial engine configuration * @returns New StreamingEngine instance * * @example * ```typescript * const engine = createStreamingEngine({ * maxConcurrentStreams: 4, * debug: true, * }); * ``` */ export declare function createStreamingEngine(config?: Partial): StreamingEngine; /** * Creates a chunk with proper defaults. * * @param data - Chunk data * @param boundaryId - Parent boundary ID * @param options - Additional chunk options * @returns Properly formatted StreamChunk */ export declare function createChunk(data: string | Uint8Array, boundaryId: string, options?: Partial>): StreamChunk; /** * Calculates a simple checksum for chunk integrity verification. * * @param data - Data to checksum * @returns Hexadecimal checksum string */ export declare function calculateChecksum(data: string | Uint8Array): string; /** * Creates a ReadableStream that emits chunks from the engine. * * @param engine - StreamingEngine instance * @param boundaryId - Boundary to stream from * @returns ReadableStream of StreamChunks * * @example * ```typescript * const stream = createReadableStream(engine, 'main-content'); * * const reader = stream.getReader(); * while (true) { * const { done, value } = await reader.read(); * if (done) break; * console.log('Received chunk:', value); * } * ``` */ export declare function createReadableStream(engine: StreamingEngine, boundaryId: string): ReadableStream; /** * Creates a TransformStream for processing chunks. * * @param transformer - Chunk transformer function * @returns TransformStream for chunk processing */ export declare function createChunkTransformStream(transformer: (chunk: StreamChunk) => StreamChunk | Promise): TransformStream;