Merge pull request #2246 from pikasTech/fix/1217-workbench-timeline-d518

fix: stabilize workbench session timeline
This commit is contained in:
Lyon
2026-06-28 17:11:00 +08:00
committed by GitHub
4 changed files with 186 additions and 11 deletions
@@ -627,6 +627,110 @@ test("workbench session messages use turn timeline instead of clustered projecti
}
});
test("workbench session messages use prompt admission time when turn seq is stale", async () => {
const sessionId = "ses_workbench_timeline_stale_turn_seq";
const turnRows = [
{ traceId: "trc_stale_turn_01", turnSeq: 1, at: "2026-06-28T08:40:42.000Z", prompt: "sentinel-01", reply: "ok-01" },
{ traceId: "trc_stale_turn_02", turnSeq: 3, at: "2026-06-28T08:43:51.000Z", prompt: "sentinel-02", reply: "ok-02" },
{ traceId: "trc_stale_turn_03", turnSeq: 2, at: "2026-06-28T08:44:07.000Z", prompt: "sentinel-03", reply: "ok-03" },
{ traceId: "trc_stale_turn_04", turnSeq: 5, at: "2026-06-28T08:44:27.000Z", prompt: "sentinel-04", reply: "ok-04" },
{ traceId: "trc_stale_turn_05", turnSeq: 4, at: "2026-06-28T08:44:44.000Z", prompt: "sentinel-05", reply: "ok-05" }
];
const facts = emptyFacts();
facts.sessions.push({
sessionId,
ownerUserId: ACTOR.id,
ownerRole: ACTOR.role,
agentId: "hwlab-code-agent",
status: "completed",
lastTraceId: turnRows.at(-1).traceId,
projectedSeq: 99,
sourceSeq: 99,
terminal: true,
sealed: true,
createdAt: turnRows[0].at,
updatedAt: "2026-06-28T08:45:00.000Z",
valuesRedacted: true
});
turnRows.forEach((row, index) => {
const userMessageId = `msg_${row.traceId}_user`;
const agentMessageId = `msg_${row.traceId}_agent`;
const finishedAt = new Date(Date.parse(row.at) + 5000).toISOString();
facts.messages.push(
{ messageId: userMessageId, sessionId, turnId: row.traceId, traceId: row.traceId, role: "user", status: "sent", text: row.prompt, projectedSeq: 100 + index, sourceSeq: 100 + index, createdAt: row.at, updatedAt: row.at, valuesRedacted: true },
{ messageId: agentMessageId, sessionId, turnId: row.traceId, traceId: row.traceId, role: "agent", status: "completed", text: row.reply, projectedSeq: 200 + index, sourceSeq: 200 + index, terminal: true, sealed: true, createdAt: finishedAt, updatedAt: finishedAt, timing: { startedAt: row.at, lastEventAt: finishedAt, finishedAt, durationMs: 5000, valuesRedacted: true }, valuesRedacted: true }
);
facts.parts.push(
{ partId: `prt_${row.traceId}_user`, messageId: userMessageId, sessionId, turnId: row.traceId, traceId: row.traceId, partIndex: 0, partType: "text", status: "sent", text: row.prompt, projectedSeq: 100 + index, updatedAt: row.at, valuesRedacted: true },
{ partId: `prt_${row.traceId}_agent`, messageId: agentMessageId, sessionId, turnId: row.traceId, traceId: row.traceId, partIndex: 0, partType: "final_response", status: "completed", text: row.reply, projectedSeq: 200 + index, terminal: true, sealed: true, updatedAt: finishedAt, valuesRedacted: true }
);
facts.turns.push({ turnId: row.traceId, sessionId, traceId: row.traceId, messageId: agentMessageId, status: "completed", projectedSeq: row.turnSeq, sourceSeq: row.turnSeq, terminal: true, sealed: true, startedAt: row.at, lastEventAt: finishedAt, finishedAt, durationMs: 5000, updatedAt: finishedAt, valuesRedacted: true });
facts.checkpoints.push({ traceId: row.traceId, sessionId, turnId: row.traceId, status: "completed", projectionStatus: "caught_up", projectionHealth: "caught-up", projectedSeq: row.turnSeq, sourceSeq: row.turnSeq, terminal: true, sealed: true, startedAt: row.at, lastEventAt: finishedAt, finishedAt, durationMs: 5000, updatedAt: finishedAt, valuesRedacted: true });
});
const runtimeStore = createRuntimeStoreFromFacts(facts);
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(sessionId)}/messages?limit=20`);
assert.equal(messages.status, 200);
assert.deepEqual(messages.body.messages.map((message) => message.text), ["sentinel-01", "ok-01", "sentinel-02", "ok-02", "sentinel-03", "ok-03", "sentinel-04", "ok-04", "sentinel-05", "ok-05"]);
assert.equal(messages.body.roleSequencePrefix, "UAUAUAUAUA");
assert.equal(messages.body.adjacentSameRoleCount, 0);
assert.match(messages.body.timelineDigest, /^sha256:[0-9a-f]{64}$/u);
const sessions = await getJson(port, `/v1/workbench/sessions?includeSessionId=${encodeURIComponent(sessionId)}`);
assert.equal(sessions.status, 200);
assert.equal(sessions.body.sessions[0].title, "sentinel-05");
assert.equal(sessions.body.sessions[0].preview, "ok-05");
assert.equal(sessions.body.sessions[0].firstUserMessagePreview, "sentinel-01");
} finally {
await new Promise((resolve, reject) => server.close((error) => error ? reject(error) : resolve()));
}
});
test("workbench session list hides shadowed AgentRun alias rows", async () => {
const sessionId = "ses_fake_echo_canonical";
const aliasSessionId = "ses_agentrun_fake_echo_canonical";
const traceId = "trc_fake_echo_alias";
const facts = emptyFacts();
facts.sessions.push(
{ sessionId: aliasSessionId, ownerUserId: ACTOR.id, ownerRole: ACTOR.role, status: "failed", lastTraceId: traceId, projectedSeq: 20, sourceSeq: 20, updatedAt: "2026-06-28T08:52:18.000Z", valuesRedacted: true },
{ sessionId, ownerUserId: ACTOR.id, ownerRole: ACTOR.role, status: "completed", lastTraceId: traceId, projectedSeq: 10, sourceSeq: 10, updatedAt: "2026-06-28T08:52:10.000Z", valuesRedacted: true }
);
facts.messages.push(
{ messageId: "msg_fake_user", sessionId, turnId: traceId, traceId, role: "user", status: "sent", text: "ECHO sentinel-01", projectedSeq: 1, sourceSeq: 1, createdAt: "2026-06-28T08:52:00.000Z", updatedAt: "2026-06-28T08:52:00.000Z", valuesRedacted: true },
{ messageId: "msg_fake_agent", sessionId, turnId: traceId, traceId, role: "agent", status: "completed", text: "ECHO sentinel-01 ok", projectedSeq: 2, sourceSeq: 2, terminal: true, sealed: true, createdAt: "2026-06-28T08:52:05.000Z", updatedAt: "2026-06-28T08:52:05.000Z", valuesRedacted: true }
);
facts.parts.push(
{ partId: "prt_fake_user", messageId: "msg_fake_user", sessionId, turnId: traceId, traceId, partIndex: 0, partType: "text", status: "sent", text: "ECHO sentinel-01", projectedSeq: 1, valuesRedacted: true },
{ partId: "prt_fake_agent", messageId: "msg_fake_agent", sessionId, turnId: traceId, traceId, partIndex: 0, partType: "final_response", status: "completed", text: "ECHO sentinel-01 ok", projectedSeq: 2, terminal: true, sealed: true, valuesRedacted: true }
);
const runtimeStore = createRuntimeStoreFromFacts(facts);
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 sessions = await getJson(port, "/v1/workbench/sessions?limit=10");
assert.equal(sessions.status, 200);
assert.deepEqual(sessions.body.sessions.map((session) => session.sessionId), [sessionId]);
assert.equal(sessions.body.sessions[0].title, "ECHO sentinel-01");
assert.equal(sessions.body.sessions[0].preview, "ECHO sentinel-01 ok");
} 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 = {
@@ -2193,6 +2297,10 @@ async function waitForCondition(predicate, timeoutMs = 500) {
function createDurableFactsRuntimeStore({ sessions = [], queryError = null, queries = [] } = {}) {
const facts = emptyFacts();
for (const item of sessions) mergeFacts(facts, buildDurableFactsForSession(item));
return createRuntimeStoreFromFacts(facts, { queryError, queries });
}
function createRuntimeStoreFromFacts(facts, { queryError = null, queries = [] } = {}) {
return {
queries,
async queryWorkbenchFacts(params = {}) {
+70 -5
View File
@@ -1,5 +1,5 @@
/*
* 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
* SPEC: PJ2026-0104010803 Workbench唯一投影 draft-2026-06-20-p2-terminal-outbox-recovery; draft-2026-06-27-p0-read-model-timeline-contract; draft-2026-06-28-p0-d518-session-timeline-consistency; 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 sync、billing finalize 或 workspace repair。
*/
import { createHash } from "node:crypto";
@@ -890,9 +890,22 @@ async function queryFactsForSessionSummaries(readModel, sessions = []) {
}
function visibleFactSessions(facts, actor) {
return factArray(facts?.sessions)
return hideShadowedAgentRunAliasSessions(factArray(facts?.sessions)
.filter((session) => canActorReadFactSession(session, actor))
.sort(compareFactSessionsDesc);
.sort(compareFactSessionsDesc));
}
function hideShadowedAgentRunAliasSessions(sessions = []) {
const canonicalTraceIds = new Set(sessions
.filter((session) => !isAgentRunAliasSessionId(factSessionId(session)))
.map(factLastTraceId)
.filter(Boolean));
if (canonicalTraceIds.size === 0) return sessions;
return sessions.filter((session) => !isAgentRunAliasSessionId(factSessionId(session)) || !canonicalTraceIds.has(factLastTraceId(session)));
}
function isAgentRunAliasSessionId(sessionId) {
return /^ses_agentrun_/u.test(textValue(sessionId));
}
function canActorReadFactSession(session, actor) {
@@ -1553,7 +1566,8 @@ function compareFactMessagesAsc(left, right) {
function compareFactMessagesForSessionAsc(left, right, facts = {}) {
const leftKey = factMessageTimelineKey(left, facts);
const rightKey = factMessageTimelineKey(right, facts);
return compareOptionalNumberAsc(leftKey.turnSeq, rightKey.turnSeq)
return compareOptionalTimestampAsc(leftKey.timelineAnchorAt, rightKey.timelineAnchorAt)
|| compareOptionalNumberAsc(leftKey.turnSeq, rightKey.turnSeq)
|| compareTimestampAsc(leftKey.turnStartedAt, rightKey.turnStartedAt)
|| compareText(leftKey.traceId, rightKey.traceId)
|| compareOptionalNumberAsc(leftKey.roleRank, rightKey.roleRank)
@@ -1566,8 +1580,10 @@ function factMessageTimelineKey(message, facts = {}) {
const traceId = safeTraceId(message?.traceId) ?? null;
const turn = traceId ? factTurnForTrace(facts, traceId) : null;
const checkpoint = traceId ? factCheckpointForTrace(facts, traceId) : null;
const anchorMessage = traceId ? factTimelineAnchorMessageForTrace(facts, traceId) : null;
return {
traceId: traceId ?? "",
timelineAnchorAt: firstTimestampIso(anchorMessage?.createdAt, anchorMessage?.updatedAt),
turnSeq: firstFiniteNumber(factSeq(turn), factSeq(checkpoint)),
turnStartedAt: firstTimestampIso(turn?.startedAt, checkpoint?.startedAt, message?.createdAt, message?.updatedAt),
roleRank: factMessageTimelineRoleRank(message?.role),
@@ -1576,6 +1592,15 @@ function factMessageTimelineKey(message, facts = {}) {
};
}
function factTimelineAnchorMessageForTrace(facts = {}, traceId) {
const safeTrace = safeTraceId(traceId);
if (!safeTrace) return null;
const messages = factArray(facts.messages).filter((message) => message.traceId === safeTrace);
const userMessages = messages.filter((message) => normalizeRole(message?.role) === "user");
const candidates = userMessages.length > 0 ? userMessages : messages;
return [...candidates].sort((left, right) => compareOptionalTimestampAsc(left?.createdAt ?? left?.updatedAt, right?.createdAt ?? right?.updatedAt) || compareOptionalNumberAsc(factSeq(left), factSeq(right)) || compareText(left?.messageId, right?.messageId))[0] ?? null;
}
function factMessageTimelineRoleRank(role) {
const normalized = normalizeRole(role);
if (normalized === "user") return 0;
@@ -1619,6 +1644,12 @@ function compareTimestampDesc(left, right) {
return timestampMs(right) - timestampMs(left);
}
function compareOptionalTimestampAsc(left, right) {
const leftMs = optionalTimestampMs(left);
const rightMs = optionalTimestampMs(right);
return (leftMs ?? Number.MAX_SAFE_INTEGER) - (rightMs ?? Number.MAX_SAFE_INTEGER);
}
function compareNumberAsc(left, right) {
return numericValue(left, Number.MAX_SAFE_INTEGER) - numericValue(right, Number.MAX_SAFE_INTEGER);
}
@@ -1636,6 +1667,11 @@ function timestampMs(value) {
return Number.isFinite(parsed) ? parsed : 0;
}
function optionalTimestampMs(value) {
const parsed = Date.parse(String(value ?? ""));
return Number.isFinite(parsed) ? parsed : null;
}
function numericValue(value, fallback) {
const parsed = Number(value);
return Number.isFinite(parsed) ? parsed : fallback;
@@ -1737,6 +1773,10 @@ async function handleWorkbenchMessagePage(request, response, url, options, actor
cursor: offset > 0 ? cursorFromOffset(offset) : null,
nextCursor: nextOffset < messages.length ? cursorFromOffset(nextOffset) : null,
hasMore: nextOffset < messages.length,
timelineDigest: workbenchTimelineDigest(page),
roleSequencePrefix: messagesRoleSequence(page).slice(0, 32),
adjacentSameRoleCount: adjacentSameRoleCount(page),
traceIds: uniqueText(page.map((message) => message?.traceId)).slice(0, 16),
valuesRedacted: true,
secretMaterialStored: false
};
@@ -2098,7 +2138,7 @@ function emitWorkbenchSessionMessagesReadOtel(request, options = {}, fields = {}
const messages = Array.isArray(fields.messages) ? fields.messages : [];
const traceIds = uniqueText(messages.map((message) => message?.traceId)).slice(0, 8);
const businessTraceId = (safeTraceId(traceIds[0]) ?? textValue(fields.sessionId)) || "workbench_session_messages";
const roleSequence = messages.map((message) => roleSequenceSymbol(message?.role)).join("");
const roleSequence = messagesRoleSequence(messages);
void emitCodeAgentOtelSpan("session_messages_read", businessTraceId, options.env ?? process.env, {
startTimeMs: fields.startTimeMs,
status: error ? "error" : "ok",
@@ -2116,6 +2156,7 @@ function emitWorkbenchSessionMessagesReadOtel(request, options = {}, fields = {}
roleSequencePrefix: roleSequence.slice(0, 32),
consecutiveUserPrefix: consecutiveRolePrefix(messages, "user"),
adjacentSameRoleCount: adjacentSameRoleCount(messages),
timelineDigest: workbenchTimelineDigest(messages),
userCount: messages.filter((message) => normalizeRole(message?.role) === "user").length,
agentCount: messages.filter((message) => isAssistantLikeRole(message?.role)).length,
traceIds: traceIds.join(","),
@@ -2124,6 +2165,30 @@ function emitWorkbenchSessionMessagesReadOtel(request, options = {}, fields = {}
});
}
function messagesRoleSequence(messages = []) {
return messages.map((message) => roleSequenceSymbol(message?.role)).join("");
}
function workbenchTimelineDigest(messages = []) {
const payload = messages.map((message) => {
const terminal = isTerminalProjectionStatus(normalizeStatus(message?.status)) || Boolean(message?.finishedAt);
return {
messageId: message?.messageId ?? null,
role: normalizeRole(message?.role),
sessionId: message?.sessionId ?? null,
traceId: message?.traceId ?? null,
turnId: message?.turnId ?? null,
status: normalizeStatus(message?.status),
partIds: factArray(message?.parts).map((part) => part?.partId ?? null),
textHash: textValue(message?.text || message?.textPreview) ? hash(textValue(message?.text || message?.textPreview)).slice(0, 16) : null,
startedAt: message?.startedAt ?? null,
finishedAt: terminal ? message?.finishedAt ?? null : null,
durationMs: terminal && Number.isFinite(Number(message?.durationMs)) ? Math.trunc(Number(message.durationMs)) : null
};
});
return `sha256:${hash(JSON.stringify(payload))}`;
}
function sessionSummaryNeedsTitleFallback(session) {
if (!session || typeof session !== "object") return true;
return !sessionPreviewText(session);
@@ -1,4 +1,4 @@
// SPEC: PJ2026-010403 API契约 draft-2026-06-18-r1; PJ2026-010401 Web工作台 draft-2026-06-18-r1; PJ2026-0104010803 唯一投影 draft-2026-06-18-p0-unique-projection.
// SPEC: PJ2026-010403 API契约 draft-2026-06-18-r1; PJ2026-010401 Web工作台 draft-2026-06-18-r1; PJ2026-0104010803 唯一投影 draft-2026-06-18-p0-unique-projection; draft-2026-06-28-p0-d518-session-timeline-consistency.
// Responsibility: Normalized Workbench server-state reducer for session, turn, and trace caches.
import type { ChatMessage, ProjectionDiagnostic, WorkbenchSessionRecord } from "@/types";
@@ -228,7 +228,7 @@ function sealTerminalTransitionMessageTiming(existing: ChatMessage, incoming: Ch
const existingTiming = messageTimingProjection(existing);
const startedAt = existingTiming?.startedAt ?? incomingTiming?.startedAt ?? timestampOrNull(existing.createdAt) ?? null;
const finishedAt = terminalFinishedAt(startedAt, incomingTiming?.finishedAt, timestampOrNull(incoming.updatedAt), incomingTiming?.lastEventAt) ?? incomingTiming?.finishedAt ?? null;
const durationMs = positiveDuration(durationBetween(startedAt, finishedAt)) ?? positiveDuration(durationBetween(startedAt, new Date().toISOString())) ?? positiveDuration(incomingTiming?.durationMs) ?? (incomingTiming?.durationMs === 0 ? 1000 : null);
const durationMs = positiveDuration(durationBetween(startedAt, finishedAt)) ?? positiveDuration(incomingTiming?.durationMs) ?? (incomingTiming?.durationMs === 0 ? 0 : null);
if (durationMs === null) return incoming;
const timing: NonNullable<ChatMessage["timing"]> = {
...(incomingTiming ?? {}),
+6 -4
View File
@@ -1,4 +1,4 @@
// SPEC: PJ2026-010403 API契约 draft-2026-06-18-r1; PJ2026-010401 Web工作台 draft-2026-06-18-r1; PJ2026-0104010803 唯一投影 draft-2026-06-18-p0-unique-projection; PJ2026-01060505 Workbench Performance draft-2026-06-17-p0.
// SPEC: PJ2026-010403 API契约 draft-2026-06-18-r1; PJ2026-010401 Web工作台 draft-2026-06-18-r1; PJ2026-0104010803 唯一投影 draft-2026-06-18-p0-unique-projection; draft-2026-06-28-p0-d518-session-timeline-consistency; PJ2026-01060505 Workbench Performance draft-2026-06-17-p0.
// Responsibility: Session-first Workbench state orchestration for selection, turn admission, and trace lifecycle rendering.
import { computed, nextTick, ref } from "vue";
@@ -2127,7 +2127,8 @@ function messageTimingPatchForMerge(message: ChatMessage, value: unknown): Parti
if (incomingDuration !== null) return messageTimingPatchFromProjection({ ...(existingTiming ?? {}), ...incomingTiming, durationMs: incomingDuration, valuesRedacted: incomingTiming.valuesRedacted !== false && existingTiming?.valuesRedacted !== false });
const runningDuration = isTerminalMessageStatus(message.status) ? null : runningDurationFromTiming(existingTiming);
if (runningDuration !== null) {
const finishedAt = incomingTiming.finishedAt ?? incomingTiming.lastEventAt ?? new Date().toISOString();
const finishedAt = incomingTiming.finishedAt ?? incomingTiming.lastEventAt ?? null;
if (!finishedAt) return terminalMessageTimingPatchForNormalize(value);
return messageTimingPatchFromProjection({ ...(incomingTiming ?? {}), startedAt: existingTiming?.startedAt ?? incomingTiming.startedAt ?? null, lastEventAt: incomingTiming.lastEventAt ?? finishedAt, finishedAt, durationMs: runningDuration, valuesRedacted: incomingTiming.valuesRedacted !== false && existingTiming?.valuesRedacted !== false });
}
return terminalMessageTimingPatchForNormalize(value);
@@ -2180,7 +2181,8 @@ function messageTerminalFromRunningSealPatchForProjectionMerge(message: ChatMess
if (runningDuration === null) return incomingPatch;
const incomingDuration = firstPositiveFiniteNumber(incomingTiming?.durationMs) ?? positiveDurationBetween(incomingTiming?.startedAt, incomingTiming?.finishedAt) ?? positiveDurationBetween(incomingTiming?.startedAt, incomingTiming?.lastEventAt);
const durationMs = Math.max(incomingDuration ?? 0, runningDuration);
const finishedAt = incomingTiming?.finishedAt ?? incomingTiming?.lastEventAt ?? new Date().toISOString();
const finishedAt = incomingTiming?.finishedAt ?? incomingTiming?.lastEventAt ?? null;
if (!finishedAt) return incomingPatch;
const timing = {
...(incomingTiming ?? {}),
startedAt: previousTiming?.startedAt ?? incomingTiming?.startedAt ?? null,
@@ -2210,7 +2212,7 @@ function terminalMessageTimingPatchForNormalize(value: unknown): Partial<ChatMes
?? positiveDurationBetween(timing.startedAt, timing.finishedAt)
?? positiveDurationBetween(timing.startedAt, timing.lastEventAt)
?? positiveDurationBetween(timing.startedAt, record?.updatedAt)
?? (timing.durationMs === 0 ? 1000 : null);
?? (timing.durationMs === 0 ? 0 : null);
return messageTimingPatchFromProjection({ ...timing, durationMs });
}