Spaces:
Runtime error
Runtime error
File size: 5,474 Bytes
cd8bd0a | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 | import http from "node:http";
import net from "node:net";
import { randomUUID } from "node:crypto";
import { createResponsesWsProxy } from "./responses-ws-proxy.mjs";
import { ensurePeerStampToken, wrapRequestListenerWithPeerStamp } from "./peer-stamp.mjs";
import { maybeHandleWebdav } from "./webdav-handler.mjs";
import methodGuard from "./http-method-guard.cjs";
const originalCreateServer = http.createServer.bind(http);
const proxiesByPort = new Map();
const { wrapRequestListenerWithMethodGuard } = methodGuard;
process.env.OMNIROUTE_WS_BRIDGE_SECRET ||= randomUUID();
// Per-process secret proving the trusted peer-IP stamp came from this server.
ensurePeerStampToken();
function getPort(server) {
const address = server.address?.();
if (address && typeof address === "object" && typeof address.port === "number") {
return address.port;
}
const rawPort = process.env.PORT || process.env.DASHBOARD_PORT || "3000";
const parsed = Number.parseInt(rawPort, 10);
return Number.isFinite(parsed) && parsed > 0 ? parsed : 3000;
}
function getProxy(server) {
const port = getPort(server);
const existing = proxiesByPort.get(port);
if (existing) return existing;
const proxy = createResponsesWsProxy({
baseUrl: `http://127.0.0.1:${port}`,
bridgeSecret: process.env.OMNIROUTE_WS_BRIDGE_SECRET,
});
proxiesByPort.set(port, proxy);
return proxy;
}
function proxyLiveWs(req, socket, head) {
const targetPort = parseInt(process.env.LIVE_WS_PORT || "20129", 10);
const targetSocket = net.connect(targetPort, "127.0.0.1", () => {
let rawRequest = `${req.method} ${req.url} HTTP/${req.httpVersion}\r\n`;
for (const [key, val] of Object.entries(req.headers)) {
if (Array.isArray(val)) {
for (const v of val) rawRequest += `${key}: ${v}\r\n`;
} else {
rawRequest += `${key}: ${val}\r\n`;
}
}
rawRequest += "\r\n";
targetSocket.write(rawRequest);
if (head && head.length > 0) targetSocket.write(head);
targetSocket.pipe(socket);
socket.pipe(targetSocket);
});
targetSocket.on("error", () => !socket.destroyed && socket.destroy());
socket.on("error", () => !targetSocket.destroyed && targetSocket.destroy());
}
function wrapUpgradeListener(server, listener) {
return async function responsesWsAwareUpgrade(req, socket, head) {
try {
const url = new URL(req.url || "/", `http://${req.headers.host || "localhost"}`);
if (url.pathname === "/live-ws" || url.pathname.startsWith("/live-ws")) {
proxyLiveWs(req, socket, head);
return;
}
const handled = await getProxy(server).handleUpgrade(req, socket, head);
if (handled) return;
return listener.call(this, req, socket, head);
} catch (error) {
if (!socket.destroyed) {
socket.destroy(error instanceof Error ? error : undefined);
}
console.error("[Responses WS] Upgrade handling failed:", error);
}
};
}
/**
* Wrap a request listener so WebDAV requests at /api/v1/webdav are handled
* before the peer-stamp/Next.js layer sees them.
* Returns true if the request was handled; the wrapped listener is never called.
*/
function wrapRequestListenerWithWebdav(listener) {
return async function webdavAwareRequestHandler(req, res) {
try {
const handled = await maybeHandleWebdav(req, res);
if (handled) return;
} catch {
// Never block a request on WebDAV errors — fall through to Next
}
return listener.call(this, req, res);
};
}
http.createServer = function createServerWithResponsesWs(...args) {
// Next's standalone server.js may pass its request listener directly to
// createServer; wrap it so the real TCP peer IP is stamped before Next runs.
const lastFnIdx = args.map((a) => typeof a === "function").lastIndexOf(true);
if (lastFnIdx >= 0) {
// Method guard runs before Next because Next 16 rejects TRACE while constructing requests.
args[lastFnIdx] = wrapRequestListenerWithMethodGuard(
wrapRequestListenerWithWebdav(wrapRequestListenerWithPeerStamp(args[lastFnIdx]))
);
}
const server = originalCreateServer(...args);
const originalOn = server.on.bind(server);
const originalAddListener = server.addListener.bind(server);
server.on = function patchedOn(eventName, listener) {
if (eventName === "upgrade" && typeof listener === "function") {
return originalOn(eventName, wrapUpgradeListener(server, listener));
}
// …or it may attach the handler via server.on("request"): wrap that too.
if (eventName === "request" && typeof listener === "function") {
return originalOn(
eventName,
wrapRequestListenerWithMethodGuard(
wrapRequestListenerWithWebdav(wrapRequestListenerWithPeerStamp(listener))
)
);
}
return originalOn(eventName, listener);
};
server.addListener = function patchedAddListener(eventName, listener) {
if (eventName === "upgrade" && typeof listener === "function") {
return originalAddListener(eventName, wrapUpgradeListener(server, listener));
}
if (eventName === "request" && typeof listener === "function") {
return originalAddListener(
eventName,
wrapRequestListenerWithMethodGuard(
wrapRequestListenerWithWebdav(wrapRequestListenerWithPeerStamp(listener))
)
);
}
return originalAddListener(eventName, listener);
};
return server;
};
await import("./server.js");
|