fix: seal workbench terminal message turns
This commit is contained in:
@@ -1177,6 +1177,69 @@ test("workbench read model projects terminal result atomically across session, m
|
||||
}
|
||||
});
|
||||
|
||||
test("workbench read model seals turn from terminal message projection when checkpoint is stale", async () => {
|
||||
const sessionId = "ses_workbench_terminal_message_stale_checkpoint";
|
||||
const traceId = "trc_terminal_message_stale_checkpoint";
|
||||
const finalText = "terminal message projection must close the visible turn";
|
||||
const startedAt = "2026-06-29T22:09:32.000Z";
|
||||
const finishedAt = "2026-06-29T22:09:47.000Z";
|
||||
const facts = emptyFacts();
|
||||
facts.sessions.push({
|
||||
sessionId,
|
||||
ownerUserId: ACTOR.id,
|
||||
ownerRole: ACTOR.role,
|
||||
agentId: "hwlab-code-agent",
|
||||
status: "running",
|
||||
lastTraceId: traceId,
|
||||
projectedSeq: 33,
|
||||
sourceSeq: 36,
|
||||
createdAt: startedAt,
|
||||
updatedAt: finishedAt,
|
||||
valuesRedacted: true
|
||||
});
|
||||
facts.messages.push(
|
||||
{ messageId: "msg_stale_user", sessionId, turnId: traceId, traceId, role: "user", status: "sent", text: "hi", projectedSeq: 31, sourceSeq: 31, createdAt: startedAt, updatedAt: startedAt, valuesRedacted: true },
|
||||
{ messageId: "msg_stale_agent", sessionId, turnId: traceId, traceId, role: "agent", status: "completed", text: finalText, projectedSeq: 36, sourceSeq: 36, terminal: true, sealed: true, timing: { startedAt, lastEventAt: finishedAt, finishedAt, durationMs: 15000, valuesRedacted: true }, startedAt, lastEventAt: finishedAt, finishedAt, durationMs: 15000, createdAt: startedAt, updatedAt: finishedAt, valuesRedacted: true }
|
||||
);
|
||||
facts.parts.push(
|
||||
{ partId: "prt_stale_user", messageId: "msg_stale_user", sessionId, turnId: traceId, traceId, partIndex: 0, partType: "text", status: "sent", text: "hi", projectedSeq: 31, sourceSeq: 31, updatedAt: startedAt, valuesRedacted: true },
|
||||
{ partId: "prt_stale_agent", messageId: "msg_stale_agent", sessionId, turnId: traceId, traceId, partIndex: 0, partType: "final_response", status: "completed", text: finalText, projectedSeq: 36, sourceSeq: 36, terminal: true, sealed: true, updatedAt: finishedAt, valuesRedacted: true }
|
||||
);
|
||||
facts.turns.push({ turnId: traceId, sessionId, traceId, messageId: "msg_stale_agent", status: "running", projectedSeq: 33, sourceSeq: 33, terminal: false, sealed: false, startedAt, lastEventAt: finishedAt, updatedAt: finishedAt, valuesRedacted: true });
|
||||
facts.checkpoints.push({ traceId, sessionId, turnId: traceId, runId: "run_stale_checkpoint", commandId: "cmd_stale_checkpoint", status: "running", projectionStatus: "projecting", projectionHealth: "degraded", projectedSeq: 33, sourceSeq: 33, terminal: false, sealed: false, startedAt, lastEventAt: finishedAt, 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 sessions = await getJson(port, `/v1/workbench/sessions?includeSessionId=${encodeURIComponent(sessionId)}`);
|
||||
assert.equal(sessions.status, 200);
|
||||
assert.equal(sessions.body.sessions[0].status, "completed");
|
||||
assert.equal(sessions.body.sessions[0].running, false);
|
||||
assert.equal(sessions.body.sessions[0].terminal, true);
|
||||
|
||||
const messages = await getJson(port, `/v1/workbench/sessions/${encodeURIComponent(sessionId)}/messages?limit=10`);
|
||||
assert.equal(messages.status, 200);
|
||||
const assistant = messages.body.messages.find((message) => message.role === "agent");
|
||||
assert.equal(assistant.status, "completed");
|
||||
assert.equal(assistant.text, finalText);
|
||||
|
||||
const turn = await getJson(port, `/v1/workbench/turns/${encodeURIComponent(traceId)}`);
|
||||
assert.equal(turn.status, 200);
|
||||
assert.equal(turn.body.turn.status, "completed");
|
||||
assert.equal(turn.body.turn.running, false);
|
||||
assert.equal(turn.body.turn.terminal, true);
|
||||
assert.equal(turn.body.turn.finalResponse.text, finalText);
|
||||
} finally {
|
||||
await new Promise((resolve, reject) => server.close((error) => error ? reject(error) : resolve()));
|
||||
}
|
||||
});
|
||||
|
||||
test("workbench read model does not expose trace-only memory without visible session or result owner", async () => {
|
||||
const traceStore = createCodeAgentTraceStore();
|
||||
const traceId = "trc_workbench_trace_only";
|
||||
|
||||
@@ -951,16 +951,16 @@ function factSessionSummary(session, facts = {}) {
|
||||
const checkpoint = traceId ? factCheckpointForTrace(facts, traceId) : null;
|
||||
const trace = traceId ? factTraceSnapshot(facts, traceId) : null;
|
||||
const messages = factMessagesForSession(session, facts);
|
||||
const terminalMessage = traceId ? factTerminalMessageForTrace(messages, traceId) : null;
|
||||
const traceStatus = normalizeTerminalStatus(trace?.status);
|
||||
const checkpointStatus = normalizeTerminalStatus(checkpoint?.status);
|
||||
const messageTerminalStatus = factTerminalStatusFromMessageProjection(terminalMessage);
|
||||
const turnStatus = normalizeStatus(turn?.status);
|
||||
const launchContext = compactLaunchContext(session?.sessionJson?.launchContext);
|
||||
// Terminal checkpoint/trace evidence is the turn lifecycle authority. A
|
||||
// stale running turn row must not keep the visible timer running after the
|
||||
// terminal event has already been projected.
|
||||
const status = normalizeStatus(checkpointStatus ?? traceStatus ?? turnStatus ?? session?.status);
|
||||
const timingSource = checkpointStatus ? checkpoint : turn ?? checkpoint ?? (traceId ? null : session);
|
||||
const timing = factTimingProjection(timingSource, status);
|
||||
const status = normalizeStatus(checkpointStatus ?? traceStatus ?? messageTerminalStatus ?? turnStatus ?? session?.status);
|
||||
const timing = traceId
|
||||
? factCombinedTimingProjection(status, checkpoint, messageTerminalStatus ? terminalMessage : null, turn, trace)
|
||||
: factTimingProjection(session, status);
|
||||
const title = sessionTitleFromMessages(messages);
|
||||
const preview = sessionPreviewFromMessages(messages);
|
||||
return {
|
||||
@@ -1188,15 +1188,16 @@ function factMessageDto(message, parts = [], facts = {}) {
|
||||
const checkpoint = traceId ? factCheckpointForTrace(facts, traceId) : null;
|
||||
const messageStatus = normalizeStatus(message?.status);
|
||||
const assistantLike = isAssistantLikeRole(role);
|
||||
const checkpointStatus = assistantLike ? normalizeTerminalStatus(checkpoint?.status) : null;
|
||||
const status = normalizeStatus(checkpointStatus ?? (assistantLike ? turn?.status : null) ?? messageStatus);
|
||||
const timing = assistantLike
|
||||
? factCombinedTimingProjection(status, checkpoint, turn, message)
|
||||
: factTimingProjection(message, status);
|
||||
const legacyText = projectionText(message?.text, message?.content, message?.message, message?.finalResponse);
|
||||
const normalizedParts = parts.length > 0
|
||||
? [...parts].sort(compareFactPartsAsc).map((part) => factPartDto(part, messageId, traceId)).filter(Boolean)
|
||||
: !assistantLike && legacyText ? [partFact({ type: "text", text: legacyText, status: message?.status }, 0, messageId, traceId)] : [];
|
||||
const checkpointStatus = assistantLike ? normalizeTerminalStatus(checkpoint?.status) : null;
|
||||
const messageTerminalStatus = assistantLike ? factTerminalStatusFromMessageProjection(message, normalizedParts) : null;
|
||||
const status = normalizeStatus(checkpointStatus ?? messageTerminalStatus ?? (assistantLike ? turn?.status : null) ?? messageStatus);
|
||||
const timing = assistantLike
|
||||
? factCombinedTimingProjection(status, checkpoint, messageTerminalStatus ? { ...message, parts: normalizedParts } : null, turn, message)
|
||||
: factTimingProjection(message, status);
|
||||
const text = factMessageAuthorityText({ role, status, parts: normalizedParts, legacyText });
|
||||
return {
|
||||
messageId,
|
||||
@@ -1245,6 +1246,31 @@ function factMessageFinalResponseText(message) {
|
||||
return firstFactPartText(message.parts, "final_response");
|
||||
}
|
||||
|
||||
function factTerminalMessageForTrace(messages = [], traceId = null) {
|
||||
const safeTrace = safeTraceId(traceId);
|
||||
return [...factArray(messages)].reverse().find((message) => {
|
||||
if (!message || !isAssistantLikeRole(message.role)) return false;
|
||||
if (safeTrace && message.traceId !== safeTrace) return false;
|
||||
return Boolean(factTerminalStatusFromMessageProjection(message));
|
||||
}) ?? null;
|
||||
}
|
||||
|
||||
function factTerminalStatusFromMessageProjection(message = null, parts = null) {
|
||||
if (!message || !isAssistantLikeRole(message.role)) return null;
|
||||
const sourceParts = Array.isArray(parts) ? parts : factArray(message.parts);
|
||||
if (!firstFactPartText(sourceParts, "final_response")) return null;
|
||||
const messageStatus = normalizeTerminalStatus(message.status);
|
||||
if (messageStatus) return messageStatus;
|
||||
for (const part of sourceParts) {
|
||||
if (textValue(part?.partType ?? part?.type) !== "final_response") continue;
|
||||
if (!projectionText(part?.text, part?.content, part?.message)) continue;
|
||||
const partStatus = normalizeTerminalStatus(part?.status);
|
||||
if (partStatus) return partStatus;
|
||||
if (part?.terminal === true || part?.sealed === true) return "completed";
|
||||
}
|
||||
return message.terminal === true || message.sealed === true ? "completed" : null;
|
||||
}
|
||||
|
||||
function factPartDto(part, messageId, traceId) {
|
||||
const partId = safePartId(part?.partId) || textValue(part?.partId) || `prt_${hash(`${messageId}:${part?.partIndex ?? 0}:${part?.partType ?? part?.type ?? "text"}`).slice(0, 24)}`;
|
||||
const text = projectionText(part?.text, part?.content, part?.message);
|
||||
@@ -1277,14 +1303,16 @@ function factTurnSnapshot({ turn = null, session = null, facts = {}, traceId, tu
|
||||
const messages = session ? factMessagesForSession(session, facts) : [];
|
||||
const userMessage = messages.find((message) => message.role === "user") ?? null;
|
||||
const assistantMessage = [...messages].reverse().find((message) => isAssistantLikeRole(message.role)) ?? null;
|
||||
const terminalMessage = safeTrace ? factTerminalMessageForTrace(messages, safeTrace) : null;
|
||||
const checkpoint = safeTrace ? factCheckpointForTrace(facts, safeTrace) : null;
|
||||
const trace = factTraceSnapshot(facts, safeTrace);
|
||||
const checkpointStatus = normalizeTerminalStatus(checkpoint?.status);
|
||||
const traceStatus = normalizeTerminalStatus(trace.status);
|
||||
const status = normalizeStatus(checkpointStatus ?? traceStatus ?? turn?.status ?? session?.status);
|
||||
const messageTerminalStatus = factTerminalStatusFromMessageProjection(terminalMessage);
|
||||
const status = normalizeStatus(checkpointStatus ?? traceStatus ?? messageTerminalStatus ?? turn?.status ?? session?.status);
|
||||
const terminal = isTerminalProjectionStatus(status);
|
||||
const assistantText = terminal ? factMessageFinalResponseText(assistantMessage) : null;
|
||||
const timing = factCombinedTimingProjection(status, checkpoint, turn, trace);
|
||||
const timing = factCombinedTimingProjection(status, checkpoint, messageTerminalStatus ? terminalMessage : null, turn, trace);
|
||||
return {
|
||||
turnId: resolvedTurnId,
|
||||
traceId: safeTrace,
|
||||
|
||||
@@ -1479,10 +1479,11 @@ function nonBlockingProjection(projection: ProjectionDiagnostic | null): Project
|
||||
const traceStatus = normalizedStatusText(trace.status ?? snapshot.status) ?? null;
|
||||
const clearCompletedDiagnostics = shouldClearCompletedTurnDiagnostics(traceStatus, null);
|
||||
const runnerTrace = clearCompletedDiagnostics ? clearRunnerTraceTransientDiagnostics(mergedRunnerTrace) : mergedRunnerTrace;
|
||||
const terminal = (snapshot as AgentChatResultResponse).terminal === true || isTerminalMessageStatus(traceStatus);
|
||||
rememberTraceAuthority(runnerTrace);
|
||||
const error = clearCompletedDiagnostics ? null : message.role === "agent" ? normalizeAgentError(runnerTrace.error ?? message.error) : normalizeAgentError(message.error);
|
||||
const projection = clearCompletedDiagnostics ? nonBlockingProjection(trace.projection ?? null) : trace.projection ?? runnerTrace.projection ?? message.projection ?? null;
|
||||
return { ...message, ...messageTimingPatchForMerge(message, trace), runnerTrace, error: clearCompletedDiagnostics ? null : error ?? message.error ?? null, projection, projectionStatus: projection?.projectionStatus ?? null, projectionHealth: projection?.projectionHealth ?? null, blocker: projection?.blocker ?? null, updatedAt: new Date().toISOString() };
|
||||
return { ...message, ...messageTimingPatchForMerge(message, trace), ...messageStatusPatchForTerminalMerge(message, traceStatus, terminal), runnerTrace, error: clearCompletedDiagnostics ? null : error ?? message.error ?? null, projection, projectionStatus: projection?.projectionStatus ?? null, projectionHealth: projection?.projectionHealth ?? null, blocker: projection?.blocker ?? null, updatedAt: new Date().toISOString() };
|
||||
}));
|
||||
if (!matchedMessage && ownerSessionId === activeSessionId.value) {
|
||||
void refreshRealtimeSessionMessages(ownerSessionId, `realtime-trace-gap:${traceId}`);
|
||||
@@ -1524,7 +1525,7 @@ function nonBlockingProjection(projection: ProjectionDiagnostic | null): Project
|
||||
const error = resultError ?? (terminal || clearCompletedDiagnostics ? null : normalizeAgentError(runnerTrace?.error ?? message.error));
|
||||
const projection = clearCompletedDiagnostics ? nonBlockingProjection(resultProjection) : resultProjection ?? (terminal ? null : runnerTrace.projection ?? message.projection ?? null);
|
||||
const agentRun = agentRunFromResult(result, runnerTrace) ?? agentRunFromMessage(message);
|
||||
return { ...message, ...messageTimingPatchForMerge(message, result), runnerTrace, error, projection, projectionStatus: projection?.projectionStatus ?? null, projectionHealth: projection?.projectionHealth ?? null, blocker: projection?.blocker ?? null, agentRun: agentRun ?? undefined, updatedAt: new Date().toISOString() };
|
||||
return { ...message, ...messageTimingPatchForMerge(message, result), ...messageStatusPatchForTerminalMerge(message, resultStatus, terminal), runnerTrace, error, projection, projectionStatus: projection?.projectionStatus ?? null, projectionHealth: projection?.projectionHealth ?? null, blocker: projection?.blocker ?? null, agentRun: agentRun ?? undefined, updatedAt: new Date().toISOString() };
|
||||
}));
|
||||
rememberTurnStatus(traceId, result);
|
||||
markWorkbenchTraceProjected(traceId);
|
||||
|
||||
Reference in New Issue
Block a user