fix: emit opencode event stream otel on client close

This commit is contained in:
UniDesk Codex
2026-06-30 14:40:15 +08:00
parent a1e9a5200a
commit d42007c04e
3 changed files with 162 additions and 9 deletions
+15 -9
View File
@@ -76,6 +76,14 @@ export function proxyCloudApiRequest({
let timeout = null;
let streamTransformEnded = false;
const endStreamTransform = () => {
if (!streamTransformEnded && streamTransform?.end) {
streamTransformEnded = true;
return streamTransform.end();
}
return null;
};
const settle = (callback, value) => {
if (settled) return;
settled = true;
@@ -107,10 +115,7 @@ export function proxyCloudApiRequest({
writeStreamOutput(response, streamTransform?.write ? streamTransform.write(chunk) : chunk);
});
upstreamResponse.on("end", () => {
if (!streamTransformEnded && streamTransform?.end) {
streamTransformEnded = true;
writeStreamOutput(response, streamTransform.end());
}
writeStreamOutput(response, endStreamTransform());
response.end();
settle(resolve, { statusCode: upstreamStatusCode, streaming: true, streamTransformStats: streamTransform?.stats });
});
@@ -123,14 +128,15 @@ export function proxyCloudApiRequest({
settle(reject, error);
});
upstreamResponse.on("close", () => {
if (!streamTransformEnded && streamTransform?.end) {
streamTransformEnded = true;
writeStreamOutput(response, streamTransform.end());
}
writeStreamOutput(response, endStreamTransform());
if (!response.writableEnded) response.end();
settle(resolve, { statusCode: upstreamStatusCode, streaming: true, streamTransformStats: streamTransform?.stats });
});
response.on("close", () => upstream.destroy());
response.on("close", () => {
endStreamTransform();
upstream.destroy();
settle(resolve, { statusCode: upstreamStatusCode, streaming: true, errorCode: "client_closed", streamTransformStats: streamTransform?.stats });
});
return;
}
@@ -144,6 +144,62 @@ test("cloud web proxy does not reject after SSE headers are already sent", async
}
});
test("cloud web proxy resolves stream transform stats when the client closes SSE", async () => {
const upstream = createServer((request, response) => {
request.resume();
response.writeHead(200, { "content-type": "text/event-stream; charset=utf-8" });
response.write('data: {"directory":"/workspace","payload":{"type":"session.status"}}\n\n');
});
await listen(upstream);
let proxyResult = null;
const transform = {
stats: { enabled: true, writes: 0, ended: false },
write(chunk) {
this.stats.writes += 1;
return chunk;
},
end() {
this.stats.ended = true;
return "";
}
};
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: true,
streamTransform: transform
}).then((result) => {
proxyResult = result;
}).catch((error) => response.destroy(error));
});
await listen(proxy);
try {
const response = await fetch(`${serverUrl(proxy)}/global/event`, {
headers: { accept: "text/event-stream" }
});
assert.equal(response.status, 200);
const reader = response.body.getReader();
const first = await reader.read();
assert.equal(first.done, false);
await reader.cancel();
await waitFor(() => proxyResult !== null, "expected proxy result after client close");
assert.equal(proxyResult.statusCode, 200);
assert.equal(proxyResult.streaming, true);
assert.equal(proxyResult.errorCode, "client_closed");
assert.deepEqual(proxyResult.streamTransformStats, { enabled: true, writes: 1, ended: true });
} 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();
@@ -189,6 +245,15 @@ function close(server) {
});
}
async function waitFor(predicate, message, { timeoutMs = 500, intervalMs = 10 } = {}) {
const deadline = Date.now() + timeoutMs;
while (Date.now() < deadline) {
if (predicate()) return;
await new Promise((resolve) => setTimeout(resolve, intervalMs));
}
assert.fail(message);
}
function serverUrl(server) {
const address = server.address();
return `http://127.0.0.1:${address.port}`;
@@ -825,6 +825,88 @@ test("cloud web OpenCode proxy emits bounded OTLP span", async () => {
}
});
test("cloud web OpenCode event stream emits rewrite OTLP stats when the client closes", async () => {
const otelBodies = [];
const otel = createServer(async (request, response) => {
let body = "";
for await (const chunk of request) body += chunk;
otelBodies.push(JSON.parse(body));
response.writeHead(200, { "content-type": "application/json" });
response.end("{}\n");
});
await listen(otel);
const restoreEnv = withEnv({ OTEL_EXPORTER_OTLP_TRACES_ENDPOINT: `${serverUrl(otel)}/v1/traces` });
const cloudApi = createServer((request, response) => {
request.resume();
response.writeHead(200, { "content-type": "application/json" });
response.end(JSON.stringify({ authenticated: request.headers.cookie === "hwlab_session=session-a" }));
});
const opencodeResponses = [];
const opencode = createServer((request, response) => {
request.resume();
opencodeResponses.push(response);
response.writeHead(200, { "content-type": "text/event-stream; charset=utf-8" });
response.write('data: {"directory":"/workspace","payload":{"type":"session.status","properties":{"sessionID":"ses_1"}}}\n\n');
});
await listen(cloudApi);
await listen(opencode);
const cloudWeb = createCloudWebServer({
serviceId: "hwlab-cloud-web",
roots: [],
cloudApiBaseUrl: serverUrl(cloudApi),
cloudApiProxyTimeoutMs: 1000,
opencodeUpstreamUrl: serverUrl(opencode),
opencodeProxyHost: "127.0.0.1",
opencodeProxyTimeoutMs: 1000,
opencodeUsername: "oc_user",
opencodePassword: "oc_password",
opencodeEventDirectoryRewrite: { from: "/workspace", to: "/" },
healthPayload: () => ({ status: "ok" }),
sendJson(response, statusCode, body) {
const payload = JSON.stringify(body);
response.writeHead(statusCode, { "content-type": "application/json", "content-length": Buffer.byteLength(payload) });
response.end(payload);
}
});
await listen(cloudWeb);
try {
const response = await fetch(`${serverUrl(cloudWeb)}/global/event`, {
headers: { accept: "text/event-stream", cookie: "hwlab_session=session-a" }
});
assert.equal(response.status, 200);
const reader = response.body.getReader();
const first = await reader.read();
assert.equal(first.done, false);
const text = new TextDecoder().decode(first.value);
assert.match(text, /"directory":"\/"/u);
await reader.cancel();
await waitFor(() => otelBodies.length > 0, "expected OpenCode event stream OTel span");
const spans = otelBodies.flatMap((body) => body.resourceSpans?.flatMap((resourceSpan) => resourceSpan.scopeSpans?.flatMap((scopeSpan) => scopeSpan.spans ?? []) ?? []) ?? []);
const span = spans.find((item) => item.name === "opencode.proxy.request");
assert.ok(span, "expected opencode.proxy.request span");
const attrs = new Map((span.attributes ?? []).map((entry) => [entry.key, entry.value?.stringValue ?? entry.value?.intValue ?? entry.value?.boolValue]));
assert.equal(attrs.get("http.route"), "/global/event");
assert.equal(attrs.get("opencode.proxy.streaming"), true);
assert.equal(attrs.get("opencode.proxy.sse.directory_rewrite_enabled"), true);
assert.equal(attrs.get("opencode.proxy.sse.directory_rewrite_from"), "/workspace");
assert.equal(attrs.get("opencode.proxy.sse.directory_rewrite_to"), "/");
assert.equal(attrs.get("opencode.proxy.sse.directory_rewrite_data_lines"), "1");
assert.equal(attrs.get("opencode.proxy.sse.directory_rewrite_events"), "1");
assert.equal(attrs.get("opencode.proxy.sse.directory_rewrite_json_errors"), "0");
} finally {
for (const response of opencodeResponses) response.destroy();
restoreEnv();
await close(cloudWeb);
await close(opencode);
await close(cloudApi);
await close(otel);
}
});
test("cloud web OpenCode proxy accepts short-lived tickets minted by the shell", async () => {
const restoreEnv = withEnv({
HWLAB_CLOUD_WEB_DISPLAY_TIME_ZONE: "Asia/Shanghai",