439 lines
17 KiB
TypeScript
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;
|
|
}
|