Files
pikasTech-HWLAB/internal/cloud/workbench-debug-fake-sse.ts
2026-07-09 15:46:28 +02:00

439 lines
17 KiB
TypeScript

/*
* 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;
}