Merge pull request #2271 from pikasTech/fix/2255-workbench-freeze
fix: 收敛 Workbench terminal message 后的 turn 状态
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 () => {
|
test("workbench read model does not expose trace-only memory without visible session or result owner", async () => {
|
||||||
const traceStore = createCodeAgentTraceStore();
|
const traceStore = createCodeAgentTraceStore();
|
||||||
const traceId = "trc_workbench_trace_only";
|
const traceId = "trc_workbench_trace_only";
|
||||||
|
|||||||
@@ -951,16 +951,16 @@ function factSessionSummary(session, facts = {}) {
|
|||||||
const checkpoint = traceId ? factCheckpointForTrace(facts, traceId) : null;
|
const checkpoint = traceId ? factCheckpointForTrace(facts, traceId) : null;
|
||||||
const trace = traceId ? factTraceSnapshot(facts, traceId) : null;
|
const trace = traceId ? factTraceSnapshot(facts, traceId) : null;
|
||||||
const messages = factMessagesForSession(session, facts);
|
const messages = factMessagesForSession(session, facts);
|
||||||
|
const terminalMessage = traceId ? factTerminalMessageForTrace(messages, traceId) : null;
|
||||||
const traceStatus = normalizeTerminalStatus(trace?.status);
|
const traceStatus = normalizeTerminalStatus(trace?.status);
|
||||||
const checkpointStatus = normalizeTerminalStatus(checkpoint?.status);
|
const checkpointStatus = normalizeTerminalStatus(checkpoint?.status);
|
||||||
|
const messageTerminalStatus = factTerminalStatusFromMessageProjection(terminalMessage);
|
||||||
const turnStatus = normalizeStatus(turn?.status);
|
const turnStatus = normalizeStatus(turn?.status);
|
||||||
const launchContext = compactLaunchContext(session?.sessionJson?.launchContext);
|
const launchContext = compactLaunchContext(session?.sessionJson?.launchContext);
|
||||||
// Terminal checkpoint/trace evidence is the turn lifecycle authority. A
|
const status = normalizeStatus(checkpointStatus ?? traceStatus ?? messageTerminalStatus ?? turnStatus ?? session?.status);
|
||||||
// stale running turn row must not keep the visible timer running after the
|
const timing = traceId
|
||||||
// terminal event has already been projected.
|
? factCombinedTimingProjection(status, checkpoint, messageTerminalStatus ? terminalMessage : null, turn, trace)
|
||||||
const status = normalizeStatus(checkpointStatus ?? traceStatus ?? turnStatus ?? session?.status);
|
: factTimingProjection(session, status);
|
||||||
const timingSource = checkpointStatus ? checkpoint : turn ?? checkpoint ?? (traceId ? null : session);
|
|
||||||
const timing = factTimingProjection(timingSource, status);
|
|
||||||
const title = sessionTitleFromMessages(messages);
|
const title = sessionTitleFromMessages(messages);
|
||||||
const preview = sessionPreviewFromMessages(messages);
|
const preview = sessionPreviewFromMessages(messages);
|
||||||
return {
|
return {
|
||||||
@@ -1188,15 +1188,16 @@ function factMessageDto(message, parts = [], facts = {}) {
|
|||||||
const checkpoint = traceId ? factCheckpointForTrace(facts, traceId) : null;
|
const checkpoint = traceId ? factCheckpointForTrace(facts, traceId) : null;
|
||||||
const messageStatus = normalizeStatus(message?.status);
|
const messageStatus = normalizeStatus(message?.status);
|
||||||
const assistantLike = isAssistantLikeRole(role);
|
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 legacyText = projectionText(message?.text, message?.content, message?.message, message?.finalResponse);
|
||||||
const normalizedParts = parts.length > 0
|
const normalizedParts = parts.length > 0
|
||||||
? [...parts].sort(compareFactPartsAsc).map((part) => factPartDto(part, messageId, traceId)).filter(Boolean)
|
? [...parts].sort(compareFactPartsAsc).map((part) => factPartDto(part, messageId, traceId)).filter(Boolean)
|
||||||
: !assistantLike && legacyText ? [partFact({ type: "text", text: legacyText, status: message?.status }, 0, messageId, traceId)] : [];
|
: !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 });
|
const text = factMessageAuthorityText({ role, status, parts: normalizedParts, legacyText });
|
||||||
return {
|
return {
|
||||||
messageId,
|
messageId,
|
||||||
@@ -1245,6 +1246,31 @@ function factMessageFinalResponseText(message) {
|
|||||||
return firstFactPartText(message.parts, "final_response");
|
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) {
|
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 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);
|
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 messages = session ? factMessagesForSession(session, facts) : [];
|
||||||
const userMessage = messages.find((message) => message.role === "user") ?? null;
|
const userMessage = messages.find((message) => message.role === "user") ?? null;
|
||||||
const assistantMessage = [...messages].reverse().find((message) => isAssistantLikeRole(message.role)) ?? 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 checkpoint = safeTrace ? factCheckpointForTrace(facts, safeTrace) : null;
|
||||||
const trace = factTraceSnapshot(facts, safeTrace);
|
const trace = factTraceSnapshot(facts, safeTrace);
|
||||||
const checkpointStatus = normalizeTerminalStatus(checkpoint?.status);
|
const checkpointStatus = normalizeTerminalStatus(checkpoint?.status);
|
||||||
const traceStatus = normalizeTerminalStatus(trace.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 terminal = isTerminalProjectionStatus(status);
|
||||||
const assistantText = terminal ? factMessageFinalResponseText(assistantMessage) : null;
|
const assistantText = terminal ? factMessageFinalResponseText(assistantMessage) : null;
|
||||||
const timing = factCombinedTimingProjection(status, checkpoint, turn, trace);
|
const timing = factCombinedTimingProjection(status, checkpoint, messageTerminalStatus ? terminalMessage : null, turn, trace);
|
||||||
return {
|
return {
|
||||||
turnId: resolvedTurnId,
|
turnId: resolvedTurnId,
|
||||||
traceId: safeTrace,
|
traceId: safeTrace,
|
||||||
|
|||||||
@@ -1479,10 +1479,11 @@ function nonBlockingProjection(projection: ProjectionDiagnostic | null): Project
|
|||||||
const traceStatus = normalizedStatusText(trace.status ?? snapshot.status) ?? null;
|
const traceStatus = normalizedStatusText(trace.status ?? snapshot.status) ?? null;
|
||||||
const clearCompletedDiagnostics = shouldClearCompletedTurnDiagnostics(traceStatus, null);
|
const clearCompletedDiagnostics = shouldClearCompletedTurnDiagnostics(traceStatus, null);
|
||||||
const runnerTrace = clearCompletedDiagnostics ? clearRunnerTraceTransientDiagnostics(mergedRunnerTrace) : mergedRunnerTrace;
|
const runnerTrace = clearCompletedDiagnostics ? clearRunnerTraceTransientDiagnostics(mergedRunnerTrace) : mergedRunnerTrace;
|
||||||
|
const terminal = (snapshot as AgentChatResultResponse).terminal === true || isTerminalMessageStatus(traceStatus);
|
||||||
rememberTraceAuthority(runnerTrace);
|
rememberTraceAuthority(runnerTrace);
|
||||||
const error = clearCompletedDiagnostics ? null : message.role === "agent" ? normalizeAgentError(runnerTrace.error ?? message.error) : normalizeAgentError(message.error);
|
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;
|
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) {
|
if (!matchedMessage && ownerSessionId === activeSessionId.value) {
|
||||||
void refreshRealtimeSessionMessages(ownerSessionId, `realtime-trace-gap:${traceId}`);
|
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 error = resultError ?? (terminal || clearCompletedDiagnostics ? null : normalizeAgentError(runnerTrace?.error ?? message.error));
|
||||||
const projection = clearCompletedDiagnostics ? nonBlockingProjection(resultProjection) : resultProjection ?? (terminal ? null : runnerTrace.projection ?? message.projection ?? null);
|
const projection = clearCompletedDiagnostics ? nonBlockingProjection(resultProjection) : resultProjection ?? (terminal ? null : runnerTrace.projection ?? message.projection ?? null);
|
||||||
const agentRun = agentRunFromResult(result, runnerTrace) ?? agentRunFromMessage(message);
|
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);
|
rememberTurnStatus(traceId, result);
|
||||||
markWorkbenchTraceProjected(traceId);
|
markWorkbenchTraceProjected(traceId);
|
||||||
|
|||||||
Reference in New Issue
Block a user