diff --git a/deploy/deploy.yaml b/deploy/deploy.yaml index 0a656ae6..bac286b9 100644 --- a/deploy/deploy.yaml +++ b/deploy/deploy.yaml @@ -585,6 +585,8 @@ lanes: HWLAB_CLOUD_WEB_OPENCODE_UPSTREAM_URL: http://opencode-server.hwlab-v03.svc.cluster.local:4096 HWLAB_CLOUD_WEB_OPENCODE_USERNAME: secretRef:hwlab-opencode-server-auth/username HWLAB_CLOUD_WEB_OPENCODE_PASSWORD: secretRef:hwlab-opencode-server-auth/password + HWLAB_CLOUD_WEB_OPENCODE_EVENT_DIRECTORY_FROM: /workspace + HWLAB_CLOUD_WEB_OPENCODE_EVENT_DIRECTORY_TO: / OTEL_EXPORTER_OTLP_TRACES_ENDPOINT: http://otel-collector.platform-infra.svc.cluster.local:4318/v1/traces OTEL_SERVICE_NAME: hwlab-cloud-web observable: true diff --git a/internal/dev-entrypoint/cloud-web-proxy.mjs b/internal/dev-entrypoint/cloud-web-proxy.mjs index d3cdadf4..8ccfb6e5 100644 --- a/internal/dev-entrypoint/cloud-web-proxy.mjs +++ b/internal/dev-entrypoint/cloud-web-proxy.mjs @@ -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)) { diff --git a/internal/dev-entrypoint/cloud-web-runtime.mjs b/internal/dev-entrypoint/cloud-web-runtime.mjs index 01f7f6c4..85acc0d9 100644 --- a/internal/dev-entrypoint/cloud-web-runtime.mjs +++ b/internal/dev-entrypoint/cloud-web-runtime.mjs @@ -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"; diff --git a/internal/dev-entrypoint/cloud-web-runtime.test.mjs b/internal/dev-entrypoint/cloud-web-runtime.test.mjs index 4b50f194..ff4a24b2 100644 --- a/internal/dev-entrypoint/cloud-web-runtime.test.mjs +++ b/internal/dev-entrypoint/cloud-web-runtime.test.mjs @@ -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) => {