/** * MediaStreamHandler * * Manages MediaStream connections including: * - Stream setup for sending/receiving media * - Track lifecycle management (add/remove) * - Frame activity monitoring for stale connection detection * - Persistent stream abstraction across reconnections */ import { PeerConnectionManager, PeerConnectionManagerOptions } from './peer-connection-manager'; import { IceServersProvider, TwinTransport, TurnServerConfig, PhygridMediaStream, MediaStreamOptions, MediaTrackKind, MissingMediaTrackError, PeerConnectionConfig, ExtendedMediaStreamTrack, SignalingSink, } from './types'; import { DEFAULT_WEBRTC_OPTIONS } from './webrtc-manager'; const FRAME_INACTIVITY_TIMEOUT = 5000; // 5 seconds without frames = stale // Liveness backstop for a SENDONLY peer (a responder serving media to a viewer). // Such a peer receives no media, so the frame monitor never runs; and // @roamhq/wrtc does not reliably transition pc.connectionState when the remote // viewer silently vanishes (page reload/crash). We instead watch inbound // transport activity via getStats — a live remote keeps sending STUN consent + // RTCP back, so a prolonged freeze means it's gone. const OUTBOUND_LIVENESS_CHECK_INTERVAL_MS = 5000; const OUTBOUND_STALE_TIMEOUT_MS = 20000; // inbound silent this long ⇒ remote gone export interface MediaStreamHandlerOptions 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; /** Media options */ mediaOptions?: MediaStreamOptions; /** 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; /** * When set, this peer's outbound answer/ice is handed to the sink instead of * published to the transport (one-shot HTTP-signaled answer). See * PeerConnectionConfig.signalingSink. */ signalingSink?: SignalingSink; } export class MediaStreamHandler { private targetTwinId: string; private channelName: string; private isInitiator: boolean; private transport: TwinTransport; private options: Required< Omit< MediaStreamHandlerOptions, 'mediaOptions' | 'peerId' | 'selfManagedSignaling' | 'iceServersProvider' | 'signalingSink' > > & { mediaOptions: MediaStreamOptions; }; private peerId?: string; private selfManagedSignaling: boolean; private iceServersProvider?: IceServersProvider; private signalingSink?: SignalingSink; private pcManager: PeerConnectionManager | null = null; private remoteStream: MediaStream | null = null; private localStream: MediaStream | null = null; private isClosed = false; private isReceiving = false; private isConnecting = false; // Track management private tracks: MediaStreamTrack[] = []; private trackLastFrameTime: Map = new Map(); private frameActivityInterval: NodeJS.Timeout | null = null; // Outbound (sendonly) liveness backstop state private outboundLivenessInterval: NodeJS.Timeout | null = null; private lastInboundTotal = 0; private lastInboundProgressAt = 0; private everReceivedInbound = false; // Event listeners private trackListeners: Set<(track: MediaStreamTrack) => void> = new Set(); private frameListeners: Set<(frameData: 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; // Guard against duplicate connect() calls (e.g. React StrictMode double-effect) private connectPromise: Promise | null = null; constructor( targetTwinId: string, isInitiator: boolean, transport: TwinTransport, options: MediaStreamHandlerOptions = {}, 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, mediaOptions: { direction: options.mediaOptions?.direction ?? (isInitiator ? 'recvonly' : 'sendrecv'), localStream: options.mediaOptions?.localStream, kinds: options.mediaOptions?.kinds, }, }; this.peerId = options.peerId; this.selfManagedSignaling = options.selfManagedSignaling ?? false; this.iceServersProvider = options.iceServersProvider; this.signalingSink = options.signalingSink; // Create local stream for non-initiator if not provided if (!isInitiator && !this.options.mediaOptions.localStream) { this.localStream = new MediaStream(); this.options.mediaOptions.localStream = this.localStream; } else { this.localStream = this.options.mediaOptions.localStream ?? null; } } /** * 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; } /** Whether close() has been called. A closed handler cannot reconnect. */ isHandlerClosed(): boolean { return this.isClosed; } /** * Connect and return a PhygridMediaStream. * Idempotent: if already connecting, returns the same promise. */ async connect(): Promise { if (this.isClosed) { throw new Error('MediaStreamHandler has been closed'); } // Return existing connect promise to prevent duplicate PeerConnectionManagers if (this.connectPromise) { return this.connectPromise; } this.connectPromise = this.doConnect(); return this.connectPromise; } private async doConnect(): Promise { this.isConnecting = true; // Channel prefix includes name for unique signaling per channel const channelPrefix = `media-${this.channelName}`; const pcConfig: PeerConnectionConfig = { targetTwinId: this.targetTwinId, isInitiator: this.isInitiator, connectionType: 'mediastream', 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, signalingSink: this.signalingSink, onConnected: () => this.handlePeerConnected(), onDisconnected: () => this.handlePeerDisconnected(), onError: (err) => this.handleError(err), onReconnecting: (attempt) => { this.isConnecting = true; this.isReceiving = false; // pc is being rebuilt — pause the backstop; handlePeerConnected re-arms it. this.stopOutboundLivenessMonitor(); this.onReconnectingCallback?.(attempt); }, onReconnected: (attempt) => this.handleReconnected(attempt), // Setup transceivers/tracks BEFORE offer is created (critical for ICE gathering) onPeerConnectionCreated: (pc) => { if (this.isInitiator) { if (this.localStream && this.localStream.getTracks().length > 0) { // Initiator with local stream: Add actual tracks before offer this.localStream.getTracks().forEach((track) => { pc.addTrack(track, this.localStream!); }); } else { // Initiator without local stream: add one receiving transceiver per // requested kind. Default ['video'] preserves the legacy behaviour // for every caller that predates MediaStreamOptions.kinds. const kinds = this.options.mediaOptions.kinds ?? ['video']; for (const kind of kinds) { pc.addTransceiver(kind, { direction: this.options.mediaOptions.direction }); } } } else if (this.localStream) { // Responder: Add local tracks before answer is sent this.localStream.getTracks().forEach((track) => { pc.addTrack(track, this.localStream!); }); } // Handle incoming tracks (both sides) pc.ontrack = (event) => { // Keep a SINGLE stable remoteStream object across reconnects: an app that // captured getStream() as a