Merge pull request #2337 from pikasTech/fix/1415-terminal-body
fix: seal workbench failed terminal body
This commit is contained in:
@@ -2556,8 +2556,6 @@ function codeAgentCompatProjectionPayload(payload = {}, context = {}) {
|
||||
|
||||
function codeAgentResultPollTerminalReady(body = {}) {
|
||||
if (body?.terminal !== true) return false;
|
||||
const status = normalizeTurnStatus(body.status);
|
||||
if (status !== "completed") return true;
|
||||
return Boolean(messageAuthorityTextValue(body.finalResponse ?? body.reply ?? body.terminalEvidence?.finalResponse));
|
||||
}
|
||||
|
||||
@@ -2606,7 +2604,7 @@ export function codeAgentTurnStatusPayload({ traceId, result, snapshot, resultPo
|
||||
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 terminalSealBlocked = terminalStatus === "completed" && !codeAgentCompletedTurnFinalText(finalResponse, resultObject, snapshotObject);
|
||||
const terminalSealBlocked = Boolean(terminalStatus && isTurnTerminalStatus(terminalStatus) && !codeAgentTerminalFinalText(finalResponse, resultObject, snapshotObject));
|
||||
const rawStatus = normalizeTurnStatus(
|
||||
terminalStatus,
|
||||
resultObject?.agentRun?.commandState,
|
||||
@@ -2662,7 +2660,7 @@ export function codeAgentTurnStatusPayload({ traceId, result, snapshot, resultPo
|
||||
};
|
||||
}
|
||||
|
||||
function codeAgentCompletedTurnFinalText(finalResponse = null, resultObject = null, snapshotObject = null) {
|
||||
function codeAgentTerminalFinalText(finalResponse = null, resultObject = null, snapshotObject = null) {
|
||||
const candidates = [
|
||||
finalResponse,
|
||||
resultObject?.finalResponse,
|
||||
@@ -2679,7 +2677,10 @@ function codeAgentCompletedTurnFinalText(finalResponse = null, resultObject = nu
|
||||
|
||||
function codeAgentAuthoritativeTerminalStatus(resultObject, snapshotObject, traceId) {
|
||||
const sealedPayload = codeAgentPayloadHasSealedFinalResponse(resultObject ?? {}) ? resultObject : codeAgentPayloadHasSealedFinalResponse(snapshotObject ?? {}) ? snapshotObject : null;
|
||||
if (sealedPayload) return "completed";
|
||||
if (sealedPayload) {
|
||||
const sealedStatus = normalizeTurnStatus(sealedPayload?.finalResponse?.status, sealedPayload?.terminalEvidence?.finalResponse?.status, sealedPayload?.agentRun?.terminalStatus, sealedPayload?.agentRun?.status, sealedPayload?.status);
|
||||
return sealedStatus && isTurnTerminalStatus(sealedStatus) ? sealedStatus : "completed";
|
||||
}
|
||||
const snapshotStatus = normalizeTurnStatus(snapshotObject?.terminalEvidence?.traceSummary?.terminalStatus, snapshotObject?.terminalEvidence?.agentRun?.terminalStatus, snapshotObject?.terminalEvidence?.status, snapshotObject?.status);
|
||||
if (snapshotStatus && isTurnTerminalStatus(snapshotStatus) && codeAgentSnapshotHasTerminalAuthority(snapshotObject)) return snapshotStatus;
|
||||
const agentRunStatus = normalizeTurnStatus(resultObject?.agentRun?.terminalStatus, resultObject?.agentRun?.commandState, resultObject?.agentRun?.status, resultObject?.agentRun?.runStatus);
|
||||
|
||||
@@ -289,6 +289,70 @@ test("code agent turn status keeps completed AgentRun without final response uns
|
||||
assert.equal(payload.finalResponse, null);
|
||||
});
|
||||
|
||||
test("code agent turn status keeps failed AgentRun without final response unsealed", () => {
|
||||
const traceId = "trc_code_agent_failed_without_final";
|
||||
const payload = codeAgentTurnStatusPayload({
|
||||
traceId,
|
||||
result: {
|
||||
traceId,
|
||||
status: "failed",
|
||||
agentRun: {
|
||||
runId: "run_code_agent_failed_without_final",
|
||||
commandId: "cmd_code_agent_failed_without_final",
|
||||
status: "failed",
|
||||
terminalStatus: "failed"
|
||||
},
|
||||
runnerTrace: {
|
||||
traceId,
|
||||
status: "failed",
|
||||
events: [
|
||||
{ seq: 1, type: "assistant_message", status: "running", message: "progress only" },
|
||||
{ seq: 2, type: "terminal_status", terminal: true, terminalStatus: "failed" }
|
||||
]
|
||||
}
|
||||
},
|
||||
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, "failed");
|
||||
assert.equal(payload.terminalSealBlocked, true);
|
||||
assert.equal(payload.waitingFor, "final_response");
|
||||
assert.equal(payload.finalResponse, null);
|
||||
});
|
||||
|
||||
test("code agent turn status seals failed AgentRun when final response is authoritative", () => {
|
||||
const traceId = "trc_code_agent_failed_with_final";
|
||||
const payload = codeAgentTurnStatusPayload({
|
||||
traceId,
|
||||
result: {
|
||||
traceId,
|
||||
status: "failed",
|
||||
finalResponse: { text: "Workbench terminal failed: provider-stream-disconnected", status: "failed" },
|
||||
agentRun: {
|
||||
runId: "run_code_agent_failed_with_final",
|
||||
commandId: "cmd_code_agent_failed_with_final",
|
||||
status: "failed",
|
||||
terminalStatus: "failed"
|
||||
}
|
||||
},
|
||||
snapshot: null,
|
||||
resultPollError: null,
|
||||
refreshError: null,
|
||||
options: { env: {} }
|
||||
});
|
||||
assert.equal(payload.status, "failed");
|
||||
assert.equal(payload.running, false);
|
||||
assert.equal(payload.terminal, true);
|
||||
assert.equal(payload.terminalSealBlocked, false);
|
||||
assert.equal(payload.finalResponse.text, "Workbench terminal failed: provider-stream-disconnected");
|
||||
});
|
||||
|
||||
test("code agent turn status seals completed AgentRun when finalText is authoritative", () => {
|
||||
const traceId = "trc_code_agent_terminal_with_final_text";
|
||||
const payload = codeAgentTurnStatusPayload({
|
||||
@@ -1317,6 +1381,67 @@ test("workbench read model seals turn from terminal message projection when chec
|
||||
}
|
||||
});
|
||||
|
||||
test("workbench read model exposes failed terminal body when final response part is missing", async () => {
|
||||
const sessionId = "ses_workbench_failed_terminal_missing_part";
|
||||
const traceId = "trc_failed_terminal_missing_part";
|
||||
const startedAt = "2026-06-29T22:19:32.000Z";
|
||||
const finishedAt = "2026-06-29T22:19:47.000Z";
|
||||
const facts = emptyFacts();
|
||||
facts.sessions.push({
|
||||
sessionId,
|
||||
ownerUserId: ACTOR.id,
|
||||
ownerRole: ACTOR.role,
|
||||
agentId: "hwlab-code-agent",
|
||||
status: "failed",
|
||||
lastTraceId: traceId,
|
||||
projectedSeq: 40,
|
||||
sourceSeq: 40,
|
||||
terminal: true,
|
||||
sealed: true,
|
||||
createdAt: startedAt,
|
||||
updatedAt: finishedAt,
|
||||
valuesRedacted: true
|
||||
});
|
||||
facts.messages.push(
|
||||
{ messageId: "msg_failed_missing_part_user", sessionId, turnId: traceId, traceId, role: "user", status: "sent", text: "hi", projectedSeq: 39, sourceSeq: 39, createdAt: startedAt, updatedAt: startedAt, valuesRedacted: true },
|
||||
{ messageId: "msg_failed_missing_part_agent", sessionId, turnId: traceId, traceId, role: "agent", status: "failed", text: "", projectedSeq: 40, sourceSeq: 40, terminal: true, sealed: true, timing: { startedAt, lastEventAt: finishedAt, finishedAt, durationMs: 15000, valuesRedacted: true }, startedAt, lastEventAt: finishedAt, finishedAt, durationMs: 15000, createdAt: startedAt, updatedAt: finishedAt, valuesRedacted: true }
|
||||
);
|
||||
facts.turns.push({ turnId: traceId, sessionId, traceId, messageId: "msg_failed_missing_part_agent", status: "failed", projectedSeq: 40, sourceSeq: 40, terminal: true, sealed: true, finalResponse: null, diagnostic: { blocker: { code: "provider-stream-disconnected", message: "provider stream disconnected" }, projectionStatus: "caught_up", projectionHealth: "healthy", valuesRedacted: true }, startedAt, lastEventAt: finishedAt, finishedAt, durationMs: 15000, updatedAt: finishedAt, valuesRedacted: true });
|
||||
facts.checkpoints.push({ traceId, sessionId, turnId: traceId, runId: "run_failed_missing_part", commandId: "cmd_failed_missing_part", status: "failed", projectionStatus: "caught_up", projectionHealth: "healthy", projectedSeq: 40, sourceSeq: 40, terminal: true, sealed: true, diagnostic: { blocker: { code: "provider-stream-disconnected", message: "provider stream disconnected" }, valuesRedacted: true }, startedAt, lastEventAt: finishedAt, finishedAt, durationMs: 15000, updatedAt: finishedAt, valuesRedacted: true });
|
||||
const runtimeStore = createRuntimeStoreFromFacts(facts);
|
||||
const accessController = {
|
||||
async ensureBootstrap() {},
|
||||
async authenticate() { return { ok: true, actor: ACTOR, session: { id: "uss_workbench_reader" } }; }
|
||||
};
|
||||
const server = createCloudApiServer({ accessController, runtimeStore });
|
||||
await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve));
|
||||
|
||||
try {
|
||||
const { port } = server.address();
|
||||
const sessions = await getJson(port, `/v1/workbench/sessions?includeSessionId=${encodeURIComponent(sessionId)}`);
|
||||
assert.equal(sessions.status, 200);
|
||||
assert.equal(sessions.body.sessions[0].status, "failed");
|
||||
assert.equal(sessions.body.sessions[0].terminal, true);
|
||||
|
||||
const messages = await getJson(port, `/v1/workbench/sessions/${encodeURIComponent(sessionId)}/messages?limit=10`);
|
||||
assert.equal(messages.status, 200);
|
||||
const assistant = messages.body.messages.find((message) => message.role === "agent");
|
||||
assert.equal(assistant.status, "failed");
|
||||
assert.match(assistant.text, /provider-stream-disconnected/u);
|
||||
assert.equal(assistant.parts[0].type, "final_response");
|
||||
assert.match(assistant.parts[0].text, /provider-stream-disconnected/u);
|
||||
|
||||
const turn = await getJson(port, `/v1/workbench/turns/${encodeURIComponent(traceId)}`);
|
||||
assert.equal(turn.status, 200);
|
||||
assert.equal(turn.body.turn.status, "failed");
|
||||
assert.equal(turn.body.turn.running, false);
|
||||
assert.equal(turn.body.turn.terminal, true);
|
||||
assert.match(turn.body.turn.finalResponse.text, /provider-stream-disconnected/u);
|
||||
} finally {
|
||||
await new Promise((resolve, reject) => server.close((error) => error ? reject(error) : resolve()));
|
||||
}
|
||||
});
|
||||
|
||||
test("workbench read model does not expose trace-only memory without visible session or result owner", async () => {
|
||||
const traceStore = createCodeAgentTraceStore();
|
||||
const traceId = "trc_workbench_trace_only";
|
||||
|
||||
@@ -15,7 +15,7 @@ import {
|
||||
} from "./server-http-utils.ts";
|
||||
import { createWorkbenchReadModel } from "./workbench-read-model.ts";
|
||||
import { createWorkbenchRuntimeClient } from "./workbench-runtime-client.ts";
|
||||
import { durableTraceStatus, RUNNING_STATUSES, TERMINAL_STATUSES } from "./workbench-turn-projection.ts";
|
||||
import { durableTraceStatus, RUNNING_STATUSES, terminalFinalResponse, TERMINAL_STATUSES } from "./workbench-turn-projection.ts";
|
||||
import { emitCodeAgentOtelSpan, emitHttpServerRequestSpan } from "./otel-trace.ts";
|
||||
|
||||
const DEFAULT_PAGE_LIMIT = 50;
|
||||
@@ -1189,10 +1189,17 @@ function factMessageDto(message, parts = [], facts = {}) {
|
||||
const messageStatus = normalizeStatus(message?.status);
|
||||
const assistantLike = isAssistantLikeRole(role);
|
||||
const legacyText = projectionText(message?.text, message?.content, message?.message, message?.finalResponse);
|
||||
const normalizedParts = parts.length > 0
|
||||
const baseParts = parts.length > 0
|
||||
? [...parts].sort(compareFactPartsAsc).map((part) => factPartDto(part, messageId, traceId)).filter(Boolean)
|
||||
: !assistantLike && legacyText ? [partFact({ type: "text", text: legacyText, status: message?.status }, 0, messageId, traceId)] : [];
|
||||
const checkpointStatus = assistantLike ? normalizeTerminalStatus(checkpoint?.status) : null;
|
||||
const preliminaryStatus = normalizeStatus(checkpointStatus ?? (assistantLike ? turn?.status : null) ?? messageStatus);
|
||||
const syntheticTerminalFinalResponse = assistantLike && !factPartsHaveFinalResponse(baseParts)
|
||||
? factSyntheticTerminalFinalResponse(preliminaryStatus, traceId, message, checkpoint, turn)
|
||||
: null;
|
||||
const normalizedParts = syntheticTerminalFinalResponse
|
||||
? [...baseParts, partFact({ type: "final_response", text: syntheticTerminalFinalResponse.text, status: syntheticTerminalFinalResponse.status }, baseParts.length, messageId, traceId)]
|
||||
: baseParts;
|
||||
const messageTerminalStatus = assistantLike ? factTerminalStatusFromMessageProjection(message, normalizedParts) : null;
|
||||
const status = normalizeStatus(checkpointStatus ?? messageTerminalStatus ?? (assistantLike ? turn?.status : null) ?? messageStatus);
|
||||
const timing = assistantLike
|
||||
@@ -1228,6 +1235,16 @@ function factMessageAuthorityText({ role, status, parts = [], legacyText = null
|
||||
return firstFactPartText(parts, "final_response") ?? "";
|
||||
}
|
||||
|
||||
function factPartsHaveFinalResponse(parts = []) {
|
||||
return parts.some((part) => part?.type === "final_response" && Boolean(projectionText(part?.text)));
|
||||
}
|
||||
|
||||
function factSyntheticTerminalFinalResponse(status, traceId = null, ...records) {
|
||||
const terminalStatus = normalizeTerminalStatus(status);
|
||||
if (!terminalStatus || terminalStatus === "completed") return null;
|
||||
return terminalFinalResponse(terminalStatus, { traceId, status: terminalStatus, records }, { evidence: records });
|
||||
}
|
||||
|
||||
function isTerminalProjectionStatus(status) {
|
||||
return TERMINAL_STATUSES.has(status) && !RUNNING_STATUSES.has(status);
|
||||
}
|
||||
|
||||
@@ -273,6 +273,95 @@ test("workbench projection writer keeps completed AgentRun without final respons
|
||||
assert.equal(facts.checkpoints[0].sealed, false);
|
||||
});
|
||||
|
||||
test("workbench projection writer seals failed AgentRun turns with failure final response", 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:03:00.000Z"
|
||||
};
|
||||
}
|
||||
};
|
||||
|
||||
await writeWorkbenchProjectionSession({
|
||||
accessController,
|
||||
runtimeStore,
|
||||
traceId: "trc_writer_terminal_failed",
|
||||
ownerUserId: "usr_writer",
|
||||
ownerRole: "user",
|
||||
sessionId: "ses_writer_terminal_failed",
|
||||
projectId: "prj_writer",
|
||||
conversationId: "cnv_writer_terminal_failed",
|
||||
threadId: "thread-writer-terminal-failed",
|
||||
status: "failed",
|
||||
payload: {
|
||||
traceId: "trc_writer_terminal_failed",
|
||||
status: "failed",
|
||||
error: { code: "provider-stream-disconnected", message: "provider stream disconnected" },
|
||||
runnerTrace: {
|
||||
traceId: "trc_writer_terminal_failed",
|
||||
status: "failed",
|
||||
startedAt: "2026-06-20T11:02:30.000Z",
|
||||
lastEventAt: "2026-06-20T11:03:00.000Z",
|
||||
finishedAt: "2026-06-20T11:03:00.000Z",
|
||||
events: [
|
||||
{ seq: 1, type: "assistant_message", status: "running", message: "progress only" },
|
||||
{ seq: 2, type: "result", status: "failed", terminal: true, label: "agentrun:terminal:failed", errorCode: "provider-stream-disconnected" }
|
||||
],
|
||||
eventCount: 2
|
||||
},
|
||||
agentRun: { runId: "run_writer_terminal_failed", commandId: "cmd_writer_terminal_failed", status: "failed", terminalStatus: "failed", failureKind: "provider-stream-disconnected", lastSeq: 2 },
|
||||
updatedAt: "2026-06-20T11:03:00.000Z"
|
||||
},
|
||||
session: {
|
||||
sessionStatus: "failed",
|
||||
messages: [
|
||||
{ messageId: "msg_writer_failed_user", role: "user", text: "question", status: "sent", turnId: "trc_writer_terminal_failed", traceId: "trc_writer_terminal_failed" },
|
||||
{ messageId: "msg_writer_failed_agent", role: "agent", text: "", status: "failed", turnId: "trc_writer_terminal_failed", traceId: "trc_writer_terminal_failed" }
|
||||
]
|
||||
}
|
||||
});
|
||||
|
||||
assert.equal(factWrites.length, 1);
|
||||
const facts = factWrites[0].params.facts;
|
||||
const agentMessage = facts.messages.find((message) => message.messageId === "msg_writer_failed_agent");
|
||||
const finalPart = facts.parts.find((part) => part.messageId === "msg_writer_failed_agent" && part.partType === "final_response");
|
||||
assert.equal(facts.sessions[0].status, "failed");
|
||||
assert.equal(facts.sessions[0].terminal, true);
|
||||
assert.equal(facts.sessions[0].sealed, true);
|
||||
assert.equal(agentMessage.status, "failed");
|
||||
assert.equal(agentMessage.terminal, true);
|
||||
assert.equal(agentMessage.sealed, true);
|
||||
assert.match(agentMessage.text, /provider-stream-disconnected/u);
|
||||
assert.equal(finalPart?.status, "failed");
|
||||
assert.match(finalPart?.text ?? "", /provider-stream-disconnected/u);
|
||||
assert.equal(facts.turns[0].status, "failed");
|
||||
assert.equal(facts.turns[0].terminal, true);
|
||||
assert.equal(facts.turns[0].sealed, true);
|
||||
assert.match(facts.turns[0].finalResponse.text, /provider-stream-disconnected/u);
|
||||
assert.equal(facts.turns[0].diagnostic.projectionStatus, "caught-up");
|
||||
assert.equal(facts.turns[0].diagnostic.projectionHealth, "healthy");
|
||||
assert.equal(facts.turns[0].diagnostic.blocker, null);
|
||||
assert.equal(facts.checkpoints[0].projectionStatus, "caught_up");
|
||||
assert.equal(facts.checkpoints[0].terminal, true);
|
||||
assert.equal(facts.checkpoints[0].sealed, true);
|
||||
assert.match(facts.checkpoints[0].finalResponse.text, /provider-stream-disconnected/u);
|
||||
});
|
||||
|
||||
test("workbench projection writer does not synthesize terminal duration from updatedAt", async () => {
|
||||
const factWrites = [];
|
||||
const runtimeStore = {
|
||||
|
||||
@@ -7,7 +7,7 @@ import { createHash } from "node:crypto";
|
||||
import { defaultCodeAgentTraceStore } from "./code-agent-trace-store.ts";
|
||||
import { emitCodeAgentOtelSpan } from "./otel-trace.ts";
|
||||
import { safeTraceId } from "./server-http-utils.ts";
|
||||
import { createWorkbenchTurnProjection, normalizeWorkbenchStatus, projectionDiagnostics, TERMINAL_STATUSES } from "./workbench-turn-projection.ts";
|
||||
import { createWorkbenchTurnProjection, normalizeWorkbenchStatus, projectionDiagnostics, terminalFinalResponse, TERMINAL_STATUSES } from "./workbench-turn-projection.ts";
|
||||
|
||||
export async function writeWorkbenchProjectionSession({ accessController, runtimeStore = null, traceStore = defaultCodeAgentTraceStore, traceId, ownerUserId, ownerRole = null, sessionId = null, projectId = null, conversationId = null, threadId = null, status = "active", session = {}, payload = null, params = {}, preserveLastTraceId = false } = {}) {
|
||||
if (!ownerUserId || typeof accessController?.recordAgentSessionOwner !== "function") return null;
|
||||
@@ -178,7 +178,8 @@ export async function writeWorkbenchProjectionEvent({ runtimeStore = null, event
|
||||
if (projectedSeq <= 0) throw new Error("Workbench projection allocator returned an invalid projectedSeq.");
|
||||
const eventId = isAssistantMessageTraceEvent(event) ? stableFactId("wte", { traceId, sourceEventId }) : textValue(event.id) || stableFactId("wte", { traceId, sourceEventId });
|
||||
const previousFinalText = finalResponseTextValue(resolvedPreviousCheckpoint?.finalResponse, resolvedPreviousCheckpoint?.assistantText, resolvedPreviousCheckpoint?.finalText);
|
||||
const messageFinalText = terminal ? finalResponseTextValue(event.finalResponse, event.assistantText, event.reply, previousFinalText) : null;
|
||||
const eventTerminalFinalResponse = terminal ? terminalFinalResponse(status, event, { evidence: event, finalResponse: event.finalResponse }) : null;
|
||||
const messageFinalText = terminal ? finalResponseTextValue(event.finalResponse, event.assistantText, event.reply, eventTerminalFinalResponse, previousFinalText) : null;
|
||||
const userMessageFact = sessionId && !suppressedAfterSeal && resolvedPreviousCheckpoint?.userMessage ? normalizeMessageFact(resolvedPreviousCheckpoint.userMessage, 0, { traceId, sessionId, turnId, terminal: false, terminalStatus: "sent", finalText: null, timestamp: projectedAt, timing }) : null;
|
||||
const messageFact = sessionId && !suppressedAfterSeal ? normalizeMessageFact({
|
||||
messageId: eventProjectionAssistantMessageId(traceId, event),
|
||||
|
||||
@@ -119,22 +119,22 @@ 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 sealedTerminal = turn.terminal === true && 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;
|
||||
const source = sealedTerminal ? null : projectionDiagnosticSource(trace);
|
||||
const rawBlocker = sealedTerminal ? null : refreshError ?? source?.blocker ?? trace?.blocker ?? result?.blocker ?? result?.error ?? null;
|
||||
const retryingProviderInterruption = !sealedTerminal && turn.running === true ? retryableProviderInterruptionEvidence(rawBlocker) : null;
|
||||
const blocker = retryingProviderInterruption ? null : rawBlocker;
|
||||
const hasProjectionInput = hasTraceProjection(trace) || Boolean(result || result?.agentRun);
|
||||
const sourceStatus = sealedCompleted ? null : normalizeProjectionStatus(source?.projectionStatus ?? trace?.projectionStatus);
|
||||
const sourceStatus = sealedTerminal ? null : normalizeProjectionStatus(source?.projectionStatus ?? trace?.projectionStatus);
|
||||
const effectiveSourceStatus = retryingProviderInterruption && sourceStatus === "blocked" ? "projecting" : sourceStatus;
|
||||
const status = sealedCompleted ? "caught-up" : waitingFor ? "projecting" : effectiveSourceStatus ?? (blocker ? "blocked" : turn.terminal ? "caught-up" : hasProjectionInput ? "projecting" : "unknown");
|
||||
const status = sealedTerminal ? "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 sourceHealth = sealedTerminal ? null : normalizeProjectionHealth(source?.projectionHealth ?? trace?.projectionHealth);
|
||||
const effectiveSourceHealth = retryingProviderInterruption && (sourceHealth === "degraded" || sourceHealth === "unavailable" || sourceHealth === "stalled") ? "projecting" : sourceHealth;
|
||||
const projectionHealth = sealedCompleted ? "healthy" : waitingFor ? "projecting" : effectiveSourceHealth
|
||||
const projectionHealth = sealedTerminal ? "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);
|
||||
const staleMs = sealedTerminal ? null : projectionStaleMs(source?.staleMs ?? trace?.staleMs, turn.updatedAt ?? trace?.updatedAt ?? result?.updatedAt);
|
||||
return {
|
||||
projectionStatus: status,
|
||||
projectionHealth,
|
||||
@@ -197,20 +197,22 @@ function terminalTurnEvidence({ result = null, traceTerminal = null } = {}) {
|
||||
return traceTerminal;
|
||||
}
|
||||
|
||||
function terminalFinalResponse(status, result = null, traceTerminal = null) {
|
||||
if (normalizeWorkbenchStatus(status) === "canceled") return canonicalCancelFinalResponse(result, traceTerminal);
|
||||
export function terminalFinalResponse(status, result = null, traceTerminal = null) {
|
||||
const normalizedStatus = normalizeWorkbenchStatus(status);
|
||||
if (normalizedStatus === "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),
|
||||
status: normalizedStatus,
|
||||
traceId: textValue(result?.traceId ?? traceTerminal?.finalResponse?.traceId ?? traceTerminal?.evidence?.traceId) || null,
|
||||
source: "result-final-response-evidence",
|
||||
valuesPrinted: false
|
||||
};
|
||||
}
|
||||
if (normalizedStatus !== "completed") return terminalFailureFinalResponse(normalizedStatus, result, traceTerminal);
|
||||
return null;
|
||||
}
|
||||
|
||||
@@ -223,6 +225,42 @@ function canonicalCancelFinalResponse(result = null, traceTerminal = null) {
|
||||
};
|
||||
}
|
||||
|
||||
function terminalFailureFinalResponse(status, result = null, traceTerminal = null) {
|
||||
const normalizedStatus = normalizeWorkbenchStatus(status);
|
||||
if (!normalizedStatus || normalizedStatus === "completed" || normalizedStatus === "unknown") return null;
|
||||
const failure = terminalFailureEvidence(result, traceTerminal?.evidence, traceTerminal);
|
||||
const failureKind = normalizedProjectionFailureKind(failure?.failureKind ?? failure?.errorCode ?? failure?.code ?? failure?.name ?? failure?.status ?? normalizedStatus);
|
||||
const label = failureKind && failureKind !== "unknown" ? failureKind : normalizedStatus;
|
||||
return {
|
||||
text: `Workbench terminal ${normalizedStatus}: ${label}`,
|
||||
status: normalizedStatus,
|
||||
traceId: textValue(result?.traceId ?? traceTerminal?.finalResponse?.traceId ?? traceTerminal?.evidence?.traceId) || null,
|
||||
source: "terminal-failure-diagnostic",
|
||||
failureKind: label,
|
||||
valuesPrinted: false,
|
||||
valuesRedacted: true
|
||||
};
|
||||
}
|
||||
|
||||
function terminalFailureEvidence(...values) {
|
||||
for (const value of values) {
|
||||
if (Array.isArray(value)) {
|
||||
const nested = terminalFailureEvidence(...value);
|
||||
if (nested) return nested;
|
||||
continue;
|
||||
}
|
||||
const record = objectValue(value);
|
||||
if (!record) continue;
|
||||
const kind = normalizedProjectionFailureKind(record.failureKind ?? record.errorCode ?? record.code ?? record.name);
|
||||
if (kind && kind !== "unknown") return record;
|
||||
for (const nested of [record.error, record.blocker, record.diagnostic, record.terminalEvidence, record.providerTrace, record.agentRun, record.payload, record.evidence, record.records]) {
|
||||
const found = terminalFailureEvidence(nested);
|
||||
if (found) return found;
|
||||
}
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
function resultHasTerminalAuthority(result = null, traceTerminal = null) {
|
||||
if (!result || typeof result !== "object") return false;
|
||||
if (traceTerminal) return true;
|
||||
|
||||
@@ -2,7 +2,10 @@ import assert from "node:assert/strict";
|
||||
import { test } from "bun:test";
|
||||
|
||||
import {
|
||||
messageHasTerminalResponse,
|
||||
messageNeedsTerminalDiagnostics,
|
||||
terminalMessagePatchFromTurnResult,
|
||||
traceHasTerminalResponse,
|
||||
turnResultIsTerminalForMerge,
|
||||
turnResultStatusForMerge
|
||||
} from "./workbench-message-projection-runtime";
|
||||
@@ -47,3 +50,36 @@ test("turn result merge seals completed result with final response", () => {
|
||||
assert.equal((patch?.finalResponse as any)?.text, "final answer");
|
||||
assert.equal(patch?.text, "final answer");
|
||||
});
|
||||
|
||||
test("turn result merge seals failed result with failure final response", () => {
|
||||
const traceId = "trc_frontend_terminal_failed_with_body";
|
||||
const message = { id: "msg_frontend_terminal_failed_with_body", role: "agent", status: "running", traceId } as any;
|
||||
const result = {
|
||||
traceId,
|
||||
status: "failed",
|
||||
terminal: true,
|
||||
finalResponse: { text: "Workbench terminal failed: provider-stream-disconnected", status: "failed" }
|
||||
} as any;
|
||||
|
||||
const patch = terminalMessagePatchFromTurnResult(message, result);
|
||||
assert.equal(turnResultStatusForMerge(result), "failed");
|
||||
assert.equal(turnResultIsTerminalForMerge(result), true);
|
||||
assert.equal(patch?.status, "failed");
|
||||
assert.equal((patch?.finalResponse as any)?.status, "failed");
|
||||
assert.equal((patch?.finalResponse as any)?.text, "Workbench terminal failed: provider-stream-disconnected");
|
||||
assert.equal(patch?.text, "Workbench terminal failed: provider-stream-disconnected");
|
||||
});
|
||||
|
||||
test("terminal response helpers treat failed body as sealed authority", () => {
|
||||
const traceId = "trc_frontend_terminal_failed_helper";
|
||||
const sealed = { id: "msg_frontend_terminal_failed_helper", role: "agent", status: "failed", traceId, text: "Workbench terminal failed: provider-stream-disconnected", finalResponse: { text: "Workbench terminal failed: provider-stream-disconnected" } } as any;
|
||||
const unsealed = { id: "msg_frontend_terminal_failed_helper_unsealed", role: "agent", status: "failed", traceId } as any;
|
||||
|
||||
assert.equal(messageHasTerminalResponse(sealed), true);
|
||||
assert.equal(messageNeedsTerminalDiagnostics(sealed), false);
|
||||
assert.equal(traceHasTerminalResponse(traceId, [sealed]), true);
|
||||
|
||||
assert.equal(messageHasTerminalResponse(unsealed), false);
|
||||
assert.equal(messageNeedsTerminalDiagnostics(unsealed), true);
|
||||
assert.equal(traceHasTerminalResponse(traceId, [unsealed]), false);
|
||||
});
|
||||
|
||||
@@ -325,20 +325,21 @@ export function terminalMessagePatchFromTurnResult(message: ChatMessage, result:
|
||||
}
|
||||
|
||||
function terminalMessageBodyPatchFromTurnResult(message: ChatMessage, result: AgentChatResultResponse, resultStatus: string | null): Partial<ChatMessage> {
|
||||
if (normalizedStatusText(resultStatus) !== "completed") return {};
|
||||
const status = normalizedStatusText(resultStatus);
|
||||
if (!isTerminalMessageStatus(status)) return {};
|
||||
const finalText = terminalFinalResponseTextFromTurnResult(result);
|
||||
if (!finalText) return {};
|
||||
const traceId = firstNonEmptyString(result.traceId, message.traceId, message.runnerTrace?.traceId) ?? null;
|
||||
return {
|
||||
finalResponse: {
|
||||
text: finalText,
|
||||
status: "completed",
|
||||
status,
|
||||
traceId,
|
||||
sealed: true,
|
||||
source: "turn-result",
|
||||
valuesRedacted: true
|
||||
},
|
||||
text: projectedAgentMessageText({ status: "completed", finalText })
|
||||
text: projectedAgentMessageText({ status, finalText })
|
||||
};
|
||||
}
|
||||
|
||||
@@ -448,7 +449,7 @@ export function agentErrorFromApiFailure(result: ApiResult<unknown>, projection:
|
||||
export function messageNeedsTerminalDiagnostics(message: ChatMessage): boolean {
|
||||
if (message.role !== "agent") return false;
|
||||
if (!isTerminalMessageStatus(message.status)) return false;
|
||||
if (messageHasCompletedFinalResponse(message)) return false;
|
||||
if (messageHasTerminalResponse(message)) return false;
|
||||
const traceId = firstNonEmptyString(message.traceId, message.runnerTrace?.traceId);
|
||||
if (!traceId) return false;
|
||||
const agentRun = agentRunFromMessage(message);
|
||||
@@ -475,15 +476,15 @@ function messageHasActiveTraceForHydration(message: ChatMessage): boolean {
|
||||
|| isTraceActiveStatus(trace?.traceStatus);
|
||||
}
|
||||
|
||||
export function traceHasCompletedFinalResponse(traceId: string | null | undefined, source: ChatMessage[]): boolean {
|
||||
export function traceHasTerminalResponse(traceId: string | null | undefined, source: ChatMessage[]): boolean {
|
||||
const id = firstNonEmptyString(traceId);
|
||||
if (!id) return false;
|
||||
return source.some((message) => firstNonEmptyString(message.traceId, message.runnerTrace?.traceId) === id && messageHasCompletedFinalResponse(message));
|
||||
return source.some((message) => firstNonEmptyString(message.traceId, message.runnerTrace?.traceId) === id && messageHasTerminalResponse(message));
|
||||
}
|
||||
|
||||
export function messageHasCompletedFinalResponse(message: ChatMessage | null | undefined): boolean {
|
||||
export function messageHasTerminalResponse(message: ChatMessage | null | undefined): boolean {
|
||||
if (!message || message.role !== "agent") return false;
|
||||
if (normalizedStatusText(message.status) !== "completed") return false;
|
||||
if (!isTerminalMessageStatus(message.status)) return false;
|
||||
return Boolean(firstNonEmptyString(message.text, messageText((message as Record<string, unknown>).content), finalResponseText((message as Record<string, unknown>).finalResponse)));
|
||||
}
|
||||
|
||||
|
||||
@@ -18,6 +18,21 @@ test("server state reducer keeps completed message without final response unseal
|
||||
assert.notEqual(selectSessionStatusAuthority(state)[sessionId]?.status, "completed");
|
||||
});
|
||||
|
||||
test("server state reducer keeps failed message without terminal body unsealed", () => {
|
||||
const sessionId = "ses_state_failed_without_body";
|
||||
const traceId = "trc_state_failed_without_body";
|
||||
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_failed_without_body", 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_failed_without_body", role: "agent", status: "failed", 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, "failed");
|
||||
});
|
||||
|
||||
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";
|
||||
@@ -30,3 +45,45 @@ test("server state reducer seals completed message when final response exists",
|
||||
assert.equal(message.text, "final answer");
|
||||
assert.equal(selectSessionStatusAuthority(state)[sessionId]?.status, "completed");
|
||||
});
|
||||
|
||||
test("server state reducer keeps sealed terminal messages when a stale partial snapshot arrives", () => {
|
||||
const sessionId = "ses_state_partial_snapshot";
|
||||
const traceId = "trc_state_partial_snapshot";
|
||||
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_partial_user_1", role: "user", status: "sent", text: "one", sessionId, traceId: "trc_state_partial_1" } as any,
|
||||
{ id: "msg_state_partial_agent_1", role: "agent", status: "completed", text: "one done", sessionId, traceId: "trc_state_partial_1" } as any,
|
||||
{ id: "msg_state_partial_user_2", role: "user", status: "sent", text: "two", sessionId, traceId } as any,
|
||||
{ id: "msg_state_partial_agent_2", role: "agent", status: "completed", text: "final answer", finalResponse: { text: "final answer" }, sessionId, traceId, finishedAt: "2026-07-01T00:00:04.000Z", durationMs: 4000 } as any
|
||||
] });
|
||||
state = reduceWorkbenchServerState(state, { type: "session.messages", sessionId, messages: [
|
||||
{ id: "msg_state_partial_user_1", role: "user", status: "sent", text: "one", sessionId, traceId: "trc_state_partial_1" } as any,
|
||||
{ id: "msg_state_partial_agent_1", role: "agent", status: "completed", text: "one done", sessionId, traceId: "trc_state_partial_1" } as any
|
||||
] });
|
||||
|
||||
const messages = selectActiveMessages(state, sessionId);
|
||||
assert.equal(messages.length, 4);
|
||||
const terminal = messages.find((message) => message.id === "msg_state_partial_agent_2") as any;
|
||||
assert.equal(terminal.status, "completed");
|
||||
assert.equal(terminal.text, "final answer");
|
||||
assert.equal(terminal.finalResponse.text, "final answer");
|
||||
assert.equal(selectSessionStatusAuthority(state)[sessionId]?.status, "completed");
|
||||
});
|
||||
|
||||
test("server state reducer promotes failed terminal message authority and ignores running list rollback", () => {
|
||||
const sessionId = "ses_state_failed_terminal_authority";
|
||||
const traceId = "trc_state_failed_terminal_authority";
|
||||
let state = createWorkbenchServerState();
|
||||
state = reduceWorkbenchServerState(state, { type: "session.detail", session: { sessionId, status: "running", lastTraceId: traceId, messages: [] } as any });
|
||||
state = reduceWorkbenchServerState(state, { type: "session.messages", sessionId, messages: [
|
||||
{ id: "msg_state_failed_user", role: "user", status: "sent", text: "one", sessionId, traceId } as any,
|
||||
{ id: "msg_state_failed_agent", role: "agent", status: "failed", text: "Workbench terminal failed: provider-stream-disconnected", finalResponse: { text: "Workbench terminal failed: provider-stream-disconnected" }, sessionId, traceId, finishedAt: "2026-07-01T00:00:05.000Z", durationMs: 5000 } as any
|
||||
] });
|
||||
|
||||
assert.equal(selectSessionStatusAuthority(state)[sessionId]?.status, "failed");
|
||||
state = reduceWorkbenchServerState(state, { type: "session.list", sessions: [{ sessionId, status: "running", lastTraceId: traceId, messages: [] } as any] });
|
||||
|
||||
assert.equal(selectSessionStatusAuthority(state)[sessionId]?.status, "failed");
|
||||
assert.equal(selectActiveMessages(state, sessionId).find((message) => message.id === "msg_state_failed_agent")?.status, "failed");
|
||||
});
|
||||
|
||||
@@ -188,7 +188,7 @@ function messageMatchesSnapshot(existing: ChatMessage, incoming: ChatMessage): b
|
||||
}
|
||||
|
||||
function mergeMessageSnapshot(existing: ChatMessage, incoming: ChatMessage): ChatMessage {
|
||||
const normalizedIncoming = normalizeUnsealedCompletedMessage(incoming);
|
||||
const normalizedIncoming = normalizeUnsealedTerminalMessage(incoming);
|
||||
const merged = {
|
||||
...existing,
|
||||
...normalizedIncoming,
|
||||
@@ -200,17 +200,23 @@ function mergeMessageSnapshot(existing: ChatMessage, incoming: ChatMessage): Cha
|
||||
}
|
||||
|
||||
function mergeMessageList(existing: ChatMessage[], incoming: ChatMessage[]): ChatMessage[] {
|
||||
return incoming.map((rawMessage) => {
|
||||
const message = normalizeUnsealedCompletedMessage(rawMessage);
|
||||
const previous = existing.find((item) => messageMatchesSnapshot(item, message));
|
||||
return previous ? sealExistingMessageTiming(previous, message) : message;
|
||||
});
|
||||
const merged = existing.map((message) => normalizeUnsealedTerminalMessage(message));
|
||||
for (const rawMessage of incoming) {
|
||||
const message = normalizeUnsealedTerminalMessage(rawMessage);
|
||||
const index = merged.findIndex((item) => messageMatchesSnapshot(item, message));
|
||||
if (index >= 0) {
|
||||
merged[index] = sealExistingMessageTiming(merged[index], message);
|
||||
} else {
|
||||
merged.push(message);
|
||||
}
|
||||
}
|
||||
return merged;
|
||||
}
|
||||
|
||||
function normalizeUnsealedCompletedMessage(message: ChatMessage): ChatMessage {
|
||||
function normalizeUnsealedTerminalMessage(message: ChatMessage): ChatMessage {
|
||||
if (message.role !== "agent") return message;
|
||||
if (normalizedMessageStatus(message.status) !== "completed") return message;
|
||||
if (messageHasCompletedFinalResponse(message)) return message;
|
||||
if (!isTerminalMessageStatus(message.status)) return message;
|
||||
if (messageHasTerminalResponse(message)) return message;
|
||||
const timing = message.timing ? { ...message.timing, finishedAt: null, durationMs: null } : message.timing;
|
||||
return {
|
||||
...message,
|
||||
@@ -394,11 +400,12 @@ function sessionStatusAuthorityFromDetail(session: WorkbenchSessionRecord): Sess
|
||||
}
|
||||
|
||||
function sessionStatusAuthorityFromMessages(sessionId: string, messages: ChatMessage[]): SessionStatusAuthority | null {
|
||||
const message = [...messages].reverse().find((item) => messageHasCompletedFinalResponse(item));
|
||||
const message = [...messages].reverse().find((item) => messageHasTerminalResponse(item));
|
||||
if (!message) return null;
|
||||
const status = canonicalTerminalMessageStatus(message.status);
|
||||
return {
|
||||
sessionId,
|
||||
status: "completed",
|
||||
status,
|
||||
updatedAt: textValue(message.updatedAt) ?? textValue(message.finishedAt) ?? textValue(message.lastEventAt),
|
||||
lastTraceId: textValue(message.traceId) ?? textValue(message.runnerTrace?.traceId),
|
||||
projection: projectionFromMessageRecord(message),
|
||||
@@ -407,7 +414,7 @@ function sessionStatusAuthorityFromMessages(sessionId: string, messages: ChatMes
|
||||
}
|
||||
|
||||
function mergeSessionStatusAuthority(existing: SessionStatusAuthority | undefined, incoming: SessionStatusAuthority): SessionStatusAuthority {
|
||||
if (existing?.status === "completed" && incoming.status !== "completed" && isSameTraceAuthority(existing, incoming) && isRunningSessionStatus(incoming.status)) return existing;
|
||||
if (isTerminalSessionStatusAuthority(existing) && !isTerminalSessionStatusAuthority(incoming) && isSameTraceAuthority(existing, incoming) && isRunningSessionStatus(incoming.status)) return existing;
|
||||
return {
|
||||
...(existing ?? {}),
|
||||
...incoming,
|
||||
@@ -416,6 +423,11 @@ function mergeSessionStatusAuthority(existing: SessionStatusAuthority | undefine
|
||||
};
|
||||
}
|
||||
|
||||
function isTerminalSessionStatusAuthority(session: SessionStatusAuthority | undefined): boolean {
|
||||
if (!session) return false;
|
||||
return isTerminalMessageStatus(session.status) && !isRunningSessionStatus(session.status);
|
||||
}
|
||||
|
||||
function mergeTurnStatusAuthority(existing: TurnStatusAuthority | undefined, incoming: TurnStatusAuthority): TurnStatusAuthority {
|
||||
if (existing && isTerminalTurnStatusAuthority(existing) && !isTerminalTurnStatusAuthority(incoming) && isRunningSessionStatus(incoming.status)) return existing;
|
||||
return {
|
||||
@@ -431,10 +443,7 @@ function isTerminalTurnStatusAuthority(turn: TurnStatusAuthority): boolean {
|
||||
}
|
||||
|
||||
function messageIsTerminalForState(message: ChatMessage): boolean {
|
||||
const status = normalizedMessageStatus(message.status);
|
||||
if (!isTerminalMessageStatus(status)) return false;
|
||||
if (status !== "completed") return true;
|
||||
return messageHasCompletedFinalResponse(message);
|
||||
return messageHasTerminalResponse(message);
|
||||
}
|
||||
|
||||
function isSameTraceAuthority(left: SessionStatusAuthority, right: SessionStatusAuthority): boolean {
|
||||
@@ -447,12 +456,18 @@ function isRunningSessionStatus(value: unknown): boolean {
|
||||
return ["", "pending", "running", "accepted", "queued", "dispatching", "streaming", "processing", "retrying", "busy", "creating"].includes(normalizedMessageStatus(value));
|
||||
}
|
||||
|
||||
function messageHasCompletedFinalResponse(message: ChatMessage): boolean {
|
||||
function messageHasTerminalResponse(message: ChatMessage): boolean {
|
||||
if (message.role !== "agent") return false;
|
||||
if (normalizedMessageStatus(message.status) !== "completed") return false;
|
||||
const status = normalizedMessageStatus(message.status);
|
||||
if (!isTerminalMessageStatus(status)) return false;
|
||||
return Boolean(messageFinalResponseText(message));
|
||||
}
|
||||
|
||||
function canonicalTerminalMessageStatus(value: unknown): string {
|
||||
const status = normalizedMessageStatus(value);
|
||||
return status === "cancelled" ? "canceled" : status;
|
||||
}
|
||||
|
||||
function messageFinalResponseText(message: ChatMessage): string | null {
|
||||
return textValue(message.text) ?? textValue(message.content) ?? nestedTextValue((message as Record<string, unknown>).finalResponse);
|
||||
}
|
||||
|
||||
@@ -18,6 +18,36 @@ test("composer treats completed message without final response as unsealed runni
|
||||
assert.equal(composer.targetTraceId, traceId);
|
||||
});
|
||||
|
||||
test("composer treats failed message without terminal body as unsealed running turn", () => {
|
||||
const traceId = "trc_session_failed_without_body";
|
||||
const sessionId = "ses_session_failed_without_body";
|
||||
const message = { id: "msg_session_failed_without_body", role: "agent", status: "failed", traceId, sessionId } as any;
|
||||
const composer = resolveComposerState({
|
||||
messages: [message],
|
||||
activeSessionId: sessionId,
|
||||
chatPending: true,
|
||||
currentRequest: { traceId, sessionId, status: "running" },
|
||||
turnStatusAuthority: { [traceId]: { traceId, sessionId, status: "failed", terminal: true, running: false } as any }
|
||||
});
|
||||
assert.equal(composer.submitMode, "steer");
|
||||
assert.equal(composer.targetTraceId, traceId);
|
||||
});
|
||||
|
||||
test("composer treats failed message with terminal body as sealed terminal turn", () => {
|
||||
const traceId = "trc_session_failed_with_body";
|
||||
const sessionId = "ses_session_failed_with_body";
|
||||
const message = { id: "msg_session_failed_with_body", role: "agent", status: "failed", traceId, sessionId, text: "Workbench terminal failed: provider-stream-disconnected" } as any;
|
||||
const composer = resolveComposerState({
|
||||
messages: [message],
|
||||
activeSessionId: sessionId,
|
||||
chatPending: true,
|
||||
currentRequest: { traceId, sessionId, status: "running" },
|
||||
turnStatusAuthority: { [traceId]: { traceId, sessionId, status: "failed", terminal: true, running: false } as any }
|
||||
});
|
||||
assert.equal(composer.submitMode, "turn");
|
||||
assert.equal(composer.targetTraceId, null);
|
||||
});
|
||||
|
||||
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";
|
||||
@@ -27,3 +57,13 @@ test("cancel target remains available for completed message until final response
|
||||
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);
|
||||
});
|
||||
|
||||
test("cancel target remains available for failed message until terminal body exists", () => {
|
||||
const traceId = "trc_cancel_failed_without_body";
|
||||
const sessionId = "ses_cancel_failed_without_body";
|
||||
const unsealed = { id: "msg_cancel_failed_without_body", role: "agent", status: "failed", traceId, sessionId } as any;
|
||||
const sealed = { ...unsealed, id: "msg_cancel_failed_with_body", text: "Workbench terminal failed: provider-stream-disconnected" } 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);
|
||||
});
|
||||
|
||||
@@ -84,8 +84,8 @@ export function resolveComposerState(input: { messages: ChatMessage[]; sessions?
|
||||
const activeByMessage = latestMessage?.role === "agent" && isActiveStatus(latestMessage.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 terminalTurn = turn?.terminal === true || isTerminalStatus(turnStatus);
|
||||
const terminalByTurn = terminalTurn && messageHasTerminalResponse(latestMessage);
|
||||
const activeByStatus = !terminalByMessage && !terminalByTurn && (turn?.running === true || isActiveStatus(turnStatus) || activeByRequest || activeByMessage);
|
||||
const terminal = terminalByMessage || terminalByTurn;
|
||||
const active = activeSession(input.sessions ?? [], sessionId);
|
||||
@@ -285,14 +285,11 @@ function latestAgentMessage(messages: ChatMessage[] | undefined): ChatMessage |
|
||||
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);
|
||||
return isTerminalStatus(status) && messageHasTerminalResponse(message);
|
||||
}
|
||||
|
||||
function messageHasCompletedFinalResponse(message: ChatMessage | null | undefined): boolean {
|
||||
function messageHasTerminalResponse(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));
|
||||
}
|
||||
|
||||
@@ -31,7 +31,7 @@ import {
|
||||
firstFiniteNumber,
|
||||
isTerminalMessageStatus,
|
||||
isTraceActiveStatus,
|
||||
messageHasCompletedFinalResponse,
|
||||
messageHasTerminalResponse,
|
||||
messageNeedsTerminalDiagnostics,
|
||||
messageNeedsTraceHydration,
|
||||
messageStatusPatchForTerminalMerge,
|
||||
@@ -57,7 +57,7 @@ import {
|
||||
terminalMessageTimingPatchForNormalize,
|
||||
turnResultIsTerminalForMerge,
|
||||
turnResultStatusForMerge,
|
||||
traceHasCompletedFinalResponse,
|
||||
traceHasTerminalResponse,
|
||||
traceHasEvents,
|
||||
traceResultHasTerminalEvidence,
|
||||
traceSnapshotError
|
||||
@@ -560,7 +560,7 @@ export const useWorkbenchStore = defineStore("workbench", () => {
|
||||
);
|
||||
if (!response.ok || !response.data) {
|
||||
if (shouldSuppressTransientWorkbenchReadFailure(response)) return;
|
||||
if (traceHasCompletedFinalResponse(traceId, messages.value)) return;
|
||||
if (traceHasTerminalResponse(traceId, messages.value)) return;
|
||||
applyProjectionDiagnostic(traceId, projectionDiagnosticFromApiFailure(response, { code: "message_projection_refresh_failed", message: response.error ?? "消息投影刷新失败,主消息正文保持上一份 canonical projection。", health: response.status === 0 ? "unavailable" : "degraded" }));
|
||||
return;
|
||||
}
|
||||
@@ -634,15 +634,11 @@ export const useWorkbenchStore = defineStore("workbench", () => {
|
||||
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;
|
||||
return message ? messageHasTerminalResponse(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;
|
||||
return messageHasTerminalResponse(message);
|
||||
}
|
||||
|
||||
async function refreshTurnStatusByTraceId(traceId: string | null | undefined, options: { force?: boolean } = {}): Promise<void> {
|
||||
@@ -655,7 +651,7 @@ export const useWorkbenchStore = defineStore("workbench", () => {
|
||||
return;
|
||||
}
|
||||
if (shouldSuppressTransientWorkbenchReadFailure(response)) return;
|
||||
if (traceHasCompletedFinalResponse(id, messages.value)) return;
|
||||
if (traceHasTerminalResponse(id, messages.value)) return;
|
||||
applyProjectionDiagnostic(id, projectionDiagnosticFromApiFailure(response, { code: "turn_status_poll_failed", message: response.error ?? "状态更新失败,当前 turn 状态暂不可见。", health: response.status === 0 ? "unavailable" : "degraded" }));
|
||||
}
|
||||
|
||||
@@ -844,7 +840,7 @@ export const useWorkbenchStore = defineStore("workbench", () => {
|
||||
const result = await fetchTraceHydrationPage(traceId, afterProjectedSeq, { force: options.force });
|
||||
if (!result.ok || !result.data) {
|
||||
if (shouldSuppressTransientWorkbenchReadFailure(result)) return;
|
||||
if (messageHasCompletedFinalResponse(message)) return;
|
||||
if (messageHasTerminalResponse(message)) return;
|
||||
applyProjectionDiagnostic(traceId, projectionDiagnosticFromApiFailure(result, { code: "trace_hydration_failed", message: result.error ?? "Trace 更新超时,运行记录暂不可见。", health: result.status === 0 ? "unavailable" : "degraded" }));
|
||||
return;
|
||||
}
|
||||
@@ -1118,12 +1114,7 @@ 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 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);
|
||||
const terminal = messageHasTerminalResponse(message);
|
||||
if (terminal) {
|
||||
clearActiveTraceRestGapFill(id);
|
||||
return;
|
||||
@@ -1483,9 +1474,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, messageText((message as Record<string, unknown>).content), finalResponseText((message as Record<string, unknown>).finalResponse)));
|
||||
return messageHasTerminalResponse(message);
|
||||
}
|
||||
|
||||
function failTrace(traceId: string, message: string): void {
|
||||
@@ -1848,10 +1837,10 @@ async function sealRestoredActiveTurnMessages(source: ChatMessage[]): Promise<Ch
|
||||
|
||||
function messageNeedsRestoredTurnSeal(message: ChatMessage): boolean {
|
||||
if (message.role !== "agent") return false;
|
||||
if (isTerminalMessageStatus(message.status)) return false;
|
||||
if (messageHasTerminalResponse(message)) return false;
|
||||
const traceId = firstNonEmptyString(message.traceId, message.runnerTrace?.traceId);
|
||||
if (!traceId) return false;
|
||||
return isTraceActiveStatus(message.status) || message.traceAutoLifecycle === "running";
|
||||
return isTraceActiveStatus(message.status) || isTerminalMessageStatus(message.status) || message.traceAutoLifecycle === "running";
|
||||
}
|
||||
|
||||
function providerThreadIdForRequest(threadId: string | null | undefined): string | null {
|
||||
@@ -1896,8 +1885,7 @@ function normalizeChatMessage(message: ChatMessage): ChatMessage {
|
||||
function activeTraceIdFromMessages(messages: ChatMessage[], turnStatusAuthority: Record<string, TurnStatusAuthority>): string | null {
|
||||
for (const message of [...messages].reverse()) {
|
||||
if (message.role !== "agent") continue;
|
||||
if (messageHasCompletedFinalResponse(message)) continue;
|
||||
if (isTerminalMessageStatus(message.status)) continue;
|
||||
if (messageHasTerminalResponse(message)) continue;
|
||||
const traceId = firstNonEmptyString(message.traceId, message.runnerTrace?.traceId);
|
||||
if (!traceId) continue;
|
||||
const turn = turnStatusAuthority[traceId];
|
||||
|
||||
Reference in New Issue
Block a user