fix: harden E2EE key exchange — election, rotation, queuing, error handling

- Use deterministic key holder election (lowest user_id) instead of
  Map insertion order which is not guaranteed to match join order
- Use parseUserId() instead of raw parseInt() for LiveKit identity parsing
- Add concurrent key rotation guard (_rotatingKey flag) to prevent
  races when multiple participants leave in rapid succession
- Queue voice_e2ee_announce messages that arrive before ECDH keypair
  is ready; drain after keypair generation in connectAndSetup
- Propagate decryption failures to roomKeyResolver so connectAndSetup
  unblocks with an error instead of hanging
- Reject (not resolve) roomKeyResolver on voice leave for proper cleanup
- Convert dynamic await import("@lib/e2eeCrypto") to static imports
- Add VOICE_E2EE_ANNOUNCE/OFFER to protocolTypes.ts enum constants
- Use typed S.VOICE_E2EE_* constants in dispatcher instead of string casts
- Add payload size limits for encrypted_key (1024) and iv (128) on server

https://claude.ai/code/session_01KKo3RwjdmcNzkgXNfUkgNT
This commit is contained in:
Claude
2026-04-04 19:13:51 +00:00
parent 0c4d9f702c
commit b89a7efa6a
4 changed files with 86 additions and 28 deletions
+2 -2
View File
@@ -389,13 +389,13 @@ export function wireDispatcher(ws: WsClient): DispatcherCleanup {
// ── Voice E2EE (client-side ECDH key exchange) ────────
unsubs.push(
ws.on("voice_e2ee_announce" as S, (payload: { user_id: number; public_key: string }) => {
ws.on(S.VOICE_E2EE_ANNOUNCE, (payload) => {
void handleE2EEAnnounce(payload.user_id, payload.public_key);
}),
);
unsubs.push(
ws.on("voice_e2ee_offer" as S, (payload: { from_user_id: number; encrypted_key: string; iv: string }) => {
ws.on(S.VOICE_E2EE_OFFER, (payload) => {
void handleE2EEOffer(payload.from_user_id, payload.encrypted_key, payload.iv);
}),
);
+74 -26
View File
@@ -15,6 +15,15 @@ import { createLogger } from "@lib/logger";
import { invoke } from "@tauri-apps/api/core";
import { AudioPipeline } from "@lib/audioPipeline";
import { AudioElements } from "@lib/audioElements";
import {
generateECDHKeyPair,
exportPublicKey,
importPublicKey,
generateRoomKey,
roomKeyToBase64,
wrapRoomKey,
unwrapRoomKey,
} from "@lib/e2eeCrypto";
import { DeviceManager } from "@lib/deviceManager";
import {
type VideoTrackDeps,
@@ -136,8 +145,13 @@ export class LiveKitSession {
private _peerPublicKeys: Map<number, CryptoKey> = new Map();
/** True if this client is the key holder (longest-present participant). */
private _isKeyHolder = false;
/** Resolver for non-key-holders waiting to receive the room key via offer. */
/** Resolver/rejector for non-key-holders waiting to receive the room key via offer. */
private _roomKeyResolver: (() => void) | null = null;
private _roomKeyRejector: ((err: Error) => void) | null = null;
/** Guard: true while a key rotation is in progress (prevents concurrent rotations). */
private _rotatingKey = false;
/** Announces that arrived before our ECDH keypair was ready. Drained after keypair init. */
private _pendingAnnounces: Array<{ userId: number; publicKeyBase64: string }> = [];
// --- State transition (single writer) ---
@@ -424,7 +438,6 @@ export class LiveKitSession {
// E2EE keys are exchanged client-side via ECDH. On reconnect, the
// room key is still in _roomKey from the previous session. Re-apply it.
if (this._roomKey) {
const { roomKeyToBase64 } = await import("@lib/e2eeCrypto");
// eslint-disable-next-line no-await-in-loop -- must set key before connect
await this._e2eeKeyProvider.setKey(roomKeyToBase64(this._roomKey));
}
@@ -794,16 +807,24 @@ export class LiveKitSession {
// ── Client-side E2EE key exchange (ECDH) ──────────────────────────
// Generate a fresh ECDH keypair for this session.
const { generateECDHKeyPair, exportPublicKey, generateRoomKey, roomKeyToBase64 } =
await import("@lib/e2eeCrypto");
this._ecdhKeyPair = await generateECDHKeyPair();
this._peerPublicKeys.clear();
const myPubKeyBase64 = await exportPublicKey(this._ecdhKeyPair.publicKey);
// Drain any announces that arrived before our keypair was ready.
// These are from existing participants whose public keys the server
// relayed during voice_join sync.
const queued = this._pendingAnnounces.splice(0);
for (const { userId: qId, publicKeyBase64: qKey } of queued) {
const peerKey = await importPublicKey(qKey);
this._peerPublicKeys.set(qId, peerKey);
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, we're first.
// sync — if we received none (including queued), we're first.
const existingPeerCount = this._peerPublicKeys.size;
this._isKeyHolder = existingPeerCount === 0;
@@ -816,8 +837,9 @@ export class LiveKitSession {
// Wait for the key holder to send us the room key via voice_e2ee_offer.
// This promise resolves when handleE2EEOffer() sets _roomKey.
log.info("E2EE: waiting for room key from key holder", { channelId });
const roomKeyPromise = new Promise<void>((resolve) => {
const roomKeyPromise = new Promise<void>((resolve, reject) => {
this._roomKeyResolver = resolve;
this._roomKeyRejector = reject;
});
// Don't block forever — timeout after 10 seconds.
const timeout = new Promise<void>((_, reject) =>
@@ -829,6 +851,7 @@ export class LiveKitSession {
log.warn("E2EE: key exchange timed out, proceeding without E2EE", { channelId });
}
this._roomKeyResolver = null;
this._roomKeyRejector = null;
}
// Announce our public key so existing participants (and the key holder)
@@ -1076,8 +1099,13 @@ export class LiveKitSession {
* the room key to them.
*/
async handleE2EEAnnounce(userId: number, publicKeyBase64: string): Promise<void> {
// Queue if our keypair isn't ready yet (announce arrived during connectAndSetup).
if (!this._ecdhKeyPair) {
this._pendingAnnounces.push({ userId, publicKeyBase64 });
log.info("E2EE: queued announce (keypair not ready)", { userId });
return;
}
try {
const { importPublicKey, wrapRoomKey } = await import("@lib/e2eeCrypto");
const peerKey = await importPublicKey(publicKeyBase64);
this._peerPublicKeys.set(userId, peerKey);
log.info("E2EE: received peer public key", { userId });
@@ -1110,7 +1138,6 @@ export class LiveKitSession {
ivBase64: string,
): Promise<void> {
try {
const { unwrapRoomKey, roomKeyToBase64 } = await import("@lib/e2eeCrypto");
const peerKey = this._peerPublicKeys.get(fromUserId);
if (!peerKey) {
log.warn("E2EE: received offer from unknown peer", { fromUserId });
@@ -1134,22 +1161,30 @@ export class LiveKitSession {
if (this._roomKeyResolver) {
this._roomKeyResolver();
this._roomKeyResolver = null;
this._roomKeyRejector = null;
}
} catch (err) {
log.error("E2EE: failed to handle offer", err);
// Propagate decryption failure so the waiting connectAndSetup unblocks.
if (this._roomKeyRejector) {
this._roomKeyRejector(err instanceof Error ? err : new Error(String(err)));
this._roomKeyResolver = null;
this._roomKeyRejector = null;
}
}
}
/**
* Handle a participant leaving the voice channel. If we become the new key
* holder, rotate the room key and distribute to remaining peers.
*
* Key holder election: the participant with the lowest user ID among remaining
* participants is elected. This is deterministic and does not depend on Map
* insertion order (which is not guaranteed to match server join order).
*/
async handleParticipantLeft(userId: number): Promise<void> {
this._peerPublicKeys.delete(userId);
// Determine if we should become the new key holder.
// Key holder is the longest-present participant — determined by voice state
// list order. The first user in the voiceUsers map for our channel is the holder.
const channelId = this._currentChannelId;
if (!channelId) return;
@@ -1157,24 +1192,30 @@ export class LiveKitSession {
const channelUsers = state.voiceUsers.get(channelId);
if (!channelUsers || channelUsers.size === 0) return;
// The first user in the map is the longest-present (server sends in join order).
const firstUserId = channelUsers.keys().next().value;
// Get our own user ID from the ws client.
// If we're the first user remaining, we're the new key holder.
const wasKeyHolder = this._isKeyHolder;
// We need to know our own user ID — derive from the room's local participant.
const myUserId = this._room?.localParticipant?.identity
? parseInt(this._room.localParticipant.identity, 10)
: null;
// Elect key holder: lowest user_id among remaining participants.
let lowestUserId = Infinity;
for (const uid of channelUsers.keys()) {
if (uid < lowestUserId) lowestUserId = uid;
}
if (myUserId !== null && firstUserId === myUserId && !wasKeyHolder) {
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;
if (myUserId !== 0 && lowestUserId === myUserId && !wasKeyHolder) {
// Prevent concurrent rotations (e.g. two participants leave in rapid succession).
if (this._rotatingKey) {
log.warn("E2EE: key rotation already in progress, skipping", { userId, channelId });
return;
}
this._rotatingKey = true;
this._isKeyHolder = true;
log.info("E2EE: became key holder after participant left", { userId, channelId });
// Rotate the room key — generate a new one and distribute to all remaining peers.
try {
const { generateRoomKey, roomKeyToBase64, wrapRoomKey } =
await import("@lib/e2eeCrypto");
this._roomKey = generateRoomKey();
await this._e2eeKeyProvider.setKey(roomKeyToBase64(this._roomKey));
log.info("E2EE: rotated room key", { channelId });
@@ -1198,6 +1239,8 @@ export class LiveKitSession {
}
} catch (err) {
log.error("E2EE: failed to rotate room key", err);
} finally {
this._rotatingKey = false;
}
}
}
@@ -1208,10 +1251,15 @@ export class LiveKitSession {
this._roomKey = null;
this._peerPublicKeys.clear();
this._isKeyHolder = false;
if (this._roomKeyResolver) {
this._roomKeyResolver();
this._roomKeyResolver = null;
this._rotatingKey = false;
this._pendingAnnounces.length = 0;
// Reject (not resolve) so waiting connectAndSetup sees a failure, not a
// silent success with no room key.
if (this._roomKeyRejector) {
this._roomKeyRejector(new Error("Voice session ended"));
}
this._roomKeyResolver = null;
this._roomKeyRejector = null;
}
/** Retry microphone permission after being in listen-only mode. */
@@ -39,6 +39,8 @@ export const ServerMessageType = {
PONG: "pong",
DM_CHANNEL_OPEN: "dm_channel_open",
DM_CHANNEL_CLOSE: "dm_channel_close",
VOICE_E2EE_ANNOUNCE: "voice_e2ee_announce",
VOICE_E2EE_OFFER: "voice_e2ee_offer",
} as const;
export type ServerMessageTypeValue = (typeof ServerMessageType)[keyof typeof ServerMessageType];
@@ -66,6 +68,8 @@ export const ClientMessageType = {
PING: "ping",
// Extension (not in protocol-schema.json but used in practice)
VOICE_TOKEN_REFRESH: "voice_token_refresh",
VOICE_E2EE_ANNOUNCE: "voice_e2ee_announce",
VOICE_E2EE_OFFER: "voice_e2ee_offer",
} as const;
export type ClientMessageTypeValue = (typeof ClientMessageType)[keyof typeof ClientMessageType];
+6
View File
@@ -62,6 +62,12 @@ func (h *Hub) handleVoiceE2EEOffer(_ context.Context, c *Client, payload json.Ra
c.sendMsg(buildErrorMsg(ErrCodeBadPayload, "target_user_id, encrypted_key, and iv are required"))
return
}
// Size limits: AES-256-GCM encrypted 32-byte key ≈ 64 base64 chars + 16-byte
// auth tag. 1024 chars is generous. IV is 12 bytes = 16 base64 chars.
if len(p.EncryptedKey) > 1024 || len(p.IV) > 128 {
c.sendMsg(buildErrorMsg(ErrCodeBadPayload, "encrypted_key or iv too large"))
return
}
// Verify the target is in the same voice channel.
h.mu.RLock()