/* * SPEC: PJ2026-0104010803 Workbench事件流可见性 draft-2026-07-09-p0-kafka-authority. * Responsibility: debug-only SSE passthrough for HWLAB Kafka event envelopes. */ import { sendJson } from "./server-http-utils.ts"; import { openKafkaEventStream } from "./kafka-event-bridge.ts"; const CONTRACT_VERSION = "workbench-debug-kafka-sse-v1"; const DEFAULT_STREAM = "hwlab"; const ALLOWED_STREAMS = new Set(["hwlab", "agentrun"]); export async function handleWorkbenchKafkaSseDebugHttp(request, response, url, options = {}) { const route = routeSuffix(url.pathname); try { if (route === "") { if (request.method !== "GET") return methodNotAllowed(response, "GET"); return sendJson(response, 200, describeKafkaSseDebug(url, options.env ?? process.env)); } if (route === "/events") { if (request.method !== "GET") return methodNotAllowed(response, "GET"); return openKafkaDebugSse(request, response, url, options); } return sendJson(response, 404, debugError("workbench_debug_kafka_sse_route_not_found", "Workbench debug Kafka SSE route is not implemented.", { route })); } catch (error) { options.logger?.warn?.({ event: "workbench_debug_kafka_sse_failed", route, errorName: error?.name ?? "Error", message: error instanceof Error ? error.message : String(error ?? "unknown"), valuesRedacted: true }); if (response.headersSent) return writeSse(response, "hwlab.kafka.error", { ok: false, error: errorMessagePayload(error) }); return sendJson(response, 500, debugError("workbench_debug_kafka_sse_failed", "Workbench debug Kafka SSE request failed.")); } } function describeKafkaSseDebug(url, env) { const stream = streamFromUrl(url); return { ok: true, contractVersion: CONTRACT_VERSION, stream, topic: topicForStream(stream, env), filters: filtersFromUrl(url), eventsRoute: `/v1/workbench/debug/kafka-sse/events${url.search || ""}`, valuesPrinted: false }; } async function openKafkaDebugSse(request, response, url, options) { const env = options.env ?? process.env; const stream = streamFromUrl(url); const filters = filtersFromUrl(url); response.writeHead(200, { "content-type": "text/event-stream; charset=utf-8", "cache-control": "no-store, no-transform", connection: "keep-alive", "x-accel-buffering": "no", "x-content-type-options": "nosniff" }); if (typeof response.flushHeaders === "function") response.flushHeaders(); writeSse(response, "hwlab.kafka.connected", { ok: true, contractVersion: CONTRACT_VERSION, stream, topic: topicForStream(stream, env), filters, serverSentAt: new Date().toISOString(), valuesPrinted: false }); let kafkaStream = null; let closed = false; const close = () => { if (closed) return; closed = true; void kafkaStream?.stop?.(); }; request.on("close", close); try { kafkaStream = await openKafkaEventStream({ env, stream, fromBeginning: url.searchParams.get("fromBeginning") === "1" || url.searchParams.get("fromBeginning") === "true", ...filters, kafkaFactory: options.kafkaFactory, onEvent: async (record) => { writeSse(response, "hwlab.kafka.event", { ok: true, contractVersion: CONTRACT_VERSION, stream, topic: record.topic, partition: record.partition, offset: record.offset, key: record.key, timestamp: record.timestamp, value: record.value, serverSentAt: new Date().toISOString(), valuesPrinted: false }); }, onError: (error) => writeSse(response, "hwlab.kafka.error", { ok: false, error: errorMessagePayload(error), valuesPrinted: false }) }); } catch (error) { writeSse(response, "hwlab.kafka.error", { ok: false, error: errorMessagePayload(error), valuesPrinted: false }); } } function writeSse(response, eventName, payload) { if (response.destroyed || response.writableEnded) return false; try { response.write(`event: ${eventName}\n`); const id = sseEventId(payload); if (id) response.write(`id: ${id}\n`); response.write(`data: ${JSON.stringify(payload)}\n\n`); return true; } catch { return false; } } function routeSuffix(pathname) { return String(pathname || "").replace(/^\/v1\/workbench\/debug\/kafka-sse/u, ""); } function streamFromUrl(url) { const stream = String(url.searchParams.get("stream") || DEFAULT_STREAM).trim(); return ALLOWED_STREAMS.has(stream) ? stream : DEFAULT_STREAM; } function filtersFromUrl(url) { return compactObject({ traceId: safeId(url.searchParams.get("traceId") || url.searchParams.get("trace-id")), sessionId: safeId(url.searchParams.get("sessionId") || url.searchParams.get("session-id")), runId: safeId(url.searchParams.get("runId") || url.searchParams.get("run-id")), commandId: safeId(url.searchParams.get("commandId") || url.searchParams.get("command-id")) }); } function topicForStream(stream, env) { if (stream === "agentrun") return textValue(env.HWLAB_KAFKA_AGENTRUN_EVENT_TOPIC) || "agentrun.event.v1"; return textValue(env.HWLAB_KAFKA_EVENT_TOPIC) || "hwlab.event.v1"; } function methodNotAllowed(response, allow) { response.setHeader("allow", allow); return sendJson(response, 405, debugError("method_not_allowed", `Method not allowed; expected ${allow}.`)); } function debugError(code, message, extra = {}) { return { ok: false, error: { code, message, ...extra }, valuesPrinted: false }; } function errorMessagePayload(error) { return { code: "kafka_sse_error", message: error instanceof Error ? error.message : String(error ?? "unknown") }; } function safeId(value) { const text = textValue(value); return text && /^[A-Za-z0-9_.:-]{3,220}$/u.test(text) ? text : null; } function textValue(value) { const text = typeof value === "string" ? value.trim() : value === null || value === undefined ? "" : String(value).trim(); return text.length > 0 ? text : null; } function compactObject(value) { return Object.fromEntries(Object.entries(value).filter(([, entry]) => textValue(entry))); } function sseEventId(payload) { const topic = textValue(payload?.topic); const partition = textValue(payload?.partition); const offset = textValue(payload?.offset); return topic && partition && offset ? `${topic}:${partition}:${offset}` : null; }