Files
pikasTech-HWLAB/internal/dev-entrypoint/cloud-web-proxy.test.mjs
T

188 lines
6.7 KiB
JavaScript

import assert from "node:assert/strict";
import { createServer } from "node:http";
import test from "node:test";
import {
isCloudWebSseRoute,
proxyCloudApiRequest,
upstreamRequestHeaders
} from "./cloud-web-proxy.mjs";
test("cloud web proxy preserves Code Agent short-connection headers", () => {
const headers = upstreamRequestHeaders({
headers: {
accept: "application/json",
"content-type": "application/json",
prefer: "respond-async",
"x-hwlab-short-connection": "1",
"x-trace-id": "trc_proxy_short_connection",
"x-request-id": "req_proxy_short_connection",
authorization: "Bearer must-not-forward"
}
}, "{\"message\":\"hello\"}");
assert.equal(headers.accept, "application/json");
assert.equal(headers["content-type"], "application/json");
assert.equal(headers.prefer, "respond-async");
assert.equal(headers["x-hwlab-short-connection"], "1");
assert.equal(headers["x-trace-id"], "trc_proxy_short_connection");
assert.equal(headers["x-request-id"], "req_proxy_short_connection");
assert.equal(headers.authorization, undefined);
assert.equal(headers["content-length"], Buffer.byteLength("{\"message\":\"hello\"}"));
});
test("cloud web proxy streams SSE chunks without fixed content-length", async () => {
const upstream = createServer((request, response) => {
assert.equal(request.headers["x-trace-id"], "trc_stream_proxy");
response.writeHead(200, {
"content-type": "text/event-stream; charset=utf-8",
"cache-control": "no-store",
"x-accel-buffering": "no"
});
setTimeout(() => response.write("event: snapshot\ndata: {\"events\":[{\"seq\":1}]}\n\n"), 40);
setTimeout(() => response.write("event: runnerTrace\ndata: {\"snapshot\":{\"events\":[{\"seq\":1},{\"seq\":2}]}}\n\n"), 140);
setTimeout(() => response.end(), 500);
});
await listen(upstream);
const proxy = createServer((request, response) => {
const url = new URL(request.url || "/", "http://cloud-web.local");
proxyCloudApiRequest({
target: new URL(url.pathname + url.search, serverUrl(upstream)),
request,
response,
timeoutMs: 1000,
forceStream: isCloudWebSseRoute(url.pathname)
}).catch((error) => response.destroy(error));
});
await listen(proxy);
try {
const startedAt = Date.now();
const response = await fetch(`${serverUrl(proxy)}/v1/agent/chat/trace/trc_stream_proxy/stream`, {
headers: {
"x-trace-id": "trc_stream_proxy",
accept: "text/event-stream"
}
});
assert.equal(response.status, 200);
assert.equal(response.headers.get("content-length"), null);
assert.match(response.headers.get("content-type") || "", /^text\/event-stream/iu);
assert.equal(response.headers.get("cache-control"), "no-store");
assert.equal(response.headers.get("x-accel-buffering"), "no");
const reader = response.body.getReader();
const decoder = new TextDecoder();
let text = "";
let firstEventAt = null;
let secondEventAt = null;
while (!secondEventAt) {
const { value, done } = await reader.read();
if (done) break;
text += decoder.decode(value, { stream: true });
const eventCount = (text.match(/^event:/gmu) || []).length;
if (eventCount >= 1 && firstEventAt === null) firstEventAt = Date.now() - startedAt;
if (eventCount >= 2) secondEventAt = Date.now() - startedAt;
}
await reader.cancel().catch(() => {});
assert.ok(firstEventAt !== null && firstEventAt < 250, `first SSE event was not flushed promptly: ${firstEventAt}ms`);
assert.ok(secondEventAt !== null && secondEventAt < 450, `second SSE event was not flushed promptly: ${secondEventAt}ms`);
assert.ok(secondEventAt > firstEventAt, "event count must grow at two distinct times");
} finally {
await close(proxy);
await close(upstream);
}
});
test("cloud web proxy does not reject after SSE headers are already sent", async () => {
const upstream = createServer((request, response) => {
request.resume();
response.writeHead(200, { "content-type": "text/event-stream; charset=utf-8" });
response.write("event: snapshot\ndata: {\"events\":[{\"seq\":1}]}\n\n");
setTimeout(() => response.destroy(new Error("simulated upstream stream failure")), 40);
});
await listen(upstream);
let rejected = false;
const proxy = createServer((request, response) => {
const url = new URL(request.url || "/", "http://cloud-web.local");
proxyCloudApiRequest({
target: new URL(url.pathname + url.search, serverUrl(upstream)),
request,
response,
timeoutMs: 1000,
forceStream: isCloudWebSseRoute(url.pathname)
}).catch((error) => {
rejected = true;
if (!response.headersSent) response.destroy(error);
});
});
await listen(proxy);
try {
const response = await fetch(`${serverUrl(proxy)}/v1/agent/chat/trace/trc_stream_proxy_error/stream`, {
headers: { accept: "text/event-stream" }
});
assert.equal(response.status, 200);
const reader = response.body.getReader();
await reader.read();
await reader.read().catch(() => {});
await new Promise((resolve) => setTimeout(resolve, 80));
assert.equal(rejected, false, "stream errors after headers must not bubble into JSON error writers");
} finally {
await close(proxy);
await close(upstream);
}
});
test("cloud web proxy buffers normal JSON and sets content-length", async () => {
const upstream = createServer((request, response) => {
request.resume();
setTimeout(() => {
response.writeHead(200, { "content-type": "application/json; charset=utf-8" });
response.end(JSON.stringify({ ok: true }));
}, 20);
});
await listen(upstream);
const proxy = createServer((request, response) => {
proxyCloudApiRequest({
target: new URL(request.url || "/", serverUrl(upstream)),
request,
response,
body: "",
timeoutMs: 1000
}).catch((error) => response.destroy(error));
});
await listen(proxy);
try {
const response = await fetch(`${serverUrl(proxy)}/v1`);
assert.equal(response.status, 200);
assert.equal(response.headers.get("content-length"), String(Buffer.byteLength("{\"ok\":true}")));
assert.deepEqual(await response.json(), { ok: true });
} finally {
await close(proxy);
await close(upstream);
}
});
function listen(server) {
return new Promise((resolve, reject) => {
server.once("error", reject);
server.listen(0, "127.0.0.1", resolve);
});
}
function close(server) {
return new Promise((resolve, reject) => {
server.close((error) => (error ? reject(error) : resolve()));
});
}
function serverUrl(server) {
const address = server.address();
return `http://127.0.0.1:${address.port}`;
}