/** * WebRTCManager * * Main entry point for WebRTC functionality. Provides: * - Simple API for DataChannel and MediaStream connections * - Event system with opt-in verbose mode * - Connection management and lifecycle handling * - Automatic reconnection (transparent to developers) */ import { DataChannelHandler } from './data-channel-handler'; import { MediaStreamHandler } from './media-stream-handler'; import { ResponderFanout } from './responder-fanout'; import { ensureWebRTCGlobals, isWebRTCAvailable } from './webrtc-globals'; import { IceServersProvider, TwinTransport, WebRTCManagerOptions, WebRTCEvent, WebRTCEventData, WebRTCEventCallback, PhygridDataChannel, PhygridMediaStream, MediaStreamOptions, MediaStreamResponderOptions, DataChannelResponderOptions, MediaStreamCallback, DataChannelCallback, NoMediaResponderError, } from './types'; /** * Resolved (defaulted) options, minus iceServersProvider — the provider is a * function with no meaningful default, so it is held separately on the manager * rather than forced into this fully-required shape. */ type ResolvedWebRTCOptions = Required>; export const DEFAULT_WEBRTC_OPTIONS: ResolvedWebRTCOptions = { verbose: false, useStun: true, stunServers: ['stun:stun.l.google.com:19302', 'stun:stun1.l.google.com:19302'], turnServers: (() => { const url = process.env.PHYSTACK_TURN_URL; if (!url) return []; const username = process.env.PHYSTACK_TURN_USERNAME || ''; const credential = process.env.PHYSTACK_TURN_CREDENTIAL || ''; return [{ urls: url, username, credential }]; })(), // Always 'all': TURN is a relay FALLBACK, not the only transport. Forcing // 'relay' (the previous behaviour when a TURN URL was set) disabled host/srflx // candidates and pushed all same-network traffic through the relay. Callers // that genuinely want relay-only (e.g. a cross-NAT verification harness) can // still pass iceTransportPolicy: 'relay' explicitly. iceTransportPolicy: 'all', connectionTimeout: 15000, initialRetryDelay: 1000, maxRetryDelay: 30000, }; export class WebRTCManager { private transport: TwinTransport; private options: ResolvedWebRTCOptions; // Optional async source for the full ICE server list (see // WebRTCManagerOptions.iceServersProvider). Held separately from `options` // because it has no default, and threaded into every handler/peer it creates. private iceServersProvider?: IceServersProvider; // Active connections by "{targetTwinId}:{channelName}" private dataChannelHandlers: Map = new Map(); private mediaStreamHandlers: Map = new Map(); // Active responder fan-out dispatchers (one-to-many). Tracked so close() can // tear them down along with all their per-peer sessions. private responderFanouts: Set> = new Set(); // Media fan-outs keyed by "{targetTwinId}:{channelName}", so answerMediaOffer // can locate the armed responder for a one-shot HTTP-signaled offer. A subset // of responderFanouts (which also holds data-channel fan-outs for teardown). private mediaResponderFanouts: Map> = new Map(); // Event listeners private eventListeners: Map>> = new Map(); private isInitialized = false; private isClosed = false; constructor(transport: TwinTransport, options: WebRTCManagerOptions = {}) { this.transport = transport; this.options = { verbose: options.verbose ?? DEFAULT_WEBRTC_OPTIONS.verbose, 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.iceServersProvider = options.iceServersProvider; } // =========================================================================== // DataChannel API // =========================================================================== /** * Create a DataChannel connection to a target twin. * Returns a PhygridDataChannel for sending/receiving messages. * @param targetTwinId - The twin ID to connect to * @param channelName - Optional channel name for multiple channels to same peer (default: 'default') */ async createDataChannel(targetTwinId: string, channelName: string = 'default'): Promise { await this.ensureInitialized(); if (this.isClosed) { throw new Error('WebRTCManager has been closed'); } // Key includes channel name to support multiple channels per peer const handlerKey = `${targetTwinId}:${channelName}`; // Check for existing connection const existing = this.dataChannelHandlers.get(handlerKey); if (existing) { // Return a new channel to the same target return existing.connect(); } const handler = new DataChannelHandler( targetTwinId, true, // isInitiator this.transport, { useStun: this.options.useStun, stunServers: this.options.stunServers, turnServers: this.options.turnServers, iceTransportPolicy: this.options.iceTransportPolicy, iceServersProvider: this.iceServersProvider, connectionTimeout: this.options.connectionTimeout, initialRetryDelay: this.options.initialRetryDelay, maxRetryDelay: this.options.maxRetryDelay, }, channelName, ); this.setupDataChannelHandlerCallbacks(handler, targetTwinId); this.dataChannelHandlers.set(handlerKey, handler); try { return await handler.connect(); } catch (error) { this.dataChannelHandlers.delete(handlerKey); throw error; } } // =========================================================================== // MediaStream API // =========================================================================== /** * Create a MediaStream connection to a target twin. * Returns a PhygridMediaStream for sending/receiving media. * @param targetTwinId - The twin ID to connect to * @param options - Optional MediaStream options * @param channelName - Optional channel name for multiple streams to same peer (default: 'default') */ async createMediaStream( targetTwinId: string, options?: MediaStreamOptions, channelName: string = 'default', ): Promise { await this.ensureInitialized(); if (this.isClosed) { throw new Error('WebRTCManager has been closed'); } // Key includes channel name to support multiple streams per peer const handlerKey = `${targetTwinId}:${channelName}`; // Check for existing connection const existing = this.mediaStreamHandlers.get(handlerKey); if (existing && !existing.isHandlerClosed()) { return existing.connect(); } // A per-stream close() never evicts the handler from this map, so a closed // handler lingers here. Reusing it would throw "MediaStreamHandler has been // closed" — drop the stale entry and build a fresh connection instead (this // is what makes a close-then-reopen on the same twin+channel work). if (existing) { this.mediaStreamHandlers.delete(handlerKey); } const handler = new MediaStreamHandler( targetTwinId, true, // isInitiator this.transport, { useStun: this.options.useStun, stunServers: this.options.stunServers, turnServers: this.options.turnServers, iceTransportPolicy: this.options.iceTransportPolicy, iceServersProvider: this.iceServersProvider, connectionTimeout: this.options.connectionTimeout, initialRetryDelay: this.options.initialRetryDelay, maxRetryDelay: this.options.maxRetryDelay, mediaOptions: options, }, channelName, ); this.setupMediaStreamHandlerCallbacks(handler, targetTwinId); this.mediaStreamHandlers.set(handlerKey, handler); try { return await handler.connect(); } catch (error) { this.mediaStreamHandlers.delete(handlerKey); throw error; } } // =========================================================================== // Responder API (one-to-many fan-out) // =========================================================================== /** * Arm a media responder: one local source fanned out to many simultaneous * peers/viewers, each as its own peer connection. The callback fires ONCE PER * CONNECTED PEER with that peer's own stream and peer id (N=1 is the * single-viewer case). * * Each peer gets a FRESH local stream from options.createLocalStream() — give * each peer connection its own track from the shared source rather than reusing * one track object, and never stop the shared source when one peer leaves. * * @returns a stop() function that tears down the responder and all peers. */ async acceptMediaStream( targetTwinId: string, options: MediaStreamResponderOptions, callback: MediaStreamCallback, ): Promise<() => void> { await this.ensureInitialized(); if (this.isClosed) { throw new Error('WebRTCManager has been closed'); } const channelName = options.channelName ?? 'default'; const direction = options.direction ?? 'sendonly'; const channelPrefix = `media-${channelName}`; // A sending responder with no source factory connects but delivers no media — // a silent black stream. Warn loudly rather than let it look "connected". if ((direction === 'sendonly' || direction === 'sendrecv') && !options.createLocalStream) { console.warn( '[WebRTCManager] acceptMediaStream: direction sends media but no createLocalStream provided — ' + 'peers will connect but receive no track. Pass options.createLocalStream.', ); } const fanout = new ResponderFanout({ transport: this.transport, targetTwinId, channelPrefix, maxPeers: options.maxPeers, createSession: (peerId, replyToTwinId, onClosed, signalingSink) => { const localStream = options.createLocalStream?.(); const handler = new MediaStreamHandler( replyToTwinId, false, // responder this.transport, { useStun: this.options.useStun, stunServers: this.options.stunServers, turnServers: this.options.turnServers, iceTransportPolicy: this.options.iceTransportPolicy, iceServersProvider: this.iceServersProvider, connectionTimeout: this.options.connectionTimeout, initialRetryDelay: this.options.initialRetryDelay, maxRetryDelay: this.options.maxRetryDelay, mediaOptions: { direction, localStream }, peerId, selfManagedSignaling: true, signalingSink, }, channelName, ); handler.setCallbacks({ onConnected: () => this.emit('connected', { targetTwinId, connectionType: 'mediastream' }), onDisconnected: () => { this.emit('disconnected', { targetTwinId, connectionType: 'mediastream' }); onClosed(); }, onError: (error) => this.emit('error', { targetTwinId, error }), }); return handler; }, onPeer: (stream, peerId, info) => { try { callback(stream, peerId, info); } catch (err) { console.error('[WebRTCManager] Error in media stream callback:', err); } }, }); const fanoutKey = `${targetTwinId}:${channelName}`; this.responderFanouts.add(fanout); this.mediaResponderFanouts.set(fanoutKey, fanout); await fanout.start(); return () => { fanout.close(); this.responderFanouts.delete(fanout); // Only drop the map entry if it still points at THIS fan-out — a re-arm on // the same twin+channel may have already replaced it. if (this.mediaResponderFanouts.get(fanoutKey) === fanout) { this.mediaResponderFanouts.delete(fanoutKey); } }; } /** * Answer a SINGLE WebRTC offer delivered out-of-band (from an HTTP request body, * not the socket transport) against an already-armed media responder, and return * the matching non-trickle answer SDP (WHEP). The viewer then connects using the * publisher's relay candidates embedded in the answer — no further signaling. * * The media responder must already be armed for this twin + channel via * acceptMediaStream (PeripheralTwinInstance.onMediaStream); otherwise this throws * NoMediaResponderError. If the responder is at its maxPeers cap, the underlying * fan-out throws PeerCapReachedError. * * @param targetTwinId - the twin whose armed media responder should answer * @param offerSdp - the viewer's offer SDP * @param options.viewerId - id distinguishing this viewer from other peers * @param options.channelName - media channel (default: 'default') * @returns the gathering-complete answer SDP */ async answerMediaOffer( targetTwinId: string, offerSdp: string, options: { viewerId: string; channelName?: string }, ): Promise { await this.ensureInitialized(); if (this.isClosed) { throw new Error('WebRTCManager has been closed'); } const channelName = options.channelName ?? 'default'; const fanoutKey = `${targetTwinId}:${channelName}`; const fanout = this.mediaResponderFanouts.get(fanoutKey); if (!fanout) { throw new NoMediaResponderError(targetTwinId, channelName); } return fanout.ingestExternalOffer(options.viewerId, offerSdp); } /** * Arm a data-channel responder: many simultaneous initiators, each as its own * peer connection and data channel. The callback fires ONCE PER CONNECTED PEER * with that peer's own channel and peer id (N=1 is the single-peer case). * * @returns a stop() function that tears down the responder and all peers. */ async acceptDataChannel( targetTwinId: string, callback: DataChannelCallback, options: DataChannelResponderOptions = {}, ): Promise<() => void> { await this.ensureInitialized(); if (this.isClosed) { throw new Error('WebRTCManager has been closed'); } const channelName = options.channelName ?? 'default'; const channelPrefix = `dc-${channelName}`; const fanout = new ResponderFanout({ transport: this.transport, targetTwinId, channelPrefix, maxPeers: options.maxPeers, createSession: (peerId, replyToTwinId, onClosed) => { const handler = new DataChannelHandler( replyToTwinId, false, // responder this.transport, { useStun: this.options.useStun, stunServers: this.options.stunServers, turnServers: this.options.turnServers, iceTransportPolicy: this.options.iceTransportPolicy, iceServersProvider: this.iceServersProvider, connectionTimeout: this.options.connectionTimeout, initialRetryDelay: this.options.initialRetryDelay, maxRetryDelay: this.options.maxRetryDelay, peerId, selfManagedSignaling: true, }, channelName, ); handler.setCallbacks({ onConnected: () => this.emit('connected', { targetTwinId, connectionType: 'datachannel' }), onDisconnected: () => { this.emit('disconnected', { targetTwinId, connectionType: 'datachannel' }); onClosed(); }, onError: (error) => this.emit('error', { targetTwinId, error }), }); return handler; }, onPeer: (channel, peerId, info) => { try { callback(channel, peerId, info); } catch (err) { console.error('[WebRTCManager] Error in peer data-channel callback:', err); } }, }); this.responderFanouts.add(fanout); await fanout.start(); return () => { fanout.close(); this.responderFanouts.delete(fanout); }; } // =========================================================================== // Event API // =========================================================================== /** * Subscribe to WebRTC events. * * Standard events (always available): * - 'connected': Connection established * - 'disconnected': Connection lost * - 'error': An error occurred * * Verbose events (when options.verbose is true): * - 'reconnecting': Attempting to reconnect * - 'reconnected': Successfully reconnected * - 'ice-state-change': ICE connection state changed * - 'signaling-state-change': Signaling state changed * - 'ice-candidate': ICE candidate received */ on(event: E, callback: WebRTCEventCallback): void { if (!this.eventListeners.has(event)) { this.eventListeners.set(event, new Set()); } this.eventListeners.get(event)!.add(callback); } /** * Unsubscribe from WebRTC events. */ off(event: E, callback: WebRTCEventCallback): void { const listeners = this.eventListeners.get(event); if (listeners) { listeners.delete(callback); } } // =========================================================================== // Lifecycle // =========================================================================== /** * Check if WebRTC is available in the current environment. */ static isAvailable(): boolean { return isWebRTCAvailable(); } /** * Get active DataChannel connection count. */ getActiveDataChannelCount(): number { return this.dataChannelHandlers.size; } /** * Get active MediaStream connection count. */ getActiveMediaStreamCount(): number { return this.mediaStreamHandlers.size; } /** * Close all connections and clean up resources. */ close(): void { if (this.isClosed) return; this.isClosed = true; // Close all DataChannel handlers this.dataChannelHandlers.forEach((handler) => { try { handler.close(); } catch (err) { console.error('[WebRTCManager] Error closing DataChannel handler:', err); } }); this.dataChannelHandlers.clear(); // Close all MediaStream handlers this.mediaStreamHandlers.forEach((handler) => { try { handler.close(); } catch (err) { console.error('[WebRTCManager] Error closing MediaStream handler:', err); } }); this.mediaStreamHandlers.clear(); // Close all responder fan-out dispatchers (and their per-peer sessions) this.responderFanouts.forEach((fanout) => { try { fanout.close(); } catch (err) { console.error('[WebRTCManager] Error closing responder fan-out:', err); } }); this.responderFanouts.clear(); this.mediaResponderFanouts.clear(); // Clear listeners this.eventListeners.clear(); } // =========================================================================== // Private Methods // =========================================================================== private async ensureInitialized(): Promise { if (this.isInitialized) return; const available = await ensureWebRTCGlobals(); if (!available) { console.warn('WebRTC native module not available — WebRTC features disabled'); return; } this.isInitialized = true; } private setupDataChannelHandlerCallbacks(handler: DataChannelHandler, targetTwinId: string): void { handler.setCallbacks({ onConnected: () => { this.emit('connected', { targetTwinId, connectionType: 'datachannel' }); }, onDisconnected: () => { this.emit('disconnected', { targetTwinId, connectionType: 'datachannel' }); }, onError: (error) => { this.emit('error', { targetTwinId, error }); }, onReconnecting: (attempt) => { if (this.options.verbose) { this.emit('reconnecting', { targetTwinId, attempt }); } }, onReconnected: (attempt) => { if (this.options.verbose) { this.emit('reconnected', { targetTwinId, attempt }); } }, }); } private setupMediaStreamHandlerCallbacks(handler: MediaStreamHandler, targetTwinId: string): void { handler.setCallbacks({ onConnected: () => { this.emit('connected', { targetTwinId, connectionType: 'mediastream' }); }, onDisconnected: () => { this.emit('disconnected', { targetTwinId, connectionType: 'mediastream' }); }, onError: (error) => { this.emit('error', { targetTwinId, error }); }, onReconnecting: (attempt) => { if (this.options.verbose) { this.emit('reconnecting', { targetTwinId, attempt }); } }, onReconnected: (attempt) => { if (this.options.verbose) { this.emit('reconnected', { targetTwinId, attempt }); } }, }); } private emit(event: E, data: WebRTCEventData[E]): void { const listeners = this.eventListeners.get(event); if (!listeners) return; listeners.forEach((callback) => { try { callback(data); } catch (err) { console.error(`[WebRTCManager] Error in ${event} event listener:`, err); } }); } } // Re-export types for convenience export type { TwinTransport, WebRTCManagerOptions, IceServersProvider, IceServersResult, WebRTCEvent, PhygridDataChannel, PhygridMediaStream, MediaStreamOptions, } from './types';