feat: increased turn cred duration, and fixed gateway crashing

This commit is contained in:
=
2025-03-09 21:15:10 +05:30
parent f68602280e
commit b055cda64d
2 changed files with 53 additions and 39 deletions
+52 -38
View File
@@ -96,6 +96,7 @@ export const pingGatewayAndVerify = async ({
error: err as Error error: err as Error
}); });
}); });
for (let attempt = 1; attempt <= maxRetries; attempt += 1) { for (let attempt = 1; attempt <= maxRetries; attempt += 1) {
try { try {
const stream = quicClient.connection.newStream("bidi"); const stream = quicClient.connection.newStream("bidi");
@@ -108,17 +109,13 @@ export const pingGatewayAndVerify = async ({
const { value, done } = await reader.read(); const { value, done } = await reader.read();
if (done) { if (done) {
throw new BadRequestError({ throw new Error("Gateway closed before receiving PONG");
message: "Gateway closed before receiving PONG"
});
} }
const response = Buffer.from(value).toString(); const response = Buffer.from(value).toString();
if (response !== "PONG\n" && response !== "PONG") { if (response !== "PONG\n" && response !== "PONG") {
throw new BadRequestError({ throw new Error(`Failed to Ping. Unexpected response: ${response}`);
message: `Failed to Ping. Unexpected response: ${response}`
});
} }
reader.releaseLock(); reader.releaseLock();
@@ -146,6 +143,7 @@ interface TProxyServer {
server: net.Server; server: net.Server;
port: number; port: number;
cleanup: () => Promise<void>; cleanup: () => Promise<void>;
getProxyError: () => string;
} }
const setupProxyServer = async ({ const setupProxyServer = async ({
@@ -170,6 +168,7 @@ const setupProxyServer = async ({
error: err as Error error: err as Error
}); });
}); });
const proxyErrorMsg = [""];
return new Promise((resolve, reject) => { return new Promise((resolve, reject) => {
const server = net.createServer(); const server = net.createServer();
@@ -185,31 +184,33 @@ const setupProxyServer = async ({
const forwardWriter = stream.writable.getWriter(); const forwardWriter = stream.writable.getWriter();
await forwardWriter.write(Buffer.from(`FORWARD-TCP ${targetHost}:${targetPort}\n`)); await forwardWriter.write(Buffer.from(`FORWARD-TCP ${targetHost}:${targetPort}\n`));
forwardWriter.releaseLock(); forwardWriter.releaseLock();
/* eslint-disable @typescript-eslint/no-misused-promises */
// Set up bidirectional copy // Set up bidirectional copy
const setupCopy = async () => { const setupCopy = () => {
// Client to QUIC // Client to QUIC
// eslint-disable-next-line // eslint-disable-next-line
(async () => { (async () => {
try { const writer = stream.writable.getWriter();
const writer = stream.writable.getWriter();
// Create a handler for client data // Create a handler for client data
clientConn.on("data", async (chunk) => { clientConn.on("data", (chunk) => {
await writer.write(chunk); writer.write(chunk).catch((err) => {
proxyErrorMsg.push((err as Error)?.message);
}); });
});
// Handle client connection close // Handle client connection close
clientConn.on("end", async () => { clientConn.on("end", () => {
await writer.close(); writer.close().catch((err) => {
logger.error(err);
}); });
});
clientConn.on("error", async (err) => { clientConn.on("error", (clientConnErr) => {
await writer.abort(err); writer.abort(clientConnErr?.message).catch((err) => {
proxyErrorMsg.push((err as Error)?.message);
}); });
} catch (err) { });
clientConn.destroy();
}
})(); })();
// QUIC to Client // QUIC to Client
@@ -238,15 +239,18 @@ const setupProxyServer = async ({
} }
} }
} catch (err) { } catch (err) {
proxyErrorMsg.push((err as Error)?.message);
clientConn.destroy(); clientConn.destroy();
} }
})(); })();
}; };
await setupCopy();
// setupCopy();
// Handle connection closure // Handle connection closure
clientConn.on("close", async () => { clientConn.on("close", () => {
await stream.destroy(); stream.destroy().catch((err) => {
proxyErrorMsg.push((err as Error)?.message);
});
}); });
const cleanup = async () => { const cleanup = async () => {
@@ -254,13 +258,18 @@ const setupProxyServer = async ({
await stream.destroy(); await stream.destroy();
}; };
clientConn.on("error", (err) => { clientConn.on("error", (clientConnErr) => {
logger.error(err, "Client socket error"); logger.error(clientConnErr, "Client socket error");
void cleanup(); cleanup().catch((err) => {
reject(err); logger.error(err, "Client conn cleanup");
});
}); });
clientConn.on("end", cleanup); clientConn.on("end", () => {
cleanup().catch((err) => {
logger.error(err, "Client conn end");
});
});
} catch (err) { } catch (err) {
logger.error(err, "Failed to establish target connection:"); logger.error(err, "Failed to establish target connection:");
clientConn.end(); clientConn.end();
@@ -272,12 +281,12 @@ const setupProxyServer = async ({
reject(err); reject(err);
}); });
server.on("close", async () => { server.on("close", () => {
await quicClient?.destroy(); quicClient?.destroy().catch((err) => {
logger.error(err, "Failed to destroy quic client");
});
}); });
/* eslint-enable */
server.listen(0, () => { server.listen(0, () => {
const address = server.address(); const address = server.address();
if (!address || typeof address === "string") { if (!address || typeof address === "string") {
@@ -293,7 +302,8 @@ const setupProxyServer = async ({
cleanup: async () => { cleanup: async () => {
server.close(); server.close();
await quicClient?.destroy(); await quicClient?.destroy();
} },
getProxyError: () => proxyErrorMsg.join(",")
}); });
}); });
}); });
@@ -316,7 +326,7 @@ export const withGatewayProxy = async (
const { relayHost, relayPort, targetHost, targetPort, tlsOptions, identityId, orgId } = options; const { relayHost, relayPort, targetHost, targetPort, tlsOptions, identityId, orgId } = options;
// Setup the proxy server // Setup the proxy server
const { port, cleanup } = await setupProxyServer({ const { port, cleanup, getProxyError } = await setupProxyServer({
targetHost, targetHost,
targetPort, targetPort,
relayPort, relayPort,
@@ -330,8 +340,12 @@ export const withGatewayProxy = async (
// Execute the callback with the allocated port // Execute the callback with the allocated port
await callback(port); await callback(port);
} catch (err) { } catch (err) {
logger.error(err, "Failed to proxy"); const proxyErrorMessage = getProxyError();
throw new BadRequestError({ message: (err as Error)?.message }); if (proxyErrorMessage) {
logger.error(new Error(proxyErrorMessage), "Failed to proxy");
}
logger.error(err, "Failed to do gateway");
throw new BadRequestError({ message: proxyErrorMessage || (err as Error)?.message });
} finally { } finally {
// Ensure cleanup happens regardless of success or failure // Ensure cleanup happens regardless of success or failure
await cleanup(); await cleanup();
+1 -1
View File
@@ -1,6 +1,6 @@
import crypto from "node:crypto"; import crypto from "node:crypto";
const TURN_TOKEN_TTL = 60 * 60 * 1000; // 24 hours in milliseconds const TURN_TOKEN_TTL = 24 * 60 * 60 * 1000; // 24 hours in milliseconds
export const getTurnCredentials = (id: string, authSecret: string, ttl = TURN_TOKEN_TTL) => { export const getTurnCredentials = (id: string, authSecret: string, ttl = TURN_TOKEN_TTL) => {
const timestamp = Math.floor((Date.now() + ttl) / 1000); const timestamp = Math.floor((Date.now() + ttl) / 1000);
const username = `${timestamp}:${id}`; const username = `${timestamp}:${id}`;