diff --git a/Client/tauri-client/src/lib/livekitSession.ts b/Client/tauri-client/src/lib/livekitSession.ts index fa7a73d7..2482ea10 100644 --- a/Client/tauri-client/src/lib/livekitSession.ts +++ b/Client/tauri-client/src/lib/livekitSession.ts @@ -10,6 +10,7 @@ import { leaveVoiceChannel, setListenOnly, } from "@stores/voice.store"; +import { authStore } from "@stores/auth.store"; import { loadPref } from "@components/settings/helpers"; import { createLogger } from "@lib/logger"; import { invoke } from "@tauri-apps/api/core"; @@ -150,6 +151,9 @@ export class LiveKitSession { private _roomKeyRejector: ((err: Error) => void) | null = null; /** Guard: true while a key rotation is in progress (prevents concurrent rotations). */ private _rotatingKey = false; + /** Monotonic counter incremented on every key rotation. handleE2EEOffer captures the + * epoch before async work and discards the result if epoch changed (stale offer). */ + private _e2eeEpoch = 0; /** Announces that arrived before our ECDH keypair was ready. Drained after keypair init. */ private _pendingAnnounces: Array<{ userId: number; publicKeyBase64: string }> = []; @@ -435,12 +439,24 @@ export class LiveKitSession { return; } - // E2EE keys are exchanged client-side via ECDH. On reconnect, the - // room key is still in _roomKey from the previous session. Re-apply it. + // E2EE: Regenerate ECDH keypair for the new session (forward secrecy) + // and re-announce so other participants can re-wrap the room key for us. + // If we still have the room key from before disconnect, re-apply it now + // so audio works immediately; the key holder will send a fresh offer if + // the key was rotated during our absence. + // eslint-disable-next-line no-await-in-loop -- must set up E2EE before connect + this._ecdhKeyPair = await generateECDHKeyPair(); + this._peerPublicKeys.clear(); if (this._roomKey) { // eslint-disable-next-line no-await-in-loop -- must set key before connect await this._e2eeKeyProvider.setKey(roomKeyToBase64(this._roomKey)); } + // eslint-disable-next-line no-await-in-loop -- must export before connect + const reconnectPubKey = await exportPublicKey(this._ecdhKeyPair.publicKey); + this.ws?.send({ + type: "voice_e2ee_announce", + payload: { public_key: reconnectPubKey }, + }); // eslint-disable-next-line no-await-in-loop -- sequential reconnect: must connect before restoring state await newRoom.connect(resolvedUrl, token); @@ -821,15 +837,22 @@ export class LiveKitSession { log.info("E2EE: drained queued announce", { userId: qId }); } - // Determine if we are the key holder (first in channel = no existing - // voice_e2ee_announce messages received before this point). - // The server sends existing participants' public keys during voice_join - // sync — if we received none (including queued), we're first. - const existingPeerCount = this._peerPublicKeys.size; - this._isKeyHolder = existingPeerCount === 0; + // Determine key holder by lowest user ID in the channel (deterministic, + // avoids race when two users join simultaneously — both would otherwise + // see zero peers and elect themselves key holder). + const myId = authStore.getState().user?.id ?? 0; + const channelVoiceUsers = voiceStore.getState().voiceUsers.get(channelId); + let lowestInChannel = myId; + if (channelVoiceUsers) { + for (const uid of channelVoiceUsers.keys()) { + if (uid < lowestInChannel) lowestInChannel = uid; + } + } + this._isKeyHolder = myId !== 0 && lowestInChannel === myId; if (this._isKeyHolder) { // We're the first participant — generate the room key. + this._e2eeEpoch++; this._roomKey = generateRoomKey(); await this._e2eeKeyProvider.setKey(roomKeyToBase64(this._roomKey)); log.info("E2EE: key holder — generated room key", { channelId }); @@ -842,13 +865,17 @@ export class LiveKitSession { this._roomKeyRejector = reject; }); // Don't block forever — timeout after 10 seconds. - const timeout = new Promise((_, reject) => - setTimeout(() => reject(new Error("E2EE key exchange timeout")), 10_000), - ); + let timeoutId: ReturnType | null = null; + const timeout = new Promise((_, reject) => { + timeoutId = setTimeout(() => reject(new Error("E2EE key exchange timeout")), 10_000); + }); try { await Promise.race([roomKeyPromise, timeout]); } catch { log.warn("E2EE: key exchange timed out, proceeding without E2EE", { channelId }); + this.onErrorCallback?.("End-to-end encryption unavailable — key exchange timed out"); + } finally { + if (timeoutId !== null) clearTimeout(timeoutId); } this._roomKeyResolver = null; this._roomKeyRejector = null; @@ -1106,16 +1133,32 @@ export class LiveKitSession { return; } try { + // Deduplicate: if we already have a key for this user with the same base64, + // skip the import. If the key changed, log a warning (could be a reconnect + // or a protocol violation). + const existingKey = this._peerPublicKeys.get(userId); const peerKey = await importPublicKey(publicKeyBase64); + if (existingKey) { + const existingB64 = await exportPublicKey(existingKey); + if (existingB64 === publicKeyBase64) { + log.debug("E2EE: duplicate announce from same key, ignoring", { userId }); + return; + } + log.warn("E2EE: peer public key changed (reconnect or key rotation?)", { userId }); + } this._peerPublicKeys.set(userId, peerKey); log.info("E2EE: received peer public key", { userId }); // If we're the key holder and have a room key, wrap it for the new peer. - if (this._isKeyHolder && this._roomKey && this._ecdhKeyPair) { + // Capture keypair + roomKey before async work to avoid null dereference if + // clearE2EEState() runs concurrently. + const keypair = this._ecdhKeyPair; + const currentRoomKey = this._roomKey; + if (this._isKeyHolder && currentRoomKey && keypair) { const { encryptedKey, iv } = await wrapRoomKey( - this._ecdhKeyPair.privateKey, + keypair.privateKey, peerKey, - this._roomKey, + currentRoomKey, ); this.ws?.send({ type: "voice_e2ee_offer", @@ -1143,17 +1186,33 @@ export class LiveKitSession { log.warn("E2EE: received offer from unknown peer", { fromUserId }); return; } - if (!this._ecdhKeyPair) { + const keypair = this._ecdhKeyPair; + if (!keypair) { log.warn("E2EE: received offer but no ECDH keypair"); return; } - this._roomKey = await unwrapRoomKey( - this._ecdhKeyPair.privateKey, + // Capture epoch before async work — if a key rotation occurs during + // unwrap, the epoch will have advanced and we discard this stale result. + const epochBefore = this._e2eeEpoch; + + const unwrapped = await unwrapRoomKey( + keypair.privateKey, peerKey, encryptedKeyBase64, ivBase64, ); + + if (this._e2eeEpoch !== epochBefore) { + log.info("E2EE: discarding stale offer (epoch changed during unwrap)", { + fromUserId, + epochBefore, + epochNow: this._e2eeEpoch, + }); + return; + } + + this._roomKey = unwrapped; await this._e2eeKeyProvider.setKey(roomKeyToBase64(this._roomKey)); log.info("E2EE: room key received and applied", { fromUserId }); @@ -1199,10 +1258,7 @@ export class LiveKitSession { } const wasKeyHolder = this._isKeyHolder; - // Use parseUserId for safe parsing of LiveKit identity strings. - const myUserId = this._room?.localParticipant?.identity - ? parseUserId(this._room.localParticipant.identity) - : 0; + const myUserId = authStore.getState().user?.id ?? 0; if (myUserId !== 0 && lowestUserId === myUserId && !wasKeyHolder) { // Prevent concurrent rotations (e.g. two participants leave in rapid succession). @@ -1216,15 +1272,20 @@ export class LiveKitSession { // Rotate the room key — generate a new one and distribute to all remaining peers. try { + this._e2eeEpoch++; this._roomKey = generateRoomKey(); await this._e2eeKeyProvider.setKey(roomKeyToBase64(this._roomKey)); - log.info("E2EE: rotated room key", { channelId }); + log.info("E2EE: rotated room key", { channelId, epoch: this._e2eeEpoch }); - // Wrap and send the new key to all remaining peers. - if (this._ecdhKeyPair) { - for (const [peerId, peerKey] of this._peerPublicKeys) { + // Snapshot peers before async loop — new peers that arrive during + // wrapping are handled by the post-rotation check below. + const keypair = this._ecdhKeyPair; + const peersSnapshot = new Map(this._peerPublicKeys); + + if (keypair) { + for (const [peerId, peerKey] of peersSnapshot) { const { encryptedKey, iv } = await wrapRoomKey( - this._ecdhKeyPair.privateKey, + keypair.privateKey, peerKey, this._roomKey, ); @@ -1234,8 +1295,27 @@ export class LiveKitSession { }); } log.info("E2EE: distributed rotated key to peers", { - peerCount: this._peerPublicKeys.size, + peerCount: peersSnapshot.size, }); + + // H3: Check for peers that arrived during the rotation loop and + // send them the new key too. + if (keypair === this._ecdhKeyPair && this._roomKey) { + for (const [peerId, peerKey] of this._peerPublicKeys) { + if (!peersSnapshot.has(peerId)) { + const { encryptedKey, iv } = await wrapRoomKey( + keypair.privateKey, + peerKey, + this._roomKey, + ); + this.ws?.send({ + type: "voice_e2ee_offer", + payload: { target_user_id: peerId, encrypted_key: encryptedKey, iv }, + }); + log.info("E2EE: sent rotated key to late-arriving peer", { peerId }); + } + } + } } } catch (err) { log.error("E2EE: failed to rotate room key", err); @@ -1252,6 +1332,7 @@ export class LiveKitSession { this._peerPublicKeys.clear(); this._isKeyHolder = false; this._rotatingKey = false; + this._e2eeEpoch = 0; this._pendingAnnounces.length = 0; // Reject (not resolve) so waiting connectAndSetup sees a failure, not a // silent success with no room key. diff --git a/Server/ws/voice_e2ee.go b/Server/ws/voice_e2ee.go index 3b871b7e..0dde9dab 100644 --- a/Server/ws/voice_e2ee.go +++ b/Server/ws/voice_e2ee.go @@ -2,6 +2,7 @@ package ws import ( "context" + "encoding/base64" "encoding/json" "log/slog" ) @@ -32,6 +33,10 @@ func (h *Hub) handleVoiceE2EEAnnounce(_ context.Context, c *Client, payload json c.sendMsg(buildErrorMsg(ErrCodeBadPayload, "public_key too large")) return } + if _, err := base64.StdEncoding.DecodeString(p.PublicKey); err != nil { + c.sendMsg(buildErrorMsg(ErrCodeBadPayload, "public_key is not valid base64")) + return + } // Store the public key on the client for later retrieval by new joiners. c.setE2EEPubKey(p.PublicKey) @@ -68,16 +73,28 @@ func (h *Hub) handleVoiceE2EEOffer(_ context.Context, c *Client, payload json.Ra c.sendMsg(buildErrorMsg(ErrCodeBadPayload, "encrypted_key or iv too large")) return } + if _, err := base64.StdEncoding.DecodeString(p.EncryptedKey); err != nil { + c.sendMsg(buildErrorMsg(ErrCodeBadPayload, "encrypted_key is not valid base64")) + return + } + if _, err := base64.StdEncoding.DecodeString(p.IV); err != nil { + c.sendMsg(buildErrorMsg(ErrCodeBadPayload, "iv is not valid base64")) + return + } - // Verify the target is in the same voice channel. + // Verify the target is in the same voice channel — lookup and channel + // check must be atomic (under the same lock hold) to prevent TOCTOU races + // where the target leaves between lookup and the channel comparison. h.mu.RLock() target, ok := h.clients[p.TargetUserID] - h.mu.RUnlock() if !ok { + h.mu.RUnlock() c.sendMsg(buildErrorMsg(ErrCodeBadPayload, "target user not connected")) return } - if target.getVoiceChID() != voiceChID { + targetChID := target.getVoiceChID() + h.mu.RUnlock() + if targetChID != voiceChID { c.sendMsg(buildErrorMsg(ErrCodeForbidden, "target user not in your voice channel")) return }