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
This commit is contained in:
Josh Lehman 2026-09-08 16:43:28 -07:00 • committed by GitHub
parent 7f7aeb6990
commit f42d3fb000
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
10 changed files with 853 additions and 85 deletions

View file

@ -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),
);
}

View file

@ -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()

View file

@ -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<T extends { status: string }>(
runs: Readonly<Record<string, T>>,
): Readonly<Record<string, T>> {
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]);
}

View file

@ -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 });

View file

@ -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<string, unknown>) {
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,
]);
});
});

View file

@ -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<string, SessionProjectionRun> = { ...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<Record<string, SessionProjectionRun>>,
): Readonly<Record<string, SessionProjectionRun>> {
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 =

View file

@ -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", () => {

View file

@ -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();

View file

@ -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.";

View file

@ -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"),