diff --git a/apps/shared/OpenClawKit/Tests/OpenClawKitTests/ChatViewModelOutboxTests.swift b/apps/shared/OpenClawKit/Tests/OpenClawKitTests/ChatViewModelOutboxTests.swift index baea3a6d1d12..66cb0a33188f 100644 --- a/apps/shared/OpenClawKit/Tests/OpenClawKitTests/ChatViewModelOutboxTests.swift +++ b/apps/shared/OpenClawKit/Tests/OpenClawKitTests/ChatViewModelOutboxTests.swift @@ -4,62 +4,16 @@ import OpenClawKit import Testing @testable import OpenClawChatUI -extension OpenClawChatSQLiteTranscriptCache { - fileprivate func markCommandAwaitingConfirmation(id: String) async -> OpenClawChatOutboxUpdateResult { - guard let command = await self.loadCommands().first(where: { $0.id == id }) else { return .missing } - return await self.markCommandAwaitingConfirmation(id: id, attemptVersion: command.attemptVersion) - } - - fileprivate func markCommandFailedIfPresent( - id: String, - retryCount: Int, - lastError: String?) async -> OpenClawChatOutboxUpdateResult - { - guard let command = await self.loadCommands().first(where: { $0.id == id }) else { return .missing } - return await self.markCommandFailedIfPresent( - id: id, - attemptVersion: command.attemptVersion, - retryCount: retryCount, - lastError: lastError) - } - - private func markCommandRetriedIfPresent( - id: String, - agentID: String?, - deliverySessionKey: String, - routingContract: String) async -> OpenClawChatOutboxUpdateResult - { - guard let command = await self.loadCommands().first(where: { $0.id == id }) else { return .missing } - return await self.markCommandRetriedIfPresent( - id: id, - expectation: OpenClawChatOutboxRetryExpectation( - attemptVersion: command.attemptVersion, - retryCount: command.retryCount, - lastError: command.lastError), - agentID: agentID, - deliverySessionKey: deliverySessionKey, - routingContract: routingContract, - replacementID: nil) - } - - fileprivate func confirmCommand(id: String) async -> OpenClawChatOutboxUpdateResult { - guard let command = await self.loadCommands().first(where: { $0.id == id }) else { return .missing } - return await self.confirmCommand(id: id, attemptVersion: command.attemptVersion) - } -} - -private func makeOutboxDatabaseDirectory() throws -> URL { - let dir = FileManager.default.temporaryDirectory - .appendingPathComponent("chat-outbox-tests-\(UUID().uuidString)", isDirectory: true) - try FileManager.default.createDirectory(at: dir, withIntermediateDirectories: true) - return dir -} - -private func makeOutboxStore( - databaseDirectoryURL: URL, - gatewayID: String) throws -> OpenClawChatSQLiteTranscriptCache +private func makeOutboxStore() throws -> ( + store: OpenClawChatSQLiteTranscriptCache, + databases: OpenClawClientDatabases, + directory: URL) { - try OpenClawClientDatabases(directoryURL: databaseDirectoryURL).store(gatewayID: gatewayID) + let directory = FileManager.default.temporaryDirectory + .appendingPathComponent("chat-outbox-tests-\(UUID().uuidString)", isDirectory: true) + try FileManager.default.createDirectory(at: directory, withIntermediateDirectories: true) + let databases = try OpenClawClientDatabases(directoryURL: directory) + return (databases.store(gatewayID: "gw-test"), databases, directory) } extension OpenClawChatSQLiteTranscriptCache { @@ -122,16 +76,8 @@ private actor OutboxTransportState { let sessionListStarted = DeleteGate() var staleHistoryRows: [AnyCodable]? - func setHeldSendGate(_ gate: DeleteGate?) { - self.heldSendGate = gate - } - - func setCommandListGate(_ gate: DeleteGate?) { - self.commandListGate = gate - } - - func waitUntilCommandListStarted() async { - await self.commandListStarted.wait() + func update(_ operation: @Sendable (isolated OutboxTransportState) -> Void) { + operation(self) } func awaitCommandListGate() async { @@ -141,10 +87,6 @@ private actor OutboxTransportState { } } - func setSessionListGate(_ gate: DeleteGate?) { - self.sessionListGate = gate - } - func awaitSessionListGate() async { await self.sessionListStarted.open() if let sessionListGate { @@ -152,10 +94,6 @@ private actor OutboxTransportState { } } - func setStaleHistoryRows(_ rows: [AnyCodable]?) { - self.staleHistoryRows = rows - } - var sentIdempotencyKeys: [String] = [] var sentMessages: [String] = [] var sentSessionKeys: [String] = [] @@ -168,51 +106,11 @@ private actor OutboxTransportState { self.sendFails = sendFails } - func setHistoryFails(_ fails: Bool) { - self.historyFails = fails - } - - func setSessionListFails(_ fails: Bool) { - self.sessionListFails = fails - } - func recordHistoryRequest(agentID: String?) { self.historyRequestCount += 1 self.historyRequestAgentIDs.append(agentID) } - func setHealthy(_ healthy: Bool) { - self.healthy = healthy - } - - func setSessionRoutingContract(_ contract: String) { - self.sessionRoutingContract = contract - } - - func replaceRoute() { - self.routeGeneration += 1 - } - - func setSendFails(_ fails: Bool) { - self.sendFails = fails - } - - func setSendFailsAfterRecording(_ fails: Bool) { - self.sendFailsAfterRecording = fails - } - - func setSendRejects(_ rejects: Bool) { - self.sendRejects = rejects - } - - func setSendResponseErrors(_ rejects: Bool) { - self.sendResponseErrors = rejects - } - - func setSendRoutingChanged(_ changed: Bool) { - self.sendRoutingChanged = changed - } - func recordSend( sessionKey: String, agentID: String?, @@ -259,7 +157,7 @@ private final class OutboxTestTransport: @unchecked Sendable, OpenClawChatTransp } func goOnline() async { - await self.state.setHealthy(true) + await self.state.update { $0.healthy = true } self.continuation.yield(.health(ok: true)) } @@ -357,7 +255,7 @@ private final class OutboxTestTransport: @unchecked Sendable, OpenClawChatTransp if let gate = await state.heldSendGate { // One-shot: only the first send is held so tests can pin the // window where the flush is mid-drain. - await self.state.setHeldSendGate(nil) + await self.state.update { $0.heldSendGate = nil } await gate.wait() } if let expectedRoute, await state.routeGeneration != expectedRoute { @@ -545,24 +443,27 @@ private func sendWhileOffline(_ vm: OpenClawChatViewModel, text: String) async t } } -@MainActor -private func queuedStateCount(_ vm: OpenClawChatViewModel) -> Int { - vm.outboxStatesByMessageID.count -} +/// Protocol delegation plus switches makes race windows deterministic without copying the store contract. +private actor ScriptedOutbox: OpenClawChatCommandOutbox { + enum Forwarding { case full, minimal, holdingCancellation } -/// Forwarding outbox that can delay `loadCommands`, making restore-vs-send -/// interleavings deterministic in tests. -private actor DelayingOutbox: OpenClawChatCommandOutbox { private nonisolated let base: OpenClawChatSQLiteTranscriptCache + private let forwarding: Forwarding private var loadDelayNanoseconds: UInt64 = 0 private var enqueueRelease: DeleteGate? private var recoveryAvailable = true private var terminalWritesAvailable = true + private var captured = DeleteGate() + private var snapshotRelease = DeleteGate() + private var shouldHoldNextLoad = false private let enqueueStarted = DeleteGate() private let recoveryAttempted = DeleteGate() + private let canceled = DeleteGate() + private let cancellationRelease = DeleteGate() - init(base: OpenClawChatSQLiteTranscriptCache) { + init(base: OpenClawChatSQLiteTranscriptCache, forwarding: Forwarding = .full) { self.base = base + self.forwarding = forwarding } nonisolated func changes() -> AsyncStream { @@ -598,6 +499,28 @@ private actor DelayingOutbox: OpenClawChatCommandOutbox { self.enqueueRelease = nil } + func holdNextLoad() { + self.captured = DeleteGate() + self.snapshotRelease = DeleteGate() + self.shouldHoldNextLoad = true + } + + func waitUntilSnapshotCaptured() async { + await self.captured.wait() + } + + func releaseSnapshot() async { + await self.snapshotRelease.open() + } + + func waitUntilCanceled() async { + await self.canceled.wait() + } + + func releaseCancellation() async { + await self.cancellationRelease.open() + } + func enqueueCommand(_ command: OpenClawChatOutboxCommand) async -> Bool { if let enqueueRelease { await self.enqueueStarted.open() @@ -607,17 +530,30 @@ private actor DelayingOutbox: OpenClawChatCommandOutbox { } func loadCommands() async -> [OpenClawChatOutboxCommand] { - if self.loadDelayNanoseconds > 0 { - try? await Task.sleep(nanoseconds: self.loadDelayNanoseconds) - } - return await self.base.loadCommands() + await self.delayLoad() + let commands = await self.base.loadCommands() + await self.finishHeldLoad() + return commands } func loadCommandsIfAvailable() async -> [OpenClawChatOutboxCommand]? { - if self.loadDelayNanoseconds > 0 { - try? await Task.sleep(nanoseconds: self.loadDelayNanoseconds) + await self.delayLoad() + guard let commands = await self.base.loadCommandsIfAvailable() else { return nil } + await self.finishHeldLoad() + return commands + } + + private func delayLoad() async { + guard self.loadDelayNanoseconds > 0 else { return } + try? await Task.sleep(nanoseconds: self.loadDelayNanoseconds) + } + + private func finishHeldLoad() async { + if self.shouldHoldNextLoad { + self.shouldHoldNextLoad = false + await self.captured.open() + await self.snapshotRelease.wait() } - return await self.base.loadCommandsIfAvailable() } @discardableResult @@ -662,7 +598,12 @@ private actor DelayingOutbox: OpenClawChatCommandOutbox { } func cancelCommand(id: String) async -> OpenClawChatOutboxUpdateResult { - await self.base.cancelCommand(id: id) + let result = await self.base.cancelCommand(id: id) + if self.forwarding == .holdingCancellation { + await self.canceled.open() + await self.cancellationRelease.wait() + } + return result } func confirmCommand(id: String, attemptVersion: Int) async -> OpenClawChatOutboxUpdateResult { @@ -670,15 +611,18 @@ private actor DelayingOutbox: OpenClawChatCommandOutbox { } func branchState(for scope: OpenClawChatOutboxScope) async -> OpenClawChatOutboxBranchState? { - await self.base.branchState(for: scope) + guard self.forwarding == .full else { return nil } + return await self.base.branchState(for: scope) } func beginBranchSwitch(_ scope: OpenClawChatOutboxScope) async -> Bool { - await self.base.beginBranchSwitch(scope) + guard self.forwarding == .full else { return false } + return await self.base.beginBranchSwitch(scope) } func cancelBranchSwitch(_ scope: OpenClawChatOutboxScope) async -> Bool { - await self.base.cancelBranchSwitch(scope) + guard self.forwarding == .full else { return false } + return await self.base.cancelBranchSwitch(scope) } func reconcileBranchScope( @@ -689,7 +633,8 @@ private actor DelayingOutbox: OpenClawChatCommandOutbox { activeTranscriptEntryIDs: Set, lastError: String) async -> [OpenClawChatOutboxCommand]? { - await self.base.reconcileBranchScope( + guard self.forwarding == .full else { return nil } + return await self.base.reconcileBranchScope( scope, previousState: previousState, activeLeafEntryID: activeLeafEntryID, branchLeafEntryIDs: branchLeafEntryIDs, activeTranscriptEntryIDs: activeTranscriptEntryIDs, @@ -701,7 +646,11 @@ private actor DelayingOutbox: OpenClawChatCommandOutbox { activeLeafEntryID: String, lastError: String) async -> [OpenClawChatOutboxCommand]? { - await self.base.confirmBranchChange(scope, activeLeafEntryID: activeLeafEntryID, lastError: lastError) + guard self.forwarding == .full else { return nil } + return await self.base.confirmBranchChange( + scope, + activeLeafEntryID: activeLeafEntryID, + lastError: lastError) } func updateLastActiveLeafEntryID( @@ -709,203 +658,18 @@ private actor DelayingOutbox: OpenClawChatCommandOutbox { expectedEpoch: Int, for scope: OpenClawChatOutboxScope) async -> Bool { - await self.base.updateLastActiveLeafEntryID(leafEntryID, expectedEpoch: expectedEpoch, for: scope) - } -} - -/// Returns one already-read command snapshot only after the test releases it, -/// reproducing a restore that resumes after another view canceled the row. -private actor SnapshotHoldingOutbox: OpenClawChatCommandOutbox { - private nonisolated let base: OpenClawChatSQLiteTranscriptCache - private var captured = DeleteGate() - private var release = DeleteGate() - private var shouldHoldNextLoad = false - - init(base: OpenClawChatSQLiteTranscriptCache) { - self.base = base - } - - func waitUntilSnapshotCaptured() async { - await self.captured.wait() - } - - func holdNextLoad() { - self.captured = DeleteGate() - self.release = DeleteGate() - self.shouldHoldNextLoad = true - } - - func releaseSnapshot() async { - await self.release.open() - } - - nonisolated func changes() -> AsyncStream { - self.base.changes() - } - - func enqueueCommand(_ command: OpenClawChatOutboxCommand) async -> Bool { - await self.base.enqueueCommand(command) - } - - func loadCommands() async -> [OpenClawChatOutboxCommand] { - let commands = await base.loadCommands() - if self.shouldHoldNextLoad { - self.shouldHoldNextLoad = false - await self.captured.open() - await self.release.wait() - } - return commands - } - - func loadCommandsIfAvailable() async -> [OpenClawChatOutboxCommand]? { - guard let commands = await base.loadCommandsIfAvailable() else { return nil } - if self.shouldHoldNextLoad { - self.shouldHoldNextLoad = false - await self.captured.open() - await self.release.wait() - } - return commands - } - - @discardableResult - func recoverInterruptedSends() async -> Bool { - await self.base.recoverInterruptedSends() - } - - func claimNextCommand() async -> OpenClawChatOutboxCommand? { - await self.base.claimNextCommand() - } - - func markCommandQueued( - id: String, - attemptVersion: Int, - retryCount: Int, - lastError: String?) async -> OpenClawChatOutboxUpdateResult - { - await self.base.markCommandQueued( - id: id, - attemptVersion: attemptVersion, - retryCount: retryCount, - lastError: lastError) - } - - func markCommandAwaitingConfirmation(id: String, attemptVersion: Int) async -> OpenClawChatOutboxUpdateResult { - await self.base.markCommandAwaitingConfirmation(id: id, attemptVersion: attemptVersion) - } - - func markCommandFailedIfPresent( - id: String, attemptVersion: Int, - retryCount: Int, - lastError: String?) async -> OpenClawChatOutboxUpdateResult - { - await self.base.markCommandFailedIfPresent( - id: id, - attemptVersion: attemptVersion, - retryCount: retryCount, - lastError: lastError) - } - - func cancelCommand(id: String) async -> OpenClawChatOutboxUpdateResult { - await self.base.cancelCommand(id: id) - } - - func confirmCommand(id: String, attemptVersion: Int) async -> OpenClawChatOutboxUpdateResult { - await self.base.confirmCommand(id: id, attemptVersion: attemptVersion) - } -} - -/// Holds a completed durable cancellation before its result returns to the -/// MainActor, making late canonical-proof ordering deterministic. -private actor CancellationHoldingOutbox: OpenClawChatCommandOutbox { - private nonisolated let base: OpenClawChatSQLiteTranscriptCache - private let canceled = DeleteGate() - private let release = DeleteGate() - - init(base: OpenClawChatSQLiteTranscriptCache) { - self.base = base - } - - func waitUntilCanceled() async { - await self.canceled.wait() - } - - func releaseCancellation() async { - await self.release.open() - } - - nonisolated func changes() -> AsyncStream { - self.base.changes() - } - - func enqueueCommand(_ command: OpenClawChatOutboxCommand) async -> Bool { - await self.base.enqueueCommand(command) - } - - func loadCommands() async -> [OpenClawChatOutboxCommand] { - await self.base.loadCommands() - } - - func loadCommandsIfAvailable() async -> [OpenClawChatOutboxCommand]? { - await self.base.loadCommandsIfAvailable() - } - - @discardableResult - func recoverInterruptedSends() async -> Bool { - await self.base.recoverInterruptedSends() - } - - func claimNextCommand() async -> OpenClawChatOutboxCommand? { - await self.base.claimNextCommand() - } - - func markCommandQueued( - id: String, - attemptVersion: Int, - retryCount: Int, - lastError: String?) async -> OpenClawChatOutboxUpdateResult - { - await self.base.markCommandQueued( - id: id, - attemptVersion: attemptVersion, - retryCount: retryCount, - lastError: lastError) - } - - func markCommandAwaitingConfirmation(id: String, attemptVersion: Int) async -> OpenClawChatOutboxUpdateResult { - await self.base.markCommandAwaitingConfirmation(id: id, attemptVersion: attemptVersion) - } - - func markCommandFailedIfPresent( - id: String, attemptVersion: Int, - retryCount: Int, - lastError: String?) async -> OpenClawChatOutboxUpdateResult - { - await self.base.markCommandFailedIfPresent( - id: id, - attemptVersion: attemptVersion, - retryCount: retryCount, - lastError: lastError) - } - - func cancelCommand(id: String) async -> OpenClawChatOutboxUpdateResult { - let result = await base.cancelCommand(id: id) - await self.canceled.open() - await self.release.wait() - return result - } - - func confirmCommand(id: String, attemptVersion: Int) async -> OpenClawChatOutboxUpdateResult { - await self.base.confirmCommand(id: id, attemptVersion: attemptVersion) + guard self.forwarding == .full else { return false } + return await self.base.updateLastActiveLeafEntryID( + leafEntryID, + expectedEpoch: expectedEpoch, + for: scope) } } struct ChatViewModelOutboxTests { @Test func `offline send queues durably and renders queued row`() async throws { - let databaseDirectory = try makeOutboxDatabaseDirectory() + let (store, _, databaseDirectory) = try makeOutboxStore() defer { try? FileManager.default.removeItem(at: databaseDirectory) } - let store = try makeOutboxStore( - databaseDirectoryURL: databaseDirectory, - gatewayID: "gw-test") let transport = OutboxTestTransport(healthy: false) let vm = await makeOutboxViewModel(transport: transport, outbox: store) #expect(await MainActor.run { vm.supportsOfflineTextOutbox }) @@ -959,11 +723,8 @@ struct ChatViewModelOutboxTests { } @Test func `unsupported gateway keeps queued work and surfaces upgrade action`() async throws { - let databaseDirectory = try makeOutboxDatabaseDirectory() + let (store, _, databaseDirectory) = try makeOutboxStore() defer { try? FileManager.default.removeItem(at: databaseDirectory) } - let store = try makeOutboxStore( - databaseDirectoryURL: databaseDirectory, - gatewayID: "gw-test") let message = OpenClawChatTransportUpgradeMessage.routingContract let transport = OutboxTestTransport( healthy: false, @@ -982,11 +743,8 @@ struct ChatViewModelOutboxTests { } @Test func `offline queue persists the effective thinking level`() async throws { - let databaseDirectory = try makeOutboxDatabaseDirectory() + let (store, _, databaseDirectory) = try makeOutboxStore() defer { try? FileManager.default.removeItem(at: databaseDirectory) } - let store = try makeOutboxStore( - databaseDirectoryURL: databaseDirectory, - gatewayID: "gw-test") let sessions = [outboxSessionEntry(key: "main", thinkingLevels: ["off"])] let transport = OutboxTestTransport(healthy: false, sessions: sessions) let vm = await makeOutboxViewModel(transport: transport, outbox: store) @@ -1004,11 +762,8 @@ struct ChatViewModelOutboxTests { } @Test func `inert outbox does not capability gate healthy live chat`() async throws { - let databaseDirectory = try makeOutboxDatabaseDirectory() + let (store, _, databaseDirectory) = try makeOutboxStore() defer { try? FileManager.default.removeItem(at: databaseDirectory) } - let store = try makeOutboxStore( - databaseDirectoryURL: databaseDirectory, - gatewayID: "gw-test") let transport = OutboxTestTransport( healthy: true, routeUnavailableReason: OpenClawChatTransportUpgradeMessage.routingContract) @@ -1037,11 +792,8 @@ struct ChatViewModelOutboxTests { } @Test func `legacy transport preserves its untargeted session key`() async throws { - let databaseDirectory = try makeOutboxDatabaseDirectory() + let (store, _, databaseDirectory) = try makeOutboxStore() defer { try? FileManager.default.removeItem(at: databaseDirectory) } - let store = try makeOutboxStore( - databaseDirectoryURL: databaseDirectory, - gatewayID: "gw-test") let transport = OutboxTestTransport(healthy: false, requiresRoutingContract: false) let vm = await makeOutboxViewModel( transport: transport, @@ -1057,6 +809,7 @@ struct ChatViewModelOutboxTests { #expect(command.agentID == nil) #expect(await store.markCommandFailedIfPresent( id: command.id, + attemptVersion: command.attemptVersion, retryCount: 1, lastError: "legacy failure") == .updated) @@ -1089,11 +842,8 @@ struct ChatViewModelOutboxTests { } @Test func `reserved unknown session stays unscoped in durable delivery`() async throws { - let databaseDirectory = try makeOutboxDatabaseDirectory() + let (store, _, databaseDirectory) = try makeOutboxStore() defer { try? FileManager.default.removeItem(at: databaseDirectory) } - let store = try makeOutboxStore( - databaseDirectoryURL: databaseDirectory, - gatewayID: "gw-test") let transport = OutboxTestTransport(healthy: false) let vm = await makeOutboxViewModel( transport: transport, @@ -1117,11 +867,8 @@ struct ChatViewModelOutboxTests { @Test(arguments: ["global", "main"]) func `mutable alias queued turn keeps its original agent target`(_ sessionKey: String) async throws { - let databaseDirectory = try makeOutboxDatabaseDirectory() + let (store, _, databaseDirectory) = try makeOutboxStore() defer { try? FileManager.default.removeItem(at: databaseDirectory) } - let store = try makeOutboxStore( - databaseDirectoryURL: databaseDirectory, - gatewayID: "gw-test") let offlineTransport = OutboxTestTransport(healthy: false) let agentAView = await makeOutboxViewModel( transport: offlineTransport, @@ -1160,11 +907,8 @@ struct ChatViewModelOutboxTests { } @Test func `reconnect waits for canonical session metadata before flushing Ultra`() async throws { - let databaseDirectory = try makeOutboxDatabaseDirectory() + let (store, _, databaseDirectory) = try makeOutboxStore() defer { try? FileManager.default.removeItem(at: databaseDirectory) } - let store = try makeOutboxStore( - databaseDirectoryURL: databaseDirectory, - gatewayID: "gw-test") #expect(await store.updateLastActiveLeafEntryID( "leaf-new", expectedEpoch: 0, @@ -1191,7 +935,7 @@ struct ChatViewModelOutboxTests { outboxSessionEntry(key: "agent:beta:main", thinkingLevels: ["off", "ultra"]), ] let transport = OutboxTestTransport(healthy: false, sessions: sessions) - await transport.state.setSessionListGate(listGate) + await transport.state.update { $0.sessionListGate = listGate } let vm = await makeOutboxViewModel( transport: transport, outbox: store, @@ -1215,17 +959,14 @@ struct ChatViewModelOutboxTests { } @Test func `reconnect retries metadata after a transient session list failure`() async throws { - let databaseDirectory = try makeOutboxDatabaseDirectory() + let (store, _, databaseDirectory) = try makeOutboxStore() defer { try? FileManager.default.removeItem(at: databaseDirectory) } - let store = try makeOutboxStore( - databaseDirectoryURL: databaseDirectory, - gatewayID: "gw-test") #expect(await store.enqueueCommand(outboxTestCommand( id: "retry-metadata", text: "send after retry", createdAt: Date().timeIntervalSince1970))) let transport = OutboxTestTransport(healthy: false) - await transport.state.setSessionListFails(true) + await transport.state.update { $0.sessionListFails = true } let vm = await makeOutboxViewModel(transport: transport, outbox: store) await MainActor.run { vm.load() } @@ -1236,7 +977,7 @@ struct ChatViewModelOutboxTests { } #expect(await store.loadCommands().map(\.status) == [.queued]) - await transport.state.setSessionListFails(false) + await transport.state.update { $0.sessionListFails = false } await transport.goOnline() try await waitUntil("next healthy transition retries metadata and drains") { await store.loadCommands().isEmpty @@ -1245,11 +986,8 @@ struct ChatViewModelOutboxTests { } @Test func `unscoped opaque peer ID preserves case in its durable target`() async throws { - let databaseDirectory = try makeOutboxDatabaseDirectory() + let (store, _, databaseDirectory) = try makeOutboxStore() defer { try? FileManager.default.removeItem(at: databaseDirectory) } - let store = try makeOutboxStore( - databaseDirectoryURL: databaseDirectory, - gatewayID: "gw-test") let sessionKey = "Matrix:Channel:!MixedRoomAbCdEf:example.org" let transport = OutboxTestTransport(healthy: false) let vm = await makeOutboxViewModel( @@ -1271,13 +1009,10 @@ struct ChatViewModelOutboxTests { } @Test func `changed default agent parks command visibly for retry`() async throws { - let databaseDirectory = try makeOutboxDatabaseDirectory() + let (store, _, databaseDirectory) = try makeOutboxStore() defer { try? FileManager.default.removeItem(at: databaseDirectory) } - let store = try makeOutboxStore( - databaseDirectoryURL: databaseDirectory, - gatewayID: "gw-test") let oldTransport = OutboxTestTransport(healthy: false) - await oldTransport.state.setSessionRoutingContract("per-sender|main|agent-a") + await oldTransport.state.update { $0.sessionRoutingContract = "per-sender|main|agent-a" } let oldView = await makeOutboxViewModel( transport: oldTransport, outbox: store, @@ -1287,7 +1022,7 @@ struct ChatViewModelOutboxTests { try await sendWhileOffline(oldView, text: "old main target") let newTransport = OutboxTestTransport(healthy: false) - await newTransport.state.setSessionRoutingContract("per-sender|main|agent-b") + await newTransport.state.update { $0.sessionRoutingContract = "per-sender|main|agent-b" } let newView = await makeOutboxViewModel( transport: newTransport, outbox: store, @@ -1320,17 +1055,14 @@ struct ChatViewModelOutboxTests { } @Test func `atomic gateway routing rejection parks without retrying`() async throws { - let databaseDirectory = try makeOutboxDatabaseDirectory() + let (store, _, databaseDirectory) = try makeOutboxStore() defer { try? FileManager.default.removeItem(at: databaseDirectory) } - let store = try makeOutboxStore( - databaseDirectoryURL: databaseDirectory, - gatewayID: "gw-test") let transport = OutboxTestTransport(healthy: false) let vm = await makeOutboxViewModel(transport: transport, outbox: store) await MainActor.run { vm.load() } try await sendWhileOffline(vm, text: "do not cross config reload") - await transport.state.setSendRoutingChanged(true) + await transport.state.update { $0.sendRoutingChanged = true } await transport.goOnline() try await waitUntil("atomic routing rejection parks") { await store.loadCommands().first?.status == .failed @@ -1342,11 +1074,8 @@ struct ChatViewModelOutboxTests { } @Test func `failed alias row remains reachable after owner change`() async throws { - let databaseDirectory = try makeOutboxDatabaseDirectory() + let (store, _, databaseDirectory) = try makeOutboxStore() defer { try? FileManager.default.removeItem(at: databaseDirectory) } - let store = try makeOutboxStore( - databaseDirectoryURL: databaseDirectory, - gatewayID: "gw-test") #expect(await store.enqueueCommand(OpenClawChatOutboxCommand( id: "c-old-failure", sessionKey: "main", @@ -1378,11 +1107,8 @@ struct ChatViewModelOutboxTests { } @Test func `ownerless global retry stays failed without a selected agent`() async throws { - let databaseDirectory = try makeOutboxDatabaseDirectory() + let (store, _, databaseDirectory) = try makeOutboxStore() defer { try? FileManager.default.removeItem(at: databaseDirectory) } - let store = try makeOutboxStore( - databaseDirectoryURL: databaseDirectory, - gatewayID: "gw-test") #expect(await store.enqueueCommand(OpenClawChatOutboxCommand( id: "c-ownerless", sessionKey: "global", @@ -1416,12 +1142,9 @@ struct ChatViewModelOutboxTests { } @Test func `unavailable recovery keeps the live send FIFO gate closed`() async throws { - let databaseDirectory = try makeOutboxDatabaseDirectory() + let (store, _, databaseDirectory) = try makeOutboxStore() defer { try? FileManager.default.removeItem(at: databaseDirectory) } - let store = try makeOutboxStore( - databaseDirectoryURL: databaseDirectory, - gatewayID: "gw-test") - let outbox = DelayingOutbox(base: store) + let outbox = ScriptedOutbox(base: store) await outbox.setRecoveryAvailable(false) let transport = OutboxTestTransport(healthy: false) let vm = await makeOutboxViewModel(transport: transport, outbox: outbox) @@ -1433,11 +1156,8 @@ struct ChatViewModelOutboxTests { } @Test func `reconnect flushes queued commands in order with their idempotency keys`() async throws { - let databaseDirectory = try makeOutboxDatabaseDirectory() + let (store, _, databaseDirectory) = try makeOutboxStore() defer { try? FileManager.default.removeItem(at: databaseDirectory) } - let store = try makeOutboxStore( - databaseDirectoryURL: databaseDirectory, - gatewayID: "gw-test") let transport = OutboxTestTransport(healthy: false) // One facade routes cache and outbox operations to their separate files. let vm = await makeOutboxViewModel(transport: transport, outbox: store, transcriptCache: store) @@ -1465,7 +1185,7 @@ struct ChatViewModelOutboxTests { await MainActor.run { vm.sessionId == "sess-live" } } #expect(await userTexts(vm) == ["first", "second"]) - #expect(await MainActor.run { queuedStateCount(vm) } == 0) + #expect(await MainActor.run { vm.outboxStatesByMessageID.count } == 0) // The refreshed gateway history is cached for cold offline browsing. let cached = await store.loadTranscript(sessionKey: "main", agentID: "main") @@ -1476,11 +1196,8 @@ struct ChatViewModelOutboxTests { } @Test func `overlapping view models share one atomic FIFO sender`() async throws { - let databaseDirectory = try makeOutboxDatabaseDirectory() + let (store, _, databaseDirectory) = try makeOutboxStore() defer { try? FileManager.default.removeItem(at: databaseDirectory) } - let store = try makeOutboxStore( - databaseDirectoryURL: databaseDirectory, - gatewayID: "gw-test") let now = Date().timeIntervalSince1970 #expect(await store.enqueueCommand(outboxTestCommand(id: "c-1", text: "first", createdAt: now))) #expect(await store.enqueueCommand(outboxTestCommand(id: "c-2", text: "second", createdAt: now + 1))) @@ -1488,7 +1205,7 @@ struct ChatViewModelOutboxTests { let firstVM = await makeOutboxViewModel(transport: transport, outbox: store) let secondVM = await makeOutboxViewModel(transport: transport, outbox: store) - await transport.state.setHealthy(true) + await transport.state.update { $0.healthy = true } await MainActor.run { firstVM.load() secondVM.load() @@ -1504,11 +1221,8 @@ struct ChatViewModelOutboxTests { } @Test func `assistant reply for a flushed run lands via the external-run final event`() async throws { - let databaseDirectory = try makeOutboxDatabaseDirectory() + let (store, _, databaseDirectory) = try makeOutboxStore() defer { try? FileManager.default.removeItem(at: databaseDirectory) } - let store = try makeOutboxStore( - databaseDirectoryURL: databaseDirectory, - gatewayID: "gw-test") let transport = OutboxTestTransport(healthy: false) let vm = await makeOutboxViewModel(transport: transport, outbox: store) @@ -1524,7 +1238,7 @@ struct ChatViewModelOutboxTests { // satisfied by the event path, not a lucky history refresh. (The // scripted history never contains assistant rows, so leaving it on // would wipe the appended final with an incomplete snapshot.) - await transport.state.setHistoryFails(true) + await transport.state.update { $0.historyFails = true } // Flushed runs are intentionally not in pendingRuns; the reply is // delivered through the session-scoped external-run final branch. @@ -1551,11 +1265,8 @@ struct ChatViewModelOutboxTests { } @Test func `acknowledged turn stays in client state until history confirms it`() async throws { - let databaseDirectory = try makeOutboxDatabaseDirectory() + let (store, _, databaseDirectory) = try makeOutboxStore() defer { try? FileManager.default.removeItem(at: databaseDirectory) } - let store = try makeOutboxStore( - databaseDirectoryURL: databaseDirectory, - gatewayID: "gw-test") let transport = OutboxTestTransport(healthy: false) let vm = await makeOutboxViewModel(transport: transport, outbox: store, transcriptCache: store) @@ -1564,7 +1275,7 @@ struct ChatViewModelOutboxTests { // chat.send ACKs before user-turn persistence. With history still // unreachable, the row must remain durable and non-replayable. - await transport.state.setHistoryFails(true) + await transport.state.update { $0.historyFails = true } await transport.goOnline() try await waitUntil("acknowledgement awaits history") { await store.loadCommands().map(\.status) == [.awaitingConfirmation] @@ -1577,7 +1288,7 @@ struct ChatViewModelOutboxTests { vm.messages.contains { vm.outboxState(for: $0.id) == .confirming } }) - await transport.state.setHistoryFails(false) + await transport.state.update { $0.historyFails = false } await MainActor.run { vm.refresh() } try await waitUntil("canonical history confirms send") { await store.loadCommands().isEmpty @@ -1585,17 +1296,17 @@ struct ChatViewModelOutboxTests { } @Test func `healthy restore reconciles a previously acknowledged turn`() async throws { - let databaseDirectory = try makeOutboxDatabaseDirectory() + let (store, _, databaseDirectory) = try makeOutboxStore() defer { try? FileManager.default.removeItem(at: databaseDirectory) } - let store = try makeOutboxStore( - databaseDirectoryURL: databaseDirectory, - gatewayID: "gw-test") #expect(await store.enqueueCommand(outboxTestCommand( id: "c-awaiting", text: "already acknowledged", createdAt: Date().timeIntervalSince1970))) - #expect(await store.claimNextCommand()?.id == "c-awaiting") - #expect(await store.markCommandAwaitingConfirmation(id: "c-awaiting") == .updated) + let claimed = await store.claimNextCommand() + #expect(claimed?.id == "c-awaiting") + #expect(await store.markCommandAwaitingConfirmation( + id: "c-awaiting", + attemptVersion: claimed?.attemptVersion ?? 0) == .updated) let transport = OutboxTestTransport(healthy: true) await transport.state.recordSend( @@ -1614,16 +1325,13 @@ struct ChatViewModelOutboxTests { } @Test func `delete race preserves a turn already proven by canonical history`() async throws { - let databaseDirectory = try makeOutboxDatabaseDirectory() + let (store, _, databaseDirectory) = try makeOutboxStore() defer { try? FileManager.default.removeItem(at: databaseDirectory) } - let store = try makeOutboxStore( - databaseDirectoryURL: databaseDirectory, - gatewayID: "gw-test") #expect(await store.enqueueCommand(outboxTestCommand( id: "c-delivered", text: "delivered already", createdAt: Date().timeIntervalSince1970))) - let holdingOutbox = CancellationHoldingOutbox(base: store) + let holdingOutbox = ScriptedOutbox(base: store, forwarding: .holdingCancellation) let transport = OutboxTestTransport(healthy: false) let vm = await makeOutboxViewModel( transport: transport, @@ -1659,11 +1367,8 @@ struct ChatViewModelOutboxTests { } @Test func `offline local slash command keeps its draft and skips transport`() async throws { - let databaseDirectory = try makeOutboxDatabaseDirectory() + let (store, _, databaseDirectory) = try makeOutboxStore() defer { try? FileManager.default.removeItem(at: databaseDirectory) } - let store = try makeOutboxStore( - databaseDirectoryURL: databaseDirectory, - gatewayID: "gw-test") let transport = OutboxTestTransport(healthy: false) let vm = await makeOutboxViewModel(transport: transport, outbox: store) await MainActor.run { @@ -1682,20 +1387,20 @@ struct ChatViewModelOutboxTests { } @Test func `lagging history does not recursively refresh an acknowledged turn`() async throws { - let databaseDirectory = try makeOutboxDatabaseDirectory() + let (store, _, databaseDirectory) = try makeOutboxStore() defer { try? FileManager.default.removeItem(at: databaseDirectory) } - let store = try makeOutboxStore( - databaseDirectoryURL: databaseDirectory, - gatewayID: "gw-test") #expect(await store.enqueueCommand(outboxTestCommand( id: "c-lagging", text: "not persisted yet", createdAt: Date().timeIntervalSince1970))) - #expect(await store.claimNextCommand()?.id == "c-lagging") - #expect(await store.markCommandAwaitingConfirmation(id: "c-lagging") == .updated) + let claimed = await store.claimNextCommand() + #expect(claimed?.id == "c-lagging") + #expect(await store.markCommandAwaitingConfirmation( + id: "c-lagging", + attemptVersion: claimed?.attemptVersion ?? 0) == .updated) let transport = OutboxTestTransport(healthy: true) - await transport.state.setStaleHistoryRows([]) + await transport.state.update { $0.staleHistoryRows = [] } let vm = await makeOutboxViewModel(transport: transport, outbox: store) await MainActor.run { vm.load() } try await waitUntil("bootstrap history settles") { @@ -1715,11 +1420,8 @@ struct ChatViewModelOutboxTests { } @Test func `gateway rejections burn attempts then fail terminally and support tap retry`() async throws { - let databaseDirectory = try makeOutboxDatabaseDirectory() + let (store, _, databaseDirectory) = try makeOutboxStore() defer { try? FileManager.default.removeItem(at: databaseDirectory) } - let store = try makeOutboxStore( - databaseDirectoryURL: databaseDirectory, - gatewayID: "gw-test") let transport = OutboxTestTransport(healthy: false) let vm = await makeOutboxViewModel(transport: transport, outbox: store) @@ -1727,7 +1429,7 @@ struct ChatViewModelOutboxTests { try await sendWhileOffline(vm, text: "doomed") // Gateway is reachable again but rejects the run on every attempt. - await transport.state.setSendRejects(true) + await transport.state.update { $0.sendRejects = true } await transport.goOnline() try await waitUntil("command failed after max attempts") { @@ -1747,7 +1449,7 @@ struct ChatViewModelOutboxTests { // Tap-to-retry resets attempts; with the gateway accepting again the // command now flushes and the row disappears. - await transport.state.setSendRejects(false) + await transport.state.update { $0.sendRejects = false } let failedMessageID = try #require(await MainActor.run { vm.messages.first { vm.outboxState(for: $0.id)?.isFailed == true }?.id }) @@ -1759,11 +1461,8 @@ struct ChatViewModelOutboxTests { } @Test func `unavailable terminal write drops health instead of advancing FIFO`() async throws { - let databaseDirectory = try makeOutboxDatabaseDirectory() + let (store, _, databaseDirectory) = try makeOutboxStore() defer { try? FileManager.default.removeItem(at: databaseDirectory) } - let store = try makeOutboxStore( - databaseDirectoryURL: databaseDirectory, - gatewayID: "gw-test") #expect(await store.updateLastActiveLeafEntryID( "leaf-new", expectedEpoch: 0, @@ -1782,10 +1481,10 @@ struct ChatViewModelOutboxTests { status: .queued, retryCount: OpenClawChatViewModel.maxOutboxSendAttempts - 1, lastError: "rejected"))) - let outbox = DelayingOutbox(base: store) + let outbox = ScriptedOutbox(base: store) await outbox.setTerminalWritesAvailable(false) let transport = OutboxTestTransport(healthy: false) - await transport.state.setSendRejects(true) + await transport.state.update { $0.sendRejects = true } let vm = await makeOutboxViewModel(transport: transport, outbox: outbox) await MainActor.run { vm.load() } @@ -1802,17 +1501,14 @@ struct ChatViewModelOutboxTests { } @Test func `gateway response errors are definitive and burn retry attempts`() async throws { - let databaseDirectory = try makeOutboxDatabaseDirectory() + let (store, _, databaseDirectory) = try makeOutboxStore() defer { try? FileManager.default.removeItem(at: databaseDirectory) } - let store = try makeOutboxStore( - databaseDirectoryURL: databaseDirectory, - gatewayID: "gw-test") let transport = OutboxTestTransport(healthy: false) let vm = await makeOutboxViewModel(transport: transport, outbox: store) await MainActor.run { vm.load() } try await sendWhileOffline(vm, text: "definitively rejected") - await transport.state.setSendResponseErrors(true) + await transport.state.update { $0.sendResponseErrors = true } await transport.goOnline() try await waitUntil("response error exhausts retry budget") { @@ -1824,13 +1520,10 @@ struct ChatViewModelOutboxTests { } @Test func `definitive live-send rejection restores draft without queueing`() async throws { - let databaseDirectory = try makeOutboxDatabaseDirectory() + let (store, _, databaseDirectory) = try makeOutboxStore() defer { try? FileManager.default.removeItem(at: databaseDirectory) } - let store = try makeOutboxStore( - databaseDirectoryURL: databaseDirectory, - gatewayID: "gw-test") let transport = OutboxTestTransport(healthy: true) - await transport.state.setSendResponseErrors(true) + await transport.state.update { $0.sendResponseErrors = true } let vm = await makeOutboxViewModel(transport: transport, outbox: store) await MainActor.run { vm.load() } @@ -1852,11 +1545,8 @@ struct ChatViewModelOutboxTests { } @Test func `ambiguous live transport failure requires explicit retry`() async throws { - let databaseDirectory = try makeOutboxDatabaseDirectory() + let (store, _, databaseDirectory) = try makeOutboxStore() defer { try? FileManager.default.removeItem(at: databaseDirectory) } - let store = try makeOutboxStore( - databaseDirectoryURL: databaseDirectory, - gatewayID: "gw-test") // Health reads true, but the actual send path is down: the send gate // is bypassed and the transport error must preserve instead of losing // the optimistic turn. @@ -1895,7 +1585,7 @@ struct ChatViewModelOutboxTests { // Connectivity recovery only reconciles history; it cannot replay an // unproven send. A user retry creates the new delivery intent. - await transport.state.setSendFails(false) + await transport.state.update { $0.sendFails = false } await transport.goOnline() try await Task.sleep(nanoseconds: 50_000_000) #expect(await store.loadCommands().map(\.status) == [.failed]) @@ -1910,17 +1600,14 @@ struct ChatViewModelOutboxTests { } @Test func `lost queued send ack reconciles history without replay`() async throws { - let databaseDirectory = try makeOutboxDatabaseDirectory() + let (store, _, databaseDirectory) = try makeOutboxStore() defer { try? FileManager.default.removeItem(at: databaseDirectory) } - let store = try makeOutboxStore( - databaseDirectoryURL: databaseDirectory, - gatewayID: "gw-test") let transport = OutboxTestTransport(healthy: false) let vm = await makeOutboxViewModel(transport: transport, outbox: store) await MainActor.run { vm.load() } try await sendWhileOffline(vm, text: "accepted before disconnect") - await transport.state.setSendFailsAfterRecording(true) + await transport.state.update { $0.sendFailsAfterRecording = true } await transport.goOnline() try await waitUntil("gateway accepted before ack loss") { @@ -1934,11 +1621,8 @@ struct ChatViewModelOutboxTests { } @Test func `tap retry refreshes createdAt so an expired command can resend`() async throws { - let databaseDirectory = try makeOutboxDatabaseDirectory() + let (store, _, databaseDirectory) = try makeOutboxStore() defer { try? FileManager.default.removeItem(at: databaseDirectory) } - let store = try makeOutboxStore( - databaseDirectoryURL: databaseDirectory, - gatewayID: "gw-test") // A command that sat offline past the staleness bound. let staleCreatedAt = Date().timeIntervalSince1970 - OpenClawChatSQLiteTranscriptCache.outboxCommandMaxAge - 60 @@ -1984,11 +1668,8 @@ struct ChatViewModelOutboxTests { } @Test func `flush gates captured thinking using the queued session metadata`() async throws { - let databaseDirectory = try makeOutboxDatabaseDirectory() + let (store, _, databaseDirectory) = try makeOutboxStore() defer { try? FileManager.default.removeItem(at: databaseDirectory) } - let store = try makeOutboxStore( - databaseDirectoryURL: databaseDirectory, - gatewayID: "gw-test") let now = Date().timeIntervalSince1970 #expect(await store.enqueueCommand( OpenClawChatOutboxCommand( @@ -2036,11 +1717,8 @@ struct ChatViewModelOutboxTests { } @Test func `gateway history caches a flushed background-session turn`() async throws { - let databaseDirectory = try makeOutboxDatabaseDirectory() + let (store, _, databaseDirectory) = try makeOutboxStore() defer { try? FileManager.default.removeItem(at: databaseDirectory) } - let store = try makeOutboxStore( - databaseDirectoryURL: databaseDirectory, - gatewayID: "gw-test") // Queued for a session the user is no longer viewing. #expect(await store.enqueueCommand( OpenClawChatOutboxCommand( @@ -2070,11 +1748,8 @@ struct ChatViewModelOutboxTests { } @Test func `background canonical alias event confirms by idempotency key`() async throws { - let databaseDirectory = try makeOutboxDatabaseDirectory() + let (store, _, databaseDirectory) = try makeOutboxStore() defer { try? FileManager.default.removeItem(at: databaseDirectory) } - let store = try makeOutboxStore( - databaseDirectoryURL: databaseDirectory, - gatewayID: "gw-test") #expect(await store.enqueueCommand(OpenClawChatOutboxCommand( id: "c-alias", sessionKey: "main", @@ -2087,8 +1762,11 @@ struct ChatViewModelOutboxTests { status: .queued, retryCount: 0, lastError: nil))) - #expect(await store.claimNextCommand()?.id == "c-alias") - #expect(await store.markCommandAwaitingConfirmation(id: "c-alias") == .updated) + let claimed = await store.claimNextCommand() + #expect(claimed?.id == "c-alias") + #expect(await store.markCommandAwaitingConfirmation( + id: "c-alias", + attemptVersion: claimed?.attemptVersion ?? 0) == .updated) let transport = OutboxTestTransport(healthy: false) let vm = await makeOutboxViewModel(transport: transport, outbox: store, transcriptCache: store) await MainActor.run { @@ -2121,11 +1799,8 @@ struct ChatViewModelOutboxTests { } @Test func `full queue refuses enqueue and keeps the draft`() async throws { - let databaseDirectory = try makeOutboxDatabaseDirectory() + let (store, _, databaseDirectory) = try makeOutboxStore() defer { try? FileManager.default.removeItem(at: databaseDirectory) } - let store = try makeOutboxStore( - databaseDirectoryURL: databaseDirectory, - gatewayID: "gw-test") for index in 0..