From 19ea9d0962095176d4cf623b3e60119515962f7a Mon Sep 17 00:00:00 2001 From: UniDesk Codex Date: Wed, 1 Jul 2026 19:41:49 +0800 Subject: [PATCH] fix: enforce workbench terminal seal invariant --- docs/reference/cloud-workbench.md | 2 +- internal/cloud/server-code-agent-http.ts | 44 +++++++++- internal/cloud/server-workbench-http.test.ts | 78 ++++++++++++++++- .../cloud/workbench-projection-writer.test.ts | 83 +++++++++++++++++++ internal/cloud/workbench-projection-writer.ts | 8 +- internal/cloud/workbench-turn-projection.ts | 32 +++++-- internal/db/runtime-store.test.ts | 34 +++++++- internal/db/runtime-store.ts | 36 +++++++- ...rkbench-message-projection-runtime.test.ts | 49 +++++++++++ .../workbench-message-projection-runtime.ts | 29 ++++++- .../src/stores/workbench-server-state.test.ts | 32 +++++++ .../src/stores/workbench-server-state.ts | 42 ++++++++-- .../src/stores/workbench-session.test.ts | 29 +++++++ .../src/stores/workbench-session.ts | 26 +++++- web/hwlab-cloud-web/src/stores/workbench.ts | 65 ++++++++++----- 15 files changed, 531 insertions(+), 58 deletions(-) create mode 100644 web/hwlab-cloud-web/src/stores/workbench-message-projection-runtime.test.ts create mode 100644 web/hwlab-cloud-web/src/stores/workbench-server-state.test.ts create mode 100644 web/hwlab-cloud-web/src/stores/workbench-session.test.ts diff --git a/docs/reference/cloud-workbench.md b/docs/reference/cloud-workbench.md index c1f2f6d7..45cfc994 100644 --- a/docs/reference/cloud-workbench.md +++ b/docs/reference/cloud-workbench.md @@ -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。 diff --git a/internal/cloud/server-code-agent-http.ts b/internal/cloud/server-code-agent-http.ts index bf42de08..595e4d78 100644 --- a/internal/cloud/server-code-agent-http.ts +++ b/internal/cloud/server-code-agent-http.ts @@ -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) { diff --git a/internal/cloud/server-workbench-http.test.ts b/internal/cloud/server-workbench-http.test.ts index 815f0a42..dc4cb744 100644 --- a/internal/cloud/server-workbench-http.test.ts +++ b/internal/cloud/server-workbench-http.test.ts @@ -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 () => { diff --git a/internal/cloud/workbench-projection-writer.test.ts b/internal/cloud/workbench-projection-writer.test.ts index 6e8e7d58..895f0123 100644 --- a/internal/cloud/workbench-projection-writer.test.ts +++ b/internal/cloud/workbench-projection-writer.test.ts @@ -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 = { diff --git a/internal/cloud/workbench-projection-writer.ts b/internal/cloud/workbench-projection-writer.ts index 1e965314..9ed4dbca 100644 --- a/internal/cloud/workbench-projection-writer.ts +++ b/internal/cloud/workbench-projection-writer.ts @@ -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(); diff --git a/internal/cloud/workbench-turn-projection.ts b/internal/cloud/workbench-turn-projection.ts index b9a74172..b185445d 100644 --- a/internal/cloud/workbench-turn-projection.ts +++ b/internal/cloud/workbench-turn-projection.ts @@ -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; } diff --git a/internal/db/runtime-store.test.ts b/internal/db/runtime-store.test.ts index 0b623cfe..d46f4a98 100644 --- a/internal/db/runtime-store.test.ts +++ b/internal/db/runtime-store.test.ts @@ -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 }); diff --git a/internal/db/runtime-store.ts b/internal/db/runtime-store.ts index 226aae21..2434241d 100644 --- a/internal/db/runtime-store.ts +++ b/internal/db/runtime-store.ts @@ -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, ""); diff --git a/web/hwlab-cloud-web/src/stores/workbench-message-projection-runtime.test.ts b/web/hwlab-cloud-web/src/stores/workbench-message-projection-runtime.test.ts new file mode 100644 index 00000000..895e189b --- /dev/null +++ b/web/hwlab-cloud-web/src/stores/workbench-message-projection-runtime.test.ts @@ -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"); +}); diff --git a/web/hwlab-cloud-web/src/stores/workbench-message-projection-runtime.ts b/web/hwlab-cloud-web/src/stores/workbench-message-projection-runtime.ts index 5f5cfd82..2ddf2742 100644 --- a/web/hwlab-cloud-web/src/stores/workbench-message-projection-runtime.ts +++ b/web/hwlab-cloud-web/src/stores/workbench-message-projection-runtime.ts @@ -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 | 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).finalResponse))); + return Boolean(firstNonEmptyString(message.text, messageText((message as Record).content), finalResponseText((message as Record).finalResponse))); } export function traceHasEvents(trace: ChatMessage["runnerTrace"]): boolean { diff --git a/web/hwlab-cloud-web/src/stores/workbench-server-state.test.ts b/web/hwlab-cloud-web/src/stores/workbench-server-state.test.ts new file mode 100644 index 00000000..7a90bb9b --- /dev/null +++ b/web/hwlab-cloud-web/src/stores/workbench-server-state.test.ts @@ -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"); +}); diff --git a/web/hwlab-cloud-web/src/stores/workbench-server-state.ts b/web/hwlab-cloud-web/src/stores/workbench-server-state.ts index 70c78cc5..04c33968 100644 --- a/web/hwlab-cloud-web/src/stores/workbench-server-state.ts +++ b/web/hwlab-cloud-web/src/stores/workbench-server-state.ts @@ -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 = 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).finalResponse); + return textValue(message.text) ?? textValue(message.content) ?? nestedTextValue((message as Record).finalResponse); } function projectionFromMessageRecord(message: ChatMessage): ProjectionDiagnostic | null { diff --git a/web/hwlab-cloud-web/src/stores/workbench-session.test.ts b/web/hwlab-cloud-web/src/stores/workbench-session.test.ts new file mode 100644 index 00000000..3dcda5bf --- /dev/null +++ b/web/hwlab-cloud-web/src/stores/workbench-session.test.ts @@ -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); +}); diff --git a/web/hwlab-cloud-web/src/stores/workbench-session.ts b/web/hwlab-cloud-web/src/stores/workbench-session.ts index 744c3eaf..39f41a20 100644 --- a/web/hwlab-cloud-web/src/stores/workbench-session.ts +++ b/web/hwlab-cloud-web/src/stores/workbench-session.ts @@ -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).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); diff --git a/web/hwlab-cloud-web/src/stores/workbench.ts b/web/hwlab-cloud-web/src/stores/workbench.ts index a9f45599..2e0834f0 100644 --- a/web/hwlab-cloud-web/src/stores/workbench.ts +++ b/web/hwlab-cloud-web/src/stores/workbench.ts @@ -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 { 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; - 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).finalResponse))); + return Boolean(firstNonEmptyString(message.text, messageText((message as Record).content), finalResponseText((message as Record).finalResponse))); } function failTrace(traceId: string, message: string): void {