From f0525f75da633cb4fba41453db6312538b4fe68b Mon Sep 17 00:00:00 2001 From: Phuc Nguyen Date: Mon, 7 Sep 2026 15:01:40 +0700 Subject: [PATCH] Align mobile auth and session signaling with server contracts --- apps/mobile/README.md | 30 ++ apps/mobile/app/(app)/join.tsx | 39 +- apps/mobile/app/(auth)/signup.tsx | 14 +- apps/mobile/src/contexts/AuthContext.test.tsx | 219 ++++++++ apps/mobile/src/contexts/AuthContext.tsx | 85 +++- apps/mobile/src/hooks/useWebRTCHost.test.ts | 368 +++++++++++++- apps/mobile/src/hooks/useWebRTCHost.ts | 242 ++++++--- apps/mobile/src/hooks/useWebRTCViewer.test.ts | 478 +++++++++++++++++- apps/mobile/src/hooks/useWebRTCViewer.ts | 220 ++++++-- apps/mobile/src/lib/api.test.ts | 34 +- apps/mobile/src/lib/api.ts | 7 +- apps/mobile/src/lib/api/auth.test.ts | 182 ++++++- apps/mobile/src/lib/api/auth.ts | 97 +++- apps/mobile/src/lib/api/sessions.test.ts | 57 ++- apps/mobile/src/lib/api/sessions.ts | 33 +- .../src/lib/auth-session-edge-cases.test.ts | 134 +++++ .../mobile/src/lib/auth-session-races.test.ts | 112 ++++ apps/mobile/src/lib/auth-session.test.ts | 286 +++++++++++ apps/mobile/src/lib/auth-session.ts | 273 ++++++++++ apps/mobile/src/lib/event-source.test.ts | 20 +- apps/mobile/src/lib/event-source.ts | 16 +- apps/mobile/src/lib/secure-storage.test.ts | 3 + apps/mobile/src/lib/secure-storage.ts | 2 +- .../src/test/fixtures/server-contracts.ts | 128 +++++ 24 files changed, 2837 insertions(+), 242 deletions(-) create mode 100644 apps/mobile/src/contexts/AuthContext.test.tsx create mode 100644 apps/mobile/src/lib/auth-session-edge-cases.test.ts create mode 100644 apps/mobile/src/lib/auth-session-races.test.ts create mode 100644 apps/mobile/src/lib/auth-session.test.ts create mode 100644 apps/mobile/src/lib/auth-session.ts create mode 100644 apps/mobile/src/test/fixtures/server-contracts.ts diff --git a/apps/mobile/README.md b/apps/mobile/README.md index f21e0d81..6a2f9e32 100644 --- a/apps/mobile/README.md +++ b/apps/mobile/README.md @@ -125,3 +125,33 @@ an App Group, and matching Apple signing entitlements. Those are a separate nati not treat a successful JavaScript bundle as proof that iOS broadcasting is configured. The PairUX app config currently enables the `voip` background mode; remove it or add the matching CallKit/PushKit flow before an App Store submission. + +## Zero-credit validation checklist + +Session reads and refreshes fail closed for the current app process if secure-store deletion +fails during logout. A failed login does not unblock the old tokens; a successfully committed +login does. This also applies when the auth provider remounts. It is not a guarantee of persistent +deletion across an app restart when the OS storage operation failed; verify that failure mode +and recovery on a device before treating it as a release guarantee. + +Run every step below before considering an EAS cloud build. Each one is free and catches a class +of defect that a green cloud build would only package. + +1. `pnpm --filter @pairux/shared-types build` — the mobile app compiles against the built types. +2. `pnpm check:mobile` from the repository root — lint, typecheck, the full unit suite, and the + production bundle verifier. Run it uncached (`--force`) when validating a release candidate. +3. Contract check: any change touching `/api` calls must be validated against the actual route + handler in `apps/web/src/app/api/**` — response envelopes are `{ data, error }`, auth expiry + is in seconds, and the join lookup returns its payload without a wrapper. Update + `src/test/fixtures/server-contracts.ts` from the route source, never from memory. +4. Identity check: signaling tests must keep the authenticated user id, the participant row id, + and the SSE `subscriberId` distinct. A test that reuses one id for all three can pass while + the live flow deadlocks. +5. `pnpm --filter @pairux/mobile build -- --no-install --clean` — clean native prebuild, then + `pnpm --filter @pairux/mobile verify:android-screen-share` against the generated project. +6. If a local Android SDK is available: `pnpm mobile:android` on a device/emulator, then smoke + the supported flow with two distinct accounts — login, join-code lookup, join, host offer / + viewer answer, audio/chat, kill-network reconnect, clean leave. +7. Record anything not exercised (real TURN traversal, background/foreground on physical + hardware, iOS ReplayKit) as an explicit device-only gap instead of assuming the build proves + it. diff --git a/apps/mobile/app/(app)/join.tsx b/apps/mobile/app/(app)/join.tsx index dcf6b5c1..5ea3f415 100644 --- a/apps/mobile/app/(app)/join.tsx +++ b/apps/mobile/app/(app)/join.tsx @@ -12,8 +12,7 @@ import { Platform, } from 'react-native'; import { useRouter } from 'expo-router'; -import type { Session } from '@pairux/shared-types'; -import { sessionApi } from '@/lib/api/sessions'; +import { sessionApi, isScheduledLookup, type JoinLookupResult } from '@/lib/api/sessions'; export default function JoinScreen() { const router = useRouter(); @@ -21,7 +20,7 @@ export default function JoinScreen() { const [displayName, setDisplayName] = useState(''); const [error, setError] = useState(''); const [loading, setLoading] = useState(false); - const [lookupResult, setLookupResult] = useState(null); + const [lookupResult, setLookupResult] = useState(null); const [lookingUp, setLookingUp] = useState(false); async function handleLookup() { @@ -38,8 +37,11 @@ export default function JoinScreen() { return; } - if (result.data?.session) { - setLookupResult(result.data.session); + // The route returns the session (or scheduled meeting) payload directly. + if (result.data) { + setLookupResult(result.data); + } else { + setError('Session not found or has ended'); } } catch { setError('Failed to look up session'); @@ -62,10 +64,18 @@ export default function JoinScreen() { } if (result.data) { + // The participant row carries the authoritative session id. + const sessionId: string = + (result.data.session_id as string | undefined) ?? + (lookupResult && !isScheduledLookup(lookupResult) ? lookupResult.id : ''); + if (!sessionId) { + setError('Joined, but the server did not return a session ID'); + return; + } router.push({ pathname: '/(app)/session/[id]', params: { - id: lookupResult?.id ?? '', + id: sessionId, role: 'viewer', participantId: result.data.id, }, @@ -129,11 +139,24 @@ export default function JoinScreen() { {/* Lookup result */} - {lookupResult ? ( + {lookupResult && isScheduledLookup(lookupResult) ? ( + + Scheduled meeting + {lookupResult.title} + + Starts {new Date(lookupResult.scheduled_at).toLocaleString()} ( + {lookupResult.duration_minutes} min) + + + This meeting hasn't started yet. Try again once the host goes live. + + + ) : lookupResult ? ( Session found - Status: {lookupResult.status} | Code: {lookupResult.join_code} + Status: {lookupResult.status} | Code: {lookupResult.join_code} | Participants:{' '} + {lookupResult.participant_count} {/* Display name input */} diff --git a/apps/mobile/app/(auth)/signup.tsx b/apps/mobile/app/(auth)/signup.tsx index b2ec6801..cc8456e6 100644 --- a/apps/mobile/app/(auth)/signup.tsx +++ b/apps/mobile/app/(auth)/signup.tsx @@ -29,7 +29,7 @@ export default function SignupScreen() { const [showPassword, setShowPassword] = useState(false); const [error, setError] = useState(''); const [loading, setLoading] = useState(false); - const [success, setSuccess] = useState(false); + const [success, setSuccess] = useState<{ needsConfirmation: boolean } | null>(null); async function handleSignup() { setError(''); @@ -62,6 +62,7 @@ export default function SignupScreen() { const result = await signup({ email: email.trim(), password, + confirmPassword, firstName: firstName.trim(), lastName: lastName.trim(), }); @@ -69,7 +70,7 @@ export default function SignupScreen() { if (result.error) { setError(result.error); } else { - setSuccess(true); + setSuccess({ needsConfirmation: result.needsConfirmation ?? true }); } } catch { setError('An unexpected error occurred'); @@ -87,10 +88,13 @@ export default function SignupScreen() { className="mb-6 h-16 w-16" resizeMode="contain" /> - Check your email + + {success.needsConfirmation ? 'Check your email' : 'Account created'} + - We've sent a confirmation email to {email}. Please verify your email address to - continue. + {success.needsConfirmation + ? `We've sent a confirmation email to ${email}. Please verify your email address to continue.` + : 'Your account is ready. You can sign in now.'} { diff --git a/apps/mobile/src/contexts/AuthContext.test.tsx b/apps/mobile/src/contexts/AuthContext.test.tsx new file mode 100644 index 00000000..bc186d01 --- /dev/null +++ b/apps/mobile/src/contexts/AuthContext.test.tsx @@ -0,0 +1,219 @@ +import { describe, it, expect, vi, beforeEach } from 'vitest'; +import { renderHook, waitFor, act } from '@testing-library/react'; +import { AuthProvider, useAuth } from './AuthContext'; +import { isAuthExpired, type StoredAuth } from '@/lib/secure-storage'; +import { + readAuthSession as getStoredAuth, + refreshAuthSession, + clearAuthSession, +} from '@/lib/auth-session'; +import { authApi } from '@/lib/api/auth'; +import { AUTH_USER_ID } from '../test/fixtures/server-contracts'; + +vi.mock('@/lib/secure-storage'); +vi.mock('@/lib/auth-session'); +vi.mock('@/lib/api/auth', () => ({ + authApi: { + login: vi.fn(), + signup: vi.fn(), + logout: vi.fn(), + getSession: vi.fn(), + }, +})); + +const storedAuth: StoredAuth = { + accessToken: 'access-token-1', + refreshToken: 'refresh-token-1', + expiresAt: Date.now() + 3600000, + user: { id: AUTH_USER_ID, email: 'user@example.com' }, +}; + +function deferred() { + let resolve!: (value: T | PromiseLike) => void; + const promise = new Promise((resolvePromise) => { + resolve = resolvePromise; + }); + return { promise, resolve }; +} + +function renderAuth() { + return renderHook(() => useAuth(), { + wrapper: ({ children }: { children: React.ReactNode }) => ( + {children} + ), + }); +} + +describe('AuthContext session restore', () => { + beforeEach(() => { + vi.clearAllMocks(); + }); + + it('restores a valid stored session', async () => { + vi.mocked(getStoredAuth).mockResolvedValue(storedAuth); + vi.mocked(isAuthExpired).mockReturnValue(false); + + const { result } = renderAuth(); + + await waitFor(() => expect(result.current.isLoading).toBe(false)); + expect(result.current.isAuthenticated).toBe(true); + expect(result.current.user).toEqual(storedAuth.user); + expect(refreshAuthSession).not.toHaveBeenCalled(); + }); + + it('refreshes an expired stored session instead of forcing a re-login', async () => { + vi.mocked(getStoredAuth).mockResolvedValue(storedAuth); + vi.mocked(isAuthExpired).mockReturnValue(true); + vi.mocked(refreshAuthSession).mockResolvedValue({ + auth: { ...storedAuth, accessToken: 'access-token-2' }, + }); + + const { result } = renderAuth(); + + await waitFor(() => expect(result.current.isLoading).toBe(false)); + expect(refreshAuthSession).toHaveBeenCalledTimes(1); + expect(result.current.isAuthenticated).toBe(true); + expect(result.current.user).toEqual(storedAuth.user); + expect(clearAuthSession).not.toHaveBeenCalled(); + }); + + it('signs out only when the server definitively rejects the refresh token', async () => { + vi.mocked(getStoredAuth).mockResolvedValue(storedAuth); + vi.mocked(isAuthExpired).mockReturnValue(true); + vi.mocked(refreshAuthSession).mockResolvedValue({ auth: null, failure: 'rejected' }); + + const { result } = renderAuth(); + + await waitFor(() => expect(result.current.isLoading).toBe(false)); + expect(result.current.isAuthenticated).toBe(false); + expect(clearAuthSession).toHaveBeenCalledTimes(1); + }); + + it('keeps the stored session when the refresh fails transiently (offline start)', async () => { + vi.mocked(getStoredAuth).mockResolvedValue(storedAuth); + vi.mocked(isAuthExpired).mockReturnValue(true); + vi.mocked(refreshAuthSession).mockResolvedValue({ auth: null, failure: 'transient' }); + + const { result } = renderAuth(); + + await waitFor(() => expect(result.current.isLoading).toBe(false)); + // A network hiccup is not a logout: the refresh token may still be good + expect(clearAuthSession).not.toHaveBeenCalled(); + expect(result.current.isAuthenticated).toBe(true); + expect(result.current.user).toEqual(storedAuth.user); + }); + + it('leaves auth state alone when the restore refresh was superseded', async () => { + vi.mocked(getStoredAuth).mockResolvedValue(storedAuth); + vi.mocked(isAuthExpired).mockReturnValue(true); + vi.mocked(refreshAuthSession).mockResolvedValue({ auth: null, failure: 'superseded' }); + + const { result } = renderAuth(); + + await waitFor(() => expect(result.current.isLoading).toBe(false)); + // A concurrent login/logout owns the state; restore must not clear it + expect(clearAuthSession).not.toHaveBeenCalled(); + expect(result.current.user).toBeNull(); + }); + + it('stays signed out with no stored session', async () => { + vi.mocked(getStoredAuth).mockResolvedValue(null); + + const { result } = renderAuth(); + + await waitFor(() => expect(result.current.isLoading).toBe(false)); + expect(result.current.isAuthenticated).toBe(false); + expect(refreshAuthSession).not.toHaveBeenCalled(); + }); +}); + +describe('AuthContext lifecycle races', () => { + beforeEach(() => { + vi.clearAllMocks(); + }); + + it('finishes loading if a login supersedes restore but fails', async () => { + const restore = deferred(); + vi.mocked(getStoredAuth).mockReturnValueOnce(restore.promise); + vi.mocked(authApi.login).mockResolvedValueOnce({ error: 'Invalid credentials' }); + const { result } = renderAuth(); + await act(async () => { + await result.current.login('fixture@example.com', 'fixture-password'); + }); + expect(result.current.isLoading).toBe(false); + await act(async () => { + restore.resolve(null); + }); + expect(result.current.isLoading).toBe(false); + expect(result.current.user).toBeNull(); + }); + + it('does not resurrect the user when a login resolves after a logout', async () => { + vi.mocked(getStoredAuth).mockResolvedValue(null); + vi.mocked(authApi.logout).mockResolvedValue(undefined); + const loginGate = deferred<{ data?: StoredAuth; error?: string }>(); + vi.mocked(authApi.login).mockReturnValueOnce(loginGate.promise); + + const { result } = renderAuth(); + await waitFor(() => expect(result.current.isLoading).toBe(false)); + + let loginPromise!: Promise<{ error?: string }>; + act(() => { + loginPromise = result.current.login('new@example.com', 'password'); + }); + await act(async () => { + await result.current.logout(); + }); + await act(async () => { + loginGate.resolve({ + data: { + ...storedAuth, + user: { id: 'late-user', email: 'late@example.com' }, + }, + }); + await loginPromise; + }); + + expect(result.current.user).toBeNull(); + expect(result.current.isAuthenticated).toBe(false); + }); + + it('does not let a slow restore overwrite a fresh login', async () => { + const restoreGate = deferred(); + vi.mocked(getStoredAuth).mockReturnValueOnce(restoreGate.promise); + vi.mocked(isAuthExpired).mockReturnValue(false); + vi.mocked(authApi.login).mockResolvedValue({ + data: { ...storedAuth, user: { id: 'fresh-user', email: 'fresh@example.com' } }, + }); + + const { result } = renderAuth(); + await act(async () => { + await result.current.login('fresh@example.com', 'password'); + }); + expect(result.current.user?.id).toBe('fresh-user'); + + await act(async () => { + // The pre-login stored session finally loads — it is stale now + restoreGate.resolve(storedAuth); + }); + await waitFor(() => expect(result.current.isLoading).toBe(false)); + + expect(result.current.user?.id).toBe('fresh-user'); + }); + + it('signs out immediately even while the logout network call is pending', async () => { + vi.mocked(getStoredAuth).mockResolvedValue(storedAuth); + vi.mocked(isAuthExpired).mockReturnValue(false); + vi.mocked(authApi.logout).mockReturnValue(new Promise(() => undefined)); + + const { result } = renderAuth(); + await waitFor(() => expect(result.current.isAuthenticated).toBe(true)); + + act(() => { + void result.current.logout(); + }); + + expect(result.current.user).toBeNull(); + expect(result.current.isAuthenticated).toBe(false); + }); +}); diff --git a/apps/mobile/src/contexts/AuthContext.tsx b/apps/mobile/src/contexts/AuthContext.tsx index 4420fecc..dd0c1666 100644 --- a/apps/mobile/src/contexts/AuthContext.tsx +++ b/apps/mobile/src/contexts/AuthContext.tsx @@ -4,16 +4,15 @@ * Provides login/signup/logout actions and persists auth state * via expo-secure-store. Auto-restores session on mount. */ -import React, { createContext, useContext, useState, useEffect, useCallback } from 'react'; +import React, { createContext, useContext, useState, useEffect, useCallback, useRef } from 'react'; import { - getStoredAuth, isAuthExpired, - clearStoredAuth, storeCredentials, getStoredCredentials, clearStoredCredentials, type StoredCredentials, } from '@/lib/secure-storage'; +import { clearAuthSession, readAuthSession, refreshAuthSession } from '@/lib/auth-session'; import { authApi } from '@/lib/api/auth'; interface AuthUser { @@ -32,9 +31,10 @@ interface AuthContextValue { signup: (params: { email: string; password: string; + confirmPassword: string; firstName: string; lastName: string; - }) => Promise<{ error?: string; success?: boolean }>; + }) => Promise<{ error?: string; success?: boolean; needsConfirmation?: boolean }>; logout: () => Promise; } @@ -44,22 +44,48 @@ export function AuthProvider({ children }: { children: React.ReactNode }) { const [user, setUser] = useState(null); const [isLoading, setIsLoading] = useState(true); const [rememberMe, setRememberMe] = useState(false); + // Bumped when a login or logout starts so a slower restore (or an older + // login) that resolves afterwards cannot publish stale auth state. + const authOpSeqRef = useRef(0); + const mountedRef = useRef(true); + + useEffect(() => { + mountedRef.current = true; + return () => { + mountedRef.current = false; + authOpSeqRef.current += 1; + }; + }, []); // Restore session from secure storage on mount useEffect(() => { + const seq = authOpSeqRef.current; + const canPublish = () => mountedRef.current && authOpSeqRef.current === seq; async function restore() { try { - const stored = await getStoredAuth(); - if (stored && !isAuthExpired(stored)) { - setUser(stored.user); - } else if (stored) { - // Token expired — clear it - await clearStoredAuth(); + const stored = await readAuthSession(); + if (!stored) return; + if (!isAuthExpired(stored)) { + if (canPublish()) setUser(stored.user); + return; } + // Token expired — refresh instead of forcing a re-login + const { auth, failure } = await refreshAuthSession(); + if (auth) { + if (canPublish()) setUser(auth.user); + } else if (failure === 'rejected') { + // The server refused the refresh token; this session is dead. + if (canPublish()) await clearAuthSession(); + } else if (failure === 'transient') { + // Offline start or server hiccup: keep the stored session (the + // refresh token may still be good) and retry per API request. + if (canPublish()) setUser(stored.user); + } + // 'superseded'/'signed-out': a login or logout already owns the state. } catch { - // Failed to read storage + // Failed to read storage — stay signed out without destroying tokens. } finally { - setIsLoading(false); + if (canPublish()) setIsLoading(false); } } void restore(); @@ -67,16 +93,22 @@ export function AuthProvider({ children }: { children: React.ReactNode }) { const login = useCallback( async (email: string, password: string) => { + const seq = ++authOpSeqRef.current; const result = await authApi.login(email, password); if (result.error) { + if (mountedRef.current && authOpSeqRef.current === seq) setIsLoading(false); return { error: result.error }; } if (result.data) { - setUser(result.data.user); - if (rememberMe) { - await storeCredentials({ email, password }); - } else { - await clearStoredCredentials(); + // Publish only if no logout/newer login started while we awaited. + if (mountedRef.current && authOpSeqRef.current === seq) { + setUser(result.data.user); + setIsLoading(false); + if (rememberMe) { + await storeCredentials({ email, password }); + } else { + await clearStoredCredentials(); + } } } return {}; @@ -85,12 +117,18 @@ export function AuthProvider({ children }: { children: React.ReactNode }) { ); const signup = useCallback( - async (params: { email: string; password: string; firstName: string; lastName: string }) => { + async (params: { + email: string; + password: string; + confirmPassword: string; + firstName: string; + lastName: string; + }) => { const result = await authApi.signup(params); if (result.error) { return { error: result.error }; } - return { success: true }; + return { success: true, needsConfirmation: result.data?.needsConfirmation ?? true }; }, [] ); @@ -100,8 +138,15 @@ export function AuthProvider({ children }: { children: React.ReactNode }) { }, []); const logout = useCallback(async () => { + // Sign out locally first; the epoch bump inside authApi.logout -> + // clearAuthSession invalidates in-flight refreshes/logins, and the + // server-side revoke is fire-and-forget. + authOpSeqRef.current += 1; + if (mountedRef.current) { + setUser(null); + setIsLoading(false); + } await authApi.logout(); - setUser(null); }, []); return ( diff --git a/apps/mobile/src/hooks/useWebRTCHost.test.ts b/apps/mobile/src/hooks/useWebRTCHost.test.ts index d804a0f5..28649c1f 100644 --- a/apps/mobile/src/hooks/useWebRTCHost.test.ts +++ b/apps/mobile/src/hooks/useWebRTCHost.test.ts @@ -3,11 +3,12 @@ import { renderHook, act, waitFor } from '@testing-library/react'; import { Platform } from 'react-native'; import { useWebRTCHost } from './useWebRTCHost'; import { createEventSource } from '../lib/event-source'; -import { getStoredAuth } from '../lib/secure-storage'; +import { getValidAccessToken } from '../lib/auth-session'; import { mediaDevices } from 'react-native-webrtc'; import type { MediaStream, MediaStreamTrack } from 'react-native-webrtc'; -import { emitAppStateChange } from '../test/setup'; +import { emitAppStateChange, mockPeerConnections } from '../test/setup'; import { runAndroidNativePrompt } from '../lib/android-native-prompt'; +import { hostConnectedEventData, HOST_USER_ID } from '../test/fixtures/server-contracts'; function deferred() { let resolve!: (value: T | PromiseLike) => void; @@ -21,14 +22,8 @@ vi.mock('../config', () => ({ API_BASE_URL: 'https://pairux.com', })); -vi.mock('../lib/secure-storage', () => ({ - getStoredAuth: vi.fn().mockResolvedValue({ - accessToken: 'test-token', - refreshToken: 'refresh', - expiresAt: Date.now() + 3600000, - user: { id: 'host-1', email: 'host@example.com' }, - }), - isAuthExpired: vi.fn().mockReturnValue(false), +vi.mock('../lib/auth-session', () => ({ + getValidAccessToken: vi.fn().mockResolvedValue('test-token'), })); const { mockClose, mockAddEventListener, mockEventSources } = vi.hoisted(() => ({ @@ -41,7 +36,7 @@ const { mockClose, mockAddEventListener, mockEventSources } = vi.hoisted(() => ( })); vi.mock('../lib/event-source', () => ({ - createEventSource: vi.fn(() => { + createEventSource: vi.fn((_url: string, _options?: { headers?: Record }) => { const listeners = new Map void>(); const source = { listeners, @@ -62,6 +57,7 @@ describe('useWebRTCHost', () => { beforeEach(() => { vi.clearAllMocks(); mockEventSources.length = 0; + vi.mocked(getValidAccessToken).mockResolvedValue('test-token'); vi.mocked(fetch).mockResolvedValue({ ok: true, text: async () => 'ok', @@ -121,12 +117,13 @@ describe('useWebRTCHost', () => { }); expect(createEventSource).toHaveBeenCalledWith( - expect.stringContaining('/api/sessions/session-1/signal/stream') + expect.stringContaining('/api/sessions/session-1/signal/stream'), + expect.anything() ); }); it('should set error when not authenticated', async () => { - vi.mocked(getStoredAuth).mockResolvedValueOnce(null); + vi.mocked(getValidAccessToken).mockResolvedValueOnce(null); const { result } = renderHook(() => useWebRTCHost({ @@ -675,4 +672,349 @@ describe('useWebRTCHost', () => { vi.useRealTimers(); } }); + + it('sends the SSE bearer token in the Authorization header, not the URL', async () => { + const { result } = renderHook(() => + useWebRTCHost({ + sessionId: 'session-1', + hostId: 'host-1', + }) + ); + + await act(async () => { + await result.current.startHosting(); + }); + + const [url, options] = vi.mocked(createEventSource).mock.calls[0] as [ + string, + { headers?: Record } | undefined, + ]; + expect(url).not.toContain('token='); + expect(options).toEqual({ headers: { Authorization: 'Bearer test-token' } }); + }); + + it('signs offers with the server-assigned subscriberId, not the local hostId', async () => { + const { result } = renderHook(() => + useWebRTCHost({ + sessionId: 'session-1', + hostId: 'host-local-id', + }) + ); + + await act(async () => { + await result.current.startHosting(); + }); + const source = mockEventSources[0]; + expect(source).toBeDefined(); + + act(() => { + // The server identifies the authenticated host by its auth user id + source.listeners.get('connected')?.({ data: hostConnectedEventData() }); + }); + + await act(async () => { + source.listeners.get('presence-join')?.({ + data: JSON.stringify({ presences: [{ user_id: 'viewer-sub-1', role: 'viewer' }] }), + }); + await Promise.resolve(); + await Promise.resolve(); + await Promise.resolve(); + }); + + const requestBody = (init?: RequestInit): string => + typeof init?.body === 'string' ? init.body : ''; + const offerCall = vi + .mocked(fetch) + .mock.calls.find(([, init]) => requestBody(init).includes('"type":"offer"')); + expect(offerCall).toBeDefined(); + const body = requestBody(offerCall?.[1]); + expect(body).toContain(`"senderId":"${HOST_USER_ID}"`); + expect(body).toContain('"targetId":"viewer-sub-1"'); + expect(body).not.toContain('host-local-id'); + }); + + it('ignores answers addressed to a different subscriber', async () => { + const { result } = renderHook(() => + useWebRTCHost({ + sessionId: 'session-1', + hostId: 'host-local-id', + }) + ); + + await act(async () => { + await result.current.startHosting(); + }); + const source = mockEventSources[0]; + + act(() => { + source.listeners.get('connected')?.({ data: hostConnectedEventData() }); + }); + await act(async () => { + source.listeners.get('presence-join')?.({ + data: JSON.stringify({ presences: [{ user_id: 'viewer-sub-1', role: 'viewer' }] }), + }); + await Promise.resolve(); + await Promise.resolve(); + }); + await waitFor(() => expect(result.current.viewerCount).toBe(1)); + + const viewer = result.current.viewers.get('viewer-sub-1'); + expect(viewer).toBeDefined(); + (viewer!.peerConnection as unknown as { signalingState: string }).signalingState = + 'have-local-offer'; + + // Targeted at some other subscriber: must not touch our peer connection + await act(async () => { + source.listeners.get('signal')?.({ + data: JSON.stringify({ + type: 'answer', + sdp: 'answer-sdp', + senderId: 'viewer-sub-1', + targetId: 'someone-else', + timestamp: Date.now(), + }), + }); + await Promise.resolve(); + }); + expect(viewer!.peerConnection.setRemoteDescription).not.toHaveBeenCalled(); + + // Targeted at our server-assigned subscriber id: processed + await act(async () => { + source.listeners.get('signal')?.({ + data: JSON.stringify({ + type: 'answer', + sdp: 'answer-sdp', + senderId: 'viewer-sub-1', + targetId: HOST_USER_ID, + timestamp: Date.now(), + }), + }); + await Promise.resolve(); + }); + expect(viewer!.peerConnection.setRemoteDescription).toHaveBeenCalledWith({ + type: 'answer', + sdp: 'answer-sdp', + }); + }); + + it('posts the host liveness heartbeat while hosting and stops with it', async () => { + vi.useFakeTimers(); + vi.setSystemTime(new Date('2026-08-28T00:00:00Z')); + + const heartbeatCalls = () => + vi + .mocked(fetch) + .mock.calls.filter( + ([url]) => typeof url === 'string' && url.includes('/api/sessions/session-1/heartbeat') + ).length; + + try { + const { result } = renderHook(() => + useWebRTCHost({ + sessionId: 'session-1', + hostId: 'host-1', + }) + ); + + await act(async () => { + await result.current.startHosting(); + }); + expect(heartbeatCalls()).toBe(0); + + await act(async () => { + mockEventSources[0]?.listeners.get('connected')?.({ data: hostConnectedEventData() }); + await vi.advanceTimersByTimeAsync(0); + }); + // Stamped immediately so the room appears live without a 30s wait + expect(heartbeatCalls()).toBe(1); + expect(fetch).toHaveBeenCalledWith( + 'https://pairux.com/api/sessions/session-1/heartbeat', + expect.objectContaining({ + method: 'POST', + headers: { Authorization: 'Bearer test-token' }, + }) + ); + + await act(async () => { + await vi.advanceTimersByTimeAsync(30_000); + }); + expect(heartbeatCalls()).toBe(2); + + act(() => { + result.current.stopHosting(); + }); + await act(async () => { + await vi.advanceTimersByTimeAsync(60_000); + }); + expect(heartbeatCalls()).toBe(2); + } finally { + vi.useRealTimers(); + } + }); + + describe('credential lifecycle guards', () => { + const requestBodies = () => + vi + .mocked(fetch) + .mock.calls.map(([, init]) => (typeof init?.body === 'string' ? init.body : '')); + + it('drops signals instead of posting once no valid token can be resolved', async () => { + const { result } = renderHook(() => + useWebRTCHost({ + sessionId: 'session-1', + hostId: 'host-1', + }) + ); + await act(async () => { + await result.current.startHosting(); + }); + const source = mockEventSources[0]; + act(() => { + source.listeners.get('connected')?.({ data: hostConnectedEventData() }); + }); + + // The account signs out: refresh yields nothing from here on + vi.mocked(getValidAccessToken).mockResolvedValue(null); + + await act(async () => { + source.listeners.get('presence-join')?.({ + data: JSON.stringify({ presences: [{ user_id: 'viewer-sub-1', role: 'viewer' }] }), + }); + await Promise.resolve(); + await Promise.resolve(); + await Promise.resolve(); + }); + + // The offer for the joining viewer must not be posted with any + // previously cached credential + expect(requestBodies().some((body) => body.includes('"type":"offer"'))).toBe(false); + }); + + it('does not signal after hosting stops while the token refresh is in flight', async () => { + const { result } = renderHook(() => + useWebRTCHost({ + sessionId: 'session-1', + hostId: 'host-1', + }) + ); + await act(async () => { + await result.current.startHosting(); + }); + const source = mockEventSources[0]; + act(() => { + source.listeners.get('connected')?.({ data: hostConnectedEventData() }); + }); + + const tokenGate = deferred(); + vi.mocked(getValidAccessToken).mockReturnValue(tokenGate.promise); + + await act(async () => { + source.listeners.get('presence-join')?.({ + data: JSON.stringify({ presences: [{ user_id: 'viewer-sub-1', role: 'viewer' }] }), + }); + await Promise.resolve(); + await Promise.resolve(); + }); + + act(() => { + result.current.stopHosting(); + }); + + await act(async () => { + tokenGate.resolve('late-token'); + await Promise.resolve(); + await Promise.resolve(); + }); + + expect(requestBodies().some((body) => body.includes('"type":"offer"'))).toBe(false); + const headers = vi + .mocked(fetch) + .mock.calls.map(([, init]) => JSON.stringify(init?.headers ?? {})); + expect(headers.some((header) => header.includes('late-token'))).toBe(false); + }); + }); + + it('keeps the published screen-share stream across an automatic SSE-error restart', async () => { + vi.useFakeTimers(); + try { + const { result } = renderHook(() => + useWebRTCHost({ + sessionId: 'session-1', + hostId: 'host-1', + }) + ); + await act(async () => { + await result.current.startHosting(); + }); + const source = mockEventSources[0]; + act(() => { + source.listeners.get('connected')?.({ data: hostConnectedEventData() }); + }); + await act(async () => { + source.listeners.get('presence-join')?.({ + data: JSON.stringify({ presences: [{ user_id: 'viewer-sub-1', role: 'viewer' }] }), + }); + await vi.advanceTimersByTimeAsync(0); + }); + expect(result.current.viewerCount).toBe(1); + + 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); + }); + + // Transient transport error: the hook rebuilds the SSE after 3s + await act(async () => { + source.listeners.get('error')?.({ data: '' }); + await vi.advanceTimersByTimeAsync(3_100); + }); + expect(mockEventSources).toHaveLength(2); + const rebuiltSource = mockEventSources[1]; + + act(() => { + rebuiltSource.listeners.get('connected')?.({ data: hostConnectedEventData() }); + }); + const pcCountBefore = mockPeerConnections.length; + await act(async () => { + rebuiltSource.listeners.get('presence-join')?.({ + data: JSON.stringify({ presences: [{ user_id: 'viewer-sub-1', role: 'viewer' }] }), + }); + await vi.advanceTimersByTimeAsync(0); + }); + expect(mockPeerConnections.length).toBe(pcCountBefore + 1); + const rebuiltPc = mockPeerConnections[mockPeerConnections.length - 1]; + + // The still-live capture re-attaches: no re-publish, no new capture + // permission prompt, and the caller-owned track is never stopped + expect(rebuiltPc.addTrack).toHaveBeenCalledWith(screenTrack, screenStream); + expect(screenTrack.stop).not.toHaveBeenCalled(); + + // A user-initiated stop still clears the published stream + act(() => { + result.current.stopHosting(); + }); + await act(async () => { + await result.current.startHosting(); + }); + const freshSource = mockEventSources[2]; + act(() => { + freshSource.listeners.get('connected')?.({ data: hostConnectedEventData() }); + }); + await act(async () => { + freshSource.listeners.get('presence-join')?.({ + data: JSON.stringify({ presences: [{ user_id: 'viewer-sub-1', role: 'viewer' }] }), + }); + await vi.advanceTimersByTimeAsync(0); + }); + const freshPc = mockPeerConnections[mockPeerConnections.length - 1]; + expect(freshPc.addTrack).not.toHaveBeenCalledWith(screenTrack, expect.anything()); + } finally { + vi.useRealTimers(); + } + }); }); diff --git a/apps/mobile/src/hooks/useWebRTCHost.ts b/apps/mobile/src/hooks/useWebRTCHost.ts index 6b4d18fe..8b820bc4 100644 --- a/apps/mobile/src/hooks/useWebRTCHost.ts +++ b/apps/mobile/src/hooks/useWebRTCHost.ts @@ -31,7 +31,7 @@ import { isAndroidNativePromptActive, subscribeAndroidNativePrompt, } from '../lib/android-native-prompt'; -import { getStoredAuth, isAuthExpired } from '../lib/secure-storage'; +import { getValidAccessToken } from '../lib/auth-session'; import { createEventSource, type SSEConnection } from '../lib/event-source'; // RN WebRTC's RTCDataChannel type (differs from browser global) @@ -70,6 +70,11 @@ const _BITRATE_PRESETS: Record = { const STATS_INTERVAL = 30000; const HEARTBEAT_TIMEOUT = 75000; const HEARTBEAT_WATCHDOG_INTERVAL = 15000; +// Delay before rebuilding the SSE transport (with a fresh token) after an error +const SSE_RECONNECT_DELAY = 3000; +// Cadence of the host liveness ping that keeps a public room on /live +// (matches the desktop host in CapturePreview.tsx) +const HOST_LIVENESS_INTERVAL = 30000; const DEFAULT_ICE_SERVERS: RTCIceServer[] = [ { urls: 'stun:stun.l.google.com:19302' }, @@ -97,6 +102,16 @@ interface SignalMessage { timestamp: number; } +interface StopHostingOptions { + /** + * Keep the caller-owned published stream reference so an automatic + * transport restart re-attaches the live screen-share tracks to the + * rebuilt peer connections. Without this, every transient SSE drop would + * blank the share and force a brand-new native capture permission flow. + */ + preservePublishedStream?: boolean; +} + interface UseWebRTCHostOptions { sessionId: string; hostId: string; @@ -145,9 +160,10 @@ export function useWebRTCHost({ const viewersRef = useRef>(new Map()); const statsIntervalRef = useRef | null>(null); const heartbeatWatchdogRef = useRef | null>(null); + const livenessIntervalRef = useRef | null>(null); + const sseReconnectTimerRef = 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); @@ -162,13 +178,17 @@ export function useWebRTCHost({ const iceServersRef = useRef(DEFAULT_ICE_SERVERS); const pendingCandidatesRef = useRef>(new Map()); + // The SSE `connected` event carries the server-assigned subscriber ID; the + // stream only delivers answers/ICE targeted at that ID, so it must be the + // senderId on every offer this host posts. Falls back to the hostId prop. + const hostSignalIdRef = useRef(hostId); 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); + const stopHostingRef = useRef<((options?: StopHostingOptions) => void) | undefined>(undefined); onControlRequestRef.current = onControlRequest; onInputReceivedRef.current = onInputReceived; onViewerJoinedRef.current = onViewerJoined; @@ -181,16 +201,22 @@ export function useWebRTCHost({ // Send signal via API const sendSignal = useCallback( - async (signal: SignalMessage): Promise => { + async (signal: SignalMessage, generation: number): Promise => { try { - const headers: Record = { 'Content-Type': 'application/json' }; - if (authTokenRef.current) { - headers.Authorization = `Bearer ${authTokenRef.current}`; + // Resolve the token per request so signaling keeps working after the + // access token expires mid-session (single-flight refresh). No token + // means signed out — never post with a stale cached credential, and + // never post for a generation that stopped while the token resolved. + const token = await getValidAccessToken(); + if (!isCurrentGeneration(generation)) return false; + if (!token) { + console.error('[WebRTCHost] Not authenticated; dropping signal'); + return false; } const response = await fetch(`${API_BASE_URL}/api/sessions/${sessionId}/signal`, { method: 'POST', - headers, + headers: { 'Content-Type': 'application/json', Authorization: `Bearer ${token}` }, body: JSON.stringify(signal), }); @@ -204,7 +230,7 @@ export function useWebRTCHost({ return false; } }, - [sessionId] + [isCurrentGeneration, sessionId] ); const renegotiateViewer = useCallback( @@ -229,18 +255,21 @@ export function useWebRTCHost({ throw new Error(`Failed to create an SDP offer for viewer ${viewer.id}`); } - const sent = await sendSignal({ - type: 'offer', - sdp: offer.sdp, - senderId: hostId, - targetId: viewer.id, - timestamp: Date.now(), - }); + const sent = await sendSignal( + { + type: 'offer', + sdp: offer.sdp, + senderId: hostSignalIdRef.current, + targetId: viewer.id, + timestamp: Date.now(), + }, + generation + ); if (!sent) { throw new Error(`Failed to signal viewer ${viewer.id}`); } }, - [hostId, isCurrentGeneration, sendSignal] + [isCurrentGeneration, sendSignal] ); // Report usage stats @@ -289,14 +318,15 @@ export function useWebRTCHost({ } }); - const headers: Record = { 'Content-Type': 'application/json' }; - if (authTokenRef.current) { - headers.Authorization = `Bearer ${authTokenRef.current}`; + const token = await getValidAccessToken(); + if (!isCurrentGeneration(generation) || viewersRef.current.get(viewer.id) !== viewer) { + return; } + if (!token) continue; await fetch(`${API_BASE_URL}/api/sessions/${sessionId}/stats`, { method: 'POST', - headers, + headers: { 'Content-Type': 'application/json', Authorization: `Bearer ${token}` }, body: JSON.stringify({ participantId: hostId, role: 'host', @@ -378,20 +408,23 @@ export function useWebRTCHost({ if (!isCurrentGeneration(generation)) return; if (offer.sdp) { - await sendSignal({ - type: 'offer', - sdp: offer.sdp, - senderId: hostId, - targetId: otherId, - timestamp: Date.now(), - }); + await sendSignal( + { + type: 'offer', + sdp: offer.sdp, + senderId: hostSignalIdRef.current, + targetId: otherId, + timestamp: Date.now(), + }, + generation + ); } } catch (err) { console.error(`[WebRTCHost] Failed to relay audio to ${otherId}:`, err); } } }, - [hostId, isCurrentGeneration, sendSignal] + [isCurrentGeneration, sendSignal] ); // Create peer connection for a viewer @@ -452,13 +485,16 @@ export function useWebRTCHost({ return; } if (event.candidate) { - void sendSignal({ - type: 'ice-candidate', - candidate: event.candidate.toJSON(), - senderId: hostId, - targetId: viewerId, - timestamp: Date.now(), - }); + void sendSignal( + { + type: 'ice-candidate', + candidate: event.candidate.toJSON(), + senderId: hostSignalIdRef.current, + targetId: viewerId, + timestamp: Date.now(), + }, + generation + ); } }); @@ -558,7 +594,7 @@ export function useWebRTCHost({ return pc; }, - [hostId, handleDataChannelMessage, isCurrentGeneration, sendSignal, relayAudioToOtherViewers] + [handleDataChannelMessage, isCurrentGeneration, sendSignal, relayAudioToOtherViewers] ); // Remove viewer @@ -583,7 +619,7 @@ export function useWebRTCHost({ const handleViewerJoin = useCallback( async (viewerId: string, generation: number) => { if (!isCurrentGeneration(generation)) return; - if (viewerId === hostId) return; + if (viewerId === hostSignalIdRef.current || viewerId === hostId) return; if (viewersRef.current.has(viewerId)) return; console.log('[WebRTCHost] Viewer joining:', viewerId); @@ -615,13 +651,16 @@ export function useWebRTCHost({ if (!isCurrentGeneration(generation) || viewersRef.current.get(viewerId) !== viewer) return; if (offer.sdp) { - await sendSignal({ - type: 'offer', - sdp: offer.sdp, - senderId: hostId, - targetId: viewerId, - timestamp: Date.now(), - }); + await sendSignal( + { + type: 'offer', + sdp: offer.sdp, + senderId: hostSignalIdRef.current, + targetId: viewerId, + timestamp: Date.now(), + }, + generation + ); } } catch (err) { console.error('[WebRTCHost] Failed to create offer:', err); @@ -637,7 +676,11 @@ export function useWebRTCHost({ const handleSignalMessage = useCallback( async (signal: SignalMessage, generation: number) => { if (!isCurrentGeneration(generation)) return; - if (signal.targetId && signal.targetId !== hostId) return; + // Only process signals addressed to this host (viewer answers/ICE are + // targeted at the senderId we used in the offer) and never our own + // broadcasts echoed back. + if (signal.senderId === hostSignalIdRef.current) return; + if (signal.targetId && signal.targetId !== hostSignalIdRef.current) return; const viewerId = signal.senderId; @@ -692,7 +735,7 @@ export function useWebRTCHost({ } } }, - [hostId, isCurrentGeneration] + [isCurrentGeneration] ); // Toggle host microphone @@ -711,6 +754,27 @@ export function useWebRTCHost({ setMicEnabled(newEnabled); }, []); + // Liveness heartbeat: while actively hosting, ping the server so a published + // room shows as live on pairux.com/live — and, critically, so it falls OFF + // /live automatically when this app is closed/killed (the pings just stop). + // Same route and cadence as the desktop host (CapturePreview.tsx). + const sendLivenessHeartbeat = useCallback( + async (generation: number) => { + if (!isCurrentGeneration(generation)) return; + try { + const token = await getValidAccessToken(); + if (!token || !isCurrentGeneration(generation)) return; + await fetch(`${API_BASE_URL}/api/sessions/${sessionId}/heartbeat`, { + method: 'POST', + headers: { Authorization: `Bearer ${token}` }, + }); + } catch { + // best effort — a missed ping just delays the live/offline flip + } + }, + [isCurrentGeneration, sessionId] + ); + // Start hosting const startHosting = useCallback(async () => { if (!mountedRef.current) { @@ -729,16 +793,18 @@ export function useWebRTCHost({ console.log('[WebRTCHost] Starting hosting for session:', sessionId); - // Get auth token from secure storage + // Get a valid access token, refreshing through /api/auth/refresh if the + // stored one has expired. + let sseToken: string; try { - const stored = await getStoredAuth(); + const token = await getValidAccessToken(); if (!isCurrentGeneration(generation)) return; - if (!stored || isAuthExpired(stored)) { + if (!token) { isStartingRef.current = false; setError('Not authenticated. Please log in again.'); return; } - authTokenRef.current = stored.accessToken; + sseToken = token; } catch (err) { console.error('[WebRTCHost] Failed to get auth token:', err); if (isCurrentGeneration(generation)) { @@ -772,16 +838,16 @@ export function useWebRTCHost({ setMicEnabled(false); } - // Build SSE URL + // Build SSE URL. The bearer token goes in the Authorization header + // (react-native-sse supports headers, unlike browser EventSource) so + // credentials never appear in proxy/request logs via the query string. const sseParams = new URLSearchParams({ participantId: hostId, }); - if (authTokenRef.current) { - sseParams.set('token', authTokenRef.current); - } - const sseUrl = `${API_BASE_URL}/api/sessions/${sessionId}/signal/stream?${sseParams.toString()}`; - const eventSource = createEventSource(sseUrl); + const eventSource = createEventSource(sseUrl, { + headers: { Authorization: `Bearer ${sseToken}` }, + }); if (!isCurrentGeneration(generation)) { eventSource.close(); return; @@ -801,13 +867,31 @@ export function useWebRTCHost({ setError(null); try { - const data = JSON.parse(event.data) as { iceServers?: RTCIceServer[] }; + const data = JSON.parse(event.data) as { + subscriberId?: string; + iceServers?: RTCIceServer[]; + }; + // Viewers answer to the senderId in our offers; the stream only + // delivers signals targeted at this server-assigned ID. + hostSignalIdRef.current = data.subscriberId ?? hostId; if (data.iceServers && data.iceServers.length > 0) { iceServersRef.current = data.iceServers; } } catch { - // Use default ICE servers + // Use default ICE servers and the local host identity + hostSignalIdRef.current = hostId; } + + // Start the liveness ping now that the room is actually being hosted. + // Stamp immediately so the room appears live without a 30s wait. + if (livenessIntervalRef.current) { + clearInterval(livenessIntervalRef.current); + } + void sendLivenessHeartbeat(generation); + livenessIntervalRef.current = setInterval( + () => void sendLivenessHeartbeat(generation), + HOST_LIVENESS_INTERVAL + ); }); eventSource.addEventListener('heartbeat', () => { @@ -832,7 +916,7 @@ export function useWebRTCHost({ presences: { user_id: string; role: string }[]; }; for (const presence of presences) { - if (presence.role === 'viewer' && presence.user_id !== hostId) { + if (presence.role === 'viewer' && presence.user_id !== hostSignalIdRef.current) { void handleViewerJoin(presence.user_id, generation); } } @@ -860,6 +944,15 @@ export function useWebRTCHost({ console.error('[WebRTCHost] SSE error'); isStartingRef.current = false; setError('Connection to server lost. Reconnecting...'); + // The wrapper disables the library's stale-header auto-retry, so + // rebuild the transport ourselves with a freshly refreshed token. + if (sseReconnectTimerRef.current) return; + sseReconnectTimerRef.current = setTimeout(() => { + sseReconnectTimerRef.current = null; + if (!isCurrentEventSource()) return; + stopHostingRef.current?.({ preservePublishedStream: true }); + void startHostingRef.current?.(); + }, SSE_RECONNECT_DELAY); }); // Start stats reporting @@ -871,7 +964,7 @@ export function useWebRTCHost({ if (Date.now() - lastHeartbeatAtRef.current <= HEARTBEAT_TIMEOUT) return; setError('Connection heartbeat timed out. Reconnecting...'); - stopHostingRef.current?.(); + stopHostingRef.current?.({ preservePublishedStream: true }); void startHostingRef.current?.(); }, HEARTBEAT_WATCHDOG_INTERVAL); }, [ @@ -882,12 +975,15 @@ export function useWebRTCHost({ isCurrentGeneration, removeViewer, reportStats, + sendLivenessHeartbeat, ]); startHostingRef.current = startHosting; - // Stop hosting - const stopHosting = useCallback(() => { + // Stop hosting. The automatic restart paths pass preservePublishedStream + // so the still-live capture tracks re-attach to the rebuilt connections; + // user-initiated stops, background suspends, and unmount clear everything. + const stopHostingInternal = useCallback((options?: StopHostingOptions) => { console.log('[WebRTCHost] Stopping hosting'); generationRef.current += 1; isStartingRef.current = false; @@ -902,6 +998,16 @@ export function useWebRTCHost({ heartbeatWatchdogRef.current = null; } + if (livenessIntervalRef.current) { + clearInterval(livenessIntervalRef.current); + livenessIntervalRef.current = null; + } + + if (sseReconnectTimerRef.current) { + clearTimeout(sseReconnectTimerRef.current); + sseReconnectTimerRef.current = null; + } + if (eventSourceRef.current) { eventSourceRef.current.close(); eventSourceRef.current = null; @@ -914,7 +1020,9 @@ export function useWebRTCHost({ publishedStreamSendersRef.current.clear(); pendingCandidatesRef.current.clear(); publishedStreamVersionRef.current += 1; - localStreamRef.current = null; + if (!options?.preservePublishedStream) { + localStreamRef.current = null; + } if (hostMicStreamRef.current) { hostMicStreamRef.current.getTracks().forEach((track) => { @@ -933,7 +1041,11 @@ export function useWebRTCHost({ } }, []); - stopHostingRef.current = stopHosting; + stopHostingRef.current = stopHostingInternal; + + const stopHosting = useCallback(() => { + stopHostingInternal(); + }, [stopHostingInternal]); const removePublishedSenders = useCallback((viewer: ViewerConnection): boolean => { const publishedSenders = publishedStreamSendersRef.current.get(viewer.id); diff --git a/apps/mobile/src/hooks/useWebRTCViewer.test.ts b/apps/mobile/src/hooks/useWebRTCViewer.test.ts index 5571b1cd..2ed5d266 100644 --- a/apps/mobile/src/hooks/useWebRTCViewer.test.ts +++ b/apps/mobile/src/hooks/useWebRTCViewer.test.ts @@ -2,23 +2,24 @@ import { describe, it, expect, vi, beforeEach } from 'vitest'; 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 { getValidAccessToken } from '../lib/auth-session'; import { mediaDevices } from 'react-native-webrtc'; import type { MediaStream, MediaStreamTrack } from 'react-native-webrtc'; import { emitAppStateChange, mockPeerConnections } from '../test/setup'; +import { + viewerConnectedEventData, + AUTH_USER_ID, + PARTICIPANT_ROW_ID, + HOST_USER_ID, + OTHER_VIEWER_ID, +} from '../test/fixtures/server-contracts'; vi.mock('../config', () => ({ API_BASE_URL: 'https://pairux.com', })); -vi.mock('../lib/secure-storage', () => ({ - getStoredAuth: vi.fn().mockResolvedValue({ - accessToken: 'test-token', - refreshToken: 'refresh', - expiresAt: Date.now() + 3600000, - user: { id: 'viewer-1', email: 'viewer@example.com' }, - }), - isAuthExpired: vi.fn().mockReturnValue(false), +vi.mock('../lib/auth-session', () => ({ + getValidAccessToken: vi.fn().mockResolvedValue('test-token'), })); const { mockClose, mockAddEventListener, mockEventSources } = vi.hoisted(() => ({ @@ -31,7 +32,7 @@ const { mockClose, mockAddEventListener, mockEventSources } = vi.hoisted(() => ( })); vi.mock('../lib/event-source', () => ({ - createEventSource: vi.fn(() => { + createEventSource: vi.fn((_url: string, _options?: { headers?: Record }) => { const listeners = new Map void>(); const source = { listeners, @@ -48,10 +49,19 @@ vi.mock('../lib/event-source', () => ({ }), })); +function deferred() { + let resolve!: (value: T | PromiseLike) => void; + const promise = new Promise((resolvePromise) => { + resolve = resolvePromise; + }); + return { promise, resolve }; +} + describe('useWebRTCViewer', () => { beforeEach(() => { vi.clearAllMocks(); mockEventSources.length = 0; + vi.mocked(getValidAccessToken).mockResolvedValue('test-token'); vi.mocked(fetch).mockResolvedValue({ ok: true, text: async () => 'ok', @@ -106,7 +116,8 @@ describe('useWebRTCViewer', () => { }); expect(createEventSource).toHaveBeenCalledWith( - expect.stringContaining('/api/sessions/session-1/signal/stream') + expect.stringContaining('/api/sessions/session-1/signal/stream'), + expect.anything() ); }); @@ -123,11 +134,12 @@ describe('useWebRTCViewer', () => { }); expect(createEventSource).toHaveBeenCalledWith( - expect.stringContaining('participantId=viewer-42') + expect.stringContaining('participantId=viewer-42'), + expect.anything() ); }); - it('should include auth token in SSE URL params', async () => { + it('sends the SSE bearer token in the Authorization header, not the URL', async () => { renderHook(() => useWebRTCViewer({ sessionId: 'session-1', @@ -139,7 +151,12 @@ describe('useWebRTCViewer', () => { await new Promise((r) => setTimeout(r, 0)); }); - expect(createEventSource).toHaveBeenCalledWith(expect.stringContaining('token=test-token')); + const [url, options] = vi.mocked(createEventSource).mock.calls[0] as [ + string, + { headers?: Record } | undefined, + ]; + expect(url).not.toContain('token='); + expect(options).toEqual({ headers: { Authorization: 'Bearer test-token' } }); }); it('echoes the host negotiation ID in its answer', async () => { @@ -170,6 +187,7 @@ describe('useWebRTCViewer', () => { type: 'offer', sdp: 'mobile-host-offer', senderId: 'host-1', + targetId: 'viewer-1', negotiationId: 'mobile-offer-1', timestamp: Date.now(), }), @@ -189,7 +207,7 @@ describe('useWebRTCViewer', () => { }); it('should set error when not authenticated', async () => { - vi.mocked(getStoredAuth).mockResolvedValueOnce(null); + vi.mocked(getValidAccessToken).mockResolvedValueOnce(null); const { result } = renderHook(() => useWebRTCViewer({ @@ -349,6 +367,7 @@ describe('useWebRTCViewer', () => { type: 'offer', sdp: 'mobile-host-offer', senderId: 'host-1', + targetId: 'viewer-1', negotiationId: 'inactive-offer-1', timestamp: Date.now(), }), @@ -665,4 +684,433 @@ describe('useWebRTCViewer', () => { vi.useRealTimers(); } }); + + describe('signaling identity and cross-viewer guards', () => { + const requestBody = (init?: RequestInit): string => + typeof init?.body === 'string' ? init.body : ''; + + /** Connects the SSE stream and delivers one targeted host offer. */ + async function connectAndReceiveHostOffer() { + renderHook(() => + useWebRTCViewer({ + sessionId: 'session-1', + participantId: PARTICIPANT_ROW_ID, + }) + ); + await act(async () => { + await Promise.resolve(); + await Promise.resolve(); + }); + const source = mockEventSources[0]; + expect(source).toBeDefined(); + + act(() => { + source.listeners.get('connected')?.({ data: viewerConnectedEventData() }); + }); + const peer = mockPeerConnections[0]; + expect(peer).toBeDefined(); + + await act(async () => { + source.listeners.get('signal')?.({ + data: JSON.stringify({ + type: 'offer', + sdp: 'host-offer-sdp', + senderId: HOST_USER_ID, + targetId: AUTH_USER_ID, + timestamp: Date.now(), + }), + }); + await Promise.resolve(); + await Promise.resolve(); + await Promise.resolve(); + }); + + return { source, peer }; + } + + it('adopts the server-assigned subscriberId for outgoing answers', async () => { + const { peer } = await connectAndReceiveHostOffer(); + + 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(`"senderId":"${AUTH_USER_ID}"`), + }) + ) + ); + // The answer goes back to the host, not broadcast + const answerCall = vi + .mocked(fetch) + .mock.calls.find(([, init]) => requestBody(init).includes('"type":"answer"')); + expect(answerCall).toBeDefined(); + expect(requestBody(answerCall?.[1])).toContain(`"targetId":"${HOST_USER_ID}"`); + // The participant row id must not leak into signaling identities + expect(requestBody(answerCall?.[1])).not.toContain(PARTICIPANT_ROW_ID); + }); + + it("ignores another viewer's broadcast ICE-restart offer", async () => { + const { source, peer } = await connectAndReceiveHostOffer(); + expect(peer.setRemoteDescription).toHaveBeenCalledTimes(1); + + await act(async () => { + source.listeners.get('signal')?.({ + data: JSON.stringify({ + type: 'offer', + sdp: 'other-viewer-restart-offer', + senderId: OTHER_VIEWER_ID, + // no targetId: exactly how a viewer ICE restart is broadcast + timestamp: Date.now(), + }), + }); + await Promise.resolve(); + await Promise.resolve(); + await Promise.resolve(); + }); + + expect(peer.setRemoteDescription).toHaveBeenCalledTimes(1); + }); + + it("drops another viewer's broadcast ICE candidates but accepts the host's targeted ones", async () => { + const { source, peer } = await connectAndReceiveHostOffer(); + + await act(async () => { + source.listeners.get('signal')?.({ + data: JSON.stringify({ + type: 'ice-candidate', + candidate: { candidate: 'other-viewer-candidate' }, + senderId: OTHER_VIEWER_ID, + timestamp: Date.now(), + }), + }); + source.listeners.get('signal')?.({ + data: JSON.stringify({ + type: 'ice-candidate', + candidate: { candidate: 'host-candidate' }, + senderId: HOST_USER_ID, + targetId: AUTH_USER_ID, + timestamp: Date.now(), + }), + }); + await Promise.resolve(); + await Promise.resolve(); + await Promise.resolve(); + }); + + expect(peer.addIceCandidate).toHaveBeenCalledTimes(1); + expect(peer.addIceCandidate).toHaveBeenCalledWith( + expect.objectContaining({ candidate: 'host-candidate' }) + ); + }); + + it('targets outgoing ICE candidates at the host with the subscriber identity', async () => { + const { peer } = await connectAndReceiveHostOffer(); + + const iceListener = peer.addEventListener.mock.calls.find( + ([eventName]) => eventName === 'icecandidate' + )?.[1] as ((event: { candidate: { toJSON: () => unknown } | null }) => void) | undefined; + expect(iceListener).toBeDefined(); + + await act(async () => { + iceListener?.({ candidate: { toJSON: () => ({ candidate: 'local-candidate' }) } }); + await Promise.resolve(); + await Promise.resolve(); + }); + + const candidateCall = vi + .mocked(fetch) + .mock.calls.find(([, init]) => requestBody(init).includes('"type":"ice-candidate"')); + expect(candidateCall).toBeDefined(); + const body = requestBody(candidateCall?.[1]); + expect(body).toContain(`"senderId":"${AUTH_USER_ID}"`); + expect(body).toContain(`"targetId":"${HOST_USER_ID}"`); + }); + + it('refreshes the access token before opening the SSE stream', async () => { + vi.mocked(getValidAccessToken).mockResolvedValueOnce('refreshed-token'); + + renderHook(() => + useWebRTCViewer({ + sessionId: 'session-1', + participantId: PARTICIPANT_ROW_ID, + }) + ); + await act(async () => { + await Promise.resolve(); + await Promise.resolve(); + }); + + const [, options] = vi.mocked(createEventSource).mock.calls[0] as [ + string, + { headers?: Record } | undefined, + ]; + expect(options).toEqual({ headers: { Authorization: 'Bearer refreshed-token' } }); + }); + + it("rejects the host's candidates explicitly targeted at another viewer", async () => { + const { source, peer } = await connectAndReceiveHostOffer(); + + await act(async () => { + source.listeners.get('signal')?.({ + data: JSON.stringify({ + type: 'ice-candidate', + candidate: { candidate: 'for-other-viewer' }, + senderId: HOST_USER_ID, + targetId: OTHER_VIEWER_ID, + timestamp: Date.now(), + }), + }); + await Promise.resolve(); + await Promise.resolve(); + }); + + expect(peer.addIceCandidate).not.toHaveBeenCalled(); + }); + + it("ignores a known viewer's offer even when it targets us", async () => { + renderHook(() => + useWebRTCViewer({ + sessionId: 'session-1', + participantId: PARTICIPANT_ROW_ID, + }) + ); + await act(async () => { + await Promise.resolve(); + await Promise.resolve(); + }); + const source = mockEventSources[0]; + + act(() => { + source.listeners.get('connected')?.({ data: viewerConnectedEventData() }); + // Server presence marks who is host and who is a fellow viewer + source.listeners.get('presence-join')?.({ + data: JSON.stringify({ + presences: [ + { user_id: HOST_USER_ID, role: 'host' }, + { user_id: OTHER_VIEWER_ID, role: 'viewer' }, + ], + }), + }); + }); + const peer = mockPeerConnections[0]; + + await act(async () => { + source.listeners.get('signal')?.({ + data: JSON.stringify({ + type: 'offer', + sdp: 'viewer-crafted-offer', + senderId: OTHER_VIEWER_ID, + targetId: AUTH_USER_ID, + timestamp: Date.now(), + }), + }); + await Promise.resolve(); + await Promise.resolve(); + await Promise.resolve(); + }); + expect(peer.setRemoteDescription).not.toHaveBeenCalled(); + + // The real host still negotiates normally afterwards + await act(async () => { + source.listeners.get('signal')?.({ + data: JSON.stringify({ + type: 'offer', + sdp: 'host-offer-sdp', + senderId: HOST_USER_ID, + targetId: AUTH_USER_ID, + timestamp: Date.now(), + }), + }); + await Promise.resolve(); + await Promise.resolve(); + await Promise.resolve(); + }); + expect(peer.setRemoteDescription).toHaveBeenCalledTimes(1); + expect(peer.setRemoteDescription).toHaveBeenCalledWith({ + type: 'offer', + sdp: 'host-offer-sdp', + }); + }); + + it('does not let an unknown sender displace the host negotiation', async () => { + const { source, peer } = await connectAndReceiveHostOffer(); + expect(peer.setRemoteDescription).toHaveBeenCalledTimes(1); + + await act(async () => { + source.listeners.get('signal')?.({ + data: JSON.stringify({ + type: 'offer', + sdp: 'takeover-offer', + senderId: 'intruder-9999', + targetId: AUTH_USER_ID, + timestamp: Date.now(), + }), + }); + await Promise.resolve(); + await Promise.resolve(); + await Promise.resolve(); + }); + + expect(peer.setRemoteDescription).toHaveBeenCalledTimes(1); + }); + + it("drains only the negotiated host's early candidates after the offer arrives", async () => { + renderHook(() => + useWebRTCViewer({ + sessionId: 'session-1', + participantId: PARTICIPANT_ROW_ID, + }) + ); + await act(async () => { + await Promise.resolve(); + await Promise.resolve(); + }); + const source = mockEventSources[0]; + act(() => { + source.listeners.get('connected')?.({ data: viewerConnectedEventData() }); + }); + const peer = mockPeerConnections[0]; + + await act(async () => { + source.listeners.get('signal')?.({ + data: JSON.stringify({ + type: 'ice-candidate', + candidate: { candidate: 'intruder-early' }, + senderId: 'intruder-9999', + targetId: AUTH_USER_ID, + timestamp: Date.now(), + }), + }); + source.listeners.get('signal')?.({ + data: JSON.stringify({ + type: 'ice-candidate', + candidate: { candidate: 'host-early' }, + senderId: HOST_USER_ID, + targetId: AUTH_USER_ID, + timestamp: Date.now(), + }), + }); + await Promise.resolve(); + await Promise.resolve(); + }); + expect(peer.addIceCandidate).not.toHaveBeenCalled(); + + await act(async () => { + source.listeners.get('signal')?.({ + data: JSON.stringify({ + type: 'offer', + sdp: 'host-offer-sdp', + senderId: HOST_USER_ID, + targetId: AUTH_USER_ID, + timestamp: Date.now(), + }), + }); + await Promise.resolve(); + await Promise.resolve(); + await Promise.resolve(); + }); + + await waitFor(() => expect(peer.addIceCandidate).toHaveBeenCalledTimes(1)); + expect(peer.addIceCandidate).toHaveBeenCalledWith( + expect.objectContaining({ candidate: 'host-early' }) + ); + }); + }); + + describe('credential lifecycle guards', () => { + const requestBodies = () => + vi + .mocked(fetch) + .mock.calls.map(([, init]) => (typeof init?.body === 'string' ? init.body : '')); + + it('drops the answer instead of posting once no valid token can be resolved', async () => { + renderHook(() => + useWebRTCViewer({ + sessionId: 'session-1', + participantId: PARTICIPANT_ROW_ID, + }) + ); + await act(async () => { + await Promise.resolve(); + await Promise.resolve(); + }); + const source = mockEventSources[0]; + act(() => { + source.listeners.get('connected')?.({ data: viewerConnectedEventData() }); + }); + + // The account signs out: refresh yields nothing from here on + vi.mocked(getValidAccessToken).mockResolvedValue(null); + + await act(async () => { + source.listeners.get('signal')?.({ + data: JSON.stringify({ + type: 'offer', + sdp: 'host-offer-sdp', + senderId: HOST_USER_ID, + targetId: AUTH_USER_ID, + timestamp: Date.now(), + }), + }); + await Promise.resolve(); + await Promise.resolve(); + await Promise.resolve(); + }); + + expect(requestBodies().some((body) => body.includes('"type":"answer"'))).toBe(false); + }); + + it('does not post the answer when disconnected while the token refresh is in flight', async () => { + const { result } = renderHook(() => + useWebRTCViewer({ + sessionId: 'session-1', + participantId: PARTICIPANT_ROW_ID, + }) + ); + await act(async () => { + await Promise.resolve(); + await Promise.resolve(); + }); + const source = mockEventSources[0]; + act(() => { + source.listeners.get('connected')?.({ data: viewerConnectedEventData() }); + }); + + const tokenGate = deferred(); + vi.mocked(getValidAccessToken).mockReturnValue(tokenGate.promise); + + await act(async () => { + source.listeners.get('signal')?.({ + data: JSON.stringify({ + type: 'offer', + sdp: 'host-offer-sdp', + senderId: HOST_USER_ID, + targetId: AUTH_USER_ID, + timestamp: Date.now(), + }), + }); + await Promise.resolve(); + await Promise.resolve(); + await Promise.resolve(); + }); + + act(() => { + result.current.disconnect(); + }); + + await act(async () => { + tokenGate.resolve('late-token'); + await Promise.resolve(); + await Promise.resolve(); + }); + + expect(requestBodies().some((body) => body.includes('"type":"answer"'))).toBe(false); + const headers = vi + .mocked(fetch) + .mock.calls.map(([, init]) => JSON.stringify(init?.headers ?? {})); + expect(headers.some((header) => header.includes('late-token'))).toBe(false); + }); + }); }); diff --git a/apps/mobile/src/hooks/useWebRTCViewer.ts b/apps/mobile/src/hooks/useWebRTCViewer.ts index 51789c74..920043c5 100644 --- a/apps/mobile/src/hooks/useWebRTCViewer.ts +++ b/apps/mobile/src/hooks/useWebRTCViewer.ts @@ -29,7 +29,7 @@ import { markTrackAsSpeech, } from '@pairux/shared-types'; import { API_BASE_URL } from '../config'; -import { getStoredAuth, isAuthExpired } from '../lib/secure-storage'; +import { getValidAccessToken } from '../lib/auth-session'; import { createEventSource, type SSEConnection } from '../lib/event-source'; // RN WebRTC's RTCDataChannel type (differs from browser global) @@ -54,6 +54,8 @@ const STATS_INTERVAL = 30000; const STATS_DISPLAY_INTERVAL = 2000; const HEARTBEAT_TIMEOUT = 75000; const HEARTBEAT_WATCHDOG_INTERVAL = 15000; +// Delay before rebuilding the SSE transport (with a fresh token) after an error +const SSE_RECONNECT_DELAY = 3000; const DEFAULT_ICE_SERVERS: RTCIceServer[] = [ { urls: 'stun:stun.l.google.com:19302' }, @@ -120,10 +122,10 @@ export function useWebRTCViewer({ const micStreamRef = useRef(null); const eventSourceRef = useRef(null); const dataChannelRef = useRef(null); - const authTokenRef = useRef(null); const statsIntervalRef = useRef | null>(null); const statsReportIntervalRef = useRef | null>(null); const heartbeatWatchdogRef = useRef | null>(null); + const sseReconnectTimerRef = useRef | null>(null); const lastHeartbeatAtRef = useRef(0); const reconnectAttemptsRef = useRef(0); const inputSequenceRef = useRef(0); @@ -137,8 +139,21 @@ export function useWebRTCViewer({ const maxReconnectAttempts = 3; const iceServersRef = useRef(DEFAULT_ICE_SERVERS); - const pendingCandidatesRef = useRef([]); + // Candidates that arrive before the host's offer, kept with their sender so + // draining can discard anything not from the negotiated host. + const pendingCandidatesRef = useRef<{ senderId: string; candidate: RTCIceCandidateInit }[]>([]); const signalQueueRef = useRef>(Promise.resolve()); + // The SSE `connected` event carries the server-assigned subscriber ID; the + // stream only delivers signals targeted at that ID, so every outgoing + // senderId must use it (for an authenticated user it is the auth user id, + // not the participant row id). Matches the desktop viewer. + const signalSenderIdRef = useRef(participantId); + // The host's signaling ID, learned from the first offer addressed to us. + const hostSenderIdRef = useRef(null); + // Roles by subscriber id from the server's presence events — the one + // authoritative signal for who the session host is (the server assigns + // role 'host' only to the session's host_user_id). + const peerRolesRef = useRef>(new Map()); const handleConnectionFailureRef = useRef< ((generation: number, pc: RTCPeerConnection) => Promise) | undefined @@ -170,16 +185,22 @@ export function useWebRTCViewer({ // Send signal via API const sendSignal = useCallback( - async (signal: SignalMessage) => { + async (signal: SignalMessage, generation: number) => { try { - const headers: Record = { 'Content-Type': 'application/json' }; - if (authTokenRef.current) { - headers.Authorization = `Bearer ${authTokenRef.current}`; + // Resolve the token per request so signaling keeps working after the + // access token expires mid-session (single-flight refresh). No token + // means signed out — never post with a stale cached credential, and + // never post for a generation that stopped while the token resolved. + const token = await getValidAccessToken(); + if (!isCurrentGeneration(generation)) return; + if (!token) { + console.error('[WebRTCViewer] Not authenticated; dropping signal'); + return; } const response = await fetch(`${API_BASE_URL}/api/sessions/${sessionId}/signal`, { method: 'POST', - headers, + headers: { 'Content-Type': 'application/json', Authorization: `Bearer ${token}` }, body: JSON.stringify(signal), }); @@ -190,7 +211,7 @@ export function useWebRTCViewer({ console.error('[WebRTCViewer] Error sending signal:', err); } }, - [sessionId] + [isCurrentGeneration, sessionId] ); // Handle data channel messages @@ -396,14 +417,13 @@ export function useWebRTCViewer({ } }); - const headers: Record = { 'Content-Type': 'application/json' }; - if (authTokenRef.current) { - headers.Authorization = `Bearer ${authTokenRef.current}`; - } + const token = await getValidAccessToken(); + if (!isCurrentGeneration(generation) || peerConnectionRef.current !== pc) return; + if (!token) return; await fetch(`${API_BASE_URL}/api/sessions/${sessionId}/stats`, { method: 'POST', - headers, + headers: { 'Content-Type': 'application/json', Authorization: `Bearer ${token}` }, body: JSON.stringify({ participantId, role: 'viewer', @@ -436,9 +456,33 @@ export function useWebRTCViewer({ const pc = peerConnectionRef.current; if (!pc) return; + // Sender/target guards: the SSE stream forwards untargeted broadcasts to + // every subscriber, so another viewer's ICE-restart offer or ICE + // candidates would otherwise stomp this peer connection. + const selfId = signalSenderIdRef.current; + if (message.senderId === selfId) return; + try { switch (message.type) { case 'offer': { + // The host addresses every offer to a specific viewer; an + // untargeted offer is another viewer's ICE-restart broadcast. + if (message.targetId !== selfId) return; + const senderRole = peerRolesRef.current.get(message.senderId); + // Presence marks the session host authoritatively — a sender the + // server called a viewer can never host-negotiate with us, even + // when it targets us explicitly. + if (senderRole === 'viewer') return; + if ( + hostSenderIdRef.current !== null && + message.senderId !== hostSenderIdRef.current && + senderRole !== 'host' + ) { + // An unknown sender must not displace an ongoing negotiation. + return; + } + hostSenderIdRef.current = message.senderId; + if (pc.signalingState !== 'stable') { console.warn( `[WebRTCViewer] Received offer in ${pc.signalingState} state — rolling back` @@ -456,22 +500,27 @@ export function useWebRTCViewer({ if (!isCurrentGeneration(generation) || peerConnectionRef.current !== pc) return; if (answer.sdp) { - await sendSignal({ - type: 'answer', - sdp: answer.sdp, - senderId: participantId, - targetId: message.senderId, - ...(message.negotiationId ? { negotiationId: message.negotiationId } : {}), - timestamp: Date.now(), - }); + await sendSignal( + { + type: 'answer', + sdp: answer.sdp, + senderId: selfId, + targetId: message.senderId, + ...(message.negotiationId ? { negotiationId: message.negotiationId } : {}), + timestamp: Date.now(), + }, + generation + ); } + if (!isCurrentGeneration(generation) || peerConnectionRef.current !== pc) return; - // Drain buffered ICE candidates + // Drain buffered ICE candidates from the negotiated host only const pending = pendingCandidatesRef.current; if (pending.length > 0) { pendingCandidatesRef.current = []; - for (const candidate of pending) { - await pc.addIceCandidate(new RTCIceCandidate(candidate)); + for (const entry of pending) { + if (entry.senderId !== hostSenderIdRef.current) continue; + await pc.addIceCandidate(new RTCIceCandidate(entry.candidate)); if (!isCurrentGeneration(generation) || peerConnectionRef.current !== pc) return; } } @@ -479,9 +528,22 @@ export function useWebRTCViewer({ } case 'ice-candidate': { + const targetedAtUs = message.targetId === selfId; + // A candidate explicitly addressed to another subscriber is never + // ours — even when it comes from the host we negotiate with. + if (message.targetId && !targetedAtUs) return; + // Viewers never feed candidates into our host connection. + if (peerRolesRef.current.get(message.senderId) === 'viewer') return; + const knownHost = hostSenderIdRef.current; + if (knownHost !== null && message.senderId !== knownHost) return; + if (!targetedAtUs && knownHost === null) return; + if (message.candidate?.candidate) { if (!pc.remoteDescription) { - pendingCandidatesRef.current.push(message.candidate); + pendingCandidatesRef.current.push({ + senderId: message.senderId, + candidate: message.candidate, + }); } else { await pc.addIceCandidate(new RTCIceCandidate(message.candidate)); } @@ -496,7 +558,7 @@ export function useWebRTCViewer({ } } }, - [isCurrentGeneration, participantId, sendSignal] + [isCurrentGeneration, sendSignal] ); // Serialize signal processing @@ -536,12 +598,18 @@ export function useWebRTCViewer({ if (!isCurrentGeneration(generation) || peerConnectionRef.current !== pc) return; if (offer.sdp) { - await sendSignal({ - type: 'offer', - sdp: offer.sdp, - senderId: participantId, - timestamp: Date.now(), - }); + await sendSignal( + { + type: 'offer', + sdp: offer.sdp, + senderId: signalSenderIdRef.current, + // Address the restart at the host so it cannot disturb the + // other viewers' peer connections. + ...(hostSenderIdRef.current ? { targetId: hostSenderIdRef.current } : {}), + timestamp: Date.now(), + }, + generation + ); } } catch { if (isCurrentGeneration(generation)) { @@ -557,7 +625,7 @@ export function useWebRTCViewer({ connectionFailureInFlightRef.current = false; } }, - [isCurrentGeneration, participantId, sendSignal] + [isCurrentGeneration, sendSignal] ); handleConnectionFailureRef.current = handleConnectionFailure; @@ -590,12 +658,17 @@ export function useWebRTCViewer({ 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(), - }); + void sendSignal( + { + type: 'ice-candidate', + candidate: event.candidate.toJSON(), + senderId: signalSenderIdRef.current, + // Address candidates at the host so other viewers never see them. + ...(hostSenderIdRef.current ? { targetId: hostSenderIdRef.current } : {}), + timestamp: Date.now(), + }, + generation + ); } }); @@ -652,7 +725,7 @@ export function useWebRTCViewer({ return pc; }, - [isCurrentGeneration, participantId, sendSignal, setupDataChannel] + [isCurrentGeneration, sendSignal, setupDataChannel] ); // Disconnect @@ -675,6 +748,11 @@ export function useWebRTCViewer({ heartbeatWatchdogRef.current = null; } + if (sseReconnectTimerRef.current) { + clearTimeout(sseReconnectTimerRef.current); + sseReconnectTimerRef.current = null; + } + if (dataChannelRef.current) { dataChannelRef.current.close(); dataChannelRef.current = null; @@ -702,6 +780,8 @@ export function useWebRTCViewer({ lastHeartbeatAtRef.current = 0; pendingCandidatesRef.current = []; signalQueueRef.current = Promise.resolve(); + hostSenderIdRef.current = null; + peerRolesRef.current = new Map(); remoteStreamRef.current = null; if (mountedRef.current) { setRemoteStream(null); @@ -753,16 +833,18 @@ export function useWebRTCViewer({ console.log('[WebRTCViewer] Starting viewer for session:', sessionId); - // Get auth token from secure storage + // Get a valid access token, refreshing through /api/auth/refresh if the + // stored one has expired (e.g. reconnecting after a long time away). + let sseToken: string; try { - const stored = await getStoredAuth(); + const token = await getValidAccessToken(); if (!isCurrentGeneration(generation)) return; - if (!stored || isAuthExpired(stored)) { + if (!token) { isConnectingRef.current = false; setError('Not authenticated. Please log in again.'); return; } - authTokenRef.current = stored.accessToken; + sseToken = token; } catch (err) { console.error('[WebRTCViewer] Failed to get auth token:', err); if (isCurrentGeneration(generation)) { @@ -797,14 +879,14 @@ export function useWebRTCViewer({ setMicEnabled(false); } - // Build SSE URL + // Build SSE URL. The bearer token goes in the Authorization header + // (react-native-sse supports headers, unlike browser EventSource) so + // credentials never appear in proxy/request logs via the query string. const sseParams = new URLSearchParams({ participantId }); - if (authTokenRef.current) { - sseParams.set('token', authTokenRef.current); - } - const sseUrl = `${API_BASE_URL}/api/sessions/${sessionId}/signal/stream?${sseParams.toString()}`; - const eventSource = createEventSource(sseUrl); + const eventSource = createEventSource(sseUrl, { + headers: { Authorization: `Bearer ${sseToken}` }, + }); if (!isCurrentGeneration(generation)) { eventSource.close(); return; @@ -824,13 +906,30 @@ export function useWebRTCViewer({ setError(null); try { - const data = JSON.parse(event.data) as { iceServers?: RTCIceServer[] }; + const data = JSON.parse(event.data) as { + subscriberId?: string; + iceServers?: RTCIceServer[]; + }; + // The server filters targeted signals by this ID, so it must be the + // senderId on everything we post (auth user id, not participant id). + signalSenderIdRef.current = data.subscriberId ?? participantId; + if (data.subscriberId && data.subscriberId !== participantId) { + console.log( + '[WebRTCViewer] Using server-assigned signaling senderId:', + data.subscriberId + ); + } if (data.iceServers && data.iceServers.length > 0) { iceServersRef.current = data.iceServers; } } catch { - // Use default ICE servers + // Use default ICE servers and the local participant identity + signalSenderIdRef.current = participantId; } + // The host will re-offer on this fresh subscription; re-learn who it + // is, along with the presence roles of every subscriber. + hostSenderIdRef.current = null; + peerRolesRef.current = new Map(); // A reconnect can emit another connected event on the same SSE object. // Retire the old peer before attaching a replacement. @@ -886,6 +985,12 @@ export function useWebRTCViewer({ presences: { user_id: string; role: string }[]; }; for (const presence of presences) { + if (presence.user_id) { + peerRolesRef.current.set( + presence.user_id, + presence.role === 'host' ? 'host' : 'viewer' + ); + } if (presence.role === 'host') { console.log('[WebRTCViewer] Host is present:', presence.user_id); } @@ -916,6 +1021,15 @@ export function useWebRTCViewer({ console.error('[WebRTCViewer] SSE error'); isConnectingRef.current = false; setError('Connection to server lost. Reconnecting...'); + // The wrapper disables the library's stale-header auto-retry, so + // rebuild the transport ourselves with a freshly refreshed token. + if (sseReconnectTimerRef.current) return; + sseReconnectTimerRef.current = setTimeout(() => { + sseReconnectTimerRef.current = null; + if (!isCurrentEventSource()) return; + disconnectRef.current?.(); + void initializeRef.current?.(); + }, SSE_RECONNECT_DELAY); }); statsIntervalRef.current = setInterval( diff --git a/apps/mobile/src/lib/api.test.ts b/apps/mobile/src/lib/api.test.ts index d4b7b3c5..156cbfba 100644 --- a/apps/mobile/src/lib/api.test.ts +++ b/apps/mobile/src/lib/api.test.ts @@ -35,9 +35,41 @@ describe('api', () => { expect(token).toBeNull(); }); - it('should return null when auth is expired', async () => { + it('refreshes an expired token through /api/auth/refresh', async () => { vi.mocked(secureStorage.getStoredAuth).mockResolvedValue(mockAuth); vi.mocked(secureStorage.isAuthExpired).mockReturnValue(true); + vi.mocked(fetch).mockResolvedValue({ + ok: true, + json: async () => ({ + data: { + session: { + accessToken: 'refreshed-token', + refreshToken: 'refresh-token-2', + expiresAt: 1_795_000_000, + }, + }, + }), + } as Response); + + const token = await getAuthToken(); + expect(fetch).toHaveBeenCalledWith( + 'https://pairux.com/api/auth/refresh', + expect.objectContaining({ + method: 'POST', + body: JSON.stringify({ refreshToken: 'refresh-token' }), + }) + ); + expect(token).toBe('refreshed-token'); + }); + + it('returns null when the expired token cannot be refreshed', async () => { + vi.mocked(secureStorage.getStoredAuth).mockResolvedValue(mockAuth); + vi.mocked(secureStorage.isAuthExpired).mockReturnValue(true); + vi.mocked(fetch).mockResolvedValue({ + ok: false, + status: 401, + json: async () => ({ error: 'Invalid Refresh Token' }), + } as Response); const token = await getAuthToken(); expect(token).toBeNull(); diff --git a/apps/mobile/src/lib/api.ts b/apps/mobile/src/lib/api.ts index 79e299a6..9e1d2fa2 100644 --- a/apps/mobile/src/lib/api.ts +++ b/apps/mobile/src/lib/api.ts @@ -8,7 +8,7 @@ * with Bearer token from secure store. */ import { API_BASE_URL } from '../config'; -import { getStoredAuth, isAuthExpired } from './secure-storage'; +import { getValidAccessToken } from './auth-session'; export interface ApiResponse { data?: T; @@ -16,9 +16,8 @@ export interface ApiResponse { } export async function getAuthToken(): Promise { - const stored = await getStoredAuth(); - if (!stored || isAuthExpired(stored)) return null; - return stored.accessToken; + // Refreshes through /api/auth/refresh when the stored token is expired. + return getValidAccessToken(); } export async function apiRequest( diff --git a/apps/mobile/src/lib/api/auth.test.ts b/apps/mobile/src/lib/api/auth.test.ts index ab57fa19..5cb8aeb9 100644 --- a/apps/mobile/src/lib/api/auth.test.ts +++ b/apps/mobile/src/lib/api/auth.test.ts @@ -1,45 +1,60 @@ import { describe, it, expect, vi, beforeEach } from 'vitest'; import { authApi } from './auth'; +import { clearAuthSession } from '../auth-session'; import * as secureStorage from '../secure-storage'; +import { + loginSuccessEnvelope, + signupNeedsConfirmationEnvelope, + signupConfirmedEnvelope, + SUPABASE_EXPIRES_AT_SECONDS, + AUTH_USER_ID, +} from '../../test/fixtures/server-contracts'; vi.mock('../secure-storage'); vi.mock('../../config', () => ({ API_BASE_URL: 'https://pairux.com', })); +function deferred() { + let resolve!: (value: T | PromiseLike) => void; + const promise = new Promise((resolvePromise) => { + resolve = resolvePromise; + }); + return { promise, resolve }; +} + describe('authApi', () => { beforeEach(() => { vi.clearAllMocks(); }); describe('login', () => { - it('should send login request and store auth data', async () => { - const loginResponse = { - accessToken: 'token-123', - refreshToken: 'refresh-123', - expiresAt: Date.now() + 3600000, - user: { id: 'user-1', email: 'test@example.com' }, - }; - + it('parses the wrapped { data: { user, session } } envelope and stores ms expiry', async () => { vi.mocked(fetch).mockResolvedValue({ ok: true, - json: async () => loginResponse, + json: async () => loginSuccessEnvelope, } as Response); - const result = await authApi.login('test@example.com', 'password'); + const result = await authApi.login('user@example.com', 'password'); expect(fetch).toHaveBeenCalledWith('https://pairux.com/api/auth/login', { method: 'POST', headers: { 'Content-Type': 'application/json' }, - body: JSON.stringify({ email: 'test@example.com', password: 'password' }), + body: JSON.stringify({ email: 'user@example.com', password: 'password' }), }); expect(secureStorage.storeAuth).toHaveBeenCalledWith({ - accessToken: 'token-123', - refreshToken: 'refresh-123', - expiresAt: loginResponse.expiresAt, - user: { id: 'user-1', email: 'test@example.com' }, + accessToken: 'access-token-1', + refreshToken: 'refresh-token-1', + // Supabase returns expires_at in seconds; storage keeps milliseconds + expiresAt: SUPABASE_EXPIRES_AT_SECONDS * 1000, + user: { id: AUTH_USER_ID, email: 'user@example.com' }, }); - expect(result.data).toEqual(expect.objectContaining({ accessToken: 'token-123' })); + expect(result.data).toEqual( + expect.objectContaining({ + accessToken: 'access-token-1', + expiresAt: SUPABASE_EXPIRES_AT_SECONDS * 1000, + }) + ); }); it('should return error for failed login', async () => { @@ -53,24 +68,94 @@ describe('authApi', () => { expect(secureStorage.storeAuth).not.toHaveBeenCalled(); }); + it('rejects an OK response without the expected envelope', async () => { + vi.mocked(fetch).mockResolvedValue({ + ok: true, + // Flat legacy shape — not what the server sends + json: async () => ({ accessToken: 'x', refreshToken: 'y', expiresAt: 1, user: {} }), + } as Response); + + const result = await authApi.login('user@example.com', 'password'); + expect(result.error).toBe('Invalid response from server'); + expect(secureStorage.storeAuth).not.toHaveBeenCalled(); + }); + it('should return network error on fetch failure', async () => { vi.mocked(fetch).mockRejectedValue(new Error('Network error')); - const result = await authApi.login('test@example.com', 'password'); + const result = await authApi.login('user@example.com', 'password'); expect(result.error).toBe('Network error'); }); + + it('rejects a login envelope with an empty access token', async () => { + vi.mocked(fetch).mockResolvedValue({ + ok: true, + json: async () => ({ + data: { + user: loginSuccessEnvelope.data.user, + session: { accessToken: '', refreshToken: 'r', expiresAt: SUPABASE_EXPIRES_AT_SECONDS }, + }, + }), + } as Response); + + const result = await authApi.login('user@example.com', 'password'); + expect(result.error).toBe('Invalid response from server'); + expect(secureStorage.storeAuth).not.toHaveBeenCalled(); + }); + + it('rejects a login envelope with a missing expiry (would store NaN)', async () => { + vi.mocked(fetch).mockResolvedValue({ + ok: true, + json: async () => ({ + data: { + user: loginSuccessEnvelope.data.user, + session: { accessToken: 'a', refreshToken: 'r' }, + }, + }), + } as Response); + + const result = await authApi.login('user@example.com', 'password'); + expect(result.error).toBe('Invalid response from server'); + expect(secureStorage.storeAuth).not.toHaveBeenCalled(); + }); + + it('fails instead of storing when a logout wins the race against the login', async () => { + const gate = deferred(); + vi.mocked(fetch).mockReturnValueOnce(gate.promise); + + const pending = authApi.login('user@example.com', 'password'); + await clearAuthSession(); + gate.resolve({ ok: true, json: async () => loginSuccessEnvelope } as Response); + + const result = await pending; + expect(result.error).toBe('Sign-in was interrupted. Please try again.'); + expect(secureStorage.storeAuth).not.toHaveBeenCalled(); + }); + + it('returns an error when the login session cannot be persisted', async () => { + vi.mocked(fetch).mockResolvedValue({ + ok: true, + json: async () => loginSuccessEnvelope, + } as Response); + vi.mocked(secureStorage.storeAuth).mockRejectedValueOnce(new Error('keystore unavailable')); + + const result = await authApi.login('user@example.com', 'password'); + expect(result.error).toBe('keystore unavailable'); + expect(result.data).toBeUndefined(); + }); }); describe('signup', () => { - it('should send signup request', async () => { + it('sends confirmPassword to the server and parses the wrapped envelope', async () => { vi.mocked(fetch).mockResolvedValue({ ok: true, - json: async () => ({ message: 'Account created' }), + json: async () => signupNeedsConfirmationEnvelope, } as Response); const result = await authApi.signup({ email: 'new@example.com', - password: 'password', + password: 'Password1', + confirmPassword: 'Password1', firstName: 'Jane', lastName: 'Doe', }); @@ -80,27 +165,51 @@ describe('authApi', () => { headers: { 'Content-Type': 'application/json' }, body: JSON.stringify({ email: 'new@example.com', - password: 'password', + password: 'Password1', + confirmPassword: 'Password1', firstName: 'Jane', lastName: 'Doe', }), }); - expect(result.data?.message).toBe('Account created'); + expect(result.data).toEqual({ + message: 'Check your email to confirm your account', + needsConfirmation: true, + }); + }); + + it('reports when the account is created without email confirmation', async () => { + vi.mocked(fetch).mockResolvedValue({ + ok: true, + json: async () => signupConfirmedEnvelope, + } as Response); + + const result = await authApi.signup({ + email: 'new@example.com', + password: 'Password1', + confirmPassword: 'Password1', + firstName: 'Jane', + lastName: 'Doe', + }); + expect(result.data).toEqual({ + message: 'Account created successfully', + needsConfirmation: false, + }); }); it('should return error for failed signup', async () => { vi.mocked(fetch).mockResolvedValue({ ok: false, - json: async () => ({ error: 'Email already in use' }), + json: async () => ({ error: 'Passwords do not match' }), } as Response); const result = await authApi.signup({ email: 'existing@example.com', - password: 'password', + password: 'Password1', + confirmPassword: 'Password2', firstName: 'Jane', lastName: 'Doe', }); - expect(result.error).toBe('Email already in use'); + expect(result.error).toBe('Passwords do not match'); }); }); @@ -110,7 +219,7 @@ describe('authApi', () => { accessToken: 'token', refreshToken: 'refresh', expiresAt: Date.now() + 3600000, - user: { id: '1', email: 'test@example.com' }, + user: { id: AUTH_USER_ID, email: 'user@example.com' }, }); vi.mocked(secureStorage.isAuthExpired).mockReturnValue(false); @@ -130,5 +239,26 @@ describe('authApi', () => { await authApi.logout(); expect(secureStorage.clearStoredAuth).toHaveBeenCalled(); }); + + it('completes locally with the captured token even when the revoke request hangs', async () => { + vi.mocked(secureStorage.getStoredAuth).mockResolvedValue({ + accessToken: 'expired-access', + refreshToken: 'refresh', + expiresAt: Date.now() - 1000, + user: { id: AUTH_USER_ID, email: 'user@example.com' }, + }); + vi.mocked(secureStorage.isAuthExpired).mockReturnValue(true); + // The revoke request never answers; local logout must not wait for it + vi.mocked(fetch).mockReturnValue(new Promise(() => undefined)); + + await authApi.logout(); + + expect(secureStorage.clearStoredAuth).toHaveBeenCalledTimes(1); + expect(fetch).toHaveBeenCalledTimes(1); + const [url, init] = vi.mocked(fetch).mock.calls[0]; + expect(url).toBe('https://pairux.com/api/auth/logout'); + // The token is captured as-is: signing out must never mint new tokens + expect((init?.headers as Record).Authorization).toBe('Bearer expired-access'); + }); }); }); diff --git a/apps/mobile/src/lib/api/auth.ts b/apps/mobile/src/lib/api/auth.ts index e791a212..9f6fa732 100644 --- a/apps/mobile/src/lib/api/auth.ts +++ b/apps/mobile/src/lib/api/auth.ts @@ -3,24 +3,43 @@ * * Handles login, signup, and logout via the cloud API. * Stores/clears tokens in secure storage. + * + * The server wraps every success payload as `{ data: ... }` + * (apps/web/src/lib/api.ts successResponse) and returns Supabase's + * `expires_at` in SECONDS; storage keeps milliseconds. */ import { API_BASE_URL } from '../../config'; -import { storeAuth, clearStoredAuth, type StoredAuth } from '../secure-storage'; +import type { StoredAuth } from '../secure-storage'; +import { + beginAuthMutation, + endAuthSession, + commitLoginSession, + parseSessionEnvelope, +} from '../auth-session'; import { apiRequest } from '../api'; interface LoginResponse { - accessToken: string; - refreshToken: string; - expiresAt: number; - user: { id: string; email: string }; + data?: { + user: { id: string; email: string }; + session: { accessToken: string; refreshToken: string; expiresAt: number }; + }; + error?: string; } interface SignupResponse { - message: string; + data?: { + user?: { id: string; email?: string }; + message: string; + needsConfirmation: boolean; + }; + error?: string; } export const authApi = { async login(email: string, password: string): Promise<{ data?: StoredAuth; error?: string }> { + // A newer login intent or logout invalidates this attempt, even before + // the newer network request has completed. + const startedEpoch = beginAuthMutation(); try { const response = await fetch(`${API_BASE_URL}/api/auth/login`, { method: 'POST', @@ -28,20 +47,31 @@ export const authApi = { body: JSON.stringify({ email, password }), }); - const data = (await response.json()) as LoginResponse & { error?: string }; + const data = (await response.json()) as LoginResponse; if (!response.ok) { return { error: data.error ?? 'Failed to sign in' }; } + // Defensive: the fields are non-optional in the type, but an OK + // response from a proxy or an older server may not carry them — + // and empty/NaN token fields must never reach secure storage. + const payload = data.data as Partial> | undefined; + const session = parseSessionEnvelope(payload?.session); + const user = payload?.user as Partial<{ id: string; email: string }> | undefined; + if (!session || typeof user?.id !== 'string' || user.id.length === 0) { + return { error: 'Invalid response from server' }; + } + const auth: StoredAuth = { - accessToken: data.accessToken, - refreshToken: data.refreshToken, - expiresAt: data.expiresAt, - user: data.user, + ...session, + user: { id: user.id, email: typeof user.email === 'string' ? user.email : email }, }; - await storeAuth(auth); + const committed = await commitLoginSession(auth, startedEpoch); + if (!committed) { + return { error: 'Sign-in was interrupted. Please try again.' }; + } return { data: auth }; } catch (error) { return { error: error instanceof Error ? error.message : 'Network error' }; @@ -51,9 +81,10 @@ export const authApi = { async signup(params: { email: string; password: string; + confirmPassword: string; firstName: string; lastName: string; - }): Promise<{ data?: SignupResponse; error?: string }> { + }): Promise<{ data?: { message: string; needsConfirmation: boolean }; error?: string }> { try { const response = await fetch(`${API_BASE_URL}/api/auth/signup`, { method: 'POST', @@ -61,25 +92,51 @@ export const authApi = { body: JSON.stringify(params), }); - const data = (await response.json()) as SignupResponse & { error?: string }; + const data = (await response.json()) as SignupResponse; if (!response.ok) { return { error: data.error ?? 'Failed to sign up' }; } - return { data }; + if (!data.data) { + return { error: 'Invalid response from server' }; + } + + return { + data: { + message: data.data.message, + needsConfirmation: data.data.needsConfirmation, + }, + }; } catch (error) { return { error: error instanceof Error ? error.message : 'Network error' }; } }, async logout(): Promise { - try { - await apiRequest('/api/auth/logout', { method: 'POST' }); - } catch { - // Best-effort — always clear local tokens + // Capture the current token as-is for the best-effort server-side + // revoke. Never refresh here: signing out must not mint new tokens, + // and local logout must complete even when the network stalls. + const previous = await endAuthSession(); + const token = previous?.accessToken; + + if (token) { + const controller = new AbortController(); + const timer = setTimeout(() => { + controller.abort(); + }, 10000); + void fetch(`${API_BASE_URL}/api/auth/logout`, { + method: 'POST', + headers: { 'Content-Type': 'application/json', Authorization: `Bearer ${token}` }, + signal: controller.signal, + }) + .catch(() => { + // Best-effort revoke; the local session is already gone. + }) + .finally(() => { + clearTimeout(timer); + }); } - await clearStoredAuth(); }, async getSession(): Promise<{ data?: { user: { id: string; email: string } }; error?: string }> { diff --git a/apps/mobile/src/lib/api/sessions.test.ts b/apps/mobile/src/lib/api/sessions.test.ts index 0201b602..e33895c5 100644 --- a/apps/mobile/src/lib/api/sessions.test.ts +++ b/apps/mobile/src/lib/api/sessions.test.ts @@ -1,6 +1,14 @@ import { describe, it, expect, vi, beforeEach } from 'vitest'; -import { sessionApi } from './sessions'; +import { sessionApi, isScheduledLookup } from './sessions'; import * as secureStorage from '../secure-storage'; +import { + joinLookupLiveEnvelope, + joinLookupScheduledEnvelope, + joinParticipantEnvelope, + AUTH_USER_ID, + PARTICIPANT_ROW_ID, + SESSION_ID, +} from '../../test/fixtures/server-contracts'; vi.mock('../secure-storage'); vi.mock('../../config', () => ({ @@ -100,7 +108,7 @@ describe('sessionApi', () => { it('should POST to join a session by code', async () => { vi.mocked(fetch).mockResolvedValue({ ok: true, - json: async () => ({ data: { id: 'participant-1', user_id: 'user-1' } }), + json: async () => joinParticipantEnvelope, } as Response); const result = await sessionApi.join('ABC123', 'Jane'); @@ -111,15 +119,18 @@ describe('sessionApi', () => { body: JSON.stringify({ displayName: 'Jane' }), }) ); - expect(result.data).toEqual({ id: 'participant-1', user_id: 'user-1' }); + // The participant row id is NOT the authenticated user id + expect(result.data?.id).toBe(PARTICIPANT_ROW_ID); + expect(result.data?.user_id).toBe(AUTH_USER_ID); + expect(result.data?.session_id).toBe(SESSION_ID); }); }); describe('lookup', () => { - it('should GET session info by join code', async () => { + it('parses the direct session payload the route returns (no wrapper)', async () => { vi.mocked(fetch).mockResolvedValue({ ok: true, - json: async () => ({ data: { session: { id: 'session-1' } } }), + json: async () => joinLookupLiveEnvelope, } as Response); const result = await sessionApi.lookup('ABC123'); @@ -131,7 +142,41 @@ describe('sessionApi', () => { }), }) ); - expect(result.data?.session).toEqual({ id: 'session-1' }); + expect(result.data).toBeDefined(); + if (!result.data || isScheduledLookup(result.data)) { + throw new Error('Expected a live session lookup result'); + } + expect(result.data.id).toBe(SESSION_ID); + expect(result.data.join_code).toBe('ABC123'); + expect(result.data.status).toBe('active'); + expect(result.data.participant_count).toBe(2); + }); + + it('recognizes the scheduled-meeting payload', async () => { + vi.mocked(fetch).mockResolvedValue({ + ok: true, + json: async () => joinLookupScheduledEnvelope, + } as Response); + + const result = await sessionApi.lookup('ABC123'); + expect(result.data).toBeDefined(); + if (!result.data || !isScheduledLookup(result.data)) { + throw new Error('Expected a scheduled lookup result'); + } + expect(result.data.title).toBe('Team Standup'); + expect(result.data.scheduled_at).toBe('2026-09-08T10:00:00.000Z'); + }); + + it('surfaces the 404 error for an unknown code', async () => { + vi.mocked(fetch).mockResolvedValue({ + ok: false, + status: 404, + json: async () => ({ error: 'Session not found or has ended' }), + } as Response); + + const result = await sessionApi.lookup('ZZZZZZ'); + expect(result.error).toBe('Session not found or has ended'); + expect(result.data).toBeUndefined(); }); }); }); diff --git a/apps/mobile/src/lib/api/sessions.ts b/apps/mobile/src/lib/api/sessions.ts index 7150451d..0d6a292d 100644 --- a/apps/mobile/src/lib/api/sessions.ts +++ b/apps/mobile/src/lib/api/sessions.ts @@ -2,10 +2,41 @@ * Session API module. * * Port of apps/desktop/src/renderer/lib/api.ts sessionApi. + * + * GET /api/sessions/join/[joinCode] returns the session payload directly + * (or a scheduled-meeting payload discriminated by `scheduled: true`), + * matching the desktop session:lookup handler in + * apps/desktop/src/main/ipc/session.ts — there is no `{ session }` wrapper. */ import type { Session, SessionParticipant } from '@pairux/shared-types'; import { apiRequest } from '../api'; +export interface JoinLookupSession { + id: string; + join_code: string; + status: string; + settings: { quality?: string; allowControl?: boolean; maxParticipants?: number }; + created_at: string; + participant_count: number; +} + +export interface JoinLookupScheduled { + scheduled: true; + id: string; + join_code: string; + title: string; + description: string | null; + scheduled_at: string; + duration_minutes: number; + invitees: { name: string | null; rsvp_status: string }[]; +} + +export type JoinLookupResult = JoinLookupSession | JoinLookupScheduled; + +export function isScheduledLookup(result: JoinLookupResult): result is JoinLookupScheduled { + return 'scheduled' in result && result.scheduled; +} + export const sessionApi = { async create(settings?: { allowGuestControl?: boolean; maxParticipants?: number }) { return apiRequest('/api/sessions', { @@ -38,6 +69,6 @@ export const sessionApi = { }, async lookup(joinCode: string) { - return apiRequest<{ session: Session }>(`/api/sessions/join/${joinCode}`); + return apiRequest(`/api/sessions/join/${joinCode}`); }, }; diff --git a/apps/mobile/src/lib/auth-session-edge-cases.test.ts b/apps/mobile/src/lib/auth-session-edge-cases.test.ts new file mode 100644 index 00000000..eef6d332 --- /dev/null +++ b/apps/mobile/src/lib/auth-session-edge-cases.test.ts @@ -0,0 +1,134 @@ +import { beforeEach, describe, expect, it, vi } from 'vitest'; +import * as SecureStore from 'expo-secure-store'; +import { + beginAuthMutation, + clearAuthSession, + commitLoginSession, + getValidAccessToken, + parseSessionEnvelope, + readAuthSession, + refreshAuthSession, +} from './auth-session'; +import { storeAuth } from './secure-storage'; +import { authApi } from './api/auth'; +import { getStoredAuth } from './secure-storage'; + +function deferred() { + let resolve!: (value: T) => void; + const promise = new Promise((done) => { + resolve = done; + }); + return { promise, resolve }; +} + +const auth = { + accessToken: 'fixture-access', + refreshToken: 'fixture-refresh', + expiresAt: Date.now() + 3600000, + user: { id: 'fixture-user', email: 'fixture@example.com' }, +}; + +describe('session read boundaries', () => { + beforeEach(async () => { + await clearAuthSession(); + await storeAuth(auth); + vi.clearAllMocks(); + }); + + it('does not return an access token from a storage read completed after logout', async () => { + const gate = deferred(); + vi.mocked(SecureStore.getItemAsync).mockImplementationOnce(() => gate.promise); + const pending = getValidAccessToken(); + await vi.waitFor(() => expect(SecureStore.getItemAsync).toHaveBeenCalledTimes(1)); + await clearAuthSession(); + gate.resolve(JSON.stringify(auth)); + expect(await pending).toBeNull(); + }); + + it('treats a refresh rate limit as temporary, not a revoked session', async () => { + vi.mocked(fetch).mockResolvedValueOnce({ ok: false, status: 429 } as Response); + expect(await refreshAuthSession()).toMatchObject({ auth: null, failure: 'transient' }); + }); + + it('rejects an expiry that overflows during conversion to milliseconds', () => { + expect( + parseSessionEnvelope({ + accessToken: 'fixture-access', + refreshToken: 'fixture-refresh', + expiresAt: Number.MAX_VALUE, + }) + ).toBeNull(); + }); + + it('lets the newest login intent win even when the earlier response arrives first', async () => { + const first = deferred(); + const second = deferred(); + vi.mocked(fetch).mockReturnValueOnce(first.promise).mockReturnValueOnce(second.promise); + const loginA = authApi.login('a@example.com', 'fixture-password'); + const loginB = authApi.login('b@example.com', 'fixture-password'); + const response = (id: string) => + ({ + ok: true, + json: async () => ({ + data: { + user: { id, email: `${id}@example.com` }, + session: { + accessToken: `${id}-access`, + refreshToken: `${id}-refresh`, + expiresAt: Math.floor(Date.now() / 1000) + 3600, + }, + }, + }), + }) as Response; + first.resolve(response('a')); + expect((await loginA).data).toBeUndefined(); + second.resolve(response('b')); + expect((await loginB).data?.user.id).toBe('b'); + expect((await getStoredAuth())?.user.id).toBe('b'); + }); + + it('does not return a token while the logout deletion is pending', async () => { + const gate = deferred(); + const originalDelete = vi.mocked(SecureStore.deleteItemAsync).getMockImplementation()!; + vi.mocked(SecureStore.deleteItemAsync).mockImplementationOnce(async (key) => { + await gate.promise; + return originalDelete(key); + }); + const logout = clearAuthSession(); + await vi.waitFor(() => expect(SecureStore.deleteItemAsync).toHaveBeenCalledTimes(1)); + const reading = getValidAccessToken(); + gate.resolve(undefined); + await logout; + expect(await reading).toBeNull(); + }); + + it('blocks stored tokens and refresh after logout deletion fails', async () => { + vi.mocked(SecureStore.deleteItemAsync).mockRejectedValueOnce(new Error('Storage unavailable')); + await expect(clearAuthSession()).rejects.toThrow('Storage unavailable'); + expect(await getStoredAuth()).toEqual(auth); + expect(await readAuthSession()).toBeNull(); + expect(await getValidAccessToken()).toBeNull(); + expect(await refreshAuthSession()).toEqual({ auth: null, failure: 'signed-out' }); + expect(fetch).not.toHaveBeenCalled(); + }); + + it('does not unblock a failed logout just because another login was attempted', async () => { + vi.mocked(SecureStore.deleteItemAsync).mockRejectedValueOnce(new Error('Storage unavailable')); + await expect(clearAuthSession()).rejects.toThrow(); + vi.mocked(fetch).mockResolvedValueOnce({ + ok: false, + status: 401, + json: async () => ({ error: 'Invalid credentials' }), + } as Response); + expect((await authApi.login('bad@example.com', 'bad-password')).error).toBeDefined(); + expect(await getValidAccessToken()).toBeNull(); + }); + + it('allows a successfully committed login after logout deletion failed', async () => { + vi.mocked(SecureStore.deleteItemAsync).mockRejectedValueOnce(new Error('Storage unavailable')); + await expect(clearAuthSession()).rejects.toThrow(); + const replacement = { ...auth, accessToken: 'replacement-access' }; + expect(await commitLoginSession(replacement, beginAuthMutation())).toBe(true); + expect(await getValidAccessToken()).toBe('replacement-access'); + }); +}); diff --git a/apps/mobile/src/lib/auth-session-races.test.ts b/apps/mobile/src/lib/auth-session-races.test.ts new file mode 100644 index 00000000..79ef99c7 --- /dev/null +++ b/apps/mobile/src/lib/auth-session-races.test.ts @@ -0,0 +1,112 @@ +import { beforeEach, describe, expect, it, vi } from 'vitest'; +import * as SecureStore from 'expo-secure-store'; +import { secureStoreData } from '../test/setup'; +import { authApi } from './api/auth'; +import { clearAuthSession, getValidAccessToken, refreshAuthToken } from './auth-session'; +import { getStoredAuth, storeAuth } from './secure-storage'; + +function deferred() { + let resolve!: (value: T) => void; + const promise = new Promise((done) => { + resolve = done; + }); + return { promise, resolve }; +} + +const oldAuth = { + accessToken: 'old-access', + refreshToken: 'old-refresh', + expiresAt: 1, + user: { id: 'old-user', email: 'old@example.com' }, +}; +const freshSession = { + accessToken: 'refreshed-old-access', + refreshToken: 'refreshed-old-refresh', + expiresAt: Math.floor(Date.now() / 1000) + 3600, +}; +const response = (data: unknown) => ({ ok: true, json: async () => ({ data }) }) as Response; + +describe('session mutation races', () => { + beforeEach(async () => { + await clearAuthSession(); + await storeAuth(oldAuth); + vi.clearAllMocks(); + }); + + it('does not let an old refresh overwrite a completed login for a different user', async () => { + const gate = deferred(); + vi.mocked(fetch).mockReturnValueOnce(gate.promise); + const pending = refreshAuthToken(); + await vi.waitFor(() => expect(fetch).toHaveBeenCalledTimes(1)); + vi.mocked(fetch).mockResolvedValueOnce( + response({ + user: { id: 'new-user', email: 'new@example.com' }, + session: { ...freshSession, accessToken: 'new-access', refreshToken: 'new-refresh' }, + }) + ); + const login = await authApi.login('new@example.com', 'fixture-password'); + expect(login.error).toBeUndefined(); + gate.resolve(response({ session: freshSession })); + await pending; + expect((await getStoredAuth())?.user.id).toBe('new-user'); + expect((await getStoredAuth())?.accessToken).toBe('new-access'); + }); + + it('serializes logout behind an already-started secure storage write', async () => { + const writeGate = deferred(); + vi.mocked(SecureStore.setItemAsync).mockImplementationOnce(async (key, value) => { + await writeGate.promise; + secureStoreData.set(key, value); + }); + vi.mocked(fetch).mockResolvedValueOnce(response({ session: freshSession })); + const pending = refreshAuthToken(); + await vi.waitFor(() => expect(SecureStore.setItemAsync).toHaveBeenCalledTimes(1)); + const clearing = clearAuthSession(); + writeGate.resolve(undefined); + await Promise.all([pending, clearing]); + expect(await getStoredAuth()).toBeNull(); + }); + + it.each([false, true])( + 'discards a superseded native login write when the newer login succeeds=%s', + async (newerSucceeds) => { + await clearAuthSession(); + vi.clearAllMocks(); + const writeGate = deferred(); + vi.mocked(SecureStore.setItemAsync).mockImplementationOnce(async (key, value) => { + await writeGate.promise; + secureStoreData.set(key, value); + }); + const loginResponse = (id: string) => + response({ + user: { id, email: `${id}@example.com` }, + session: { ...freshSession, accessToken: `${id}-access`, refreshToken: `${id}-refresh` }, + }); + vi.mocked(fetch).mockResolvedValueOnce(loginResponse('older')); + const older = authApi.login('older@example.com', 'fixture-password'); + await vi.waitFor(() => expect(SecureStore.setItemAsync).toHaveBeenCalledTimes(1)); + vi.mocked(fetch).mockResolvedValueOnce( + newerSucceeds + ? loginResponse('newer') + : ({ + ok: false, + status: 401, + json: async () => ({ error: 'Invalid credentials' }), + } as Response) + ); + const newer = authApi.login('newer@example.com', 'fixture-password'); + writeGate.resolve(undefined); + expect((await older).error).toBeDefined(); + const result = await newer; + if (newerSucceeds) { + expect(result.data?.user.id).toBe('newer'); + expect((await getStoredAuth())?.user.id).toBe('newer'); + expect(await getValidAccessToken()).toBe('newer-access'); + } else { + expect(result.error).toBeDefined(); + expect(await getStoredAuth()).toBeNull(); + expect(await getValidAccessToken()).toBeNull(); + } + } + ); +}); diff --git a/apps/mobile/src/lib/auth-session.test.ts b/apps/mobile/src/lib/auth-session.test.ts new file mode 100644 index 00000000..904a149c --- /dev/null +++ b/apps/mobile/src/lib/auth-session.test.ts @@ -0,0 +1,286 @@ +import { describe, it, expect, vi, beforeEach } from 'vitest'; +import * as SecureStore from 'expo-secure-store'; +import { + getValidAccessToken, + getValidAuth, + refreshAuthToken, + refreshAuthSession, + clearAuthSession, + commitLoginSession, + beginAuthMutation, +} from './auth-session'; +import { storeAuth, getStoredAuth, type StoredAuth } from './secure-storage'; +import { + refreshSuccessEnvelope, + SUPABASE_EXPIRES_AT_SECONDS, + AUTH_USER_ID, +} from '../test/fixtures/server-contracts'; + +vi.mock('../config', () => ({ + API_BASE_URL: 'https://pairux.com', +})); + +function makeAuth(overrides: Partial = {}): StoredAuth { + return { + accessToken: 'access-token-1', + refreshToken: 'refresh-token-1', + expiresAt: Date.now() + 3600000, + user: { id: AUTH_USER_ID, email: 'user@example.com' }, + ...overrides, + }; +} + +function deferred() { + let resolve!: (value: T | PromiseLike) => void; + const promise = new Promise((resolvePromise) => { + resolve = resolvePromise; + }); + return { promise, resolve }; +} + +describe('auth-session', () => { + beforeEach(async () => { + vi.clearAllMocks(); + await clearAuthSession(); + }); + + describe('getValidAccessToken', () => { + it('returns the stored token without a refresh when not expired', async () => { + await storeAuth(makeAuth()); + + const token = await getValidAccessToken(); + expect(token).toBe('access-token-1'); + expect(fetch).not.toHaveBeenCalled(); + }); + + it('returns null when signed out', async () => { + const token = await getValidAccessToken(); + expect(token).toBeNull(); + expect(fetch).not.toHaveBeenCalled(); + }); + + it('refreshes an expired token via /api/auth/refresh and stores ms expiry', async () => { + await storeAuth(makeAuth({ expiresAt: Date.now() - 1000 })); + vi.mocked(fetch).mockResolvedValue({ + ok: true, + json: async () => refreshSuccessEnvelope, + } as Response); + + const token = await getValidAccessToken(); + + expect(fetch).toHaveBeenCalledWith('https://pairux.com/api/auth/refresh', { + method: 'POST', + headers: { 'Content-Type': 'application/json' }, + body: JSON.stringify({ refreshToken: 'refresh-token-1' }), + }); + expect(token).toBe('access-token-2'); + + const stored = await getStoredAuth(); + expect(stored).toEqual({ + accessToken: 'access-token-2', + refreshToken: 'refresh-token-2', + // Supabase returns expires_at in seconds; storage keeps milliseconds + expiresAt: SUPABASE_EXPIRES_AT_SECONDS * 1000, + user: { id: AUTH_USER_ID, email: 'user@example.com' }, + }); + }); + + it('deduplicates concurrent refreshes into a single request', async () => { + await storeAuth(makeAuth({ expiresAt: Date.now() - 1000 })); + const gate = deferred(); + vi.mocked(fetch).mockReturnValue(gate.promise); + + const [first, second] = [getValidAccessToken(), getValidAccessToken()]; + gate.resolve({ + ok: true, + json: async () => refreshSuccessEnvelope, + } as Response); + + expect(await first).toBe('access-token-2'); + expect(await second).toBe('access-token-2'); + expect(fetch).toHaveBeenCalledTimes(1); + }); + + it('returns null and keeps stored auth when the refresh is rejected', async () => { + const expired = makeAuth({ expiresAt: Date.now() - 1000 }); + await storeAuth(expired); + vi.mocked(fetch).mockResolvedValue({ + ok: false, + status: 401, + json: async () => ({ error: 'Invalid Refresh Token' }), + } as Response); + + const token = await getValidAccessToken(); + expect(token).toBeNull(); + // A failed refresh must not log the user out by itself + expect(await getStoredAuth()).toEqual(expired); + }); + + it('returns null when the refresh envelope has no session', async () => { + await storeAuth(makeAuth({ expiresAt: Date.now() - 1000 })); + vi.mocked(fetch).mockResolvedValue({ + ok: true, + json: async () => ({ data: {} }), + } as Response); + + expect(await getValidAccessToken()).toBeNull(); + }); + }); + + describe('logout race', () => { + it('discards a refresh that resolves after clearAuthSession', async () => { + await storeAuth(makeAuth({ expiresAt: Date.now() - 1000 })); + const gate = deferred(); + vi.mocked(fetch).mockReturnValue(gate.promise); + + const pending = refreshAuthToken(); + await clearAuthSession(); + gate.resolve({ + ok: true, + json: async () => refreshSuccessEnvelope, + } as Response); + + expect(await pending).toBeNull(); + // The logout wins: no tokens may be resurrected + expect(await getStoredAuth()).toBeNull(); + }); + }); + + describe('getValidAuth', () => { + it('returns the refreshed auth with the original user identity', async () => { + await storeAuth(makeAuth({ expiresAt: Date.now() - 1000 })); + vi.mocked(fetch).mockResolvedValue({ + ok: true, + json: async () => refreshSuccessEnvelope, + } as Response); + + const auth = await getValidAuth(); + expect(auth?.user).toEqual({ id: AUTH_USER_ID, email: 'user@example.com' }); + expect(auth?.accessToken).toBe('access-token-2'); + }); + }); + + describe('refreshAuthSession failure classification', () => { + it("reports 'rejected' for a definitive 4xx and keeps stored auth", async () => { + const expired = makeAuth({ expiresAt: Date.now() - 1000 }); + await storeAuth(expired); + vi.mocked(fetch).mockResolvedValue({ + ok: false, + status: 401, + json: async () => ({ error: 'Invalid Refresh Token' }), + } as Response); + + expect(await refreshAuthSession()).toEqual({ auth: null, failure: 'rejected' }); + expect(await getStoredAuth()).toEqual(expired); + }); + + it("reports 'transient' for a 5xx so callers do not destroy the session", async () => { + const expired = makeAuth({ expiresAt: Date.now() - 1000 }); + await storeAuth(expired); + vi.mocked(fetch).mockResolvedValue({ + ok: false, + status: 503, + json: async () => ({ error: 'Service unavailable' }), + } as Response); + + expect(await refreshAuthSession()).toEqual({ auth: null, failure: 'transient' }); + expect(await getStoredAuth()).toEqual(expired); + }); + + it("reports 'transient' when the network fails", async () => { + const expired = makeAuth({ expiresAt: Date.now() - 1000 }); + await storeAuth(expired); + vi.mocked(fetch).mockRejectedValue(new Error('offline')); + + expect(await refreshAuthSession()).toEqual({ auth: null, failure: 'transient' }); + expect(await getStoredAuth()).toEqual(expired); + }); + + it("reports 'signed-out' without touching the network when nothing is stored", async () => { + expect(await refreshAuthSession()).toEqual({ auth: null, failure: 'signed-out' }); + expect(fetch).not.toHaveBeenCalled(); + }); + + it('does not claim success when the secure-store write fails', async () => { + await storeAuth(makeAuth({ expiresAt: Date.now() - 1000 })); + vi.mocked(fetch).mockResolvedValue({ + ok: true, + json: async () => refreshSuccessEnvelope, + } as Response); + vi.mocked(SecureStore.setItemAsync).mockRejectedValueOnce(new Error('keystore unavailable')); + + expect(await refreshAuthSession()).toEqual({ auth: null, failure: 'transient' }); + // The mutation queue must survive the rejected write + await expect(clearAuthSession()).resolves.toBeUndefined(); + expect(await getStoredAuth()).toBeNull(); + }); + }); + + describe('envelope validation', () => { + const invalidSessions: [string, unknown][] = [ + ['an empty accessToken', { accessToken: '', refreshToken: 'r', expiresAt: 100 }], + ['an empty refreshToken', { accessToken: 'a', refreshToken: '', expiresAt: 100 }], + ['a missing expiresAt', { accessToken: 'a', refreshToken: 'r' }], + ['a NaN expiresAt', { accessToken: 'a', refreshToken: 'r', expiresAt: NaN }], + ['an Infinity expiresAt', { accessToken: 'a', refreshToken: 'r', expiresAt: Infinity }], + ['a non-string token', { accessToken: 42, refreshToken: 'r', expiresAt: 100 }], + ]; + + it.each(invalidSessions)('never stores a refresh payload with %s', async (_label, session) => { + const expired = makeAuth({ expiresAt: Date.now() - 1000 }); + await storeAuth(expired); + vi.mocked(fetch).mockResolvedValue({ + ok: true, + json: async () => ({ data: { session } }), + } as Response); + + expect(await refreshAuthSession()).toEqual({ auth: null, failure: 'transient' }); + expect(await getStoredAuth()).toEqual(expired); + }); + }); + + describe('session supersession', () => { + it('starts a fresh refresh for the new session instead of reusing a superseded one', async () => { + await storeAuth(makeAuth({ expiresAt: Date.now() - 1000, refreshToken: 'old-rt' })); + const gate = deferred(); + vi.mocked(fetch).mockReturnValueOnce(gate.promise); + const first = refreshAuthToken(); + await vi.waitFor(() => expect(fetch).toHaveBeenCalledTimes(1)); + + await clearAuthSession(); + await storeAuth( + makeAuth({ + expiresAt: Date.now() - 1000, + refreshToken: 'new-rt', + user: { id: 'user-b', email: 'b@example.com' }, + }) + ); + vi.mocked(fetch).mockResolvedValueOnce({ + ok: true, + json: async () => refreshSuccessEnvelope, + } as Response); + + const second = await refreshAuthToken(); + expect(second?.user.id).toBe('user-b'); + expect(fetch).toHaveBeenCalledTimes(2); + expect(vi.mocked(fetch).mock.calls[1]?.[1]?.body).toBe( + JSON.stringify({ refreshToken: 'new-rt' }) + ); + + gate.resolve({ ok: true, json: async () => refreshSuccessEnvelope } as Response); + expect(await first).toBeNull(); + const stored = await getStoredAuth(); + expect(stored?.user.id).toBe('user-b'); + expect(stored?.accessToken).toBe('access-token-2'); + }); + + it('discards a login commit when a logout started after the login', async () => { + const epoch = beginAuthMutation(); + await clearAuthSession(); + + const committed = await commitLoginSession(makeAuth(), epoch); + expect(committed).toBe(false); + expect(await getStoredAuth()).toBeNull(); + }); + }); +}); diff --git a/apps/mobile/src/lib/auth-session.ts b/apps/mobile/src/lib/auth-session.ts new file mode 100644 index 00000000..cb190281 --- /dev/null +++ b/apps/mobile/src/lib/auth-session.ts @@ -0,0 +1,273 @@ +/** + * Access-token lifecycle for the mobile app. + * + * Mobile port of the desktop refresh flow in + * apps/desktop/src/main/auth/secure-storage.ts: a single-flight + * POST /api/auth/refresh with the stored refresh token, converting the + * Supabase `expires_at` seconds to the millisecond timestamps kept in + * secure storage. + * + * Unlike desktop's synchronous safeStorage, expo-secure-store writes are + * async, so every token mutation is serialized through one queue and only + * committed while the session epoch it started under is still current. The + * epoch advances on logout and on every completed login, which makes the + * outcome of overlapping refreshes, logins, and logouts deterministic: a + * slow refresh can neither overwrite a newer account's tokens nor resurrect + * a session the user just signed out of. + */ +import { API_BASE_URL } from '../config'; +import { + storeAuth, + getStoredAuth, + clearStoredAuth, + isAuthExpired, + type StoredAuth, +} from './secure-storage'; + +// ── Session epoch + serialized storage mutations ────────────────── + +let sessionEpoch = 0; +// Fail closed for this process if native deletion fails. This cannot guarantee +// deletion across an app restart when the OS secure store itself is unavailable. +let storedSessionBlocked = false; + +// Every secure-store mutation runs through this queue so an in-flight native +// write can never land after (and silently undo) a logout's delete. +let mutationQueue: Promise = Promise.resolve(); + +function enqueueMutation(task: () => Promise): Promise { + const run = mutationQueue.then(task, task); + mutationQueue = run.then( + () => undefined, + () => undefined + ); + return run; +} + +/** Invalidate older login attempts before starting a new one. */ +export function beginAuthMutation(): number { + return ++sessionEpoch; +} + +/** + * Clears stored tokens. The epoch bump is synchronous, so refreshes and + * logins already in flight are invalidated before this resolves — even when + * one of them is mid-way through a native secure-store write. + */ +export async function clearAuthSession(): Promise { + await endAuthSession(); +} + +/** Capture only the outgoing session for revocation, then delete it atomically. */ +export async function endAuthSession(): Promise { + const epoch = ++sessionEpoch; + storedSessionBlocked = true; + return enqueueMutation(async () => { + let previous: StoredAuth | null = null; + try { + previous = await getStoredAuth(); + } catch { + // A failed read must not prevent local deletion. + } + await clearStoredAuth(); + if (sessionEpoch === epoch) storedSessionBlocked = false; + return previous; + }); +} + +/** + * Persist a fresh login's session. Returns false (storing nothing that + * survives) when a logout or a newer login won the race. On success the + * epoch advances, so refreshes and older logins still in flight for the + * previous session are discarded when they eventually resolve. + */ +export async function commitLoginSession(auth: StoredAuth, startedEpoch: number): Promise { + return enqueueMutation(async () => { + if (sessionEpoch !== startedEpoch) return false; + await storeAuth(auth); + if (sessionEpoch !== startedEpoch) { + // A newer login may fail without writing or deleting anything. Remove + // this superseded login inside the queue before any newer writer runs. + const cleanupEpoch = sessionEpoch; + storedSessionBlocked = true; + await clearStoredAuth(); + if (sessionEpoch === cleanupEpoch) storedSessionBlocked = false; + return false; + } + sessionEpoch += 1; + storedSessionBlocked = false; + return true; + }); +} + +// ── Server envelope validation ──────────────────────────────────── + +export interface SessionEnvelope { + accessToken: string; + refreshToken: string; + /** Milliseconds since epoch (converted from the server's seconds). */ + expiresAt: number; +} + +/** + * Validate a login/refresh `session` payload. Returns null unless both + * tokens are non-empty strings and `expiresAt` is a positive finite number + * of seconds. Supabase types `expires_at` as optional, so a proxy or server + * change can omit it — storing the resulting NaN would create a token that + * never counts as expired and therefore never refreshes. + */ +export function parseSessionEnvelope(session: unknown): SessionEnvelope | null { + if (typeof session !== 'object' || session === null) return null; + const { accessToken, refreshToken, expiresAt } = session as Record; + if (typeof accessToken !== 'string' || accessToken.trim().length === 0) return null; + if (typeof refreshToken !== 'string' || refreshToken.trim().length === 0) return null; + if (typeof expiresAt !== 'number' || !Number.isFinite(expiresAt) || expiresAt <= 0) return null; + const expiresAtMs = expiresAt * 1000; + if (!Number.isFinite(expiresAtMs) || expiresAtMs > 8640000000000000) return null; + return { accessToken, refreshToken, expiresAt: expiresAtMs }; +} + +// ── Single-flight refresh ───────────────────────────────────────── + +export type RefreshFailure = + /** No stored refresh token — the user is signed out. */ + | 'signed-out' + /** The server definitively refused the refresh token (4xx). */ + | 'rejected' + /** Network/server error or malformed payload — the token may still work. */ + | 'transient' + /** A login or logout completed while this refresh was in flight. */ + | 'superseded'; + +export interface RefreshResult { + auth: StoredAuth | null; + failure?: RefreshFailure; +} + +let refreshInFlight: { epoch: number; promise: Promise } | null = null; + +/** + * Refresh the access token using the stored refresh token. Concurrent calls + * within one session epoch share a single request; after a login or logout + * the stale in-flight promise is never handed out, so a caller for the new + * account starts a fresh refresh instead of receiving the old account's + * (discarded) outcome. + */ +export async function refreshAuthSession(): Promise { + if (refreshInFlight?.epoch === sessionEpoch) { + return refreshInFlight.promise; + } + + const entry = { epoch: sessionEpoch, promise: runRefresh(sessionEpoch) }; + refreshInFlight = entry; + try { + return await entry.promise; + } finally { + // Only retire our own entry: a newer epoch's refresh may already be here. + if (refreshInFlight === entry) { + refreshInFlight = null; + } + } +} + +async function runRefresh(startedEpoch: number): Promise { + const failureResult = (failure: RefreshFailure): RefreshResult => ({ + auth: null, + failure: sessionEpoch === startedEpoch ? failure : 'superseded', + }); + const stored = await readAuthSession(); + if (sessionEpoch !== startedEpoch) return failureResult('superseded'); + if (!stored?.refreshToken) { + return { auth: null, failure: 'signed-out' }; + } + + let session: SessionEnvelope | null = null; + try { + const response = await fetch(`${API_BASE_URL}/api/auth/refresh`, { + method: 'POST', + headers: { 'Content-Type': 'application/json' }, + body: JSON.stringify({ refreshToken: stored.refreshToken }), + }); + + if (!response.ok) { + console.error('[Auth] Token refresh failed:', response.status); + // The refresh route returns 400/401 for a missing or rejected token. + // Rate limits/timeouts (429/408) must not destroy a recoverable session. + return failureResult( + response.status === 400 || response.status === 401 ? 'rejected' : 'transient' + ); + } + + const result = (await response.json()) as { data?: { session?: unknown } }; + session = parseSessionEnvelope(result.data?.session); + } catch (error) { + console.error('[Auth] Token refresh error:', error); + return failureResult('transient'); + } + + if (!session) { + console.error('[Auth] Token refresh returned an invalid session payload'); + return failureResult('transient'); + } + + const updated: StoredAuth = { ...session, user: stored.user }; + + try { + const committed = await enqueueMutation(async () => { + // A logout/login during the request wins: discard the refreshed tokens. + if (sessionEpoch !== startedEpoch) return false; + await storeAuth(updated); + // Re-check: a logout may have arrived while the native write ran; its + // queued delete (behind us in the queue) removes what we just wrote. + return sessionEpoch === startedEpoch; + }); + if (!committed) { + return { auth: null, failure: 'superseded' }; + } + } catch (error) { + console.error('[Auth] Failed to persist refreshed session:', error); + return failureResult('transient'); + } + + return { auth: updated }; +} + +/** + * Refresh the access token using the stored refresh token. + * Returns the updated StoredAuth on success, or null on failure. + * Concurrent calls are deduplicated. + */ +export async function refreshAuthToken(): Promise { + return (await refreshAuthSession()).auth; +} + +/** Read the current session without reviving tokens from a failed logout. */ +export async function readAuthSession(): Promise { + const epoch = sessionEpoch; + await mutationQueue; + if (sessionEpoch !== epoch || storedSessionBlocked) return null; + const stored = await getStoredAuth(); + return sessionEpoch === epoch ? stored : null; +} + +/** + * Get a valid auth, refreshing the token if expired. + * Returns null if not authenticated or refresh fails. + */ +export async function getValidAuth(): Promise { + const epoch = sessionEpoch; + const stored = await readAuthSession(); + if (!stored || sessionEpoch !== epoch) return null; + + if (!isAuthExpired(stored)) return stored; + + // Token is expired — try to refresh + return refreshAuthToken(); +} + +/** Access token for API/SSE calls, refreshed when expired. */ +export async function getValidAccessToken(): Promise { + const epoch = sessionEpoch; + const auth = await getValidAuth(); + return sessionEpoch === epoch ? (auth?.accessToken ?? null) : null; +} diff --git a/apps/mobile/src/lib/event-source.test.ts b/apps/mobile/src/lib/event-source.test.ts index 1dec1865..e465327a 100644 --- a/apps/mobile/src/lib/event-source.test.ts +++ b/apps/mobile/src/lib/event-source.test.ts @@ -23,14 +23,30 @@ describe('event-source', () => { vi.clearAllMocks(); }); - it('should create an SSEConnection wrapping RNEventSource', () => { + it('should create an SSEConnection wrapping RNEventSource without auto-reconnect', () => { const connection = createEventSource('https://example.com/sse'); - expect(RNEventSource).toHaveBeenCalledWith('https://example.com/sse'); + // pollingInterval 0: reconnects are owned by the caller so each attempt + // can carry a freshly refreshed Authorization header. + expect(RNEventSource).toHaveBeenCalledWith('https://example.com/sse', { + headers: undefined, + pollingInterval: 0, + }); expect(connection).toHaveProperty('addEventListener'); expect(connection).toHaveProperty('close'); }); + it('passes request headers through to RNEventSource', () => { + createEventSource('https://example.com/sse', { + headers: { Authorization: 'Bearer token-1' }, + }); + + expect(RNEventSource).toHaveBeenCalledWith('https://example.com/sse', { + headers: { Authorization: 'Bearer token-1' }, + pollingInterval: 0, + }); + }); + it('should forward addEventListener calls to underlying EventSource', () => { 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 7b87d68f..d83daeaa 100644 --- a/apps/mobile/src/lib/event-source.ts +++ b/apps/mobile/src/lib/event-source.ts @@ -10,13 +10,25 @@ type PairUXEvents = 'connected' | 'heartbeat' | 'signal' | 'presence-join' | 'pr export type SSEEventHandler = (event: { data: string }) => void; +export interface SSEOptions { + /** Extra request headers, e.g. the Authorization bearer token. */ + headers?: Record; +} + export interface SSEConnection { addEventListener: (event: string, handler: SSEEventHandler) => void; close: () => void; } -export function createEventSource(url: string): SSEConnection { - const es = new RNEventSource(url); +export function createEventSource(url: string, options: SSEOptions = {}): SSEConnection { + // pollingInterval: 0 disables the library's built-in auto-reconnect. Its + // retries would replay the captured Authorization header long after the + // token expired, silently downgrading the stream to a guest identity — + // callers own reconnection so every attempt carries a fresh token. + const es = new RNEventSource(url, { + headers: options.headers, + pollingInterval: 0, + }); return { addEventListener(event: string, handler: SSEEventHandler) { diff --git a/apps/mobile/src/lib/secure-storage.test.ts b/apps/mobile/src/lib/secure-storage.test.ts index 8a52f4e1..38c79c14 100644 --- a/apps/mobile/src/lib/secure-storage.test.ts +++ b/apps/mobile/src/lib/secure-storage.test.ts @@ -61,6 +61,9 @@ describe('secure-storage', () => { }); describe('isAuthExpired', () => { + it.each([NaN, Infinity, -Infinity])('treats non-finite expiry %s as expired', (expiresAt) => { + expect(isAuthExpired({ ...mockAuth, expiresAt })).toBe(true); + }); it('should return false for non-expired token', () => { const auth: StoredAuth = { ...mockAuth, diff --git a/apps/mobile/src/lib/secure-storage.ts b/apps/mobile/src/lib/secure-storage.ts index be6cc08f..d949a104 100644 --- a/apps/mobile/src/lib/secure-storage.ts +++ b/apps/mobile/src/lib/secure-storage.ts @@ -35,7 +35,7 @@ export async function clearStoredAuth(): Promise { export function isAuthExpired(auth: StoredAuth): boolean { // Consider expired if within 5 minutes of expiry - return Date.now() >= auth.expiresAt - 5 * 60 * 1000; + return !Number.isFinite(auth.expiresAt) || Date.now() >= auth.expiresAt - 5 * 60 * 1000; } // --- Remembered credentials (separate from session tokens) --- diff --git a/apps/mobile/src/test/fixtures/server-contracts.ts b/apps/mobile/src/test/fixtures/server-contracts.ts new file mode 100644 index 00000000..0c15cfaa --- /dev/null +++ b/apps/mobile/src/test/fixtures/server-contracts.ts @@ -0,0 +1,128 @@ +/** + * Server-shaped response fixtures for mobile contract tests. + * + * Every shape here mirrors an actual route handler in apps/web/src/app/api: + * - auth login/signup/refresh wrap payloads as `{ data: ... }` + * (apps/web/src/lib/api.ts successResponse) and return Supabase's + * `expires_at` in SECONDS since epoch; + * - GET /api/sessions/join/[joinCode] returns the session payload directly + * (or a scheduled-meeting payload flagged `scheduled: true`); + * - POST /api/sessions/join/[joinCode] returns the session_participants row + * from the join_session RPC; + * - the signaling SSE `connected` event carries the server-assigned + * subscriberId (the auth user id for authenticated clients). + * + * The IDs are deliberately distinct so tests fail whenever code confuses the + * authenticated user id, the participant row id, and the SSE subscriber id. + */ + +export const AUTH_USER_ID = 'auth-user-1111'; +export const PARTICIPANT_ROW_ID = 'participant-row-2222'; +export const HOST_USER_ID = 'host-user-3333'; +export const OTHER_VIEWER_ID = 'other-viewer-4444'; +export const SESSION_ID = 'session-5555'; + +// Supabase expiry: seconds since epoch. Clients must store milliseconds. +export const SUPABASE_EXPIRES_AT_SECONDS = 1_795_000_000; + +export const authUser = { id: AUTH_USER_ID, email: 'user@example.com' }; + +export const loginSuccessEnvelope = { + data: { + user: authUser, + session: { + accessToken: 'access-token-1', + refreshToken: 'refresh-token-1', + expiresAt: SUPABASE_EXPIRES_AT_SECONDS, + }, + }, +}; + +export const signupNeedsConfirmationEnvelope = { + data: { + user: authUser, + message: 'Check your email to confirm your account', + needsConfirmation: true, + }, +}; + +export const signupConfirmedEnvelope = { + data: { + user: authUser, + message: 'Account created successfully', + needsConfirmation: false, + }, +}; + +export const refreshSuccessEnvelope = { + data: { + session: { + accessToken: 'access-token-2', + refreshToken: 'refresh-token-2', + expiresAt: SUPABASE_EXPIRES_AT_SECONDS, + }, + }, +}; + +export const joinLookupLiveEnvelope = { + data: { + id: SESSION_ID, + join_code: 'ABC123', + status: 'active', + settings: { allowGuestControl: false, maxParticipants: 20 }, + created_at: '2026-09-01T00:00:00.000Z', + participant_count: 2, + }, +}; + +export const joinLookupScheduledEnvelope = { + data: { + scheduled: true as const, + id: 'sched-6666', + join_code: 'ABC123', + title: 'Team Standup', + description: null, + scheduled_at: '2026-09-08T10:00:00.000Z', + duration_minutes: 30, + invitees: [{ name: 'Alice', rsvp_status: 'accepted' }], + }, +}; + +export const joinParticipantEnvelope = { + data: { + id: PARTICIPANT_ROW_ID, + session_id: SESSION_ID, + user_id: AUTH_USER_ID, + display_name: 'Phone Viewer', + role: 'viewer', + control_state: 'view-only', + joined_at: '2026-09-07T00:00:00.000Z', + left_at: null, + }, +}; + +export const sseIceServers = [ + { urls: 'turn:turn.pairux.com:3478', username: 'turn-user', credential: 'turn-pass' }, +]; + +/** SSE `connected` event payload for an authenticated viewer. */ +export function viewerConnectedEventData(overrides: Record = {}): string { + return JSON.stringify({ + sessionId: SESSION_ID, + subscriberId: AUTH_USER_ID, + isHost: false, + iceServers: sseIceServers, + ...overrides, + }); +} + +/** SSE `connected` event payload for the authenticated host. */ +export function hostConnectedEventData(overrides: Record = {}): string { + return JSON.stringify({ + sessionId: SESSION_ID, + subscriberId: HOST_USER_ID, + isHost: true, + iceServers: sseIceServers, + ...overrides, + }); +}