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");