fix(ios): avoid transient duplicate final replies (#98117)

* Fix iOS final reply dedupe

* fix(ios): scope final message reconciliation

* docs(ios): explain final message reconciliation key

---------

Co-authored-by: joshavant <830519+joshavant@users.noreply.github.com>
This commit is contained in:
ooiuuii
2026-06-30 23:05:33 -05:00
committed by GitHub
co-authored by joshavant
parent 59df350f3c
commit 5a5913a98b
2 changed files with 413 additions and 6 deletions
@@ -47,6 +47,11 @@ public final class OpenClawChatViewModel {
}
private var pendingLocalUserEchoMessageIDsByRunID: [String: UUID] = [:]
// Final chat events and durable session-message rows arrive independently.
// Keep each provisional final scoped to the run's user turn so a later identical
// answer in the same session does not adopt or suppress the wrong row.
private var runMessageScopesByRunID: [String: RunMessageScope] = [:]
private var provisionalFinalMessagesByID: [UUID: ProvisionalFinalMessage] = [:]
private var sessionGeneration: UInt64 = 0
private var bootstrapGeneration: UInt64 = 0
// A newer same-session history request only invalidates older responses after it applies.
@@ -111,6 +116,17 @@ public final class OpenClawChatViewModel {
var timestamp: Double?
}
private struct RunMessageScope {
var session: SessionSnapshot
var latestUserTurn: LatestUserTurn?
}
private struct ProvisionalFinalMessage {
var reconciliationKey: String
var runId: String?
var scope: RunMessageScope
}
private var pendingToolCallsById: [String: OpenClawChatPendingToolCall] = [:] {
didSet {
self.pendingToolCalls = self.pendingToolCallsById.values
@@ -359,6 +375,8 @@ public final class OpenClawChatViewModel {
Self.reconcileMessageIDs(previous: self.messages, incoming: incoming)
}
self.prunePendingLocalUserEchoMessageIDs()
self.pruneProvisionalFinalMessages()
self.pruneRunMessageScopes()
self.sessionId = payload.sessionId
// Incomplete refreshes can arrive before durable assistant history.
// The latest visible user turn must survive answered before it can reject older replies.
@@ -526,6 +544,14 @@ public final class OpenClawChatViewModel {
}.joined(separator: "\\u{001E}")
}
private static func finalMessageContentFingerprint(for message: OpenClawChatMessage) -> String {
message.content.map { item in
let type = (item.type ?? "text").trimmingCharacters(in: .whitespacesAndNewlines).lowercased()
let text = (item.text ?? "").trimmingCharacters(in: .whitespacesAndNewlines)
return [type, text].joined(separator: "\\u{001F}")
}.joined(separator: "\\u{001E}")
}
private static func messageIdentityKey(for message: OpenClawChatMessage) -> String? {
let role = message.role.trimmingCharacters(in: .whitespacesAndNewlines).lowercased()
guard !role.isEmpty else { return nil }
@@ -557,6 +583,104 @@ public final class OpenClawChatViewModel {
return [role, toolCallId, toolName, contentFingerprint].joined(separator: "|")
}
private static func finalMessageReconciliationKey(for message: OpenClawChatMessage) -> String? {
let role = message.role.trimmingCharacters(in: .whitespacesAndNewlines).lowercased()
guard role == "assistant" else { return nil }
// chat.final and session.message can serialize the same final row with
// different timestamps/content ids; the run user-turn scope owns safety
// for repeated same-text replies across turns.
let contentFingerprint = Self.finalMessageContentFingerprint(for: message)
let toolCallId = (message.toolCallId ?? "").trimmingCharacters(in: .whitespacesAndNewlines)
let toolName = (message.toolName ?? "").trimmingCharacters(in: .whitespacesAndNewlines)
if contentFingerprint.isEmpty, toolCallId.isEmpty, toolName.isEmpty {
return nil
}
return [role, toolCallId, toolName, contentFingerprint].joined(separator: "|")
}
private static func normalizedRunID(_ runId: String?) -> String? {
let trimmed = runId?.trimmingCharacters(in: .whitespacesAndNewlines) ?? ""
return trimmed.isEmpty ? nil : trimmed
}
private func currentRunMessageScope() -> RunMessageScope {
RunMessageScope(
session: self.currentSessionSnapshot(),
latestUserTurn: Self.latestUserTurn(in: self.messages))
}
private func runMessageScope(for runId: String?) -> RunMessageScope {
guard let runId = Self.normalizedRunID(runId),
let scope = self.runMessageScopesByRunID[runId],
self.isCurrentSession(scope.session)
else {
return self.currentRunMessageScope()
}
return scope
}
private static func isSameUserTurnBoundary(_ lhs: LatestUserTurn?, _ rhs: LatestUserTurn?) -> Bool {
switch (lhs, rhs) {
case (nil, nil):
return true
case let (lhs?, rhs?):
if let lhsKey = lhs.refreshKey, let rhsKey = rhs.refreshKey {
return lhsKey == rhsKey && lhs.occurrence == rhs.occurrence
}
return lhs.refreshKey == nil &&
rhs.refreshKey == nil &&
lhs.occurrence == rhs.occurrence &&
lhs.timestamp == rhs.timestamp
default:
return false
}
}
private static func indexAfterLatestUserTurn(
_ latestUserTurn: LatestUserTurn?,
in messages: [OpenClawChatMessage])
-> [OpenClawChatMessage].Index
{
guard let latestUserTurn else { return messages.startIndex }
if let refreshKey = latestUserTurn.refreshKey {
var occurrence = 0
for index in messages.indices {
guard self.userRefreshIdentityKey(for: messages[index]) == refreshKey else { continue }
occurrence += 1
if occurrence == latestUserTurn.occurrence {
return messages.index(after: index)
}
}
} else if let timestamp = latestUserTurn.timestamp,
let index = messages.lastIndex(where: { message in
message.role.trimmingCharacters(in: .whitespacesAndNewlines).lowercased() == "user" &&
message.timestamp == timestamp
})
{
return messages.index(after: index)
}
return messages.lastIndex(where: { message in
message.role.trimmingCharacters(in: .whitespacesAndNewlines).lowercased() == "user"
}).map { messages.index(after: $0) } ?? messages.startIndex
}
private func hasCanonicalFinalMessageMatching(
_ message: OpenClawChatMessage,
scope: RunMessageScope) -> Bool
{
guard let key = Self.finalMessageReconciliationKey(for: message) else { return false }
guard self.isCurrentSession(scope.session) else { return false }
let searchStart = Self.indexAfterLatestUserTurn(scope.latestUserTurn, in: self.messages)
guard searchStart < self.messages.endIndex else { return false }
return self.messages[searchStart...].contains { existing in
self.provisionalFinalMessagesByID[existing.id] == nil &&
Self.finalMessageReconciliationKey(for: existing) == key
}
}
private func prunePendingLocalUserEchoMessageIDs() {
guard !self.pendingLocalUserEchoMessageIDsByRunID.isEmpty else { return }
let visibleMessageIDs = Set(messages.map(\.id))
@@ -565,6 +689,26 @@ public final class OpenClawChatViewModel {
}
}
private func pruneProvisionalFinalMessages() {
guard !self.provisionalFinalMessagesByID.isEmpty else { return }
let visibleMessageIDs = Set(messages.map(\.id))
self.provisionalFinalMessagesByID = self.provisionalFinalMessagesByID.filter { entry in
visibleMessageIDs.contains(entry.key) && self.isCurrentSession(entry.value.scope.session)
}
}
private func pruneRunMessageScopes() {
self.runMessageScopesByRunID = self.runMessageScopesByRunID.filter { entry in
self.isCurrentSession(entry.value.session)
}
guard self.runMessageScopesByRunID.count > 64 else { return }
let referencedRunIDs = Set(self.pendingRuns)
.union(self.provisionalFinalMessagesByID.values.compactMap(\.runId))
self.runMessageScopesByRunID = self.runMessageScopesByRunID.filter { entry in
referencedRunIDs.contains(entry.key)
}
}
private func adoptPendingLocalUserEcho(incoming: OpenClawChatMessage) -> Bool {
guard let incomingKey = Self.userRefreshIdentityKey(for: incoming) else { return false }
guard let matchIndex = messages.lastIndex(where: { existing in
@@ -594,6 +738,43 @@ public final class OpenClawChatViewModel {
return true
}
private func adoptProvisionalFinalMessage(incoming: OpenClawChatMessage) -> Bool {
guard let incomingKey = Self.finalMessageReconciliationKey(for: incoming) else { return false }
let canonicalScope = self.currentRunMessageScope()
let searchStart = Self.indexAfterLatestUserTurn(canonicalScope.latestUserTurn, in: self.messages)
guard searchStart < self.messages.endIndex else { return false }
guard let matchIndex = messages[searchStart...].lastIndex(where: { existing in
guard let provisional = self.provisionalFinalMessagesByID[existing.id] else { return false }
return provisional.reconciliationKey == incomingKey &&
Self.isSameUserTurnBoundary(provisional.scope.latestUserTurn, canonicalScope.latestUserTurn)
}) else {
return false
}
let existing = self.messages[matchIndex]
let provisional = self.provisionalFinalMessagesByID[existing.id]
var updated = self.messages
updated[matchIndex] = OpenClawChatMessage(
id: existing.id,
role: incoming.role,
content: incoming.content,
timestamp: incoming.timestamp ?? existing.timestamp,
toolCallId: incoming.toolCallId,
toolName: incoming.toolName,
usage: incoming.usage,
stopReason: incoming.stopReason,
errorMessage: incoming.errorMessage)
self.provisionalFinalMessagesByID.removeValue(forKey: existing.id)
if let runId = provisional?.runId {
self.runMessageScopesByRunID.removeValue(forKey: runId)
}
self.messages = Self.dedupeMessages(updated)
self.pruneProvisionalFinalMessages()
self.pruneRunMessageScopes()
return true
}
private static func reconcileMessageIDs(
previous: [OpenClawChatMessage],
incoming: [OpenClawChatMessage]) -> [OpenClawChatMessage]
@@ -840,6 +1021,7 @@ public final class OpenClawChatViewModel {
content: userContent,
timestamp: userMessageTimestamp))
self.pendingLocalUserEchoMessageIDsByRunID[runId] = userMessageID
self.runMessageScopesByRunID[runId] = self.currentRunMessageScope()
// Clear input immediately for responsive UX (before network await)
self.input = ""
@@ -863,9 +1045,11 @@ public final class OpenClawChatViewModel {
+ "localRunId=\(runId) remoteRunId=\(response.runId)")
if response.runId != runId {
let pendingUserMessageID = self.pendingLocalUserEchoMessageIDsByRunID.removeValue(forKey: runId)
let runScope = self.runMessageScopesByRunID.removeValue(forKey: runId)
self.clearPendingRun(runId)
self.pendingRuns.insert(response.runId)
self.pendingLocalUserEchoMessageIDsByRunID[response.runId] = pendingUserMessageID
self.runMessageScopesByRunID[response.runId] = runScope
self.armPendingRunTimeout(runId: response.runId)
}
if response.status == "ok" {
@@ -894,6 +1078,7 @@ public final class OpenClawChatViewModel {
} catch {
guard self.isCurrentSession(sessionSnapshot) else { return }
self.removePendingLocalUserEcho(for: runId)
self.runMessageScopesByRunID.removeValue(forKey: runId)
self.clearPendingRun(runId)
self.errorText = error.localizedDescription
self.logDiagnostic(
@@ -955,6 +1140,8 @@ public final class OpenClawChatViewModel {
self.modelSelectionID = Self.defaultModelSelectionID
self.messages = []
self.pendingLocalUserEchoMessageIDsByRunID.removeAll()
self.runMessageScopesByRunID.removeAll()
self.provisionalFinalMessagesByID.removeAll()
self.sessionId = nil
self.pendingToolCallsById = [:]
self.streamingAssistantText = nil
@@ -989,6 +1176,8 @@ public final class OpenClawChatViewModel {
self.modelSelectionID = Self.defaultModelSelectionID
self.messages = []
self.pendingLocalUserEchoMessageIDsByRunID.removeAll()
self.runMessageScopesByRunID.removeAll()
self.provisionalFinalMessagesByID.removeAll()
self.sessionId = nil
self.pendingToolCallsById = [:]
self.streamingAssistantText = nil
@@ -1016,6 +1205,8 @@ public final class OpenClawChatViewModel {
return
}
self.runMessageScopesByRunID.removeAll()
self.provisionalFinalMessagesByID.removeAll()
self.startBootstrap()
}
@@ -1519,9 +1710,14 @@ public final class OpenClawChatViewModel {
if self.adoptPendingLocalUserEcho(incoming: sanitized) {
return
}
if self.adoptProvisionalFinalMessage(incoming: sanitized) {
return
}
let reconciled = Self.reconcileMessageIDs(previous: self.messages, incoming: self.messages + [sanitized])
self.messages = Self.dedupeMessages(reconciled)
self.pruneProvisionalFinalMessages()
self.pruneRunMessageScopes()
}
private func handleChatEvent(_ chat: OpenClawChatEventPayload) {
@@ -1608,8 +1804,28 @@ public final class OpenClawChatViewModel {
stopReason: "stop")
}
let runId = Self.normalizedRunID(chat.runId)
let scope = self.runMessageScope(for: runId)
guard self.isCurrentSession(scope.session) else { return }
guard let reconciliationKey = Self.finalMessageReconciliationKey(for: message) else { return }
if self.hasCanonicalFinalMessageMatching(message, scope: scope) {
if let runId {
self.runMessageScopesByRunID.removeValue(forKey: runId)
}
return
}
let reconciled = Self.reconcileMessageIDs(previous: self.messages, incoming: self.messages + [message])
self.messages = Self.dedupeMessages(reconciled)
if self.messages.contains(where: { $0.id == message.id }) {
self.provisionalFinalMessagesByID[message.id] = ProvisionalFinalMessage(
reconciliationKey: reconciliationKey,
runId: runId,
scope: scope)
}
self.pruneProvisionalFinalMessages()
self.pruneRunMessageScopes()
}
private static func isAssistantMessage(_ message: OpenClawChatMessage) -> Bool {
@@ -1761,7 +1977,7 @@ public final class OpenClawChatViewModel {
}
private func removePendingLocalUserEcho(for runId: String) {
guard let messageID = self.pendingLocalUserEchoMessageIDsByRunID[runId] else { return }
guard let messageID = pendingLocalUserEchoMessageIDsByRunID[runId] else { return }
self.messages.removeAll { $0.id == messageID }
self.pendingLocalUserEchoMessageIDsByRunID[runId] = nil
}
@@ -1852,7 +2068,7 @@ public final class OpenClawChatViewModel {
guard let lastUserIndex = messages.lastIndex(where: { $0.role.lowercased() == "user" }) else {
return nil
}
guard let refreshKey = self.userRefreshIdentityKey(for: messages[lastUserIndex]) else {
guard let refreshKey = userRefreshIdentityKey(for: messages[lastUserIndex]) else {
return LatestUserTurn(
refreshKey: nil,
occurrence: 0,
@@ -3,14 +3,32 @@ import OpenClawKit
import Testing
@testable import OpenClawChatUI
private func chatTextMessage(role: String, text: String, timestamp: Double) -> AnyCodable {
AnyCodable([
private func chatTextMessage(role: String, text: String, timestamp: Double, contentId: String? = nil) -> AnyCodable {
var content: [String: Any] = ["type": "text", "text": text]
if let contentId {
content["id"] = contentId
}
return AnyCodable([
"role": role,
"content": [["type": "text", "text": text]],
"content": [content],
"timestamp": timestamp,
])
}
private func chatTextModelMessage(role: String, text: String, timestamp: Double) -> OpenClawChatMessage {
OpenClawChatMessage(
role: role,
content: [
OpenClawChatMessageContent(
type: "text",
text: text,
mimeType: nil,
fileName: nil,
content: nil),
],
timestamp: timestamp)
}
private func chatErrorMessage(role: String, errorMessage: String, timestamp: Double) -> AnyCodable {
AnyCodable([
"role": role,
@@ -799,6 +817,175 @@ struct ChatViewModelTests {
}
}
@Test func `session message adopts provisional final event reply`() async throws {
let sessionId = "sess-main"
let now = Date().timeIntervalSince1970 * 1000
let finalRefreshGate = SessionSubscribeGate()
let historyCount = AsyncCounter()
let history = historyPayload(sessionId: sessionId)
let (transport, vm) = await makeViewModel(
historyResponses: [history, history],
requestHistoryHook: { _ in
let count = await historyCount.increment()
if count == 2 {
await finalRefreshGate.wait()
}
})
try await loadAndWaitBootstrap(vm: vm, sessionId: sessionId)
await sendUserMessage(vm, text: "hello")
try await waitUntil("pending run starts") { await MainActor.run { vm.pendingRunCount == 1 } }
let runId = try await waitForLastSentRunId(transport)
transport.emit(
.chat(
OpenClawChatEventPayload(
runId: runId,
sessionKey: "main",
state: "final",
message: chatTextMessage(
role: "assistant",
text: "dedupe me",
timestamp: now + 1,
contentId: "live-final-content"),
errorMessage: nil)))
try await waitUntil("provisional final visible once") {
await MainActor.run {
vm.messages.count(where: { msg in
msg.role == "assistant" && msg.content.first?.text == "dedupe me"
}) == 1
}
}
transport.emit(
.sessionMessage(
OpenClawSessionMessageEventPayload(
sessionKey: "agent:main:main",
message: chatTextModelMessage(role: "assistant", text: "dedupe me", timestamp: now + 2),
messageId: "msg-assistant-final",
messageSeq: 2)))
try await waitUntil("canonical session message adopted final event row") {
await MainActor.run {
let matches = vm.messages.filter { msg in
msg.role == "assistant" && msg.content.first?.text == "dedupe me"
}
return matches.count == 1 && matches.first?.timestamp == now + 2
}
}
await finalRefreshGate.release()
}
@Test func `final event does not duplicate canonical assistant session message`() async throws {
let sessionId = "sess-main"
let now = Date().timeIntervalSince1970 * 1000
let finalRefreshGate = SessionSubscribeGate()
let historyCount = AsyncCounter()
let history = historyPayload(sessionId: sessionId)
let (transport, vm) = await makeViewModel(
historyResponses: [history, history],
requestHistoryHook: { _ in
let count = await historyCount.increment()
if count == 2 {
await finalRefreshGate.wait()
}
})
try await loadAndWaitBootstrap(vm: vm, sessionId: sessionId)
await sendUserMessage(vm, text: "hello")
try await waitUntil("pending run starts") { await MainActor.run { vm.pendingRunCount == 1 } }
let runId = try await waitForLastSentRunId(transport)
transport.emit(
.sessionMessage(
OpenClawSessionMessageEventPayload(
sessionKey: "agent:main:main",
message: chatTextModelMessage(role: "assistant", text: "canonical first", timestamp: now + 2),
messageId: "msg-assistant-first",
messageSeq: 2)))
try await waitUntil("canonical assistant visible once") {
await MainActor.run {
vm.messages.count(where: { msg in
msg.role == "assistant" && msg.content.first?.text == "canonical first"
}) == 1
}
}
transport.emit(
.chat(
OpenClawChatEventPayload(
runId: runId,
sessionKey: "main",
state: "final",
message: chatTextMessage(role: "assistant", text: "canonical first", timestamp: now + 1),
errorMessage: nil)))
try await Task.sleep(nanoseconds: 50_000_000)
#expect(await MainActor.run {
let matches = vm.messages.filter { msg in
msg.role == "assistant" && msg.content.first?.text == "canonical first"
}
return matches.count == 1 && matches.first?.timestamp == now + 2
})
await finalRefreshGate.release()
}
@Test func `later identical session reply does not adopt prior turn provisional final`() async throws {
let sessionId = "sess-main"
let now = Date().timeIntervalSince1970 * 1000
let history = historyPayload(sessionId: sessionId)
let (transport, vm) = await makeViewModel(
historyResponses: [history, history],
sendMessageHook: { runId in
OpenClawChatSendResponse(runId: runId, status: "pending")
})
try await loadAndWaitBootstrap(vm: vm, sessionId: sessionId)
await sendUserMessage(vm, text: "first turn")
try await waitUntil("first pending run starts") { await MainActor.run { vm.pendingRunCount == 1 } }
let firstRunId = try await waitForLastSentRunId(transport)
transport.emit(
.chat(
OpenClawChatEventPayload(
runId: firstRunId,
sessionKey: "main",
state: "final",
message: chatTextMessage(role: "assistant", text: "OK", timestamp: now + 1),
errorMessage: nil)))
try await waitUntil("first provisional final visible") {
await MainActor.run {
vm.messages.count(where: { msg in
msg.role == "assistant" && msg.content.first?.text == "OK"
}) == 1
}
}
await sendUserMessage(vm, text: "second turn")
try await waitUntil("second pending run starts") { await MainActor.run { vm.pendingRunCount == 1 } }
transport.emit(
.sessionMessage(
OpenClawSessionMessageEventPayload(
sessionKey: "agent:main:main",
message: chatTextModelMessage(role: "assistant", text: "OK", timestamp: now + 4),
messageId: "msg-second-assistant",
messageSeq: 4)))
try await waitUntil("second identical reply appends after second user") {
await MainActor.run {
let okReplies = vm.messages.filter { msg in
msg.role == "assistant" && msg.content.first?.text == "OK"
}
return okReplies.count == 2 && vm.messages.last?.timestamp == now + 4
}
}
}
@Test func `completion wait refreshes history and clears pending run`() async throws {
let sessionId = "sess-main"
let now = (Date().timeIntervalSince1970 * 1000) + 10000
@@ -1455,7 +1642,11 @@ struct ChatViewModelTests {
}
@Test func `dedupes gateway echo of local user message`() async throws {
let (transport, vm) = await makeViewModel(historyResponses: [historyPayload()])
let (transport, vm) = await makeViewModel(
historyResponses: [historyPayload()],
sendMessageHook: { runId in
OpenClawChatSendResponse(runId: runId, status: "pending")
})
await MainActor.run { vm.load() }
try await waitUntil("bootstrap history loaded") { await MainActor.run { vm.messages.isEmpty } }