diff --git a/internal/cloud/server-workbench-http.test.ts b/internal/cloud/server-workbench-http.test.ts index 4e9950fc..ed373258 100644 --- a/internal/cloud/server-workbench-http.test.ts +++ b/internal/cloud/server-workbench-http.test.ts @@ -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, diff --git a/internal/cloud/server-workbench-http.ts b/internal/cloud/server-workbench-http.ts index 703c46ba..12254daf 100644 --- a/internal/cloud/server-workbench-http.ts +++ b/internal/cloud/server-workbench-http.ts @@ -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 sync、billing 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"; }