/** * Stream Broadcaster * * Core infrastructure for streaming to clients. * Orchestrates auth, transport, and connection management. * * Note: Event subscriptions are handled by ChannelControllers, * not directly by StreamBroadcaster. */ import type { BaseAuthAdapter } from './auth/BaseAuthAdapter'; import type { BaseTransportAdapter } from './transports/BaseTransportAdapter'; import type { StreamAuthResult, StreamBroadcasterConfig, StreamChannel, StreamManagerStats, StreamMessage, StreamTransportType } from '@plyaz/types/core'; /** * StreamBroadcaster - Core streaming infrastructure * * Features: * - Manages connections via StreamConnectionManager * - Routes messages through transport adapters (SSE, WebSocket) * - Handles auth via adapter * - Automatic stale connection cleanup * - Channel-based message broadcasting * * Note: Event subscriptions are handled by endpoint controllers * registered via StreamRegistry (e.g., FilesChannelController). * * @example * ```ts * const broadcaster = new StreamBroadcaster( * new AnonymousAuthAdapter(), * [new SSETransportAdapter()] * ); * * broadcaster.initialize(); * * // Register with StreamRegistry * await StreamRegistry.initialize({ * broadcaster, * endpoints: [ * { endpoint: FilesStreamEndpoint, config: { enabled: true } }, * ], * }); * * // In API route * const authResult = await broadcaster.authenticate(request); * const response = await broadcaster.createConnection('conn-1', 'sse', request); * broadcaster.subscribe('conn-1', 'upload:abc123', authResult); * return response; * ``` */ export declare class StreamBroadcaster { private readonly connectionManager; private readonly authAdapter; private readonly transports; private readonly maxStaleAge; private unsubscribers; private cleanupTimer?; private initialized; constructor(authAdapter: BaseAuthAdapter, transports: BaseTransportAdapter[], config?: StreamBroadcasterConfig); /** * Initialize cleanup timer * Call once on startup * * Note: Event subscriptions are now handled by endpoint controllers * registered via StreamRegistry (e.g., FilesChannelController, SystemChannelController). */ initialize(): void; /** * Authenticate a connection request * * @param request - HTTP request * @returns Authentication result */ authenticate(request: Request): Promise; /** * Check if user can subscribe to a channel * * @param authResult - Authentication result * @param channel - Channel to subscribe */ canSubscribe(authResult: StreamAuthResult, channel: StreamChannel): boolean; /** * Create a new streaming connection * * @param connectionId - Unique connection ID * @param transportType - Transport to use (sse or websocket) * @param request - HTTP request * @param authResult - Optional auth result (if already authenticated) * @returns Response (for SSE) or void (for WebSocket) */ createConnection(connectionId: string, transportType: StreamTransportType, request: Request, authResult?: StreamAuthResult): Promise; /** * Subscribe connection to a channel * * @param connectionId - Connection ID * @param channel - Channel to subscribe * @param authResult - Auth result for permission check * @returns Whether subscription succeeded */ subscribe(connectionId: string, channel: StreamChannel, authResult?: StreamAuthResult): boolean; /** * Unsubscribe connection from a channel * * @param connectionId - Connection ID * @param channel - Channel to unsubscribe */ unsubscribe(connectionId: string, channel: StreamChannel): void; /** * Close a connection * * @param connectionId - Connection to close */ closeConnection(connectionId: string): Promise; /** * Broadcast message to all connections on a channel * * @param channel - Channel to broadcast to * @param message - Message to send */ broadcastToChannel(channel: StreamChannel, message: StreamMessage): void; /** * Send message to a specific connection * * @param connectionId - Connection ID * @param message - Message to send */ sendToConnection(connectionId: string, message: StreamMessage): Promise; /** * Broadcast to all connections (system-wide) * * @param message - Message to send */ broadcastToAll(message: StreamMessage): void; /** * Get manager statistics */ getStats(): StreamManagerStats; /** * Dispose broadcaster (for shutdown) */ dispose(): void; /** * Clean up stale connections */ private cleanupStaleConnections; } //# sourceMappingURL=StreamBroadcaster.d.ts.map