fix: enforce workbench terminal seal invariant
This commit is contained in:
@@ -21,7 +21,7 @@ Workbench 投影写路径必须做到 0 隐式 fallback。Admission、projection
|
||||
|
||||
Workbench 的 trace/message/projection 运行时必须是独立模块边界。`workbench.ts` 只保留 session/route authority 校验、Pinia state commit、刷新调度和用户动作编排;trace snapshot、terminal result、message timing/status patch、agent error normalize、projection diagnostic 裁剪和 final response 文本提取等纯算法统一由 `web/hwlab-cloud-web/src/stores/workbench-message-projection-runtime.ts` 提供。新增或修复浏览器 smoke 时不得把这些 helper 重新散写回 store 或组件,也不得通过 reload、repair、localStorage truth、GET read-through、测试专用后门或删除 guard 来绕过真实投影问题。web-probe origin、视口、采样、命令超时、provider/lane 和报警阈值只从选中 node/lane 的受控 YAML/source-of-truth 进入验证命令,不在 SPEC 或前端 runtime 中写第二份数值。
|
||||
|
||||
Workbench terminal 三态必须作为一个不可拆分的投影不变量维护:session rail 状态、turn 卡片状态和 final response object 要么同时表现为完成且 final response 存在,要么同时表现为运行且 final response 不存在。持久 read model 是 session、messages、turn、trace tuple 的 source of truth;前端 `workbench-server-state` 只能做归一化缓存和单调合并,不能把 stale running 刷新覆盖到同一 trace 的 terminal authority 上。`session.status`、`turn.status`、`session.list/detail/messages`、`message.snapshot`、projection page merge、REST trace hydration 和 SSE trace snapshot 都必须保留已 sealed terminal 的 `status`、`traceAutoLifecycle`、`text`、`finalResponse` 和 terminal timing;新 trace 的 running 可以让同一 session 进入下一轮运行,但同 trace 的 running/non-terminal 只能作为旧事件丢弃或合并为 runnerTrace 证据。判断 final response 是否存在时不得只看 `text`,`finalResponse` 对象本体也属于 Workbench read model 合同;新增修复应优先补 `workbench-message-projection-runtime` 或 `workbench-server-state` 的最小单元测试,而不是用 UI 特判、reload 或额外 fallback 掩盖投影分叉。
|
||||
Workbench terminal 三态必须作为一个不可拆分的投影不变量维护:session rail 状态、turn 卡片状态和 final response object 要么同时表现为完成且 final response 存在,要么同时表现为运行且 final response 不存在。持久 read model 是 session、messages、turn、trace tuple 的 source of truth;前端 `workbench-server-state` 只能做归一化缓存和单调合并,不能把 stale running 刷新覆盖到同一 trace 的 terminal authority 上。`session.status`、`turn.status`、`session.list/detail/messages`、`message.snapshot`、projection page merge、REST trace hydration 和 SSE trace snapshot 都必须保留已 sealed terminal 的 `status`、`traceAutoLifecycle`、`text`、`finalResponse` 和 terminal timing;新 trace 的 running 可以让同一 session 进入下一轮运行,但同 trace 的 running/non-terminal 只能作为旧事件丢弃或合并为 runnerTrace 证据。判断 final response 是否存在时必须能从 `text`、`content`、`finalResponse` 或 `reply/finalText` 等权威字段提取非空文本,空对象、进度 assistant trace 文本或仅有 AgentRun completed 状态都不能 seal completed;`/v1/agent/turns/:traceId`、Workbench read model、projection writer、runtime store invariant 和前端缓存层必须共用这个 terminal seal gate。新增修复应优先补 `workbench-message-projection-runtime`、`workbench-server-state` 或对应后端 read-model 的最小单元测试,而不是用 UI 特判、reload 或额外 fallback 掩盖投影分叉。
|
||||
|
||||
MDTODO 发起 Workbench 执行时,HWPOD 执行上下文的唯一权威来源是 Project Management source registry。Workbench Launch 服务端必须通过 `taskRef -> sourceId/fileRef -> source` 解析 `launchContext.executionContext`,并把同一份 `contextFingerprint` 写入 session owner、Workbench facts、project-management link 和 OTel span;浏览器传入的 HWPOD 字段只能作为任务元数据,不能作为权威执行上下文。`sourceKind=hwpod-workspace` 但缺少 `hwpodId`、`nodeId` 或 `workspaceRootRef` 时,launch 必须显式失败,不能创建“看似成功但无法执行”的空 session。
|
||||
|
||||
|
||||
@@ -2598,7 +2598,7 @@ async function getCodeAgentSessionByTraceId(traceId, options = {}) {
|
||||
return await options.accessController.getAgentSessionByTraceId(safeId);
|
||||
}
|
||||
|
||||
function codeAgentTurnStatusPayload({ traceId, result, snapshot, resultPollError, refreshError, options }) {
|
||||
export function codeAgentTurnStatusPayload({ traceId, result, snapshot, resultPollError, refreshError, options }) {
|
||||
const resultObject = result && typeof result === "object" ? result : null;
|
||||
const snapshotObject = snapshot && typeof snapshot === "object" ? snapshot : null;
|
||||
const lifecycle = codeAgentTurnLifecycleFields(traceId, resultObject ?? snapshotObject ?? {});
|
||||
@@ -2606,7 +2606,8 @@ function codeAgentTurnStatusPayload({ traceId, result, snapshot, resultPollError
|
||||
const lastEvent = events.at(-1) ?? null;
|
||||
const finalResponse = resultObject?.finalResponse ?? snapshotObject?.finalResponse ?? snapshotObject?.terminalEvidence?.finalResponse ?? codeAgentFinalResponseEvidence(resultObject ?? snapshotObject ?? {}, traceId);
|
||||
const terminalStatus = codeAgentAuthoritativeTerminalStatus(resultObject, snapshotObject, traceId);
|
||||
const status = normalizeTurnStatus(
|
||||
const terminalSealBlocked = terminalStatus === "completed" && !codeAgentCompletedTurnFinalText(finalResponse, resultObject, snapshotObject);
|
||||
const rawStatus = normalizeTurnStatus(
|
||||
terminalStatus,
|
||||
resultObject?.agentRun?.commandState,
|
||||
resultObject?.agentRun?.status,
|
||||
@@ -2618,6 +2619,7 @@ function codeAgentTurnStatusPayload({ traceId, result, snapshot, resultPollError
|
||||
resultObject?.agentRun?.terminalStatus,
|
||||
snapshotObject?.terminalEvidence?.traceSummary?.terminalStatus
|
||||
);
|
||||
const status = terminalSealBlocked ? "running" : rawStatus;
|
||||
const found = Boolean(resultObject || (snapshotObject && snapshotObject.status !== "missing") || snapshotObject?.persisted === true);
|
||||
const running = isTurnRunningStatus(status);
|
||||
const terminal = isTurnTerminalStatus(status);
|
||||
@@ -2638,7 +2640,10 @@ function codeAgentTurnStatusPayload({ traceId, result, snapshot, resultPollError
|
||||
updatedAt: textValue(resultObject?.updatedAt ?? resultObject?.agentRun?.updatedAt ?? snapshotObject?.updatedAt ?? timingFields.lastEventAt ?? timingFields.startedAt) || null,
|
||||
...timingFields,
|
||||
lastEventLabel: textValue(snapshotObject?.lastEventLabel ?? runnerTrace?.lastEventLabel ?? lastEvent?.label ?? lastEvent?.type) || null,
|
||||
waitingFor: textValue(snapshotObject?.waitingFor ?? runnerTrace?.waitingFor) || null,
|
||||
waitingFor: terminalSealBlocked ? "final_response" : textValue(snapshotObject?.waitingFor ?? runnerTrace?.waitingFor) || null,
|
||||
terminalObserved: Boolean(terminalStatus),
|
||||
terminalObservedStatus: terminalStatus ?? null,
|
||||
terminalSealBlocked,
|
||||
resultUrl: `/v1/agent/chat/result/${encodeURIComponent(traceId)}`,
|
||||
traceUrl: `/v1/agent/traces/${encodeURIComponent(traceId)}`,
|
||||
turnUrl: `/v1/agent/turns/${encodeURIComponent(traceId)}`,
|
||||
@@ -2657,6 +2662,21 @@ function codeAgentTurnStatusPayload({ traceId, result, snapshot, resultPollError
|
||||
};
|
||||
}
|
||||
|
||||
function codeAgentCompletedTurnFinalText(finalResponse = null, resultObject = null, snapshotObject = null) {
|
||||
const candidates = [
|
||||
finalResponse,
|
||||
resultObject?.finalResponse,
|
||||
snapshotObject?.finalResponse,
|
||||
snapshotObject?.terminalEvidence?.finalResponse,
|
||||
resultObject?.terminalEvidence?.finalResponse
|
||||
];
|
||||
for (const value of candidates) {
|
||||
const text = conversationText(value);
|
||||
if (text) return text;
|
||||
}
|
||||
return codeAgentFinalResponseText(resultObject ?? {}) || codeAgentFinalResponseText(snapshotObject ?? {});
|
||||
}
|
||||
|
||||
function codeAgentAuthoritativeTerminalStatus(resultObject, snapshotObject, traceId) {
|
||||
const sealedPayload = codeAgentPayloadHasSealedFinalResponse(resultObject ?? {}) ? resultObject : codeAgentPayloadHasSealedFinalResponse(snapshotObject ?? {}) ? snapshotObject : null;
|
||||
if (sealedPayload) return "completed";
|
||||
@@ -3971,7 +3991,23 @@ function codeAgentFinalResponseEvidence(payload = {}, traceId = null) {
|
||||
}
|
||||
|
||||
function codeAgentFinalResponseText(payload = {}) {
|
||||
return conversationText(payload?.finalResponse?.text ?? payload?.finalResponse?.content ?? payload?.reply?.content ?? payload?.message?.content ?? payload?.assistantText);
|
||||
for (const value of [
|
||||
payload?.finalResponse?.text,
|
||||
payload?.finalResponse?.content,
|
||||
payload?.finalResponse,
|
||||
payload?.reply?.content,
|
||||
payload?.reply,
|
||||
payload?.message?.content,
|
||||
payload?.message,
|
||||
payload?.assistantText,
|
||||
payload?.finalText,
|
||||
payload?.text,
|
||||
payload?.content
|
||||
]) {
|
||||
const text = conversationText(value);
|
||||
if (text) return text;
|
||||
}
|
||||
return "";
|
||||
}
|
||||
|
||||
function codeAgentTraceSummaryEvidence(payload = {}, traceId = null, finalResponse = null) {
|
||||
|
||||
@@ -8,7 +8,7 @@ import { createCloudApiServer } from "./server.ts";
|
||||
import { createBackendPerformanceStore } from "./backend-performance.ts";
|
||||
import { createCodeAgentTraceStore } from "./code-agent-trace-store.ts";
|
||||
import { classifyWorkbenchReadModelFailure } from "./server-workbench-http.ts";
|
||||
import { createCodeAgentChatResultStore } from "./server-code-agent-http.ts";
|
||||
import { codeAgentTurnStatusPayload, createCodeAgentChatResultStore } from "./server-code-agent-http.ts";
|
||||
import { createWorkbenchTurnProjection, durableTraceStatus, projectionDiagnostics, traceTerminalEvidence } from "./workbench-turn-projection.ts";
|
||||
|
||||
const ACTOR = { id: "usr_workbench_reader", username: "reader", displayName: "Reader", role: "user", status: "active" };
|
||||
@@ -238,10 +238,82 @@ test("workbench turn projection keeps progress-only assistant trace text out of
|
||||
eventCount: 2
|
||||
};
|
||||
const projection = createWorkbenchTurnProjection({ traceId, result: { traceId, status: "completed" }, trace });
|
||||
assert.equal(projection.status, "completed");
|
||||
assert.equal(projection.terminal, true);
|
||||
const diagnostics = projectionDiagnostics({ traceId, result: { traceId, status: "completed" }, trace, projection });
|
||||
assert.equal(projection.status, "running");
|
||||
assert.equal(projection.running, true);
|
||||
assert.equal(projection.terminal, false);
|
||||
assert.equal(projection.terminalObserved, true);
|
||||
assert.equal(projection.terminalObservedStatus, "completed");
|
||||
assert.equal(projection.waitingFor, "final_response");
|
||||
assert.equal(projection.finalResponse, null);
|
||||
assert.equal(projection.assistantText, null);
|
||||
assert.equal(diagnostics.projectionStatus, "projecting");
|
||||
assert.equal(diagnostics.projectionHealth, "projecting");
|
||||
assert.equal(diagnostics.waitingFor, "final_response");
|
||||
});
|
||||
|
||||
test("code agent turn status keeps completed AgentRun without final response unsealed", () => {
|
||||
const traceId = "trc_code_agent_terminal_without_final";
|
||||
const payload = codeAgentTurnStatusPayload({
|
||||
traceId,
|
||||
result: {
|
||||
traceId,
|
||||
status: "completed",
|
||||
agentRun: {
|
||||
runId: "run_code_agent_terminal_without_final",
|
||||
commandId: "cmd_code_agent_terminal_without_final",
|
||||
status: "completed",
|
||||
terminalStatus: "completed"
|
||||
},
|
||||
runnerTrace: {
|
||||
traceId,
|
||||
status: "completed",
|
||||
events: [
|
||||
{ seq: 1, type: "assistant_message", status: "running", message: "progress only" },
|
||||
{ seq: 2, type: "terminal_status", terminal: true, terminalStatus: "completed" }
|
||||
]
|
||||
}
|
||||
},
|
||||
snapshot: null,
|
||||
resultPollError: null,
|
||||
refreshError: null,
|
||||
options: { env: {} }
|
||||
});
|
||||
assert.equal(payload.status, "running");
|
||||
assert.equal(payload.running, true);
|
||||
assert.equal(payload.terminal, false);
|
||||
assert.equal(payload.terminalObserved, true);
|
||||
assert.equal(payload.terminalObservedStatus, "completed");
|
||||
assert.equal(payload.terminalSealBlocked, true);
|
||||
assert.equal(payload.waitingFor, "final_response");
|
||||
assert.equal(payload.finalResponse, null);
|
||||
});
|
||||
|
||||
test("code agent turn status seals completed AgentRun when finalText is authoritative", () => {
|
||||
const traceId = "trc_code_agent_terminal_with_final_text";
|
||||
const payload = codeAgentTurnStatusPayload({
|
||||
traceId,
|
||||
result: {
|
||||
traceId,
|
||||
status: "completed",
|
||||
finalText: "authoritative final",
|
||||
agentRun: {
|
||||
runId: "run_code_agent_terminal_with_final_text",
|
||||
commandId: "cmd_code_agent_terminal_with_final_text",
|
||||
status: "completed",
|
||||
terminalStatus: "completed"
|
||||
}
|
||||
},
|
||||
snapshot: null,
|
||||
resultPollError: null,
|
||||
refreshError: null,
|
||||
options: { env: {} }
|
||||
});
|
||||
assert.equal(payload.status, "completed");
|
||||
assert.equal(payload.running, false);
|
||||
assert.equal(payload.terminal, true);
|
||||
assert.equal(payload.terminalSealBlocked, false);
|
||||
assert.equal(payload.finalResponse.text, "authoritative final");
|
||||
});
|
||||
|
||||
test("workbench trace event page exposes projectedSeq cursor range from durable facts", async () => {
|
||||
|
||||
@@ -190,6 +190,89 @@ test("workbench projection writer commits terminal owner evidence as sealed dura
|
||||
assert.equal(facts.checkpoints[0].timing.finishedAt, "2026-06-20T11:00:00.000Z");
|
||||
});
|
||||
|
||||
test("workbench projection writer keeps completed AgentRun without final response unsealed", async () => {
|
||||
const factWrites = [];
|
||||
const runtimeStore = {
|
||||
async writeWorkbenchFacts(params, requestMeta) {
|
||||
factWrites.push({ params, requestMeta });
|
||||
return { written: true, facts: params.facts };
|
||||
}
|
||||
};
|
||||
const accessController = {
|
||||
async recordAgentSessionOwner(input) {
|
||||
return {
|
||||
id: input.sessionId,
|
||||
projectId: input.projectId,
|
||||
ownerUserId: input.ownerUserId,
|
||||
conversationId: input.conversationId,
|
||||
threadId: input.threadId,
|
||||
lastTraceId: input.traceId,
|
||||
status: input.status,
|
||||
session: input.session,
|
||||
updatedAt: "2026-06-20T11:02:00.000Z"
|
||||
};
|
||||
}
|
||||
};
|
||||
|
||||
await writeWorkbenchProjectionSession({
|
||||
accessController,
|
||||
runtimeStore,
|
||||
traceId: "trc_writer_terminal_missing_final",
|
||||
ownerUserId: "usr_writer",
|
||||
ownerRole: "user",
|
||||
sessionId: "ses_writer_terminal_missing_final",
|
||||
projectId: "prj_writer",
|
||||
conversationId: "cnv_writer_terminal_missing_final",
|
||||
threadId: "thread-writer-terminal-missing-final",
|
||||
status: "completed",
|
||||
payload: {
|
||||
traceId: "trc_writer_terminal_missing_final",
|
||||
status: "completed",
|
||||
runnerTrace: {
|
||||
traceId: "trc_writer_terminal_missing_final",
|
||||
status: "completed",
|
||||
startedAt: "2026-06-20T11:01:30.000Z",
|
||||
lastEventAt: "2026-06-20T11:02:00.000Z",
|
||||
finishedAt: "2026-06-20T11:02:00.000Z",
|
||||
events: [
|
||||
{ seq: 1, type: "assistant_message", status: "running", message: "progress only" },
|
||||
{ seq: 2, type: "result", status: "completed", terminal: true, label: "agentrun:terminal:completed" }
|
||||
],
|
||||
eventCount: 2
|
||||
},
|
||||
agentRun: { runId: "run_writer_terminal_missing_final", commandId: "cmd_writer_terminal_missing_final", status: "completed", terminalStatus: "completed", lastSeq: 2 },
|
||||
updatedAt: "2026-06-20T11:02:00.000Z"
|
||||
},
|
||||
session: {
|
||||
sessionStatus: "completed",
|
||||
messages: [
|
||||
{ messageId: "msg_writer_missing_final_user", role: "user", text: "question", status: "sent", turnId: "trc_writer_terminal_missing_final", traceId: "trc_writer_terminal_missing_final" },
|
||||
{ messageId: "msg_writer_missing_final_agent", role: "agent", text: "", status: "completed", turnId: "trc_writer_terminal_missing_final", traceId: "trc_writer_terminal_missing_final" }
|
||||
]
|
||||
}
|
||||
});
|
||||
|
||||
assert.equal(factWrites.length, 1);
|
||||
const facts = factWrites[0].params.facts;
|
||||
const agentMessage = facts.messages.find((message) => message.messageId === "msg_writer_missing_final_agent");
|
||||
assert.equal(facts.sessions[0].status, "running");
|
||||
assert.equal(facts.sessions[0].terminal, false);
|
||||
assert.equal(facts.sessions[0].sealed, false);
|
||||
assert.equal(agentMessage.status, "running");
|
||||
assert.equal(agentMessage.terminal, false);
|
||||
assert.equal(agentMessage.sealed, false);
|
||||
assert.equal(facts.parts.some((part) => part.partType === "final_response"), false);
|
||||
assert.equal(facts.turns[0].status, "running");
|
||||
assert.equal(facts.turns[0].terminal, false);
|
||||
assert.equal(facts.turns[0].sealed, false);
|
||||
assert.equal(facts.turns[0].finalResponse, null);
|
||||
assert.equal(facts.turns[0].diagnostic.projectionStatus, "projecting");
|
||||
assert.equal(facts.turns[0].diagnostic.waitingFor, "final_response");
|
||||
assert.equal(facts.checkpoints[0].projectionStatus, "projecting");
|
||||
assert.equal(facts.checkpoints[0].terminal, false);
|
||||
assert.equal(facts.checkpoints[0].sealed, false);
|
||||
});
|
||||
|
||||
test("workbench projection writer does not synthesize terminal duration from updatedAt", async () => {
|
||||
const factWrites = [];
|
||||
const runtimeStore = {
|
||||
|
||||
@@ -528,6 +528,7 @@ function buildWorkbenchProjectionFacts({ traceId = null, ownerUserId = null, own
|
||||
const projection = createWorkbenchTurnProjection({ traceId: safeId, result: payload, session: { id: resolvedSessionId, status: normalizedStatus, session }, trace: payload?.runnerTrace ?? null });
|
||||
const projectedStatus = normalizeWorkbenchStatus(projection.status ?? normalizedStatus);
|
||||
const terminal = projection.terminal === true;
|
||||
const terminalSealBlocked = projection.terminalSealBlocked === true;
|
||||
const terminalStatus = terminal ? (TERMINAL_STATUSES.has(projectedStatus) ? projectedStatus : normalizedStatus) : projectedStatus;
|
||||
const timing = terminalTimingAtLeastProjectedAt(projectionTimingForStatus(projection.timing, terminal), terminal);
|
||||
const timingAuthorityIssue = terminalTimingAuthorityIssue(timing, { terminal, traceId: safeId, source: "facts", status: terminalStatus, label: payload?.lastEventLabel, sourceSeq: projection.lastProjectedSeq });
|
||||
@@ -547,6 +548,7 @@ function buildWorkbenchProjectionFacts({ traceId = null, ownerUserId = null, own
|
||||
terminal,
|
||||
terminalStatus,
|
||||
finalText,
|
||||
terminalSealBlocked,
|
||||
timestamp,
|
||||
timing
|
||||
});
|
||||
@@ -870,15 +872,15 @@ function hasTerminalResultTraceEvent(events = []) {
|
||||
});
|
||||
}
|
||||
|
||||
function normalizeMessageFact(message = {}, index, { traceId, sessionId, conversationId, threadId, terminal, terminalStatus = "completed", finalText = null, timestamp, timing: contextTiming }) {
|
||||
function normalizeMessageFact(message = {}, index, { traceId, sessionId, conversationId, threadId, terminal, terminalStatus = "completed", finalText = null, terminalSealBlocked = false, timestamp, timing: contextTiming }) {
|
||||
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";
|
||||
const status = appliesToContextTrace && terminal ? contextTerminalStatus : normalizeWorkbenchStatus(message.status ?? (terminal && role !== "user" ? contextTerminalStatus : role === "user" ? "sent" : "running"));
|
||||
const messageTerminal = role !== "user" && TERMINAL_STATUSES.has(status);
|
||||
const status = appliesToContextTrace && terminalSealBlocked ? "running" : appliesToContextTrace && terminal ? contextTerminalStatus : normalizeWorkbenchStatus(message.status ?? (terminal && role !== "user" ? contextTerminalStatus : role === "user" ? "sent" : "running"));
|
||||
const messageTerminal = role !== "user" && !(appliesToContextTrace && terminalSealBlocked) && TERMINAL_STATUSES.has(status);
|
||||
const messageTiming = normalizeTimingProjection(message.timing) ?? normalizeTimingProjection(message);
|
||||
const contextTimingProjection = appliesToContextTrace ? normalizeTimingProjection(contextTiming) : null;
|
||||
const timing = role === "user" ? messageTiming ?? emptyTimingProjection() : contextTimingProjection ?? messageTiming ?? emptyTimingProjection();
|
||||
|
||||
@@ -13,10 +13,14 @@ export function createWorkbenchTurnProjection({ turnId = null, traceId = null, r
|
||||
const traceTerminal = traceTerminalEvidence(trace);
|
||||
const terminalEvidence = terminalTurnEvidence({ result, traceTerminal });
|
||||
const activeEvidence = activeTurnEvidence({ result, session, trace });
|
||||
const status = terminalEvidence?.status ?? activeEvidence?.status ?? "unknown";
|
||||
const evidenceStatus = terminalEvidence?.status ?? activeEvidence?.status ?? "unknown";
|
||||
const terminalObserved = Boolean(terminalEvidence && TERMINAL_STATUSES.has(evidenceStatus) && !RUNNING_STATUSES.has(evidenceStatus));
|
||||
const terminalFinalText = terminalObserved ? projectionText(terminalEvidence?.finalResponse) : null;
|
||||
const terminalSealBlocked = terminalObserved && evidenceStatus === "completed" && !terminalFinalText;
|
||||
const status = terminalSealBlocked ? normalizeActiveStatus(activeEvidence?.status ?? "running") : evidenceStatus;
|
||||
const running = RUNNING_STATUSES.has(status);
|
||||
const terminal = Boolean(terminalEvidence && TERMINAL_STATUSES.has(status) && !running);
|
||||
const finalText = terminal ? projectionText(terminalEvidence?.finalResponse) : null;
|
||||
const terminal = Boolean(terminalObserved && !terminalSealBlocked && TERMINAL_STATUSES.has(status) && !running);
|
||||
const finalText = terminal ? terminalFinalText : null;
|
||||
const agentRun = objectValue(result?.agentRun ?? trace?.agentRun);
|
||||
const lastEvent = traceLastEvent(trace);
|
||||
const timing = createWorkbenchTurnTimingProjection({ result, session, trace, status, terminal });
|
||||
@@ -26,6 +30,10 @@ export function createWorkbenchTurnProjection({ turnId = null, traceId = null, r
|
||||
status,
|
||||
running,
|
||||
terminal,
|
||||
terminalObserved,
|
||||
terminalObservedStatus: terminalObserved ? evidenceStatus : null,
|
||||
terminalSealBlocked,
|
||||
waitingFor: terminalSealBlocked ? "final_response" : null,
|
||||
source: terminalEvidence?.source ?? activeEvidence?.source ?? null,
|
||||
terminalEvidence,
|
||||
finalResponse: finalText ? { text: finalText, status, traceId: projectionTraceId, valuesPrinted: false } : null,
|
||||
@@ -112,6 +120,7 @@ export function createWorkbenchTurnTimingProjection({ result = null, session = n
|
||||
export function projectionDiagnostics({ traceId = null, projection = null, result = null, trace = null, refreshError = null } = {}) {
|
||||
const turn = projection ?? createWorkbenchTurnProjection({ traceId, result, trace });
|
||||
const sealedCompleted = turn.terminal === true && turn.status === "completed" && Boolean(turn.finalResponse?.text);
|
||||
const waitingFor = textValue(turn.waitingFor) || null;
|
||||
const source = sealedCompleted ? null : projectionDiagnosticSource(trace);
|
||||
const rawBlocker = sealedCompleted ? null : refreshError ?? source?.blocker ?? trace?.blocker ?? result?.blocker ?? result?.error ?? null;
|
||||
const retryingProviderInterruption = !sealedCompleted && turn.running === true ? retryableProviderInterruptionEvidence(rawBlocker) : null;
|
||||
@@ -119,11 +128,11 @@ export function projectionDiagnostics({ traceId = null, projection = null, resul
|
||||
const hasProjectionInput = hasTraceProjection(trace) || Boolean(result || result?.agentRun);
|
||||
const sourceStatus = sealedCompleted ? null : normalizeProjectionStatus(source?.projectionStatus ?? trace?.projectionStatus);
|
||||
const effectiveSourceStatus = retryingProviderInterruption && sourceStatus === "blocked" ? "projecting" : sourceStatus;
|
||||
const status = sealedCompleted ? "caught-up" : effectiveSourceStatus ?? (blocker ? "blocked" : turn.terminal ? "caught-up" : hasProjectionInput ? "projecting" : "unknown");
|
||||
const status = sealedCompleted ? "caught-up" : waitingFor ? "projecting" : effectiveSourceStatus ?? (blocker ? "blocked" : turn.terminal ? "caught-up" : hasProjectionInput ? "projecting" : "unknown");
|
||||
const diagnostic = blocker ? diagnosticBlocker(blocker) : null;
|
||||
const sourceHealth = sealedCompleted ? null : normalizeProjectionHealth(source?.projectionHealth ?? trace?.projectionHealth);
|
||||
const effectiveSourceHealth = retryingProviderInterruption && (sourceHealth === "degraded" || sourceHealth === "unavailable" || sourceHealth === "stalled") ? "projecting" : sourceHealth;
|
||||
const projectionHealth = sealedCompleted ? "healthy" : effectiveSourceHealth
|
||||
const projectionHealth = sealedCompleted ? "healthy" : waitingFor ? "projecting" : effectiveSourceHealth
|
||||
?? projectionHealthFor({ status, turn, hasProjectionInput, blocker: diagnostic });
|
||||
const staleMs = sealedCompleted ? null : projectionStaleMs(source?.staleMs ?? trace?.staleMs, turn.updatedAt ?? trace?.updatedAt ?? result?.updatedAt);
|
||||
return {
|
||||
@@ -135,6 +144,9 @@ export function projectionDiagnostics({ traceId = null, projection = null, resul
|
||||
staleMs,
|
||||
blocker: diagnostic,
|
||||
retryingProviderInterruption,
|
||||
waitingFor,
|
||||
terminalObserved: turn.terminalObserved === true,
|
||||
terminalObservedStatus: turn.terminalObservedStatus ?? null,
|
||||
updatedAt: turn.updatedAt ?? null,
|
||||
valuesRedacted: true
|
||||
};
|
||||
@@ -189,6 +201,16 @@ function terminalFinalResponse(status, result = null, traceTerminal = null) {
|
||||
if (normalizeWorkbenchStatus(status) === "canceled") return canonicalCancelFinalResponse(result, traceTerminal);
|
||||
const direct = traceTerminal?.finalResponse ?? result?.finalResponse;
|
||||
if (direct) return direct;
|
||||
const resultText = projectionText(result?.assistantText, result?.finalText, result?.reply);
|
||||
if (resultText) {
|
||||
return {
|
||||
text: resultText,
|
||||
status: normalizeWorkbenchStatus(status),
|
||||
traceId: textValue(result?.traceId ?? traceTerminal?.finalResponse?.traceId ?? traceTerminal?.evidence?.traceId) || null,
|
||||
source: "result-final-response-evidence",
|
||||
valuesPrinted: false
|
||||
};
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
|
||||
@@ -1006,6 +1006,38 @@ test("configured postgres runtime persists and queries Workbench durable facts",
|
||||
assert.ok(queryClient.calls.some((call) => call.sql.startsWith("SELECT session_json FROM workbench_sessions WHERE session_id = $1")));
|
||||
});
|
||||
|
||||
test("configured postgres runtime rejects completed Workbench turn facts without final response", async () => {
|
||||
const queryClient = createFakePostgresClient({ migrationReady: true });
|
||||
const store = createConfiguredCloudRuntimeStore({
|
||||
env: {
|
||||
HWLAB_CLOUD_RUNTIME_ADAPTER: "postgres",
|
||||
HWLAB_CLOUD_RUNTIME_DURABLE: "true"
|
||||
},
|
||||
dbUrl: "postgres://hwlab_redacted@db.example.invalid:5432/hwlab",
|
||||
queryClient,
|
||||
now: () => "2026-06-20T10:10:00.000Z"
|
||||
});
|
||||
|
||||
await assert.rejects(
|
||||
() => store.writeWorkbenchFacts({
|
||||
facts: {
|
||||
turns: [{
|
||||
turnId: "turn_fact_missing_final",
|
||||
sessionId: "ses_fact_missing_final",
|
||||
traceId: "trc_fact_missing_final",
|
||||
messageId: "msg_fact_missing_final",
|
||||
status: "completed",
|
||||
terminal: true,
|
||||
sealed: true,
|
||||
finalResponse: null,
|
||||
updatedAt: "2026-06-20T10:10:01.000Z"
|
||||
}]
|
||||
}
|
||||
}),
|
||||
(error) => error?.code === "workbench_terminal_final_response_missing"
|
||||
);
|
||||
});
|
||||
|
||||
test("configured postgres runtime appends Workbench aggregate events and outbox in one transaction", async () => {
|
||||
const queryClient = createFakePostgresClient({ migrationReady: true });
|
||||
const store = createConfiguredCloudRuntimeStore({
|
||||
@@ -1020,7 +1052,7 @@ test("configured postgres runtime appends Workbench aggregate events and outbox
|
||||
|
||||
const write = await store.writeWorkbenchFacts({ facts: {
|
||||
traceEvents: [{ id: "wte_event_stream_1", traceId: "trc_event_stream", sessionId: "ses_event_stream", turnId: "turn_event_stream", sourceSeq: 1, sourceEventId: "source-event-1", projectedSeq: 1, eventType: "assistant", occurredAt: "2026-06-24T12:00:01.000Z" }],
|
||||
turns: [{ turnId: "turn_event_stream", sessionId: "ses_event_stream", traceId: "trc_event_stream", status: "completed", sourceSeq: 2, sourceEventId: "source-terminal-2", projectedSeq: 2, terminal: true, sealed: true, updatedAt: "2026-06-24T12:00:02.000Z" }]
|
||||
turns: [{ turnId: "turn_event_stream", sessionId: "ses_event_stream", traceId: "trc_event_stream", status: "completed", finalResponse: { text: "event stream done" }, sourceSeq: 2, sourceEventId: "source-terminal-2", projectedSeq: 2, terminal: true, sealed: true, updatedAt: "2026-06-24T12:00:02.000Z" }]
|
||||
} });
|
||||
const outbox = await store.readWorkbenchProjectionOutbox({ afterSeq: 0, traceId: "trc_event_stream", limit: 10 });
|
||||
|
||||
|
||||
@@ -2673,6 +2673,11 @@ function normalizeWorkbenchTurnFact(value, requestMeta = {}, now) {
|
||||
const sourceSeq = nonNegativeInteger(input.sourceSeq);
|
||||
const projectedSeq = nonNegativeInteger(input.projectedSeq ?? sourceSeq);
|
||||
const timestamp = timestampOr(input.updatedAt ?? input.createdAt, now);
|
||||
const status = textOr(input.status, "unknown");
|
||||
const terminal = Boolean(input.terminal);
|
||||
const sealed = Boolean(input.sealed);
|
||||
const finalResponse = input.finalResponse ?? null;
|
||||
assertWorkbenchTurnFinalResponseInvariant({ status, terminal, sealed, finalResponse });
|
||||
return pruneUndefined({
|
||||
...input,
|
||||
id: turnId,
|
||||
@@ -2680,13 +2685,13 @@ function normalizeWorkbenchTurnFact(value, requestMeta = {}, now) {
|
||||
sessionId,
|
||||
traceId,
|
||||
messageId: textOr(input.messageId ?? requestMeta.messageId, "") || null,
|
||||
status: textOr(input.status, "unknown"),
|
||||
status,
|
||||
projectedSeq,
|
||||
sourceSeq,
|
||||
sourceEventId: textOr(input.sourceEventId ?? input.eventId, "") || null,
|
||||
terminal: Boolean(input.terminal),
|
||||
sealed: Boolean(input.sealed),
|
||||
finalResponse: input.finalResponse ?? null,
|
||||
terminal,
|
||||
sealed,
|
||||
finalResponse,
|
||||
diagnostic: normalizeJsonObject(input.diagnostic ?? input.projection ?? input.blocker),
|
||||
createdAt: timestampOr(input.createdAt, timestamp),
|
||||
updatedAt: timestamp,
|
||||
@@ -2694,6 +2699,29 @@ function normalizeWorkbenchTurnFact(value, requestMeta = {}, now) {
|
||||
});
|
||||
}
|
||||
|
||||
function assertWorkbenchTurnFinalResponseInvariant({ status, terminal, sealed, finalResponse } = {}) {
|
||||
if (normalizeWorkbenchFactStatus(status) !== "completed") return;
|
||||
if (!terminal && !sealed) return;
|
||||
if (workbenchFactFinalResponseText(finalResponse)) return;
|
||||
const error = new Error("completed Workbench turn facts must include an authoritative final response before sealing");
|
||||
error.code = "workbench_terminal_final_response_missing";
|
||||
error.layer = "workbench-read-model";
|
||||
error.category = "projection-invariant";
|
||||
throw error;
|
||||
}
|
||||
|
||||
function normalizeWorkbenchFactStatus(value) {
|
||||
return textOr(value, "").toLowerCase().replace(/_/gu, "-");
|
||||
}
|
||||
|
||||
function workbenchFactFinalResponseText(value) {
|
||||
if (value && typeof value === "object" && !Array.isArray(value)) {
|
||||
return workbenchFactFinalResponseText(value.text ?? value.content ?? value.message ?? value.summary ?? value.preview);
|
||||
}
|
||||
const text = textOr(value, "");
|
||||
return text && text !== "[object Object]" ? text : null;
|
||||
}
|
||||
|
||||
function normalizeWorkbenchTraceEventFact(value, requestMeta = {}, now) {
|
||||
const input = normalizeJsonObject(value);
|
||||
const traceId = textOr(input.traceId ?? requestMeta.traceId, "");
|
||||
|
||||
@@ -0,0 +1,49 @@
|
||||
import assert from "node:assert/strict";
|
||||
import { test } from "bun:test";
|
||||
|
||||
import {
|
||||
terminalMessagePatchFromTurnResult,
|
||||
turnResultIsTerminalForMerge,
|
||||
turnResultStatusForMerge
|
||||
} from "./workbench-message-projection-runtime";
|
||||
|
||||
test("turn result merge keeps completed result without final response unsealed", () => {
|
||||
const traceId = "trc_frontend_terminal_without_final";
|
||||
const message = { id: "msg_frontend_terminal_without_final", role: "agent", status: "running", traceId } as any;
|
||||
const result = {
|
||||
traceId,
|
||||
status: "completed",
|
||||
terminal: true,
|
||||
terminalSealBlocked: true,
|
||||
waitingFor: "final_response",
|
||||
runnerTrace: {
|
||||
traceId,
|
||||
events: [
|
||||
{ seq: 1, type: "assistant_message", status: "running", message: "progress only" },
|
||||
{ seq: 2, type: "terminal_status", terminal: true, terminalStatus: "completed" }
|
||||
]
|
||||
}
|
||||
} as any;
|
||||
|
||||
assert.equal(turnResultStatusForMerge(result), "running");
|
||||
assert.equal(turnResultIsTerminalForMerge(result), false);
|
||||
assert.equal(terminalMessagePatchFromTurnResult(message, result), null);
|
||||
});
|
||||
|
||||
test("turn result merge seals completed result with final response", () => {
|
||||
const traceId = "trc_frontend_terminal_with_final";
|
||||
const message = { id: "msg_frontend_terminal_with_final", role: "agent", status: "running", traceId } as any;
|
||||
const result = {
|
||||
traceId,
|
||||
status: "completed",
|
||||
terminal: true,
|
||||
finalResponse: { text: "final answer" }
|
||||
} as any;
|
||||
|
||||
const patch = terminalMessagePatchFromTurnResult(message, result);
|
||||
assert.equal(turnResultStatusForMerge(result), "completed");
|
||||
assert.equal(turnResultIsTerminalForMerge(result), true);
|
||||
assert.equal(patch?.status, "completed");
|
||||
assert.equal((patch?.finalResponse as any)?.text, "final answer");
|
||||
assert.equal(patch?.text, "final answer");
|
||||
});
|
||||
@@ -272,9 +272,32 @@ export function messageStatusPatchForTerminalMerge(message: ChatMessage, resultS
|
||||
return {};
|
||||
}
|
||||
|
||||
export function turnResultTerminalSealBlocked(result: AgentChatResultResponse | TraceSnapshot | null | undefined): boolean {
|
||||
const record = recordValue(result);
|
||||
if (!record) return false;
|
||||
const status = normalizedStatusText(record.status) ?? null;
|
||||
const explicitBlocked = record.terminalSealBlocked === true || firstNonEmptyString(record.waitingFor) === "final_response";
|
||||
if (!explicitBlocked && status !== "completed") return false;
|
||||
if (terminalFinalResponseTextFromTurnResult(record as AgentChatResultResponse)) return false;
|
||||
return explicitBlocked || status === "completed";
|
||||
}
|
||||
|
||||
export function turnResultStatusForMerge(result: AgentChatResultResponse | TraceSnapshot | null | undefined): string | null {
|
||||
const record = recordValue(result);
|
||||
if (!record) return null;
|
||||
if (turnResultTerminalSealBlocked(result)) return "running";
|
||||
return normalizedStatusText(record.status) ?? null;
|
||||
}
|
||||
|
||||
export function turnResultIsTerminalForMerge(result: AgentChatResultResponse | TraceSnapshot | null | undefined): boolean {
|
||||
const record = recordValue(result);
|
||||
if (!record || turnResultTerminalSealBlocked(result)) return false;
|
||||
return record.terminal === true || isTerminalMessageStatus(normalizedStatusText(record.status));
|
||||
}
|
||||
|
||||
export function terminalMessagePatchFromTurnResult(message: ChatMessage, result: AgentChatResultResponse): Partial<ChatMessage> | null {
|
||||
const resultStatus = normalizedStatusText(result.status) ?? null;
|
||||
const terminal = result.terminal === true || isTerminalMessageStatus(resultStatus);
|
||||
const resultStatus = turnResultStatusForMerge(result);
|
||||
const terminal = turnResultIsTerminalForMerge(result);
|
||||
if (!terminal) return null;
|
||||
const resultError = normalizeAgentError(result.error ?? null);
|
||||
const resultProjection = projectionFromResult(result);
|
||||
@@ -461,7 +484,7 @@ export function traceHasCompletedFinalResponse(traceId: string | null | undefine
|
||||
export function messageHasCompletedFinalResponse(message: ChatMessage | null | undefined): boolean {
|
||||
if (!message || message.role !== "agent") return false;
|
||||
if (normalizedStatusText(message.status) !== "completed") return false;
|
||||
return Boolean(firstNonEmptyString(message.text, finalResponseText((message as Record<string, unknown>).finalResponse)));
|
||||
return Boolean(firstNonEmptyString(message.text, messageText((message as Record<string, unknown>).content), finalResponseText((message as Record<string, unknown>).finalResponse)));
|
||||
}
|
||||
|
||||
export function traceHasEvents(trace: ChatMessage["runnerTrace"]): boolean {
|
||||
|
||||
@@ -0,0 +1,32 @@
|
||||
import assert from "node:assert/strict";
|
||||
import { test } from "bun:test";
|
||||
|
||||
import { createWorkbenchServerState, reduceWorkbenchServerState, selectActiveMessages, selectSessionStatusAuthority } from "./workbench-server-state";
|
||||
|
||||
test("server state reducer keeps completed message without final response unsealed", () => {
|
||||
const sessionId = "ses_state_completed_without_final";
|
||||
const traceId = "trc_state_completed_without_final";
|
||||
let state = createWorkbenchServerState();
|
||||
state = reduceWorkbenchServerState(state, { type: "session.detail", session: { sessionId, status: "running", messages: [] } as any });
|
||||
state = reduceWorkbenchServerState(state, { type: "session.messages", sessionId, messages: [{ id: "msg_state_completed_without_final", role: "agent", status: "running", traceId, sessionId, startedAt: "2026-07-01T00:00:00.000Z" } as any] });
|
||||
state = reduceWorkbenchServerState(state, { type: "session.messages", sessionId, messages: [{ id: "msg_state_completed_without_final", role: "agent", status: "completed", traceId, sessionId, startedAt: "2026-07-01T00:00:00.000Z", finishedAt: "2026-07-01T00:00:02.000Z", durationMs: 2000 } as any] });
|
||||
|
||||
const [message] = selectActiveMessages(state, sessionId);
|
||||
assert.equal(message.status, "running");
|
||||
assert.equal(message.finishedAt, null);
|
||||
assert.equal(message.durationMs, null);
|
||||
assert.notEqual(selectSessionStatusAuthority(state)[sessionId]?.status, "completed");
|
||||
});
|
||||
|
||||
test("server state reducer seals completed message when final response exists", () => {
|
||||
const sessionId = "ses_state_completed_with_final";
|
||||
const traceId = "trc_state_completed_with_final";
|
||||
let state = createWorkbenchServerState();
|
||||
state = reduceWorkbenchServerState(state, { type: "session.detail", session: { sessionId, status: "running", messages: [] } as any });
|
||||
state = reduceWorkbenchServerState(state, { type: "session.messages", sessionId, messages: [{ id: "msg_state_completed_with_final", role: "agent", status: "completed", traceId, sessionId, text: "final answer", finishedAt: "2026-07-01T00:00:02.000Z", durationMs: 2000 } as any] });
|
||||
|
||||
const [message] = selectActiveMessages(state, sessionId);
|
||||
assert.equal(message.status, "completed");
|
||||
assert.equal(message.text, "final answer");
|
||||
assert.equal(selectSessionStatusAuthority(state)[sessionId]?.status, "completed");
|
||||
});
|
||||
@@ -188,26 +188,45 @@ function messageMatchesSnapshot(existing: ChatMessage, incoming: ChatMessage): b
|
||||
}
|
||||
|
||||
function mergeMessageSnapshot(existing: ChatMessage, incoming: ChatMessage): ChatMessage {
|
||||
const normalizedIncoming = normalizeUnsealedCompletedMessage(incoming);
|
||||
const merged = {
|
||||
...existing,
|
||||
...incoming,
|
||||
runnerTrace: incoming.runnerTrace ?? existing.runnerTrace ?? null,
|
||||
...normalizedIncoming,
|
||||
runnerTrace: normalizedIncoming.runnerTrace ?? existing.runnerTrace ?? null,
|
||||
traceAutoLifecycle: existing.traceAutoLifecycle,
|
||||
updatedAt: incoming.updatedAt ?? existing.updatedAt
|
||||
updatedAt: normalizedIncoming.updatedAt ?? existing.updatedAt
|
||||
};
|
||||
return sealExistingMessageTiming(existing, merged);
|
||||
}
|
||||
|
||||
function mergeMessageList(existing: ChatMessage[], incoming: ChatMessage[]): ChatMessage[] {
|
||||
return incoming.map((message) => {
|
||||
return incoming.map((rawMessage) => {
|
||||
const message = normalizeUnsealedCompletedMessage(rawMessage);
|
||||
const previous = existing.find((item) => messageMatchesSnapshot(item, message));
|
||||
return previous ? sealExistingMessageTiming(previous, message) : message;
|
||||
});
|
||||
}
|
||||
|
||||
function normalizeUnsealedCompletedMessage(message: ChatMessage): ChatMessage {
|
||||
if (message.role !== "agent") return message;
|
||||
if (normalizedMessageStatus(message.status) !== "completed") return message;
|
||||
if (messageHasCompletedFinalResponse(message)) return message;
|
||||
const timing = message.timing ? { ...message.timing, finishedAt: null, durationMs: null } : message.timing;
|
||||
return {
|
||||
...message,
|
||||
status: "running",
|
||||
traceAutoLifecycle: message.traceAutoLifecycle ?? "running",
|
||||
timing,
|
||||
finishedAt: null,
|
||||
durationMs: null
|
||||
};
|
||||
}
|
||||
|
||||
function sealExistingMessageTiming(existing: ChatMessage, incoming: ChatMessage): ChatMessage {
|
||||
if (!isTerminalMessageStatus(existing.status) && isTerminalMessageStatus(incoming.status)) return sealTerminalTransitionMessageTiming(existing, incoming);
|
||||
if (isTerminalMessageStatus(existing.status)) return sealExistingTerminalMessageTiming(existing, incoming);
|
||||
const existingTerminal = messageIsTerminalForState(existing);
|
||||
const incomingTerminal = messageIsTerminalForState(incoming);
|
||||
if (!existingTerminal && incomingTerminal) return sealTerminalTransitionMessageTiming(existing, incoming);
|
||||
if (existingTerminal) return sealExistingTerminalMessageTiming(existing, incoming);
|
||||
return sealExistingRunningMessageTiming(existing, incoming);
|
||||
}
|
||||
|
||||
@@ -228,7 +247,7 @@ function sealExistingRunningMessageTiming(existing: ChatMessage, incoming: ChatM
|
||||
}
|
||||
|
||||
function sealExistingTerminalMessageTiming(existing: ChatMessage, incoming: ChatMessage): ChatMessage {
|
||||
if (!isTerminalMessageStatus(existing.status)) return incoming;
|
||||
if (!messageIsTerminalForState(existing)) return incoming;
|
||||
const timing = sealedTerminalTimingProjection(existing, incoming);
|
||||
const patch: Partial<ChatMessage> = sealedTerminalMessagePatch(existing);
|
||||
if (timing && timing.durationMs != null) Object.assign(patch, {
|
||||
@@ -411,6 +430,13 @@ function isTerminalTurnStatusAuthority(turn: TurnStatusAuthority): boolean {
|
||||
return turn.terminal === true || (turn.running !== true && isTerminalMessageStatus(turn.status));
|
||||
}
|
||||
|
||||
function messageIsTerminalForState(message: ChatMessage): boolean {
|
||||
const status = normalizedMessageStatus(message.status);
|
||||
if (!isTerminalMessageStatus(status)) return false;
|
||||
if (status !== "completed") return true;
|
||||
return messageHasCompletedFinalResponse(message);
|
||||
}
|
||||
|
||||
function isSameTraceAuthority(left: SessionStatusAuthority, right: SessionStatusAuthority): boolean {
|
||||
const leftTrace = textValue(left.lastTraceId);
|
||||
const rightTrace = textValue(right.lastTraceId);
|
||||
@@ -428,7 +454,7 @@ function messageHasCompletedFinalResponse(message: ChatMessage): boolean {
|
||||
}
|
||||
|
||||
function messageFinalResponseText(message: ChatMessage): string | null {
|
||||
return textValue(message.text) ?? nestedTextValue((message as Record<string, unknown>).finalResponse);
|
||||
return textValue(message.text) ?? textValue(message.content) ?? nestedTextValue((message as Record<string, unknown>).finalResponse);
|
||||
}
|
||||
|
||||
function projectionFromMessageRecord(message: ChatMessage): ProjectionDiagnostic | null {
|
||||
|
||||
@@ -0,0 +1,29 @@
|
||||
import assert from "node:assert/strict";
|
||||
import { test } from "bun:test";
|
||||
|
||||
import { resolveCancelableAgentMessage, resolveComposerState } from "./workbench-session";
|
||||
|
||||
test("composer treats completed message without final response as unsealed running turn", () => {
|
||||
const traceId = "trc_session_completed_without_final";
|
||||
const sessionId = "ses_session_completed_without_final";
|
||||
const message = { id: "msg_session_completed_without_final", role: "agent", status: "completed", traceId, sessionId } as any;
|
||||
const composer = resolveComposerState({
|
||||
messages: [message],
|
||||
activeSessionId: sessionId,
|
||||
chatPending: true,
|
||||
currentRequest: { traceId, sessionId, status: "running" },
|
||||
turnStatusAuthority: { [traceId]: { traceId, sessionId, status: "completed", terminal: true, running: false } as any }
|
||||
});
|
||||
assert.equal(composer.submitMode, "steer");
|
||||
assert.equal(composer.targetTraceId, traceId);
|
||||
});
|
||||
|
||||
test("cancel target remains available for completed message until final response exists", () => {
|
||||
const traceId = "trc_cancel_completed_without_final";
|
||||
const sessionId = "ses_cancel_completed_without_final";
|
||||
const unsealed = { id: "msg_cancel_completed_without_final", role: "agent", status: "completed", traceId, sessionId } as any;
|
||||
const sealed = { ...unsealed, id: "msg_cancel_completed_with_final", text: "final answer" } as any;
|
||||
|
||||
assert.equal(resolveCancelableAgentMessage({ messages: [unsealed], targetTraceId: traceId, targetSessionId: sessionId, turnStatusAuthority: { [traceId]: { traceId, sessionId, status: "running", running: true, terminal: false } as any } })?.id, unsealed.id);
|
||||
assert.equal(resolveCancelableAgentMessage({ messages: [sealed], targetTraceId: traceId, targetSessionId: sessionId, turnStatusAuthority: { [traceId]: { traceId, sessionId, status: "running", running: true, terminal: false } as any } }), null);
|
||||
});
|
||||
@@ -82,9 +82,12 @@ export function resolveComposerState(input: { messages: ChatMessage[]; sessions?
|
||||
const turn = activeTraceId ? input.turnStatusAuthority?.[activeTraceId] : null;
|
||||
const activeByRequest = Boolean(input.chatPending && currentRequest && isActiveStatus(currentRequest.status ?? "running"));
|
||||
const activeByMessage = latestMessage?.role === "agent" && isActiveStatus(latestMessage.status);
|
||||
const terminalByMessage = latestMessage?.role === "agent" && isTerminalStatus(latestMessage.status);
|
||||
const activeByStatus = !terminalByMessage && (turn?.running === true || isActiveStatus(turn?.status) || activeByRequest || activeByMessage);
|
||||
const terminal = terminalByMessage || turn?.terminal === true || isTerminalStatus(turn?.status);
|
||||
const terminalByMessage = messageIsTerminalForSession(latestMessage);
|
||||
const turnStatus = normalizeSessionStatus(turn?.status);
|
||||
const turnCompletedMissingFinal = (turn?.terminal === true || turnStatus === "completed") && !messageHasCompletedFinalResponse(latestMessage);
|
||||
const terminalByTurn = !turnCompletedMissingFinal && (turn?.terminal === true || isTerminalStatus(turnStatus));
|
||||
const activeByStatus = !terminalByMessage && !terminalByTurn && (turn?.running === true || isActiveStatus(turnStatus) || activeByRequest || activeByMessage);
|
||||
const terminal = terminalByMessage || terminalByTurn;
|
||||
const active = activeSession(input.sessions ?? [], sessionId);
|
||||
const effectiveSessionId = firstNonEmptyString(turn?.sessionId, currentRequest?.sessionId, active?.sessionId, sessionId);
|
||||
const threadId = firstNonEmptyString(turn?.threadId, currentRequest?.threadId, active?.threadId);
|
||||
@@ -103,7 +106,7 @@ export function resolveCancelableAgentMessage(input: { messages: ChatMessage[];
|
||||
if (message.role !== "agent") continue;
|
||||
if (firstNonEmptyString(message.traceId, message.runnerTrace?.traceId) !== targetTraceId) continue;
|
||||
if (!messageBelongsToCancelTarget(message, input)) continue;
|
||||
if (isTerminalStatus(message.status)) continue;
|
||||
if (messageIsTerminalForSession(message)) continue;
|
||||
return message;
|
||||
}
|
||||
return null;
|
||||
@@ -279,6 +282,21 @@ function latestAgentMessage(messages: ChatMessage[] | undefined): ChatMessage |
|
||||
return [...(messages ?? [])].reverse().find((message) => message.role === "agent") ?? null;
|
||||
}
|
||||
|
||||
function messageIsTerminalForSession(message: ChatMessage | null | undefined): boolean {
|
||||
if (!message || message.role !== "agent") return false;
|
||||
const status = normalizeSessionStatus(message.status);
|
||||
if (!isTerminalStatus(status)) return false;
|
||||
if (status !== "completed") return true;
|
||||
return messageHasCompletedFinalResponse(message);
|
||||
}
|
||||
|
||||
function messageHasCompletedFinalResponse(message: ChatMessage | null | undefined): boolean {
|
||||
if (!message || message.role !== "agent") return false;
|
||||
if (normalizeSessionStatus(message.status) !== "completed") return false;
|
||||
const finalResponse = recordValue((message as Record<string, unknown>).finalResponse);
|
||||
return Boolean(firstNonEmptyString(message.text, message.content, finalResponse?.text, finalResponse?.content, finalResponse?.message));
|
||||
}
|
||||
|
||||
function resolveSessionTabStatus(session: WorkbenchSessionRecord, authority: SessionStatusAuthority | null | undefined): string {
|
||||
void session;
|
||||
const authorityStatus = normalizeSessionStatus(authority?.status);
|
||||
|
||||
@@ -55,6 +55,8 @@ import {
|
||||
shouldClearCompletedTurnDiagnostics,
|
||||
shouldSuppressTransientWorkbenchReadFailure,
|
||||
terminalMessageTimingPatchForNormalize,
|
||||
turnResultIsTerminalForMerge,
|
||||
turnResultStatusForMerge,
|
||||
traceHasCompletedFinalResponse,
|
||||
traceHasEvents,
|
||||
traceResultHasTerminalEvidence,
|
||||
@@ -518,9 +520,9 @@ export const useWorkbenchStore = defineStore("workbench", () => {
|
||||
if (!id) return false;
|
||||
if (currentRequest.value?.traceId === id) return true;
|
||||
const turn = turnStatusAuthority.value[id];
|
||||
if (turn?.terminal === true || isTerminalMessageStatus(turn?.status)) return false;
|
||||
const message = [...messages.value].reverse().find((item) => firstNonEmptyString(item.traceId, item.runnerTrace?.traceId) === id) ?? null;
|
||||
if (isTerminalMessageStatus(message?.status)) return false;
|
||||
if (traceAuthorityIsSealed(turn, message)) return false;
|
||||
if (messageIsSealedTerminal(message)) return false;
|
||||
return isTraceActiveStatus(turn?.status) || isTraceActiveStatus(message?.status);
|
||||
}
|
||||
|
||||
@@ -623,19 +625,33 @@ export const useWorkbenchStore = defineStore("workbench", () => {
|
||||
const id = firstNonEmptyString(traceId);
|
||||
if (!id) return false;
|
||||
const turn = turnStatusAuthority.value[id];
|
||||
if (turn?.terminal === true || isTerminalMessageStatus(turn?.status)) return false;
|
||||
const message = [...source].reverse().find((item) => firstNonEmptyString(item.traceId, item.runnerTrace?.traceId) === id) ?? null;
|
||||
if (isTerminalMessageStatus(message?.status)) return false;
|
||||
if (traceAuthorityIsSealed(turn, message)) return false;
|
||||
if (messageIsSealedTerminal(message)) return false;
|
||||
return true;
|
||||
}
|
||||
|
||||
function traceAuthorityIsSealed(turn: TurnStatusAuthority | undefined, message: ChatMessage | null): boolean {
|
||||
const status = normalizedStatusText(turn?.status);
|
||||
if (!(turn?.terminal === true || isTerminalMessageStatus(status))) return false;
|
||||
if (status !== "completed") return true;
|
||||
return message ? messageHasCompletedFinalResponse(message) : false;
|
||||
}
|
||||
|
||||
function messageIsSealedTerminal(message: ChatMessage | null | undefined): boolean {
|
||||
const status = normalizedStatusText(message?.status);
|
||||
if (!isTerminalMessageStatus(status)) return false;
|
||||
if (status !== "completed") return true;
|
||||
return message ? messageHasCompletedFinalResponse(message) : false;
|
||||
}
|
||||
|
||||
async function refreshTurnStatusByTraceId(traceId: string | null | undefined, options: { force?: boolean } = {}): Promise<void> {
|
||||
const id = firstNonEmptyString(traceId);
|
||||
if (!id) return;
|
||||
const response = await fetchWorkbenchTurnStatus(id, shouldUseActivityTimeoutForTrace(id), { force: options.force });
|
||||
if (response.ok && response.data) {
|
||||
applyTurnStatusSnapshot(id, response.data);
|
||||
if (response.data.terminal === true || isTerminalMessageStatus(response.data.status)) completeTrace(id, response.data, { forceRead: options.force });
|
||||
if (turnResultIsTerminalForMerge(response.data as AgentChatResultResponse)) completeTrace(id, response.data, { forceRead: options.force });
|
||||
return;
|
||||
}
|
||||
if (shouldSuppressTransientWorkbenchReadFailure(response)) return;
|
||||
@@ -654,9 +670,9 @@ export const useWorkbenchStore = defineStore("workbench", () => {
|
||||
function rememberTurnStatus(traceId: string, result: AgentChatResultResponse | TraceSnapshot): void {
|
||||
const id = firstNonEmptyString(result.traceId, traceId);
|
||||
if (!id) return;
|
||||
const status = normalizedStatusText(result.status) ?? null;
|
||||
const status = turnResultStatusForMerge(result as AgentChatResultResponse);
|
||||
const running = (result as AgentChatResultResponse).running === true || isTraceActiveStatus(status);
|
||||
const terminal = (result as AgentChatResultResponse).terminal === true || isTerminalMessageStatus(status);
|
||||
const terminal = turnResultIsTerminalForMerge(result as AgentChatResultResponse);
|
||||
reduceServerState({
|
||||
type: "turn.status",
|
||||
turn: {
|
||||
@@ -675,8 +691,8 @@ export const useWorkbenchStore = defineStore("workbench", () => {
|
||||
function syncTurnStatusToMessage(traceId: string, result: AgentChatResultResponse | TraceSnapshot): void {
|
||||
const authoritySessionId = traceResultSessionId(result);
|
||||
updateTraceMessages(traceId, authoritySessionId, (message) => {
|
||||
const resultStatus = normalizedStatusText((result as AgentChatResultResponse).status) ?? null;
|
||||
const terminal = (result as AgentChatResultResponse).terminal === true || isTerminalMessageStatus(resultStatus);
|
||||
const resultStatus = turnResultStatusForMerge(result as AgentChatResultResponse);
|
||||
const terminal = turnResultIsTerminalForMerge(result as AgentChatResultResponse);
|
||||
const resultError = normalizeAgentError((result as AgentChatResultResponse).error ?? null);
|
||||
const resultProjection = projectionFromResult(result as AgentChatResultResponse);
|
||||
const mergedRunnerTrace = mergeTerminalResultTrace(message.runnerTrace, result as AgentChatResultResponse);
|
||||
@@ -768,7 +784,7 @@ export const useWorkbenchStore = defineStore("workbench", () => {
|
||||
applyTurnStatusSnapshot(canonicalTraceId, response.data);
|
||||
publishWorkbenchProjectionSignal(sessionId, canonicalTraceId, "submit-admitted");
|
||||
scheduleSessionListRefresh(sessionId, runtimePolicy.sessionListTerminalRefreshDelayMs);
|
||||
if ((response.data as AgentChatResultResponse).terminal === true || isTerminalMessageStatus(response.data.status)) {
|
||||
if (turnResultIsTerminalForMerge(response.data as AgentChatResultResponse)) {
|
||||
completeTrace(canonicalTraceId, response.data as AgentChatResultResponse);
|
||||
return true;
|
||||
}
|
||||
@@ -912,7 +928,7 @@ export const useWorkbenchStore = defineStore("workbench", () => {
|
||||
if (!ownerSessionId) return;
|
||||
const events = Array.isArray(result.events) ? result.events : Array.isArray(result.traceEvents) ? result.traceEvents : [];
|
||||
const activityLabel = firstNonEmptyString(result.lastEventLabel, result.status);
|
||||
if (ownerSessionId === activeSessionId.value && (events.length > 0 || result.terminal === true || isTerminalMessageStatus(result.status))) recordActivity(`trace:${activityLabel ?? "hydrated"}`);
|
||||
if (ownerSessionId === activeSessionId.value && (events.length > 0 || turnResultIsTerminalForMerge(result))) recordActivity(`trace:${activityLabel ?? "hydrated"}`);
|
||||
markWorkbenchTraceEventsReceived({ traceId, events, transport: "rest_gap" });
|
||||
updateSessionMessages(ownerSessionId, (source) => source.map((message) => {
|
||||
if (!messageMatchesTraceAuthority(message, traceId, authoritySessionId, ownerSessionId)) return message;
|
||||
@@ -947,8 +963,8 @@ export const useWorkbenchStore = defineStore("workbench", () => {
|
||||
lastEventLabel: result.lastEventLabel ?? undefined,
|
||||
updatedAt: new Date().toISOString()
|
||||
} as NonNullable<ChatMessage["runnerTrace"]>;
|
||||
const resultStatus = normalizedStatusText(result.status) ?? null;
|
||||
const terminal = result.terminal === true || isTerminalMessageStatus(resultStatus);
|
||||
const resultStatus = turnResultStatusForMerge(result);
|
||||
const terminal = turnResultIsTerminalForMerge(result);
|
||||
const resultError = normalizeAgentError(result.error ?? null);
|
||||
const resultProjection = projectionFromResult(result);
|
||||
const mergedRunnerTrace = mergeRunnerTrace(message.runnerTrace, nextTrace);
|
||||
@@ -961,10 +977,10 @@ export const useWorkbenchStore = defineStore("workbench", () => {
|
||||
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() };
|
||||
}));
|
||||
markWorkbenchTraceProjected(traceId);
|
||||
if (traceResultHasTerminalEvidence(result) && result.terminal !== true && !isTerminalMessageStatus(result.status)) {
|
||||
if (traceResultHasTerminalEvidence(result) && !turnResultIsTerminalForMerge(result)) {
|
||||
void refreshTerminalTraceFromRest(traceId, "trace-hydration-terminal-evidence");
|
||||
}
|
||||
if (result.terminal === true || isTerminalMessageStatus(result.status)) {
|
||||
if (turnResultIsTerminalForMerge(result)) {
|
||||
rememberTurnStatus(traceId, result);
|
||||
if (ownerSessionId === activeSessionId.value) {
|
||||
chatPending.value = false;
|
||||
@@ -1102,7 +1118,12 @@ export const useWorkbenchStore = defineStore("workbench", () => {
|
||||
const ownerMessages = ownerSessionId ? serverState.value.messagesBySessionId[ownerSessionId] ?? [] : messages.value;
|
||||
const message = latestMessageForTrace(id, ownerMessages);
|
||||
const turn = turnStatusAuthority.value[id];
|
||||
const terminal = turn?.terminal === true || isTerminalMessageStatus(turn?.status) || isTerminalMessageStatus(message?.status);
|
||||
const turnStatus = normalizedStatusText(turn?.status);
|
||||
const messageStatus = normalizedStatusText(message?.status);
|
||||
const messageCompletedWithFinal = messageStatus === "completed" && message ? messageHasCompletedFinalResponse(message) : false;
|
||||
const messageTerminal = isTerminalMessageStatus(messageStatus) && (messageStatus !== "completed" || messageCompletedWithFinal);
|
||||
const turnCompletedMissingFinal = (turn?.terminal === true || turnStatus === "completed") && !messageCompletedWithFinal;
|
||||
const terminal = !turnCompletedMissingFinal && (turn?.terminal === true || isTerminalMessageStatus(turnStatus) || messageTerminal);
|
||||
if (terminal) {
|
||||
clearActiveTraceRestGapFill(id);
|
||||
return;
|
||||
@@ -1339,7 +1360,7 @@ export const useWorkbenchStore = defineStore("workbench", () => {
|
||||
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);
|
||||
const terminal = turnResultIsTerminalForMerge({ ...(snapshot as AgentChatResultResponse), status: traceStatus } as AgentChatResultResponse);
|
||||
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;
|
||||
@@ -1375,8 +1396,8 @@ export const useWorkbenchStore = defineStore("workbench", () => {
|
||||
if (ownerSessionId === activeSessionId.value) recordActivity(`trace:terminal:${firstNonEmptyString(result.lastEventLabel, result.status, "completed") ?? "completed"}`);
|
||||
updateSessionMessages(ownerSessionId, (source) => source.map((message) => {
|
||||
if (!messageMatchesTraceAuthority(message, traceId, authoritySessionId, ownerSessionId)) return message;
|
||||
const resultStatus = normalizedStatusText(result.status) ?? null;
|
||||
const terminal = result.terminal === true || isTerminalMessageStatus(resultStatus);
|
||||
const resultStatus = turnResultStatusForMerge(result);
|
||||
const terminal = turnResultIsTerminalForMerge(result);
|
||||
const resultError = normalizeAgentError(result.error ?? null);
|
||||
const resultProjection = projectionFromResult(result);
|
||||
const mergedRunnerTrace = mergeTerminalResultTrace(message.runnerTrace, result);
|
||||
@@ -1416,8 +1437,8 @@ export const useWorkbenchStore = defineStore("workbench", () => {
|
||||
if (!ownerSessionId) return;
|
||||
updateSessionMessages(ownerSessionId, (source) => source.map((message) => {
|
||||
if (!messageMatchesTraceAuthority(message, traceId, authoritySessionId, ownerSessionId)) return message;
|
||||
const resultStatus = normalizedStatusText(result.status) ?? null;
|
||||
const terminal = result.terminal === true || isTerminalMessageStatus(resultStatus);
|
||||
const resultStatus = turnResultStatusForMerge(result);
|
||||
const terminal = turnResultIsTerminalForMerge(result);
|
||||
const resultError = normalizeAgentError(result.error ?? null);
|
||||
const resultProjection = projectionFromResult(result);
|
||||
const mergedRunnerTrace = mergeTerminalResultTrace(message.runnerTrace, result);
|
||||
@@ -1464,7 +1485,7 @@ export const useWorkbenchStore = defineStore("workbench", () => {
|
||||
function messageHasSealedTerminalResult(message: ChatMessage | null): boolean {
|
||||
if (!message) return false;
|
||||
if (!isTerminalMessageStatus(message.status)) return false;
|
||||
return Boolean(firstNonEmptyString(message.text, finalResponseText((message as Record<string, unknown>).finalResponse)));
|
||||
return Boolean(firstNonEmptyString(message.text, messageText((message as Record<string, unknown>).content), finalResponseText((message as Record<string, unknown>).finalResponse)));
|
||||
}
|
||||
|
||||
function failTrace(traceId: string, message: string): void {
|
||||
|
||||
Reference in New Issue
Block a user