From f42d3fb000bf7fb7af6e6f1bbb714f56caea6c19 Mon Sep 17 00:00:00 2001 From: Josh Lehman Date: Tue, 8 Sep 2026 16:43:28 -0700 Subject: [PATCH] fix(ui): prevent duplicate final replies after history hydration (#142422) * fix(ui): reconcile hydrated terminal finals Refs openclaw-827 * perf(ui): skip settled projection matches * perf(ui): keep terminal reconciliation bounded * fix(ui): preserve partial terminal recovery * fix(ui): match enriched terminal identity * fix(ui): require durable terminal evidence * fix(gateway-client): reconcile unmarked cli finals * fix(gateway-client): preserve unmarked final reconciliation * fix(gateway-client): retain inferred terminal recovery * fix(gateway-client): retire removed terminal recovery Preserve tentative live ordering only while its matched history row remains. Retire recovery on authoritative removal, filtering, or explicit terminal confirmation. Refs: openclaw-116 --- .../src/session-projection-final-identity.ts | 222 +++++++++++++ .../src/session-projection-message-content.ts | 6 + .../src/session-projection-run-retention.ts | 20 ++ .../session-projection-run-terminal.test.ts | 42 ++- ...projection-terminal-reconciliation.test.ts | 292 ++++++++++++++++++ .../gateway-client/src/session-projection.ts | 151 +++++---- ui/src/app/vite-config.node.test.ts | 10 +- ui/src/e2e/chat-live-final-order.e2e.test.ts | 108 +++++++ ui/src/pages/chat/chat-state.test.ts | 86 ++++++ ui/vite.config.ts | 1 + 10 files changed, 853 insertions(+), 85 deletions(-) create mode 100644 packages/gateway-client/src/session-projection-final-identity.ts create mode 100644 packages/gateway-client/src/session-projection-run-retention.ts create mode 100644 packages/gateway-client/src/session-projection-terminal-reconciliation.test.ts diff --git a/packages/gateway-client/src/session-projection-final-identity.ts b/packages/gateway-client/src/session-projection-final-identity.ts new file mode 100644 index 000000000000..399438ecde1e --- /dev/null +++ b/packages/gateway-client/src/session-projection-final-identity.ts @@ -0,0 +1,222 @@ +/** Terminal identity rules used to reconcile live and durable assistant projections. */ + +import { asNullableRecord as readRecord } from "@openclaw/normalization-core/record-coerce"; +import { stableStringify } from "@openclaw/normalization-core/stable-stringify"; +import { + hasDisplayableSessionMessage, + readSessionMessageDisplayContent, +} from "./session-projection-message-content.js"; +import { + readSessionMessageIdentity, + type SessionMessageIdentity, +} from "./session-projection-message-identity.js"; + +type TerminalProjectionEntry = { + message: unknown; + identity: SessionMessageIdentity | null; + live: boolean; +}; + +type TerminalProjectionRun = { + message?: unknown; + status: string; + acceptedFinalMessageIdentities?: readonly string[]; +}; + +function readPersistedFinalIdentity(message: unknown): string | null { + const identity = readSessionMessageIdentity(message); + if (identity?.externalSource) { + return `import:${identity.role}:${identity.externalSource}`; + } + if (identity?.id && !identity.isImported) { + return `id:${identity.role}:${identity.id}`; + } + if (identity?.sequence !== null && identity?.sequence !== undefined) { + return `seq:${identity.role}:${identity.sequence}`; + } + return null; +} + +function hasCompatiblePersistedFinalIdentity(currentMessage: unknown, incomingMessage: unknown) { + const current = readSessionMessageIdentity(currentMessage); + const incoming = readSessionMessageIdentity(incomingMessage); + if (!current || !incoming || current.role !== incoming.role) { + return false; + } + if (current.isImported || incoming.isImported) { + if (!current.isImported || !incoming.isImported) { + return false; + } + if (current.externalSource && incoming.externalSource) { + return current.externalSource === incoming.externalSource; + } + return ( + current.sequence !== null && + incoming.sequence !== null && + current.sequence === incoming.sequence + ); + } + if (current.id && incoming.id) { + return current.id === incoming.id; + } + return ( + current.sequence !== null && + incoming.sequence !== null && + current.sequence === incoming.sequence + ); +} + +function readFinalContentIdentity(message: unknown): string | null { + const display = readSessionMessageDisplayContent(message); + if (!display.text && !display.hasNonText) { + return null; + } + const identity = readSessionMessageIdentity(message); + const record = readRecord(message); + const metadata = readRecord(record?.["__openclaw"]); + try { + return `content:${stableStringify([ + identity?.role ?? "assistant", + display.text, + display.hasNonText ? (record?.content ?? null) : null, + metadata?.media ?? null, + identity?.isImported + ? [ + metadata?.importedFrom ?? null, + metadata?.cliSessionId ?? null, + metadata?.externalId ?? null, + ] + : null, + ])}`; + } catch { + return null; + } +} + +function hasTerminalStopReason(message: unknown): boolean { + const stopReason = readRecord(message)?.stopReason; + return ( + stopReason === "stop" || + stopReason === "length" || + stopReason === "error" || + stopReason === "aborted" || + stopReason === "end_turn" + ); +} + +function hasCompletedRunSnapshotContext( + entry: TerminalProjectionEntry, + snapshot: readonly TerminalProjectionEntry[], + runId: string | null, +): boolean { + if (!runId || entry.identity?.runId !== runId) { + return false; + } + const entryIndex = snapshot.indexOf(entry); + if (entryIndex < 0) { + return false; + } + const hasEarlierUser = snapshot + .slice(0, entryIndex) + .some((candidate) => candidate.identity?.role === "user" && candidate.identity.runId === runId); + const hasLaterAssistant = snapshot + .slice(entryIndex + 1) + .some( + (candidate) => candidate.identity?.role === "assistant" && candidate.identity.runId === runId, + ); + return hasEarlierUser && !hasLaterAssistant; +} + +/** Read stable persisted identity first, falling back to canonical display content. */ +export function readSessionProjectionFinalMessageIdentity(message: unknown): string | null { + if (!hasDisplayableSessionMessage(message)) { + return null; + } + return readPersistedFinalIdentity(message) ?? readFinalContentIdentity(message); +} + +/** Check whether a displayable terminal may recover a prior empty terminal. */ +export function canRecoverSessionProjectionFinal( + currentMessage: unknown, + incomingMessage: unknown, +): boolean { + if (hasDisplayableSessionMessage(currentMessage)) { + return false; + } + const currentIdentity = readPersistedFinalIdentity(currentMessage); + return ( + currentIdentity === null || hasCompatiblePersistedFinalIdentity(currentMessage, incomingMessage) + ); +} + +/** Check whether a run has already accepted the same terminal reply. */ +export function hasSessionProjectionAcceptedFinal( + run: TerminalProjectionRun | undefined, + message: unknown, +): boolean { + const identity = readSessionProjectionFinalMessageIdentity(message); + return Boolean( + identity && + run && + (run.acceptedFinalMessageIdentities?.includes(identity) || + readSessionProjectionFinalMessageIdentity(run.message) === identity), + ); +} + +/** Match an unsequenced live terminal to exactly one durable same-run terminal row. */ +export function findUniqueSnapshotTerminalMatch( + current: TerminalProjectionEntry, + matches: readonly TerminalProjectionEntry[], + run: TerminalProjectionRun | undefined, + snapshot: readonly TerminalProjectionEntry[], +): { entry: TerminalProjectionEntry; inferred: boolean } | null { + if ( + !current.live || + current.identity?.role !== "assistant" || + current.identity.id || + current.identity.sequence !== null || + !run || + run.status === "streaming" + ) { + return null; + } + const terminalContent = readFinalContentIdentity(current.message); + if (!terminalContent || readFinalContentIdentity(run.message) !== terminalContent) { + return null; + } + const durableTerminalMatches = matches.filter((entry) => { + const metadata = readRecord(readRecord(entry.message)?.["__openclaw"]); + return ( + (metadata?.runTerminal === true || + (entry.identity?.runId === current.identity?.runId && + hasTerminalStopReason(entry.message)) || + hasCompletedRunSnapshotContext(entry, snapshot, current.identity?.runId ?? null)) && + readFinalContentIdentity(entry.message) === terminalContent + ); + }); + const entry = durableTerminalMatches.length === 1 ? durableTerminalMatches[0] : undefined; + if (!entry) { + return null; + } + const metadata = readRecord(readRecord(entry.message)?.["__openclaw"]); + return { + entry, + inferred: metadata?.runTerminal !== true && !hasTerminalStopReason(entry.message), + }; +} + +/** Check whether ordinary single-match promotion needs terminal-content verification. */ +export function isUnsequencedLiveTerminal( + current: TerminalProjectionEntry, + run: TerminalProjectionRun | undefined, +): boolean { + return Boolean( + current.live && + current.identity?.role === "assistant" && + !current.identity.id && + current.identity.sequence === null && + run && + run.status !== "streaming" && + readFinalContentIdentity(current.message) === readFinalContentIdentity(run.message), + ); +} diff --git a/packages/gateway-client/src/session-projection-message-content.ts b/packages/gateway-client/src/session-projection-message-content.ts index f1e765a44093..b75bb64744f0 100644 --- a/packages/gateway-client/src/session-projection-message-content.ts +++ b/packages/gateway-client/src/session-projection-message-content.ts @@ -35,6 +35,12 @@ export function readSessionMessageDisplayContent(message: unknown): { return { text: fallback ?? texts.join("\n"), hasNonText, usesFallbackText: fallback !== null }; } +/** Check whether a projected message has text or another displayable block. */ +export function hasDisplayableSessionMessage(message: unknown): boolean { + const { text, hasNonText } = readSessionMessageDisplayContent(message); + return Boolean(text) || hasNonText; +} + function normalizeChatErrorComparisonText(text: string): string { return text .trim() diff --git a/packages/gateway-client/src/session-projection-run-retention.ts b/packages/gateway-client/src/session-projection-run-retention.ts new file mode 100644 index 000000000000..8bd867557933 --- /dev/null +++ b/packages/gateway-client/src/session-projection-run-retention.ts @@ -0,0 +1,20 @@ +/** Bounded retention for completed and active session projection runs. */ + +const MAX_TRACKED_SESSION_RUNS = 200; +const RETAINED_SESSION_RUNS = 150; + +/** Retain every active run and the newest completed runs within the projection bound. */ +export function retainSessionProjectionRuns( + runs: Readonly>, +): Readonly> { + const entries = Object.entries(runs); + if (entries.length <= MAX_TRACKED_SESSION_RUNS) { + return runs; + } + const active = entries.filter(([, run]) => run.status === "streaming"); + const terminal = entries.filter(([, run]) => run.status !== "streaming"); + const terminalLimit = Math.max(0, RETAINED_SESSION_RUNS - active.length); + const retainedTerminal = terminalLimit > 0 ? terminal.slice(-terminalLimit) : []; + // Live streams are never expendable; completed runs are retained by completion order. + return Object.fromEntries([...active, ...retainedTerminal]); +} diff --git a/packages/gateway-client/src/session-projection-run-terminal.test.ts b/packages/gateway-client/src/session-projection-run-terminal.test.ts index 3a34254c9ae6..95859016d5ce 100644 --- a/packages/gateway-client/src/session-projection-run-terminal.test.ts +++ b/packages/gateway-client/src/session-projection-run-terminal.test.ts @@ -132,9 +132,12 @@ describe("session run terminal bookkeeping", () => { ).toBe(failed); }); - it("upgrades an empty completed final exactly once without reopening the run", () => { - const emptyMessage = createMessage("assistant", ""); - const deliveredMessage = createMessage("assistant", "eventual final"); + it("upgrades an identified empty completed final exactly once without reopening the run", () => { + const emptyMessage = createMessage("assistant", "", { id: "assistant-final", seq: 7 }); + const deliveredMessage = createMessage("assistant", "eventual final", { + id: "assistant-final", + seq: 7, + }); let state = reduceSessionProjection(createSessionProjection(primaryScope), { type: "runTerminal", runId: "run-1", @@ -142,6 +145,17 @@ describe("session run terminal bookkeeping", () => { message: emptyMessage, }); expect(hasSessionProjectionAcceptedFinal(state.runs["run-1"], emptyMessage)).toBe(false); + const mismatchedMessage = createMessage("assistant", "wrong final", { + id: "different-assistant-final", + seq: 7, + }); + state = reduceSessionProjection(state, { + type: "runTerminal", + runId: "run-1", + status: "completed", + message: mismatchedMessage, + }); + expect(state.runs["run-1"]?.message).toBe(emptyMessage); state = reduceSessionProjection(state, { type: "runTerminal", runId: "run-1", @@ -168,6 +182,28 @@ describe("session run terminal bookkeeping", () => { expect(reduceSessionProjection(acceptedLaterFinal, laterEvent)).toBe(acceptedLaterFinal); }); + it("recovers a sequence-identified empty final when persisted metadata adds an ID", () => { + const emptyMessage = createMessage("assistant", "", { seq: 7 }); + const deliveredMessage = createMessage("assistant", "eventual final", { + id: "assistant-final", + seq: 7, + }); + let state = reduceSessionProjection(createSessionProjection(primaryScope), { + type: "runTerminal", + runId: "run-1", + status: "completed", + message: emptyMessage, + }); + state = reduceSessionProjection(state, { + type: "runTerminal", + runId: "run-1", + status: "completed", + message: deliveredMessage, + }); + + expect(state.runs["run-1"]?.message).toBe(deliveredMessage); + }); + it("accepts distinct same-run persisted finals and ignores the later final's replay", () => { const first = createMessage("assistant", "first final", { id: "assistant-a", seq: 4 }); const second = createMessage("assistant", "second final", { id: "assistant-b", seq: 5 }); diff --git a/packages/gateway-client/src/session-projection-terminal-reconciliation.test.ts b/packages/gateway-client/src/session-projection-terminal-reconciliation.test.ts new file mode 100644 index 000000000000..e498370f462e --- /dev/null +++ b/packages/gateway-client/src/session-projection-terminal-reconciliation.test.ts @@ -0,0 +1,292 @@ +import { describe, expect, it } from "vitest"; +import { + createSessionProjection, + projectLiveSessionMessage, + reconcileSessionProjectionSnapshot, + reduceSessionProjection, + type SessionProjectionScope, +} from "./session-projection.js"; + +const scope: SessionProjectionScope = { + sessionKey: "agent:main:shared", + sessionId: "session-1", + agentId: "main", + lifecycleRevision: 1, + activeLeafEntryId: "leaf-1", +}; + +function createAssistantMessage(text: string, metadata?: Record) { + return { + role: "assistant", + content: [{ type: "text", text }], + ...(metadata ? { __openclaw: metadata } : {}), + }; +} + +describe("terminal snapshot reconciliation", () => { + it("promotes the actual terminal when history contains an earlier same-run tool boundary", () => { + const runId = "tool-heavy-run"; + const toolBoundary = { + role: "assistant", + content: [ + { type: "text", text: "Checking the repository." }, + { type: "toolCall", id: "read-1", name: "read", arguments: { path: "AGENTS.md" } }, + ], + __openclaw: { id: "assistant-tool-boundary", seq: 2, runId }, + }; + const synthetic = createAssistantMessage("The repair is complete."); + const persisted = { + role: "assistant", + content: [{ text: "The repair is complete.", type: "text" }], + __openclaw: { id: "assistant-final", seq: 4, runId, runTerminal: true }, + }; + let state = reduceSessionProjection(createSessionProjection(scope), { + type: "runTerminal", + runId, + status: "completed", + message: synthetic, + }); + state = projectLiveSessionMessage(state, structuredClone(synthetic), { runId }); + + expect( + reconcileSessionProjectionSnapshot(state, [toolBoundary, persisted], scope).messages, + ).toEqual([toolBoundary, persisted]); + }); + + it("promotes an unmarked same-run CLI terminal with a terminal stop reason", () => { + const runId = "cli-run"; + const synthetic = createAssistantMessage("The CLI repair is complete."); + const persisted = { + role: "assistant", + api: "cli", + content: [{ text: "The CLI repair is complete.", type: "text" }], + idempotencyKey: `cli-assistant:${runId}`, + stopReason: "stop", + __openclaw: { id: "assistant-final", seq: 4 }, + }; + let state = reduceSessionProjection(createSessionProjection(scope), { + type: "runTerminal", + runId, + status: "completed", + message: synthetic, + }); + state = projectLiveSessionMessage(state, structuredClone(synthetic), { runId }); + + expect(reconcileSessionProjectionSnapshot(state, [persisted], scope).messages).toEqual([ + persisted, + ]); + }); + + it("promotes a trailing unmarked terminal after the durable same-run user turn", () => { + const runId = "browser-run"; + const user = { + role: "user", + content: [{ text: "Please finish the repair.", type: "text" }], + __openclaw: { id: "user-prompt", idempotencyKey: `${runId}:user`, seq: 1 }, + }; + const synthetic = createAssistantMessage("The browser repair is complete."); + const persisted = createAssistantMessage("The browser repair is complete.", { + id: "assistant-final", + seq: 2, + runId, + }); + let state = reduceSessionProjection(createSessionProjection(scope), { + type: "runTerminal", + runId, + status: "completed", + message: synthetic, + }); + state = projectLiveSessionMessage(state, structuredClone(synthetic), { runId }); + + expect(reconcileSessionProjectionSnapshot(state, [user, persisted], scope).messages).toEqual([ + user, + persisted, + ]); + }); + + it("promotes equivalent string and block terminal content", () => { + const runId = "string-terminal-run"; + const user = { + role: "user", + content: [{ text: "Please finish the repair.", type: "text" }], + __openclaw: { id: "user-prompt", idempotencyKey: `${runId}:user`, seq: 1 }, + }; + const persisted = createAssistantMessage("The repair is complete.", { + id: "assistant-final", + seq: 2, + runId, + }); + let state = reduceSessionProjection(createSessionProjection(scope), { + type: "runTerminal", + runId, + status: "completed", + message: "The repair is complete.", + }); + state = projectLiveSessionMessage(state, "The repair is complete.", { runId }); + + expect(reconcileSessionProjectionSnapshot(state, [user, persisted], scope).messages).toEqual([ + user, + persisted, + ]); + }); + + it("restores an inferred terminal when later history reveals a tool boundary", () => { + const runId = "partial-history-run"; + const user = { + role: "user", + content: [{ text: "Please inspect the repository.", type: "text" }], + __openclaw: { id: "user-prompt", idempotencyKey: `${runId}:user`, seq: 1 }, + }; + const synthetic = createAssistantMessage("Still working."); + const earlier = createAssistantMessage("Still working.", { + id: "assistant-earlier", + seq: 2, + runId, + }); + const laterToolBoundary = { + role: "assistant", + content: [ + { type: "text", text: "Checking another file." }, + { type: "toolCall", id: "read-2", name: "read", arguments: { path: "src/index.ts" } }, + ], + __openclaw: { id: "assistant-tool-boundary", seq: 3, runId }, + }; + let state = reduceSessionProjection(createSessionProjection(scope), { + type: "runTerminal", + runId, + status: "completed", + message: synthetic, + }); + state = projectLiveSessionMessage(state, synthetic, { runId }); + state = reconcileSessionProjectionSnapshot(state, [user, earlier], scope); + expect(state.messages).toEqual([user, earlier]); + + expect( + reconcileSessionProjectionSnapshot(state, [user, earlier, laterToolBoundary], scope).messages, + ).toEqual([user, earlier, laterToolBoundary, synthetic]); + }); + + it.each(["inferred", "runTerminal", "stopReason"])( + "does not restore a removed or filtered %s terminal", + (evidence) => { + const runId = "retired-terminal-run"; + const user = { + role: "user", + content: "Finish the task.", + __openclaw: { id: "user", seq: 1, runId }, + }; + const synthetic = createAssistantMessage("The task is complete."); + const persisted = { + ...createAssistantMessage("The task is complete.", { + id: "final", + seq: 2, + runId, + ...(evidence === "runTerminal" ? { runTerminal: true } : {}), + }), + ...(evidence === "stopReason" ? { stopReason: "stop" } : {}), + }; + let state = reduceSessionProjection(createSessionProjection(scope), { + type: "runTerminal", + runId, + status: "completed", + message: synthetic, + }); + state = projectLiveSessionMessage(state, synthetic, { runId }); + state = reconcileSessionProjectionSnapshot(state, [user, persisted], scope); + expect(state.messages).toEqual([user, persisted]); + + const removed = reconcileSessionProjectionSnapshot(state, [user], scope); + expect(removed.messages).toEqual([user]); + expect(reconcileSessionProjectionSnapshot(removed, [user], scope).messages).toEqual([user]); + const filtered = reconcileSessionProjectionSnapshot(state, [user, persisted], scope, { + shouldIncludeMessage: (message) => message === user, + }); + expect(filtered.messages).toEqual([user]); + expect(reconcileSessionProjectionSnapshot(filtered, [user], scope).messages).toEqual([user]); + }, + ); + + it("retains an unsequenced terminal when matching content precedes a later tool boundary", () => { + const runId = "partial-history-run"; + const user = { + role: "user", + content: [{ text: "Please inspect the repository.", type: "text" }], + __openclaw: { id: "user-prompt", idempotencyKey: `${runId}:user`, seq: 1 }, + }; + const synthetic = createAssistantMessage("Still working."); + const earlier = createAssistantMessage("Still working.", { + id: "assistant-earlier", + seq: 2, + runId, + }); + const laterToolBoundary = { + role: "assistant", + content: [ + { type: "text", text: "Checking another file." }, + { type: "toolCall", id: "read-2", name: "read", arguments: { path: "src/index.ts" } }, + ], + __openclaw: { id: "assistant-tool-boundary", seq: 3, runId }, + }; + let state = reduceSessionProjection(createSessionProjection(scope), { + type: "runTerminal", + runId, + status: "completed", + message: synthetic, + }); + state = projectLiveSessionMessage(state, synthetic, { runId }); + + expect( + reconcileSessionProjectionSnapshot(state, [user, earlier, laterToolBoundary], scope).messages, + ).toEqual([user, earlier, laterToolBoundary, synthetic]); + }); + + it("retains an unsequenced terminal when partial history has one unmarked same-content row", () => { + const runId = "partial-tool-history-run"; + const synthetic = createAssistantMessage("The repair is complete."); + const earlier = createAssistantMessage("The repair is complete.", { + id: "assistant-earlier", + seq: 2, + runId, + }); + let state = reduceSessionProjection(createSessionProjection(scope), { + type: "runTerminal", + runId, + status: "completed", + message: synthetic, + }); + state = projectLiveSessionMessage(state, synthetic, { runId }); + + expect(reconcileSessionProjectionSnapshot(state, [earlier], scope).messages).toEqual([ + earlier, + synthetic, + ]); + }); + + it("retains an unsequenced terminal when multiple same-run rows have terminal content", () => { + const runId = "ambiguous-run"; + const synthetic = createAssistantMessage("The repair is complete."); + const first = createAssistantMessage("The repair is complete.", { + id: "assistant-first", + seq: 2, + runId, + }); + const second = createAssistantMessage("The repair is complete.", { + id: "assistant-second", + seq: 3, + runId, + }); + let state = reduceSessionProjection(createSessionProjection(scope), { + type: "runTerminal", + runId, + status: "completed", + message: synthetic, + }); + state = projectLiveSessionMessage(state, synthetic, { runId }); + + expect(reconcileSessionProjectionSnapshot(state, [first, second], scope).messages).toEqual([ + first, + second, + synthetic, + ]); + }); +}); diff --git a/packages/gateway-client/src/session-projection.ts b/packages/gateway-client/src/session-projection.ts index 2080db83e746..937ffd0269ff 100644 --- a/packages/gateway-client/src/session-projection.ts +++ b/packages/gateway-client/src/session-projection.ts @@ -2,8 +2,15 @@ import { asNullableRecord as readRecord } from "@openclaw/normalization-core/record-coerce"; import { + canRecoverSessionProjectionFinal, + hasSessionProjectionAcceptedFinal, + findUniqueSnapshotTerminalMatch, + isUnsequencedLiveTerminal, + readSessionProjectionFinalMessageIdentity, +} from "./session-projection-final-identity.js"; +import { + hasDisplayableSessionMessage, isSessionProjectionErrorMessage, - readSessionMessageDisplayContent, } from "./session-projection-message-content.js"; import { normalizeSessionProjectionRunId, @@ -13,6 +20,11 @@ import { type SessionMessageEnvelope, type SessionMessageIdentity, } from "./session-projection-message-identity.js"; +import { retainSessionProjectionRuns } from "./session-projection-run-retention.js"; +export { + hasSessionProjectionAcceptedFinal, + readSessionProjectionFinalMessageIdentity, +} from "./session-projection-final-identity.js"; export { reduceSessionProjectionRunEvent, type SessionProjectionGatewayRunEvent, @@ -57,6 +69,10 @@ export type SessionProjectionRun = { status: SessionProjectionRunStatus; message?: unknown; acceptedFinalMessageIdentities?: readonly string[]; + inferredSnapshotTerminal?: { + entry: SessionProjectionEntry; + matchedIdentity: SessionMessageIdentity; + }; stopReason?: string; errorKind?: string; errorMessage?: string; @@ -79,8 +95,6 @@ export type SessionProjectionState = { hasTransportGap: boolean; }; -const MAX_TRACKED_SESSION_RUNS = 200; -const RETAINED_SESSION_RUNS = 150; const MAX_ACCEPTED_FINAL_MESSAGES_PER_RUN = 32; const SESSION_PROJECTION_SCOPE_KEYS = [ "sessionKey", @@ -440,93 +454,66 @@ export function reconcileSessionProjectionSnapshot( return createSessionProjection(scope, visibleMessages); } let entries = createProjectionEntries(visibleMessages); + const runs: Record = { ...state.runs }; for (const current of state.entries) { if ( (!current.live && !current.pending) || - options.shouldIncludeMessage?.(current.message) === false || - entries.filter((entry) => entryMatches(entry, current, true)).length === 1 + options.shouldIncludeMessage?.(current.message) === false ) { continue; } - entries = insertEntry(entries, current, state.runs); + const matches = entries.filter((entry) => entryMatches(entry, current, true)); + const run = current.identity?.runId ? runs[current.identity.runId] : undefined; + const terminalMatch = findUniqueSnapshotTerminalMatch(current, matches, run, entries); + if ((matches.length === 1 && !isUnsequencedLiveTerminal(current, run)) || terminalMatch) { + if ( + terminalMatch?.inferred && + terminalMatch.entry.identity && + current.identity?.runId && + run + ) { + // Tentative history matches retain their original live ordering until confirmed. + runs[current.identity.runId] = { + ...run, + inferredSnapshotTerminal: { + entry: current, + matchedIdentity: terminalMatch.entry.identity, + }, + }; + } + continue; + } + entries = insertEntry(entries, current, runs); + } + for (const [runId, run] of Object.entries(runs)) { + const inferred = run.inferredSnapshotTerminal; + if (!inferred) { + continue; + } + // Removal and visibility policy retire a candidate; neither contradicts its identity. + const candidateRemains = entries.some((entry) => + sameTranscriptIdentity(entry.identity, inferred.matchedIdentity), + ); + const visible = options.shouldIncludeMessage?.(inferred.entry.message) !== false; + const matches = entries.filter((entry) => entryMatches(entry, inferred.entry, true)); + const terminalMatch = findUniqueSnapshotTerminalMatch(inferred.entry, matches, run, entries); + if (candidateRemains && visible && terminalMatch?.inferred) { + continue; + } + if (candidateRemains && visible && !terminalMatch) { + entries = insertEntry(entries, inferred.entry, runs); + } + const { inferredSnapshotTerminal: _inferred, ...settledRun } = run; + runs[runId] = settledRun; } return { ...withEntries(state, entries), + runs, scope: { ...state.scope, ...scope }, hasTransportGap: false, }; } -function hasDisplayableSessionMessage(message: unknown): boolean { - const { text, hasNonText } = readSessionMessageDisplayContent(message); - return Boolean(text) || hasNonText; -} - -export function readSessionProjectionFinalMessageIdentity(message: unknown): string | null { - const display = readSessionMessageDisplayContent(message); - if (!display.text && !display.hasNonText) { - return null; - } - const identity = readSessionMessageIdentity(message); - if (identity?.externalSource) { - return `import:${identity.role}:${identity.externalSource}`; - } - if (identity?.id && !identity.isImported) { - return `id:${identity.role}:${identity.id}`; - } - if (identity?.sequence !== null && identity?.sequence !== undefined) { - return `seq:${identity.role}:${identity.sequence}`; - } - const record = readRecord(message); - const metadata = readRecord(record?.["__openclaw"]); - try { - return `content:${JSON.stringify([ - identity?.role ?? "assistant", - typeof message === "string" ? message : (record?.content ?? null), - metadata?.media ?? null, - identity?.isImported - ? [ - metadata?.importedFrom ?? null, - metadata?.cliSessionId ?? null, - metadata?.externalId ?? null, - ] - : null, - ...(display.usesFallbackText ? [record?.text] : []), - ])}`; - } catch { - return null; - } -} - -/** Replayed finals are recognized against this run's bounded canonical terminal history. */ -export function hasSessionProjectionAcceptedFinal( - run: SessionProjectionRun | undefined, - message: unknown, -): boolean { - const identity = readSessionProjectionFinalMessageIdentity(message); - return Boolean( - identity && - run && - (run.acceptedFinalMessageIdentities?.includes(identity) || - readSessionProjectionFinalMessageIdentity(run.message) === identity), - ); -} - -function retainSessionProjectionRuns( - runs: Readonly>, -): Readonly> { - const entries = Object.entries(runs); - if (entries.length <= MAX_TRACKED_SESSION_RUNS) { - return runs; - } - const active = entries.filter(([, run]) => run.status === "streaming"); - const terminal = entries.filter(([, run]) => run.status !== "streaming"); - const terminalLimit = Math.max(0, RETAINED_SESSION_RUNS - active.length); - const retainedTerminal = terminalLimit > 0 ? terminal.slice(-terminalLimit) : []; - // Live streams are never expendable; completed runs are retained by completion order. - return Object.fromEntries([...active, ...retainedTerminal]); -} - function updateRun( state: SessionProjectionState, incoming: SessionProjectionRun, @@ -554,16 +541,18 @@ function updateRun( if (current && current.status !== "streaming" && !resumesErrorProjection) { const incomingFinalIdentity = readSessionProjectionFinalMessageIdentity(incoming.message); const incomingIsFinal = incoming.status === "completed" || incoming.status === "yielded"; - const canRecoverFinal = - !hasDisplayableSessionMessage(current.message) || - (current.acceptedFinalMessageIdentities?.length ?? 0) > 0; + const currentHasDisplayableMessage = hasDisplayableSessionMessage(current.message); + const canAcceptFinal = currentHasDisplayableMessage + ? current.status === incoming.status || + (current.acceptedFinalMessageIdentities?.length ?? 0) > 0 + : canRecoverSessionProjectionFinal(current.message, incoming.message); const acceptFinal = incomingIsFinal && - (current.status === incoming.status || canRecoverFinal) && + canAcceptFinal && incomingFinalIdentity !== null && !hasSessionProjectionAcceptedFinal(current, incoming.message); // Distinct valid finals are remembered; the first delivered reply remains immutable. - const recoverMessage = acceptFinal && !hasDisplayableSessionMessage(current.message); + const recoverMessage = acceptFinal && !currentHasDisplayableMessage; const recoverError = readNonemptyString(current.errorMessage) === null && incomingErrorMessage !== null; const updateTerminalSequence = diff --git a/ui/src/app/vite-config.node.test.ts b/ui/src/app/vite-config.node.test.ts index 7c8bd2084812..fee8ed49eed3 100644 --- a/ui/src/app/vite-config.node.test.ts +++ b/ui/src/app/vite-config.node.test.ts @@ -400,6 +400,9 @@ describe("Control UI Vite config", () => { const resultAliasIndex = aliases.findIndex( (alias) => alias.find === "@openclaw/normalization-core/result", ); + const stableStringifyAliasIndex = aliases.findIndex( + (alias) => alias.find === "@openclaw/normalization-core/stable-stringify", + ); const rootAliasIndex = aliases.findIndex( (alias) => alias.find === "@openclaw/normalization-core", ); @@ -407,8 +410,13 @@ describe("Control UI Vite config", () => { find: "@openclaw/normalization-core/result", replacement: path.join(repoRoot, "packages/normalization-core/src/result.ts"), }); + expect(aliases[stableStringifyAliasIndex]).toEqual({ + find: "@openclaw/normalization-core/stable-stringify", + replacement: path.join(repoRoot, "packages/normalization-core/src/stable-stringify.ts"), + }); expect(resultAliasIndex).toBeGreaterThanOrEqual(0); - expect(rootAliasIndex).toBeGreaterThan(resultAliasIndex); + expect(stableStringifyAliasIndex).toBeGreaterThanOrEqual(0); + expect(rootAliasIndex).toBeGreaterThan(stableStringifyAliasIndex); }); it("uses Node package resolution for external packages inherited by worktrees", () => { diff --git a/ui/src/e2e/chat-live-final-order.e2e.test.ts b/ui/src/e2e/chat-live-final-order.e2e.test.ts index 7ae5c8dd11d3..bb0b07c42bdd 100644 --- a/ui/src/e2e/chat-live-final-order.e2e.test.ts +++ b/ui/src/e2e/chat-live-final-order.e2e.test.ts @@ -258,6 +258,114 @@ suite.define(() => { } }); + it("hydrates one terminal when same-run tool history overlaps final persistence", async () => { + const context = await suite.newBrowserContext(createControlUiE2eContextOptions()); + const page = await context.newPage(); + const runId = "tool-heavy-run"; + const promptText = "Inspect the repository."; + const boundaryText = "Checking the repository."; + const finalText = "The repair is complete."; + const prompt = { + role: "user", + content: [{ type: "text", text: promptText }], + __openclaw: { id: "prompt", idempotencyKey: `${runId}:user`, seq: 1 }, + timestamp: Date.now(), + }; + const toolBoundary = { + role: "assistant", + content: [ + { type: "text", text: boundaryText }, + { type: "toolCall", id: "read-1", name: "read", arguments: { path: "AGENTS.md" } }, + ], + __openclaw: { id: "assistant-tool-boundary", runId, seq: 2 }, + timestamp: Date.now(), + }; + const persistedFinal = { + role: "assistant", + content: [{ type: "text", text: finalText }], + __openclaw: { id: "assistant-final", runId, runTerminal: true, seq: 4 }, + timestamp: Date.now(), + }; + + try { + const gateway = await installMockGateway(page, { + deferredMethods: ["chat.history"], + methodResponses: { + "chat.history": { + messages: [], + sessionId: "session:agent:main:main", + sessionInfo: { + activeRunIds: [], + hasActiveRun: false, + key: "main", + kind: "direct", + status: "done", + updatedAt: Date.now(), + }, + }, + }, + }); + await page.goto(`${suite.server.baseUrl}chat`); + await gateway.waitForRequest("chat.startup"); + await gateway.emitChatFinal({ runId, text: finalText }); + await page.locator(".chat-thread-inner").getByText(finalText, { exact: true }).waitFor(); + + await gateway.setHistoryMessages([prompt, toolBoundary, persistedFinal]); + const historyBefore = (await gateway.getRequests("chat.history")).length; + await gateway.emitGatewayEvent("session.message", { + activeRunIds: [], + hasActiveRun: false, + message: prompt, + messageId: "prompt", + messageSeq: 1, + session: { + activeRunIds: [], + hasActiveRun: false, + key: "main", + kind: "direct", + status: "done", + updatedAt: Date.now(), + }, + sessionKey: "main", + }); + await gateway.waitForRequest("chat.history", { after: historyBefore }); + await gateway.resolveDeferred("chat.history", { + messages: [prompt, toolBoundary, persistedFinal], + sessionId: "session:agent:main:main", + sessionInfo: { + activeRunIds: [], + hasActiveRun: false, + key: "main", + kind: "direct", + lastRunId: runId, + status: "done", + updatedAt: Date.now(), + }, + }); + + const thread = page.locator(".chat-thread-inner"); + await thread.getByText(boundaryText, { exact: true }).waitFor(); + await expect + .poll(() => + thread.locator(".chat-group.assistant .chat-bubble", { hasText: finalText }).count(), + ) + .toBe(1); + await expect + .poll(() => + thread.evaluate( + (element, texts) => { + const rows = Array.from(element.querySelectorAll(".chat-bubble")); + return texts.map((text) => rows.findIndex((row) => row.textContent?.includes(text))); + }, + [promptText, boundaryText, finalText], + ), + ) + .toEqual([0, 1, 2]); + } finally { + await suite.closeBrowserContext(context); + } + }); + it("keeps durable turns ordered when a live final arrives before transcript events", async () => { const context = await suite.newBrowserContext(createControlUiE2eContextOptions()); const page = await context.newPage(); diff --git a/ui/src/pages/chat/chat-state.test.ts b/ui/src/pages/chat/chat-state.test.ts index d8ff35bf3191..0dca76bceedd 100644 --- a/ui/src/pages/chat/chat-state.test.ts +++ b/ui/src/pages/chat/chat-state.test.ts @@ -1174,6 +1174,92 @@ describe("canonical session message recovery", () => { ]); }); + it("hydrates one final when an earlier same-run tool boundary overlaps terminal persistence", async () => { + const runId = "tool-heavy-run"; + const prompt = { + role: "user", + content: [{ type: "text", text: "Inspect the repository." }], + __openclaw: { id: "prompt", idempotencyKey: `${runId}:user`, seq: 1 }, + }; + const toolBoundary = { + role: "assistant", + content: [ + { type: "text", text: "Checking the repository." }, + { type: "toolCall", id: "read-1", name: "read", arguments: { path: "AGENTS.md" } }, + ], + __openclaw: { id: "assistant-tool-boundary", seq: 2, runId }, + }; + const persistedFinal = { + role: "assistant", + content: [{ type: "text", text: "The repair is complete." }], + __openclaw: { id: "assistant-final", seq: 4, runId, runTerminal: true }, + }; + const request = vi.fn().mockResolvedValue({ + messages: [prompt, toolBoundary, persistedFinal], + sessionId: "selected-session", + sessionInfo: { + key: "agent:main:main", + kind: "direct", + updatedAt: 1, + hasActiveRun: false, + activeRunIds: [], + lastRunId: runId, + status: "done", + }, + }); + const { state } = createSessionEventState({ + chatMessages: [prompt], + chatHistoryPagination: { hasMore: false }, + chatRunId: runId, + chatStream: null, + chatStreamSegments: [], + chatToolMessages: [], + client: { request } as unknown as GatewayBrowserClient, + }); + + handlePageGatewayEvent(state, { + type: "event", + event: "chat", + payload: { + sessionKey: state.sessionKey, + runId, + state: "final", + message: { + role: "assistant", + content: [{ type: "text", text: "The repair is complete." }], + }, + }, + }); + handlePageGatewayEvent(state, { + type: "event", + event: "session.message", + payload: { + sessionKey: state.sessionKey, + message: prompt, + messageId: "prompt", + messageSeq: 1, + hasActiveRun: false, + activeRunIds: [], + session: { + key: state.sessionKey, + kind: "direct", + status: "done", + hasActiveRun: false, + activeRunIds: [], + updatedAt: 1, + }, + }, + }); + await loadChatHistory(state); + + expect(state.chatMessages).toEqual([prompt, toolBoundary, persistedFinal]); + expect(renderedTranscript(state)).toEqual([ + { role: "user", text: "Inspect the repository." }, + { role: "tool", text: "Checking the repository." }, + { role: "assistant", text: "The repair is complete." }, + ]); + }); + it("keeps the owned local prompt before an early durable reply after placement abandonment", () => { const runId = "local-placement-run-2"; const promptText = "Resume locally 2."; diff --git a/ui/vite.config.ts b/ui/vite.config.ts index 9bc0790e04d3..d8edce6f3bea 100644 --- a/ui/vite.config.ts +++ b/ui/vite.config.ts @@ -323,6 +323,7 @@ export function resolveSourcePackageAliasesForVite(): ControlUiViteAlias[] { sourcePackageAlias("normalization-core", "phone-presentation"), sourcePackageAlias("normalization-core", "record-coerce"), sourcePackageAlias("normalization-core", "result"), + sourcePackageAlias("normalization-core", "stable-stringify"), sourcePackageAlias("normalization-core", "string-coerce"), sourcePackageAlias("normalization-core", "string-normalization"), sourcePackageAlias("normalization-core", "utf16-slice"),