diff --git a/internal/cloud/server-workbench-http.test.ts b/internal/cloud/server-workbench-http.test.ts index 27ffbb74..712b291e 100644 --- a/internal/cloud/server-workbench-http.test.ts +++ b/internal/cloud/server-workbench-http.test.ts @@ -645,6 +645,131 @@ test("workbench read model recovers trace events from durable projection without } }); +test("workbench read model canonicalizes duplicate assistant facts for the same trace", async () => { + const traceStore = createCodeAgentTraceStore(); + const traceId = "trc_workbench_duplicate_agent_message"; + const preferredMessageId = "msg_workbench_duplicate_agent_message_agent"; + const duplicateMessageId = "msg_3ae9b11ac11c2dcf816980916f1152b0"; + const finalText = "duplicate history final response"; + const session = { + id: "ses_workbench_duplicate_agent_message", + projectId: "prj_hwpod_workbench", + agentId: "hwlab-code-agent", + status: "completed", + startedAt: "2026-06-17T01:20:00.000Z", + endedAt: "2026-06-17T01:20:16.000Z", + ownerUserId: ACTOR.id, + conversationId: "cnv_workbench_duplicate_agent_message", + threadId: "thread-workbench-duplicate-agent-message", + lastTraceId: traceId, + updatedAt: "2026-06-17T01:20:16.000Z", + session: { + sessionStatus: "completed", + lastTraceId: traceId, + messages: [ + { messageId: "msg_workbench_duplicate_agent_message_user", role: "user", text: "duplicate ping", traceId, createdAt: "2026-06-17T01:20:00.000Z" }, + { messageId: preferredMessageId, role: "agent", text: "", traceId, status: "completed", createdAt: "2026-06-17T01:20:16.000Z", updatedAt: "2026-06-17T01:20:16.000Z" } + ], + valuesRedacted: true, + secretMaterialStored: false + } + }; + const facts = buildDurableFactsForSession({ + session, + events: [ + { seq: 1, type: "request", status: "accepted", label: "request:accepted", createdAt: "2026-06-17T01:20:00.000Z" }, + { seq: 2, type: "result", status: "completed", label: "result:completed", terminal: true, createdAt: "2026-06-17T01:20:16.000Z" } + ], + status: "completed", + lastProjectedSeq: 2 + }); + facts.turns[0].messageId = duplicateMessageId; + facts.messages.push({ + messageId: duplicateMessageId, + sessionId: session.id, + turnId: traceId, + traceId, + role: "agent", + status: "completed", + projectedSeq: 9, + sourceSeq: 9, + sourceEventId: `${session.id}:message:duplicate-hash`, + terminal: true, + sealed: true, + text: finalText, + timing: { + startedAt: "2026-06-17T01:20:00.000Z", + lastEventAt: "2026-06-17T01:21:48.722Z", + finishedAt: "2026-06-17T01:21:48.722Z", + durationMs: 108722, + status: "completed", + valuesRedacted: true + }, + startedAt: "2026-06-17T01:20:00.000Z", + lastEventAt: "2026-06-17T01:21:48.722Z", + finishedAt: "2026-06-17T01:21:48.722Z", + durationMs: 108722, + createdAt: "2026-06-17T01:20:16.000Z", + updatedAt: "2026-06-17T01:21:48.722Z", + valuesRedacted: true + }); + facts.parts.push({ + partId: "prt_workbench_duplicate_agent_message_hash_final", + messageId: duplicateMessageId, + sessionId: session.id, + turnId: traceId, + traceId, + partIndex: 0, + partType: "final_response", + status: "completed", + text: finalText, + projectedSeq: 9, + sourceSeq: 9, + sourceEventId: `${session.id}:part:duplicate-hash`, + terminal: true, + sealed: true, + updatedAt: "2026-06-17T01:21:48.722Z", + valuesRedacted: true + }); + const runtimeStore = { + async queryWorkbenchFacts(params = {}) { + const filtered = filterFacts(facts, params); + return { + facts: filtered, + count: Object.values(filtered).reduce((sum, rows) => sum + rows.length, 0), + persistence: { adapter: "test-durable-workbench-facts", durable: true } + }; + } + }; + const accessController = { + store: { + async getAgentSession(sessionId) { return sessionId === session.id ? session : null; }, + async getAgentSessionByTraceId(requestTraceId) { return requestTraceId === traceId ? session : null; } + }, + async ensureBootstrap() {}, + async authenticate() { return { ok: true, actor: ACTOR, session: { id: "uss_workbench_reader" } }; } + }; + const server = createCloudApiServer({ accessController, traceStore, workbenchRuntime: runtimeStore, codeAgentChatResults: createCodeAgentChatResultStore() }); + 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=10`); + assert.equal(messages.status, 200); + assert.equal(messages.body.total, 2); + assert.deepEqual(messages.body.messages.map((message) => message.role), ["user", "agent"]); + const assistant = messages.body.messages[1]; + assert.equal(assistant.messageId, preferredMessageId); + assert.equal(assistant.text, finalText); + assert.equal(assistant.parts.length, 1); + assert.equal(assistant.parts[0].messageId, preferredMessageId); + assert.equal(assistant.parts[0].type, "final_response"); + assert.equal(assistant.parts[0].text, finalText); + } finally { + await new Promise((resolve, reject) => server.close((error) => error ? reject(error) : resolve())); + } +}); + test("workbench read model projects terminal result atomically across session, messages, and turn", async () => { const traceStore = createCodeAgentTraceStore(); const results = createCodeAgentChatResultStore(); diff --git a/internal/cloud/server-workbench-http.ts b/internal/cloud/server-workbench-http.ts index 930c55e4..9f652f9a 100644 --- a/internal/cloud/server-workbench-http.ts +++ b/internal/cloud/server-workbench-http.ts @@ -986,13 +986,81 @@ function factMessagesForSession(session, facts = {}) { existing.push(part); partsByMessageId.set(messageId, existing); } - return factArray(facts.messages) - .filter((message) => message.sessionId === sessionId) - .sort(compareFactMessagesAsc) - .map((message) => factMessageDto(message, partsByMessageId.get(message.messageId) ?? [], facts)) + const messageGroups = canonicalFactMessageGroupsForSession( + factArray(facts.messages) + .filter((message) => message.sessionId === sessionId) + .sort(compareFactMessagesAsc), + partsByMessageId, + facts + ); + return messageGroups + .sort((left, right) => compareFactMessagesAsc(left.message, right.message)) + .map((group) => factMessageDto(group.message, group.parts, facts)) .filter(Boolean); } +function canonicalFactMessageGroupsForSession(messages = [], partsByMessageId = new Map(), facts = {}) { + const groups = new Map(); + for (const message of messages) { + const key = factMessageCanonicalGroupKey(message); + const existing = groups.get(key) ?? []; + existing.push(message); + groups.set(key, existing); + } + return [...groups.values()].map((groupMessages) => { + const message = selectCanonicalFactMessage(groupMessages, partsByMessageId, facts); + const aliasIds = uniqueText(groupMessages.map((item) => item?.messageId)); + const parts = aliasIds.flatMap((messageId) => partsByMessageId.get(messageId) ?? []); + return { message, parts }; + }); +} + +function factMessageCanonicalGroupKey(message) { + const messageId = textValue(message?.messageId) || textValue(message?.id); + const role = textValue(message?.role) || "agent"; + const traceId = safeTraceId(message?.traceId); + const sessionId = factSessionId(message) ?? textValue(message?.sessionId) ?? "session"; + if (traceId && isAssistantLikeRole(role)) return `assistant:${sessionId}:${traceId}`; + return `message:${messageId || sessionId}:${role}:${traceId || "none"}`; +} + +function selectCanonicalFactMessage(messages = [], partsByMessageId = new Map(), facts = {}) { + if (messages.length <= 1) return messages[0] ?? null; + const traceId = safeTraceId(messages[0]?.traceId); + const preferredMessageId = traceId ? workbenchLifecycleMessageId(traceId, "agent") : null; + const turnMessageId = traceId ? textValue(factTurnForTrace(facts, traceId)?.messageId) || null : null; + return [...messages].sort((left, right) => { + const score = factMessageCanonicalScore(right, partsByMessageId, preferredMessageId, turnMessageId) - factMessageCanonicalScore(left, partsByMessageId, preferredMessageId, turnMessageId); + return score || compareFactRecordsDesc(left, right); + })[0] ?? messages[0] ?? null; +} + +function factMessageCanonicalScore(message, partsByMessageId = new Map(), preferredMessageId = null, turnMessageId = null) { + const messageId = textValue(message?.messageId) || textValue(message?.id); + let score = 0; + if (preferredMessageId && messageId === preferredMessageId) score += 10000; + if (turnMessageId && messageId === turnMessageId) score += 1000; + if (factMessageHasFinalResponsePart(message, partsByMessageId)) score += 100; + if (message?.sealed === true || message?.terminal === true) score += 50; + if (isTerminalProjectionStatus(normalizeStatus(message?.status))) score += 25; + score += Math.min(10, factSeq(message) ?? 0); + return score; +} + +function factMessageHasFinalResponsePart(message, partsByMessageId = new Map()) { + const messageId = textValue(message?.messageId) || textValue(message?.id); + if (!messageId) return false; + return (partsByMessageId.get(messageId) ?? []).some((part) => (textValue(part?.partType ?? part?.type) === "final_response") && Boolean(projectionText(part?.text, part?.content, part?.message))); +} + +function workbenchLifecycleMessageId(traceId, role = "agent") { + const traceSuffix = (safeTraceId(traceId) || String(traceId || "trace")) + .replace(/^trc_/u, "") + .replace(/[^A-Za-z0-9_.:-]/gu, "_") + .slice(0, 48) || "trace"; + return `msg_${traceSuffix}_${role === "user" ? "user" : "agent"}`; +} + function factMessageDto(message, parts = [], facts = {}) { const messageId = safeMessageId(message?.messageId) || textValue(message?.messageId); if (!messageId) return null; diff --git a/internal/cloud/workbench-projection-writer.test.ts b/internal/cloud/workbench-projection-writer.test.ts index fa05122d..17ccaa68 100644 --- a/internal/cloud/workbench-projection-writer.test.ts +++ b/internal/cloud/workbench-projection-writer.test.ts @@ -627,6 +627,85 @@ test("workbench projection event commit writes trace event and checkpoint facts" assert.equal(factWrites[1].params.facts.checkpoints[0].lastEventAt, "2026-06-20T11:02:00.000Z"); }); +test("workbench projection writer reuses lifecycle assistant message id for terminal session projection", async () => { + const runtimeStore = createCloudRuntimeStore({ now: () => "2026-06-20T11:10:00.000Z" }); + const traceId = "trc_writer_terminal_same_message_id"; + const expectedAssistantMessageId = "msg_writer_terminal_same_message_id_agent"; + const accessController = { + async recordAgentSessionOwner(input) { + return { + id: input.sessionId, + projectId: input.projectId, + ownerUserId: input.ownerUserId, + ownerRole: input.ownerRole, + conversationId: input.conversationId, + threadId: input.threadId, + lastTraceId: input.traceId, + status: input.status, + session: input.session, + updatedAt: "2026-06-20T11:10:05.000Z" + }; + } + }; + + await writeWorkbenchProjectionEvent({ + runtimeStore, + event: { + traceId, + sessionId: "ses_writer_terminal_same_message_id", + sourceSeq: 10, + type: "backend", + status: "running", + label: "agentrun:backend:running", + startedAt: "2026-06-20T11:09:55.000Z", + createdAt: "2026-06-20T11:09:55.000Z" + } + }); + await writeWorkbenchProjectionSession({ + accessController, + runtimeStore, + traceId, + ownerUserId: "usr_writer", + ownerRole: "user", + sessionId: "ses_writer_terminal_same_message_id", + projectId: "prj_writer", + conversationId: "cnv_writer_terminal_same_message_id", + threadId: "thread-writer-terminal-same-message-id", + status: "completed", + payload: { + traceId, + status: "completed", + assistantText: "same lifecycle final", + finalResponse: { text: "same lifecycle final", status: "completed", traceId }, + runnerTrace: { + traceId, + status: "completed", + startedAt: "2026-06-20T11:09:55.000Z", + lastEventAt: "2026-06-20T11:10:05.000Z", + finishedAt: "2026-06-20T11:10:05.000Z", + elapsedMs: 10000 + }, + agentRun: { runId: "run_writer_terminal_same_message_id", commandId: "cmd_writer_terminal_same_message_id", status: "completed", terminalStatus: "completed", lastSeq: 11 }, + updatedAt: "2026-06-20T11:10:05.000Z" + }, + session: { + sessionStatus: "completed", + lastTraceId: traceId, + messages: [ + { messageId: "msg_writer_terminal_same_message_id_user", role: "user", text: "same lifecycle", traceId, status: "sent" } + ], + finalResponse: { text: "same lifecycle final", status: "completed", traceId } + } + }); + + const loaded = runtimeStore.queryWorkbenchFacts({ traceId, families: ["messages", "parts", "turns"], limit: 20 }); + const assistantMessages = loaded.facts.messages.filter((message) => message.role !== "user"); + const terminalTurn = loaded.facts.turns.find((turn) => turn.terminal === true); + assert.deepEqual(assistantMessages.map((message) => message.messageId), [expectedAssistantMessageId]); + assert.equal(terminalTurn.messageId, expectedAssistantMessageId); + assert.equal(loaded.facts.parts.some((part) => part.messageId === expectedAssistantMessageId && part.partType === "final_response" && part.text === "same lifecycle final"), true); +}); + test("workbench projection event writer allocates projectedSeq idempotently by source event identity", async () => { const runtimeStore = createCloudRuntimeStore({ now: () => "2026-06-20T12:00:00.000Z" }); diff --git a/internal/cloud/workbench-projection-writer.ts b/internal/cloud/workbench-projection-writer.ts index 40879f4e..147269ca 100644 --- a/internal/cloud/workbench-projection-writer.ts +++ b/internal/cloud/workbench-projection-writer.ts @@ -296,13 +296,23 @@ export async function writeWorkbenchProjectionEvent({ runtimeStore = null, event } function eventProjectionAssistantMessageId(traceId, event = {}) { - const explicit = textValue(event.messageId ?? event.assistantMessageId ?? event.assistantMessage?.messageId); + return eventProjectionMessageId(traceId, "agent", event); +} + +function eventProjectionUserMessageId(traceId, event = {}) { + return eventProjectionMessageId(traceId, "user", event); +} + +function eventProjectionMessageId(traceId, role, event = {}) { + const explicit = role === "user" + ? textValue(event.userMessageId ?? event.userMessage?.messageId) + : textValue(event.messageId ?? event.assistantMessageId ?? event.assistantMessage?.messageId); if (explicit) return explicit; const traceSuffix = (safeTraceId(traceId) || String(traceId || "trace")) .replace(/^trc_/u, "") .replace(/[^A-Za-z0-9_.:-]/gu, "_") .slice(0, 48) || "trace"; - return `msg_${traceSuffix}_agent`; + return `msg_${traceSuffix}_${role === "user" ? "user" : "agent"}`; } async function writeWorkbenchProjectionFacts({ runtimeStore = null, traceStore = defaultCodeAgentTraceStore, traceId = null, ownerUserId = null, ownerRole = null, sessionId = null, projectId = null, conversationId = null, threadId = null, status = "active", session = {}, payload = null, params = {} } = {}) { @@ -607,7 +617,7 @@ function workbenchProjectionInputMessages({ session = {}, payload = {}, params = if (Array.isArray(payload?.messages) && payload.messages.length > 0) return payload.messages; if (!traceId) return []; const turnId = textValue(payload?.turnId) || traceId; - const assistantMessageId = textValue(payload?.assistantMessageId ?? payload?.messageId ?? payload?.assistantMessage?.messageId) || stableFactId("msg", { traceId, role: "agent" }); + const assistantMessageId = eventProjectionAssistantMessageId(traceId, payload); const createdAt = timestampValue(timestamp ?? payload?.createdAt ?? payload?.updatedAt); const messages = []; const userMessage = workbenchProjectionUserMessage({ payload, params, traceId, turnId, timestamp: createdAt }); @@ -636,7 +646,7 @@ function workbenchProjectionUserMessage({ payload = {}, params = {}, traceId = n if (!traceId || !userText) return null; const createdAt = timestampValue(timestamp ?? payload?.createdAt ?? payload?.updatedAt); return { - messageId: textValue(payload?.userMessageId ?? payload?.userMessage?.messageId) || stableFactId("msg", { traceId, role: "user" }), + messageId: eventProjectionUserMessageId(traceId, payload), role: "user", status: "sent", text: userText, @@ -651,7 +661,7 @@ function workbenchProjectionUserMessage({ payload = {}, params = {}, traceId = n function normalizeMessages(messages, context) { const normalized = Array.isArray(messages) ? messages.map((message, index) => normalizeMessageFact(message, index, context)).filter(Boolean) : []; if (context?.terminal && context?.finalText && context?.traceId && !normalized.some((message) => message.role !== "user" && message.traceId === context.traceId)) { - const projected = normalizeMessageFact({ role: "agent", status: context.terminalStatus ?? "completed", text: context.finalText, traceId: context.traceId, source: "workbench-terminal-projection" }, normalized.length, context); + const projected = normalizeMessageFact({ messageId: eventProjectionAssistantMessageId(context.traceId), role: "agent", status: context.terminalStatus ?? "completed", text: context.finalText, traceId: context.traceId, source: "workbench-terminal-projection" }, normalized.length, context); if (projected) normalized.push(projected); } return normalized; @@ -688,9 +698,9 @@ function messagePartFacts(message = {}, { finalText = null } = {}) { } function normalizeMessageFact(message = {}, index, { traceId, sessionId, conversationId, threadId, terminal, terminalStatus = "completed", finalText = null, timestamp, timing: contextTiming }) { - const messageId = textValue(message.messageId ?? message.id) || stableFactId("msg", { traceId, index, role: message.role }); - if (!sessionId || !messageId) return null; const role = textValue(message.role) || "agent"; + const messageId = textValue(message.messageId ?? message.id) || (role === "user" ? eventProjectionUserMessageId(traceId, message) : eventProjectionAssistantMessageId(traceId, message)); + if (!sessionId || !messageId) return null; const messageTraceId = safeTraceId(message.traceId ?? message.turnId); const appliesToContextTrace = role !== "user" && Boolean(traceId && (!messageTraceId || messageTraceId === traceId)); const contextTerminalStatus = normalizeWorkbenchStatus(terminalStatus) || "completed";