Start the implementation

This commit is contained in:
2025-10-11 14:05:44 +03:00
Unverified
parent 1b2f825a74
commit caf595e968
10 changed files with 1210 additions and 178 deletions
+434 -48
View File
@@ -1,21 +1,25 @@
import { getAuthHeaders } from "@/core/api/authApi";
import type { CallSignalingMessage, IceServersResponse } from "@/core/types";
import { request } from "@/core/websocket";
import { wrapCallSessionKeyForRecipient, unwrapCallSessionKeyFromSender } from "./encryption";
import { wrapCallSessionKeyForRecipient, unwrapCallSessionKeyFromSender, rotateCallSessionKey } from "./encryption";
import { fetchUserPublicKey } from "@/core/api/dmApi";
import { importAesGcmKey } from "@/utils/crypto/symmetric";
import E2EEWorker from "./e2eeWorker?worker";
import { rotateCallSessionKey, createSharedSecretAndDeriveSessionKey } from "./encryption";
export interface WebRTCCall {
peerConnection: RTCPeerConnection;
localStream: MediaStream | null;
remoteStream: MediaStream | null;
localVideoStream: MediaStream | null;
screenShareStream: MediaStream | null;
isInitiator: boolean;
remoteUserId: number;
remoteUsername: string;
isEnding?: boolean;
isMuted?: boolean;
isLocalVideoEnabled: boolean;
isScreenSharing: boolean;
isNegotiating?: boolean;
// Insertable Streams E2EE
sessionKey?: Uint8Array | null;
sessionCryptoKey?: CryptoKey | null;
@@ -28,6 +32,10 @@ export interface WebRTCCall {
export let authToken: string | null = null;
export let onCallStateChange: ((userId: number, state: string) => void) | null = null;
export let onRemoteStream: ((userId: number, stream: MediaStream) => void) | null = null;
export let onLocalVideoStream: ((userId: number, stream: MediaStream | null) => void) | null = null;
export let onRemoteVideoStream: ((userId: number, stream: MediaStream | null) => void) | null = null;
export let onLocalScreenShare: ((userId: number, stream: MediaStream | null) => void) | null = null;
export let onRemoteScreenShare: ((userId: number, stream: MediaStream | null) => void) | null = null;
const calls: Map<number, WebRTCCall> = new Map();
export function setAuthToken(token: string) {
@@ -42,6 +50,22 @@ export function setRemoteStreamHandler(handler: (userId: number, stream: MediaSt
onRemoteStream = handler;
}
export function setLocalVideoStreamHandler(handler: (userId: number, stream: MediaStream | null) => void) {
onLocalVideoStream = handler;
}
export function setRemoteVideoStreamHandler(handler: (userId: number, stream: MediaStream | null) => void) {
onRemoteVideoStream = handler;
}
export function setLocalScreenShareHandler(handler: (userId: number, stream: MediaStream | null) => void) {
onLocalScreenShare = handler;
}
export function setRemoteScreenShareHandler(handler: (userId: number, stream: MediaStream | null) => void) {
onRemoteScreenShare = handler;
}
async function sendSignalingMessage(message: CallSignalingMessage) {
if (!authToken) {
throw new Error("No auth token available");
@@ -99,10 +123,14 @@ async function createPeerConnection(userId: number): Promise<RTCPeerConnection>
peerConnection,
localStream: null,
remoteStream: null,
localVideoStream: null,
screenShareStream: null,
isInitiator: false,
remoteUserId: userId,
remoteUsername: "",
isMuted: false,
isLocalVideoEnabled: false,
isScreenSharing: false,
sessionKey: null,
sessionCryptoKey: null,
sessionId: crypto.randomUUID()
@@ -144,17 +172,103 @@ async function createPeerConnection(userId: number): Promise<RTCPeerConnection>
});
peerConnection.addEventListener("signalingstatechange", () => {
// Signaling state changed
console.log("Signaling state changed:", peerConnection.signalingState);
});
// Handle renegotiation when tracks are added/removed
peerConnection.addEventListener("negotiationneeded", async () => {
try {
console.log("Negotiation needed for user", userId);
const call = calls.get(userId);
if (!call) {
console.log("Skipping renegotiation - call not found");
return;
}
// Prevent multiple simultaneous negotiations
if (call.isNegotiating) {
console.log("Already negotiating, skipping");
return;
}
// Skip if we're in "stable" state and haven't finished the initial handshake
if (peerConnection.signalingState !== "stable") {
console.log("Skipping renegotiation - signaling state is", peerConnection.signalingState);
return;
}
call.isNegotiating = true;
console.log("Creating new offer for renegotiation (signalingState:", peerConnection.signalingState + ")");
const offer = await peerConnection.createOffer();
await peerConnection.setLocalDescription(offer);
console.log("Sending renegotiation offer to user", userId);
await sendSignalingMessage({
type: "call_offer",
fromUserId: 0,
toUserId: userId,
data: offer
});
call.isNegotiating = false;
} catch (error) {
console.error("Failed to handle negotiation:", error);
const call = calls.get(userId);
if (call) {
call.isNegotiating = false;
}
}
});
// Handle remote stream
peerConnection.addEventListener("track", (event) => {
peerConnection.addEventListener("track", async (event) => {
console.log("Received track:", event.track.kind, "from user", userId, "stream ID:", event.streams[0]?.id);
const [remoteStream] = event.streams;
const call = calls.get(userId);
if (call) {
call.remoteStream = remoteStream;
if (onRemoteStream) {
onRemoteStream(userId, remoteStream);
if (call && remoteStream) {
const track = event.track;
// Apply E2EE transform to the receiver for this new track if session key is available
// Currently disabled for video tracks due to counter synchronization issues
if (call.sessionKey && window.RTCRtpScriptTransform && track.kind === "audio") {
try {
const key = await importAesGcmKey(call.sessionKey);
const receiver = call.peerConnection.getReceivers().find(r => r.track === track);
if (receiver) {
console.log(`Applying decrypt transform to newly received ${track.kind} track`);
// @ts-ignore
receiver.transform = new RTCRtpScriptTransform(new E2EEWorker(), { key, mode: 'decrypt', sessionId: call.sessionId });
}
} catch (error) {
console.error("Failed to apply E2EE to received track:", error);
}
}
// Determine stream type based on track kind and stream ID
const streamId = remoteStream.id;
// Check if this is a screen share stream (we'll use a convention: screen share streams have "screen" in their ID)
if (streamId.includes("screen")) {
console.log("Detected screen share track, notifying handler");
// Handle remote screen share
if (onRemoteScreenShare) {
onRemoteScreenShare(userId, remoteStream);
}
} else if (track.kind === "video") {
console.log("Detected video track, notifying handler");
// Handle remote video
if (onRemoteVideoStream) {
onRemoteVideoStream(userId, remoteStream);
}
} else if (track.kind === "audio") {
console.log("Detected audio track, notifying handler");
// Handle remote audio (existing behavior)
call.remoteStream = remoteStream;
if (onRemoteStream) {
onRemoteStream(userId, remoteStream);
}
}
}
});
@@ -168,10 +282,10 @@ async function createPeerConnection(userId: number): Promise<RTCPeerConnection>
onCallStateChange(userId, peerConnection.connectionState);
}
// Clean up if connection failed or closed
// Clean up only on permanent failures
// Don't end on "disconnected" - ICE can recover from temporary disconnections
if (peerConnection.connectionState === "failed" ||
peerConnection.connectionState === "closed" ||
peerConnection.connectionState === "disconnected") {
peerConnection.connectionState === "closed") {
// Only send end call message if we're not already cleaning up
const call = calls.get(userId);
if (call && !call.isEnding) {
@@ -270,20 +384,37 @@ export async function sendWrappedCallSessionKey(userId: number, sessionKey: Uint
async function applyE2EETransforms(call: WebRTCCall): Promise<void> {
try {
// @ts-ignore
if (!call.sessionKey || !window.RTCRtpScriptTransform) return;
if (!call.sessionKey || !window.RTCRtpScriptTransform) {
console.log("Skipping E2EE transforms - session key or RTCRtpScriptTransform not available");
return;
}
console.log("Applying E2EE transforms for call", call.remoteUserId);
const key = await importAesGcmKey(call.sessionKey);
call.sessionCryptoKey = key;
const receiver = call.peerConnection.getReceivers().find(r => r.track && r.track.kind === 'audio');
if (receiver) {
// @ts-ignore
receiver.transform = new RTCRtpScriptTransform(new E2EEWorker(), { key, mode: 'decrypt' });
// Apply to audio receivers only (video E2EE requires synchronized counters)
const receivers = call.peerConnection.getReceivers();
for (const receiver of receivers) {
if (receiver.track && receiver.track.kind === "audio") {
console.log(`Applying decrypt transform to ${receiver.track.kind} receiver`);
// @ts-ignore
receiver.transform = new RTCRtpScriptTransform(new E2EEWorker(), { key, mode: 'decrypt', sessionId: call.sessionId });
}
}
const sender = call.peerConnection.getSenders().find(s => s.track && s.track.kind === 'audio');
if (sender) {
// @ts-ignore
sender.transform = new RTCRtpScriptTransform(new E2EEWorker(), { key, mode: 'encrypt' });
// Apply to audio senders only (video E2EE requires synchronized counters)
const senders = call.peerConnection.getSenders();
for (const sender of senders) {
if (sender.track && sender.track.kind === "audio") {
console.log(`Applying encrypt transform to ${sender.track.kind} sender`);
// @ts-ignore
sender.transform = new RTCRtpScriptTransform(new E2EEWorker(), { key, mode: 'encrypt', sessionId: call.sessionId });
}
}
} catch {}
} catch (error) {
console.error("Failed to apply E2EE transforms:", error);
}
}
export async function setSessionKey(userId: number, keyBytes: Uint8Array): Promise<void> {
@@ -341,19 +472,15 @@ export async function receiveWrappedSessionKey(fromUserId: number, wrappedPayloa
if (!senderPublicKey) return;
if (!wrappedPayload || !sessionKeyHash) return;
// First unwrap the session key from the encrypted payload (for validation)
await unwrapCallSessionKeyFromSender(senderPublicKey, {
// Unwrap the session key from the encrypted payload
const unwrappedSessionKey = await unwrapCallSessionKeyFromSender(senderPublicKey, {
salt: wrappedPayload.salt,
iv2: wrappedPayload.iv2,
wrapped: wrappedPayload.wrapped
});
// Then derive the actual session key from the shared secret
const call = calls.get(fromUserId);
const isInitiator = call?.isInitiator ?? false;
const derivedSessionKey = await createSharedSecretAndDeriveSessionKey(senderPublicKey, sessionKeyHash, isInitiator);
await setSessionKey(fromUserId, derivedSessionKey.key);
// Use the unwrapped session key directly (both sides should have the same key)
await setSessionKey(fromUserId, unwrappedSessionKey);
} catch (e) {
console.error("Failed to unwrap session key:", e);
}
@@ -369,10 +496,12 @@ export async function acceptCall(userId: number): Promise<boolean> {
if (!call) return false;
}
// Get user media and attach
const localStream = await navigator.mediaDevices.getUserMedia({ audio: true, video: false });
call.localStream = localStream;
localStream.getTracks().forEach(track => call!.peerConnection.addTrack(track, localStream));
// Get user media and attach (only if not already attached)
if (!call.localStream) {
const localStream = await navigator.mediaDevices.getUserMedia({ audio: true, video: false });
call.localStream = localStream;
localStream.getTracks().forEach(track => call!.peerConnection.addTrack(track, localStream));
}
// Notify initiator that callee accepted; initiator will generate offer
await sendSignalingMessage({
@@ -440,6 +569,10 @@ export async function onRemoteAccepted(userId: number): Promise<void> {
}
try {
// Small delay to ensure remote peer finishes processing the accept
// This prevents race conditions where our offer arrives before they're ready
await new Promise(resolve => setTimeout(resolve, 100));
// Create offer
const offer = await call.peerConnection.createOffer();
await call.peerConnection.setLocalDescription(offer);
@@ -459,39 +592,88 @@ export async function onRemoteAccepted(userId: number): Promise<void> {
async function createE2EETransform(sessionKey: NonNullable<WebRTCCall['sessionKey']>, peerConnection: RTCPeerConnection, sessionId?: string): Promise<void> {
try {
if (sessionKey && window.RTCRtpScriptTransform) {
const key = await importAesGcmKey(sessionKey);
const receiver = peerConnection.getReceivers().find(r => r.track && r.track.kind === 'audio');
if (receiver) {
// @ts-ignore
if (!sessionKey || !window.RTCRtpScriptTransform) {
console.log("Skipping E2EE transform in createE2EETransform - not supported or no session key");
return;
}
console.log("Creating E2EE transform with sessionId:", sessionId);
const key = await importAesGcmKey(sessionKey);
// Apply to audio receivers only (video E2EE requires synchronized counters)
const receivers = peerConnection.getReceivers();
for (const receiver of receivers) {
if (receiver.track && receiver.track.kind === "audio") {
console.log(`Applying decrypt transform to ${receiver.track.kind} in createE2EETransform`);
// @ts-ignore
receiver.transform = new RTCRtpScriptTransform(new E2EEWorker(), { key, mode: 'decrypt', sessionId });
}
const sender = peerConnection.getSenders().find(s => s.track && s.track.kind === 'audio');
if (sender) {
}
// Apply to audio senders only (video E2EE requires synchronized counters)
const senders = peerConnection.getSenders();
for (const sender of senders) {
if (sender.track && sender.track.kind === "audio") {
console.log(`Applying encrypt transform to ${sender.track.kind} in createE2EETransform`);
// @ts-ignore
sender.transform = new RTCRtpScriptTransform(new E2EEWorker(), { key, mode: 'encrypt', sessionId });
}
}
} catch (error) {
console.error("Failed to create E2EE transform:", error);
throw error;
// Don't throw - let the call continue without E2EE
}
}
export async function handleCallOffer(userId: number, offer: RTCSessionDescriptionInit): Promise<void> {
const call = calls.get(userId);
let call = calls.get(userId);
console.log("handleCallOffer called for user", userId, "offer type:", offer.type);
// Handle race condition - offer might arrive before peer connection is created
if (!call) {
throw new Error("No call found for offer");
console.log("No call found for offer, creating peer connection (race condition handling)");
await createPeerConnection(userId);
call = calls.get(userId);
if (!call) {
throw new Error("Failed to create call for offer");
}
}
try {
// Ensure we have local media before answering
if (!call.localStream) {
console.log("Getting local media for answer");
try {
const localStream = await navigator.mediaDevices.getUserMedia({ audio: true, video: false });
call.localStream = localStream;
localStream.getTracks().forEach(track => call!.peerConnection.addTrack(track, localStream));
} catch (mediaError) {
console.error("Failed to get local media:", mediaError);
// Continue anyway - we can still receive media
}
}
console.log("Setting remote description with", offer.sdp?.split('\n').filter(l => l.includes('m=')).join(', '));
// Set remote description
await call.peerConnection.setRemoteDescription(offer);
console.log("Creating answer...");
// Create answer
const answer = await call.peerConnection.createAnswer();
await call.peerConnection.setLocalDescription(answer);
// Attach transforms on callee side if session key set and insertable streams supported
await createE2EETransform(call.sessionKey!, call.peerConnection, call.sessionId);
console.log("Answer created with", answer.sdp?.split('\n').filter(l => l.includes('m=')).join(', '));
// Attach transforms on callee side if session key is available
// If not available yet, setSessionKey will apply them when it arrives
if (call.sessionKey) {
await createE2EETransform(call.sessionKey, call.peerConnection, call.sessionId);
} else {
console.log("Session key not yet available in handleCallOffer - will apply transforms when key arrives");
}
// Send answer to remote peer
await sendSignalingMessage({
@@ -500,6 +682,8 @@ export async function handleCallOffer(userId: number, offer: RTCSessionDescripti
toUserId: userId,
data: answer
});
console.log("Answer sent successfully");
} catch (error) {
console.error("Failed to handle offer:", error);
throw error;
@@ -514,8 +698,17 @@ export async function handleCallAnswer(userId: number, answer: RTCSessionDescrip
try {
await call.peerConnection.setRemoteDescription(answer);
// After signaling completes, attach transforms on initiator side if supported
await createE2EETransform(call.sessionKey!, call.peerConnection, call.sessionId);
// Reset negotiating flag
call.isNegotiating = false;
// Attach transforms on initiator side if session key is available
// If not available yet, setSessionKey will apply them when it arrives
if (call.sessionKey) {
await createE2EETransform(call.sessionKey, call.peerConnection, call.sessionId);
} else {
console.log("Session key not yet available in handleCallAnswer - will apply transforms when key arrives");
}
} catch (error) {
console.error("Failed to handle answer:", error);
throw error;
@@ -523,9 +716,10 @@ export async function handleCallAnswer(userId: number, answer: RTCSessionDescrip
}
export async function handleIceCandidate(userId: number, candidate: RTCIceCandidateInit): Promise<void> {
const call = calls.get(userId);
let call = calls.get(userId);
if (!call) {
console.warn("No call found for ICE candidate from user", userId);
console.warn("No call found for ICE candidate from user", userId, "- might arrive before connection setup");
// Don't create peer connection here - ICE candidates will be gathered again after connection is established
return;
}
@@ -616,6 +810,188 @@ export function getCall(userId: number): WebRTCCall | undefined {
return calls.get(userId);
}
export async function toggleVideo(userId: number): Promise<boolean> {
const call = calls.get(userId);
if (!call) {
return false;
}
if (!call.isLocalVideoEnabled) {
// Enable video
try {
const videoStream = await navigator.mediaDevices.getUserMedia({
video: true,
audio: false
});
call.localVideoStream = videoStream;
call.isLocalVideoEnabled = true;
// Add video track to peer connection
const videoTrack = videoStream.getVideoTracks()[0];
call.peerConnection.addTrack(videoTrack, videoStream);
console.log("Video track added successfully");
console.log("Current senders:", call.peerConnection.getSenders().map(s => s.track?.kind));
console.log("Current transceivers:", call.peerConnection.getTransceivers().map(t => ({
sender: t.sender.track?.kind,
receiver: t.receiver.track?.kind,
direction: t.direction
})));
// Note: E2EE for video is temporarily disabled for testing
// Will re-enable with proper counter synchronization
console.log("Video sent without E2EE (will implement synchronized encryption)");
// Notify local video stream handler
if (onLocalVideoStream) {
console.log("Calling onLocalVideoStream handler with stream:", videoStream);
onLocalVideoStream(userId, videoStream);
} else {
console.warn("onLocalVideoStream handler is not set!");
}
// Send signaling message to notify remote peer
console.log("Sending call_video_toggle with enabled: true");
await sendSignalingMessage({
type: "call_video_toggle",
fromUserId: 0,
toUserId: userId,
data: { enabled: true }
});
console.log("Video enabled successfully");
return true;
} catch (error) {
console.error("Failed to enable video:", error);
return false;
}
} else {
// Disable video
if (call.localVideoStream) {
call.localVideoStream.getTracks().forEach(track => {
track.stop();
// Remove track from peer connection
const senders = call.peerConnection.getSenders();
const videoSender = senders.find(s => s.track === track);
if (videoSender) {
call.peerConnection.removeTrack(videoSender);
}
});
call.localVideoStream = null;
}
call.isLocalVideoEnabled = false;
// Notify local video stream handler
if (onLocalVideoStream) {
onLocalVideoStream(userId, null);
}
// Send signaling message to notify remote peer
await sendSignalingMessage({
type: "call_video_toggle",
fromUserId: 0,
toUserId: userId,
data: { enabled: false }
});
return false;
}
}
export async function toggleScreenShare(userId: number): Promise<boolean> {
const call = calls.get(userId);
if (!call) {
return false;
}
if (!call.isScreenSharing) {
// Enable screen sharing
try {
// @ts-ignore - getDisplayMedia might not be in all TypeScript versions
const screenStream = await navigator.mediaDevices.getDisplayMedia({
video: true,
audio: false
});
// Set a special ID to identify screen share streams
Object.defineProperty(screenStream, "id", {
value: `screen-${crypto.randomUUID()}`,
writable: false
});
call.screenShareStream = screenStream;
call.isScreenSharing = true;
// Add screen share track to peer connection
const videoTrack = screenStream.getVideoTracks()[0];
// Handle when user stops sharing via browser UI
videoTrack.addEventListener("ended", () => {
toggleScreenShare(userId);
});
call.peerConnection.addTrack(videoTrack, screenStream);
console.log("Screen share track added successfully");
// Note: E2EE for screen share is temporarily disabled for testing
// Will re-enable with proper counter synchronization
console.log("Screen share sent without E2EE (will implement synchronized encryption)");
// Notify local screen share handler
if (onLocalScreenShare) {
onLocalScreenShare(userId, screenStream);
}
// Send signaling message to notify remote peer
await sendSignalingMessage({
type: "call_screen_share_toggle",
fromUserId: 0,
toUserId: userId,
data: { enabled: true }
});
return true;
} catch (error) {
console.error("Failed to enable screen sharing:", error);
return false;
}
} else {
// Disable screen sharing
if (call.screenShareStream) {
call.screenShareStream.getTracks().forEach(track => {
track.stop();
// Remove track from peer connection
const senders = call.peerConnection.getSenders();
const screenSender = senders.find(s => s.track === track);
if (screenSender) {
call.peerConnection.removeTrack(screenSender);
}
});
call.screenShareStream = null;
}
call.isScreenSharing = false;
// Notify local screen share handler
if (onLocalScreenShare) {
onLocalScreenShare(userId, null);
}
// Send signaling message to notify remote peer
await sendSignalingMessage({
type: "call_screen_share_toggle",
fromUserId: 0,
toUserId: userId,
data: { enabled: false }
});
return false;
}
}
export function cleanupCall(userId: number): void {
const call = calls.get(userId);
if (call) {
@@ -634,6 +1010,16 @@ export function cleanupCall(userId: number): void {
call.localStream.getTracks().forEach(track => track.stop());
}
// Stop local video stream
if (call.localVideoStream) {
call.localVideoStream.getTracks().forEach(track => track.stop());
}
// Stop screen share stream
if (call.screenShareStream) {
call.screenShareStream.getTracks().forEach(track => track.stop());
}
calls.delete(userId);
}
}