Merge pull request #2291 from pikasTech/issue-2274-opencode-event-directory
fix: rewrite opencode event directory for live ui
This commit is contained in:
@@ -585,6 +585,8 @@ lanes:
|
|||||||
HWLAB_CLOUD_WEB_OPENCODE_UPSTREAM_URL: http://opencode-server.hwlab-v03.svc.cluster.local:4096
|
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_USERNAME: secretRef:hwlab-opencode-server-auth/username
|
||||||
HWLAB_CLOUD_WEB_OPENCODE_PASSWORD: secretRef:hwlab-opencode-server-auth/password
|
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_EXPORTER_OTLP_TRACES_ENDPOINT: http://otel-collector.platform-infra.svc.cluster.local:4318/v1/traces
|
||||||
OTEL_SERVICE_NAME: hwlab-cloud-web
|
OTEL_SERVICE_NAME: hwlab-cloud-web
|
||||||
observable: true
|
observable: true
|
||||||
|
|||||||
@@ -67,12 +67,14 @@ export function proxyCloudApiRequest({
|
|||||||
body = "",
|
body = "",
|
||||||
timeoutMs,
|
timeoutMs,
|
||||||
forceStream = false,
|
forceStream = false,
|
||||||
extraResponseHeaders = {}
|
extraResponseHeaders = {},
|
||||||
|
streamTransform = null
|
||||||
}) {
|
}) {
|
||||||
return new Promise((resolve, reject) => {
|
return new Promise((resolve, reject) => {
|
||||||
let timedOut = false;
|
let timedOut = false;
|
||||||
let settled = false;
|
let settled = false;
|
||||||
let timeout = null;
|
let timeout = null;
|
||||||
|
let streamTransformEnded = false;
|
||||||
|
|
||||||
const settle = (callback, value) => {
|
const settle = (callback, value) => {
|
||||||
if (settled) return;
|
if (settled) return;
|
||||||
@@ -102,11 +104,15 @@ export function proxyCloudApiRequest({
|
|||||||
response.flushHeaders?.();
|
response.flushHeaders?.();
|
||||||
upstreamResponse.on("data", (chunk) => {
|
upstreamResponse.on("data", (chunk) => {
|
||||||
armTimeout();
|
armTimeout();
|
||||||
response.write(chunk);
|
writeStreamOutput(response, streamTransform?.write ? streamTransform.write(chunk) : chunk);
|
||||||
});
|
});
|
||||||
upstreamResponse.on("end", () => {
|
upstreamResponse.on("end", () => {
|
||||||
|
if (!streamTransformEnded && streamTransform?.end) {
|
||||||
|
streamTransformEnded = true;
|
||||||
|
writeStreamOutput(response, streamTransform.end());
|
||||||
|
}
|
||||||
response.end();
|
response.end();
|
||||||
settle(resolve, { statusCode: upstreamStatusCode, streaming: true });
|
settle(resolve, { statusCode: upstreamStatusCode, streaming: true, streamTransformStats: streamTransform?.stats });
|
||||||
});
|
});
|
||||||
upstreamResponse.on("error", (error) => {
|
upstreamResponse.on("error", (error) => {
|
||||||
if (response.headersSent) {
|
if (response.headersSent) {
|
||||||
@@ -117,8 +123,12 @@ export function proxyCloudApiRequest({
|
|||||||
settle(reject, error);
|
settle(reject, error);
|
||||||
});
|
});
|
||||||
upstreamResponse.on("close", () => {
|
upstreamResponse.on("close", () => {
|
||||||
|
if (!streamTransformEnded && streamTransform?.end) {
|
||||||
|
streamTransformEnded = true;
|
||||||
|
writeStreamOutput(response, streamTransform.end());
|
||||||
|
}
|
||||||
if (!response.writableEnded) response.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());
|
response.on("close", () => upstream.destroy());
|
||||||
return;
|
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 = {}) {
|
function mergeResponseHeaders(baseHeaders = {}, extraHeaders = {}) {
|
||||||
const result = { ...baseHeaders };
|
const result = { ...baseHeaders };
|
||||||
for (const [key, value] of Object.entries(extraHeaders)) {
|
for (const [key, value] of Object.entries(extraHeaders)) {
|
||||||
|
|||||||
@@ -63,7 +63,8 @@ export function createCloudWebServer({
|
|||||||
max: 2400000
|
max: 2400000
|
||||||
}),
|
}),
|
||||||
opencodeUsername = optionalRuntimeConfigEnv("HWLAB_CLOUD_WEB_OPENCODE_USERNAME") || "",
|
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) => {
|
return createServer(async (request, response) => {
|
||||||
attachRequestErrorHandlers({ request, response, serviceId });
|
attachRequestErrorHandlers({ request, response, serviceId });
|
||||||
@@ -78,7 +79,7 @@ export function createCloudWebServer({
|
|||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
if (isOpencodeProxyRequest(request, url, opencodeProxyHost)) {
|
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;
|
return;
|
||||||
}
|
}
|
||||||
if (await handleCloudWebAuth({ request, response, url, cloudApiBaseUrl, cloudApiProxyTimeoutMs, serviceId, sendJson })) {
|
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) {
|
function validateDisplayTimeZone(timeZone) {
|
||||||
try {
|
try {
|
||||||
new Intl.DateTimeFormat("en-US", { timeZone }).format(new Date(0));
|
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 traceContext = cloudWebTraceContext(request);
|
||||||
const startedAtMs = Date.now();
|
const startedAtMs = Date.now();
|
||||||
if (!cloudApiBaseUrl) {
|
if (!cloudApiBaseUrl) {
|
||||||
@@ -733,6 +799,7 @@ async function proxyOpencodeRequest({ request, response, url, cloudApiBaseUrl, c
|
|||||||
}
|
}
|
||||||
|
|
||||||
const target = opencodeTargetUrl(url, opencodeUpstreamUrl);
|
const target = opencodeTargetUrl(url, opencodeUpstreamUrl);
|
||||||
|
const streamTransform = opencodeEventDirectoryTransformForTarget(target, opencodeEventDirectoryRewrite);
|
||||||
const extraResponseHeaders = {
|
const extraResponseHeaders = {
|
||||||
...cloudWebTraceHeaders(traceContext, serviceId),
|
...cloudWebTraceHeaders(traceContext, serviceId),
|
||||||
...(ticketAuth.ticket ? { "set-cookie": opencodeTicketSetCookie(ticketAuth.ticket) } : {})
|
...(ticketAuth.ticket ? { "set-cookie": opencodeTicketSetCookie(ticketAuth.ticket) } : {})
|
||||||
@@ -748,7 +815,8 @@ async function proxyOpencodeRequest({ request, response, url, cloudApiBaseUrl, c
|
|||||||
response,
|
response,
|
||||||
body: request.method === "GET" || request.method === "HEAD" ? "" : body,
|
body: request.method === "GET" || request.method === "HEAD" ? "" : body,
|
||||||
timeoutMs: opencodeProxyTimeoutMs,
|
timeoutMs: opencodeProxyTimeoutMs,
|
||||||
extraResponseHeaders
|
extraResponseHeaders,
|
||||||
|
streamTransform
|
||||||
});
|
});
|
||||||
emitOpencodeProxySpanAsync({
|
emitOpencodeProxySpanAsync({
|
||||||
traceContext,
|
traceContext,
|
||||||
@@ -761,7 +829,8 @@ async function proxyOpencodeRequest({ request, response, url, cloudApiBaseUrl, c
|
|||||||
errorCode: proxyResult?.errorCode || "",
|
errorCode: proxyResult?.errorCode || "",
|
||||||
ticketAuth,
|
ticketAuth,
|
||||||
sessionStatusCode: session.statusCode || 0,
|
sessionStatusCode: session.statusCode || 0,
|
||||||
streaming: proxyResult?.streaming === true
|
streaming: proxyResult?.streaming === true,
|
||||||
|
streamTransformStats: proxyResult?.streamTransformStats
|
||||||
});
|
});
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
if (response.headersSent || response.writableEnded) {
|
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);
|
const endpoint = otelTracesEndpoint(process.env);
|
||||||
if (!endpoint || !traceContext?.traceId || !traceContext?.spanId) return;
|
if (!endpoint || !traceContext?.traceId || !traceContext?.spanId) return;
|
||||||
const endedAtMs = Date.now();
|
const endedAtMs = Date.now();
|
||||||
@@ -944,6 +1013,7 @@ function emitOpencodeProxySpanAsync({ traceContext, serviceId, request, url, tar
|
|||||||
"opencode.proxy.ticket_accepted": Boolean(ticketAuth?.ticket),
|
"opencode.proxy.ticket_accepted": Boolean(ticketAuth?.ticket),
|
||||||
"opencode.proxy.session_status_code": sessionStatusCode || undefined,
|
"opencode.proxy.session_status_code": sessionStatusCode || undefined,
|
||||||
"opencode.proxy.streaming": streaming === true,
|
"opencode.proxy.streaming": streaming === true,
|
||||||
|
...opencodeStreamTransformStatsAttributes(streamTransformStats),
|
||||||
valuesPrinted: false
|
valuesPrinted: false
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
@@ -961,6 +1031,18 @@ function emitOpencodeProxySpanAsync({ traceContext, serviceId, request, url, tar
|
|||||||
}, 0);
|
}, 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) {
|
function opencodeProxyRoute(pathname) {
|
||||||
const path = String(pathname || "/");
|
const path = String(pathname || "/");
|
||||||
if (path === "/global/event") return "/global/event";
|
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 () => {
|
test("cloud web OpenCode proxy emits bounded OTLP span", async () => {
|
||||||
const otelBodies = [];
|
const otelBodies = [];
|
||||||
const otel = createServer(async (request, response) => {
|
const otel = createServer(async (request, response) => {
|
||||||
|
|||||||
Reference in New Issue
Block a user