Fix circular import issue

This commit is contained in:
Owen
2026-08-17 11:30:27 -04:00
parent 7d6f3a2925
commit 749abdcb56
6 changed files with 80 additions and 67 deletions
+3 -9
View File
@@ -22,14 +22,11 @@ import {
} from "@server/db"; } from "@server/db";
import config from "@server/lib/config"; import config from "@server/lib/config";
import { setHostMeta } from "@server/lib/hostMeta"; import { setHostMeta } from "@server/lib/hostMeta";
import { initTelemetryClient } from "@server/lib/telemetry";
import { TraefikConfigManager } from "@server/lib/traefik/TraefikConfigManager"; import { TraefikConfigManager } from "@server/lib/traefik/TraefikConfigManager";
import { initCleanup } from "#dynamic/cleanup"; import { initCleanup } from "#dynamic/cleanup";
import { startSchedulers } from "#dynamic/startSchedulers";
import license from "#dynamic/license/license"; 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 { fetchServerIp } from "@server/lib/serverIpService";
import { startRebuildQueueProcessor } from "@server/lib/rebuildClientAssociations";
import { initAiModelCatalog } from "@server/lib/aiModelCatalog"; import { initAiModelCatalog } from "@server/lib/aiModelCatalog";
async function startServers() { async function startServers() {
@@ -44,13 +41,10 @@ async function startServers() {
await fetchServerIp(); await fetchServerIp();
initTelemetryClient();
initLogCleanupInterval();
initAcmeCertSync();
startRebuildQueueProcessor();
await initAiModelCatalog(); await initAiModelCatalog();
startSchedulers();
// Start all servers // Start all servers
const apiServer = createApiServer(); const apiServer = createApiServer();
const internalServer = createInternalServer(); const internalServer = createInternalServer();
@@ -16,7 +16,7 @@ import { db, exitNodes, newts, sites } from "@server/db";
import { eq } from "drizzle-orm"; import { eq } from "drizzle-orm";
import logger from "@server/logger"; import logger from "@server/logger";
import redisManager from "#private/lib/redis"; 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 INITIAL_DELAY_MS = 15 * 1000; // 15 seconds before first check
const CHECK_INTERVAL_MS = 10 * 1000; // Check every 10 seconds const CHECK_INTERVAL_MS = 10 * 1000; // Check every 10 seconds
@@ -150,47 +150,47 @@ async function processPendingReconnects(): Promise<void> {
`Exit node ${exitNodeId} is reachable. Sending newt/wg/reconnect to connected newts.` `Exit node ${exitNodeId} is reachable. Sending newt/wg/reconnect to connected newts.`
); );
await sendReconnectToNewts(exitNodeId); // await sendReconnectToNewts(exitNodeId);
await removePending(exitNodeId); await removePending(exitNodeId);
} }
} }
async function sendReconnectToNewts(exitNodeId: number): Promise<void> { // async function sendReconnectToNewts(exitNodeId: number): Promise<void> {
try { // try {
const connectedNewts = await db // const connectedNewts = await db
.select({ newtId: newts.newtId }) // .select({ newtId: newts.newtId })
.from(newts) // .from(newts)
.innerJoin(sites, eq(newts.siteId, sites.siteId)) // .innerJoin(sites, eq(newts.siteId, sites.siteId))
.where(eq(sites.exitNodeId, exitNodeId)); // .where(eq(sites.exitNodeId, exitNodeId));
if (connectedNewts.length === 0) { // if (connectedNewts.length === 0) {
logger.debug( // logger.debug(
`No newts found for exit node ${exitNodeId}, nothing to reconnect` // `No newts found for exit node ${exitNodeId}, nothing to reconnect`
); // );
return; // return;
} // }
logger.info( // logger.info(
`Sending newt/wg/reconnect to ${connectedNewts.length} newt(s) for exit node ${exitNodeId}` // `Sending newt/wg/reconnect to ${connectedNewts.length} newt(s) for exit node ${exitNodeId}`
); // );
const reconnectMessage = { // const reconnectMessage = {
type: "newt/wg/reconnect", // type: "newt/wg/reconnect",
data: {} // data: {}
}; // };
await Promise.allSettled( // await Promise.allSettled(
connectedNewts.map(({ newtId }) => // connectedNewts.map(({ newtId }) =>
sendToClient(newtId, reconnectMessage) // sendToClient(newtId, reconnectMessage)
) // )
); // );
} catch (error) { // } catch (error) {
logger.error( // logger.error(
`Failed to send reconnect messages for exit node ${exitNodeId}`, // `Failed to send reconnect messages for exit node ${exitNodeId}`,
{ error } // { error }
); // );
} // }
} // }
async function removePending(exitNodeId: number): Promise<void> { async function removePending(exitNodeId: number): Promise<void> {
pendingReconnects.delete(exitNodeId); pendingReconnects.delete(exitNodeId);
+6 -11
View File
@@ -13,22 +13,17 @@
import { import {
handleRemoteExitNodeRegisterMessage, handleRemoteExitNodeRegisterMessage,
handleRemoteExitNodePingMessage, handleRemoteExitNodePingMessage
startRemoteExitNodeOfflineChecker,
startExitNodeReconnectScheduler
} from "#private/routers/remoteExitNode"; } from "#private/routers/remoteExitNode";
import { MessageHandler } from "@server/routers/ws"; import { MessageHandler } from "@server/routers/ws";
import { build } from "@server/build"; import {
import { handleConnectionLogMessage, handleRequestLogMessage } from "#private/routers/newt"; handleConnectionLogMessage,
handleRequestLogMessage
} from "#private/routers/newt";
export const messageHandlers: Record<string, MessageHandler> = { export const messageHandlers: Record<string, MessageHandler> = {
"remoteExitNode/register": handleRemoteExitNodeRegisterMessage, "remoteExitNode/register": handleRemoteExitNodeRegisterMessage,
"remoteExitNode/ping": handleRemoteExitNodePingMessage, "remoteExitNode/ping": handleRemoteExitNodePingMessage,
"newt/access-log": handleConnectionLogMessage, "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
}
+12
View File
@@ -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();
}
-13
View File
@@ -1,4 +1,3 @@
import { build } from "@server/build";
import { import {
handleNewtRegisterMessage, handleNewtRegisterMessage,
handleReceiveBandwidthMessage, handleReceiveBandwidthMessage,
@@ -8,15 +7,12 @@ import {
handleNewtExitNodesRequestMessage, handleNewtExitNodesRequestMessage,
handleApplyBlueprintMessage, handleApplyBlueprintMessage,
handleNewtPingMessage, handleNewtPingMessage,
startNewtOfflineChecker,
handleNewtDisconnectingMessage handleNewtDisconnectingMessage
} from "../newt"; } from "../newt";
import { startPingAccumulator } from "../newt/pingAccumulator";
import { import {
handleOlmRegisterMessage, handleOlmRegisterMessage,
handleOlmRelayMessage, handleOlmRelayMessage,
handleOlmPingMessage, handleOlmPingMessage,
startOlmOfflineChecker,
handleOlmServerPeerAddMessage, handleOlmServerPeerAddMessage,
handleOlmUnRelayMessage, handleOlmUnRelayMessage,
handleOlmDisconnectingMessage, handleOlmDisconnectingMessage,
@@ -52,12 +48,3 @@ export const messageHandlers: Record<string, MessageHandler> = {
"newt/healthcheck/status": handleHealthcheckStatusMessage, "newt/healthcheck/status": handleHealthcheckStatusMessage,
"ws/round-trip/complete": handleRoundTripMessage "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
}
+25
View File
@@ -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();
}