import { useState, useCallback, useRef, useEffect } from 'react'; import { Room, RoomEvent, Track, Participant, RemoteParticipant, RemoteTrackPublication, RemoteAudioTrack, ConnectionState, ConnectionQuality, LocalAudioTrack, LocalTrackPublication, DisconnectReason, } from 'livekit-client'; import { getApiForOrigin, getChannelOrigin, getMyUserIdForOrigin, useSpaceStore } from '../stores/spaceStore'; import { wsSend } from './useWebSocket'; import { useVoiceStore } from '../stores/voiceStore'; import { useAuthStore } from '../stores/authStore'; import { useUIStore } from '../stores/uiStore'; import type { User } from '@backspace/shared'; import { broadcastVoiceStatus } from '../utils/voice'; import { consumeIntentionalCameraOff, markIntentionalCameraOff } from '../utils/voiceActions'; import { AudioManager } from '../audio/AudioManager'; import { SpeakingDetector } from '../audio/SpeakingDetector'; import { CAMERA_OVERDRIVE, buildScreenShareOptions, applyOverdrive, startScreenShare, stopScreenShare, handleScreenShareUnpublished, resolveNativeOverdrive, } from '../utils/screenShare'; import { parseStreamWatch } from '../utils/streamWatchProtocol'; import { getMediaStreamTrack } from '../utils/livekitInternals'; import { deactivate as deactivateHwOverdrive } from '../utils/hwOverdrive'; let _activeRoom: Room | null = null; let _publishedScreenShareCodec: 'vp9' | 'h264' | null = null; export function getActiveRoom(): Room | null { return _activeRoom; } export interface ParticipantInfo { identity: string; userId: string; username: string; homeUserId: string | null; isMuted: boolean; isDeafened: boolean; isCameraOn: boolean; isScreenSharing: boolean; isLocal: boolean; audioTrack: MediaStreamTrack | null; videoTrack: MediaStreamTrack | null; screenTrack: MediaStreamTrack | null; screenAudioTrack: MediaStreamTrack | null; lkVideoTrack: Track | null; // LiveKit Track for attach/detach (adaptive stream) lkScreenTrack: Track | null; // LiveKit Track for attach/detach (adaptive stream) cachedUser: User | null; // Hydrated User from member lookup, carried forward across space switches } export interface UserTile { kind: 'user'; key: string; // participant.identity participant: ParticipantInfo; videoTrack: MediaStreamTrack | null; // camera only audioTrack: MediaStreamTrack | null; // mic lkVideoTrack: Track | null; // LiveKit Track for attach/detach } export interface StreamTile { kind: 'stream'; key: string; // `${identity}:stream` participant: ParticipantInfo; screenTrack: MediaStreamTrack | null; screenAudioTrack: MediaStreamTrack | null; lkScreenTrack: Track | null; // LiveKit Track for attach/detach } export type GridTile = UserTile | StreamTile; export function deriveGridTiles(participants: ParticipantInfo[]): GridTile[] { const tiles: GridTile[] = []; for (const p of participants) { const hasLiveVideo = p.isCameraOn && p.videoTrack?.readyState === 'live'; tiles.push({ kind: 'user', key: p.identity, participant: p, videoTrack: hasLiveVideo ? p.videoTrack : null, audioTrack: p.audioTrack, lkVideoTrack: hasLiveVideo ? p.lkVideoTrack : null, }); if (p.isScreenSharing) { tiles.push({ kind: 'stream', key: `${p.identity}:stream`, participant: p, screenTrack: p.screenTrack, screenAudioTrack: p.screenAudioTrack, lkScreenTrack: p.lkScreenTrack, }); } } return tiles; } export function setStreamSubscription(room: Room | null, targetIdentity: string, subscribed: boolean) { if (!room) return; const rp = room.remoteParticipants.get(targetIdentity); if (!rp) return; rp.trackPublications.forEach((pub) => { if (pub.source === Track.Source.ScreenShare || pub.source === Track.Source.ScreenShareAudio) { (pub as RemoteTrackPublication).setSubscribed(subscribed); } }); } export function setCameraSubscription(room: Room | null, targetIdentity: string, subscribed: boolean) { if (!room) return; const rp = room.remoteParticipants.get(targetIdentity); if (!rp) return; rp.trackPublications.forEach((pub) => { if (pub.source === Track.Source.Camera) { (pub as RemoteTrackPublication).setSubscribed(subscribed); } }); } export function parseIdentity(identity: string): { userId: string; username: string } { const parts = identity.split(':'); return { userId: parts[0] ?? identity, username: parts[1] ?? identity }; } let _connectGeneration = 0; /** Null-safe Room.disconnect() wrapper — lets the SDK tear down its own internals cleanly. */ function destroyRoom(room: Room | null): Promise | void { if (!room) return; return room.disconnect(); } /** * Ensures a fresh microphone track from the AudioManager pipeline is published * to the supplied room. If the existing publication is already current (live * MediaStreamTrack matching the latest AudioManager streamGeneration), this is * a no-op apart from un-muting. Otherwise the stale track is unpublished and a * cloned destination-node track is published in its place. * * Extracted from the syncMic effect so the input-track-loss recovery path can * call it directly — without relying on syncMic's React dep array catching a * change that never re-renders the hook. */ async function republishMicrophone(r: Room, lastMicGenRef: { current: number }): Promise { const audioManager = AudioManager.getInstance(); const currentGen = audioManager.getStreamGeneration(); const micPub = r.localParticipant.getTrackPublications() .find(p => p.source === Track.Source.Microphone); if (micPub?.track) { // Track already published — check if it's still current if (micPub.track.mediaStreamTrack?.readyState === 'live' && lastMicGenRef.current === currentGen) { // Current and live — just unmute if needed if (micPub.isMuted) { await r.localParticipant.setMicrophoneEnabled(true); } return; } // Track is stale (device or constraint change) — replace it await r.localParticipant.unpublishTrack(micPub.track as LocalAudioTrack); } // Publish fresh track from AudioManager pipeline const audioTrack = audioManager.getFreshTrack(); if (!audioTrack) return; console.log('[LiveKit] Publishing fresh microphone track (gen:', currentGen, ')'); await r.localParticipant.publishTrack(audioTrack, { name: 'microphone', source: Track.Source.Microphone, }); lastMicGenRef.current = currentGen; } export function useLiveKit() { const [room, setRoom] = useState(null); const [isConnected, setIsConnected] = useState(false); const [isConnecting, setIsConnecting] = useState(false); const [connectionState, setConnectionState] = useState(ConnectionState.Disconnected); const [connectedChannelId, setConnectedChannelId] = useState(null); const [connectionError, setConnectionError] = useState(null); const roomRef = useRef(null); const connectedChannelRef = useRef(null); const switchCameraGenRef = useRef(0); const isMuted = useVoiceStore((s) => s.isMuted); const isDeafened = useVoiceStore((s) => s.isDeafened); const isCameraOn = useVoiceStore((s) => s.isCameraOn); const isScreenSharing = useVoiceStore((s) => s.isScreenSharing); const screenShareConfig = useVoiceStore((s) => s.screenShareConfig); const hwOverdrive = useVoiceStore((s) => s.hwOverdrive); const voiceUserStates = useVoiceStore((s) => s.voiceUserStates); const spaceMutedUserIds = useVoiceStore((s) => s.spaceMutedUserIds); const spaceDeafenedUserIds = useVoiceStore((s) => s.spaceDeafenedUserIds); const permissionMutedUserIds = useVoiceStore((s) => s.permissionMutedUserIds); const inputVolume = useVoiceStore((s) => s.inputVolume); const inputDeviceId = useVoiceStore((s) => s.inputDeviceId); const cameraDeviceId = useVoiceStore((s) => s.cameraDeviceId); const echoCancellation = useVoiceStore((s) => s.echoCancellation); const noiseSuppression = useVoiceStore((s) => s.noiseSuppression); const autoGainControl = useVoiceStore((s) => s.autoGainControl); const rnnoiseEnabled = useVoiceStore((s) => s.rnnoiseEnabled); const lastMicGenRef = useRef(0); const updateParticipants = useCallback(() => { const r = roomRef.current; if (!r) return; // Carry-forward: snapshot previous participants for cachedUser preservation const prevParticipants = useVoiceStore.getState().participants; const prevCacheMap = new Map(); for (const prev of prevParticipants) { prevCacheMap.set(prev.identity, prev.cachedUser); } const allParticipants: ParticipantInfo[] = []; const processParticipant = (p: Participant, isLocal: boolean) => { if (!p.identity) return; const { userId: rawId, username } = parseIdentity(p.identity); // Resolve identity: for federated calls rawId may be homeUserId from another instance. // Check DM members for a user whose homeUserId matches. let userId = rawId; const activeDmCall = useVoiceStore.getState().activeDmCall; if (activeDmCall) { const dmChannels = useSpaceStore.getState().dmChannels; const dmChannel = dmChannels.find(d => d.id === activeDmCall.dmChannelId); if (dmChannel) { const match = dmChannel.members.find(m => m.homeUserId === rawId || m.id === rawId); if (match) userId = match.id; } } const memberMatch = useSpaceStore.getState().members.find(m => m.userId === userId); let cachedUser: User | null; let homeUserId: string | null; if (memberMatch) { // Fresh data available — use and update cache cachedUser = memberMatch.user as User; homeUserId = memberMatch.user.homeUserId ?? null; } else if (isLocal) { // Local user safety net — authStore is always available cachedUser = useAuthStore.getState().user; homeUserId = cachedUser?.homeUserId ?? null; } else { // Space switched — carry forward from previous cycle cachedUser = prevCacheMap.get(p.identity) ?? null; homeUserId = cachedUser?.homeUserId ?? null; } let audioTrack: MediaStreamTrack | null = null; let videoTrack: MediaStreamTrack | null = null; let screenTrack: MediaStreamTrack | null = null; let screenAudioTrack: MediaStreamTrack | null = null; let lkVideoTrack: Track | null = null; let lkScreenTrack: Track | null = null; let hasScreenSharePublication = false; let hasCameraPublication = false; p.trackPublications.forEach((pub) => { // Detect screen share publication even if unsubscribed if (pub.source === Track.Source.ScreenShare) hasScreenSharePublication = true; // Detect camera publication even if unsubscribed (for unwatched cameras) if (pub.source === Track.Source.Camera) hasCameraPublication = true; const track = pub.track; if (!track) return; // Strict check: Track must be subscribed AND not muted to be considered "active" if (pub.isMuted) return; if (!isLocal && !pub.isSubscribed) return; const mt = track.mediaStreamTrack; if (!mt || mt.readyState !== 'live') return; if (pub.source === Track.Source.Microphone) audioTrack = mt; else if (pub.source === Track.Source.Camera && p.isCameraEnabled) { videoTrack = mt; lkVideoTrack = track; } else if (pub.source === Track.Source.ScreenShare) { screenTrack = mt; lkScreenTrack = track; } else if (pub.source === Track.Source.ScreenShareAudio) screenAudioTrack = mt; }); const userState = useVoiceStore.getState().voiceUserStates.get(userId); let isPartDeafened = false; let isPartMuted = !p.isMicrophoneEnabled; if (isLocal) { // Compute effective state: user intent || server enforcement const vs = useVoiceStore.getState(); const cvId = vs.currentVoiceChannelId; const localOrigin = cvId ? getChannelOrigin(cvId) : ''; const localMyId = cvId ? getMyUserIdForOrigin(localOrigin) : undefined; const localSpaceId = cvId ? useSpaceStore.getState().channelToSpaceMap.get(cvId) : null; const localKey = (localSpaceId && localMyId) ? `${localSpaceId}:${localMyId}` : ''; isPartMuted = vs.isMuted || vs.spaceMutedUserIds.has(localKey) || vs.permissionMutedUserIds.has(localKey); isPartDeafened = vs.isDeafened || vs.spaceDeafenedUserIds.has(localKey); } else { isPartDeafened = userState?.isDeafened ?? useVoiceStore.getState().deafenedUserIds.has(userId); if (userState) isPartMuted = userState.isMuted; } allParticipants.push({ identity: p.identity, userId, username, homeUserId, cachedUser, isMuted: isPartMuted, isDeafened: isPartDeafened, isCameraOn: hasCameraPublication && p.isCameraEnabled, // True even when unsubscribed isScreenSharing: hasScreenSharePublication, // True even when unsubscribed isLocal, audioTrack, videoTrack, screenTrack, screenAudioTrack, lkVideoTrack, lkScreenTrack, }); }; processParticipant(r.localParticipant, true); r.remoteParticipants.forEach((p) => processParticipant(p, false)); useVoiceStore.getState().setParticipants(allParticipants); SpeakingDetector.getInstance().syncTracks(allParticipants); }, []); const handleDataReceived = useCallback((payload: Uint8Array, participant?: RemoteParticipant) => { // Try the stream_watch protocol first (typed parser; returns null on non-matches). if (participant) { const sw = parseStreamWatch(payload); if (sw) { // sw.target is the streamer's bare userId. participant.identity is the // viewer's full LiveKit identity ("userId:username"); we key by identity // so ParticipantDisconnected can evict cleanly. useVoiceStore.getState().recordStreamWatch(sw.target, participant.identity, sw.watching); return; } } try { const text = new TextDecoder().decode(payload); const msg = JSON.parse(text); if (msg.type === 'deafen' && participant) { const { userId } = parseIdentity(participant.identity); useVoiceStore.getState().setUserDeafened(userId, msg.deafened === true); updateParticipants(); } } catch { } }, [updateParticipants]); // Handle Input Device & Mute Logic via AudioManager // Mute uses setMicrophoneEnabled(false) to keep the track published (silence frames) // instead of unpublishTrack() which tears down the WebRTC transport. // This preserves the Web Audio pipeline for future AudioWorklet nodes (e.g. RNNoise). useEffect(() => { const r = roomRef.current; if (!r || !isConnected) return; // Compute effective mute/deafen: user intent || server enforcement const vs = useVoiceStore.getState(); const cvId = vs.currentVoiceChannelId; const effOrigin = cvId ? getChannelOrigin(cvId) : ''; const effMyId = cvId ? getMyUserIdForOrigin(effOrigin) : undefined; const effSpaceId = cvId ? useSpaceStore.getState().channelToSpaceMap.get(cvId) : null; const effKey = (effSpaceId && effMyId) ? `${effSpaceId}:${effMyId}` : ''; const effectiveMuted = isMuted || spaceMutedUserIds.has(effKey) || permissionMutedUserIds.has(effKey); const effectiveDeafened = isDeafened || spaceDeafenedUserIds.has(effKey); const syncMic = async () => { try { const audioManager = AudioManager.getInstance(); // Sync voice processing settings to AudioManager audioManager.setVoiceProcessing({ echoCancellation, noiseSuppression, autoGainControl }); await audioManager.setRnnoiseEnabled(rnnoiseEnabled); const micPub = r.localParticipant.getTrackPublications() .find(p => p.source === Track.Source.Microphone); // If effectively muted or deafened, mute the track in-place (keep it published) if (effectiveMuted || effectiveDeafened) { if (micPub?.track && !micPub.isMuted) { await r.localParticipant.setMicrophoneEnabled(false); } return; } // Not muted — ensure the AudioManager pipeline is on the right device // and at the right volume, then republish if the published track is // stale or missing. await audioManager.setInputDevice(inputDeviceId); audioManager.setInputVolume(inputVolume); await republishMicrophone(r, lastMicGenRef); } catch (err) { console.error('[LiveKit] Failed to sync mic state:', err); } }; syncMic(); // Re-sync when AudioManager resumes const unsubscribeResume = AudioManager.getInstance().onResumed(() => { syncMic(); }); return () => { unsubscribeResume(); }; }, [isMuted, isDeafened, spaceMutedUserIds, spaceDeafenedUserIds, permissionMutedUserIds, inputDeviceId, inputVolume, isConnected, echoCancellation, noiseSuppression, autoGainControl, rnnoiseEnabled]); // Subscribe to upstream-input-track-end events from AudioManager whenever a // room is connected. The published mic track is a clone of a WebAudio // destination node and never ends on hardware loss; only the upstream // getUserMedia track does. AudioManager owns that signal — we react to it. useEffect(() => { if (!isConnected) return; const am = AudioManager.getInstance(); const subscriberRoom = roomRef.current; const unsubscribe = am.onInputTrackEnded(async () => { // Room was replaced or torn down between event emission and handler run. if (roomRef.current !== subscriberRoom || !subscriberRoom) return; const deviceId = useVoiceStore.getState().inputDeviceId; let copy = 'Microphone unavailable'; try { const probe = await navigator.mediaDevices.getUserMedia({ audio: deviceId === 'default' ? true : { deviceId: { exact: deviceId } }, }); probe.getTracks().forEach(t => t.stop()); // Probe succeeded — device is back. Re-acquire and force a republish. if (roomRef.current !== subscriberRoom) return; try { await am.setInputDevice(deviceId); if (roomRef.current !== subscriberRoom) return; await republishMicrophone(subscriberRoom, lastMicGenRef); return; } catch { copy = 'Microphone could not be restored'; } } catch (err: any) { if (err?.name === 'NotAllowedError') { copy = 'Microphone permission was revoked'; } else if (err?.name === 'NotFoundError') { if (deviceId !== 'default') { // The configured device disappeared. Fall back to default — the // store update triggers syncMic via its dep array, which calls // republishMicrophone with the freshly acquired default stream. useVoiceStore.getState().setInputDevice('default'); copy = 'Microphone disconnected — switched to system default'; } else { copy = 'Microphone disconnected'; } } } if (roomRef.current !== subscriberRoom) return; useUIStore.getState().addToast(copy, 'warning'); }); return () => { unsubscribe(); }; }, [isConnected]); // Hot-swap the camera source when cameraDeviceId changes mid-call. // Compares against the published track's actual deviceId (getSettings().deviceId) // rather than a memoised previous store value, so the null → explicit-same-device // transition is a correct no-op. useEffect(() => { const r = roomRef.current; if (!r || !isConnected || !isCameraOn) return; const camPub = r.localParticipant.getTrackPublications() .find(p => p.source === Track.Source.Camera); if (!camPub?.track) return; const currentDeviceId = camPub.track.mediaStreamTrack?.getSettings().deviceId; const target = cameraDeviceId; if (target === null) return; // "Auto" never force-switches a live publication if (currentDeviceId === target) return; // already on target const myGen = ++switchCameraGenRef.current; // Race semantics: if the effect re-fires while this IIFE is in flight // (rapid dropdown changes), the newer firing will increment the gen. // This IIFE's catch then no-ops its store rollback — the newer attempt // owns the canonical store state. (async () => { const prev = currentDeviceId ?? null; try { await r.switchActiveDevice('videoinput', target); } catch (err) { console.error('[LiveKit] Camera hot-swap failed:', err); // A newer device-switch attempt has superseded ours. Don't write stale // rollback state; let the newer attempt's outcome stand. if (myGen !== switchCameraGenRef.current) { useUIStore.getState().addToast('Could not switch camera', 'warning'); return; } const stillLive = camPub.track?.mediaStreamTrack?.readyState === 'live'; if (stillLive && prev) { useVoiceStore.getState().setCameraDeviceId(prev); } else if (stillLive) { useVoiceStore.getState().setCameraDeviceId(null); } else { // Track is dead — disable the camera entirely. Mark the flag now to // suppress the already-queued `ended` event on the dead track. markIntentionalCameraOff(); // Re-mark immediately before the disable: LiveKit may synthesize a // second `ended` event during teardown, and the flag is consumed once. // markIntentionalCameraOff is idempotent. markIntentionalCameraOff(); await r.localParticipant.setCameraEnabled(false).catch(() => {}); useVoiceStore.setState({ isCameraOn: false }); broadcastVoiceStatus(); } useUIStore.getState().addToast('Could not switch camera', 'warning'); } })(); }, [cameraDeviceId, isCameraOn, isConnected]); const connect = useCallback(async (channelId: string, isDm?: boolean) => { const storedId = isDm ? `dm-${channelId}` : channelId; if (connectedChannelRef.current === storedId && roomRef.current?.state === ConnectionState.Connected) return; // Register voice state with the WS server after LiveKit connects (not for DM calls) const registerWithServer = () => { if (isDm) return; const origin = getChannelOrigin(channelId); wsSend({ type: 'voice_join', channelId }, origin); broadcastVoiceStatus(origin); }; const gen = ++_connectGeneration; // Ensure AudioContext is created and resumed before tracks arrive await AudioManager.getInstance().resumeContext(); // 1. Reset state immediately to reflect "Loading/Switching" in UI SpeakingDetector.getInstance().clear(); setRoom(null); useVoiceStore.getState().setParticipants([]); useVoiceStore.getState().setSpeakingParticipants(new Set()); setIsConnected(false); setIsConnecting(true); setConnectionState(ConnectionState.Connecting); setConnectionError(null); setConnectedChannelId(null); // Clear this so AppLayout knows we are transitioning useVoiceStore.getState().setConnectionError(null); useVoiceStore.getState().setIsLiveKitConnected(false); useVoiceStore.getState().setConnectionQuality('unknown'); // 2. Strictly disconnect previous room (Local Ref OR Global Ref) // This handles cases where AppLayout might have remounted, losing roomRef but leaving _activeRoom alive. const roomToDisconnect = roomRef.current || _activeRoom; if (roomToDisconnect) { try { console.log('[LiveKit] Destroying previous room:', roomToDisconnect.name); await destroyRoom(roomToDisconnect); } catch (err) { console.warn('Error disconnecting from previous room:', err); } roomRef.current = null; _activeRoom = null; } try { let token: string; let url: string; // For federated calls, use the stored token from S2S relay const { federatedCallToken, federatedCallUrl, clearFederatedCallData } = useVoiceStore.getState(); if (isDm && federatedCallToken && federatedCallUrl) { token = federatedCallToken; url = federatedCallUrl; clearFederatedCallData(); } else { const client = getApiForOrigin(getChannelOrigin(channelId)); const resp = isDm ? await client.livekit.dmToken(channelId) : await client.livekit.token(channelId); token = resp.token; url = resp.url; } if (gen !== _connectGeneration) return; const newRoom = new Room({ adaptiveStream: true, dynacast: true, publishDefaults: { videoCodec: 'h264', simulcast: true } }); roomRef.current = newRoom; const guardedUpdate = () => { if (roomRef.current === newRoom) updateParticipants(); }; newRoom.on(RoomEvent.ParticipantConnected, (participant) => { guardedUpdate(); // Notify new participant of our effective deafen state const vsConn = useVoiceStore.getState(); const cvIdConn = vsConn.currentVoiceChannelId; const connOrigin = cvIdConn ? getChannelOrigin(cvIdConn) : ''; const connMyId = cvIdConn ? getMyUserIdForOrigin(connOrigin) : undefined; const connSpaceId = cvIdConn ? useSpaceStore.getState().channelToSpaceMap.get(cvIdConn) : null; const connKey = (connSpaceId && connMyId) ? `${connSpaceId}:${connMyId}` : ''; const effDeaf = vsConn.isDeafened || vsConn.spaceDeafenedUserIds.has(connKey); if (effDeaf) { const encoder = new TextEncoder(); newRoom.localParticipant.publishData( encoder.encode(JSON.stringify({ type: 'deafen', deafened: true })), { reliable: true } ).catch(() => { }); } }); newRoom.on(RoomEvent.ParticipantDisconnected, (participant: RemoteParticipant) => { useVoiceStore.getState().evictWatcher(participant.identity); guardedUpdate(); // Clean up stale WS-based voice status for the departed participant const { userId } = parseIdentity(participant.identity); useVoiceStore.getState().clearVoiceUserStatus(userId); }); newRoom.on(RoomEvent.TrackSubscribed, (track, publication, participant) => { // LiveKit auto-attaches a hidden