diff --git a/server/index.ts b/server/index.ts index 5c7102847..27c945dc4 100644 --- a/server/index.ts +++ b/server/index.ts @@ -22,14 +22,11 @@ import { } from "@server/db"; import config from "@server/lib/config"; import { setHostMeta } from "@server/lib/hostMeta"; -import { initTelemetryClient } from "@server/lib/telemetry"; import { TraefikConfigManager } from "@server/lib/traefik/TraefikConfigManager"; import { initCleanup } from "#dynamic/cleanup"; +import { startSchedulers } from "#dynamic/startSchedulers"; import license from "#dynamic/license/license"; -import { initLogCleanupInterval } from "@server/lib/cleanupLogs"; -import { initAcmeCertSync } from "@server/lib/acmeCertSync"; import { fetchServerIp } from "@server/lib/serverIpService"; -import { startRebuildQueueProcessor } from "@server/lib/rebuildClientAssociations"; import { initAiModelCatalog } from "@server/lib/aiModelCatalog"; async function startServers() { @@ -44,13 +41,10 @@ async function startServers() { await fetchServerIp(); - initTelemetryClient(); - - initLogCleanupInterval(); - initAcmeCertSync(); - startRebuildQueueProcessor(); await initAiModelCatalog(); + startSchedulers(); + // Start all servers const apiServer = createApiServer(); const internalServer = createInternalServer(); diff --git a/server/private/routers/remoteExitNode/exitNodeReconnectScheduler.ts b/server/private/routers/remoteExitNode/exitNodeReconnectScheduler.ts index 0d871583f..20bb46193 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 @@ -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/ws/messageHandlers.ts b/server/private/routers/ws/messageHandlers.ts index d91726393..b79b715b6 100644 --- a/server/private/routers/ws/messageHandlers.ts +++ b/server/private/routers/ws/messageHandlers.ts @@ -13,22 +13,17 @@ import { handleRemoteExitNodeRegisterMessage, - handleRemoteExitNodePingMessage, - startRemoteExitNodeOfflineChecker, - startExitNodeReconnectScheduler + handleRemoteExitNodePingMessage } from "#private/routers/remoteExitNode"; import { MessageHandler } from "@server/routers/ws"; -import { build } from "@server/build"; -import { handleConnectionLogMessage, handleRequestLogMessage } from "#private/routers/newt"; +import { + handleConnectionLogMessage, + handleRequestLogMessage +} from "#private/routers/newt"; export const messageHandlers: Record = { "remoteExitNode/register": handleRemoteExitNodeRegisterMessage, "remoteExitNode/ping": handleRemoteExitNodePingMessage, "newt/access-log": handleConnectionLogMessage, - "newt/request-log": handleRequestLogMessage, + "newt/request-log": handleRequestLogMessage }; - -if (build != "saas") { - startRemoteExitNodeOfflineChecker(); // this is to handle the offline check for remote exit nodes - startExitNodeReconnectScheduler(); // check pending exit node reconnects and notify newts -} diff --git a/server/private/startSchedulers.ts b/server/private/startSchedulers.ts new file mode 100644 index 000000000..ee80d5efd --- /dev/null +++ b/server/private/startSchedulers.ts @@ -0,0 +1,12 @@ +import { build } from "@server/build"; +import { startRemoteExitNodeOfflineChecker } from "./routers/remoteExitNode"; +import { startExitNodeReconnectScheduler } from "./routers/remoteExitNode/exitNodeReconnectScheduler"; +import { startSchedulers as ossStartSchedulers } from "@server/startSchedulers"; + +export function startSchedulers() { + if (build != "saas") { + startRemoteExitNodeOfflineChecker(); // this is to handle the offline check for remote exit nodes + startExitNodeReconnectScheduler(); // check pending exit node reconnects and notify newts + } + ossStartSchedulers(); +} diff --git a/server/routers/ws/messageHandlers.ts b/server/routers/ws/messageHandlers.ts index fff2bf7c4..b8ac9baa4 100644 --- a/server/routers/ws/messageHandlers.ts +++ b/server/routers/ws/messageHandlers.ts @@ -1,4 +1,3 @@ -import { build } from "@server/build"; import { handleNewtRegisterMessage, handleReceiveBandwidthMessage, @@ -8,15 +7,12 @@ import { handleNewtExitNodesRequestMessage, handleApplyBlueprintMessage, handleNewtPingMessage, - startNewtOfflineChecker, handleNewtDisconnectingMessage } from "../newt"; -import { startPingAccumulator } from "../newt/pingAccumulator"; import { handleOlmRegisterMessage, handleOlmRelayMessage, handleOlmPingMessage, - startOlmOfflineChecker, handleOlmServerPeerAddMessage, handleOlmUnRelayMessage, handleOlmDisconnectingMessage, @@ -52,12 +48,3 @@ export const messageHandlers: Record = { "newt/healthcheck/status": handleHealthcheckStatusMessage, "ws/round-trip/complete": handleRoundTripMessage }; - -// Start the ping accumulator for all builds - it batches per-site online/lastPing -// updates into periodic bulk writes, preventing connection pool exhaustion. -startPingAccumulator(); - -if (build != "saas") { - startOlmOfflineChecker(); // this is to handle the offline check for olms - startNewtOfflineChecker(); // this is to handle the offline check for newts -} diff --git a/server/startSchedulers.ts b/server/startSchedulers.ts new file mode 100644 index 000000000..4a67495cc --- /dev/null +++ b/server/startSchedulers.ts @@ -0,0 +1,25 @@ +import { build } from "@server/build"; +import { startPingAccumulator } from "./routers/newt/pingAccumulator"; +import { startOlmOfflineChecker } from "./routers/olm"; +import { startNewtOfflineChecker } from "./routers/newt"; +import { initTelemetryClient } from "@server/lib/telemetry"; +import { initLogCleanupInterval } from "@server/lib/cleanupLogs"; +import { initAcmeCertSync } from "@server/lib/acmeCertSync"; +import { startRebuildQueueProcessor } from "@server/lib/rebuildClientAssociations"; + +export function startSchedulers() { + // Start the ping accumulator for all builds - it batches per-site online/lastPing + // updates into periodic bulk writes, preventing connection pool exhaustion. + startPingAccumulator(); + + if (build != "saas") { + startOlmOfflineChecker(); // this is to handle the offline check for olms + startNewtOfflineChecker(); // this is to handle the offline check for newts + } + + initTelemetryClient(); + + initLogCleanupInterval(); + initAcmeCertSync(); + startRebuildQueueProcessor(); +}