fix: stabilize workbench assistant projection identity
This commit is contained in:
@@ -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();
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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" });
|
||||
|
||||
|
||||
@@ -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";
|
||||
|
||||
Reference in New Issue
Block a user