import { describe, expect, it } from 'bun:test'; import { WebRTCManager } from '../webrtc-manager'; import { PeerConnectionManager } from '../peer-connection-manager'; import { ResponderFanout } from '../responder-fanout'; import { MediaStreamHandler } from '../media-stream-handler'; import { NoMediaResponderError, PeerCapReachedError, type PeerConnectionConfig, type PhygridMediaStream, type TwinTransport, } from '../types'; /** * Unit tests for the one-shot (HTTP-signaled) WebRTC answer path — WHEP. * * A viewer delivers a single offer out-of-band (an HTTP request body, no socket) * and gets back one non-trickle answer SDP. The offer is fed into an ALREADY-ARMED * media responder fan-out (ResponderFanout.ingestExternalOffer), answered on a * fresh per-peer session whose outbound signaling is diverted to a per-call * signaling sink, and the answer is emitted only after ICE gathering completes (or * a cap) so it carries the publisher's relay candidates embedded. * * We drive a fake RTCPeerConnection that answers an offer and, in 'complete' mode, * emits one trickle candidate then folds it into the local SDP and fires * icegatheringstatechange. In 'stall' mode it never completes gathering, exercising * the cap fallback. Dependency-injected transports/configs throughout (no module * mocks), matching the sibling ice-servers-provider / sendonly-liveness tests. */ const OFFER_SDP = 'v=0\r\no=- 111 2 IN IP4 0.0.0.0\r\ns=-\r\n'; const BASE_ANSWER_SDP = 'v=0\r\no=- 222 2 IN IP4 0.0.0.0\r\ns=answer-stream\r\n'; const EMBEDDED_CANDIDATE_LINE = 'a=candidate:1 1 udp 2130706431 10.0.0.1 3478 typ relay\r\n'; type GatherMode = 'complete' | 'stall'; /** * Install a fake RTCPeerConnection as the global. In 'complete' mode, answering an * offer schedules (next macrotask) a single trickle candidate + gathering-complete * so the local SDP ends up with the candidate embedded; 'stall' never completes. */ function installFakePeerConnection(gather: GatherMode): { restore: () => void } { class FakeRTCPeerConnection { connectionState = 'new'; signalingState = 'stable'; iceGatheringState: RTCIceGatheringState = 'new'; onicecandidate: ((event: { candidate: unknown }) => void) | null = null; onicegatheringstatechange: (() => void) | null = null; oniceconnectionstatechange: (() => void) | null = null; onconnectionstatechange: (() => void) | null = null; ontrack: ((event: unknown) => void) | null = null; remoteDescription: { type: string; sdp: string } | null = null; localDescription: { type: string; sdp: string } | null = null; private gatheringListeners = new Set<() => void>(); addEventListener(type: string, listener: () => void): void { if (type === 'icegatheringstatechange') this.gatheringListeners.add(listener); } removeEventListener(type: string, listener: () => void): void { if (type === 'icegatheringstatechange') this.gatheringListeners.delete(listener); } addTransceiver(): void {} addTrack(): { replaceTrack: () => void } { return { replaceTrack: () => {} }; } getSenders(): unknown[] { return []; } async setRemoteDescription(desc: { type: string; sdp: string }): Promise { this.remoteDescription = desc; } async createAnswer(): Promise<{ type: string; sdp: string }> { return { type: 'answer', sdp: BASE_ANSWER_SDP }; } async setLocalDescription(desc: { type: string; sdp: string }): Promise { this.localDescription = { type: desc.type, sdp: desc.sdp }; if (desc.type === 'answer' && gather === 'complete') { setTimeout(() => this.completeGathering(), 0); } } private completeGathering(): void { // Emit one trickle candidate (the sink path must swallow this), then embed it // into the local SDP and flip gathering to 'complete'. this.onicecandidate?.({ candidate: { candidate: 'candidate:1 1 udp 1 10.0.0.1 3478 typ relay', sdpMid: '0' } }); if (this.localDescription) this.localDescription.sdp += EMBEDDED_CANDIDATE_LINE; this.iceGatheringState = 'complete'; this.onicegatheringstatechange?.(); for (const listener of this.gatheringListeners) listener(); } async getStats(): Promise> { return new Map(); } close(): void {} } const original = (globalThis as { RTCPeerConnection?: unknown }).RTCPeerConnection; (globalThis as { RTCPeerConnection?: unknown }).RTCPeerConnection = FakeRTCPeerConnection; return { restore: () => { (globalThis as { RTCPeerConnection?: unknown }).RTCPeerConnection = original; }, }; } /** A transport that records every outbound send, so tests can assert nothing hit the wire. */ function makeRecordingTransport(): TwinTransport & { sends: Array<{ target: string; data: unknown }> } { const sends: Array<{ target: string; data: unknown }> = []; return { twinId: 'cam-twin', sendMessage: async (target: string, data: unknown) => { sends.push({ target, data }); }, subscribe: async () => {}, onMessage: () => {}, offMessage: () => {}, sends, }; } /** An empty local stream — enough for a sendonly responder with no real tracks. */ const emptyLocalStream = () => ({ getTracks: () => [] }) as unknown as MediaStream; const flush = (ms: number) => new Promise((resolve) => setTimeout(resolve, ms)); describe('one-shot HTTP-signaled answer (WHEP)', () => { it('ingestExternalOffer returns an answer SDP once ICE gathering completes', async () => { const fake = installFakePeerConnection('complete'); const manager = new WebRTCManager(makeRecordingTransport()); try { await manager.acceptMediaStream('cam-twin', { createLocalStream: emptyLocalStream }, () => {}); const answerSdp = await manager.answerMediaOffer('cam-twin', OFFER_SDP, { viewerId: 'viewer-1' }); expect(answerSdp).toContain('answer'); // Gathering-complete: the candidate is embedded (non-trickle). expect(answerSdp).toContain('a=candidate'); } finally { manager.close(); fake.restore(); } }); it('falls back to the pre-gather answer at the cap when ICE gathering stalls', async () => { const fake = installFakePeerConnection('stall'); const realSetTimeout = global.setTimeout; // Collapse only the 3000ms gathering cap so the test doesn't wait it out; every // other timer (connect/connection timeouts, retry backoff) keeps its real delay. (global as { setTimeout: typeof setTimeout }).setTimeout = (( handler: (...timeoutArgs: unknown[]) => void, delay?: number, ...timeoutArgs: unknown[] ) => realSetTimeout(handler, delay === 3000 ? 0 : delay, ...timeoutArgs)) as unknown as typeof setTimeout; const manager = new WebRTCManager(makeRecordingTransport()); try { await manager.acceptMediaStream('cam-twin', { createLocalStream: emptyLocalStream }, () => {}); const answerSdp = await manager.answerMediaOffer('cam-twin', OFFER_SDP, { viewerId: 'viewer-1' }); expect(answerSdp).toContain('answer'); // Stalled gather: cap fired, so the emitted answer has no candidates embedded. expect(answerSdp).not.toContain('a=candidate'); } finally { global.setTimeout = realSetTimeout; manager.close(); fake.restore(); } }); it('throws PeerCapReachedError when the responder is already at its maxPeers cap', async () => { const fake = installFakePeerConnection('complete'); const manager = new WebRTCManager(makeRecordingTransport()); try { await manager.acceptMediaStream('cam-twin', { maxPeers: 1, createLocalStream: emptyLocalStream }, () => {}); // viewer-1 takes the one slot (it stays occupied until it connects or times out). await manager.answerMediaOffer('cam-twin', OFFER_SDP, { viewerId: 'viewer-1' }); await expect(manager.answerMediaOffer('cam-twin', OFFER_SDP, { viewerId: 'viewer-2' })).rejects.toBeInstanceOf( PeerCapReachedError, ); } finally { manager.close(); fake.restore(); } }); it('throws NoMediaResponderError when no media responder is armed for the twin+channel', async () => { const fake = installFakePeerConnection('complete'); const manager = new WebRTCManager(makeRecordingTransport()); try { await expect(manager.answerMediaOffer('cam-twin', OFFER_SDP, { viewerId: 'viewer-1' })).rejects.toBeInstanceOf( NoMediaResponderError, ); // Arming the 'default' channel does not satisfy a request for another channel. await manager.acceptMediaStream('cam-twin', { createLocalStream: emptyLocalStream }, () => {}); await expect( manager.answerMediaOffer('cam-twin', OFFER_SDP, { viewerId: 'viewer-1', channelName: 'other' }), ).rejects.toBeInstanceOf(NoMediaResponderError); } finally { manager.close(); fake.restore(); } }); it('routes the answer to the signalingSink and swallows trickle candidate sends', async () => { const fake = installFakePeerConnection('complete'); const transport = makeRecordingTransport(); const captured: Array<{ type: string; data: { sdp?: string } }> = []; const config: PeerConnectionConfig = { targetTwinId: 'direct-http', isInitiator: false, connectionType: 'mediastream', channelPrefix: 'media-default', useStun: true, stunServers: [], turnServers: [], iceTransportPolicy: 'all', selfManagedSignaling: true, signalingSink: (type, data) => captured.push({ type, data: data as { sdp?: string } }), onConnected: () => {}, onDisconnected: () => {}, onError: () => {}, onPeerConnectionCreated: () => {}, }; const manager = new PeerConnectionManager(config, transport, { connectionTimeout: 15000, initialRetryDelay: 1000, maxRetryDelay: 30000, }); try { manager.connect().catch(() => {}); // never reaches 'connected' with the fake — intentional await manager.ingestSignalingMessage({ sourceTwinId: 'direct-http', data: { type: 'media-default-direct-http:offer', peerId: 'viewer-1', data: { type: 'offer', sdp: OFFER_SDP }, }, }); await flush(50); const answers = captured.filter((message) => message.type === 'answer'); expect(answers.length).toBe(1); expect(answers[0].data.sdp).toContain('a=candidate'); // Candidates are swallowed (never routed to the sink) and nothing hit the transport. expect(captured.some((message) => message.type === 'ice')).toBe(false); expect(transport.sends.length).toBe(0); } finally { manager.close(); fake.restore(); } }); it('still evicts a viewer that takes the answer but never connects (connect-timeout)', async () => { const fake = installFakePeerConnection('complete'); const transport = makeRecordingTransport(); const closedPeers: Array<{ peerId: string; error?: Error }> = []; const fanout = new ResponderFanout({ transport, targetTwinId: 'cam-twin', channelPrefix: 'media-default', connectTimeoutMs: 40, createSession: (peerId, replyToTwinId, onClosed, signalingSink) => { const handler = new MediaStreamHandler( replyToTwinId, false, transport, { mediaOptions: { direction: 'sendonly', localStream: emptyLocalStream() }, peerId, selfManagedSignaling: true, signalingSink, }, 'default', ); handler.setCallbacks({ onDisconnected: onClosed }); return handler; }, onPeer: () => {}, onPeerClosed: (peerId, error) => closedPeers.push({ peerId, error }), }); try { await fanout.start(); const answerSdp = await fanout.ingestExternalOffer('viewer-1', OFFER_SDP); expect(answerSdp).toContain('answer'); // The viewer holds its slot immediately after answering. expect(fanout.getPeerCount()).toBe(1); // It never connects; past the connect timeout the fan-out reaps it and frees the slot. await flush(90); expect(fanout.getPeerCount()).toBe(0); expect(closedPeers.some((entry) => entry.peerId === 'viewer-1')).toBe(true); } finally { fanout.close(); fake.restore(); } }); });