/** * DataChannelHandler * * Manages RTCDataChannel lifecycle including: * - Channel creation (initiator) or reception (responder) * - Message serialization/deserialization * - Persistent channel abstraction across reconnections * - Message buffering during reconnection */ import { PeerConnectionManager, PeerConnectionManagerOptions } from './peer-connection-manager'; import { IceServersProvider, TwinTransport, TurnServerConfig, PhygridDataChannel, PeerConnectionConfig } from './types'; import { DEFAULT_WEBRTC_OPTIONS } from './webrtc-manager'; // Liveness backstop: data channels carry no frame heartbeat (unlike media), and a // reloaded peer can leave the responder pc reporting a non-'connected' state // without ever firing onconnectionstatechange. We poll pc/dc state and force a // reconnect when it stays unhealthy past the deadline. This is pc-state-based, NOT // traffic-based: an idle data channel is perfectly normal and must never be torn down. const LIVENESS_CHECK_INTERVAL_MS = 5000; // Well above connectionTimeout (15s) so it never races the normal connect/reconnect. const PC_STALE_TIMEOUT_MS = 20000; export interface DataChannelHandlerOptions extends Partial { /** Use STUN servers. Default: true */ useStun?: boolean; /** STUN server URLs */ stunServers?: string[]; /** TURN server configs for relay fallback */ turnServers?: TurnServerConfig[]; /** ICE transport policy. Default: 'all' */ iceTransportPolicy?: RTCIceTransportPolicy; /** Optional async source for the full ICE server list. See WebRTCManagerOptions.iceServersProvider. */ iceServersProvider?: IceServersProvider; /** Per-session peer id (responder fan-out: the peer this handler serves). */ peerId?: string; /** * When true, signaling is owned by an external dispatcher (responder fan-out) * which feeds messages via ingestSignalingMessage(). See PeerConnectionConfig. */ selfManagedSignaling?: boolean; } export class DataChannelHandler { private targetTwinId: string; private channelName: string; private isInitiator: boolean; private transport: TwinTransport; private options: Required>; private peerId?: string; private selfManagedSignaling: boolean; private iceServersProvider?: IceServersProvider; private pcManager: PeerConnectionManager | null = null; private dataChannel: RTCDataChannel | null = null; private isOpen = false; private isClosed = false; // Connection status private isConnecting = false; // Liveness backstop private livenessInterval: NodeJS.Timeout | null = null; private pcUnhealthySince: number | null = null; // Event listeners private messageListeners: Set<(data: any) => void> = new Set(); private closeListeners: Set<() => void> = new Set(); // Callbacks for WebRTCManager private onConnectedCallback?: () => void; private onDisconnectedCallback?: () => void; private onErrorCallback?: (error: Error) => void; private onReconnectingCallback?: (attempt: number) => void; private onReconnectedCallback?: (attempt: number) => void; constructor( targetTwinId: string, isInitiator: boolean, transport: TwinTransport, options: DataChannelHandlerOptions = {}, channelName: string = 'default', ) { this.targetTwinId = targetTwinId; this.channelName = channelName; this.isInitiator = isInitiator; this.transport = transport; this.options = { useStun: options.useStun ?? DEFAULT_WEBRTC_OPTIONS.useStun, stunServers: options.stunServers ?? DEFAULT_WEBRTC_OPTIONS.stunServers, turnServers: options.turnServers ?? DEFAULT_WEBRTC_OPTIONS.turnServers, iceTransportPolicy: options.iceTransportPolicy ?? DEFAULT_WEBRTC_OPTIONS.iceTransportPolicy, connectionTimeout: options.connectionTimeout ?? DEFAULT_WEBRTC_OPTIONS.connectionTimeout, initialRetryDelay: options.initialRetryDelay ?? DEFAULT_WEBRTC_OPTIONS.initialRetryDelay, maxRetryDelay: options.maxRetryDelay ?? DEFAULT_WEBRTC_OPTIONS.maxRetryDelay, }; this.peerId = options.peerId; this.selfManagedSignaling = options.selfManagedSignaling ?? false; this.iceServersProvider = options.iceServersProvider; } /** * Set callbacks for connection events */ setCallbacks(callbacks: { onConnected?: () => void; onDisconnected?: () => void; onError?: (error: Error) => void; onReconnecting?: (attempt: number) => void; onReconnected?: (attempt: number) => void; }): void { this.onConnectedCallback = callbacks.onConnected; this.onDisconnectedCallback = callbacks.onDisconnected; this.onErrorCallback = callbacks.onError; this.onReconnectingCallback = callbacks.onReconnecting; this.onReconnectedCallback = callbacks.onReconnected; } /** * Connect and return a PhygridDataChannel */ async connect(): Promise { if (this.isClosed) { throw new Error('DataChannelHandler has been closed'); } this.isConnecting = true; // Channel ID includes the name to support multiple channels to same peer const channelId = `dc-${this.channelName}-${this.targetTwinId}`; // Prefix for signaling messages - unique per channel name const channelPrefix = `dc-${this.channelName}`; const pcConfig: PeerConnectionConfig = { targetTwinId: this.targetTwinId, isInitiator: this.isInitiator, connectionType: 'datachannel', channelPrefix, useStun: this.options.useStun, stunServers: this.options.stunServers, turnServers: this.options.turnServers, iceTransportPolicy: this.options.iceTransportPolicy, iceServersProvider: this.iceServersProvider, peerId: this.peerId, selfManagedSignaling: this.selfManagedSignaling, onConnected: () => this.handlePeerConnected(), onDisconnected: () => this.handlePeerDisconnected(), onError: (err) => this.handleError(err), onReconnecting: (attempt) => { this.isConnecting = true; this.isOpen = false; // Clear old data channel reference - new one will be created/received this.dataChannel = null; // Pause the backstop while a reconnect is already in flight. this.stopLivenessMonitor(); this.onReconnectingCallback?.(attempt); }, onReconnected: (attempt) => this.handleReconnected(attempt), // Setup data channel/handler before signaling begins onPeerConnectionCreated: (pc) => { if (this.isInitiator) { // Initiator: Create data channel before offer is sent console.log('[DataChannelHandler] Creating data channel before offer'); this.dataChannel = pc.createDataChannel(channelId); this.dataChannel.binaryType = 'arraybuffer'; } else { // Responder: Set up ondatachannel handler BEFORE receiving offer // This ensures we don't miss the datachannel event console.log('[DataChannelHandler] Setting up ondatachannel handler'); pc.ondatachannel = (event) => { console.log('[DataChannelHandler] Received data channel from initiator'); this.dataChannel = event.channel; this.dataChannel.binaryType = 'arraybuffer'; this.attachDataChannelHandlers(this.dataChannel, () => { // No-op: Channel setup is complete. Handler registration and callbacks // are managed inside attachDataChannelHandlers. }); }; } }, }; this.pcManager = new PeerConnectionManager(pcConfig, this.transport, { connectionTimeout: this.options.connectionTimeout, initialRetryDelay: this.options.initialRetryDelay, maxRetryDelay: this.options.maxRetryDelay, }); const pc = await this.pcManager.connect(); await this.setupDataChannel(pc); return this.createPhygridDataChannel(); } /** * Feed a demultiplexed signaling message in from the responder fan-out * dispatcher (used with selfManagedSignaling). The underlying peer connection * manager buffers messages that arrive before its pc exists. */ async ingestSignalingMessage(message: any): Promise { await this.pcManager?.ingestSignalingMessage(message); } /** * Close the data channel and clean up */ close(): void { if (this.isClosed) return; this.isClosed = true; this.isOpen = false; this.stopLivenessMonitor(); if (this.dataChannel) { try { this.dataChannel.close(); } catch (err) { console.error('[DataChannelHandler] Error closing data channel:', err); } this.dataChannel = null; } if (this.pcManager) { this.pcManager.close(); this.pcManager = null; } this.notifyCloseListeners(); this.messageListeners.clear(); this.closeListeners.clear(); } // =========================================================================== // Private Methods // =========================================================================== private async setupDataChannel(pc: RTCPeerConnection): Promise { return new Promise((resolve) => { // Data channel was already created/received in onPeerConnectionCreated // Just need to attach handlers if not already done if (this.dataChannel) { this.attachDataChannelHandlers(this.dataChannel, resolve); } else if (this.isInitiator) { // Fallback: create data channel now (shouldn't happen normally) const channelId = `channel-${this.targetTwinId}`; this.dataChannel = pc.createDataChannel(channelId); this.dataChannel.binaryType = 'arraybuffer'; this.attachDataChannelHandlers(this.dataChannel, resolve); } else { // Responder: ondatachannel handler should already be set, // but it may not have fired yet - we'll resolve when it does const existingHandler = pc.ondatachannel; pc.ondatachannel = (event) => { existingHandler?.call(pc, event); resolve(); }; } }); } private attachDataChannelHandlers(dc: RTCDataChannel, onReady: () => void): void { dc.onopen = () => { this.isOpen = true; this.isConnecting = false; this.startLivenessMonitor(); this.onConnectedCallback?.(); onReady(); }; dc.onmessage = (event) => { this.handleMessage(event.data); }; dc.onclose = () => { this.isOpen = false; if (!this.isClosed) { this.onDisconnectedCallback?.(); } }; dc.onerror = (event) => { const error = new Error(`DataChannel error: ${(event as any).error?.message || 'Unknown error'}`); this.handleError(error); }; // If already open (can happen on reconnect), resolve immediately if (dc.readyState === 'open') { this.isOpen = true; this.isConnecting = false; this.startLivenessMonitor(); onReady(); } } private startLivenessMonitor(): void { if (this.livenessInterval || this.isClosed) return; this.pcUnhealthySince = null; this.livenessInterval = setInterval(() => this.checkLiveness(), LIVENESS_CHECK_INTERVAL_MS); } private stopLivenessMonitor(): void { if (this.livenessInterval) { clearInterval(this.livenessInterval); this.livenessInterval = null; } this.pcUnhealthySince = null; } private checkLiveness(): void { if (this.isClosed) { this.stopLivenessMonitor(); return; } const pc = this.pcManager?.getPeerConnection(); const dc = this.dataChannel; const pcHealthy = pc?.connectionState === 'connected'; const dcDead = !dc || dc.readyState === 'closing' || dc.readyState === 'closed'; if (!pcHealthy || dcDead) { // pc reports something other than 'connected', or the channel transport is // gone — start (or continue) the staleness clock. NOT based on traffic. this.pcUnhealthySince = this.pcUnhealthySince ?? Date.now(); if (Date.now() - this.pcUnhealthySince >= PC_STALE_TIMEOUT_MS) { this.pcUnhealthySince = null; this.stopLivenessMonitor(); if (this.selfManagedSignaling) { // Fan-out responder: it can't re-offer and has no subscription, so an // in-place rebuild would spin. Signal disconnect → the dispatcher drops // this peer and the initiator re-offers a fresh session. console.log('[DataChannelHandler] Liveness backstop: pc/dc stale, tearing down (fan-out peer)'); this.handlePeerDisconnected(); } else { console.log('[DataChannelHandler] Liveness backstop: pc/dc stale, forcing reconnect'); // Guarded against loops by isReconnecting/isShuttingDown inside the manager. this.pcManager?.reconnect(); } } } else { // Healthy — reset the clock. this.pcUnhealthySince = null; } } private handleMessage(data: string | ArrayBuffer): void { let parsedData = data; // Try to parse JSON strings if (typeof data === 'string') { try { parsedData = JSON.parse(data); } catch { // Not JSON, use as-is } } this.messageListeners.forEach((listener) => { try { listener(parsedData); } catch (err) { console.error('[DataChannelHandler] Error in message listener:', err); } }); } private handlePeerConnected(): void { // No-op: Peer connection is established but we don't need to do anything here. // DataChannel setup is handled by onPeerConnectionCreated callback. } private handlePeerDisconnected(): void { this.isOpen = false; this.onDisconnectedCallback?.(); } private handleReconnected(attempt: number): void { // Re-setup data channel after reconnection const pc = this.pcManager?.getPeerConnection(); if (pc) { this.setupDataChannel(pc).catch((err) => { console.error('[DataChannelHandler] Error re-setting up data channel:', err); }); } // Re-arm the liveness backstop after every reconnect (onReconnecting stopped // it). startLivenessMonitor is idempotent, so this is safe even if the // re-opened channel's onopen path already started it. this.startLivenessMonitor(); this.onReconnectedCallback?.(attempt); } private handleError(error: Error): void { console.error('[DataChannelHandler] DataChannel error:', error); this.onErrorCallback?.(error); } private notifyCloseListeners(): void { this.closeListeners.forEach((listener) => { try { listener(); } catch (err) { console.error('[DataChannelHandler] Error in close listener:', err); } }); } private createPhygridDataChannel(): PhygridDataChannel { return { send: (data: string | ArrayBuffer | object) => { let messageToSend: string | ArrayBuffer; if (typeof data === 'object' && !(data instanceof ArrayBuffer)) { messageToSend = JSON.stringify(data); } else { messageToSend = data as string | ArrayBuffer; } if (this.dataChannel?.readyState === 'open') { try { this.dataChannel.send(messageToSend as any); } catch (err) { console.error('[DataChannelHandler] Error sending message:', err); } } // If not open, silently drop the message - no buffering }, onMessage: (callback: (data: any) => void) => { this.messageListeners.add(callback); }, offMessage: (callback: (data: any) => void) => { this.messageListeners.delete(callback); }, onClose: (callback: () => void) => { this.closeListeners.add(callback); }, offClose: (callback: () => void) => { this.closeListeners.delete(callback); }, close: () => { this.close(); }, isOpen: () => { return this.isOpen && this.dataChannel?.readyState === 'open'; }, isConnecting: () => { return this.isConnecting; }, getTargetTwinId: () => { return this.targetTwinId; }, getChannelName: () => { return this.channelName; }, getPeerId: () => { return this.pcManager?.getPeerId() ?? this.peerId ?? null; }, }; } }