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: [
|