incognitolm commited on
Commit ·
9005248
1
Parent(s): ff954f1
Update chatStream.js
Browse files- server/chatStream.js +32 -26
server/chatStream.js
CHANGED
|
@@ -4,6 +4,7 @@ import { fileURLToPath } from "url";
|
|
| 4 |
import path from "path";
|
| 5 |
import { LIGHTNING_BASE } from "./config.js";
|
| 6 |
import WebSocket from "ws";
|
|
|
|
| 7 |
|
| 8 |
const __dirname = path.dirname(fileURLToPath(import.meta.url));
|
| 9 |
const WORKER_PATH = path.join(__dirname, "searchWorker.js");
|
|
@@ -53,61 +54,64 @@ function makeClient(accessToken, clientId) {
|
|
| 53 |
});
|
| 54 |
}
|
| 55 |
|
| 56 |
-
export async function
|
| 57 |
const wsURL =
|
| 58 |
(process.env.LIGHTNING_BASE.startsWith("https")
|
| 59 |
? process.env.LIGHTNING_BASE.replace("https", "wss")
|
| 60 |
: process.env.LIGHTNING_BASE.replace("http", "ws")) + "/ws/chat";
|
| 61 |
|
|
|
|
| 62 |
const ws = new WebSocket(wsURL);
|
| 63 |
|
| 64 |
-
// Utility to safely parse JSON
|
| 65 |
const safeParse = (str) => {
|
| 66 |
-
try { return JSON.parse(str); }
|
| 67 |
-
catch { return null; }
|
| 68 |
};
|
| 69 |
|
| 70 |
-
// Wait for open with timeout
|
| 71 |
await new Promise((resolve, reject) => {
|
| 72 |
const timer = setTimeout(() => reject(new Error("WS connection timeout")), 5000);
|
| 73 |
ws.on("open", () => {
|
|
|
|
| 74 |
clearTimeout(timer);
|
| 75 |
-
console.log("[WS] Opened connection");
|
| 76 |
resolve();
|
| 77 |
});
|
| 78 |
ws.on("error", (err) => {
|
|
|
|
| 79 |
clearTimeout(timer);
|
| 80 |
-
console.error("[WS] Connection error:", err);
|
| 81 |
reject(err);
|
| 82 |
});
|
| 83 |
});
|
| 84 |
|
| 85 |
-
// Authenticate
|
|
|
|
| 86 |
ws.send(JSON.stringify({ key: process.env.WEBSOCKET_KEY }));
|
|
|
|
| 87 |
await new Promise((resolve, reject) => {
|
| 88 |
const timer = setTimeout(() => reject(new Error("WS auth timeout")), 5000);
|
| 89 |
ws.on("message", (data) => {
|
| 90 |
-
const msg = safeParse(data.toString());
|
| 91 |
console.log("[WS RAW MESSAGE]", data.toString());
|
|
|
|
| 92 |
if (!msg) return;
|
|
|
|
| 93 |
if (msg.type === "auth" && msg.status === "ok") {
|
| 94 |
-
clearTimeout(timer);
|
| 95 |
console.log("[WS] Auth successful");
|
|
|
|
| 96 |
resolve();
|
| 97 |
}
|
| 98 |
if (msg.error) {
|
|
|
|
| 99 |
clearTimeout(timer);
|
| 100 |
-
console.error("[WS] Auth error:", msg.error);
|
| 101 |
reject(new Error(`WS auth error: ${msg.error}`));
|
| 102 |
}
|
| 103 |
});
|
| 104 |
ws.on("error", (err) => {
|
|
|
|
| 105 |
clearTimeout(timer);
|
| 106 |
reject(err);
|
| 107 |
});
|
| 108 |
});
|
| 109 |
|
| 110 |
// Send the chat body
|
|
|
|
| 111 |
ws.send(JSON.stringify({ body, headers }));
|
| 112 |
|
| 113 |
let assistantText = "";
|
|
@@ -115,7 +119,6 @@ export async function websocketChatStream(body, headers, onToken, abortSignal) {
|
|
| 115 |
let finished = false;
|
| 116 |
|
| 117 |
return new Promise((resolve, reject) => {
|
| 118 |
-
// Abort handling
|
| 119 |
if (abortSignal) {
|
| 120 |
abortSignal.addEventListener("abort", () => {
|
| 121 |
console.log("[WS] Aborted by signal");
|
|
@@ -124,31 +127,32 @@ export async function websocketChatStream(body, headers, onToken, abortSignal) {
|
|
| 124 |
});
|
| 125 |
}
|
| 126 |
|
| 127 |
-
// Handle messages
|
| 128 |
ws.on("message", (data) => {
|
| 129 |
const line = data.toString();
|
| 130 |
console.log("[WS MESSAGE RECEIVED]", line);
|
| 131 |
-
let payload = safeParse(line);
|
| 132 |
-
if (!payload) return;
|
| 133 |
|
| 134 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 135 |
if (payload.error) {
|
| 136 |
console.error("[WS PAYLOAD ERROR]", payload.error);
|
| 137 |
-
if (
|
| 138 |
return;
|
| 139 |
}
|
| 140 |
|
| 141 |
const delta = payload.choices?.[0]?.delta;
|
| 142 |
-
if (
|
| 143 |
-
|
| 144 |
-
// Stream content tokens
|
| 145 |
-
if (delta.content) {
|
| 146 |
assistantText += delta.content;
|
| 147 |
-
|
|
|
|
| 148 |
}
|
| 149 |
|
| 150 |
-
|
| 151 |
-
|
| 152 |
for (const call of delta.tool_calls) {
|
| 153 |
const entry = toolCallBuffer.get(call.index) ?? { arguments: "" };
|
| 154 |
if (call.id) entry.id = call.id;
|
|
@@ -158,7 +162,6 @@ export async function websocketChatStream(body, headers, onToken, abortSignal) {
|
|
| 158 |
}
|
| 159 |
}
|
| 160 |
|
| 161 |
-
// Finish detection
|
| 162 |
if (payload.choices?.[0]?.finish_reason) {
|
| 163 |
finished = true;
|
| 164 |
ws.close();
|
|
@@ -167,7 +170,7 @@ export async function websocketChatStream(body, headers, onToken, abortSignal) {
|
|
| 167 |
type: "function",
|
| 168 |
function: { name: t.name, arguments: t.arguments },
|
| 169 |
}));
|
| 170 |
-
console.log("[WS] Finished streaming");
|
| 171 |
resolve({ assistantText, toolCalls });
|
| 172 |
}
|
| 173 |
});
|
|
@@ -181,6 +184,9 @@ export async function websocketChatStream(body, headers, onToken, abortSignal) {
|
|
| 181 |
console.log(`[WS CLOSED] Code: ${code}, Reason: ${reason}`);
|
| 182 |
if (!finished) reject(new Error("WebSocket closed prematurely"));
|
| 183 |
});
|
|
|
|
|
|
|
|
|
|
| 184 |
});
|
| 185 |
}
|
| 186 |
|
|
|
|
| 4 |
import path from "path";
|
| 5 |
import { LIGHTNING_BASE } from "./config.js";
|
| 6 |
import WebSocket from "ws";
|
| 7 |
+
import crypto from "crypto";
|
| 8 |
|
| 9 |
const __dirname = path.dirname(fileURLToPath(import.meta.url));
|
| 10 |
const WORKER_PATH = path.join(__dirname, "searchWorker.js");
|
|
|
|
| 54 |
});
|
| 55 |
}
|
| 56 |
|
| 57 |
+
export async function websocketChatStreamVerbose(body, headers, onToken, abortSignal) {
|
| 58 |
const wsURL =
|
| 59 |
(process.env.LIGHTNING_BASE.startsWith("https")
|
| 60 |
? process.env.LIGHTNING_BASE.replace("https", "wss")
|
| 61 |
: process.env.LIGHTNING_BASE.replace("http", "ws")) + "/ws/chat";
|
| 62 |
|
| 63 |
+
console.log("[WS] Connecting to", wsURL);
|
| 64 |
const ws = new WebSocket(wsURL);
|
| 65 |
|
|
|
|
| 66 |
const safeParse = (str) => {
|
| 67 |
+
try { return JSON.parse(str); } catch { return null; }
|
|
|
|
| 68 |
};
|
| 69 |
|
|
|
|
| 70 |
await new Promise((resolve, reject) => {
|
| 71 |
const timer = setTimeout(() => reject(new Error("WS connection timeout")), 5000);
|
| 72 |
ws.on("open", () => {
|
| 73 |
+
console.log("[WS] Connection open");
|
| 74 |
clearTimeout(timer);
|
|
|
|
| 75 |
resolve();
|
| 76 |
});
|
| 77 |
ws.on("error", (err) => {
|
| 78 |
+
console.error("[WS] Connection error", err);
|
| 79 |
clearTimeout(timer);
|
|
|
|
| 80 |
reject(err);
|
| 81 |
});
|
| 82 |
});
|
| 83 |
|
| 84 |
+
// Authenticate
|
| 85 |
+
console.log("[WS] Sending auth key");
|
| 86 |
ws.send(JSON.stringify({ key: process.env.WEBSOCKET_KEY }));
|
| 87 |
+
|
| 88 |
await new Promise((resolve, reject) => {
|
| 89 |
const timer = setTimeout(() => reject(new Error("WS auth timeout")), 5000);
|
| 90 |
ws.on("message", (data) => {
|
|
|
|
| 91 |
console.log("[WS RAW MESSAGE]", data.toString());
|
| 92 |
+
const msg = safeParse(data.toString());
|
| 93 |
if (!msg) return;
|
| 94 |
+
|
| 95 |
if (msg.type === "auth" && msg.status === "ok") {
|
|
|
|
| 96 |
console.log("[WS] Auth successful");
|
| 97 |
+
clearTimeout(timer);
|
| 98 |
resolve();
|
| 99 |
}
|
| 100 |
if (msg.error) {
|
| 101 |
+
console.error("[WS] Auth error", msg.error);
|
| 102 |
clearTimeout(timer);
|
|
|
|
| 103 |
reject(new Error(`WS auth error: ${msg.error}`));
|
| 104 |
}
|
| 105 |
});
|
| 106 |
ws.on("error", (err) => {
|
| 107 |
+
console.error("[WS] Auth error event", err);
|
| 108 |
clearTimeout(timer);
|
| 109 |
reject(err);
|
| 110 |
});
|
| 111 |
});
|
| 112 |
|
| 113 |
// Send the chat body
|
| 114 |
+
console.log("[WS] Sending chat body", JSON.stringify(body, null, 2));
|
| 115 |
ws.send(JSON.stringify({ body, headers }));
|
| 116 |
|
| 117 |
let assistantText = "";
|
|
|
|
| 119 |
let finished = false;
|
| 120 |
|
| 121 |
return new Promise((resolve, reject) => {
|
|
|
|
| 122 |
if (abortSignal) {
|
| 123 |
abortSignal.addEventListener("abort", () => {
|
| 124 |
console.log("[WS] Aborted by signal");
|
|
|
|
| 127 |
});
|
| 128 |
}
|
| 129 |
|
|
|
|
| 130 |
ws.on("message", (data) => {
|
| 131 |
const line = data.toString();
|
| 132 |
console.log("[WS MESSAGE RECEIVED]", line);
|
|
|
|
|
|
|
| 133 |
|
| 134 |
+
const payload = safeParse(line);
|
| 135 |
+
if (!payload) {
|
| 136 |
+
console.warn("[WS] Failed to parse JSON:", line);
|
| 137 |
+
return;
|
| 138 |
+
}
|
| 139 |
+
|
| 140 |
+
// Log server-side errors
|
| 141 |
if (payload.error) {
|
| 142 |
console.error("[WS PAYLOAD ERROR]", payload.error);
|
| 143 |
+
if (onToken) onToken(`[ERROR] ${payload.error}`);
|
| 144 |
return;
|
| 145 |
}
|
| 146 |
|
| 147 |
const delta = payload.choices?.[0]?.delta;
|
| 148 |
+
if (delta?.content) {
|
|
|
|
|
|
|
|
|
|
| 149 |
assistantText += delta.content;
|
| 150 |
+
console.log("[WS TOKEN STREAM]", delta.content);
|
| 151 |
+
if (onToken) onToken(delta.content);
|
| 152 |
}
|
| 153 |
|
| 154 |
+
if (delta?.tool_calls) {
|
| 155 |
+
console.log("[WS TOOL CALLS]", delta.tool_calls);
|
| 156 |
for (const call of delta.tool_calls) {
|
| 157 |
const entry = toolCallBuffer.get(call.index) ?? { arguments: "" };
|
| 158 |
if (call.id) entry.id = call.id;
|
|
|
|
| 162 |
}
|
| 163 |
}
|
| 164 |
|
|
|
|
| 165 |
if (payload.choices?.[0]?.finish_reason) {
|
| 166 |
finished = true;
|
| 167 |
ws.close();
|
|
|
|
| 170 |
type: "function",
|
| 171 |
function: { name: t.name, arguments: t.arguments },
|
| 172 |
}));
|
| 173 |
+
console.log("[WS] Finished streaming", { assistantText, toolCalls });
|
| 174 |
resolve({ assistantText, toolCalls });
|
| 175 |
}
|
| 176 |
});
|
|
|
|
| 184 |
console.log(`[WS CLOSED] Code: ${code}, Reason: ${reason}`);
|
| 185 |
if (!finished) reject(new Error("WebSocket closed prematurely"));
|
| 186 |
});
|
| 187 |
+
|
| 188 |
+
ws.on("ping", () => console.log("[WS PING] Received ping"));
|
| 189 |
+
ws.on("pong", () => console.log("[WS PONG] Received pong"));
|
| 190 |
});
|
| 191 |
}
|
| 192 |
|