/* * 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"; import { defaultCodeAgentTraceStore } from "./code-agent-trace-store.ts"; import { authenticateWorkbenchRead } from "./server-workbench-read-http.ts"; import { workbenchKafkaDebugCapability } from "./workbench-kafka-debug-capability.ts"; import { completeWorkbenchKafkaDebugReplay, createWorkbenchKafkaDebugReplayTracker, markWorkbenchKafkaDebugReplayBarrierComplete, markWorkbenchKafkaDebugReplayConsumerReady, normalizeWorkbenchKafkaDebugOffsetRange, observeWorkbenchKafkaDebugReplayRecord, recordWorkbenchKafkaDebugReplayDelivery, workbenchKafkaDebugOffsetRangeWidth } from "./workbench-kafka-debug-replay-contract.ts"; const CONTRACT_VERSION = "workbench-debug-kafka-sse-v1"; const DEFAULT_STREAM = "hwlab"; const ALLOWED_STREAMS = new Set(["stdio", "agentrun", "hwlab", "hwlab-debug"]); export async function handleWorkbenchKafkaSseDebugHttp(request, response, url, options = {}) { const route = routeSuffix(url.pathname); try { const auth = await authenticateWorkbenchRead(request, response, options); if (!auth) return; if (auth.actor?.role !== "admin") { return sendJson(response, 403, debugError("admin_required", "Only admin users can access raw Workbench Kafka debug streams.")); } const stream = streamFromUrl(url); const isolatedDebug = stream === "hwlab-debug" ? isolatedDebugConfig(response, url, options.env ?? process.env, { requireReplay: route === "/events" }) : null; if (isolatedDebug === false) return; if (route === "") { if (request.method !== "GET") return methodNotAllowed(response, "GET"); return sendJson(response, 200, describeKafkaSseDebug(url, options.env ?? process.env, isolatedDebug)); } if (route === "/events") { if (request.method !== "GET") return methodNotAllowed(response, "GET"); return openKafkaDebugSse(request, response, url, options, isolatedDebug); } 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, isolatedDebug = null) { const stream = streamFromUrl(url); const requestedFilters = filtersFromUrl(url); const resolvedFilters = stream === "hwlab-debug" ? compactObject({ traceId: requestedFilters.traceId, replayId: requestedFilters.replayId }) : resolveKafkaDebugFilters(requestedFilters, { traceStore: defaultCodeAgentTraceStore }); return { ok: true, contractVersion: CONTRACT_VERSION, stream, topic: isolatedDebug?.topic ?? topicForStream(stream, env), debugIsolation: stream === "hwlab-debug", liveOnly: false, replay: stream === "hwlab-debug", replayLimit: isolatedDebug?.replayLimit ?? null, replayTimeoutMs: isolatedDebug?.replayTimeoutMs ?? null, groupIdPrefix: isolatedDebug?.groupIdPrefix ?? null, replayRequest: isolatedDebug?.replayRequest ?? null, filters: requestedFilters, resolvedFilters, eventsRoute: `/v1/workbench/debug/kafka-sse/events${url.search || ""}`, valuesPrinted: false }; } async function openKafkaDebugSse(request, response, url, options, isolatedDebug = null) { const env = options.env ?? process.env; const stream = streamFromUrl(url); const filters = filtersFromUrl(url); const isolatedReplay = stream === "hwlab-debug"; const resolvedFilters = isolatedReplay ? compactObject({ traceId: filters.traceId, replayId: filters.replayId }) : resolveKafkaDebugFilters(filters, options); const replayRequest = isolatedDebug?.replayRequest ?? replayRequestFromUrl(url); const correlatedReplay = isolatedReplay && replayRequest.mode === "correlated-v2"; const replayTracker = isolatedReplay ? createWorkbenchKafkaDebugReplayTracker({ mode: replayRequest.mode, replayId: resolvedFilters.replayId, traceId: resolvedFilters.traceId, topic: isolatedDebug.topic, offsetRange: replayRequest.offsetRange, producer: replayRequest.producer }) : null; 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(); let kafkaStream = null; let closed = false; let consumerReady = false; let replayFinished = false; let replayTimer = null; let replayCount = 0; let replayTerminalObserved = false; let barrierCompletionPending = false; const bufferedRecords = []; const close = () => { if (closed) return; closed = true; if (replayTimer) clearTimeout(replayTimer); void kafkaStream?.stop?.(); }; response.once("close", close); request.once?.("aborted", close); request.socket?.once?.("close", close); const finishReplay = (reason) => { if (!isolatedReplay || replayFinished || closed) return; replayFinished = true; if (replayTimer) clearTimeout(replayTimer); const replayResult = completeWorkbenchKafkaDebugReplay(replayTracker, { reason, terminalObserved: replayTerminalObserved }); writeSse(response, "hwlab.kafka.replay-complete", { ok: true, contractVersion: CONTRACT_VERSION, stream, topic: isolatedDebug.topic, traceId: resolvedFilters.traceId, replayComplete: true, replayMode: replayResult.mode, replayId: replayResult.replayId, phase: replayResult.phase, code: replayResult.code, result: replayResult, reason, count: replayCount, counts: replayResult.counts, rejectedByReason: replayResult.rejectedByReason, offsetRange: replayResult.offsetRange, sourceLineage: replayResult.sourceLineage, limit: isolatedDebug.replayLimit, timeoutMs: isolatedDebug.replayTimeoutMs, terminalObserved: replayTerminalObserved, barrierComplete: replayResult.barrier.completed, stopBoundary: "sse-end", clientCountsAvailable: false, valuesPrinted: false }); if (!response.writableEnded) response.end(); }; const writeRecord = (record) => { if (closed || replayFinished) return; const delivered = 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 }); if (!isolatedReplay) return; recordWorkbenchKafkaDebugReplayDelivery(replayTracker, record, delivered); replayCount = replayTracker.counts.delivered; replayTerminalObserved ||= debugRecordIsTerminal(record); if (!correlatedReplay && replayTerminalObserved) finishReplay("terminal"); else if (!correlatedReplay && replayCount >= isolatedDebug.replayLimit) finishReplay("limit"); }; if (isolatedReplay) { replayTimer = setTimeout(() => finishReplay("timeout"), isolatedDebug.replayTimeoutMs); replayTimer.unref?.(); } try { kafkaStream = await openKafkaEventStream({ env, stream, topic: isolatedDebug?.topic ?? null, groupIdPrefix: isolatedDebug?.groupIdPrefix ?? null, fromBeginning: isolatedReplay ? !correlatedReplay : url.searchParams.get("fromBeginning") === "1" || url.searchParams.get("fromBeginning") === "true", ...resolvedFilters, offsetRange: replayRequest.offsetRange, kafkaFactory: options.kafkaFactory, onRecord: isolatedReplay ? async (observation) => { observeWorkbenchKafkaDebugReplayRecord(replayTracker, observation); } : null, onBarrierComplete: correlatedReplay ? async () => { markWorkbenchKafkaDebugReplayBarrierComplete(replayTracker); if (consumerReady) finishReplay("barrier"); else barrierCompletionPending = true; } : null, onEvent: async (record) => { if (isolatedReplay && !consumerReady) { if (bufferedRecords.length < isolatedDebug.replayLimit) bufferedRecords.push(record); return; } writeRecord(record); } }); if (closed) { await kafkaStream.stop?.(); return; } if (isolatedReplay) markWorkbenchKafkaDebugReplayConsumerReady(replayTracker, { groupId: kafkaStream.groupId, seekApplied: kafkaStream.seekApplied }); writeSse(response, "hwlab.kafka.connected", { ok: true, contractVersion: CONTRACT_VERSION, consumerReady: true, groupId: kafkaStream.groupId, stream, topic: kafkaStream.topic ?? isolatedDebug?.topic ?? topicForStream(stream, env), debugIsolation: isolatedReplay, deliverySemantics: isolatedReplay ? "debug-replay" : "diagnostic", liveOnly: false, replay: isolatedReplay, replayLimit: isolatedDebug?.replayLimit ?? null, replayTimeoutMs: isolatedDebug?.replayTimeoutMs ?? null, replayId: replayTracker?.replayId ?? null, replayMode: replayTracker?.mode ?? null, offsetRange: replayTracker?.requestedOffsetRange ?? null, replayContractVersion: replayTracker?.mode === "correlated-v2" ? "workbench-kafka-debug-replay-v2" : CONTRACT_VERSION, seekApplied: kafkaStream.seekApplied === true, filters, resolvedFilters, serverSentAt: new Date().toISOString(), valuesPrinted: false }); consumerReady = true; for (const record of bufferedRecords.splice(0)) { if (replayFinished || closed) break; writeRecord(record); } if (barrierCompletionPending && !replayFinished && !closed) finishReplay("barrier"); } catch (error) { if (replayTimer) clearTimeout(replayTimer); writeSse(response, "hwlab.kafka.error", { ok: false, error: errorMessagePayload(error), valuesPrinted: false }); if (isolatedReplay) finishReplay("error"); } } function resolveKafkaDebugFilters(filters = {}, options = {}) { const traceId = textValue(filters.traceId); if (!traceId) return filters; const traceStore = options.traceStore ?? defaultCodeAgentTraceStore; const snapshot = typeof traceStore?.snapshot === "function" ? traceStore.snapshot(traceId) : null; const resolved = collectTraceLinkedIds(snapshot); const hasResolvedKafkaKey = Boolean(resolved.runId || resolved.commandId || resolved.sessionId); if (!hasResolvedKafkaKey) return filters; const resolvedSessionId = resolved.runId || resolved.commandId ? null : resolved.sessionId; return compactObject({ sessionId: filters.sessionId || resolvedSessionId, runId: filters.runId || resolved.runId, commandId: filters.commandId || resolved.commandId }); } function collectTraceLinkedIds(snapshot) { const out = { sessionId: null, runId: null, commandId: null }; const events = Array.isArray(snapshot?.events) ? snapshot.events : []; for (const event of events) collectIdsFromRecord(event, out); collectIdsFromRecord(snapshot?.lastEvent, out); return out; } function collectIdsFromRecord(record, out) { const value = record && typeof record === "object" && !Array.isArray(record) ? record : {}; const agentRun = value.agentRun && typeof value.agentRun === "object" && !Array.isArray(value.agentRun) ? value.agentRun : {}; const payload = value.payload && typeof value.payload === "object" && !Array.isArray(value.payload) ? value.payload : {}; out.sessionId ||= textValue(value.sessionId ?? value.sourceSessionId ?? agentRun.sessionId ?? payload.sessionId); out.runId ||= textValue(value.runId ?? value.sourceRunId ?? agentRun.runId ?? payload.runId); out.commandId ||= textValue(value.commandId ?? value.sourceCommandId ?? agentRun.commandId ?? payload.commandId); } 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")), replayId: safeReplayId(url.searchParams.get("replayId") || url.searchParams.get("replay-id")) }); } function topicForStream(stream, env) { if (stream === "stdio") return textValue(env.HWLAB_KAFKA_STDIO_TOPIC ?? env.AGENTRUN_KAFKA_STDIO_TOPIC) || "codex-stdio.raw.v1"; if (stream === "agentrun") return textValue(env.HWLAB_KAFKA_AGENTRUN_EVENT_TOPIC) || "agentrun.event.v1"; if (stream === "hwlab-debug") return workbenchKafkaDebugCapability(env).topic; return textValue(env.HWLAB_KAFKA_EVENT_TOPIC) || "hwlab.event.v1"; } function isolatedDebugConfig(response, url, env, options = {}) { let config; try { config = workbenchKafkaDebugCapability(env); } catch (error) { sendJson(response, Number(error?.statusCode) || 503, debugError(error?.code || "workbench_kafka_debug_capability_invalid", error instanceof Error ? error.message : "Workbench isolated Kafka debug capability is invalid.")); return false; } if (!config.enabled) { sendJson(response, 503, debugError("workbench_kafka_debug_disabled", "Workbench isolated Kafka debug capability is disabled.")); return false; } if (options.requireReplay && !safeId(url.searchParams.get("traceId") || url.searchParams.get("trace-id"))) { sendJson(response, 400, debugError("workbench_kafka_debug_trace_required", "Isolated Workbench Kafka replay requires traceId.")); return false; } if (options.requireReplay && !["1", "true"].includes(String(url.searchParams.get("fromBeginning") || "").toLowerCase())) { sendJson(response, 400, debugError("workbench_kafka_debug_replay_required", "Isolated Workbench Kafka debug requires explicit fromBeginning=true.")); return false; } const rawReplayId = url.searchParams.get("replayId") || url.searchParams.get("replay-id"); if (textValue(rawReplayId) && !safeReplayId(rawReplayId)) { sendJson(response, 400, debugError("workbench_kafka_debug_replay_id_invalid", "Replay correlation requires replayId to start with rpl_ and contain only safe identifier characters.")); return false; } const replayRequest = replayRequestFromUrl(url); if (options.requireReplay && replayRequest.offsetRangeInvalid) { sendJson(response, 400, debugError("workbench_kafka_debug_offset_range_invalid", "Replay offset barrier requires partition, firstOffset, and lastOffset with firstOffset <= lastOffset.")); return false; } if (options.requireReplay && replayRequest.mode === "correlated-v2") { if (!replayRequest.replayId || !replayRequest.offsetRange) { sendJson(response, 400, debugError("workbench_kafka_debug_v2_correlation_incomplete", "Correlated V2 replay requires replayId and partition/firstOffset/lastOffset.")); return false; } if (replayRequest.producerMetadataInvalid) { sendJson(response, 400, debugError("workbench_kafka_debug_v2_producer_metadata_invalid", "Optional producerInvoked, sourceMatched, and publishedCount must use valid boolean/integer values.")); return false; } if (replayRequest.producer.invoked === false) { sendJson(response, 409, debugError("workbench_kafka_debug_producer_not_invoked", "Correlated V2 replay cannot use producerInvoked=false.")); return false; } if (replayRequest.producer.sourceMatched === false) { sendJson(response, 409, debugError("workbench_kafka_debug_source_trace_missing", "Correlated V2 replay cannot start because the producer did not match the requested source trace.")); return false; } const rangeWidth = workbenchKafkaDebugOffsetRangeWidth(replayRequest.offsetRange); if (!rangeWidth || (replayRequest.producer.publishedCount !== null && replayRequest.producer.publishedCount !== rangeWidth)) { sendJson(response, 400, debugError("workbench_kafka_debug_producer_cardinality_mismatch", "publishedCount must equal the requested offset range width.")); return false; } if (rangeWidth > config.replayLimit) { sendJson(response, 400, debugError("workbench_kafka_debug_replay_limit_exceeded", "The correlated replay batch exceeds the YAML-owned replay limit.")); return false; } } return { ...config, replayRequest }; } function replayRequestFromUrl(url) { const rawReplayId = url.searchParams.get("replayId") || url.searchParams.get("replay-id"); const rawProducerInvoked = url.searchParams.get("producerInvoked") || url.searchParams.get("producer-invoked"); const rawSourceMatched = url.searchParams.get("sourceMatched") || url.searchParams.get("source-matched"); const rawPublishedCount = url.searchParams.get("publishedCount") || url.searchParams.get("published-count"); const rawRange = { partition: url.searchParams.get("partition"), firstOffset: url.searchParams.get("firstOffset") || url.searchParams.get("first-offset") || url.searchParams.get("startOffset") || url.searchParams.get("start-offset"), lastOffset: url.searchParams.get("lastOffset") || url.searchParams.get("last-offset") || url.searchParams.get("endOffset") || url.searchParams.get("end-offset") }; const offsetRangeProvided = Object.values(rawRange).some((entry) => textValue(entry)); const offsetRange = normalizeWorkbenchKafkaDebugOffsetRange(rawRange); const correlationProvided = Boolean(textValue(rawReplayId) || offsetRangeProvided || textValue(rawProducerInvoked) || textValue(rawSourceMatched) || textValue(rawPublishedCount)); const producerInvoked = optionalBoolean(rawProducerInvoked); const sourceMatched = optionalBoolean(rawSourceMatched); const publishedCount = optionalNonNegativeInteger(rawPublishedCount); return { mode: correlationProvided ? "correlated-v2" : "trace-only-v1", replayId: safeReplayId(rawReplayId), offsetRange, offsetRangeInvalid: offsetRangeProvided && !offsetRange, producerMetadataInvalid: Boolean( (textValue(rawProducerInvoked) && producerInvoked === null) || (textValue(rawSourceMatched) && sourceMatched === null) || (textValue(rawPublishedCount) && publishedCount === null) ), producer: { invoked: producerInvoked, sourceMatched, publishedCount }, valuesPrinted: false }; } 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 debugRecordIsTerminal(record) { const value = record?.value && typeof record.value === "object" && !Array.isArray(record.value) ? record.value : {}; const event = value.event && typeof value.event === "object" && !Array.isArray(value.event) ? value.event : {}; const eventType = textValue(value.eventType ?? event.eventType ?? event.type); return event.terminal === true || eventType === "terminal" || eventType === "result"; } function safeId(value) { const text = textValue(value); return text && /^[A-Za-z0-9_.:-]{3,220}$/u.test(text) ? text : null; } function safeReplayId(value) { const text = safeId(value); return text && /^rpl_[A-Za-z0-9_.:-]+$/u.test(text) ? text : null; } function optionalBoolean(value) { const text = textValue(value)?.toLowerCase(); if (text === "true" || text === "1") return true; if (text === "false" || text === "0") return false; return null; } function optionalNonNegativeInteger(value) { if (!textValue(value)) return null; const parsed = Number(value); return Number.isSafeInteger(parsed) && parsed >= 0 ? parsed : 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; }