Compare commits

..

3 Commits

Author SHA1 Message Date
Owen Schwartz 17375348b0 Merge pull request #3702 from fosrl/dev
1.22.2
2026-09-04 17:18:42 -04:00
Owen Schwartz e7fdbf9e85 Merge pull request #3699 from fosrl/dev
1.22.1-s.1
2026-09-04 10:38:39 -04:00
Milo Schwartz cddb5ecc3d Merge pull request #3691 from fosrl/dev
update readme
2026-09-03 16:19:47 -04:00
11 changed files with 71 additions and 159 deletions
-21
View File
@@ -227,27 +227,6 @@ export const resources = pgTable(
]
);
export const redirects = pgTable("redirects", {
redirectId: serial("redirectId").primaryKey(),
orgId: varchar("orgId")
.references(() => orgs.orgId, {
onDelete: "cascade"
})
.notNull(),
resourceId: integer("resourceId").references(() => resources.resourceId, {
onDelete: "cascade"
}),
domainId: integer("domainId").references(() => domains.domainId, {
onDelete: "cascade"
}),
niceId: text("niceId").notNull(),
name: varchar("name").notNull(),
sourcePath: varchar("sourcePath").notNull(),
destinationUrl: varchar("destinationUrl"),
permanent: boolean("permanent").notNull().default(false),
enabled: boolean("enabled").notNull().default(true)
});
export const resourceAiProviders = pgTable(
"resourceAiProviders",
{
-23
View File
@@ -243,29 +243,6 @@ export const resources = sqliteTable(
(table) => [index("idx_resources_orgId").on(table.orgId)]
);
export const redirects = sqliteTable("redirects", {
redirectId: integer("redirectId").primaryKey({ autoIncrement: true }),
orgId: text("orgId")
.references(() => orgs.orgId, {
onDelete: "cascade"
})
.notNull(),
resourceId: integer("resourceId").references(() => resources.resourceId, {
onDelete: "cascade"
}),
domainId: integer("domainId").references(() => domains.domainId, {
onDelete: "cascade"
}),
niceId: text("niceId").notNull(),
name: text("name").notNull(),
sourcePath: text("sourcePath").notNull(),
destinationUrl: text("destinationUrl"),
permanent: integer("permanent", { mode: "boolean" })
.notNull()
.default(false),
enabled: integer("enabled", { mode: "boolean" }).notNull().default(true)
});
export const resourceAiProviders = sqliteTable(
"resourceAiProviders",
{
+1 -1
View File
@@ -71,7 +71,7 @@ export async function withRetry<T>(
const jitter = Math.random() * baseDelay;
const delay = baseDelay + jitter;
logger.warn(
`Transient DB issue in ${context}, retrying attempt ${attempt}/${maxRetries} after ${delay.toFixed(0)}ms`,
`Transient DB error in ${context}, retrying attempt ${attempt}/${maxRetries} after ${delay.toFixed(0)}ms`,
{ code: error?.code ?? error?.cause?.code }
);
await new Promise((resolve) => setTimeout(resolve, delay));
+2 -2
View File
@@ -348,8 +348,8 @@ export const configSchema = z
.optional()
.pipe(z.string())
.transform((url) => url.toLowerCase()),
subnet_group: z.string().optional().default("100.89.137.0/18"),
block_size: z.number().positive().gt(0).optional().default(22),
subnet_group: z.string().optional().default("100.89.137.0/20"),
block_size: z.number().positive().gt(0).optional().default(24),
site_block_size: z
.number()
.positive()
+1 -6
View File
@@ -2478,12 +2478,7 @@ hybridRouter.post(
destinations: destinations
});
} catch (error) {
if (!(
error instanceof Error &&
error.message === "Exit node not allowed"
)) {
logger.error(error);
}
logger.error(error);
return next(
createHttpError(
HttpCode.INTERNAL_SERVER_ERROR,
@@ -1,23 +0,0 @@
/*
* This file is part of a proprietary work.
*
* Copyright (c) 2025-2026 Fossorial, Inc.
* All rights reserved.
*
* This file is licensed under the Fossorial Commercial License.
* You may not use this file except in compliance with the License.
* Unauthorized use, copying, modification, or distribution is strictly prohibited.
*
* This file is not licensed under the AGPLv3.
*/
import { EventEmitter } from "events";
export interface ExitNodeOnlineEvent {
exitNodeId: number;
endpoint: string;
}
export const EXIT_NODE_ONLINE_EVENT = "exit-node-online";
export const exitNodeEvents = new EventEmitter();
@@ -12,27 +12,11 @@
*/
import axios from "axios";
import { db, newts, sites } from "@server/db";
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 "../ws";
import {
exitNodeEvents,
EXIT_NODE_ONLINE_EVENT,
ExitNodeOnlineEvent
} from "./exitNodeEvents";
exitNodeEvents.on(
EXIT_NODE_ONLINE_EVENT,
({ exitNodeId, endpoint }: ExitNodeOnlineEvent) => {
scheduleExitNodeReconnect(exitNodeId, endpoint).catch((error) => {
logger.error("Failed to schedule exit node reconnect", {
error
});
});
}
);
// 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
@@ -42,7 +26,7 @@ const REDIS_HASH_PREFIX = "exit-node-reconnect:";
interface PendingReconnect {
startTime: number;
endpoint: string;
reachableAt: string;
}
// In-memory tracking for this node
@@ -56,15 +40,15 @@ let schedulerInterval: NodeJS.Timeout | null = null;
*/
export async function scheduleExitNodeReconnect(
exitNodeId: number,
endpoint: string
reachableAt: string
): Promise<void> {
logger.info(
`Scheduling newt reconnect for exit node ${exitNodeId} (endpoint: ${endpoint})`
`Scheduling newt reconnect for exit node ${exitNodeId} (reachableAt: ${reachableAt})`
);
const entry: PendingReconnect = {
startTime: Date.now(),
endpoint
reachableAt
};
pendingReconnects.set(exitNodeId, entry);
@@ -79,8 +63,8 @@ export async function scheduleExitNodeReconnect(
);
await redisManager.hset(
`${REDIS_HASH_PREFIX}${exitNodeId}`,
"endpoint",
endpoint
"reachableAt",
reachableAt
);
}
}
@@ -117,14 +101,14 @@ async function processPendingReconnects(): Promise<void> {
`${REDIS_HASH_PREFIX}${id}`,
"startTime"
);
const endpoint = await redisManager.hget(
const reachableAt = await redisManager.hget(
`${REDIS_HASH_PREFIX}${id}`,
"endpoint"
"reachableAt"
);
if (startTimeStr && endpoint) {
if (startTimeStr && reachableAt) {
toProcess.set(id, {
startTime: parseInt(startTimeStr, 10),
endpoint
reachableAt
});
}
}
@@ -151,7 +135,7 @@ async function processPendingReconnects(): Promise<void> {
}
// Check if the exit node HTTP endpoint is reachable
const pingUrl = `http://${entry.endpoint}/ping`;
const pingUrl = `${entry.reachableAt}/ping`;
try {
await axios.get(pingUrl, { timeout: 5000 });
} catch {
@@ -166,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);
@@ -16,7 +16,7 @@ import { MessageHandler } from "@server/routers/ws";
import { RemoteExitNode } from "@server/db";
import { eq } from "drizzle-orm";
import logger from "@server/logger";
import { exitNodeEvents, EXIT_NODE_ONLINE_EVENT } from "./exitNodeEvents";
import { scheduleExitNodeReconnect } from "./exitNodeReconnectScheduler";
/**
* Handles ping messages from clients and responds with pong
@@ -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, endpoint: exitNodes.endpoint })
.select({ online: exitNodes.online, reachableAt: exitNodes.reachableAt })
.from(exitNodes)
.where(eq(exitNodes.exitNodeId, remoteExitNode.exitNodeId))
.limit(1);
@@ -55,14 +55,12 @@ 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.endpoint
) {
exitNodeEvents.emit(EXIT_NODE_ONLINE_EVENT, {
exitNodeId: remoteExitNode.exitNodeId,
endpoint: currentExitNode.endpoint
if (currentExitNode && !currentExitNode.online && currentExitNode.reachableAt) {
scheduleExitNodeReconnect(
remoteExitNode.exitNodeId,
currentExitNode.reachableAt
).catch((error) => {
logger.error("Failed to schedule exit node reconnect", { error });
});
}
} catch (error) {
@@ -85,7 +85,7 @@ export const handleOlmServerInitAddPeerHandshake: MessageHandler = async (
);
if (!resources || resources.length === 0) {
logger.warn(
logger.error(
`handleOlmServerInitAddPeerHandshake: Resource not found`
);
await sendCancel();
@@ -94,7 +94,7 @@ export const handleOlmServerInitAddPeerHandshake: MessageHandler = async (
if (resources.length > 1) {
// error but this should not happen because the nice id cant contain a dot and the alias has to have a dot and both have to be unique within the org so there should never be multiple matches
logger.warn(
logger.error(
`handleOlmServerInitAddPeerHandshake: Multiple resources found matching the criteria`
);
return;
@@ -119,7 +119,7 @@ export const handleOlmServerInitAddPeerHandshake: MessageHandler = async (
);
if (currentResourceAssociationCaches.length === 0) {
logger.warn(
logger.error(
`handleOlmServerInitAddPeerHandshake: Client ${client.clientId} does not have access to resource ${resource.siteResourceId}`
);
await sendCancel();
@@ -127,7 +127,7 @@ export const handleOlmServerInitAddPeerHandshake: MessageHandler = async (
}
if (!resource.networkId) {
logger.warn(
logger.error(
`handleOlmServerInitAddPeerHandshake: Resource ${resource.siteResourceId} has no network`
);
await sendCancel();
@@ -141,7 +141,7 @@ export const handleOlmServerInitAddPeerHandshake: MessageHandler = async (
.where(eq(siteNetworks.networkId, resource.networkId));
if (!siteRows || siteRows.length === 0) {
logger.warn(
logger.error(
`handleOlmServerInitAddPeerHandshake: No sites found for resource ${resource.siteResourceId}`
);
await sendCancel();
@@ -164,7 +164,9 @@ export const handleOlmServerInitAddPeerHandshake: MessageHandler = async (
}
if (sitesToProcess.length === 0) {
logger.warn(`handleOlmServerInitAddPeerHandshake: No sites to process`);
logger.error(
`handleOlmServerInitAddPeerHandshake: No sites to process`
);
await sendCancel();
return;
}
@@ -191,7 +193,7 @@ export const handleOlmServerInitAddPeerHandshake: MessageHandler = async (
}
if (!site.exitNodeId) {
logger.warn(
logger.error(
`handleOlmServerInitAddPeerHandshake: Site ${site.siteId} has no exit node, skipping`
);
continue;
@@ -203,7 +205,7 @@ export const handleOlmServerInitAddPeerHandshake: MessageHandler = async (
.where(eq(exitNodes.exitNodeId, site.exitNodeId));
if (!exitNode) {
logger.warn(
logger.error(
`handleOlmServerInitAddPeerHandshake: Exit node not found for site ${site.siteId}, skipping`
);
continue;
@@ -227,7 +229,7 @@ export const handleOlmServerInitAddPeerHandshake: MessageHandler = async (
}
if (!handshakeInitiated) {
logger.warn(
logger.error(
`handleOlmServerInitAddPeerHandshake: No accessible sites with valid exit nodes found, cancelling chain`
);
await sendCancel();
+1 -1
View File
@@ -104,7 +104,7 @@ export default async function OrgLayout(props: {
subscriptionStatus = subRes.data.data;
} catch (error) {
// If subscription fetch fails, keep subscriptionStatus as null
// console.error("Failed to fetch subscription status:", error);
console.error("Failed to fetch subscription status:", error);
}
}
+1 -1
View File
@@ -14,7 +14,7 @@
"moduleResolution": "bundler",
"resolveJsonModule": true,
"isolatedModules": true,
"jsx": "react-jsx",
"jsx": "preserve",
"incremental": true,
"paths": {
"@server/*": [