fix: rewrite opencode event directory for live ui
This commit is contained in:
@@ -67,12 +67,14 @@ export function proxyCloudApiRequest({
|
||||
body = "",
|
||||
timeoutMs,
|
||||
forceStream = false,
|
||||
extraResponseHeaders = {}
|
||||
extraResponseHeaders = {},
|
||||
streamTransform = null
|
||||
}) {
|
||||
return new Promise((resolve, reject) => {
|
||||
let timedOut = false;
|
||||
let settled = false;
|
||||
let timeout = null;
|
||||
let streamTransformEnded = false;
|
||||
|
||||
const settle = (callback, value) => {
|
||||
if (settled) return;
|
||||
@@ -102,11 +104,15 @@ export function proxyCloudApiRequest({
|
||||
response.flushHeaders?.();
|
||||
upstreamResponse.on("data", (chunk) => {
|
||||
armTimeout();
|
||||
response.write(chunk);
|
||||
writeStreamOutput(response, streamTransform?.write ? streamTransform.write(chunk) : chunk);
|
||||
});
|
||||
upstreamResponse.on("end", () => {
|
||||
if (!streamTransformEnded && streamTransform?.end) {
|
||||
streamTransformEnded = true;
|
||||
writeStreamOutput(response, streamTransform.end());
|
||||
}
|
||||
response.end();
|
||||
settle(resolve, { statusCode: upstreamStatusCode, streaming: true });
|
||||
settle(resolve, { statusCode: upstreamStatusCode, streaming: true, streamTransformStats: streamTransform?.stats });
|
||||
});
|
||||
upstreamResponse.on("error", (error) => {
|
||||
if (response.headersSent) {
|
||||
@@ -117,8 +123,12 @@ export function proxyCloudApiRequest({
|
||||
settle(reject, error);
|
||||
});
|
||||
upstreamResponse.on("close", () => {
|
||||
if (!streamTransformEnded && streamTransform?.end) {
|
||||
streamTransformEnded = true;
|
||||
writeStreamOutput(response, streamTransform.end());
|
||||
}
|
||||
if (!response.writableEnded) response.end();
|
||||
settle(resolve, { statusCode: upstreamStatusCode, streaming: true });
|
||||
settle(resolve, { statusCode: upstreamStatusCode, streaming: true, streamTransformStats: streamTransform?.stats });
|
||||
});
|
||||
response.on("close", () => upstream.destroy());
|
||||
return;
|
||||
@@ -157,6 +167,15 @@ export function proxyCloudApiRequest({
|
||||
});
|
||||
}
|
||||
|
||||
function writeStreamOutput(response, output) {
|
||||
if (output === undefined || output === null || output === "") return;
|
||||
if (Array.isArray(output)) {
|
||||
for (const item of output) writeStreamOutput(response, item);
|
||||
return;
|
||||
}
|
||||
response.write(output);
|
||||
}
|
||||
|
||||
function mergeResponseHeaders(baseHeaders = {}, extraHeaders = {}) {
|
||||
const result = { ...baseHeaders };
|
||||
for (const [key, value] of Object.entries(extraHeaders)) {
|
||||
|
||||
@@ -63,7 +63,8 @@ export function createCloudWebServer({
|
||||
max: 2400000
|
||||
}),
|
||||
opencodeUsername = optionalRuntimeConfigEnv("HWLAB_CLOUD_WEB_OPENCODE_USERNAME") || "",
|
||||
opencodePassword = optionalRuntimeConfigEnv("HWLAB_CLOUD_WEB_OPENCODE_PASSWORD") || ""
|
||||
opencodePassword = optionalRuntimeConfigEnv("HWLAB_CLOUD_WEB_OPENCODE_PASSWORD") || "",
|
||||
opencodeEventDirectoryRewrite = opencodeEventDirectoryRewriteFromEnv()
|
||||
}) {
|
||||
return createServer(async (request, response) => {
|
||||
attachRequestErrorHandlers({ request, response, serviceId });
|
||||
@@ -78,7 +79,7 @@ export function createCloudWebServer({
|
||||
return;
|
||||
}
|
||||
if (isOpencodeProxyRequest(request, url, opencodeProxyHost)) {
|
||||
await proxyOpencodeRequest({ request, response, url, cloudApiBaseUrl, cloudApiProxyTimeoutMs, opencodeUpstreamUrl, opencodeProxyTimeoutMs, opencodeUsername, opencodePassword, serviceId, sendJson });
|
||||
await proxyOpencodeRequest({ request, response, url, cloudApiBaseUrl, cloudApiProxyTimeoutMs, opencodeUpstreamUrl, opencodeProxyTimeoutMs, opencodeUsername, opencodePassword, opencodeEventDirectoryRewrite, serviceId, sendJson });
|
||||
return;
|
||||
}
|
||||
if (await handleCloudWebAuth({ request, response, url, cloudApiBaseUrl, cloudApiProxyTimeoutMs, serviceId, sendJson })) {
|
||||
@@ -304,6 +305,71 @@ function opencodeProxyHostFromEnv() {
|
||||
}
|
||||
}
|
||||
|
||||
function opencodeEventDirectoryRewriteFromEnv() {
|
||||
const from = optionalRuntimeConfigEnv("HWLAB_CLOUD_WEB_OPENCODE_EVENT_DIRECTORY_FROM");
|
||||
const to = optionalRuntimeConfigEnv("HWLAB_CLOUD_WEB_OPENCODE_EVENT_DIRECTORY_TO");
|
||||
if (!from || !to || from === to) return null;
|
||||
return { from, to };
|
||||
}
|
||||
|
||||
function opencodeEventDirectoryTransformForTarget(target, rewrite) {
|
||||
if (!rewrite || target?.pathname !== "/global/event") return null;
|
||||
return opencodeEventDirectoryRewriteTransform(rewrite);
|
||||
}
|
||||
|
||||
function opencodeEventDirectoryRewriteTransform(rewrite) {
|
||||
const stats = {
|
||||
enabled: true,
|
||||
from: rewrite.from,
|
||||
to: rewrite.to,
|
||||
dataLines: 0,
|
||||
rewrittenEvents: 0,
|
||||
jsonErrors: 0
|
||||
};
|
||||
let pending = "";
|
||||
return {
|
||||
stats,
|
||||
write(chunk) {
|
||||
pending += Buffer.isBuffer(chunk) ? chunk.toString("utf8") : String(chunk ?? "");
|
||||
const output = [];
|
||||
let newlineIndex = pending.indexOf("\n");
|
||||
while (newlineIndex >= 0) {
|
||||
const line = pending.slice(0, newlineIndex + 1);
|
||||
pending = pending.slice(newlineIndex + 1);
|
||||
output.push(rewriteOpencodeEventDirectoryLine(line, rewrite, stats));
|
||||
newlineIndex = pending.indexOf("\n");
|
||||
}
|
||||
return output;
|
||||
},
|
||||
end() {
|
||||
if (!pending) return "";
|
||||
const line = pending;
|
||||
pending = "";
|
||||
return rewriteOpencodeEventDirectoryLine(line, rewrite, stats);
|
||||
}
|
||||
};
|
||||
}
|
||||
|
||||
function rewriteOpencodeEventDirectoryLine(line, rewrite, stats) {
|
||||
const match = String(line).match(/^(data:\s*)(.*?)(\r?\n)?$/u);
|
||||
if (!match) return line;
|
||||
const [, prefix, payload, lineEnding = ""] = match;
|
||||
const text = payload.trim();
|
||||
if (!text) return line;
|
||||
stats.dataLines += 1;
|
||||
try {
|
||||
const parsed = JSON.parse(text);
|
||||
if (parsed && typeof parsed === "object" && parsed.directory === rewrite.from) {
|
||||
parsed.directory = rewrite.to;
|
||||
stats.rewrittenEvents += 1;
|
||||
return `${prefix}${JSON.stringify(parsed)}${lineEnding}`;
|
||||
}
|
||||
} catch {
|
||||
stats.jsonErrors += 1;
|
||||
}
|
||||
return line;
|
||||
}
|
||||
|
||||
function validateDisplayTimeZone(timeZone) {
|
||||
try {
|
||||
new Intl.DateTimeFormat("en-US", { timeZone }).format(new Date(0));
|
||||
@@ -663,7 +729,7 @@ export async function proxyCloudApi({ request, response, url, cloudApiBaseUrl, c
|
||||
}
|
||||
}
|
||||
|
||||
async function proxyOpencodeRequest({ request, response, url, cloudApiBaseUrl, cloudApiProxyTimeoutMs, opencodeUpstreamUrl, opencodeProxyTimeoutMs, opencodeUsername, opencodePassword, serviceId, sendJson }) {
|
||||
async function proxyOpencodeRequest({ request, response, url, cloudApiBaseUrl, cloudApiProxyTimeoutMs, opencodeUpstreamUrl, opencodeProxyTimeoutMs, opencodeUsername, opencodePassword, opencodeEventDirectoryRewrite = null, serviceId, sendJson }) {
|
||||
const traceContext = cloudWebTraceContext(request);
|
||||
const startedAtMs = Date.now();
|
||||
if (!cloudApiBaseUrl) {
|
||||
@@ -733,6 +799,7 @@ async function proxyOpencodeRequest({ request, response, url, cloudApiBaseUrl, c
|
||||
}
|
||||
|
||||
const target = opencodeTargetUrl(url, opencodeUpstreamUrl);
|
||||
const streamTransform = opencodeEventDirectoryTransformForTarget(target, opencodeEventDirectoryRewrite);
|
||||
const extraResponseHeaders = {
|
||||
...cloudWebTraceHeaders(traceContext, serviceId),
|
||||
...(ticketAuth.ticket ? { "set-cookie": opencodeTicketSetCookie(ticketAuth.ticket) } : {})
|
||||
@@ -748,7 +815,8 @@ async function proxyOpencodeRequest({ request, response, url, cloudApiBaseUrl, c
|
||||
response,
|
||||
body: request.method === "GET" || request.method === "HEAD" ? "" : body,
|
||||
timeoutMs: opencodeProxyTimeoutMs,
|
||||
extraResponseHeaders
|
||||
extraResponseHeaders,
|
||||
streamTransform
|
||||
});
|
||||
emitOpencodeProxySpanAsync({
|
||||
traceContext,
|
||||
@@ -761,7 +829,8 @@ async function proxyOpencodeRequest({ request, response, url, cloudApiBaseUrl, c
|
||||
errorCode: proxyResult?.errorCode || "",
|
||||
ticketAuth,
|
||||
sessionStatusCode: session.statusCode || 0,
|
||||
streaming: proxyResult?.streaming === true
|
||||
streaming: proxyResult?.streaming === true,
|
||||
streamTransformStats: proxyResult?.streamTransformStats
|
||||
});
|
||||
} catch (error) {
|
||||
if (response.headersSent || response.writableEnded) {
|
||||
@@ -916,7 +985,7 @@ function opencodeTraceDiagnostic(traceContext, extra = {}) {
|
||||
};
|
||||
}
|
||||
|
||||
function emitOpencodeProxySpanAsync({ traceContext, serviceId, request, url, target = null, startedAtMs, statusCode = 0, errorCode = "", ticketAuth = null, sessionStatusCode = 0, streaming = false } = {}) {
|
||||
function emitOpencodeProxySpanAsync({ traceContext, serviceId, request, url, target = null, startedAtMs, statusCode = 0, errorCode = "", ticketAuth = null, sessionStatusCode = 0, streaming = false, streamTransformStats = null } = {}) {
|
||||
const endpoint = otelTracesEndpoint(process.env);
|
||||
if (!endpoint || !traceContext?.traceId || !traceContext?.spanId) return;
|
||||
const endedAtMs = Date.now();
|
||||
@@ -944,6 +1013,7 @@ function emitOpencodeProxySpanAsync({ traceContext, serviceId, request, url, tar
|
||||
"opencode.proxy.ticket_accepted": Boolean(ticketAuth?.ticket),
|
||||
"opencode.proxy.session_status_code": sessionStatusCode || undefined,
|
||||
"opencode.proxy.streaming": streaming === true,
|
||||
...opencodeStreamTransformStatsAttributes(streamTransformStats),
|
||||
valuesPrinted: false
|
||||
}
|
||||
};
|
||||
@@ -961,6 +1031,18 @@ function emitOpencodeProxySpanAsync({ traceContext, serviceId, request, url, tar
|
||||
}, 0);
|
||||
}
|
||||
|
||||
function opencodeStreamTransformStatsAttributes(stats) {
|
||||
if (!stats?.enabled) return {};
|
||||
return {
|
||||
"opencode.proxy.sse.directory_rewrite_enabled": true,
|
||||
"opencode.proxy.sse.directory_rewrite_from": stats.from,
|
||||
"opencode.proxy.sse.directory_rewrite_to": stats.to,
|
||||
"opencode.proxy.sse.directory_rewrite_data_lines": stats.dataLines,
|
||||
"opencode.proxy.sse.directory_rewrite_events": stats.rewrittenEvents,
|
||||
"opencode.proxy.sse.directory_rewrite_json_errors": stats.jsonErrors
|
||||
};
|
||||
}
|
||||
|
||||
function opencodeProxyRoute(pathname) {
|
||||
const path = String(pathname || "/");
|
||||
if (path === "/global/event") return "/global/event";
|
||||
|
||||
@@ -705,6 +705,58 @@ test("cloud web OpenCode proxy injects upstream Basic Auth without forwarding HW
|
||||
}
|
||||
});
|
||||
|
||||
test("cloud web OpenCode proxy rewrites configured event stream directory", async () => {
|
||||
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 opencode = 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","properties":{"sessionID":"ses_1","status":{"type":"idle"}}}}\n\n');
|
||||
response.write('data: {"directory":"/other","payload":{"type":"server.heartbeat","properties":{}}}\n\n');
|
||||
response.end();
|
||||
});
|
||||
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 text = await response.text();
|
||||
assert.match(text, /"directory":"\/"/u);
|
||||
assert.doesNotMatch(text, /"directory":"\/workspace"/u);
|
||||
assert.match(text, /"directory":"\/other"/u);
|
||||
} finally {
|
||||
await close(cloudWeb);
|
||||
await close(opencode);
|
||||
await close(cloudApi);
|
||||
}
|
||||
});
|
||||
|
||||
test("cloud web OpenCode proxy emits bounded OTLP span", async () => {
|
||||
const otelBodies = [];
|
||||
const otel = createServer(async (request, response) => {
|
||||
|
||||
Reference in New Issue
Block a user