/* * SPEC: PJ2026-010401080313 Workbench实时权威 draft-2026-07-09-p1-single-step-debug. * Responsibility: debug-only in-memory fake Workbench SSE queues for Cloud Web single-step reducer/UI inspection. */ import { readBody, sendJson } from "./server-http-utils.ts"; const CONTRACT_VERSION = "workbench-debug-fake-sse-v1"; const EVENT_CONTRACT_VERSION = "workbench-events-v1"; const REALTIME_AUTHORITY = "workbench-realtime-authority-v2"; const DEFAULT_QUEUE_ID = "trace-card"; const BODY_LIMIT_BYTES = 256 * 1024; const queues = new Map(); const clientsByQueueId = new Map(); export async function handleWorkbenchDebugFakeSseHttp(request, response, url, options = {}) { const route = routeSuffix(url.pathname); try { if (route === "" || route === "/sequences") { if (request.method !== "GET") return methodNotAllowed(response, "GET"); return sendJson(response, 200, describePayload(queueIdFromUrl(url))); } if (route === "/events") { if (request.method !== "GET") return methodNotAllowed(response, "GET"); return openDebugSse(request, response, url); } if (route === "/reset") { if (request.method !== "POST") return methodNotAllowed(response, "POST"); const body = await readJsonObject(request); if (!body.ok) return sendJson(response, 400, debugError("invalid_json", body.message)); const queue = resetQueue(queueIdFromBodyOrUrl(body.value, url), String(body.value.sequenceId ?? "trace-card-basic")); return sendJson(response, 200, { ok: true, contractVersion: CONTRACT_VERSION, queue: describeQueue(queue), sequences: describeSequences() }); } if (route === "/append") { if (request.method !== "POST") return methodNotAllowed(response, "POST"); const body = await readJsonObject(request); if (!body.ok) return sendJson(response, 400, debugError("invalid_json", body.message)); const queue = ensureQueue(queueIdFromBodyOrUrl(body.value, url)); const events = appendEventsFromBody(body.value); if (events.length === 0) return sendJson(response, 400, debugError("fake_sse_events_required", "events must contain at least one object.")); queue.events.push(...events.map((event, index) => normalizeRealtimeEvent(event, queue, queue.events.length + index + 1))); queue.updatedAt = new Date().toISOString(); return sendJson(response, 200, { ok: true, contractVersion: CONTRACT_VERSION, appended: events.length, queue: describeQueue(queue) }); } if (route === "/next") { if (request.method !== "POST") return methodNotAllowed(response, "POST"); const body = await readOptionalJsonObject(request); const queue = ensureQueue(queueIdFromBodyOrUrl(body.value ?? {}, url)); const delivered = deliverNext(queue); return sendJson(response, 200, { ok: true, contractVersion: CONTRACT_VERSION, delivered, queue: describeQueue(queue) }); } if (route === "/run") { if (request.method !== "POST") return methodNotAllowed(response, "POST"); const body = await readOptionalJsonObject(request); const queue = ensureQueue(queueIdFromBodyOrUrl(body.value ?? {}, url)); const delivered = deliverAll(queue); return sendJson(response, 200, { ok: true, contractVersion: CONTRACT_VERSION, delivered, queue: describeQueue(queue) }); } return sendJson(response, 404, debugError("workbench_debug_fake_sse_route_not_found", "Workbench debug fake SSE route is not implemented.", { route })); } catch (error) { options.logger?.warn?.({ event: "workbench_debug_fake_sse_failed", route, errorName: error?.name ?? "Error", message: error instanceof Error ? error.message : String(error ?? "unknown"), valuesRedacted: true }); return sendJson(response, 500, debugError("workbench_debug_fake_sse_failed", "Workbench debug fake SSE request failed.")); } } function openDebugSse(request, response, url) { const queue = ensureQueue(queueIdFromUrl(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(); const clients = clientsForQueue(queue.queueId); clients.add(response); writeSse(response, "workbench.connected", { type: "connected", contractVersion: EVENT_CONTRACT_VERSION, realtimeAuthority: REALTIME_AUTHORITY, queue: describeQueue(queue), serverSentAt: new Date().toISOString() }); request.on("close", () => clients.delete(response)); } function deliverNext(queue) { const event = queue.events[queue.cursor] ?? null; if (!event) return { ok: false, reason: "queue-empty", deliveredTo: 0 }; queue.cursor += 1; queue.updatedAt = new Date().toISOString(); return deliverEvent(queue.queueId, event); } function deliverAll(queue) { const results = []; while (queue.cursor < queue.events.length) results.push(deliverNext(queue)); return { ok: true, count: results.length, deliveredTo: results.reduce((sum, item) => sum + Number(item.deliveredTo ?? 0), 0) }; } function deliverEvent(queueId, event) { const clients = clientsForQueue(queueId); const eventName = eventNameFor(event); let deliveredTo = 0; for (const response of [...clients]) { if (response.destroyed || response.writableEnded) { clients.delete(response); continue; } if (writeSse(response, eventName, event)) deliveredTo += 1; } return { ok: true, eventName, eventType: event.type ?? null, traceId: event.traceId ?? null, deliveredTo }; } function writeSse(response, eventName, payload) { if (response.destroyed || response.writableEnded) return false; try { response.write(`event: ${eventName}\n`); const eventId = sseEventId(payload); if (eventId) response.write(`id: ${eventId}\n`); response.write(`data: ${JSON.stringify({ serverSentAt: new Date().toISOString(), ...payload })}\n\n`); return true; } catch { return false; } } function ensureQueue(queueId = DEFAULT_QUEUE_ID) { const id = safeQueueId(queueId); const existing = queues.get(id); if (existing) return existing; return resetQueue(id, "trace-card-basic"); } function resetQueue(queueId = DEFAULT_QUEUE_ID, sequenceId = "trace-card-basic") { const id = safeQueueId(queueId); const sequence = builtinSequences().find((item) => item.sequenceId === sequenceId) ?? builtinSequences()[0]; const now = new Date().toISOString(); const queue = { queueId: id, sequenceId: sequence.sequenceId, cursor: 0, events: sequence.events.map((event, index) => normalizeRealtimeEvent(event, { queueId: id }, index + 1)), createdAt: now, updatedAt: now }; queues.set(id, queue); return queue; } function normalizeRealtimeEvent(event, queue, fallbackVersion) { const source = recordValue(event) ?? {}; const sessionId = textValue(source.sessionId) || "ses_debug_fake_sse"; const threadId = textValue(source.threadId) || "thr_debug_fake_sse"; const traceId = textValue(source.traceId) || "trc_debug_fake_sse"; const type = textValue(source.type) || "trace.event"; const entity = recordValue(source.entity) ?? defaultEntity(type, source, fallbackVersion); return { sessionId, threadId, traceId, ...source, contractVersion: EVENT_CONTRACT_VERSION, realtimeAuthority: REALTIME_AUTHORITY, sessionId, threadId, traceId, type, entity, cursor: recordValue(source.cursor) ?? { traceSeq: finiteNumber(source.traceSeq ?? source.event?.projectedSeq ?? fallbackVersion), outboxSeq: finiteNumber(source.outboxSeq ?? fallbackVersion) }, debug: { source: "workbench-debug-fake-sse", queueId: queue.queueId, valuesRedacted: true } }; } function defaultEntity(type, source, version) { const traceId = textValue(source.traceId) || "trc_debug_fake_sse"; const messageId = textValue(source.message?.messageId ?? source.message?.id) || "msg_debug_agent"; const family = type === "message.snapshot" ? "messages" : type === "turn.snapshot" ? "turns" : "traceEvents"; const id = family === "messages" ? messageId : family === "turns" ? textValue(source.turn?.turnId ?? source.turn?.traceId) || traceId : `${traceId}:${version}`; return { family, id, version, outboxSeq: version, traceSeq: version, projectionRevision: `debug-${version}`, authority: REALTIME_AUTHORITY }; } function builtinSequences() { const sessionId = "ses_debug_fake_sse"; const threadId = "thr_debug_fake_sse"; const traceId = "trc_debug_fake_sse"; const t0 = "2026-07-09T08:00:00.000Z"; const events = [ traceEvent({ projectedSeq: 1, createdAt: "2026-07-09T08:00:01.000Z", label: "assistant:message", type: "assistant_message", status: "running", message: "读取 Workbench 组件和 Trace 渲染上下文。", elapsedMs: 1000 }), traceEvent({ projectedSeq: 2, createdAt: "2026-07-09T08:00:03.000Z", label: "agentrun:tool:completed", type: "tool_call", status: "completed", toolName: "commandExecution", itemId: "tool_debug_rg", command: "rg -n \"TraceTimeline|message-card\" web/hwlab-cloud-web/src", stdoutSummary: "TraceTimeline.vue\\nConversationPanel.vue\\nWorkbenchMessageCard.vue", exitCode: 0, elapsedMs: 3000 }), traceEvent({ projectedSeq: 3, createdAt: "2026-07-09T08:00:05.000Z", label: "assistant:message", type: "assistant_message", status: "running", message: "Trace 卡片已收到 tool 事件,继续等待 terminal snapshot。", elapsedMs: 5000 }), traceEvent({ projectedSeq: 4, createdAt: "2026-07-09T08:00:07.000Z", label: "assistant:completed", type: "assistant_message", status: "completed", message: "Fake SSE Trace 卡片渲染完成。", final: true, terminal: true, replyAuthority: true, elapsedMs: 7000 }), traceEvent({ projectedSeq: 5, createdAt: "2026-07-09T08:00:07.500Z", label: "turn:completed", type: "completion", status: "completed", terminal: true, elapsedMs: 7500 }) ]; return [ { sequenceId: "trace-card-basic", label: "Trace card basic", eventCount: 6, events: [ { type: "message.snapshot", sessionId, threadId, traceId, message: agentMessage({ sessionId, threadId, traceId, status: "running", createdAt: t0, events: [], eventCount: 0 }) }, ...events.slice(0, 3).map((event, index) => ({ type: "trace.event", sessionId, threadId, traceId, event, snapshot: traceSnapshot({ sessionId, threadId, traceId, status: "running", events: events.slice(0, index + 1), eventCount: index + 1, startedAt: t0 }) })), { type: "trace.event", sessionId, threadId, traceId, event: events[3], snapshot: traceSnapshot({ sessionId, threadId, traceId, status: "completed", events: events.slice(0, 4), eventCount: 4, startedAt: t0, finishedAt: "2026-07-09T08:00:07.000Z", durationMs: 7000, finalText: "Fake SSE Trace 卡片渲染完成。" }) }, { type: "message.snapshot", sessionId, threadId, traceId, message: agentMessage({ sessionId, threadId, traceId, status: "completed", createdAt: t0, updatedAt: "2026-07-09T08:00:07.000Z", text: "Fake SSE Trace 卡片渲染完成。", events, eventCount: events.length, finishedAt: "2026-07-09T08:00:07.000Z", durationMs: 7000, finalText: "Fake SSE Trace 卡片渲染完成。" }) } ] } ]; } function traceEvent(input) { return { traceId: "trc_debug_fake_sse", source: "debug-fake-sse", sourceSeq: input.projectedSeq, ...input }; } function agentMessage(input) { return { id: "msg_debug_agent", messageId: "msg_debug_agent", role: "agent", title: "Code Agent", text: input.text ?? "", status: input.status, createdAt: input.createdAt, updatedAt: input.updatedAt ?? input.createdAt, sessionId: input.sessionId, threadId: input.threadId, traceId: input.traceId, turnId: input.traceId, timing: { startedAt: input.createdAt, lastEventAt: input.updatedAt ?? input.createdAt, finishedAt: input.finishedAt ?? null, durationMs: input.durationMs ?? null, valuesRedacted: true }, runnerTrace: traceSnapshot({ sessionId: input.sessionId, threadId: input.threadId, traceId: input.traceId, status: input.status, events: input.events, eventCount: input.eventCount, startedAt: input.createdAt, lastEventAt: input.updatedAt ?? input.createdAt, finishedAt: input.finishedAt ?? null, durationMs: input.durationMs ?? null, finalText: input.finalText }) }; } function traceSnapshot(input) { return { traceId: input.traceId, sessionId: input.sessionId, threadId: input.threadId, status: input.status, events: input.events, eventCount: input.eventCount ?? input.events?.length ?? 0, fullTraceLoaded: input.status === "completed", hasMore: false, startedAt: input.startedAt, lastEventAt: input.lastEventAt ?? input.events?.at(-1)?.createdAt ?? input.startedAt, finishedAt: input.finishedAt ?? null, durationMs: input.durationMs ?? null, timing: { startedAt: input.startedAt, lastEventAt: input.lastEventAt ?? input.events?.at(-1)?.createdAt ?? input.startedAt, finishedAt: input.finishedAt ?? null, durationMs: input.durationMs ?? null, valuesRedacted: true }, finalResponse: input.finalText ? { text: input.finalText, valuesRedacted: true } : undefined, projection: { projectionStatus: "caught-up", projectionHealth: "caught-up", lastProjectedSeq: input.eventCount ?? input.events?.length ?? 0, updatedAt: input.lastEventAt ?? input.events?.at(-1)?.createdAt ?? input.startedAt, valuesRedacted: true } }; } function appendEventsFromBody(body) { const raw = Array.isArray(body.events) ? body.events : body.event ? [body.event] : []; return raw.filter((event) => event && typeof event === "object" && !Array.isArray(event)); } async function readJsonObject(request) { const body = await readBody(request, BODY_LIMIT_BYTES); try { const value = body ? JSON.parse(body) : {}; if (!value || typeof value !== "object" || Array.isArray(value)) return { ok: false, message: "body must be a JSON object" }; return { ok: true, value }; } catch (error) { return { ok: false, message: error instanceof Error ? error.message : "Invalid JSON body" }; } } async function readOptionalJsonObject(request) { if (Number(request.headers?.["content-length"] ?? 0) <= 0) return { ok: true, value: {} }; return readJsonObject(request); } function describePayload(queueId) { return { ok: true, contractVersion: CONTRACT_VERSION, sequences: describeSequences(), queue: describeQueue(ensureQueue(queueId)) }; } function describeSequences() { return builtinSequences().map((sequence) => ({ sequenceId: sequence.sequenceId, label: sequence.label, eventCount: sequence.events.length })); } function describeQueue(queue) { return { queueId: queue.queueId, sequenceId: queue.sequenceId, cursor: queue.cursor, eventCount: queue.events.length, remaining: Math.max(0, queue.events.length - queue.cursor), connectedClients: clientsForQueue(queue.queueId).size, updatedAt: queue.updatedAt }; } function clientsForQueue(queueId) { const id = safeQueueId(queueId); let clients = clientsByQueueId.get(id); if (!clients) { clients = new Set(); clientsByQueueId.set(id, clients); } return clients; } function queueIdFromBodyOrUrl(body, url) { return safeQueueId(body.queueId ?? url.searchParams.get("queueId")); } function queueIdFromUrl(url) { return safeQueueId(url.searchParams.get("queueId")); } function safeQueueId(value) { const text = String(value ?? "").trim(); return /^[A-Za-z0-9_.:-]{3,80}$/u.test(text) ? text : DEFAULT_QUEUE_ID; } function eventNameFor(event) { switch (event?.type) { case "trace.snapshot": return "workbench.trace.snapshot"; case "trace.event": return "workbench.trace.event"; case "message.snapshot": return "workbench.message.snapshot"; case "turn.snapshot": return "workbench.turn.snapshot"; case "trace.unavailable": return "workbench.trace.unavailable"; case "error": return "workbench.error"; default: return "message"; } } function sseEventId(payload) { const raw = payload?.cursor?.outboxSeq ?? payload?.outboxSeq ?? payload?.entity?.outboxSeq; const text = textValue(raw); if (!text || /[\r\n]/u.test(text)) return null; return text.slice(0, 128); } function routeSuffix(pathname) { return pathname.replace(/^\/v1\/workbench\/debug\/fake-sse/u, ""); } function methodNotAllowed(response, allowed) { return sendJson(response, 405, debugError("method_not_allowed", `Use ${allowed} for this Workbench debug route.`)); } function debugError(code, message, extra = {}) { return { ok: false, contractVersion: CONTRACT_VERSION, error: { code, message, layer: "workbench-debug-fake-sse", valuesRedacted: true }, ...extra }; } function recordValue(value) { return value && typeof value === "object" && !Array.isArray(value) ? value : null; } function textValue(value) { const text = String(value ?? "").trim(); return text || null; } function finiteNumber(value) { const number = Number(value); return Number.isFinite(number) && number >= 0 ? Math.trunc(number) : null; }