wip working with vaultclient

This commit is contained in:
Thomas Ramé
2026-04-07 16:28:23 +02:00
parent 0c21978ca4
commit 9abb560d3c
2 changed files with 135 additions and 164 deletions
@@ -1,10 +1,7 @@
/**
* STEP 2b: E2EE Manager using a custom Worker with libsodium crypto.
* Same architecture as the built-in (streams transferred to Worker), but
* using XChaCha20-Poly1305 via libsodium instead of AES-GCM via crypto.subtle.
*
* This solves the receiver-reuse problem: transferred streams survive track
* changes, so the pipe in the Worker keeps working when a participant refreshes.
* E2EE Manager using VaultClient iframe for crypto.
* Preserves codec header bytes unencrypted (required for RTP packetization).
* TransformStream runs on main thread, crypto delegated to vault iframe.
*/
import { EventEmitter } from 'events'
import { Encryption_Type } from '@livekit/protocol'
@@ -23,49 +20,41 @@ enum EncryptionEvent {
EncryptionError = 'encryptionError',
}
// Hardcoded 32-byte key — same on all participants (Step 2)
const HARDCODED_KEY = new Uint8Array([
1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 12, 13, 14, 15, 16,
17, 18, 19, 20, 21, 22, 23, 24, 25, 26, 27, 28, 29, 30, 31, 32,
])
// VP8 unencrypted header bytes (same as LiveKit FrameCryptor)
const UNENCRYPTED_BYTES = { key: 10, delta: 3, audio: 1 }
function getUnencryptedBytes(frame: RTCEncodedVideoFrame | RTCEncodedAudioFrame): number {
if (!('type' in frame)) return UNENCRYPTED_BYTES.audio
return frame.type === 'key' ? UNENCRYPTED_BYTES.key : UNENCRYPTED_BYTES.delta
}
export class VaultE2EEManager extends EventEmitter {
private vaultClient: VaultClient
private room?: Room
private encryptionEnabled = false
private _isDataChannelEncryptionEnabled = false
private worker: Worker
private encryptedKeyBytes: Uint8Array | null = null
constructor(_vaultClient: VaultClient) {
constructor(vaultClient: VaultClient) {
super()
this.worker = new Worker(
new URL('./vault-e2ee.worker.ts', import.meta.url),
{ type: 'module' },
)
this.worker.onmessage = this.onWorkerMessage
this.worker.onerror = (ev) => {
console.error('[VaultE2EE] worker error:', ev)
}
// Send init AND key immediately — key doesn't need sodium, so the Worker
// can store it before WASM loads. This avoids the race where encode/decode
// messages arrive before the key is set.
this.worker.postMessage({ kind: 'init', data: {} })
this.worker.postMessage({ kind: 'setKey', data: { key: HARDCODED_KEY, participantIdentity: '__init__' } })
this.vaultClient = vaultClient
}
get isEnabled() {
return this.encryptionEnabled
}
get isEnabled() { return this.encryptionEnabled }
get isDataChannelEncryptionEnabled() {
return this.isEnabled && this._isDataChannelEncryptionEnabled
}
set isDataChannelEncryptionEnabled(enabled: boolean) {
this._isDataChannelEncryptionEnabled = enabled
}
setEncryptedSymmetricKey(_key: ArrayBuffer): void {
// No-op for Step 2 — using hardcoded key
private freshKeyBuffer(): ArrayBuffer {
return new Uint8Array(this.encryptedKeyBytes!).buffer
}
setEncryptedSymmetricKey(key: ArrayBuffer): void {
this.encryptedKeyBytes = new Uint8Array(new Uint8Array(key))
}
setup(room: Room): void {
@@ -78,109 +67,58 @@ export class VaultE2EEManager extends EventEmitter {
setupEngine(_engine: RTCEngine): void {}
setParticipantCryptorEnabled(enabled: boolean, participantIdentity: string): void {
// Send key to worker when enabling
if (enabled) {
this.worker.postMessage({
kind: 'setKey',
data: { key: HARDCODED_KEY, participantIdentity },
})
if (
participantIdentity === this.room?.localParticipant.identity &&
this.encryptionEnabled !== enabled
) {
this.encryptionEnabled = enabled
this.emit(EncryptionEvent.ParticipantEncryptionStatusChanged, enabled, this.room!.localParticipant)
} else if (participantIdentity !== this.room?.localParticipant.identity) {
const p = this.room?.getParticipantByIdentity(participantIdentity)
if (p) this.emit(EncryptionEvent.ParticipantEncryptionStatusChanged, enabled, p)
}
this.worker.postMessage({
kind: 'enable',
data: { enabled, participantIdentity },
})
}
setSifTrailer(_trailer: Uint8Array): void {}
async encryptData(_data: Uint8Array) {
return { uuid: crypto.randomUUID(), payload: _data, iv: new Uint8Array(0), keyIndex: 0 }
async encryptData(data: Uint8Array) {
if (!this.encryptedKeyBytes) throw new Error('No key')
const r = await this.vaultClient.encryptWithKey(data.slice().buffer, this.freshKeyBuffer())
return { uuid: crypto.randomUUID(), payload: new Uint8Array(r.encryptedData).slice(), iv: new Uint8Array(0), keyIndex: 0 }
}
async handleEncryptedData(payload: Uint8Array) {
return { uuid: crypto.randomUUID(), payload }
async handleEncryptedData(payload: Uint8Array, _iv: Uint8Array, _id: string, _idx: number) {
if (!this.encryptedKeyBytes) throw new Error('No key')
const r = await this.vaultClient.decryptWithKey(payload.slice().buffer, this.freshKeyBuffer())
return { uuid: crypto.randomUUID(), payload: new Uint8Array(r.data).slice() }
}
// ── Worker messages ─────────────────────────────────────────────────
private onWorkerMessage = (ev: MessageEvent) => {
const { kind, data } = ev.data
switch (kind) {
case 'initAck':
console.info('[VaultE2EE] STEP 2b: Worker ready (libsodium)')
break
case 'enable':
if (
this.encryptionEnabled !== data.enabled &&
data.participantIdentity === this.room?.localParticipant.identity
) {
this.encryptionEnabled = data.enabled
this.emit(EncryptionEvent.ParticipantEncryptionStatusChanged, data.enabled, this.room!.localParticipant)
} else if (data.participantIdentity && data.participantIdentity !== '__init__') {
const p = this.room?.getParticipantByIdentity(data.participantIdentity)
if (p) this.emit(EncryptionEvent.ParticipantEncryptionStatusChanged, data.enabled, p)
}
break
case 'error':
this.emit(EncryptionEvent.EncryptionError, data.error, data.participantIdentity)
break
case 'pipeDead':
// Pipe ended (stream closed during disconnect). Clear E2EE_FLAG on all
// receivers so the next TrackSubscribed creates a fresh pipe.
this.room?.remoteParticipants.forEach((p) => {
p.trackPublications.forEach((pub) => {
if (pub.track?.receiver && E2EE_FLAG in pub.track.receiver) {
// @ts-expect-error
delete pub.track.receiver[E2EE_FLAG]
}
})
})
break
}
}
// ── Event listeners (same as built-in) ──────────────────────────────
// ── Events (same as built-in E2EEManager) ───────────────────────────
private setupEventListeners(room: Room): void {
room.on(RoomEvent.TrackPublished, (pub, participant) => {
this.setParticipantCryptorEnabled(
pub.trackInfo!.encryption !== Encryption_Type.NONE,
participant.identity,
)
pub.trackInfo!.encryption !== Encryption_Type.NONE, participant.identity)
})
room.on(RoomEvent.ConnectionStateChanged, (state) => {
if (state === ConnectionState.Connected) {
room.remoteParticipants.forEach((participant) => {
participant.trackPublications.forEach((pub) => {
room.remoteParticipants.forEach((p) => {
p.trackPublications.forEach((pub) => {
this.setParticipantCryptorEnabled(
pub.trackInfo!.encryption !== Encryption_Type.NONE,
participant.identity,
)
pub.trackInfo!.encryption !== Encryption_Type.NONE, p.identity)
})
})
}
})
room.on(RoomEvent.TrackUnsubscribed, (track, _, participant) => {
this.worker.postMessage({
kind: 'removeTransform',
data: { participantIdentity: participant.identity, trackId: track.mediaStreamID },
})
})
room.on(RoomEvent.TrackSubscribed, (track, pub, participant) => {
this.setupReceiver(track, participant.identity, pub.trackInfo)
room.on(RoomEvent.TrackSubscribed, (track, _pub, participant) => {
this.setupReceiver(track, participant.identity)
})
room.on(RoomEvent.SignalConnected, () => {
this.setParticipantCryptorEnabled(
room.localParticipant.isE2EEEnabled,
room.localParticipant.identity,
)
room.localParticipant.isE2EEEnabled, room.localParticipant.identity)
})
room.localParticipant.on(
@@ -191,79 +129,114 @@ export class VaultE2EEManager extends EventEmitter {
)
}
// ── Sender/Receiver — streams transferred to Worker ─────────────────
// ── Sender ──────────────────────────────────────────────────────────
private setupSender(sender: RTCRtpSender, trackId: string): void {
private setupSender(sender: RTCRtpSender, _trackId: string): void {
if (E2EE_FLAG in sender) return
if (!this.room?.localParticipant.identity) return
// @ts-expect-error
const senderStreams = sender.createEncodedStreams()
const streams = sender.createEncodedStreams()
this.worker.postMessage(
{
kind: 'encode',
data: {
readableStream: senderStreams.readable,
writableStream: senderStreams.writable,
trackId,
participantIdentity: this.room.localParticipant.identity,
},
const transformStream = new TransformStream({
transform: async (frame: RTCEncodedVideoFrame | RTCEncodedAudioFrame, controller: TransformStreamDefaultController) => {
try {
if (!this.encryptedKeyBytes || !frame.data || frame.data.byteLength === 0) {
return controller.enqueue(frame)
}
const unencryptedBytes = getUnencryptedBytes(frame)
const header = new Uint8Array(frame.data, 0, unencryptedBytes)
const payload = new Uint8Array(frame.data, unencryptedBytes)
// Vault encrypts payload → returns [nonce][ciphertext+MAC]
const { encryptedData } = await this.vaultClient.encryptWithKey(
payload.slice().buffer,
this.freshKeyBuffer(),
)
const encrypted = new Uint8Array(encryptedData)
// Frame: [header][encrypted payload]
const newData = new Uint8Array(header.byteLength + encrypted.byteLength)
newData.set(header)
newData.set(encrypted, header.byteLength)
frame.data = newData.buffer
controller.enqueue(frame)
} catch {
// Drop — never send unencrypted
}
},
[senderStreams.readable, senderStreams.writable],
)
})
streams.readable.pipeThrough(transformStream).pipeTo(streams.writable)
// @ts-expect-error
sender[E2EE_FLAG] = true
}
private setupReceiver(track: RemoteTrack, participantIdentity: string, trackInfo?: { mimeType?: string }): void {
// ── Receiver ────────────────────────────────────────────────────────
private setupReceiver(track: RemoteTrack, participantIdentity: string): void {
if (!track.receiver) return
const receiver = track.receiver
if (E2EE_FLAG in receiver) return
if (E2EE_FLAG in receiver) {
// Receiver reuse — Worker's existing pipe handles new track frames
this.worker.postMessage({
kind: 'updateCodec',
data: {
trackId: track.mediaStreamID,
participantIdentity,
codec: trackInfo?.mimeType?.split('/')[1],
},
})
return
}
// @ts-expect-error
let writable: WritableStream = receiver.writableStream
// @ts-expect-error
let readable: ReadableStream = receiver.readableStream
let writable: WritableStream
let readable: ReadableStream
try {
if (!writable || !readable) {
// @ts-expect-error
const receiverStreams = receiver.createEncodedStreams()
writable = receiverStreams.writable
readable = receiverStreams.readable
} catch {
// createEncodedStreams() already called (receiver reuse after pipe death).
// Cannot re-create streams — this receiver is stuck.
console.warn(`[VaultE2EE] Cannot create encoded streams for ${participantIdentity} (receiver reuse). Pipe unrecoverable.`)
return
const streams = receiver.createEncodedStreams()
// @ts-expect-error
receiver.writableStream = streams.writable
writable = streams.writable
// @ts-expect-error
receiver.readableStream = streams.readable
readable = streams.readable
}
this.worker.postMessage(
{
kind: 'decode',
data: {
readableStream: readable,
writableStream: writable,
trackId: track.mediaStreamID,
participantIdentity,
codec: trackInfo?.mimeType?.split('/')[1],
isReuse: false,
},
},
[readable, writable],
)
let successEmitted = false
const transformStream = new TransformStream({
transform: async (frame: RTCEncodedVideoFrame | RTCEncodedAudioFrame, controller: TransformStreamDefaultController) => {
try {
if (!this.encryptedKeyBytes || !frame.data || frame.data.byteLength === 0) {
return controller.enqueue(frame)
}
const unencryptedBytes = getUnencryptedBytes(frame)
const header = new Uint8Array(frame.data, 0, unencryptedBytes)
const encryptedPayload = new Uint8Array(frame.data, unencryptedBytes)
// Vault decrypts [nonce][ciphertext+MAC] → plaintext
const { data } = await this.vaultClient.decryptWithKey(
encryptedPayload.slice().buffer,
this.freshKeyBuffer(),
)
const plaintext = new Uint8Array(data)
// Frame: [header][plaintext]
const newData = new Uint8Array(header.byteLength + plaintext.byteLength)
newData.set(header)
newData.set(plaintext, header.byteLength)
frame.data = newData.buffer
controller.enqueue(frame)
if (!successEmitted) {
successEmitted = true
const p = this.room?.getParticipantByIdentity(participantIdentity)
if (p) this.emit(EncryptionEvent.ParticipantEncryptionStatusChanged, true, p)
}
} catch {
// Drop frame on decrypt error — keeps pipe alive
}
},
})
readable.pipeThrough(transformStream).pipeTo(writable).catch(() => {})
// @ts-expect-error
receiver[E2EE_FLAG] = true
}
@@ -199,11 +199,9 @@ export const Conference = ({
if (!encryptionEnabled || encryptionSetupComplete) return
if (useVaultE2EE) {
// Advanced mode: VaultE2EEManager handles its own Worker+KeyProvider internally
// Advanced mode: VaultE2EEManager delegates crypto to VaultClient iframe
const vaultManager = getVaultManager()
if (!vaultManager || !vaultClient) return
// setEncryptedSymmetricKey is a no-op for now (Step 1: hardcoded passphrase)
if (isAdmin) {
const existingKey = data?.encrypted_symmetric_key
if (existingKey) {