Merge pull request #2223 from pikasTech/fix/1108-dsflash-projection

fix(workbench): stabilize dsflash read model pagination
This commit is contained in:
Lyon
2026-06-27 13:41:34 +08:00
committed by GitHub
2 changed files with 247 additions and 42 deletions
+136 -16
View File
@@ -574,6 +574,122 @@ test("workbench read model exposes session, messages, turn, and trace without wr
}
});
test("workbench session messages use turn timeline instead of clustered projection writes", async () => {
const traceIds = Array.from({ length: 5 }, (_, index) => `trc_workbench_timeline_cluster_${index + 1}`);
const session = {
id: "ses_workbench_timeline_cluster",
projectId: "prj_hwpod_workbench",
agentId: "hwlab-code-agent",
status: "completed",
startedAt: "2026-06-27T00:00:00.000Z",
endedAt: "2026-06-27T00:10:00.000Z",
ownerUserId: ACTOR.id,
conversationId: "cnv_workbench_timeline_cluster",
threadId: "thread-workbench-timeline-cluster",
lastTraceId: traceIds.at(-1),
updatedAt: "2026-06-27T00:10:00.000Z",
session: {
sessionStatus: "completed",
messages: traceIds.flatMap((traceId, index) => {
const turn = index + 1;
return [
{ role: "user", text: `sentinel-${turn}`, traceId, turnId: traceId, projectedSeq: turn, sourceSeq: turn, status: "sent", createdAt: `2026-06-27T00:0${index}:00.000Z` },
{ role: "agent", text: `ok-${turn}`, traceId, turnId: traceId, projectedSeq: 100 + turn, sourceSeq: 100 + turn, status: "completed", createdAt: `2026-06-27T00:0${index}:30.000Z` }
];
}),
valuesRedacted: true,
secretMaterialStored: false
}
};
const runtimeStore = createDurableFactsRuntimeStore({ sessions: [{ session, status: "completed", lastProjectedSeq: 105 }] });
const accessController = {
async ensureBootstrap() {},
async authenticate() { return { ok: true, actor: ACTOR, session: { id: "uss_workbench_reader" } }; }
};
const server = createCloudApiServer({ accessController, runtimeStore });
await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve));
try {
const { port } = server.address();
const messages = await getJson(port, `/v1/workbench/sessions/${encodeURIComponent(session.id)}/messages?limit=20`);
assert.equal(messages.status, 200);
assert.deepEqual(messages.body.messages.map((message) => message.role), ["user", "agent", "user", "agent", "user", "agent", "user", "agent", "user", "agent"]);
assert.deepEqual(messages.body.messages.map((message) => message.text), ["sentinel-1", "ok-1", "sentinel-2", "ok-2", "sentinel-3", "ok-3", "sentinel-4", "ok-4", "sentinel-5", "ok-5"]);
const sessions = await getJson(port, `/v1/workbench/sessions?includeSessionId=${encodeURIComponent(session.id)}`);
assert.equal(sessions.status, 200);
assert.equal(sessions.body.sessions[0].title, "sentinel-5");
assert.equal(sessions.body.sessions[0].preview, "ok-5");
assert.equal(sessions.body.sessions[0].titleSource, "message-projection");
assert.equal(sessions.body.sessions[0].previewSource, "message-projection");
} finally {
await new Promise((resolve, reject) => server.close((error) => error ? reject(error) : resolve()));
}
});
test("workbench trace event continuation returns empty catch-up page instead of 404", async () => {
const traceId = "trc_workbench_trace_events_catchup";
const session = {
id: "ses_workbench_trace_events_catchup",
projectId: "prj_hwpod_workbench",
agentId: "hwlab-code-agent",
status: "running",
startedAt: "2026-06-27T01:00:00.000Z",
ownerUserId: ACTOR.id,
conversationId: "cnv_workbench_trace_events_catchup",
threadId: "thread-workbench-trace-events-catchup",
lastTraceId: traceId,
updatedAt: "2026-06-27T01:00:05.000Z",
session: {
sessionStatus: "running",
messages: [
{ role: "user", text: "start", traceId, createdAt: "2026-06-27T01:00:00.000Z" }
],
valuesRedacted: true,
secretMaterialStored: false
}
};
const runtimeStore = createDurableFactsRuntimeStore({
sessions: [{
session,
status: "running",
projectionStatus: "projecting",
projectionHealth: "projecting",
lastProjectedSeq: 5,
events: [
{ projectedSeq: 1, sourceSeq: 1, type: "request", status: "accepted", label: "request:accepted", createdAt: "2026-06-27T01:00:00.000Z" },
{ projectedSeq: 2, sourceSeq: 2, type: "tool", status: "running", label: "tool:running", createdAt: "2026-06-27T01:00:05.000Z" }
]
}]
});
const accessController = {
async ensureBootstrap() {},
async authenticate() { return { ok: true, actor: ACTOR, session: { id: "uss_workbench_reader" } }; }
};
const server = createCloudApiServer({ accessController, runtimeStore });
await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve));
try {
const { port } = server.address();
const response = await getJson(port, `/v1/workbench/traces/${encodeURIComponent(traceId)}/events?afterProjectedSeq=2&limit=100`);
assert.equal(response.status, 200);
assert.equal(response.body.status, "projecting");
assert.deepEqual(response.body.events, []);
assert.equal(response.body.range.fromProjectedSeq, null);
assert.equal(response.body.range.toProjectedSeq, null);
assert.equal(response.body.range.total, 5);
assert.equal(response.body.nextProjectedSeq, 2);
assert.equal(response.body.traceLastSeq, 5);
assert.equal(response.body.fullTraceLoaded, false);
assert.equal(response.body.projectionStatus, "projecting");
assert.equal(response.body.projectionHealth, "projecting");
assert.equal(response.body.diagnostic.code, "workbench_trace_events_missing");
assert.equal(response.body.blocker, null);
} finally {
await new Promise((resolve, reject) => server.close((error) => error ? reject(error) : resolve()));
}
});
test("workbench terminal timing ignores late projection updatedAt after sealed result (#2132)", async () => {
const traceId = "trc_workbench_terminal_late_projection_update";
const session = {
@@ -1033,7 +1149,7 @@ test("workbench trace events reports metadata gap when turn projection is visibl
}
});
test("workbench trace events reports missing event page when session projection is visible", async () => {
test("workbench trace events reports catch-up page when session projection is visible", async () => {
const traceStore = createCodeAgentTraceStore();
const traceId = "trc_workbench_trace_events_gap";
const session = {
@@ -1066,17 +1182,18 @@ test("workbench trace events reports missing event page when session projection
assert.equal(turn.body.turn.status, "completed");
const trace = await getJson(port, `/v1/workbench/traces/${encodeURIComponent(traceId)}/events?limit=10`);
assert.equal(trace.status, 404);
assert.equal(trace.body.error.code, "workbench_trace_events_missing");
assert.equal(trace.body.error.message, "Workbench trace events are missing from the durable read model.");
assert.equal(trace.body.error.sessionId, session.id);
assert.equal(trace.body.error.projectionStatus, "blocked");
assert.equal(trace.body.error.projectionHealth, "degraded");
assert.equal(trace.body.error.lastProjectedSeq, 2);
assert.equal(trace.body.error.sourceRunId, "run_workbench_trace_events_gap");
assert.equal(trace.body.error.sourceCommandId, "cmd_workbench_trace_events_gap");
assert.equal(trace.body.error.blocker.code, "workbench_trace_events_missing");
assert.notEqual(trace.body.error.message, "Workbench trace is not visible to the current actor.");
assert.equal(trace.status, 200);
assert.equal(trace.body.status, "projecting");
assert.equal(trace.body.diagnostic.code, "workbench_trace_events_missing");
assert.equal(trace.body.sessionId, session.id);
assert.equal(trace.body.projectionStatus, "projecting");
assert.equal(trace.body.projectionHealth, "projecting");
assert.equal(trace.body.lastProjectedSeq, 2);
assert.equal(trace.body.sourceRunId, "run_workbench_trace_events_gap");
assert.equal(trace.body.sourceCommandId, "cmd_workbench_trace_events_gap");
assert.equal(trace.body.blocker, null);
assert.deepEqual(trace.body.events, []);
assert.equal(trace.body.fullTraceLoaded, false);
} finally {
await new Promise((resolve, reject) => server.close((error) => error ? reject(error) : resolve()));
}
@@ -2206,15 +2323,18 @@ function normalizeTestMessages(session, { traceId, status, finalText, terminal,
const isAssistant = role === "agent" || role === "assistant";
const text = isAssistant && finalText ? finalText : String(message.text ?? message.content ?? "");
const messageStatus = isAssistant && finalText ? status : normalizeTestStatus(message.status ?? (role === "user" ? "sent" : status));
const messageTraceId = message.traceId ?? traceId;
const projectedSeq = Number.isFinite(Number(message.projectedSeq)) ? Number(message.projectedSeq) : index + 1;
const sourceSeq = Number.isFinite(Number(message.sourceSeq)) ? Number(message.sourceSeq) : projectedSeq;
return {
messageId: message.messageId ?? `msg_${session.id}_${index}`,
sessionId: session.id,
turnId: message.turnId ?? traceId,
traceId: message.traceId ?? traceId,
turnId: message.turnId ?? messageTraceId,
traceId: messageTraceId,
role,
status: messageStatus,
projectedSeq: index + 1,
sourceSeq: index + 1,
projectedSeq,
sourceSeq,
sourceEventId: `${session.id}:message:${index}`,
terminal: isAssistant && terminal,
sealed: isAssistant && terminal,
+111 -26
View File
@@ -1,5 +1,5 @@
/*
* SPEC: PJ2026-0104010803 Workbench唯一投影 draft-2026-06-20-p2-terminal-outbox-recovery; PJ2026-010403 API契约 draft-2026-06-20-p2-terminal-outbox-recovery; PJ2026-010401 Web工作台 draft-2026-06-20-p2-terminal-outbox-recovery; PJ2026-01060505 Workbench Performance draft-2026-06-17-p0
* SPEC: PJ2026-0104010803 Workbench唯一投影 draft-2026-06-20-p2-terminal-outbox-recovery; draft-2026-06-27-p0-read-model-timeline-contract; PJ2026-010403 API契约 draft-2026-06-20-p2-terminal-outbox-recovery; draft-2026-06-27-p0-workbench-read-api-contract; PJ2026-010401 Web工作台 draft-2026-06-20-p2-terminal-outbox-recovery; draft-2026-06-27-p0-workbench-read-model-contract; PJ2026-01060505 Workbench Performance draft-2026-06-17-p0
* 职责: Workbench REST read model session/message/turn/trace facts AgentRun syncbilling finalize workspace repair
*/
import { createHash } from "node:crypto";
@@ -948,10 +948,16 @@ function factSessionSummary(session, facts = {}) {
const status = normalizeStatus(checkpointStatus ?? traceStatus ?? turnStatus ?? session?.status);
const timingSource = checkpointStatus ? checkpoint : turn ?? checkpoint ?? (traceId ? null : session);
const timing = factTimingProjection(timingSource, status);
const title = sessionTitleFromMessages(messages);
const preview = sessionPreviewFromMessages(messages);
return {
sessionId,
threadId: safeOpaqueId(session?.threadId) ?? (textValue(session?.threadId) || null),
agentId: session?.agentId ?? "hwlab-code-agent",
title,
preview,
titleSource: title ? "message-projection" : null,
previewSource: preview ? "message-projection" : null,
status,
running: RUNNING_STATUSES.has(status),
terminal: TERMINAL_STATUSES.has(status) && !RUNNING_STATUSES.has(status),
@@ -1039,12 +1045,12 @@ function factMessagesForSession(session, facts = {}) {
const messageGroups = canonicalFactMessageGroupsForSession(
factArray(facts.messages)
.filter((message) => message.sessionId === sessionId)
.sort(compareFactMessagesAsc),
.sort((left, right) => compareFactMessagesForSessionAsc(left, right, facts)),
partsByMessageId,
facts
);
return messageGroups
.sort((left, right) => compareFactMessagesAsc(left.message, right.message))
.sort((left, right) => compareFactMessagesForSessionAsc(left.message, right.message, facts))
.map((group) => factMessageDto(group.message, group.parts, facts))
.filter(Boolean);
}
@@ -1544,6 +1550,55 @@ function compareFactMessagesAsc(left, right) {
return compareNumberAsc(factSeq(left), factSeq(right)) || compareTimestampAsc(left?.createdAt ?? left?.updatedAt, right?.createdAt ?? right?.updatedAt) || compareText(left?.messageId, right?.messageId);
}
function compareFactMessagesForSessionAsc(left, right, facts = {}) {
const leftKey = factMessageTimelineKey(left, facts);
const rightKey = factMessageTimelineKey(right, facts);
return compareOptionalNumberAsc(leftKey.turnSeq, rightKey.turnSeq)
|| compareTimestampAsc(leftKey.turnStartedAt, rightKey.turnStartedAt)
|| compareText(leftKey.traceId, rightKey.traceId)
|| compareOptionalNumberAsc(leftKey.roleRank, rightKey.roleRank)
|| compareTimestampAsc(leftKey.messageCreatedAt, rightKey.messageCreatedAt)
|| compareOptionalNumberAsc(leftKey.messageSeq, rightKey.messageSeq)
|| compareText(left?.messageId, right?.messageId);
}
function factMessageTimelineKey(message, facts = {}) {
const traceId = safeTraceId(message?.traceId) ?? null;
const turn = traceId ? factTurnForTrace(facts, traceId) : null;
const checkpoint = traceId ? factCheckpointForTrace(facts, traceId) : null;
return {
traceId: traceId ?? "",
turnSeq: firstFiniteNumber(factSeq(turn), factSeq(checkpoint)),
turnStartedAt: firstTimestampIso(turn?.startedAt, checkpoint?.startedAt, message?.createdAt, message?.updatedAt),
roleRank: factMessageTimelineRoleRank(message?.role),
messageCreatedAt: timestampIso(message?.createdAt ?? message?.updatedAt),
messageSeq: factSeq(message)
};
}
function factMessageTimelineRoleRank(role) {
const normalized = normalizeRole(role);
if (normalized === "user") return 0;
if (isAssistantLikeRole(normalized)) return 1;
return 2;
}
function firstFiniteNumber(...values) {
for (const value of values) {
const parsed = Number(value);
if (Number.isFinite(parsed)) return Math.trunc(parsed);
}
return null;
}
function compareOptionalNumberAsc(left, right) {
const leftNumber = Number(left);
const rightNumber = Number(right);
const leftComparable = Number.isFinite(leftNumber) ? leftNumber : Number.MAX_SAFE_INTEGER;
const rightComparable = Number.isFinite(rightNumber) ? rightNumber : Number.MAX_SAFE_INTEGER;
return leftComparable - rightComparable;
}
function compareFactPartsAsc(left, right) {
return compareNumberAsc(left?.partIndex, right?.partIndex) || compareNumberAsc(factSeq(left), factSeq(right)) || compareText(left?.partId, right?.partId);
}
@@ -1879,31 +1934,18 @@ async function handleWorkbenchTraceEventPage(request, response, url, options, ac
const traceTurn = factTurnForTrace(metadataFacts, traceId, traceId);
const turnTraceStatus = normalizeStatus(traceTurn?.status);
const traceStatus = turnTraceStatus !== "unknown" ? turnTraceStatus : durableTraceStatus(pageResult.facts.traceEvents);
const page = traceEventPageFromFacts(pageResult.facts.traceEvents, pageOptions, { total: projection.lastProjectedSeq, traceStatus });
const page = traceEventPageFromFacts(pageResult.facts.traceEvents, pageOptions, { total: projection.lastProjectedSeq, traceLastSeq: projection.lastProjectedSeq, traceStatus });
const missingTraceEvents = traceEventPageMissing(page, projection, pageOptions);
const missingTraceEventsDiagnostic = missingTraceEvents
? traceEventsReadModelBlocker("workbench_trace_events_missing", "Workbench trace events are still catching up in the durable read model.", { traceId, projection, session, turn: traceTurn, route: url.pathname })
: null;
if (missingTraceEvents) {
const blocker = traceEventsReadModelBlocker("workbench_trace_events_missing", "Workbench trace events are missing from the durable read model.", { traceId, projection, session, route: url.pathname });
attachWorkbenchReadModelDiagnosticOtel(request, { route: "/v1/workbench/traces/:id/events", code: "workbench_trace_events_missing", sessionId: factSessionId(session), traceId, count: pageResult.count });
emitWorkbenchTraceEventsReadOtel(traceId, options, {
startTimeMs: startedAt,
statusCode: 404,
phase: "events_missing",
pageOptions,
page,
projection,
sessionId: factSessionId(session),
turnId: factTurnId(traceTurn, traceId),
readCount: pageResult.count,
rawEventCount: factArray(pageResult.facts?.traceEvents).length,
metadataCount: metadata.count,
errorCode: "workbench_trace_events_missing"
});
return sendJson(response, 404, workbenchTraceEventsReadModelError(blocker, { traceId, projection, session, turn: traceTurn }));
}
const responseProjection = page.blocker ? blockedTraceEventProjection(projection, page.blocker) : projection;
const responseProjection = page.blocker ? blockedTraceEventProjection(projection, page.blocker) : missingTraceEvents ? catchingUpTraceEventProjection(projection, missingTraceEventsDiagnostic) : projection;
const body = {
ok: true,
status: page.blocker ? "blocked" : "succeeded",
status: page.blocker ? "blocked" : missingTraceEvents ? "projecting" : "succeeded",
contractVersion: "workbench-trace-events-v1",
traceId,
sessionId: factSessionId(session),
@@ -1923,6 +1965,7 @@ async function handleWorkbenchTraceEventPage(request, response, url, options, ac
lastEventAt: responseProjection.lastEventAt,
finishedAt: responseProjection.finishedAt,
durationMs: responseProjection.durationMs,
diagnostic: responseProjection.diagnostic ?? null,
valuesRedacted: true,
secretMaterialStored: false
};
@@ -1930,7 +1973,7 @@ async function handleWorkbenchTraceEventPage(request, response, url, options, ac
emitWorkbenchTraceEventsReadOtel(traceId, options, {
startTimeMs: startedAt,
statusCode: 200,
phase: page.blocker ? "blocked" : "succeeded",
phase: page.blocker ? "blocked" : missingTraceEvents ? "events_catching_up" : "succeeded",
pageOptions,
page,
projection: responseProjection,
@@ -2008,8 +2051,8 @@ function emitWorkbenchTraceEventsReadOtel(traceId, options = {}, fields = {}, er
toSeq: nonNegativeInteger(range.toProjectedSeq),
totalEvents: nonNegativeInteger(range.total ?? page.eventCount),
hasMore,
fullTraceLoaded: hasMore === null ? null : !hasMore && !page.blocker,
traceLastSeq: nonNegativeInteger(projection.lastProjectedSeq),
fullTraceLoaded: typeof page.fullTraceLoaded === "boolean" ? page.fullTraceLoaded : hasMore === null ? null : !hasMore && !page.blocker,
traceLastSeq: nonNegativeInteger(page.traceLastSeq ?? projection.lastProjectedSeq),
projectionStatus: projection.projectionStatus ?? null,
projectionHealth: projection.projectionHealth ?? null,
sourceRunId: projection.sourceRunId ?? null,
@@ -2134,7 +2177,8 @@ function traceEventPageFromFacts(sourceEvents, options, metadata = {}) {
const events = pageRows.map(factTraceEventDto).filter(Boolean);
const toProjectedSeq = events.length ? events.at(-1).projectedSeq : null;
const hasMore = rows.length > options.limit;
const total = Number.isFinite(Number(metadata.total)) ? Math.trunc(Number(metadata.total)) : hasMore ? null : Math.max(options.afterProjectedSeq, toProjectedSeq ?? options.afterProjectedSeq);
const total = Number.isFinite(Number(metadata.total)) ? Math.trunc(Number(metadata.total)) : Math.max(options.afterProjectedSeq, toProjectedSeq ?? options.afterProjectedSeq);
const traceLastSeq = Number.isFinite(Number(metadata.traceLastSeq)) ? Math.trunc(Number(metadata.traceLastSeq)) : total;
return {
events,
eventCount: total ?? Math.max(options.afterProjectedSeq, toProjectedSeq ?? options.afterProjectedSeq),
@@ -2147,6 +2191,8 @@ function traceEventPageFromFacts(sourceEvents, options, metadata = {}) {
total
},
hasMore,
fullTraceLoaded: !hasMore && (toProjectedSeq ?? options.afterProjectedSeq) >= traceLastSeq,
traceLastSeq,
nextProjectedSeq: events.length ? toProjectedSeq : options.afterProjectedSeq,
nextCursor: hasMore && toProjectedSeq ? `projected:${toProjectedSeq}` : null,
traceStatus: metadata.traceStatus ?? "unknown",
@@ -2210,6 +2256,19 @@ function blockedTraceEventProjection(projection = {}, blocker) {
};
}
function catchingUpTraceEventProjection(projection = {}, diagnostic = null) {
const status = projection.projectionStatus === "caught-up" || !projection.projectionStatus ? "projecting" : projection.projectionStatus;
const health = projection.projectionHealth === "caught-up" || projection.projectionHealth === "healthy" || !projection.projectionHealth ? "projecting" : projection.projectionHealth;
return {
...projection,
projectionStatus: status,
projectionHealth: health,
blocker: null,
diagnostic: diagnostic ? { ...diagnostic, projectionStatus: status, projectionHealth: health } : null,
valuesRedacted: true
};
}
function traceEventPageMissing(page = {}, projection = {}, options = {}) {
if (page.blocker) return false;
const expectedSeq = Number(projection?.lastProjectedSeq);
@@ -2697,6 +2756,32 @@ function firstUserPreview(messages) {
return messages.find((message) => message.role === "user")?.textPreview ?? null;
}
function latestUserPreview(messages) {
return [...messages].reverse().find((message) => message.role === "user")?.textPreview ?? null;
}
function latestValidMessagePreview(messages) {
return [...messages].reverse().find((message) => message.role === "user" || isAssistantLikeRole(message.role))?.textPreview ?? null;
}
function latestAssistantPreview(messages) {
return [...messages].reverse().find((message) => isAssistantLikeRole(message.role))?.textPreview ?? null;
}
function sessionTitleFromMessages(messages = []) {
return boundedPreviewText(latestUserPreview(messages) ?? firstUserPreview(messages) ?? latestAssistantPreview(messages));
}
function sessionPreviewFromMessages(messages = []) {
return boundedPreviewText(latestValidMessagePreview(messages) ?? firstUserPreview(messages) ?? latestAssistantPreview(messages));
}
function boundedPreviewText(value, maxLength = 160) {
const text = textValue(value).replace(/\s+/gu, " ").trim();
if (!text) return null;
return text.length > maxLength ? `${text.slice(0, maxLength - 3)}...` : text;
}
function isAssistantLikeRole(role) {
return role === "assistant" || role === "agent";
}