summary refs log tree commit diff
path: root/src/webrtc
diff options
context:
space:
mode:
Diffstat (limited to 'src/webrtc')
-rw-r--r--src/webrtc/Server.ts18
-rw-r--r--src/webrtc/events/Connection.ts9
-rw-r--r--src/webrtc/events/Message.ts9
-rw-r--r--src/webrtc/opcodes/Identify.ts29
-rw-r--r--src/webrtc/opcodes/SelectProtocol.ts27
-rw-r--r--src/webrtc/opcodes/Speaking.ts14
-rw-r--r--src/webrtc/opcodes/Video.ts86
7 files changed, 37 insertions, 155 deletions
diff --git a/src/webrtc/Server.ts b/src/webrtc/Server.ts

index 3a270e68..a20531d9 100644 --- a/src/webrtc/Server.ts +++ b/src/webrtc/Server.ts
@@ -21,13 +21,7 @@ import { closeDatabase, Config, initDatabase, initEvent } from "@spacebar/util"; import http from "http"; import ws from "ws"; import { Connection } from "./events/Connection"; -import { - loadWebRtcLibrary, - mediaServer, - WRTC_PORT_MAX, - WRTC_PORT_MIN, - WRTC_PUBLIC_IP, -} from "./util/MediaServer"; +import { loadWebRtcLibrary, mediaServer, WRTC_PORT_MAX, WRTC_PORT_MIN, WRTC_PUBLIC_IP } from "./util/MediaServer"; import { green, yellow } from "picocolors"; export class Server { @@ -36,15 +30,7 @@ export class Server { public server: http.Server; public production: boolean; - constructor({ - port, - server, - production, - }: { - port: number; - server?: http.Server; - production?: boolean; - }) { + constructor({ port, server, production }: { port: number; server?: http.Server; production?: boolean }) { this.port = port; this.production = production || false; diff --git a/src/webrtc/events/Connection.ts b/src/webrtc/events/Connection.ts
index a068a8fd..6e028c64 100644 --- a/src/webrtc/events/Connection.ts +++ b/src/webrtc/events/Connection.ts
@@ -28,11 +28,7 @@ import { onMessage } from "./Message"; // TODO: specify rate limit in config // TODO: check msg max size -export async function Connection( - this: WS.Server, - socket: WebRtcWebSocket, - request: IncomingMessage, -) { +export async function Connection(this: WS.Server, socket: WebRtcWebSocket, request: IncomingMessage) { try { socket.on("close", onClose.bind(socket)); socket.on("message", onMessage.bind(socket)); @@ -57,8 +53,7 @@ export async function Connection( socket.encoding = "json"; socket.version = Number(searchParams.get("v")) || 5; - if (socket.version < 3) - return socket.close(CLOSECODES.Unknown_error, "invalid version"); + if (socket.version < 3) return socket.close(CLOSECODES.Unknown_error, "invalid version"); setHeartbeat(socket); diff --git a/src/webrtc/events/Message.ts b/src/webrtc/events/Message.ts
index ccb88b72..93251079 100644 --- a/src/webrtc/events/Message.ts +++ b/src/webrtc/events/Message.ts
@@ -31,8 +31,7 @@ const PayloadSchema = { export async function onMessage(this: WebRtcWebSocket, buffer: Buffer) { try { const data: VoicePayload = JSON.parse(buffer.toString()); - if (data.op !== VoiceOPCodes.IDENTIFY && !this.user_id) - return this.close(CLOSECODES.Not_authenticated); + if (data.op !== VoiceOPCodes.IDENTIFY && !this.user_id) return this.close(CLOSECODES.Not_authenticated); const OPCodeHandler = OPCodeHandlers[data.op]; if (!OPCodeHandler) { @@ -42,11 +41,7 @@ export async function onMessage(this: WebRtcWebSocket, buffer: Buffer) { return; } - if ( - ![VoiceOPCodes.HEARTBEAT, VoiceOPCodes.SPEAKING].includes( - data.op as VoiceOPCodes, - ) - ) { + if (![VoiceOPCodes.HEARTBEAT, VoiceOPCodes.SPEAKING].includes(data.op as VoiceOPCodes)) { console.log("[WebRTC] Opcode " + VoiceOPCodes[data.op]); } diff --git a/src/webrtc/opcodes/Identify.ts b/src/webrtc/opcodes/Identify.ts
index 3abb26d6..8a2516c4 100644 --- a/src/webrtc/opcodes/Identify.ts +++ b/src/webrtc/opcodes/Identify.ts
@@ -17,29 +17,15 @@ */ import { CLOSECODES } from "@spacebar/gateway"; -import { - StreamSession, - VoiceState, -} from "@spacebar/util"; -import { - validateSchema, - VoiceIdentifySchema, -} from "@spacebar/schemas"; -import { - generateSsrc, - mediaServer, - Send, - VoiceOPCodes, - VoicePayload, - WebRtcWebSocket, -} from "@spacebar/webrtc"; +import { StreamSession, VoiceState } from "@spacebar/util"; +import { validateSchema, VoiceIdentifySchema } from "@spacebar/schemas"; +import { generateSsrc, mediaServer, Send, VoiceOPCodes, VoicePayload, WebRtcWebSocket } from "@spacebar/webrtc"; import { SSRCs } from "@spacebarchat/spacebar-webrtc-types"; import { subscribeToProducers } from "./Video"; export async function onIdentify(this: WebRtcWebSocket, data: VoicePayload) { clearTimeout(this.readyTimeout); - const { server_id, user_id, session_id, token, streams, video } = - validateSchema("VoiceIdentifySchema", data.d) as VoiceIdentifySchema; + const { server_id, user_id, session_id, token, streams, video } = validateSchema("VoiceIdentifySchema", data.d) as VoiceIdentifySchema; // server_id can be one of the following: a unique id for a GO Live stream, a channel id for a DM voice call, or a guild id for a guild voice channel // not sure if there's a way to determine whether a snowflake is a channel id or a guild id without checking if it exists in db @@ -92,12 +78,7 @@ export async function onIdentify(this: WebRtcWebSocket, data: VoicePayload) { this.type = type; const voiceRoomId = type === "stream" ? server_id : voiceState!.channel_id; - this.webRtcClient = await mediaServer.join( - voiceRoomId, - this.user_id, - this, - type!, - ); + this.webRtcClient = await mediaServer.join(voiceRoomId, this.user_id, this, type!); this.on("close", () => { // ice-lite media server relies on this to know when the peer went away diff --git a/src/webrtc/opcodes/SelectProtocol.ts b/src/webrtc/opcodes/SelectProtocol.ts
index 2aa841f2..a29bcacc 100644 --- a/src/webrtc/opcodes/SelectProtocol.ts +++ b/src/webrtc/opcodes/SelectProtocol.ts
@@ -16,34 +16,17 @@ along with this program. If not, see <https://www.gnu.org/licenses/>. */ import { SelectProtocolSchema, validateSchema } from "@spacebar/schemas"; -import { - VoiceOPCodes, - VoicePayload, - WebRtcWebSocket, - mediaServer, - Send, -} from "@spacebar/webrtc"; +import { VoiceOPCodes, VoicePayload, WebRtcWebSocket, mediaServer, Send } from "@spacebar/webrtc"; -export async function onSelectProtocol( - this: WebRtcWebSocket, - payload: VoicePayload, -) { +export async function onSelectProtocol(this: WebRtcWebSocket, payload: VoicePayload) { if (!this.webRtcClient) return; - const data = validateSchema( - "SelectProtocolSchema", - payload.d, - ) as SelectProtocolSchema; + const data = validateSchema("SelectProtocolSchema", payload.d) as SelectProtocolSchema; // UDP protocol not currently supported. Maybe in the future? - if (data.protocol !== "webrtc") - return this.close(4000, "only webrtc protocol supported currently"); + if (data.protocol !== "webrtc") return this.close(4000, "only webrtc protocol supported currently"); - const response = await mediaServer.onOffer( - this.webRtcClient, - data.sdp!, - data.codecs ?? [], - ); + const response = await mediaServer.onOffer(this.webRtcClient, data.sdp!, data.codecs ?? []); await Send(this, { op: VoiceOPCodes.SESSION_DESCRIPTION, diff --git a/src/webrtc/opcodes/Speaking.ts b/src/webrtc/opcodes/Speaking.ts
index bff0db97..fde63f34 100644 --- a/src/webrtc/opcodes/Speaking.ts +++ b/src/webrtc/opcodes/Speaking.ts
@@ -16,13 +16,7 @@ along with this program. If not, see <https://www.gnu.org/licenses/>. */ -import { - mediaServer, - VoiceOPCodes, - VoicePayload, - WebRtcWebSocket, - Send, -} from "../util"; +import { mediaServer, VoiceOPCodes, VoicePayload, WebRtcWebSocket, Send } from "../util"; // {"speaking":1,"delay":5,"ssrc":2805246727} @@ -30,11 +24,7 @@ export async function onSpeaking(this: WebRtcWebSocket, data: VoicePayload) { if (!this.webRtcClient) return; await Promise.all( - Array.from( - mediaServer.getClientsForRtcServer<WebRtcWebSocket>( - this.webRtcClient.voiceRoomId, - ), - ).map((client) => { + Array.from(mediaServer.getClientsForRtcServer<WebRtcWebSocket>(this.webRtcClient.voiceRoomId)).map((client) => { if (client.user_id === this.user_id) return Promise.resolve(); const ssrc = client.getOutgoingStreamSSRCsForUser(this.user_id); diff --git a/src/webrtc/opcodes/Video.ts b/src/webrtc/opcodes/Video.ts
index 1309fc4f..ed8870e4 100644 --- a/src/webrtc/opcodes/Video.ts +++ b/src/webrtc/opcodes/Video.ts
@@ -16,13 +16,7 @@ along with this program. If not, see <https://www.gnu.org/licenses/>. */ import { Stream } from "@spacebar/util"; -import { - mediaServer, - Send, - VoiceOPCodes, - VoicePayload, - WebRtcWebSocket, -} from "@spacebar/webrtc"; +import { mediaServer, Send, VoiceOPCodes, VoicePayload, WebRtcWebSocket } from "@spacebar/webrtc"; import type { WebRtcClient } from "@spacebarchat/spacebar-webrtc-types"; import { validateSchema, VoiceVideoSchema } from "@spacebar/schemas"; @@ -60,9 +54,7 @@ export async function onVideo(this: WebRtcWebSocket, payload: VoicePayload) { try { await Promise.race([ new Promise<void>((resolve, reject) => { - this.webRtcClient?.emitter.once("connected", () => - resolve(), - ); + this.webRtcClient?.emitter.once("connected", () => resolve()); }), new Promise<void>((resolve, reject) => { // Reject after 3 seconds if still not connected @@ -93,28 +85,19 @@ export async function onVideo(this: WebRtcWebSocket, payload: VoicePayload) { if (wantsToProduceAudio) { // check if we are already producing audio, if not, publish a new audio track for it if (!this.webRtcClient!.isProducingAudio()) { - console.log( - `[${this.user_id}] publishing new audio track ssrc:${d.audio_ssrc}`, - ); + console.log(`[${this.user_id}] publishing new audio track ssrc:${d.audio_ssrc}`); await this.webRtcClient.publishTrack("audio", { audio_ssrc: d.audio_ssrc, }); } // now check that all clients have subscribed to our audio - for (const client of mediaServer.getClientsForRtcServer<WebRtcWebSocket>( - voiceRoomId, - )) { + for (const client of mediaServer.getClientsForRtcServer<WebRtcWebSocket>(voiceRoomId)) { if (client.user_id === this.user_id) continue; if (!client.isSubscribedToTrack(this.user_id, "audio")) { - console.log( - `[${client.user_id}] subscribing to audio track ssrcs: ${d.audio_ssrc}`, - ); - await client.subscribeToTrack( - this.webRtcClient.user_id, - "audio", - ); + console.log(`[${client.user_id}] subscribing to audio track ssrcs: ${d.audio_ssrc}`); + await client.subscribeToTrack(this.webRtcClient.user_id, "audio"); clientsThatNeedUpdate.add(client); } @@ -125,9 +108,7 @@ export async function onVideo(this: WebRtcWebSocket, payload: VoicePayload) { this.webRtcClient!.videoStream = { ...stream, type: "video" }; // client sends "screen" on go live but expects "video" on response // check if we are already publishing video, if not, publish a new video track for it if (!this.webRtcClient!.isProducingVideo()) { - console.log( - `[${this.user_id}] publishing new video track ssrc:${d.video_ssrc}`, - ); + console.log(`[${this.user_id}] publishing new video track ssrc:${d.video_ssrc}`); await this.webRtcClient.publishTrack("video", { video_ssrc: d.video_ssrc, rtx_ssrc: d.rtx_ssrc, @@ -135,19 +116,12 @@ export async function onVideo(this: WebRtcWebSocket, payload: VoicePayload) { } // now check that all clients have subscribed to our video track - for (const client of mediaServer.getClientsForRtcServer<WebRtcWebSocket>( - voiceRoomId, - )) { + for (const client of mediaServer.getClientsForRtcServer<WebRtcWebSocket>(voiceRoomId)) { if (client.user_id === this.user_id) continue; if (!client.isSubscribedToTrack(this.user_id, "video")) { - console.log( - `[${client.user_id}] subscribing to video track ssrc: ${d.video_ssrc}`, - ); - await client.subscribeToTrack( - this.webRtcClient.user_id, - "video", - ); + console.log(`[${client.user_id}] subscribing to video track ssrc: ${d.video_ssrc}`); + await client.subscribeToTrack(this.webRtcClient.user_id, "video"); clientsThatNeedUpdate.add(client); } @@ -163,9 +137,7 @@ export async function onVideo(this: WebRtcWebSocket, payload: VoicePayload) { d: { user_id: this.user_id, // can never send audio ssrc as 0, it will mess up client state for some reason. send server generated ssrc as backup - audio_ssrc: - ssrcs.audio_ssrc ?? - this.webRtcClient!.getIncomingStreamSSRCs().audio_ssrc, + audio_ssrc: ssrcs.audio_ssrc ?? this.webRtcClient!.getIncomingStreamSSRCs().audio_ssrc, video_ssrc: ssrcs.video_ssrc ?? 0, rtx_ssrc: ssrcs.rtx_ssrc ?? 0, streams: d.streams?.map((x) => ({ @@ -181,14 +153,10 @@ export async function onVideo(this: WebRtcWebSocket, payload: VoicePayload) { } // check if we are not subscribed to producers in this server, if not, subscribe -export async function subscribeToProducers( - this: WebRtcWebSocket, -): Promise<void> { +export async function subscribeToProducers(this: WebRtcWebSocket): Promise<void> { if (!this.webRtcClient || !this.webRtcClient.webrtcConnected) return; - const clients = mediaServer.getClientsForRtcServer<WebRtcWebSocket>( - this.webRtcClient.voiceRoomId, - ); + const clients = mediaServer.getClientsForRtcServer<WebRtcWebSocket>(this.webRtcClient.voiceRoomId); await Promise.all( Array.from(clients).map(async (client) => { @@ -196,42 +164,26 @@ export async function subscribeToProducers( if (client.user_id === this.user_id) return; // cannot subscribe to self - if ( - client.isProducingAudio() && - !this.webRtcClient!.isSubscribedToTrack(client.user_id, "audio") - ) { - await this.webRtcClient!.subscribeToTrack( - client.user_id, - "audio", - ); + if (client.isProducingAudio() && !this.webRtcClient!.isSubscribedToTrack(client.user_id, "audio")) { + await this.webRtcClient!.subscribeToTrack(client.user_id, "audio"); needsUpdate = true; } - if ( - client.isProducingVideo() && - !this.webRtcClient!.isSubscribedToTrack(client.user_id, "video") - ) { - await this.webRtcClient!.subscribeToTrack( - client.user_id, - "video", - ); + if (client.isProducingVideo() && !this.webRtcClient!.isSubscribedToTrack(client.user_id, "video")) { + await this.webRtcClient!.subscribeToTrack(client.user_id, "video"); needsUpdate = true; } if (!needsUpdate) return; - const ssrcs = this.webRtcClient!.getOutgoingStreamSSRCsForUser( - client.user_id, - ); + const ssrcs = this.webRtcClient!.getOutgoingStreamSSRCsForUser(client.user_id); await Send(this, { op: VoiceOPCodes.VIDEO, d: { user_id: client.user_id, // can never send audio ssrc as 0, it will mess up client state for some reason. send server generated ssrc as backup - audio_ssrc: - ssrcs.audio_ssrc ?? - client.getIncomingStreamSSRCs().audio_ssrc, + audio_ssrc: ssrcs.audio_ssrc ?? client.getIncomingStreamSSRCs().audio_ssrc, video_ssrc: ssrcs.video_ssrc ?? 0, rtx_ssrc: ssrcs.rtx_ssrc ?? 0, streams: [