Use endpoint instead of reachableAt for remote nodes

This commit is contained in:
Owen
2026-09-07 11:55:45 -04:00
parent ea9017ac06
commit 080bcbaf97
2 changed files with 55 additions and 49 deletions
@@ -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<void> {
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<void> {
`${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<void> {
}
// 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<void> {
`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<void> {
// 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<void> {
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<void> {
pendingReconnects.delete(exitNodeId);
@@ -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) {