diff --git a/server/private/routers/remoteExitNode/exitNodeReconnectScheduler.ts b/server/private/routers/remoteExitNode/exitNodeReconnectScheduler.ts index 20bb46193..8cb3e80ec 100644 --- a/server/private/routers/remoteExitNode/exitNodeReconnectScheduler.ts +++ b/server/private/routers/remoteExitNode/exitNodeReconnectScheduler.ts @@ -16,7 +16,7 @@ import { db, exitNodes, newts, sites } from "@server/db"; import { eq } from "drizzle-orm"; import logger from "@server/logger"; import redisManager from "#private/lib/redis"; -// import { sendToClient } from "#private/routers/ws"; +import { sendToClient } from "#private/routers/ws"; const INITIAL_DELAY_MS = 15 * 1000; // 15 seconds before first check const CHECK_INTERVAL_MS = 10 * 1000; // Check every 10 seconds @@ -26,7 +26,7 @@ const REDIS_HASH_PREFIX = "exit-node-reconnect:"; interface PendingReconnect { startTime: number; - reachableAt: string; + endpoint: string; } // In-memory tracking for this node @@ -40,15 +40,15 @@ let schedulerInterval: NodeJS.Timeout | null = null; */ export async function scheduleExitNodeReconnect( exitNodeId: number, - reachableAt: string + endpoint: string ): Promise { logger.info( - `Scheduling newt reconnect for exit node ${exitNodeId} (reachableAt: ${reachableAt})` + `Scheduling newt reconnect for exit node ${exitNodeId} (endpoint: ${endpoint})` ); const entry: PendingReconnect = { startTime: Date.now(), - reachableAt + endpoint }; pendingReconnects.set(exitNodeId, entry); @@ -63,8 +63,8 @@ export async function scheduleExitNodeReconnect( ); await redisManager.hset( `${REDIS_HASH_PREFIX}${exitNodeId}`, - "reachableAt", - reachableAt + "endpoint", + endpoint ); } } @@ -101,14 +101,14 @@ async function processPendingReconnects(): Promise { `${REDIS_HASH_PREFIX}${id}`, "startTime" ); - const reachableAt = await redisManager.hget( + const endpoint = await redisManager.hget( `${REDIS_HASH_PREFIX}${id}`, - "reachableAt" + "endpoint" ); - if (startTimeStr && reachableAt) { + if (startTimeStr && endpoint) { toProcess.set(id, { startTime: parseInt(startTimeStr, 10), - reachableAt + endpoint }); } } @@ -135,7 +135,7 @@ async function processPendingReconnects(): Promise { } // Check if the exit node HTTP endpoint is reachable - const pingUrl = `${entry.reachableAt}/ping`; + const pingUrl = `http://${entry.endpoint}/ping`; try { await axios.get(pingUrl, { timeout: 5000 }); } catch { @@ -150,47 +150,47 @@ async function processPendingReconnects(): Promise { `Exit node ${exitNodeId} is reachable. Sending newt/wg/reconnect to connected newts.` ); - // await sendReconnectToNewts(exitNodeId); + await sendReconnectToNewts(exitNodeId); await removePending(exitNodeId); } } -// async function sendReconnectToNewts(exitNodeId: number): Promise { -// try { -// const connectedNewts = await db -// .select({ newtId: newts.newtId }) -// .from(newts) -// .innerJoin(sites, eq(newts.siteId, sites.siteId)) -// .where(eq(sites.exitNodeId, exitNodeId)); +async function sendReconnectToNewts(exitNodeId: number): Promise { + try { + const connectedNewts = await db + .select({ newtId: newts.newtId }) + .from(newts) + .innerJoin(sites, eq(newts.siteId, sites.siteId)) + .where(eq(sites.exitNodeId, exitNodeId)); -// if (connectedNewts.length === 0) { -// logger.debug( -// `No newts found for exit node ${exitNodeId}, nothing to reconnect` -// ); -// return; -// } + if (connectedNewts.length === 0) { + logger.debug( + `No newts found for exit node ${exitNodeId}, nothing to reconnect` + ); + return; + } -// logger.info( -// `Sending newt/wg/reconnect to ${connectedNewts.length} newt(s) for exit node ${exitNodeId}` -// ); + logger.info( + `Sending newt/wg/reconnect to ${connectedNewts.length} newt(s) for exit node ${exitNodeId}` + ); -// const reconnectMessage = { -// type: "newt/wg/reconnect", -// data: {} -// }; + const reconnectMessage = { + type: "newt/wg/reconnect", + data: {} + }; -// await Promise.allSettled( -// connectedNewts.map(({ newtId }) => -// sendToClient(newtId, reconnectMessage) -// ) -// ); -// } catch (error) { -// logger.error( -// `Failed to send reconnect messages for exit node ${exitNodeId}`, -// { error } -// ); -// } -// } + await Promise.allSettled( + connectedNewts.map(({ newtId }) => + sendToClient(newtId, reconnectMessage) + ) + ); + } catch (error) { + logger.error( + `Failed to send reconnect messages for exit node ${exitNodeId}`, + { error } + ); + } +} async function removePending(exitNodeId: number): Promise { pendingReconnects.delete(exitNodeId); diff --git a/server/private/routers/remoteExitNode/handleRemoteExitNodePingMessage.ts b/server/private/routers/remoteExitNode/handleRemoteExitNodePingMessage.ts index 10bf36d7c..14a777684 100644 --- a/server/private/routers/remoteExitNode/handleRemoteExitNodePingMessage.ts +++ b/server/private/routers/remoteExitNode/handleRemoteExitNodePingMessage.ts @@ -40,7 +40,7 @@ export const handleRemoteExitNodePingMessage: MessageHandler = async ( try { // Fetch the current state before updating so we can detect the offline→online transition const [currentExitNode] = await db - .select({ online: exitNodes.online, reachableAt: exitNodes.reachableAt }) + .select({ online: exitNodes.online, endpoint: exitNodes.endpoint }) .from(exitNodes) .where(eq(exitNodes.exitNodeId, remoteExitNode.exitNodeId)) .limit(1); @@ -55,12 +55,18 @@ export const handleRemoteExitNodePingMessage: MessageHandler = async ( .where(eq(exitNodes.exitNodeId, remoteExitNode.exitNodeId)); // If the exit node was offline and is now coming online, schedule newt reconnects - if (currentExitNode && !currentExitNode.online && currentExitNode.reachableAt) { + if ( + currentExitNode && + !currentExitNode.online && + currentExitNode.endpoint + ) { scheduleExitNodeReconnect( remoteExitNode.exitNodeId, - currentExitNode.reachableAt + currentExitNode.endpoint ).catch((error) => { - logger.error("Failed to schedule exit node reconnect", { error }); + logger.error("Failed to schedule exit node reconnect", { + error + }); }); } } catch (error) {