diff --git a/src/webrtc/events/Close.ts b/src/webrtc/events/Close.ts
index e88b2938..bf2b4e22 100644
--- a/src/webrtc/events/Close.ts
+++ b/src/webrtc/events/Close.ts
@@ -18,8 +18,8 @@
import { WebSocket } from "@spacebar/gateway";
-export async function onClose(this: WebSocket, code: number, reason: string) {
+export async function onClose(socket: WebSocket, code: number, reason: Buffer) {
console.log("[WebRTC] closed", code, reason.toString());
- this.removeAllListeners();
+ socket.rawSocket.removeAllListeners();
}
diff --git a/src/webrtc/events/Connection.ts b/src/webrtc/events/Connection.ts
index f1808399..2949ea75 100644
--- a/src/webrtc/events/Connection.ts
+++ b/src/webrtc/events/Connection.ts
@@ -28,10 +28,11 @@ 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, rawSocket: WS, request: IncomingMessage) {
+ const socket = new WebRtcWebSocket(rawSocket);
try {
- socket.on("close", onClose.bind(socket));
- socket.on("message", onMessage.bind(socket));
+ socket.rawSocket.on("close", (code, reason) => onClose(socket, code, reason));
+ socket.rawSocket.on("message", (data, isBinary) => onMessage(socket, data as Buffer));
console.log("[WebRTC] new connection", request.url);
if (process.env.WS_LOGEVENTS) {
@@ -45,7 +46,7 @@ export async function Connection(this: WS.Server, socket: WebRtcWebSocket, reque
"pong",
"unexpected-response",
].forEach((x) => {
- socket.on(x, (y) => console.log("[WebRTC]", x, y));
+ socket.rawSocket.on(x, (y) => console.log("[WebRTC]", x, y));
});
}
@@ -53,11 +54,11 @@ export async function Connection(this: WS.Server, socket: WebRtcWebSocket, reque
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.rawSocket.close(CLOSECODES.Unknown_error, "invalid version");
setHeartbeat(socket);
- socket.readyTimeout = setTimeout(() => socket.close(CLOSECODES.Session_timed_out), 1000 * 30);
+ socket.readyTimeout = setTimeout(() => socket.rawSocket.close(CLOSECODES.Session_timed_out), 1000 * 30);
await Send(socket, {
op: VoiceOPCodes.HELLO,
@@ -67,6 +68,6 @@ export async function Connection(this: WS.Server, socket: WebRtcWebSocket, reque
});
} catch (error) {
console.error("[WebRTC]", error);
- return socket.close(CLOSECODES.Unknown_error);
+ return socket.rawSocket.close(CLOSECODES.Unknown_error);
}
}
diff --git a/src/webrtc/events/Message.ts b/src/webrtc/events/Message.ts
index b5455d97..0ec2108a 100644
--- a/src/webrtc/events/Message.ts
+++ b/src/webrtc/events/Message.ts
@@ -20,10 +20,10 @@ import { CLOSECODES } from "@spacebar/gateway";
import OPCodeHandlers from "../opcodes";
import { VoiceOPCodes, VoicePayload, WebRtcWebSocket } from "../util";
-export async function onMessage(this: WebRtcWebSocket, buffer: Buffer) {
+export async function onMessage(socket: 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 && !socket.user_id) return socket.rawSocket.close(CLOSECODES.Not_authenticated);
const OPCodeHandler = OPCodeHandlers[data.op];
if (!OPCodeHandler) {
@@ -37,7 +37,7 @@ export async function onMessage(this: WebRtcWebSocket, buffer: Buffer) {
console.log("[WebRTC] Opcode " + VoiceOPCodes[data.op]);
}
- return await OPCodeHandler.call(this, data);
+ return await OPCodeHandler.call(socket, data);
} catch (error) {
console.error("[WebRTC] error", error);
// if (!this.CLOSED && this.CLOSING) return this.close(CloseCodes.Unknown_error);
|