From 2f188508e62e20594c86e2d2dbd694491fbd329f Mon Sep 17 00:00:00 2001 From: Phuc Nguyen Date: Fri, 28 Aug 2026 14:02:47 +0700 Subject: [PATCH 1/2] Stabilize mobile realtime lifecycle handling --- apps/mobile/README.md | 17 + apps/mobile/src/hooks/useChat.test.ts | 144 +++- apps/mobile/src/hooks/useChat.ts | 144 ++-- apps/mobile/src/hooks/useScreenShare.test.ts | 28 + apps/mobile/src/hooks/useScreenShare.ts | 21 +- apps/mobile/src/hooks/useWebRTCHost.test.ts | 302 +++++++- apps/mobile/src/hooks/useWebRTCHost.ts | 456 ++++++++---- apps/mobile/src/hooks/useWebRTCViewer.test.ts | 430 +++++++++++- apps/mobile/src/hooks/useWebRTCViewer.ts | 663 ++++++++++++------ apps/mobile/src/lib/event-source.test.ts | 9 + apps/mobile/src/lib/event-source.ts | 2 +- apps/mobile/src/test/setup.ts | 36 +- 12 files changed, 1839 insertions(+), 413 deletions(-) diff --git a/apps/mobile/README.md b/apps/mobile/README.md index eae78ec2..80f9f9b0 100644 --- a/apps/mobile/README.md +++ b/apps/mobile/README.md @@ -56,6 +56,23 @@ pnpm mobile:android pnpm mobile:ios ``` +## Call continuity + +The mobile host and viewer intentionally close their SSE, WebRTC peer, media, and stats resources +when the app leaves the foreground. If a call was active or still connecting, returning to the +foreground starts exactly one fresh connection through the normal authentication and signaling +path. A server heartbeat watchdog also replaces a connection that has stopped receiving SSE +heartbeats for 75 seconds. + +Reconnects preserve the user's microphone mute choice. Async media and signaling callbacks are +scoped to a connection generation, so a late callback from a backgrounded or unmounted screen +cannot restore stale peers, viewers, chat history, or microphone state. This is foreground call +recovery, not background audio support. + +An active screen share ends when the app reaches the background and must be started again after +returning. The capture hook owns that teardown so its UI cannot report a stopped native track as +still sharing. A brief iOS `inactive` transition alone does not tear down the call or screen share. + ## EAS builds Link the app to the intended Expo project and provide its UUID through `EAS_PROJECT_ID`. No diff --git a/apps/mobile/src/hooks/useChat.test.ts b/apps/mobile/src/hooks/useChat.test.ts index 7a58fa12..21398e72 100644 --- a/apps/mobile/src/hooks/useChat.test.ts +++ b/apps/mobile/src/hooks/useChat.test.ts @@ -22,14 +22,14 @@ describe('useChat', () => { vi.clearAllMocks(); }); - it('should initialize with loading state', () => { + it('should initialize without loading when disabled', () => { vi.mocked(chatApi.getHistory).mockResolvedValue({ data: { messages: [], hasMore: false }, }); const { result } = renderHook(() => useChat({ sessionId: 'session-1', enabled: false })); - expect(result.current.loading).toBe(true); + expect(result.current.loading).toBe(false); expect(result.current.messages).toEqual([]); expect(result.current.sending).toBe(false); expect(result.current.error).toBeNull(); @@ -53,7 +53,7 @@ describe('useChat', () => { const { result } = renderHook(() => useChat({ sessionId: 'session-1', enabled: false })); expect(chatApi.getHistory).not.toHaveBeenCalled(); - expect(result.current.loading).toBe(true); + expect(result.current.loading).toBe(false); }); it('should send messages', async () => { @@ -138,4 +138,142 @@ describe('useChat', () => { expect(result.current.error).toBe('Send failed'); }); + + it('ignores history returned by a previous session generation', async () => { + let resolveOldHistory!: (value: { + data: { messages: ChatMessage[]; hasMore: boolean }; + }) => void; + vi.mocked(chatApi.getHistory) + .mockReturnValueOnce( + new Promise((resolve) => { + resolveOldHistory = resolve; + }) + ) + .mockResolvedValueOnce({ + data: { + messages: [{ ...mockMessage, id: 'msg-new', session_id: 'session-2' }], + hasMore: false, + }, + }); + + const { result, rerender } = renderHook( + ({ sessionId }: { sessionId: string }) => useChat({ sessionId, enabled: true }), + { initialProps: { sessionId: 'session-1' } } + ); + + rerender({ sessionId: 'session-2' }); + await waitFor(() => { + expect(result.current.messages.map((message) => message.id)).toEqual(['msg-new']); + }); + + await act(async () => { + resolveOldHistory({ data: { messages: [mockMessage], hasMore: false } }); + await Promise.resolve(); + }); + + expect(result.current.messages.map((message) => message.id)).toEqual(['msg-new']); + }); + + it('does not overlap polling requests within the same generation', async () => { + vi.useFakeTimers(); + let resolveFirstPoll!: (value: { data: { messages: ChatMessage[]; hasMore: boolean } }) => void; + vi.mocked(chatApi.getHistory) + .mockReturnValueOnce( + new Promise((resolve) => { + resolveFirstPoll = resolve; + }) + ) + .mockResolvedValue({ data: { messages: [], hasMore: false } }); + + try { + renderHook(() => useChat({ sessionId: 'session-1', enabled: true })); + await act(async () => { + await Promise.resolve(); + }); + expect(chatApi.getHistory).toHaveBeenCalledTimes(1); + + await act(async () => { + await vi.advanceTimersByTimeAsync(2000); + }); + expect(chatApi.getHistory).toHaveBeenCalledTimes(1); + + await act(async () => { + resolveFirstPoll({ data: { messages: [], hasMore: false } }); + await Promise.resolve(); + await vi.advanceTimersByTimeAsync(2000); + }); + expect(chatApi.getHistory).toHaveBeenCalledTimes(2); + } finally { + vi.useRealTimers(); + } + }); + + it('ignores a send response after the active session changes', async () => { + let resolveSend!: (value: { data: ChatMessage }) => void; + vi.mocked(chatApi.getHistory).mockResolvedValue({ + data: { messages: [], hasMore: false }, + }); + vi.mocked(chatApi.send).mockReturnValueOnce( + new Promise((resolve) => { + resolveSend = resolve; + }) + ); + + const { result, rerender } = renderHook( + ({ sessionId }: { sessionId: string }) => useChat({ sessionId, enabled: true }), + { initialProps: { sessionId: 'session-1' } } + ); + await waitFor(() => expect(result.current.loading).toBe(false)); + + let sendPromise!: Promise; + act(() => { + sendPromise = result.current.sendMessage('old session message'); + }); + rerender({ sessionId: 'session-2' }); + await waitFor(() => expect(result.current.loading).toBe(false)); + + await act(async () => { + resolveSend({ data: mockMessage }); + await sendPromise; + }); + + expect(result.current.messages).toEqual([]); + expect(result.current.sending).toBe(false); + }); + + it('clears a stale sending state when chat is disabled and re-enabled', async () => { + let resolveSend!: (value: { data: ChatMessage }) => void; + vi.mocked(chatApi.getHistory).mockResolvedValue({ + data: { messages: [], hasMore: false }, + }); + vi.mocked(chatApi.send).mockReturnValueOnce( + new Promise((resolve) => { + resolveSend = resolve; + }) + ); + + const { result, rerender } = renderHook( + ({ enabled }: { enabled: boolean }) => useChat({ sessionId: 'session-1', enabled }), + { initialProps: { enabled: true } } + ); + await waitFor(() => expect(result.current.loading).toBe(false)); + + let sendPromise!: Promise; + act(() => { + sendPromise = result.current.sendMessage('message before disable'); + }); + expect(result.current.sending).toBe(true); + + rerender({ enabled: false }); + expect(result.current.sending).toBe(false); + rerender({ enabled: true }); + + await act(async () => { + resolveSend({ data: mockMessage }); + await sendPromise; + }); + + expect(result.current.messages).toEqual([]); + expect(result.current.sending).toBe(false); + }); }); diff --git a/apps/mobile/src/hooks/useChat.ts b/apps/mobile/src/hooks/useChat.ts index 28311882..e0dbb7c6 100644 --- a/apps/mobile/src/hooks/useChat.ts +++ b/apps/mobile/src/hooks/useChat.ts @@ -36,72 +36,122 @@ export function useChat({ const seenIdsRef = useRef>(new Set()); const pollIntervalRef = useRef | null>(null); + const pollGenerationRef = useRef(0); + const pollInFlightRef = useRef(null); + const sendOperationRef = useRef(null); + const mountedRef = useRef(true); + const previousSessionIdRef = useRef(sessionId); // Fetch messages - const fetchMessages = useCallback(async () => { - try { - const result = await chatApi.getHistory(sessionId, { limit: 100 }); + const fetchMessages = useCallback( + async (generation: number) => { + if (pollInFlightRef.current === generation) return; + pollInFlightRef.current = generation; - if (result.error) { - setError(result.error); - return; - } + try { + const result = await chatApi.getHistory(sessionId, { limit: 100 }); + if (!mountedRef.current || pollGenerationRef.current !== generation) return; - if (result.data?.messages) { - const newMessages: ChatMessage[] = []; - for (const msg of result.data.messages) { - if (!seenIdsRef.current.has(msg.id)) { - seenIdsRef.current.add(msg.id); - newMessages.push(msg); - } + if (result.error) { + setError(result.error); + return; } - if (newMessages.length > 0) { - setMessages((prev) => { - const combined = [...prev, ...newMessages]; - // Sort by timestamp ascending - combined.sort( - (a, b) => new Date(a.created_at).getTime() - new Date(b.created_at).getTime() - ); - return combined; - }); - } + if (result.data?.messages) { + const newMessages: ChatMessage[] = []; + for (const msg of result.data.messages) { + if (!seenIdsRef.current.has(msg.id)) { + seenIdsRef.current.add(msg.id); + newMessages.push(msg); + } + } - setError(null); + if (newMessages.length > 0) { + setMessages((prev) => { + const combined = [...prev, ...newMessages]; + // Sort by timestamp ascending + combined.sort( + (a, b) => new Date(a.created_at).getTime() - new Date(b.created_at).getTime() + ); + return combined; + }); + } + + setError(null); + } + } catch { + if (mountedRef.current && pollGenerationRef.current === generation) { + setError('Failed to fetch messages'); + } + } finally { + if (pollInFlightRef.current === generation) { + pollInFlightRef.current = null; + } + if (mountedRef.current && pollGenerationRef.current === generation) { + setLoading(false); + } } - } catch { - setError('Failed to fetch messages'); - } finally { - setLoading(false); - } - }, [sessionId]); + }, + [sessionId] + ); // Start polling when enabled useEffect(() => { - if (!enabled) return; + const generation = ++pollGenerationRef.current; + const sessionChanged = previousSessionIdRef.current !== sessionId; + previousSessionIdRef.current = sessionId; + + // A generation change invalidates any send started by the previous chat + // lifecycle, including enable/disable transitions within the same session. + if (sendOperationRef.current) { + sendOperationRef.current = null; + setSending(false); + } + + if (sessionChanged) { + seenIdsRef.current = new Set(); + setMessages([]); + setError(null); + setLoading(true); + } + + if (!enabled) { + setLoading(false); + return undefined; + } - void fetchMessages(); + setLoading(true); + void fetchMessages(generation); pollIntervalRef.current = setInterval(() => { - void fetchMessages(); + void fetchMessages(generation); }, POLL_INTERVAL); return () => { + if (pollGenerationRef.current === generation) { + pollGenerationRef.current += 1; + } if (pollIntervalRef.current) { clearInterval(pollIntervalRef.current); pollIntervalRef.current = null; } }; - }, [enabled, fetchMessages]); + }, [enabled, fetchMessages, sessionId]); // Send message const sendMessage = useCallback( async (content: string) => { - if (!content.trim()) return; + const trimmed = content.trim(); + if (!trimmed || sendOperationRef.current) return; + + const generation = pollGenerationRef.current; + const operation = Symbol('chat-send'); + sendOperationRef.current = operation; setSending(true); try { - const result = await chatApi.send(sessionId, content.trim(), participantId); + const result = await chatApi.send(sessionId, trimmed, participantId); + if (!mountedRef.current || pollGenerationRef.current !== generation) return; if (result.error) { setError(result.error); @@ -116,14 +166,30 @@ export function useChat({ setError(null); } catch { - setError('Failed to send message'); + if (mountedRef.current && pollGenerationRef.current === generation) { + setError('Failed to send message'); + } } finally { - setSending(false); + if (sendOperationRef.current === operation) { + sendOperationRef.current = null; + if (mountedRef.current && pollGenerationRef.current === generation) { + setSending(false); + } + } } }, [sessionId, participantId] ); + useEffect(() => { + mountedRef.current = true; + return () => { + mountedRef.current = false; + pollGenerationRef.current += 1; + sendOperationRef.current = null; + }; + }, []); + return { messages, loading, diff --git a/apps/mobile/src/hooks/useScreenShare.test.ts b/apps/mobile/src/hooks/useScreenShare.test.ts index 29b8153e..6a762497 100644 --- a/apps/mobile/src/hooks/useScreenShare.test.ts +++ b/apps/mobile/src/hooks/useScreenShare.test.ts @@ -3,6 +3,7 @@ import { beforeEach, describe, expect, it, vi } from 'vitest'; import { mediaDevices } from 'react-native-webrtc'; import type { MediaStream, MediaStreamTrack } from 'react-native-webrtc'; import { useScreenShare } from './useScreenShare'; +import { emitAppStateChange } from '../test/setup'; function deferred() { let resolve!: (value: T | PromiseLike) => void; @@ -217,6 +218,33 @@ describe('useScreenShare', () => { expect(unpublishStream).toHaveBeenCalledTimes(1); }); + it('survives an inactive interruption but stops cleanly in the background', async () => { + const capture = createCapture(); + const publishStream = vi.fn().mockResolvedValue(undefined); + const unpublishStream = vi.fn().mockResolvedValue(undefined); + vi.mocked(mediaDevices.getDisplayMedia).mockResolvedValue(capture.stream); + + const { result } = renderHook(() => useScreenShare({ publishStream, unpublishStream })); + await act(async () => { + await result.current.start(); + }); + + act(() => { + emitAppStateChange('inactive'); + }); + expect(result.current.state).toBe('active'); + expect(capture.track.stop).not.toHaveBeenCalled(); + + act(() => { + emitAppStateChange('background'); + }); + await waitFor(() => expect(result.current.state).toBe('idle')); + + expect(capture.track.stop).toHaveBeenCalledTimes(1); + expect(unpublishStream).toHaveBeenCalledTimes(1); + expect(result.current.isSharing).toBe(false); + }); + it('serializes duplicate stop requests', async () => { const capture = createCapture(); const unpublication = deferred(); diff --git a/apps/mobile/src/hooks/useScreenShare.ts b/apps/mobile/src/hooks/useScreenShare.ts index 243e8585..2e9e8420 100644 --- a/apps/mobile/src/hooks/useScreenShare.ts +++ b/apps/mobile/src/hooks/useScreenShare.ts @@ -1,4 +1,5 @@ import { useCallback, useEffect, useRef, useState } from 'react'; +import { AppState, type AppStateStatus } from 'react-native'; import { mediaDevices } from 'react-native-webrtc'; import type { MediaStream, MediaStreamTrack } from 'react-native-webrtc'; @@ -47,6 +48,7 @@ export function useScreenShare({ const stopPromiseRef = useRef | null>(null); const publicationRef = useRef | null>(null); const endedListenersRef = useRef void>>(new Map()); + const appStateRef = useRef(AppState.currentState); const transition = useCallback((nextState: ScreenShareState) => { stateRef.current = nextState; @@ -119,7 +121,7 @@ export function useScreenShare({ ); const start = useCallback(async (): Promise => { - if (stateRef.current !== 'idle' || stopPromiseRef.current) { + if (appStateRef.current !== 'active' || stateRef.current !== 'idle' || stopPromiseRef.current) { return false; } @@ -187,6 +189,23 @@ export function useScreenShare({ setError(null); }, []); + // Native capture belongs to this hook. End it explicitly when the app is + // backgrounded so the UI cannot keep claiming that a stopped track is live. + useEffect(() => { + appStateRef.current = AppState.currentState; + const subscription = AppState.addEventListener('change', (nextState) => { + const previousState = appStateRef.current; + appStateRef.current = nextState; + if (nextState === 'background' && previousState !== 'background') { + void stop(); + } + }); + + return () => { + subscription.remove(); + }; + }, [stop]); + useEffect(() => { mountedRef.current = true; return () => { diff --git a/apps/mobile/src/hooks/useWebRTCHost.test.ts b/apps/mobile/src/hooks/useWebRTCHost.test.ts index 63ee4032..4421ba45 100644 --- a/apps/mobile/src/hooks/useWebRTCHost.test.ts +++ b/apps/mobile/src/hooks/useWebRTCHost.test.ts @@ -3,7 +3,9 @@ import { renderHook, act, waitFor } from '@testing-library/react'; import { useWebRTCHost } from './useWebRTCHost'; import { createEventSource } from '../lib/event-source'; import { getStoredAuth } from '../lib/secure-storage'; +import { mediaDevices } from 'react-native-webrtc'; import type { MediaStream, MediaStreamTrack } from 'react-native-webrtc'; +import { emitAppStateChange } from '../test/setup'; vi.mock('../config', () => ({ API_BASE_URL: 'https://pairux.com', @@ -19,21 +21,37 @@ vi.mock('../lib/secure-storage', () => ({ isAuthExpired: vi.fn().mockReturnValue(false), })); -const { mockClose, mockAddEventListener } = vi.hoisted(() => ({ +const { mockClose, mockAddEventListener, mockEventSources } = vi.hoisted(() => ({ mockClose: vi.fn(), mockAddEventListener: vi.fn(), + mockEventSources: [] as { + listeners: Map void>; + close: ReturnType; + }[], })); vi.mock('../lib/event-source', () => ({ - createEventSource: vi.fn(() => ({ - addEventListener: mockAddEventListener, - close: mockClose, - })), + createEventSource: vi.fn(() => { + const listeners = new Map void>(); + const source = { + listeners, + addEventListener: vi.fn((event: string, listener: (payload: { data: string }) => void) => { + mockAddEventListener(event, listener); + listeners.set(event, listener); + }), + close: vi.fn(() => { + mockClose(); + }), + }; + mockEventSources.push(source); + return source; + }), })); describe('useWebRTCHost', () => { beforeEach(() => { vi.clearAllMocks(); + mockEventSources.length = 0; vi.mocked(fetch).mockResolvedValue({ ok: true, text: async () => 'ok', @@ -300,4 +318,278 @@ describe('useWebRTCHost', () => { false ); }); + + it('tears down once in the background and resumes exactly once when active', async () => { + const { result } = renderHook(() => + useWebRTCHost({ + sessionId: 'session-1', + hostId: 'host-1', + }) + ); + + await act(async () => { + await result.current.startHosting(); + }); + expect(createEventSource).toHaveBeenCalledTimes(1); + + act(() => { + emitAppStateChange('inactive'); + }); + expect(mockClose).not.toHaveBeenCalled(); + + act(() => { + emitAppStateChange('background'); + }); + expect(mockClose).toHaveBeenCalledTimes(1); + + await act(async () => { + emitAppStateChange('active'); + await Promise.resolve(); + await Promise.resolve(); + }); + + expect(createEventSource).toHaveBeenCalledTimes(2); + }); + + it('accepts a viewer presence event during a transient inactive window', async () => { + const { result } = renderHook(() => + useWebRTCHost({ + sessionId: 'session-1', + hostId: 'host-1', + }) + ); + + await act(async () => { + await result.current.startHosting(); + }); + const source = mockEventSources[0]; + expect(source).toBeDefined(); + + act(() => { + source.listeners.get('connected')?.({ data: '{}' }); + emitAppStateChange('inactive'); + }); + await act(async () => { + source.listeners.get('presence-join')?.({ + data: JSON.stringify({ presences: [{ user_id: 'viewer-1', role: 'viewer' }] }), + }); + await Promise.resolve(); + await Promise.resolve(); + }); + + await waitFor(() => expect(result.current.viewerCount).toBe(1)); + expect(result.current.viewers.has('viewer-1')).toBe(true); + expect(source.close).not.toHaveBeenCalled(); + }); + + it('starts on foreground when hosting was requested while inactive', async () => { + emitAppStateChange('inactive'); + const { result } = renderHook(() => + useWebRTCHost({ + sessionId: 'session-1', + hostId: 'host-1', + }) + ); + + await act(async () => { + await result.current.startHosting(); + }); + expect(createEventSource).not.toHaveBeenCalled(); + + await act(async () => { + emitAppStateChange('active'); + await Promise.resolve(); + await Promise.resolve(); + }); + + expect(createEventSource).toHaveBeenCalledTimes(1); + }); + + it('does not stop a caller-owned screen stream when hosting stops', async () => { + const { result } = renderHook(() => + useWebRTCHost({ + sessionId: 'session-1', + hostId: 'host-1', + }) + ); + await act(async () => { + await result.current.startHosting(); + }); + + const screenTrack = { + id: 'screen', + kind: 'video', + stop: vi.fn(), + } as unknown as MediaStreamTrack & { stop: ReturnType }; + const screenStream = { getTracks: () => [screenTrack] } as unknown as MediaStream; + await act(async () => { + await result.current.publishStream(screenStream); + }); + + act(() => { + result.current.stopHosting(); + }); + + expect(screenTrack.stop).not.toHaveBeenCalled(); + }); + + it('ignores presence callbacks from a retired background generation', async () => { + const { result } = renderHook(() => + useWebRTCHost({ + sessionId: 'session-1', + hostId: 'host-1', + }) + ); + + await act(async () => { + await result.current.startHosting(); + }); + const stalePresenceListener = mockEventSources[0]?.listeners.get('presence-join'); + + act(() => { + emitAppStateChange('background'); + }); + await act(async () => { + emitAppStateChange('active'); + await Promise.resolve(); + await Promise.resolve(); + }); + const currentPresenceListener = mockEventSources[1]?.listeners.get('presence-join'); + + await act(async () => { + stalePresenceListener?.({ + data: JSON.stringify({ presences: [{ user_id: 'stale-viewer', role: 'viewer' }] }), + }); + await Promise.resolve(); + }); + expect(result.current.viewerCount).toBe(0); + + act(() => { + currentPresenceListener?.({ + data: JSON.stringify({ presences: [{ user_id: 'current-viewer', role: 'viewer' }] }), + }); + }); + await waitFor(() => expect(result.current.viewerCount).toBe(1)); + expect(result.current.viewers.has('current-viewer')).toBe(true); + }); + + it('preserves the host microphone mute intent across a foreground reconnect', async () => { + const tracks: (MediaStreamTrack & { + enabled: boolean; + stop: ReturnType; + })[] = []; + const makeMicStream = (): MediaStream => { + const track = { + kind: 'audio', + enabled: true, + stop: vi.fn(), + } as unknown as MediaStreamTrack & { + enabled: boolean; + stop: ReturnType; + }; + tracks.push(track); + return { + getTracks: () => [track], + getAudioTracks: () => [track], + } as unknown as MediaStream; + }; + vi.mocked(mediaDevices.getUserMedia) + .mockResolvedValueOnce(makeMicStream()) + .mockResolvedValueOnce(makeMicStream()); + + const { result } = renderHook(() => + useWebRTCHost({ + sessionId: 'session-1', + hostId: 'host-1', + }) + ); + + await act(async () => { + await result.current.startHosting(); + }); + act(() => { + result.current.toggleMic(); + }); + expect(tracks[0]?.enabled).toBe(false); + + act(() => { + emitAppStateChange('background'); + }); + await act(async () => { + emitAppStateChange('active'); + await Promise.resolve(); + await Promise.resolve(); + }); + + expect(tracks[1]?.enabled).toBe(false); + expect(result.current.micEnabled).toBe(false); + }); + + it('stops a late microphone stream when unmounted during startup', async () => { + let resolveMic!: (stream: MediaStream) => void; + const lateTrack = { + kind: 'audio', + enabled: true, + stop: vi.fn(), + } as unknown as MediaStreamTrack & { stop: ReturnType }; + const lateStream = { + getTracks: () => [lateTrack], + getAudioTracks: () => [lateTrack], + } as unknown as MediaStream; + vi.mocked(mediaDevices.getUserMedia).mockReturnValueOnce( + new Promise((resolve) => { + resolveMic = resolve; + }) + ); + + const { result, unmount } = renderHook(() => + useWebRTCHost({ + sessionId: 'session-1', + hostId: 'host-1', + }) + ); + + let startPromise!: Promise; + await act(async () => { + startPromise = result.current.startHosting(); + await Promise.resolve(); + }); + unmount(); + + await act(async () => { + resolveMic(lateStream); + await startPromise; + }); + + expect(lateTrack.stop).toHaveBeenCalledTimes(1); + expect(createEventSource).not.toHaveBeenCalled(); + }); + + it('restarts once after the SSE heartbeat expires', async () => { + vi.useFakeTimers(); + vi.setSystemTime(new Date('2026-08-28T00:00:00Z')); + + try { + const { result } = renderHook(() => + useWebRTCHost({ + sessionId: 'session-1', + hostId: 'host-1', + }) + ); + + await act(async () => { + await result.current.startHosting(); + }); + + await act(async () => { + await vi.advanceTimersByTimeAsync(90_001); + await Promise.resolve(); + }); + + expect(mockClose).toHaveBeenCalledTimes(1); + expect(createEventSource).toHaveBeenCalledTimes(2); + } finally { + vi.useRealTimers(); + } + }); }); diff --git a/apps/mobile/src/hooks/useWebRTCHost.ts b/apps/mobile/src/hooks/useWebRTCHost.ts index 5ac9e345..ab0ea0c5 100644 --- a/apps/mobile/src/hooks/useWebRTCHost.ts +++ b/apps/mobile/src/hooks/useWebRTCHost.ts @@ -9,6 +9,7 @@ * - No HTMLAudioElement (audio handled by RN WebRTC) */ import { useState, useEffect, useRef, useCallback } from 'react'; +import { AppState, type AppStateStatus } from 'react-native'; import type { MediaStreamTrack, RTCRtpSender } from 'react-native-webrtc'; import { RTCPeerConnection, RTCIceCandidate, MediaStream, mediaDevices } from 'react-native-webrtc'; import type { @@ -63,6 +64,8 @@ const _BITRATE_PRESETS: Record = { }; const STATS_INTERVAL = 30000; +const HEARTBEAT_TIMEOUT = 75000; +const HEARTBEAT_WATCHDOG_INTERVAL = 15000; const DEFAULT_ICE_SERVERS: RTCIceServer[] = [ { urls: 'stun:stun.l.google.com:19302' }, @@ -137,9 +140,16 @@ export function useWebRTCHost({ const eventSourceRef = useRef(null); const viewersRef = useRef>(new Map()); const statsIntervalRef = useRef | null>(null); + const heartbeatWatchdogRef = useRef | null>(null); + const lastHeartbeatAtRef = useRef(0); const removeViewerRef = useRef<((viewerId: string) => void) | undefined>(undefined); const authTokenRef = useRef(null); const isStartingRef = useRef(false); + const mountedRef = useRef(true); + const generationRef = useRef(0); + const appStateRef = useRef(AppState.currentState); + const resumeOnActiveRef = useRef(false); + const micEnabledIntentRef = useRef(true); const localStreamRef = useRef(null); const hostMicStreamRef = useRef(null); const publishedStreamSendersRef = useRef>>(new Map()); @@ -150,8 +160,19 @@ export function useWebRTCHost({ const onControlRequestRef = useRef(onControlRequest); const onInputReceivedRef = useRef(onInputReceived); + const onViewerJoinedRef = useRef(onViewerJoined); + const onViewerLeftRef = useRef(onViewerLeft); + const startHostingRef = useRef<(() => Promise) | undefined>(undefined); + const stopHostingRef = useRef<(() => void) | undefined>(undefined); onControlRequestRef.current = onControlRequest; onInputReceivedRef.current = onInputReceived; + onViewerJoinedRef.current = onViewerJoined; + onViewerLeftRef.current = onViewerLeft; + + const isCurrentGeneration = useCallback( + (generation: number) => mountedRef.current && generationRef.current === generation, + [] + ); // Send signal via API const sendSignal = useCallback( @@ -182,12 +203,22 @@ export function useWebRTCHost({ ); const renegotiateViewer = useCallback( - async (viewer: ViewerConnection): Promise => { + async (viewer: ViewerConnection, generation: number): Promise => { + if (!isCurrentGeneration(generation) || viewersRef.current.get(viewer.id) !== viewer) { + return; + } + const offer = (await viewer.peerConnection.createOffer({})) as OfferAnswer; + if (!isCurrentGeneration(generation) || viewersRef.current.get(viewer.id) !== viewer) { + return; + } if (offer.sdp) { offer.sdp = tuneOpusForVoice(offer.sdp); } await viewer.peerConnection.setLocalDescription(offer); + if (!isCurrentGeneration(generation) || viewersRef.current.get(viewer.id) !== viewer) { + return; + } if (!offer.sdp) { throw new Error(`Failed to create an SDP offer for viewer ${viewer.id}`); @@ -204,79 +235,87 @@ export function useWebRTCHost({ throw new Error(`Failed to signal viewer ${viewer.id}`); } }, - [hostId, sendSignal] + [hostId, isCurrentGeneration, sendSignal] ); // Report usage stats - const reportStats = useCallback(async () => { - for (const viewer of viewersRef.current.values()) { - if (viewer.connectionState !== 'connected') continue; + const reportStats = useCallback( + async (generation: number) => { + if (!isCurrentGeneration(generation)) return; + for (const viewer of viewersRef.current.values()) { + if (!isCurrentGeneration(generation)) return; + if (viewer.connectionState !== 'connected') continue; - try { - const stats = (await viewer.peerConnection.getStats()) as Map< - string, - Record - >; - let bytesSent = 0; - let bytesReceived = 0; - let packetsSent = 0; - let packetsReceived = 0; - let packetsLost = 0; - let roundTripTime: number | undefined; - let frameRate: number | undefined; - let frameWidth: number | undefined; - let frameHeight: number | undefined; - - stats.forEach((report: Record) => { - if (report.type === 'outbound-rtp' && report.kind === 'video') { - bytesSent += (report.bytesSent as number | undefined) ?? 0; - packetsSent += (report.packetsSent as number | undefined) ?? 0; - frameRate = report.framesPerSecond as number | undefined; - frameWidth = report.frameWidth as number | undefined; - frameHeight = report.frameHeight as number | undefined; - } - if (report.type === 'inbound-rtp') { - bytesReceived += (report.bytesReceived as number | undefined) ?? 0; - packetsReceived += (report.packetsReceived as number | undefined) ?? 0; - } - if (report.type === 'remote-inbound-rtp') { - packetsLost += (report.packetsLost as number | undefined) ?? 0; + try { + const stats = (await viewer.peerConnection.getStats()) as Map< + string, + Record + >; + if (!isCurrentGeneration(generation) || viewersRef.current.get(viewer.id) !== viewer) { + return; } - if (report.type === 'candidate-pair' && report.state === 'succeeded') { - roundTripTime = ((report.currentRoundTripTime as number | undefined) ?? 0) * 1000; + let bytesSent = 0; + let bytesReceived = 0; + let packetsSent = 0; + let packetsReceived = 0; + let packetsLost = 0; + let roundTripTime: number | undefined; + let frameRate: number | undefined; + let frameWidth: number | undefined; + let frameHeight: number | undefined; + + stats.forEach((report: Record) => { + if (report.type === 'outbound-rtp' && report.kind === 'video') { + bytesSent += (report.bytesSent as number | undefined) ?? 0; + packetsSent += (report.packetsSent as number | undefined) ?? 0; + frameRate = report.framesPerSecond as number | undefined; + frameWidth = report.frameWidth as number | undefined; + frameHeight = report.frameHeight as number | undefined; + } + if (report.type === 'inbound-rtp') { + bytesReceived += (report.bytesReceived as number | undefined) ?? 0; + packetsReceived += (report.packetsReceived as number | undefined) ?? 0; + } + if (report.type === 'remote-inbound-rtp') { + packetsLost += (report.packetsLost as number | undefined) ?? 0; + } + if (report.type === 'candidate-pair' && report.state === 'succeeded') { + roundTripTime = ((report.currentRoundTripTime as number | undefined) ?? 0) * 1000; + } + }); + + const headers: Record = { 'Content-Type': 'application/json' }; + if (authTokenRef.current) { + headers.Authorization = `Bearer ${authTokenRef.current}`; } - }); - const headers: Record = { 'Content-Type': 'application/json' }; - if (authTokenRef.current) { - headers.Authorization = `Bearer ${authTokenRef.current}`; + await fetch(`${API_BASE_URL}/api/sessions/${sessionId}/stats`, { + method: 'POST', + headers, + body: JSON.stringify({ + participantId: hostId, + role: 'host', + timestamp: Date.now(), + connectionState: viewer.connectionState, + bytesSent, + bytesReceived, + packetsSent, + packetsReceived, + packetsLost, + roundTripTime, + frameRate, + frameWidth, + frameHeight, + reportInterval: STATS_INTERVAL, + }), + }); + } catch (err) { + console.error('[WebRTCHost] Failed to report stats:', err); } - - await fetch(`${API_BASE_URL}/api/sessions/${sessionId}/stats`, { - method: 'POST', - headers, - body: JSON.stringify({ - participantId: hostId, - role: 'host', - timestamp: Date.now(), - connectionState: viewer.connectionState, - bytesSent, - bytesReceived, - packetsSent, - packetsReceived, - packetsLost, - roundTripTime, - frameRate, - frameWidth, - frameHeight, - reportInterval: STATS_INTERVAL, - }), - }); - } catch (err) { - console.error('[WebRTCHost] Failed to report stats:', err); } - } - }, [sessionId, hostId]); + }, + [hostId, isCurrentGeneration, sessionId] + ); // Handle data channel messages from viewer const handleDataChannelMessage = useCallback((viewerId: string, data: string) => { @@ -309,10 +348,12 @@ export function useWebRTCHost({ // Relay audio to other viewers const relayAudioToOtherViewers = useCallback( - async (sourceViewerId: string, audioTrack: MediaStreamTrack) => { + async (sourceViewerId: string, audioTrack: MediaStreamTrack, generation: number) => { + if (!isCurrentGeneration(generation)) return; const audioStream = new MediaStream([audioTrack]); for (const [otherId, otherViewer] of viewersRef.current.entries()) { + if (!isCurrentGeneration(generation)) return; if (otherId === sourceViewerId) continue; if ( otherViewer.connectionState !== 'connected' && @@ -322,11 +363,14 @@ export function useWebRTCHost({ try { await prioritizeAudioSender(otherViewer.peerConnection.addTrack(audioTrack, audioStream)); + if (!isCurrentGeneration(generation)) return; const offer = (await otherViewer.peerConnection.createOffer({})) as OfferAnswer; + if (!isCurrentGeneration(generation)) return; // In-band FEC turns a lost packet into a duller syllable, not a gap. if (offer.sdp) offer.sdp = tuneOpusForVoice(offer.sdp); await otherViewer.peerConnection.setLocalDescription(offer); + if (!isCurrentGeneration(generation)) return; if (offer.sdp) { await sendSignal({ @@ -342,12 +386,12 @@ export function useWebRTCHost({ } } }, - [hostId, sendSignal] + [hostId, isCurrentGeneration, sendSignal] ); // Create peer connection for a viewer const createPeerConnection = useCallback( - (viewerId: string): RTCPeerConnection => { + (viewerId: string, generation: number): RTCPeerConnection => { console.log('[WebRTCHost] Creating peer connection for viewer:', viewerId); const pc = new RTCPeerConnection({ @@ -376,6 +420,12 @@ export function useWebRTCHost({ // Handle incoming audio from viewer pc.addEventListener('track', (event) => { + if ( + !isCurrentGeneration(generation) || + viewersRef.current.get(viewerId)?.peerConnection !== pc + ) { + return; + } const track = event.track; if (track?.kind === 'audio') { console.log(`[WebRTCHost] Received audio track from viewer: ${viewerId}`); @@ -383,13 +433,19 @@ export function useWebRTCHost({ if (viewer) { viewer.audioTrack = track; setViewers(new Map(viewersRef.current)); - void relayAudioToOtherViewers(viewerId, track); + void relayAudioToOtherViewers(viewerId, track, generation); } } }); // Handle ICE candidates pc.addEventListener('icecandidate', (event) => { + if ( + !isCurrentGeneration(generation) || + viewersRef.current.get(viewerId)?.peerConnection !== pc + ) { + return; + } if (event.candidate) { void sendSignal({ type: 'ice-candidate', @@ -403,6 +459,12 @@ export function useWebRTCHost({ // Handle connection state changes pc.addEventListener('connectionstatechange', () => { + if ( + !isCurrentGeneration(generation) || + viewersRef.current.get(viewerId)?.peerConnection !== pc + ) { + return; + } const viewer = viewersRef.current.get(viewerId); if (viewer) { let newState: ConnectionState; @@ -450,6 +512,12 @@ export function useWebRTCHost({ const dc = pc.createDataChannel('control', { ordered: true }); dc.addEventListener('open', () => { + if ( + !isCurrentGeneration(generation) || + viewersRef.current.get(viewerId)?.peerConnection !== pc + ) { + return; + } const viewer = viewersRef.current.get(viewerId); if (viewer) { viewer.dataChannel = dc; @@ -458,6 +526,12 @@ export function useWebRTCHost({ }); dc.addEventListener('close', () => { + if ( + !isCurrentGeneration(generation) || + viewersRef.current.get(viewerId)?.peerConnection !== pc + ) { + return; + } const viewer = viewersRef.current.get(viewerId); if (viewer) { viewer.dataChannel = null; @@ -468,42 +542,48 @@ export function useWebRTCHost({ }); dc.addEventListener('message', (event) => { + if ( + !isCurrentGeneration(generation) || + viewersRef.current.get(viewerId)?.peerConnection !== pc + ) { + return; + } handleDataChannelMessage(viewerId, typeof event.data === 'string' ? event.data : ''); }); return pc; }, - [hostId, handleDataChannelMessage, sendSignal, relayAudioToOtherViewers] + [hostId, handleDataChannelMessage, isCurrentGeneration, sendSignal, relayAudioToOtherViewers] ); // Remove viewer - const removeViewer = useCallback( - (viewerId: string) => { - const viewer = viewersRef.current.get(viewerId); - if (viewer) { - console.log('[WebRTCHost] Removing viewer:', viewerId); - viewer.peerConnection.close(); - viewersRef.current.delete(viewerId); - publishedStreamSendersRef.current.delete(viewerId); - pendingCandidatesRef.current.delete(viewerId); - setViewers(new Map(viewersRef.current)); - onViewerLeft?.(viewerId); + const removeViewer = useCallback((viewerId: string) => { + const viewer = viewersRef.current.get(viewerId); + if (viewer) { + console.log('[WebRTCHost] Removing viewer:', viewerId); + viewer.peerConnection.close(); + viewersRef.current.delete(viewerId); + publishedStreamSendersRef.current.delete(viewerId); + pendingCandidatesRef.current.delete(viewerId); + setViewers(new Map(viewersRef.current)); + if (mountedRef.current) { + onViewerLeftRef.current?.(viewerId); } - }, - [onViewerLeft] - ); + } + }, []); removeViewerRef.current = removeViewer; // Handle viewer joining const handleViewerJoin = useCallback( - async (viewerId: string) => { + async (viewerId: string, generation: number) => { + if (!isCurrentGeneration(generation)) return; if (viewerId === hostId) return; if (viewersRef.current.has(viewerId)) return; console.log('[WebRTCHost] Viewer joining:', viewerId); - const pc = createPeerConnection(viewerId); + const pc = createPeerConnection(viewerId, generation); const viewer: ViewerConnection = { id: viewerId, @@ -519,13 +599,15 @@ export function useWebRTCHost({ viewersRef.current.set(viewerId, viewer); setViewers(new Map(viewersRef.current)); - onViewerJoined?.(viewerId); + onViewerJoinedRef.current?.(viewerId); try { const offer = (await pc.createOffer({})) as OfferAnswer; + if (!isCurrentGeneration(generation) || viewersRef.current.get(viewerId) !== viewer) return; // In-band FEC turns a lost packet into a duller syllable, not a gap. if (offer.sdp) offer.sdp = tuneOpusForVoice(offer.sdp); await pc.setLocalDescription(offer); + if (!isCurrentGeneration(generation) || viewersRef.current.get(viewerId) !== viewer) return; if (offer.sdp) { await sendSignal({ @@ -538,15 +620,18 @@ export function useWebRTCHost({ } } catch (err) { console.error('[WebRTCHost] Failed to create offer:', err); - removeViewer(viewerId); + if (isCurrentGeneration(generation)) { + removeViewer(viewerId); + } } }, - [hostId, createPeerConnection, onViewerJoined, removeViewer, sendSignal] + [hostId, createPeerConnection, isCurrentGeneration, removeViewer, sendSignal] ); // Handle incoming signals const handleSignalMessage = useCallback( - async (signal: SignalMessage) => { + async (signal: SignalMessage, generation: number) => { + if (!isCurrentGeneration(generation)) return; if (signal.targetId && signal.targetId !== hostId) return; const viewerId = signal.senderId; @@ -565,6 +650,9 @@ export function useWebRTCHost({ type: 'answer', sdp: signal.sdp, }); + if (!isCurrentGeneration(generation) || viewersRef.current.get(viewerId) !== viewer) { + return; + } // Drain buffered ICE candidates const pending = pendingCandidatesRef.current.get(viewerId); @@ -572,6 +660,12 @@ export function useWebRTCHost({ pendingCandidatesRef.current.delete(viewerId); for (const candidate of pending) { await viewer.peerConnection.addIceCandidate(new RTCIceCandidate(candidate)); + if ( + !isCurrentGeneration(generation) || + viewersRef.current.get(viewerId) !== viewer + ) { + return; + } } } } @@ -593,7 +687,7 @@ export function useWebRTCHost({ } } }, - [hostId] + [hostId, isCurrentGeneration] ); // Toggle host microphone @@ -604,16 +698,28 @@ export function useWebRTCHost({ const tracks = micStream.getAudioTracks(); if (tracks.length === 0) return; - const newEnabled = !micEnabled; + const newEnabled = !micEnabledIntentRef.current; + micEnabledIntentRef.current = newEnabled; tracks.forEach((track) => { track.enabled = newEnabled; }); setMicEnabled(newEnabled); - }, [micEnabled]); + }, []); // Start hosting const startHosting = useCallback(async () => { - if (isStartingRef.current || eventSourceRef.current) return; + if (!mountedRef.current) { + return; + } + if (appStateRef.current !== 'active') { + resumeOnActiveRef.current = true; + return; + } + if (isStartingRef.current || eventSourceRef.current) { + return; + } + + const generation = ++generationRef.current; isStartingRef.current = true; console.log('[WebRTCHost] Starting hosting for session:', sessionId); @@ -621,6 +727,7 @@ export function useWebRTCHost({ // Get auth token from secure storage try { const stored = await getStoredAuth(); + if (!isCurrentGeneration(generation)) return; if (!stored || isAuthExpired(stored)) { isStartingRef.current = false; setError('Not authenticated. Please log in again.'); @@ -629,19 +736,32 @@ export function useWebRTCHost({ authTokenRef.current = stored.accessToken; } catch (err) { console.error('[WebRTCHost] Failed to get auth token:', err); - isStartingRef.current = false; - setError('Failed to authenticate. Please log in again.'); + if (isCurrentGeneration(generation)) { + isStartingRef.current = false; + setError('Failed to authenticate. Please log in again.'); + } return; } // Capture host microphone try { const micStream = await mediaDevices.getUserMedia(voiceCaptureConstraints); - markTrackAsSpeech(micStream.getAudioTracks()[0]); + if (!isCurrentGeneration(generation)) { + micStream.getTracks().forEach((track) => { + track.stop(); + }); + return; + } + const audioTrack = micStream.getAudioTracks()[0]; + markTrackAsSpeech(audioTrack); + micStream.getAudioTracks().forEach((track) => { + track.enabled = micEnabledIntentRef.current; + }); hostMicStreamRef.current = micStream; setHasMic(true); - setMicEnabled(true); + setMicEnabled(micEnabledIntentRef.current); } catch { + if (!isCurrentGeneration(generation)) return; console.warn('[WebRTCHost] No microphone available'); setHasMic(false); setMicEnabled(false); @@ -657,11 +777,20 @@ export function useWebRTCHost({ const sseUrl = `${API_BASE_URL}/api/sessions/${sessionId}/signal/stream?${sseParams.toString()}`; const eventSource = createEventSource(sseUrl); + if (!isCurrentGeneration(generation)) { + eventSource.close(); + return; + } eventSourceRef.current = eventSource; + lastHeartbeatAtRef.current = Date.now(); + + const isCurrentEventSource = () => + isCurrentGeneration(generation) && eventSourceRef.current === eventSource; eventSource.addEventListener('connected', (event) => { - if (eventSourceRef.current !== eventSource) return; + if (!isCurrentEventSource()) return; console.log('[WebRTCHost] SSE connected'); + lastHeartbeatAtRef.current = Date.now(); isStartingRef.current = false; setIsHosting(true); setError(null); @@ -676,25 +805,30 @@ export function useWebRTCHost({ } }); + eventSource.addEventListener('heartbeat', () => { + if (!isCurrentEventSource()) return; + lastHeartbeatAtRef.current = Date.now(); + }); + eventSource.addEventListener('signal', (event) => { - if (eventSourceRef.current !== eventSource) return; + if (!isCurrentEventSource()) return; try { const signal = JSON.parse(event.data) as SignalMessage; - void handleSignalMessage(signal); + void handleSignalMessage(signal, generation); } catch (err) { console.error('[WebRTCHost] Failed to parse signal:', err); } }); eventSource.addEventListener('presence-join', (event) => { - if (eventSourceRef.current !== eventSource) return; + if (!isCurrentEventSource()) return; try { const { presences } = JSON.parse(event.data) as { presences: { user_id: string; role: string }[]; }; for (const presence of presences) { if (presence.role === 'viewer' && presence.user_id !== hostId) { - void handleViewerJoin(presence.user_id); + void handleViewerJoin(presence.user_id, generation); } } } catch (err) { @@ -703,7 +837,7 @@ export function useWebRTCHost({ }); eventSource.addEventListener('presence-leave', (event) => { - if (eventSourceRef.current !== eventSource) return; + if (!isCurrentEventSource()) return; try { const { presences } = JSON.parse(event.data) as { presences: { user_id: string }[]; @@ -717,7 +851,7 @@ export function useWebRTCHost({ }); eventSource.addEventListener('error', () => { - if (eventSourceRef.current !== eventSource) return; + if (!isCurrentEventSource()) return; console.error('[WebRTCHost] SSE error'); isStartingRef.current = false; setError('Connection to server lost. Reconnecting...'); @@ -725,13 +859,32 @@ export function useWebRTCHost({ // Start stats reporting statsIntervalRef.current = setInterval(() => { - void reportStats(); + void reportStats(generation); }, STATS_INTERVAL); - }, [sessionId, hostId, handleSignalMessage, handleViewerJoin, removeViewer, reportStats]); + heartbeatWatchdogRef.current = setInterval(() => { + if (!isCurrentEventSource()) return; + if (Date.now() - lastHeartbeatAtRef.current <= HEARTBEAT_TIMEOUT) return; + + setError('Connection heartbeat timed out. Reconnecting...'); + stopHostingRef.current?.(); + void startHostingRef.current?.(); + }, HEARTBEAT_WATCHDOG_INTERVAL); + }, [ + sessionId, + hostId, + handleSignalMessage, + handleViewerJoin, + isCurrentGeneration, + removeViewer, + reportStats, + ]); + + startHostingRef.current = startHosting; // Stop hosting const stopHosting = useCallback(() => { console.log('[WebRTCHost] Stopping hosting'); + generationRef.current += 1; isStartingRef.current = false; if (statsIntervalRef.current) { @@ -739,6 +892,11 @@ export function useWebRTCHost({ statsIntervalRef.current = null; } + if (heartbeatWatchdogRef.current) { + clearInterval(heartbeatWatchdogRef.current); + heartbeatWatchdogRef.current = null; + } + if (eventSourceRef.current) { eventSourceRef.current.close(); eventSourceRef.current = null; @@ -749,14 +907,9 @@ export function useWebRTCHost({ }); viewersRef.current.clear(); publishedStreamSendersRef.current.clear(); + pendingCandidatesRef.current.clear(); publishedStreamVersionRef.current += 1; - const publishedStream = localStreamRef.current; localStreamRef.current = null; - publishedStream?.getTracks().forEach((track) => { - track.stop(); - }); - setViewers(new Map()); - setIsHosting(false); if (hostMicStreamRef.current) { hostMicStreamRef.current.getTracks().forEach((track) => { @@ -764,10 +917,19 @@ export function useWebRTCHost({ }); hostMicStreamRef.current = null; } - setMicEnabled(false); - setHasMic(false); + lastHeartbeatAtRef.current = 0; + + if (mountedRef.current) { + setViewers(new Map()); + setControllingViewer(null); + setIsHosting(false); + setMicEnabled(false); + setHasMic(false); + } }, []); + stopHostingRef.current = stopHosting; + const removePublishedSenders = useCallback((viewer: ViewerConnection): boolean => { const publishedSenders = publishedStreamSendersRef.current.get(viewer.id); if (!publishedSenders || publishedSenders.size === 0) { @@ -784,12 +946,16 @@ export function useWebRTCHost({ // Publish screen share stream const publishStream = useCallback( async (stream: MediaStream) => { + const generation = generationRef.current; + if (!isCurrentGeneration(generation)) return; + localStreamRef.current = stream; const publishVersion = ++publishedStreamVersionRef.current; try { for (const viewer of viewersRef.current.values()) { if ( + !isCurrentGeneration(generation) || publishedStreamVersionRef.current !== publishVersion || localStreamRef.current !== stream ) { @@ -826,11 +992,12 @@ export function useWebRTCHost({ } if (negotiationNeeded) { - await renegotiateViewer(viewer); + await renegotiateViewer(viewer, generation); } } } catch (publishError) { if ( + isCurrentGeneration(generation) && publishedStreamVersionRef.current === publishVersion && localStreamRef.current === stream ) { @@ -838,9 +1005,10 @@ export function useWebRTCHost({ publishedStreamVersionRef.current += 1; for (const viewer of viewersRef.current.values()) { + if (!isCurrentGeneration(generation)) return; if (!removePublishedSenders(viewer)) continue; try { - await renegotiateViewer(viewer); + await renegotiateViewer(viewer, generation); } catch (rollbackError) { console.error( `[WebRTCHost] Failed to roll back stream for ${viewer.id}:`, @@ -852,17 +1020,23 @@ export function useWebRTCHost({ throw publishError; } }, - [removePublishedSenders, renegotiateViewer] + [isCurrentGeneration, removePublishedSenders, renegotiateViewer] ); // Unpublish stream const unpublishStream = useCallback(async () => { + const generation = generationRef.current; + if (!isCurrentGeneration(generation)) return; + localStreamRef.current = null; const unpublishVersion = ++publishedStreamVersionRef.current; const failures: string[] = []; for (const viewer of viewersRef.current.values()) { - if (publishedStreamVersionRef.current !== unpublishVersion) { + if ( + !isCurrentGeneration(generation) || + publishedStreamVersionRef.current !== unpublishVersion + ) { return; } if (viewer.connectionState !== 'connected' && viewer.connectionState !== 'connecting') { @@ -871,7 +1045,7 @@ export function useWebRTCHost({ try { if (removePublishedSenders(viewer)) { - await renegotiateViewer(viewer); + await renegotiateViewer(viewer, generation); } } catch (unpublishError) { failures.push(viewer.id); @@ -882,11 +1056,40 @@ export function useWebRTCHost({ if (failures.length > 0) { throw new Error(`Failed to unpublish screen share from ${String(failures.length)} viewer(s)`); } - }, [removePublishedSenders, renegotiateViewer]); + }, [isCurrentGeneration, removePublishedSenders, renegotiateViewer]); - // Cleanup on unmount + // Suspend native media and transport while backgrounded, then restore one + // fresh host connection only if hosting was active or starting beforehand. useEffect(() => { + mountedRef.current = true; + appStateRef.current = AppState.currentState; + + const subscription = AppState.addEventListener('change', (nextState) => { + const previousState = appStateRef.current; + appStateRef.current = nextState; + + if (nextState === 'background') { + if (previousState !== 'background') { + const hadActiveHost = + isStartingRef.current || eventSourceRef.current !== null || viewersRef.current.size > 0; + resumeOnActiveRef.current = resumeOnActiveRef.current || hadActiveHost; + if (hadActiveHost) { + stopHosting(); + } + } + return; + } + + if (nextState === 'active' && previousState !== 'active' && resumeOnActiveRef.current) { + resumeOnActiveRef.current = false; + void startHostingRef.current?.(); + } + }); + return () => { + subscription.remove(); + mountedRef.current = false; + resumeOnActiveRef.current = false; stopHosting(); }; }, [stopHosting]); @@ -953,14 +1156,9 @@ export function useWebRTCHost({ if (controllingViewer === viewerId) { setControllingViewer(null); } - - viewer.peerConnection.close(); - viewersRef.current.delete(viewerId); - publishedStreamSendersRef.current.delete(viewerId); - setViewers(new Map(viewersRef.current)); - onViewerLeft?.(viewerId); + removeViewer(viewerId); }, - [controllingViewer, onViewerLeft] + [controllingViewer, removeViewer] ); // Mute/unmute viewer diff --git a/apps/mobile/src/hooks/useWebRTCViewer.test.ts b/apps/mobile/src/hooks/useWebRTCViewer.test.ts index dd594d7e..5571b1cd 100644 --- a/apps/mobile/src/hooks/useWebRTCViewer.test.ts +++ b/apps/mobile/src/hooks/useWebRTCViewer.test.ts @@ -1,8 +1,11 @@ import { describe, it, expect, vi, beforeEach } from 'vitest'; -import { renderHook, act } from '@testing-library/react'; +import { renderHook, act, waitFor } from '@testing-library/react'; import { useWebRTCViewer } from './useWebRTCViewer'; import { createEventSource } from '../lib/event-source'; import { getStoredAuth } from '../lib/secure-storage'; +import { mediaDevices } from 'react-native-webrtc'; +import type { MediaStream, MediaStreamTrack } from 'react-native-webrtc'; +import { emitAppStateChange, mockPeerConnections } from '../test/setup'; vi.mock('../config', () => ({ API_BASE_URL: 'https://pairux.com', @@ -18,21 +21,37 @@ vi.mock('../lib/secure-storage', () => ({ isAuthExpired: vi.fn().mockReturnValue(false), })); -const { mockClose, mockAddEventListener } = vi.hoisted(() => ({ +const { mockClose, mockAddEventListener, mockEventSources } = vi.hoisted(() => ({ mockClose: vi.fn(), mockAddEventListener: vi.fn(), + mockEventSources: [] as { + listeners: Map void>; + close: ReturnType; + }[], })); vi.mock('../lib/event-source', () => ({ - createEventSource: vi.fn(() => ({ - addEventListener: mockAddEventListener, - close: mockClose, - })), + createEventSource: vi.fn(() => { + const listeners = new Map void>(); + const source = { + listeners, + addEventListener: vi.fn((event: string, listener: (payload: { data: string }) => void) => { + mockAddEventListener(event, listener); + listeners.set(event, listener); + }), + close: vi.fn(() => { + mockClose(); + }), + }; + mockEventSources.push(source); + return source; + }), })); describe('useWebRTCViewer', () => { beforeEach(() => { vi.clearAllMocks(); + mockEventSources.length = 0; vi.mocked(fetch).mockResolvedValue({ ok: true, text: async () => 'ok', @@ -247,4 +266,403 @@ describe('useWebRTCViewer', () => { unmount(); expect(mockClose).toHaveBeenCalled(); }); + + it('does not reconnect when callback identities change', async () => { + const { rerender } = renderHook( + ({ onReady }: { onReady: (stream: MediaStream) => void }) => + useWebRTCViewer({ + sessionId: 'session-1', + participantId: 'viewer-1', + onStreamReady: onReady, + }), + { initialProps: { onReady: vi.fn() } } + ); + + await act(async () => { + await Promise.resolve(); + await Promise.resolve(); + }); + rerender({ onReady: vi.fn() }); + await act(async () => { + await Promise.resolve(); + }); + + expect(createEventSource).toHaveBeenCalledTimes(1); + }); + + it('tears down once in the background and resumes exactly once when active', async () => { + renderHook(() => + useWebRTCViewer({ + sessionId: 'session-1', + participantId: 'viewer-1', + }) + ); + await act(async () => { + await Promise.resolve(); + await Promise.resolve(); + }); + expect(createEventSource).toHaveBeenCalledTimes(1); + + act(() => { + emitAppStateChange('inactive'); + }); + expect(mockClose).not.toHaveBeenCalled(); + + act(() => { + emitAppStateChange('background'); + }); + expect(mockClose).toHaveBeenCalledTimes(1); + + await act(async () => { + emitAppStateChange('active'); + await Promise.resolve(); + await Promise.resolve(); + }); + + expect(createEventSource).toHaveBeenCalledTimes(2); + }); + + it('handles an offer during a transient inactive window', async () => { + renderHook(() => + useWebRTCViewer({ + sessionId: 'session-1', + participantId: 'viewer-1', + }) + ); + await act(async () => { + await Promise.resolve(); + await Promise.resolve(); + }); + const source = mockEventSources[0]; + expect(source).toBeDefined(); + + act(() => { + source.listeners.get('connected')?.({ data: '{}' }); + emitAppStateChange('inactive'); + }); + const peer = mockPeerConnections[0]; + expect(peer).toBeDefined(); + + await act(async () => { + source.listeners.get('signal')?.({ + data: JSON.stringify({ + type: 'offer', + sdp: 'mobile-host-offer', + senderId: 'host-1', + negotiationId: 'inactive-offer-1', + timestamp: Date.now(), + }), + }); + await Promise.resolve(); + await Promise.resolve(); + await Promise.resolve(); + }); + + expect(peer.setRemoteDescription).toHaveBeenCalledTimes(1); + await waitFor(() => + expect(fetch).toHaveBeenCalledWith( + 'https://pairux.com/api/sessions/session-1/signal', + expect.objectContaining({ + method: 'POST', + body: expect.stringContaining('"negotiationId":"inactive-offer-1"'), + }) + ) + ); + expect(source.close).not.toHaveBeenCalled(); + }); + + it('keeps heartbeats current throughout a long inactive window', async () => { + vi.useFakeTimers(); + vi.setSystemTime(new Date('2026-08-28T00:00:00Z')); + + try { + renderHook(() => + useWebRTCViewer({ + sessionId: 'session-1', + participantId: 'viewer-1', + }) + ); + await act(async () => { + await Promise.resolve(); + await Promise.resolve(); + }); + const source = mockEventSources[0]; + expect(source).toBeDefined(); + + act(() => { + source.listeners.get('connected')?.({ data: '{}' }); + emitAppStateChange('inactive'); + }); + + for (let elapsed = 10_000; elapsed <= 80_000; elapsed += 10_000) { + await act(async () => { + await vi.advanceTimersByTimeAsync(10_000); + source.listeners.get('heartbeat')?.({ data: '{}' }); + }); + } + + await act(async () => { + emitAppStateChange('active'); + await vi.advanceTimersByTimeAsync(16_000); + }); + + expect(source.close).not.toHaveBeenCalled(); + expect(createEventSource).toHaveBeenCalledTimes(1); + } finally { + vi.useRealTimers(); + } + }); + + it('reconnects on foreground after the watchdog expires while inactive', async () => { + vi.useFakeTimers(); + vi.setSystemTime(new Date('2026-08-28T00:00:00Z')); + + try { + renderHook(() => + useWebRTCViewer({ + sessionId: 'session-1', + participantId: 'viewer-1', + }) + ); + await act(async () => { + await Promise.resolve(); + await Promise.resolve(); + }); + const source = mockEventSources[0]; + expect(source).toBeDefined(); + + act(() => { + source.listeners.get('connected')?.({ data: '{}' }); + emitAppStateChange('inactive'); + }); + + await act(async () => { + await vi.advanceTimersByTimeAsync(90_000); + }); + + expect(source.close).toHaveBeenCalledTimes(1); + expect(createEventSource).toHaveBeenCalledTimes(1); + + await act(async () => { + emitAppStateChange('active'); + await Promise.resolve(); + await Promise.resolve(); + }); + + expect(createEventSource).toHaveBeenCalledTimes(2); + } finally { + vi.useRealTimers(); + } + }); + + it('ignores connected callbacks from a retired background generation', async () => { + renderHook(() => + useWebRTCViewer({ + sessionId: 'session-1', + participantId: 'viewer-1', + }) + ); + await act(async () => { + await Promise.resolve(); + await Promise.resolve(); + }); + const staleConnectedListener = mockEventSources[0]?.listeners.get('connected'); + + act(() => { + emitAppStateChange('background'); + }); + await act(async () => { + emitAppStateChange('active'); + await Promise.resolve(); + await Promise.resolve(); + }); + const currentConnectedListener = mockEventSources[1]?.listeners.get('connected'); + + act(() => { + staleConnectedListener?.({ data: '{}' }); + }); + expect(mockPeerConnections).toHaveLength(0); + + act(() => { + currentConnectedListener?.({ data: '{}' }); + }); + expect(mockPeerConnections).toHaveLength(1); + }); + + it('replaces the remote stream when the same SSE connection reconnects', async () => { + const onStreamReady = vi.fn(); + const onStreamEnded = vi.fn(); + const { result } = renderHook(() => + useWebRTCViewer({ + sessionId: 'session-1', + participantId: 'viewer-1', + onStreamReady, + onStreamEnded, + }) + ); + await act(async () => { + await Promise.resolve(); + await Promise.resolve(); + }); + + const connectedListener = mockEventSources[0]?.listeners.get('connected'); + const makeRemoteStream = (id: string): MediaStream => + ({ + id, + getTracks: () => [{ kind: 'video' }], + getAudioTracks: () => [], + getVideoTracks: () => [{ kind: 'video' }], + }) as unknown as MediaStream; + const firstStream = makeRemoteStream('first'); + const secondStream = makeRemoteStream('second'); + + act(() => { + connectedListener?.({ data: '{}' }); + }); + const firstPeer = mockPeerConnections[0]; + const firstTrackListener = firstPeer.addEventListener.mock.calls.find( + ([eventName]) => eventName === 'track' + )?.[1] as ((event: { streams: MediaStream[] }) => void) | undefined; + act(() => { + firstTrackListener?.({ streams: [firstStream] }); + }); + expect(result.current.remoteStream).toBe(firstStream); + + act(() => { + connectedListener?.({ data: '{}' }); + }); + expect(firstPeer.close).toHaveBeenCalledTimes(1); + expect(result.current.remoteStream).toBeNull(); + expect(onStreamEnded).toHaveBeenCalledTimes(1); + + const secondPeer = mockPeerConnections[1]; + const secondTrackListener = secondPeer.addEventListener.mock.calls.find( + ([eventName]) => eventName === 'track' + )?.[1] as ((event: { streams: MediaStream[] }) => void) | undefined; + act(() => { + secondTrackListener?.({ streams: [secondStream] }); + }); + + expect(result.current.remoteStream).toBe(secondStream); + expect(onStreamReady).toHaveBeenNthCalledWith(1, firstStream); + expect(onStreamReady).toHaveBeenNthCalledWith(2, secondStream); + }); + + it('preserves the viewer microphone mute intent across a foreground reconnect', async () => { + const tracks: (MediaStreamTrack & { + enabled: boolean; + stop: ReturnType; + })[] = []; + const makeMicStream = (): MediaStream => { + const track = { + kind: 'audio', + enabled: true, + stop: vi.fn(), + } as unknown as MediaStreamTrack & { + enabled: boolean; + stop: ReturnType; + }; + tracks.push(track); + return { + getTracks: () => [track], + getAudioTracks: () => [track], + } as unknown as MediaStream; + }; + vi.mocked(mediaDevices.getUserMedia) + .mockResolvedValueOnce(makeMicStream()) + .mockResolvedValueOnce(makeMicStream()); + + const { result } = renderHook(() => + useWebRTCViewer({ + sessionId: 'session-1', + participantId: 'viewer-1', + }) + ); + await act(async () => { + await Promise.resolve(); + await Promise.resolve(); + }); + act(() => { + result.current.toggleMic(); + }); + expect(tracks[0]?.enabled).toBe(false); + + act(() => { + emitAppStateChange('background'); + }); + await act(async () => { + emitAppStateChange('active'); + await Promise.resolve(); + await Promise.resolve(); + }); + + expect(tracks[1]?.enabled).toBe(false); + expect(result.current.micEnabled).toBe(false); + }); + + it('stops a late microphone stream when unmounted during initialization', async () => { + let resolveMic!: (stream: MediaStream) => void; + const lateTrack = { + kind: 'audio', + enabled: true, + stop: vi.fn(), + } as unknown as MediaStreamTrack & { stop: ReturnType }; + const lateStream = { + getTracks: () => [lateTrack], + getAudioTracks: () => [lateTrack], + } as unknown as MediaStream; + vi.mocked(mediaDevices.getUserMedia).mockReturnValueOnce( + new Promise((resolve) => { + resolveMic = resolve; + }) + ); + + const { unmount } = renderHook(() => + useWebRTCViewer({ + sessionId: 'session-1', + participantId: 'viewer-1', + }) + ); + await act(async () => { + await Promise.resolve(); + }); + unmount(); + + await act(async () => { + resolveMic(lateStream); + await Promise.resolve(); + await Promise.resolve(); + }); + + expect(lateTrack.stop).toHaveBeenCalledTimes(1); + expect(createEventSource).not.toHaveBeenCalled(); + }); + + it('restarts once after the SSE heartbeat expires', async () => { + vi.useFakeTimers(); + vi.setSystemTime(new Date('2026-08-28T00:00:00Z')); + + try { + renderHook(() => + useWebRTCViewer({ + sessionId: 'session-1', + participantId: 'viewer-1', + }) + ); + await act(async () => { + await Promise.resolve(); + await Promise.resolve(); + }); + + await act(async () => { + await vi.advanceTimersByTimeAsync(90_001); + await Promise.resolve(); + }); + + expect(mockClose).toHaveBeenCalledTimes(1); + expect(createEventSource).toHaveBeenCalledTimes(2); + } finally { + vi.useRealTimers(); + } + }); }); diff --git a/apps/mobile/src/hooks/useWebRTCViewer.ts b/apps/mobile/src/hooks/useWebRTCViewer.ts index e15bdaa4..51789c74 100644 --- a/apps/mobile/src/hooks/useWebRTCViewer.ts +++ b/apps/mobile/src/hooks/useWebRTCViewer.ts @@ -8,6 +8,7 @@ * - RTCPeerConnection from react-native-webrtc */ import { useState, useEffect, useRef, useCallback } from 'react'; +import { AppState, type AppStateStatus } from 'react-native'; import type { MediaStream } from 'react-native-webrtc'; import { RTCPeerConnection, RTCIceCandidate, mediaDevices } from 'react-native-webrtc'; import type { @@ -51,6 +52,8 @@ interface OfferAnswer { const STATS_INTERVAL = 30000; const STATS_DISPLAY_INTERVAL = 2000; +const HEARTBEAT_TIMEOUT = 75000; +const HEARTBEAT_WATCHDOG_INTERVAL = 15000; const DEFAULT_ICE_SERVERS: RTCIceServer[] = [ { urls: 'stun:stun.l.google.com:19302' }, @@ -120,22 +123,42 @@ export function useWebRTCViewer({ const authTokenRef = useRef(null); const statsIntervalRef = useRef | null>(null); const statsReportIntervalRef = useRef | null>(null); + const heartbeatWatchdogRef = useRef | null>(null); + const lastHeartbeatAtRef = useRef(0); const reconnectAttemptsRef = useRef(0); const inputSequenceRef = useRef(0); const isConnectingRef = useRef(false); + const connectionFailureInFlightRef = useRef(false); + const mountedRef = useRef(true); + const generationRef = useRef(0); + const appStateRef = useRef(AppState.currentState); + const resumeOnActiveRef = useRef(false); + const micEnabledIntentRef = useRef(true); const maxReconnectAttempts = 3; const iceServersRef = useRef(DEFAULT_ICE_SERVERS); const pendingCandidatesRef = useRef([]); const signalQueueRef = useRef>(Promise.resolve()); - const handleConnectionFailureRef = useRef<(() => Promise) | undefined>(undefined); + const handleConnectionFailureRef = useRef< + ((generation: number, pc: RTCPeerConnection) => Promise) | undefined + >(undefined); const onControlStateChangeRef = useRef(onControlStateChange); const onKickedRef = useRef(onKicked); + const onStreamReadyRef = useRef(onStreamReady); + const onStreamEndedRef = useRef(onStreamEnded); const disconnectRef = useRef<(() => void) | undefined>(undefined); + const initializeRef = useRef<(() => Promise) | undefined>(undefined); onControlStateChangeRef.current = onControlStateChange; onKickedRef.current = onKicked; + onStreamReadyRef.current = onStreamReady; + onStreamEndedRef.current = onStreamEnded; + + const isCurrentGeneration = useCallback( + (generation: number) => mountedRef.current && generationRef.current === generation, + [] + ); const calculateNetworkQuality = useCallback((metrics: QualityMetrics): NetworkQuality => { const { packetLoss, roundTripTime } = metrics; @@ -191,6 +214,7 @@ export function useWebRTCViewer({ onKickedRef.current?.(message.reason); break; case 'mute': { + micEnabledIntentRef.current = !message.muted; const micStream = micStreamRef.current; if (micStream) { micStream.getAudioTracks().forEach((track) => { @@ -209,27 +233,34 @@ export function useWebRTCViewer({ // Setup data channel const setupDataChannel = useCallback( - (channel: DataChannel) => { + (channel: DataChannel, generation: number) => { dataChannelRef.current = channel; + const isCurrentChannel = () => + isCurrentGeneration(generation) && dataChannelRef.current === channel; + channel.addEventListener('open', () => { + if (!isCurrentChannel()) return; setDataChannelReady(true); }); channel.addEventListener('close', () => { + if (!isCurrentChannel()) return; setDataChannelReady(false); setControlState('view-only'); }); channel.addEventListener('error', () => { + if (!isCurrentChannel()) return; setDataChannelReady(false); }); channel.addEventListener('message', (event) => { + if (!isCurrentChannel()) return; handleDataChannelMessage(typeof event.data === 'string' ? event.data : ''); }); }, - [handleDataChannelMessage] + [handleDataChannelMessage, isCurrentGeneration] ); // Control actions @@ -280,117 +311,128 @@ export function useWebRTCViewer({ ); // Collect stats for UI - const collectStats = useCallback(async () => { - const pc = peerConnectionRef.current; - if (!pc?.connectionState || pc.connectionState !== 'connected') return; + const collectStats = useCallback( + async (generation: number) => { + if (!isCurrentGeneration(generation)) return; + const pc = peerConnectionRef.current; + if (!pc?.connectionState || pc.connectionState !== 'connected') return; - try { - const stats = (await pc.getStats()) as Map>; - let bitrate = 0; - let frameRate = 0; - let packetLoss = 0; - let roundTripTime = 0; - let bytesReceived = 0; - let packetsLost = 0; - let packetsReceived = 0; - - stats.forEach((report: Record) => { - if (report.type === 'inbound-rtp' && report.kind === 'video') { - bytesReceived = (report.bytesReceived as number | undefined) ?? 0; - frameRate = (report.framesPerSecond as number | undefined) ?? 0; - packetsLost = (report.packetsLost as number | undefined) ?? 0; - packetsReceived = (report.packetsReceived as number | undefined) ?? 0; - } - if (report.type === 'candidate-pair' && report.state === 'succeeded') { - roundTripTime = ((report.currentRoundTripTime as number | undefined) ?? 0) * 1000; - } - }); + try { + const stats = (await pc.getStats()) as Map>; + if (!isCurrentGeneration(generation) || peerConnectionRef.current !== pc) return; + let bitrate = 0; + let frameRate = 0; + let packetLoss = 0; + let roundTripTime = 0; + let bytesReceived = 0; + let packetsLost = 0; + let packetsReceived = 0; + + stats.forEach((report: Record) => { + if (report.type === 'inbound-rtp' && report.kind === 'video') { + bytesReceived = (report.bytesReceived as number | undefined) ?? 0; + frameRate = (report.framesPerSecond as number | undefined) ?? 0; + packetsLost = (report.packetsLost as number | undefined) ?? 0; + packetsReceived = (report.packetsReceived as number | undefined) ?? 0; + } + if (report.type === 'candidate-pair' && report.state === 'succeeded') { + roundTripTime = ((report.currentRoundTripTime as number | undefined) ?? 0) * 1000; + } + }); - if (packetsReceived > 0) { - packetLoss = (packetsLost / (packetsReceived + packetsLost)) * 100; - } + if (packetsReceived > 0) { + packetLoss = (packetsLost / (packetsReceived + packetsLost)) * 100; + } - bitrate = bytesReceived * 8; + bitrate = bytesReceived * 8; - const metrics: QualityMetrics = { bitrate, frameRate, packetLoss, roundTripTime }; - setQualityMetrics(metrics); - setNetworkQuality(calculateNetworkQuality(metrics)); - } catch { - // Non-critical - } - }, [calculateNetworkQuality]); + const metrics: QualityMetrics = { bitrate, frameRate, packetLoss, roundTripTime }; + setQualityMetrics(metrics); + setNetworkQuality(calculateNetworkQuality(metrics)); + } catch { + // Non-critical + } + }, + [calculateNetworkQuality, isCurrentGeneration] + ); // Report stats to API - const reportStats = useCallback(async () => { - const pc = peerConnectionRef.current; - if (pc?.connectionState !== 'connected') return; + const reportStats = useCallback( + async (generation: number) => { + if (!isCurrentGeneration(generation)) return; + const pc = peerConnectionRef.current; + if (pc?.connectionState !== 'connected') return; - try { - const stats = (await pc.getStats()) as Map>; - let bytesSent = 0; - let bytesReceived = 0; - let packetsSent = 0; - let packetsReceived = 0; - let packetsLost = 0; - let roundTripTime: number | undefined; - let frameRate: number | undefined; - let frameWidth: number | undefined; - let frameHeight: number | undefined; - - stats.forEach((report: Record) => { - if (report.type === 'inbound-rtp' && report.kind === 'video') { - bytesReceived += (report.bytesReceived as number | undefined) ?? 0; - packetsReceived += (report.packetsReceived as number | undefined) ?? 0; - frameRate = report.framesPerSecond as number | undefined; - frameWidth = report.frameWidth as number | undefined; - frameHeight = report.frameHeight as number | undefined; - } - if (report.type === 'outbound-rtp') { - bytesSent += (report.bytesSent as number | undefined) ?? 0; - packetsSent += (report.packetsSent as number | undefined) ?? 0; - } - if (report.type === 'remote-inbound-rtp') { - packetsLost += (report.packetsLost as number | undefined) ?? 0; - } - if (report.type === 'candidate-pair' && report.state === 'succeeded') { - roundTripTime = ((report.currentRoundTripTime as number | undefined) ?? 0) * 1000; + try { + const stats = (await pc.getStats()) as Map>; + if (!isCurrentGeneration(generation) || peerConnectionRef.current !== pc) return; + let bytesSent = 0; + let bytesReceived = 0; + let packetsSent = 0; + let packetsReceived = 0; + let packetsLost = 0; + let roundTripTime: number | undefined; + let frameRate: number | undefined; + let frameWidth: number | undefined; + let frameHeight: number | undefined; + + stats.forEach((report: Record) => { + if (report.type === 'inbound-rtp' && report.kind === 'video') { + bytesReceived += (report.bytesReceived as number | undefined) ?? 0; + packetsReceived += (report.packetsReceived as number | undefined) ?? 0; + frameRate = report.framesPerSecond as number | undefined; + frameWidth = report.frameWidth as number | undefined; + frameHeight = report.frameHeight as number | undefined; + } + if (report.type === 'outbound-rtp') { + bytesSent += (report.bytesSent as number | undefined) ?? 0; + packetsSent += (report.packetsSent as number | undefined) ?? 0; + } + if (report.type === 'remote-inbound-rtp') { + packetsLost += (report.packetsLost as number | undefined) ?? 0; + } + if (report.type === 'candidate-pair' && report.state === 'succeeded') { + roundTripTime = ((report.currentRoundTripTime as number | undefined) ?? 0) * 1000; + } + }); + + const headers: Record = { 'Content-Type': 'application/json' }; + if (authTokenRef.current) { + headers.Authorization = `Bearer ${authTokenRef.current}`; } - }); - const headers: Record = { 'Content-Type': 'application/json' }; - if (authTokenRef.current) { - headers.Authorization = `Bearer ${authTokenRef.current}`; + await fetch(`${API_BASE_URL}/api/sessions/${sessionId}/stats`, { + method: 'POST', + headers, + body: JSON.stringify({ + participantId, + role: 'viewer', + timestamp: Date.now(), + connectionState: pc.connectionState, + bytesSent, + bytesReceived, + packetsSent, + packetsReceived, + packetsLost, + roundTripTime, + jitter: 0, + frameRate, + frameWidth, + frameHeight, + reportInterval: STATS_INTERVAL, + }), + }); + } catch (err) { + console.error('[WebRTCViewer] Stats report error:', err); } - - await fetch(`${API_BASE_URL}/api/sessions/${sessionId}/stats`, { - method: 'POST', - headers, - body: JSON.stringify({ - participantId, - role: 'viewer', - timestamp: Date.now(), - connectionState: pc.connectionState, - bytesSent, - bytesReceived, - packetsSent, - packetsReceived, - packetsLost, - roundTripTime, - jitter: 0, - frameRate, - frameWidth, - frameHeight, - reportInterval: STATS_INTERVAL, - }), - }); - } catch (err) { - console.error('[WebRTCViewer] Stats report error:', err); - } - }, [sessionId, participantId]); + }, + [isCurrentGeneration, sessionId, participantId] + ); // Process a single signaling message const processSignalMessage = useCallback( - async (message: SignalMessage) => { + async (message: SignalMessage, generation: number) => { + if (!isCurrentGeneration(generation)) return; const pc = peerConnectionRef.current; if (!pc) return; @@ -406,10 +448,12 @@ export function useWebRTCViewer({ if (!message.sdp) break; await pc.setRemoteDescription({ type: 'offer', sdp: message.sdp }); + if (!isCurrentGeneration(generation) || peerConnectionRef.current !== pc) return; const answer = (await pc.createAnswer()) as OfferAnswer; // In-band FEC turns a lost packet into a duller syllable, not a gap. if (answer.sdp) answer.sdp = tuneOpusForVoice(answer.sdp); await pc.setLocalDescription(answer); + if (!isCurrentGeneration(generation) || peerConnectionRef.current !== pc) return; if (answer.sdp) { await sendSignal({ @@ -428,6 +472,7 @@ export function useWebRTCViewer({ pendingCandidatesRef.current = []; for (const candidate of pending) { await pc.addIceCandidate(new RTCIceCandidate(candidate)); + if (!isCurrentGeneration(generation) || peerConnectionRef.current !== pc) return; } } break; @@ -446,143 +491,175 @@ export function useWebRTCViewer({ } } catch (err) { console.error('[WebRTCViewer] Error handling signal message:', err); - setError('Failed to process signaling message'); + if (isCurrentGeneration(generation)) { + setError('Failed to process signaling message'); + } } }, - [participantId, sendSignal] + [isCurrentGeneration, participantId, sendSignal] ); // Serialize signal processing const handleSignalMessage = useCallback( - (message: SignalMessage) => { - signalQueueRef.current = signalQueueRef.current.then(() => processSignalMessage(message)); + (message: SignalMessage, generation: number) => { + signalQueueRef.current = signalQueueRef.current.then(() => + processSignalMessage(message, generation) + ); }, [processSignalMessage] ); // Handle connection failure with retry - const handleConnectionFailure = useCallback(async () => { - if (reconnectAttemptsRef.current < maxReconnectAttempts) { - reconnectAttemptsRef.current++; - setConnectionState('reconnecting'); - setError( - `Connection lost. Reconnecting (${String(reconnectAttemptsRef.current)}/${String(maxReconnectAttempts)})...` - ); + const handleConnectionFailure = useCallback( + async (generation: number, pc: RTCPeerConnection) => { + if ( + connectionFailureInFlightRef.current || + !isCurrentGeneration(generation) || + peerConnectionRef.current !== pc + ) { + return; + } + connectionFailureInFlightRef.current = true; - const pc = peerConnectionRef.current; - if (pc) { - try { - const offer = (await pc.createOffer({ iceRestart: true })) as OfferAnswer; - await pc.setLocalDescription(offer); - - if (offer.sdp) { - await sendSignal({ - type: 'offer', - sdp: offer.sdp, - senderId: participantId, - timestamp: Date.now(), - }); + try { + if (reconnectAttemptsRef.current < maxReconnectAttempts) { + reconnectAttemptsRef.current++; + setConnectionState('reconnecting'); + setError( + `Connection lost. Reconnecting (${String(reconnectAttemptsRef.current)}/${String(maxReconnectAttempts)})...` + ); + + try { + const offer = (await pc.createOffer({ iceRestart: true })) as OfferAnswer; + if (!isCurrentGeneration(generation) || peerConnectionRef.current !== pc) return; + await pc.setLocalDescription(offer); + if (!isCurrentGeneration(generation) || peerConnectionRef.current !== pc) return; + + if (offer.sdp) { + await sendSignal({ + type: 'offer', + sdp: offer.sdp, + senderId: participantId, + timestamp: Date.now(), + }); + } + } catch { + if (isCurrentGeneration(generation)) { + setConnectionState('failed'); + setError('Failed to reconnect'); + } } - } catch { + } else { setConnectionState('failed'); - setError('Failed to reconnect'); + setError('Connection failed after multiple attempts'); } + } finally { + connectionFailureInFlightRef.current = false; } - } else { - setConnectionState('failed'); - setError('Connection failed after multiple attempts'); - } - }, [participantId, sendSignal]); + }, + [isCurrentGeneration, participantId, sendSignal] + ); handleConnectionFailureRef.current = handleConnectionFailure; // Create peer connection - const createPeerConnection = useCallback(() => { - const pc = new RTCPeerConnection({ - iceServers: iceServersRef.current, - iceCandidatePoolSize: 10, - }); + const createPeerConnection = useCallback( + (generation: number) => { + const pc = new RTCPeerConnection({ + iceServers: iceServersRef.current, + iceCandidatePoolSize: 10, + }); - pc.addEventListener('track', (event) => { - const stream = event.streams[0] as MediaStream | undefined; - if (!stream) return; + pc.addEventListener('track', (event) => { + if (!isCurrentGeneration(generation) || peerConnectionRef.current !== pc) return; + const stream = event.streams[0] as MediaStream | undefined; + if (!stream) return; - const current = remoteStreamRef.current; - const currentHasVideo = Boolean(current && current.getVideoTracks().length > 0); - const incomingHasVideo = stream.getVideoTracks().length > 0; + const current = remoteStreamRef.current; + const currentHasVideo = Boolean(current && current.getVideoTracks().length > 0); + const incomingHasVideo = stream.getVideoTracks().length > 0; - // Keep video stream selected if a later audio-only stream arrives. - if (!current || (!currentHasVideo && incomingHasVideo)) { - remoteStreamRef.current = stream; - setRemoteStream(stream); - onStreamReady?.(stream); - } - }); + // Keep video stream selected if a later audio-only stream arrives. + if (!current || (!currentHasVideo && incomingHasVideo)) { + remoteStreamRef.current = stream; + setRemoteStream(stream); + onStreamReadyRef.current?.(stream); + } + }); - pc.addEventListener('icecandidate', (event) => { - if (event.candidate) { - void sendSignal({ - type: 'ice-candidate', - candidate: event.candidate.toJSON(), - senderId: participantId, - timestamp: Date.now(), - }); - } - }); + pc.addEventListener('icecandidate', (event) => { + if (!isCurrentGeneration(generation) || peerConnectionRef.current !== pc) return; + if (event.candidate) { + void sendSignal({ + type: 'ice-candidate', + candidate: event.candidate.toJSON(), + senderId: participantId, + timestamp: Date.now(), + }); + } + }); - pc.addEventListener('iceconnectionstatechange', () => { - const state = pc.iceConnectionState; - switch (state) { - case 'checking': - setConnectionState('connecting'); - break; - case 'connected': - case 'completed': - setConnectionState('connected'); - setError(null); - reconnectAttemptsRef.current = 0; - // Re-apply mic priority now that negotiation is done. Some stacks - // report no encodings on a sender until then, which would have made - // the call made at addTrack() time a silent no-op. - for (const sender of pc.getSenders()) { - if (sender.track?.kind === 'audio') { - void prioritizeAudioSender(sender); + pc.addEventListener('iceconnectionstatechange', () => { + if (!isCurrentGeneration(generation) || peerConnectionRef.current !== pc) return; + const state = pc.iceConnectionState; + switch (state) { + case 'checking': + setConnectionState('connecting'); + break; + case 'connected': + case 'completed': + setConnectionState('connected'); + setError(null); + reconnectAttemptsRef.current = 0; + // Re-apply mic priority now that negotiation is done. Some stacks + // report no encodings on a sender until then, which would have made + // the call made at addTrack() time a silent no-op. + for (const sender of pc.getSenders()) { + if (sender.track?.kind === 'audio') { + void prioritizeAudioSender(sender); + } } + break; + case 'disconnected': + setConnectionState('reconnecting'); + break; + case 'failed': { + const handler = handleConnectionFailureRef.current; + if (handler) void handler(generation, pc); + break; } - break; - case 'disconnected': - setConnectionState('reconnecting'); - break; - case 'failed': { - const handler = handleConnectionFailureRef.current; - if (handler) void handler(); - break; + case 'closed': + setConnectionState('disconnected'); + remoteStreamRef.current = null; + setRemoteStream(null); + onStreamEndedRef.current?.(); + break; } - case 'closed': - setConnectionState('disconnected'); - remoteStreamRef.current = null; - setRemoteStream(null); - onStreamEnded?.(); - break; - } - }); + }); - pc.addEventListener('connectionstatechange', () => { - const handler = handleConnectionFailureRef.current; - if (pc.connectionState === 'failed' && handler) { - void handler(); - } - }); + pc.addEventListener('connectionstatechange', () => { + if (!isCurrentGeneration(generation) || peerConnectionRef.current !== pc) return; + const handler = handleConnectionFailureRef.current; + if (pc.connectionState === 'failed' && handler) { + void handler(generation, pc); + } + }); - pc.addEventListener('datachannel', (event) => { - setupDataChannel(event.channel); - }); + pc.addEventListener('datachannel', (event) => { + if (!isCurrentGeneration(generation) || peerConnectionRef.current !== pc) return; + setupDataChannel(event.channel, generation); + }); - return pc; - }, [participantId, onStreamReady, onStreamEnded, sendSignal, setupDataChannel]); + return pc; + }, + [isCurrentGeneration, participantId, sendSignal, setupDataChannel] + ); // Disconnect const disconnect = useCallback(() => { + generationRef.current += 1; + const hadRemoteStream = remoteStreamRef.current !== null; + if (statsIntervalRef.current) { clearInterval(statsIntervalRef.current); statsIntervalRef.current = null; @@ -593,6 +670,11 @@ export function useWebRTCViewer({ statsReportIntervalRef.current = null; } + if (heartbeatWatchdogRef.current) { + clearInterval(heartbeatWatchdogRef.current); + heartbeatWatchdogRef.current = null; + } + if (dataChannelRef.current) { dataChannelRef.current.close(); dataChannelRef.current = null; @@ -616,14 +698,23 @@ export function useWebRTCViewer({ } isConnectingRef.current = false; + connectionFailureInFlightRef.current = false; + lastHeartbeatAtRef.current = 0; + pendingCandidatesRef.current = []; + signalQueueRef.current = Promise.resolve(); remoteStreamRef.current = null; - setRemoteStream(null); - setConnectionState('disconnected'); - setQualityMetrics(null); - setDataChannelReady(false); - setControlState('view-only'); - setMicEnabled(false); - setHasMic(false); + if (mountedRef.current) { + setRemoteStream(null); + setConnectionState('disconnected'); + setQualityMetrics(null); + setDataChannelReady(false); + setControlState('view-only'); + setMicEnabled(false); + setHasMic(false); + if (hadRemoteStream) { + onStreamEndedRef.current?.(); + } + } }, []); disconnectRef.current = disconnect; @@ -636,16 +727,28 @@ export function useWebRTCViewer({ const tracks = micStream.getAudioTracks(); if (tracks.length === 0) return; - const newEnabled = !micEnabled; + const newEnabled = !micEnabledIntentRef.current; + micEnabledIntentRef.current = newEnabled; tracks.forEach((track) => { track.enabled = newEnabled; }); setMicEnabled(newEnabled); - }, [micEnabled]); + }, []); // Initialize connection const initialize = useCallback(async () => { - if (isConnectingRef.current || eventSourceRef.current) return; + if (!mountedRef.current) { + return; + } + if (appStateRef.current !== 'active') { + resumeOnActiveRef.current = true; + return; + } + if (isConnectingRef.current || eventSourceRef.current) { + return; + } + + const generation = ++generationRef.current; isConnectingRef.current = true; console.log('[WebRTCViewer] Starting viewer for session:', sessionId); @@ -653,6 +756,7 @@ export function useWebRTCViewer({ // Get auth token from secure storage try { const stored = await getStoredAuth(); + if (!isCurrentGeneration(generation)) return; if (!stored || isAuthExpired(stored)) { isConnectingRef.current = false; setError('Not authenticated. Please log in again.'); @@ -661,19 +765,32 @@ export function useWebRTCViewer({ authTokenRef.current = stored.accessToken; } catch (err) { console.error('[WebRTCViewer] Failed to get auth token:', err); - isConnectingRef.current = false; - setError('Failed to authenticate. Please log in again.'); + if (isCurrentGeneration(generation)) { + isConnectingRef.current = false; + setError('Failed to authenticate. Please log in again.'); + } return; } // Capture microphone try { const micStream = await mediaDevices.getUserMedia(voiceCaptureConstraints); - markTrackAsSpeech(micStream.getAudioTracks()[0]); + if (!isCurrentGeneration(generation)) { + micStream.getTracks().forEach((track) => { + track.stop(); + }); + return; + } + const audioTrack = micStream.getAudioTracks()[0]; + markTrackAsSpeech(audioTrack); + micStream.getAudioTracks().forEach((track) => { + track.enabled = micEnabledIntentRef.current; + }); micStreamRef.current = micStream; setHasMic(true); - setMicEnabled(true); + setMicEnabled(micEnabledIntentRef.current); } catch { + if (!isCurrentGeneration(generation)) return; console.warn('[WebRTCViewer] Could not access microphone'); micStreamRef.current = null; setHasMic(false); @@ -688,10 +805,20 @@ export function useWebRTCViewer({ const sseUrl = `${API_BASE_URL}/api/sessions/${sessionId}/signal/stream?${sseParams.toString()}`; const eventSource = createEventSource(sseUrl); + if (!isCurrentGeneration(generation)) { + eventSource.close(); + return; + } eventSourceRef.current = eventSource; + lastHeartbeatAtRef.current = Date.now(); + + const isCurrentEventSource = () => + isCurrentGeneration(generation) && eventSourceRef.current === eventSource; eventSource.addEventListener('connected', (event) => { + if (!isCurrentEventSource()) return; console.log('[WebRTCViewer] SSE connected'); + lastHeartbeatAtRef.current = Date.now(); isConnectingRef.current = false; setConnectionState('connecting'); setError(null); @@ -705,10 +832,27 @@ export function useWebRTCViewer({ // Use default ICE servers } - // Create peer connection after receiving ICE servers + // A reconnect can emit another connected event on the same SSE object. + // Retire the old peer before attaching a replacement. + const previousPeer = peerConnectionRef.current; + const previousStream = remoteStreamRef.current; + peerConnectionRef.current = null; + remoteStreamRef.current = null; + dataChannelRef.current?.close(); + dataChannelRef.current = null; + previousPeer?.close(); + + setRemoteStream(null); + setQualityMetrics(null); + setDataChannelReady(false); + setControlState('view-only'); + if (previousStream) { + onStreamEndedRef.current?.(); + } + pendingCandidatesRef.current = []; signalQueueRef.current = Promise.resolve(); - const pc = createPeerConnection(); + const pc = createPeerConnection(generation); peerConnectionRef.current = pc; // Add mic tracks @@ -720,16 +864,23 @@ export function useWebRTCViewer({ } }); + eventSource.addEventListener('heartbeat', () => { + if (!isCurrentEventSource()) return; + lastHeartbeatAtRef.current = Date.now(); + }); + eventSource.addEventListener('signal', (event) => { + if (!isCurrentEventSource()) return; try { const signal = JSON.parse(event.data) as SignalMessage; - handleSignalMessage(signal); + handleSignalMessage(signal, generation); } catch (err) { console.error('[WebRTCViewer] Failed to parse signal:', err); } }); eventSource.addEventListener('presence-join', (event) => { + if (!isCurrentEventSource()) return; try { const { presences } = JSON.parse(event.data) as { presences: { user_id: string; role: string }[]; @@ -745,6 +896,7 @@ export function useWebRTCViewer({ }); eventSource.addEventListener('presence-leave', (event) => { + if (!isCurrentEventSource()) return; try { const { presences } = JSON.parse(event.data) as { presences: { user_id: string; role: string }[]; @@ -760,34 +912,89 @@ export function useWebRTCViewer({ }); eventSource.addEventListener('error', () => { + if (!isCurrentEventSource()) return; console.error('[WebRTCViewer] SSE error'); isConnectingRef.current = false; setError('Connection to server lost. Reconnecting...'); }); - // Start stats collection - statsIntervalRef.current = setInterval(() => void collectStats(), STATS_DISPLAY_INTERVAL); - statsReportIntervalRef.current = setInterval(() => void reportStats(), STATS_INTERVAL); + statsIntervalRef.current = setInterval( + () => void collectStats(generation), + STATS_DISPLAY_INTERVAL + ); + statsReportIntervalRef.current = setInterval( + () => void reportStats(generation), + STATS_INTERVAL + ); + heartbeatWatchdogRef.current = setInterval(() => { + if (!isCurrentEventSource()) return; + if (Date.now() - lastHeartbeatAtRef.current <= HEARTBEAT_TIMEOUT) return; + + setError('Connection heartbeat timed out. Reconnecting...'); + disconnectRef.current?.(); + void initializeRef.current?.(); + }, HEARTBEAT_WATCHDOG_INTERVAL); }, [ sessionId, participantId, + isCurrentGeneration, handleSignalMessage, createPeerConnection, collectStats, reportStats, ]); + initializeRef.current = initialize; + // Manual reconnect const reconnect = useCallback(() => { + resumeOnActiveRef.current = false; reconnectAttemptsRef.current = 0; disconnect(); - void initialize(); - }, [disconnect, initialize]); + void initializeRef.current?.(); + }, [disconnect]); - // Initialize on mount + // Suspend native media and transport while backgrounded, then restore one + // fresh connection when the app becomes active again. useEffect(() => { - void initialize(); + mountedRef.current = true; + appStateRef.current = AppState.currentState; + + if (appStateRef.current === 'active') { + void initialize(); + } else { + resumeOnActiveRef.current = true; + } + + const subscription = AppState.addEventListener('change', (nextState) => { + const previousState = appStateRef.current; + appStateRef.current = nextState; + + if (nextState === 'background') { + if (previousState !== 'background') { + const hadActiveConnection = + isConnectingRef.current || + eventSourceRef.current !== null || + peerConnectionRef.current !== null; + resumeOnActiveRef.current = resumeOnActiveRef.current || hadActiveConnection; + if (hadActiveConnection) { + disconnect(); + } + } + return; + } + + if (nextState === 'active' && previousState !== 'active' && resumeOnActiveRef.current) { + resumeOnActiveRef.current = false; + reconnectAttemptsRef.current = 0; + void initializeRef.current?.(); + } + }); + return () => { + subscription.remove(); + mountedRef.current = false; + resumeOnActiveRef.current = false; disconnect(); }; }, [initialize, disconnect]); diff --git a/apps/mobile/src/lib/event-source.test.ts b/apps/mobile/src/lib/event-source.test.ts index ebbe4057..1dec1865 100644 --- a/apps/mobile/src/lib/event-source.test.ts +++ b/apps/mobile/src/lib/event-source.test.ts @@ -40,6 +40,15 @@ describe('event-source', () => { expect(mockAddEventListener).toHaveBeenCalled(); }); + it('should forward heartbeat listeners to the underlying EventSource', () => { + const connection = createEventSource('https://example.com/sse'); + + const handler = vi.fn(); + connection.addEventListener('heartbeat', handler); + + expect(mockAddEventListener).toHaveBeenCalledWith('heartbeat', expect.any(Function)); + }); + it('should handle error events specially', () => { const connection = createEventSource('https://example.com/sse'); diff --git a/apps/mobile/src/lib/event-source.ts b/apps/mobile/src/lib/event-source.ts index 566824e8..7b87d68f 100644 --- a/apps/mobile/src/lib/event-source.ts +++ b/apps/mobile/src/lib/event-source.ts @@ -6,7 +6,7 @@ */ import RNEventSource from 'react-native-sse'; -type PairUXEvents = 'connected' | 'signal' | 'presence-join' | 'presence-leave'; +type PairUXEvents = 'connected' | 'heartbeat' | 'signal' | 'presence-join' | 'presence-leave'; export type SSEEventHandler = (event: { data: string }) => void; diff --git a/apps/mobile/src/test/setup.ts b/apps/mobile/src/test/setup.ts index ff7b49a7..6ade3268 100644 --- a/apps/mobile/src/test/setup.ts +++ b/apps/mobile/src/test/setup.ts @@ -5,9 +5,23 @@ import * as React from 'react'; // Make React available globally for JSX globalThis.React = React; +const appStateMock = vi.hoisted(() => ({ + currentState: 'active', + listeners: new Set<(state: string) => void>(), +})); + +export function emitAppStateChange(state: 'active' | 'background' | 'inactive'): void { + appStateMock.currentState = state; + for (const listener of appStateMock.listeners) { + listener(state); + } +} + // Suppress console noise in tests const originalConsole = { ...console }; beforeEach(() => { + appStateMock.currentState = 'active'; + mockPeerConnections.length = 0; vi.stubGlobal('console', { ...originalConsole, error: vi.fn(), @@ -19,6 +33,7 @@ beforeEach(() => { afterEach(() => { cleanup(); + appStateMock.listeners.clear(); vi.unstubAllGlobals(); vi.clearAllMocks(); }); @@ -60,6 +75,19 @@ vi.mock('react-native', () => ({ ActivityIndicator: 'ActivityIndicator', KeyboardAvoidingView: 'KeyboardAvoidingView', Platform: { OS: 'ios' }, + AppState: { + get currentState() { + return appStateMock.currentState; + }, + addEventListener: vi.fn((_event: string, listener: (state: string) => void) => { + appStateMock.listeners.add(listener); + return { + remove: vi.fn(() => { + appStateMock.listeners.delete(listener); + }), + }; + }), + }, Animated: { View: 'Animated.View', Value: vi.fn(() => ({ @@ -70,6 +98,8 @@ vi.mock('react-native', () => ({ })); // ── Mock: react-native-webrtc ───────────────────────────────────── +const mockPeerConnections: MockRTCPeerConnection[] = []; + class MockRTCPeerConnection { signalingState = 'stable'; connectionState = 'new'; @@ -86,6 +116,10 @@ class MockRTCPeerConnection { setParameters: ReturnType; }[] = []; + constructor() { + mockPeerConnections.push(this); + } + createOffer = vi.fn().mockResolvedValue({ type: 'offer', sdp: 'mock-sdp' }); createAnswer = vi.fn().mockResolvedValue({ type: 'answer', sdp: 'mock-answer-sdp' }); setLocalDescription = vi.fn(async (desc: unknown) => { @@ -161,7 +195,7 @@ vi.mock('react-native-webrtc', () => ({ })); // Export for test use -export { MockRTCPeerConnection }; +export { MockRTCPeerConnection, mockPeerConnections }; // ── Mock: react-native-sse ──────────────────────────────────────── vi.mock('react-native-sse', () => { From 26ad258618378aa968f215c800f82ebb3b4193a1 Mon Sep 17 00:00:00 2001 From: Phuc Nguyen Date: Fri, 28 Aug 2026 14:18:44 +0700 Subject: [PATCH 2/2] Format web files for CI --- .../notifications/NotificationPreferences.tsx | 5 +---- apps/web/src/lib/meeting-reminders-runner.ts | 13 +++++++------ 2 files changed, 8 insertions(+), 10 deletions(-) diff --git a/apps/web/src/components/notifications/NotificationPreferences.tsx b/apps/web/src/components/notifications/NotificationPreferences.tsx index 1dcfe4c7..f4050655 100644 --- a/apps/web/src/components/notifications/NotificationPreferences.tsx +++ b/apps/web/src/components/notifications/NotificationPreferences.tsx @@ -34,10 +34,7 @@ const DEFAULT_PREFS: NotificationPrefs = { meetingReminder1Min: true, }; -const PREF_LABELS: Record< - Exclude, - string -> = { +const PREF_LABELS: Record, string> = { controlRequest: 'Control requests', chatMessage: 'Chat messages', participantJoined: 'Participant joined', diff --git a/apps/web/src/lib/meeting-reminders-runner.ts b/apps/web/src/lib/meeting-reminders-runner.ts index 4151b964..c1eea9fc 100644 --- a/apps/web/src/lib/meeting-reminders-runner.ts +++ b/apps/web/src/lib/meeting-reminders-runner.ts @@ -112,7 +112,12 @@ async function alreadySent( const byRecipient = new Map>(); for (const row of data ?? []) { - const r = row as { recipient_kind: string; recipient_key: string; channel: string; lead_minutes: number }; + const r = row as { + recipient_kind: string; + recipient_key: string; + channel: string; + lead_minutes: number; + }; const key = `${r.recipient_kind}:${r.recipient_key}:${r.channel}`; const set = byRecipient.get(key) ?? new Set(); set.add(r.lead_minutes); @@ -292,11 +297,7 @@ export async function runMeetingReminders(now: Date = new Date()): Promise