From 4c1c4a4ae8c69d75313db20f94eea2af904efabc Mon Sep 17 00:00:00 2001 From: "Jason (Json)" <263060202+fuller-stack-dev@users.noreply.github.com> Date: Fri, 28 Aug 2026 00:03:49 -0600 Subject: [PATCH] fix(ios): show sidebar sessions on multi-agent gateways (#131356) * fix(ios): scope sidebar sessions to selected agent * fix(ios): fence session roster ownership * fix(ios): partition cached sessions by agent * refactor(chat): remove unused session list overload * fix(macos): scope session rosters to selected agent * test(ios): guard canonical activity roster owner * fix(chat): fail closed for legacy scoped caches --- .../Chat/IOSGatewayChatTransport.swift | 3 +- .../ios/Sources/Design/CommandCenterTab.swift | 6 +- .../Sources/Design/IPadActivityScreen.swift | 19 ++- apps/ios/Sources/Model/NodeAppModel.swift | 70 ++++++++--- apps/ios/Sources/RootSidebarModel.swift | 27 +++-- .../GatewayConnectionControllerTests.swift | 32 ++++- apps/ios/Tests/NodeAppModelInvokeTests.swift | 8 +- apps/ios/Tests/RootTabsSourceGuardTests.swift | 36 +++++- .../Sources/OpenClaw/WebChatSwiftUI.swift | 22 +++- .../MacGatewayChatTransportMappingTests.swift | 15 +++ .../OpenClawChatUI/ChatGatewayRequest.swift | 4 + .../ChatSessionSidebarModel.swift | 10 +- .../OpenClawChatUI/ChatTranscriptCache.swift | 108 +++++++++++++++-- .../ChatTranscriptCacheContracts.swift | 18 ++- .../OpenClawChatUI/ChatTransport.swift | 10 +- .../ChatViewModel+TranscriptCache.swift | 14 ++- .../ChatViewModel+TransportEvents.swift | 4 +- .../OpenClawChatUI/ChatViewModel.swift | 2 +- .../OpenClawChatUI/ClientDatabases.swift | 34 ++++++ .../ChatGatewayRequestTests.swift | 11 +- .../ChatSessionSidebarModelTests.swift | 33 ++++++ .../ChatTranscriptCacheStoreTests.swift | 109 +++++++++++++++++- .../ChatViewModelTranscriptCacheTests.swift | 104 +++++++++++++++-- 23 files changed, 613 insertions(+), 86 deletions(-) diff --git a/apps/ios/Sources/Chat/IOSGatewayChatTransport.swift b/apps/ios/Sources/Chat/IOSGatewayChatTransport.swift index d66dcddb13a4..bed7db83eb33 100644 --- a/apps/ios/Sources/Chat/IOSGatewayChatTransport.swift +++ b/apps/ios/Sources/Chat/IOSGatewayChatTransport.swift @@ -328,7 +328,8 @@ struct IOSGatewayChatTransport: OpenClawChatTransport { let request = OpenClawChatGatewayRequests.sessionsList( limit: limit, search: search, - archived: archived) + archived: archived, + agentID: self.globalAgentId) let res = try await gateway.request(request) return try JSONDecoder().decode(OpenClawChatSessionsListResponse.self, from: res) } diff --git a/apps/ios/Sources/Design/CommandCenterTab.swift b/apps/ios/Sources/Design/CommandCenterTab.swift index 42895b0e1af4..b8c5aac7e684 100644 --- a/apps/ios/Sources/Design/CommandCenterTab.swift +++ b/apps/ios/Sources/Design/CommandCenterTab.swift @@ -1276,6 +1276,8 @@ struct CommandSessionsScreen: View { // New Group editor) alongside the fresh session list. self.knownGroups = SessionGroupStore.load() let requestsArchived = self.showArchived + let sourceGatewayID = self.appModel.chatTranscriptCacheGatewayID + let sourceAgentID = self.appModel.chatDeliveryAgentId self.isLoading = true self.loadErrorText = nil defer { self.isLoading = false } @@ -1288,7 +1290,9 @@ struct CommandSessionsScreen: View { self.sessions = roster.sessions } catch { guard requestsArchived == self.showArchived else { return } - self.sessions = requestsArchived ? [] : await self.appModel.loadCachedChatSessions() + self.sessions = requestsArchived ? [] : await self.appModel.loadCachedChatSessions( + gatewayID: sourceGatewayID, + agentID: sourceAgentID) self.loadErrorText = self.sessions.isEmpty ? "Try again after the gateway reconnects." : nil } } diff --git a/apps/ios/Sources/Design/IPadActivityScreen.swift b/apps/ios/Sources/Design/IPadActivityScreen.swift index 7396af13aa27..3f3e3abcc1ce 100644 --- a/apps/ios/Sources/Design/IPadActivityScreen.swift +++ b/apps/ios/Sources/Design/IPadActivityScreen.swift @@ -208,23 +208,20 @@ struct IPadActivityScreen: View { private func refreshSessions() async { guard self.scenePhase == .active else { return } - guard self.sessionsAvailable else { - self.sessions = await self.appModel.loadCachedChatSessions() - self.loadErrorText = nil - return - } - self.isLoading = true self.loadErrorText = nil defer { self.isLoading = false } do { - let transport = self.appModel.makeChatTransport() - let response = try await transport.listSessions(limit: CommandCenterTab.recentSessionsFetchLimit) - self.sessions = response.sessions - await self.appModel.storeCachedChatSessions(response.sessions) + let roster = try await self.appModel.loadChatSessionRoster( + limit: CommandCenterTab.recentSessionsFetchLimit) + self.sessions = roster.sessions } catch { - self.sessions = await self.appModel.loadCachedChatSessions() + let sourceGatewayID = self.appModel.chatTranscriptCacheGatewayID + let sourceAgentID = self.appModel.chatDeliveryAgentId + self.sessions = await self.appModel.loadCachedChatSessions( + gatewayID: sourceGatewayID, + agentID: sourceAgentID) self.loadErrorText = self.sessions.isEmpty ? "Try again after the gateway reconnects." : nil } } diff --git a/apps/ios/Sources/Model/NodeAppModel.swift b/apps/ios/Sources/Model/NodeAppModel.swift index e8054e14783b..3201436a1ed6 100644 --- a/apps/ios/Sources/Model/NodeAppModel.swift +++ b/apps/ios/Sources/Model/NodeAppModel.swift @@ -650,13 +650,12 @@ final class NodeAppModel { return stableID } - /// Recreation key for the chat view model. Includes the cache gateway - /// identity: switching paired gateways while the transport mode stays - /// "operator" must rebuild the view model so transcripts are never read - /// from or written under another gateway's cache scope. + /// Session-list refresh identity. Gateway and agent changes must restart + /// requests so no roster or cached projection crosses either owner. var chatViewModelIdentityID: String { let gatewayID = self.chatTranscriptCacheGatewayIdentityComponent - return "\(self.chatTransportModeID)|\(gatewayID)|\(self.chatTranscriptCacheGeneration)" + let agentID = self.chatDeliveryAgentId ?? "" + return "\(self.chatTransportModeID)|\(gatewayID)|\(agentID)|\(self.chatTranscriptCacheGeneration)" } /// Stable owner key for the long-lived chat view model. Connectivity still @@ -732,14 +731,37 @@ final class NodeAppModel { self.synchronizeTalkSessionKey() } - func loadCachedChatSessions() async -> [OpenClawChatSessionEntry] { - guard let cache = self.makeChatOfflineStore() else { return [] } - return await cache.loadSessions() + func loadCachedChatSessions( + gatewayID: String?, + agentID: String?) async -> [OpenClawChatSessionEntry] + { + guard GatewayStableIdentifier.matches(self.chatTranscriptCacheGatewayID, gatewayID), + let cache = self.makeChatOfflineStore(), + GatewayStableIdentifier.matches(cache.gatewayID, gatewayID) + else { return [] } + let sessions = await cache.loadSessions(agentID: agentID) + guard GatewayStableIdentifier.matches(self.chatTranscriptCacheGatewayID, gatewayID), + self.chatDeliveryAgentId == agentID + else { return [] } + return sessions.filter { + ChatSessionSidebarModel.isSessionInActiveAgentScope( + key: $0.key, + agentID: $0.agentId, + activeAgentID: agentID) + } } - func storeCachedChatSessions(_ sessions: [OpenClawChatSessionEntry]) async { - guard let cache = self.makeChatOfflineStore() else { return } - await cache.storeSessions(sessions) + func storeCachedChatSessions( + _ sessions: [OpenClawChatSessionEntry], + gatewayID: String?, + agentID: String?) async + { + guard GatewayStableIdentifier.matches(self.chatTranscriptCacheGatewayID, gatewayID), + self.chatDeliveryAgentId == agentID, + let cache = self.makeChatOfflineStore(), + GatewayStableIdentifier.matches(cache.gatewayID, gatewayID) + else { return } + await cache.storeSessions(sessions, agentID: agentID) } /// Delete one forgotten gateway's cache and durable state, or both @@ -1691,6 +1713,8 @@ final class NodeAppModel { } if selectedAgentChanged { self.focusedChatSessionKey = nil + self.shareDeliveryChannel = nil + self.shareDeliveryTo = nil } self.synchronizeTalkSessionKey() if let relay = ShareGatewayRelaySettings.loadConfig() { @@ -1704,6 +1728,13 @@ final class NodeAppModel { deliveryChannel: self.shareDeliveryChannel, deliveryTo: self.shareDeliveryTo)) } + if selectedAgentChanged { + // Delivery metadata belongs to the selected agent. Rehydrate it + // after the selection commit; request-time identity fences stale replies. + Task { [weak self] in + await self?.refreshShareRouteFromGateway() + } + } } func setGlobalWakeWords(_ words: [String]) async { @@ -5193,23 +5224,32 @@ extension NodeAppModel { } do { + guard let sourceGatewayID = self.chatTranscriptCacheGatewayID, + let sourceRoute = await operatorGateway.currentRoute(ifGatewayID: sourceGatewayID) + else { return } + let sourceAgentID = self.chatDeliveryAgentId + let sourceMainSessionKey = self.mainSessionKey let request = OpenClawChatGatewayRequests.sessionsList( limit: 80, search: nil, archived: false, + agentID: sourceAgentID, timeoutMs: 10000) - let response = try await operatorGateway.request(request) + let response = try await operatorGateway.request(request, ifCurrentRoute: sourceRoute) let decoded = try JSONDecoder().decode(SessionsListResult.self, from: response) - let currentKey = self.mainSessionKey let sorted = decoded.sessions.sorted { ($0.updatedAt ?? 0) > ($1.updatedAt ?? 0) } let exactMatch = sorted.first { row in - row.key == currentKey && normalize(row.lastChannel) != nil && normalize(row.lastTo) != nil + row.key == sourceMainSessionKey && normalize(row.lastChannel) != nil && normalize(row.lastTo) != nil } let selected = exactMatch let channel = normalize(selected?.lastChannel) let to = normalize(selected?.lastTo) - guard shouldApply() else { return } + guard shouldApply(), + GatewayStableIdentifier.matches(self.chatTranscriptCacheGatewayID, sourceGatewayID), + self.chatDeliveryAgentId == sourceAgentID, + self.mainSessionKey == sourceMainSessionKey + else { return } await MainActor.run { self.shareDeliveryChannel = channel self.shareDeliveryTo = to diff --git a/apps/ios/Sources/RootSidebarModel.swift b/apps/ios/Sources/RootSidebarModel.swift index c61aa984921d..26cf39d9a9ad 100644 --- a/apps/ios/Sources/RootSidebarModel.swift +++ b/apps/ios/Sources/RootSidebarModel.swift @@ -94,16 +94,19 @@ extension NodeAppModel { archived: Bool = false, allowCachedFallback: Bool = true) async throws -> ChatSessionRosterSnapshot { + let sourceGatewayID = self.chatTranscriptCacheGatewayID + let sourceAgentID = self.chatDeliveryAgentId guard self.isLocalChatFixtureEnabled || self.isOperatorGatewayConnected else { guard allowCachedFallback else { throw URLError(.notConnectedToInternet) } return await ChatSessionRosterSnapshot( - sessions: archived ? [] : self.loadCachedChatSessions(), + sessions: archived ? [] : self.loadCachedChatSessions( + gatewayID: sourceGatewayID, + agentID: sourceAgentID), isCached: true, isComplete: false) } do { - let sourceGatewayID = self.chatTranscriptCacheGatewayID let snapshot: ChatSessionRosterSnapshot if self.isLocalChatFixtureEnabled { let response = try await self.makeChatTransport().listSessions(limit: limit, archived: archived) @@ -124,11 +127,14 @@ extension NodeAppModel { limit: limit, search: nil, archived: archived, + agentID: sourceAgentID, offset: offset) let data = try await self.operatorSession.request(request, ifCurrentRoute: route) return try JSONDecoder().decode(OpenClawChatSessionsListResponse.self, from: data) } - guard GatewayStableIdentifier.matches(self.chatTranscriptCacheGatewayID, sourceGatewayID) else { + guard GatewayStableIdentifier.matches(self.chatTranscriptCacheGatewayID, sourceGatewayID), + self.chatDeliveryAgentId == sourceAgentID + else { throw CancellationError() } } @@ -136,9 +142,13 @@ extension NodeAppModel { if !archived { // An interrupted page must not replace a more complete offline roster. if snapshot.isComplete { - await self.storeCachedChatSessions(snapshot.sessions) + await self.storeCachedChatSessions( + snapshot.sessions, + gatewayID: sourceGatewayID, + agentID: sourceAgentID) if let sourceGatewayID, - !GatewayStableIdentifier.matches(self.chatTranscriptCacheGatewayID, sourceGatewayID) + !GatewayStableIdentifier.matches(self.chatTranscriptCacheGatewayID, sourceGatewayID) || + self.chatDeliveryAgentId != sourceAgentID { throw CancellationError() } @@ -149,7 +159,9 @@ extension NodeAppModel { throw CancellationError() } catch { guard allowCachedFallback, !archived else { throw error } - let cached = await self.loadCachedChatSessions() + let cached = await self.loadCachedChatSessions( + gatewayID: sourceGatewayID, + agentID: sourceAgentID) guard !cached.isEmpty else { throw error } return ChatSessionRosterSnapshot(sessions: cached, isCached: true, isComplete: false) } @@ -460,7 +472,8 @@ final class RootSidebarModel { case let .sessionObserver(digest): self.sessions = ChatSessionSidebarModel.applying( observerDigest: digest, - to: self.sessions) + to: self.sessions, + activeAgentId: appModel.chatDeliveryAgentId) case .seqGap: await self.refreshSessions(appModel: appModel) return true diff --git a/apps/ios/Tests/GatewayConnectionControllerTests.swift b/apps/ios/Tests/GatewayConnectionControllerTests.swift index 233ef2f59677..c5878c054ba0 100644 --- a/apps/ios/Tests/GatewayConnectionControllerTests.swift +++ b/apps/ios/Tests/GatewayConnectionControllerTests.swift @@ -2043,7 +2043,7 @@ private func waitUntil( token: nil, password: nil, sessionKey: "main") - defaults.set(try JSONEncoder().encode(otherRelay), forKey: "share.gatewayRelay.config.v1") + try defaults.set(JSONEncoder().encode(otherRelay), forKey: "share.gatewayRelay.config.v1") let mismatched = try #require(ShareGatewayRelaySettings.loadConfig()) #expect(mismatched.token == nil) #expect(mismatched.password == nil) @@ -2638,7 +2638,7 @@ private func waitUntil( useTLS: false, lastConnectedAtMs: nil)) let appModel = NodeAppModel() - let session = OpenClawChatSessionEntry( + var session = OpenClawChatSessionEntry( key: "agent:main:a", kind: nil, displayName: "Gateway A session", @@ -2658,12 +2658,34 @@ private func waitUntil( modelProvider: nil, model: nil, contextTokens: nil) + session.agentId = "main" + var matchingBare = session + matchingBare.key = "shared-tool" + var ownerlessPrefixed = session + ownerlessPrefixed.key = "agent:main:legacy" + ownerlessPrefixed.agentId = nil + appModel.gatewayDefaultAgentId = "main" - await appModel.storeCachedChatSessions([session]) + await appModel.storeCachedChatSessions( + [session, matchingBare, ownerlessPrefixed], + gatewayID: gatewayA, + agentID: "main") + var workGlobal = session + workGlobal.key = "global" + workGlobal.agentId = "work" + appModel.selectedAgentId = "work" + await appModel.storeCachedChatSessions([workGlobal], gatewayID: gatewayA, agentID: "work") _ = GatewaySettingsStore.setActiveGateway(stableID: gatewayB) - #expect(await appModel.loadCachedChatSessions().isEmpty) + #expect(await appModel.loadCachedChatSessions(gatewayID: gatewayB, agentID: "work").isEmpty) _ = GatewaySettingsStore.setActiveGateway(stableID: gatewayA) - #expect(await appModel.loadCachedChatSessions() == [session]) + appModel.selectedAgentId = "main" + #expect(await appModel.loadCachedChatSessions(gatewayID: gatewayA, agentID: "main") == [ + session, + matchingBare, + ownerlessPrefixed, + ]) + appModel.selectedAgentId = "work" + #expect(await appModel.loadCachedChatSessions(gatewayID: gatewayA, agentID: "work") == [workGlobal]) } private static func makeNodeOptions( diff --git a/apps/ios/Tests/NodeAppModelInvokeTests.swift b/apps/ios/Tests/NodeAppModelInvokeTests.swift index f9ce763ab48c..15e5bdb5fc54 100644 --- a/apps/ios/Tests/NodeAppModelInvokeTests.swift +++ b/apps/ios/Tests/NodeAppModelInvokeTests.swift @@ -1167,18 +1167,24 @@ private final class TimingOutDeviceStatusService: DeviceStatusServicing { #expect(appModel.chatDeliveryAgentId == nil) } - @Test @MainActor func `chat delivery owner requires persisted or gateway ownership`() { + @Test @MainActor func `chat delivery owner and refresh identity follow gateway ownership`() { let appModel = NodeAppModel() + let ownerlessIdentity = appModel.chatViewModelIdentityID #expect(appModel.chatDeliveryAgentId == nil) appModel.gatewayDefaultAgentId = " Agent-A " + let defaultIdentity = appModel.chatViewModelIdentityID #expect(appModel.chatDeliveryAgentId == "agent-a") + #expect(defaultIdentity != ownerlessIdentity) appModel.setSelectedAgentId(" Agent-B ") + let selectedIdentity = appModel.chatViewModelIdentityID #expect(appModel.chatDeliveryAgentId == "agent-b") + #expect(selectedIdentity != defaultIdentity) appModel.openChat(sessionKey: "agent:Agent-C:incident") #expect(appModel.chatDeliveryAgentId == "agent-c") + #expect(appModel.chatViewModelIdentityID != selectedIdentity) } @Test @MainActor func `init preserves saved talk mode preference`() { diff --git a/apps/ios/Tests/RootTabsSourceGuardTests.swift b/apps/ios/Tests/RootTabsSourceGuardTests.swift index c1ceec9c18c7..c44172c1faf3 100644 --- a/apps/ios/Tests/RootTabsSourceGuardTests.swift +++ b/apps/ios/Tests/RootTabsSourceGuardTests.swift @@ -103,6 +103,39 @@ struct RootTabsSourceGuardTests { #expect(refreshID.contains("self.scenePhase == .active")) } + @Test func `session roster lifecycle stays on one gateway and agent owner`() throws { + let sidebarSource = try String(contentsOf: Self.rootSidebarModelSourceURL(), encoding: .utf8) + let appModelSource = try String(contentsOf: Self.nodeAppModelSourceURL(), encoding: .utf8) + let rosterLoad = try Self.extract( + sidebarSource, + from: "func loadChatSessionRoster(", + to: "@MainActor\n@Observable") + let observerApply = try Self.extract( + sidebarSource, + from: "private func handleSessionEvent(", + to: "func reportSessionError(") + let refreshIdentity = try Self.extract( + appModelSource, + from: "var chatViewModelIdentityID: String", + to: "var chatViewModelOwnerID: String") + let shareRoute = try Self.extract( + appModelSource, + from: "private func refreshShareRouteFromGateway(", + to: "func runSharePipelineSelfTest()") + + #expect(rosterLoad.contains("let sourceAgentID = self.chatDeliveryAgentId")) + #expect(rosterLoad.contains("agentID: sourceAgentID")) + #expect(rosterLoad.contains("self.chatDeliveryAgentId == sourceAgentID")) + #expect(rosterLoad.contains("self.chatDeliveryAgentId != sourceAgentID")) + #expect(observerApply.contains("activeAgentId: appModel.chatDeliveryAgentId")) + #expect(refreshIdentity.contains("self.chatDeliveryAgentId")) + #expect(shareRoute.contains("let sourceAgentID = self.chatDeliveryAgentId")) + #expect(shareRoute.contains("let sourceMainSessionKey = self.mainSessionKey")) + #expect(shareRoute.contains("ifCurrentRoute: sourceRoute")) + #expect(shareRoute.contains("self.chatDeliveryAgentId == sourceAgentID")) + #expect(shareRoute.contains("self.mainSessionKey == sourceMainSessionKey")) + } + @Test func `sidebar dashboard keeps per field last known good values and drains cron pages`() throws { let source = try String(contentsOf: Self.rootSidebarModelSourceURL(), encoding: .utf8) let dashboardCommit = try Self.extract( @@ -289,7 +322,8 @@ struct RootTabsSourceGuardTests { let projectSource = try String(contentsOf: Self.xcodeProjectSourceURL(), encoding: .utf8) #expect(activitySource.contains("struct IPadActivityScreen: View")) - #expect(activitySource.contains("self.appModel.makeChatTransport()")) + #expect(activitySource.contains("self.appModel.loadChatSessionRoster(")) + #expect(!activitySource.contains("self.appModel.makeChatTransport()")) #expect(appModelSource.contains("return IOSGatewayChatTransport(")) #expect(appModelSource.contains("globalAgentId: self.chatDeliveryAgentId")) #expect(!appModelSource.contains("defaultAgentId: self.gatewayDefaultAgentId")) diff --git a/apps/macos/Sources/OpenClaw/WebChatSwiftUI.swift b/apps/macos/Sources/OpenClaw/WebChatSwiftUI.swift index 40e9bfdfb1f5..2004b816e4ea 100644 --- a/apps/macos/Sources/OpenClaw/WebChatSwiftUI.swift +++ b/apps/macos/Sources/OpenClaw/WebChatSwiftUI.swift @@ -277,7 +277,7 @@ struct MacGatewayChatTransport: OpenClawChatTransport { search: String?, archived: Bool) async throws -> OpenClawChatSessionsListResponse { - let request = OpenClawChatGatewayRequests.sessionsList( + let request = self.sessionsListRequest( limit: limit, search: search, archived: archived) @@ -309,6 +309,18 @@ struct MacGatewayChatTransport: OpenClawChatTransport { sessions: decoded.sessions) } + func sessionsListRequest( + limit: Int?, + search: String?, + archived: Bool) -> OpenClawChatGatewayRequest + { + OpenClawChatGatewayRequests.sessionsList( + limit: limit, + search: search, + archived: archived, + agentID: self.routingIdentity.currentAgentID()) + } + func listChildSessions(parentKey: String) async throws -> [OpenClawChatSessionEntry] { try await self.listChildSessions(parentKey: parentKey, serverLease: nil) } @@ -1149,6 +1161,11 @@ final class WebChatSwiftUIWindowController: NSObject, NSWindowDelegate { cachedDefaultAgentID: routingIdentity.defaultAgentID) (transport as? MacGatewayChatTransport)? .updateDefaultGlobalAgentID(effectiveAgentID) + // Keep request and cache ownership in lockstep before the + // persistence await can admit a roster refresh. + vm.syncDeliveryIdentity( + activeAgentId: effectiveAgentID, + sessionRoutingContract: routingIdentity.contract) if let store = transcriptCache as? OpenClawChatSQLiteTranscriptCache, !usesPrimaryAppRuntime || store.gatewayID == MacChatTranscriptCache.currentGatewayID(), let persistedIdentity = OpenClawChatSessionRoutingIdentity( @@ -1156,9 +1173,6 @@ final class WebChatSwiftUIWindowController: NSObject, NSWindowDelegate { { await store.storeSessionRoutingIdentity(persistedIdentity) } - vm.syncDeliveryIdentity( - activeAgentId: effectiveAgentID, - sessionRoutingContract: routingIdentity.contract) } } } diff --git a/apps/macos/Tests/OpenClawIPCTests/MacGatewayChatTransportMappingTests.swift b/apps/macos/Tests/OpenClawIPCTests/MacGatewayChatTransportMappingTests.swift index 3eced12e03e3..32507fda6629 100644 --- a/apps/macos/Tests/OpenClawIPCTests/MacGatewayChatTransportMappingTests.swift +++ b/apps/macos/Tests/OpenClawIPCTests/MacGatewayChatTransportMappingTests.swift @@ -42,6 +42,21 @@ struct MacGatewayChatTransportMappingTests { agentID: nil)) } + @Test func `session list request follows the current routing agent`() { + let transport = MacGatewayChatTransport(defaultGlobalAgentID: " Agent-A ") + + let first = transport.sessionsListRequest(limit: 50, search: nil, archived: false) + #expect(first.params["agentId"]?.value as? String == "agent-a") + + transport.updateDefaultGlobalAgentID("Agent-B") + let second = transport.sessionsListRequest(limit: nil, search: "recent", archived: true) + #expect(second.params["agentId"]?.value as? String == "agent-b") + + let unowned = MacGatewayChatTransport() + .sessionsListRequest(limit: nil, search: nil, archived: false) + #expect(unowned.params["agentId"] == nil) + } + @Test func `fixed connection does not inherit app wide cache routing`() async throws { let url = try #require(URL(string: "wss://fixed.example")) let connection = GatewayConnection(configProvider: { diff --git a/apps/shared/OpenClawKit/Sources/OpenClawChatUI/ChatGatewayRequest.swift b/apps/shared/OpenClawKit/Sources/OpenClawChatUI/ChatGatewayRequest.swift index efae00c07f7d..8b73ca7ba4e1 100644 --- a/apps/shared/OpenClawKit/Sources/OpenClawChatUI/ChatGatewayRequest.swift +++ b/apps/shared/OpenClawKit/Sources/OpenClawChatUI/ChatGatewayRequest.swift @@ -192,6 +192,7 @@ public enum OpenClawChatGatewayRequests { limit: Int?, search: String?, archived: Bool, + agentID: String? = nil, includeGlobal: Bool = true, includeUnknown: Bool = false, activeMinutes: Int? = nil, @@ -204,6 +205,9 @@ public enum OpenClawChatGatewayRequests { "includeGlobal": AnyCodable(includeGlobal), "includeUnknown": AnyCodable(includeUnknown), ] + if let agentID = normalized(agentID) { + params["agentId"] = AnyCodable(agentID) + } if let limit { params["limit"] = AnyCodable(limit) } diff --git a/apps/shared/OpenClawKit/Sources/OpenClawChatUI/ChatSessionSidebarModel.swift b/apps/shared/OpenClawKit/Sources/OpenClawChatUI/ChatSessionSidebarModel.swift index b92ff86fd904..b63d77bd5392 100644 --- a/apps/shared/OpenClawKit/Sources/OpenClawChatUI/ChatSessionSidebarModel.swift +++ b/apps/shared/OpenClawKit/Sources/OpenClawChatUI/ChatSessionSidebarModel.swift @@ -502,9 +502,17 @@ public enum ChatSessionSidebarModel { status != "running" } - static func isSessionInActiveAgentScope(key: String, activeAgentID: String?) -> Bool { + public static func isSessionInActiveAgentScope( + key: String, + agentID: String? = nil, + activeAgentID: String?) -> Bool + { let normalizedAgent = activeAgentID?.trimmingCharacters(in: .whitespacesAndNewlines).lowercased() ?? "" guard !normalizedAgent.isEmpty else { return true } + // Gateway row ownership outranks ambiguous bare/global keys. Missing + // metadata is a shipped legacy-cache state and keeps key-only behavior. + let rowAgent = agentID?.trimmingCharacters(in: .whitespacesAndNewlines).lowercased() + if let rowAgent, !rowAgent.isEmpty, rowAgent != normalizedAgent { return false } let parts = key.split(separator: ":", maxSplits: 2, omittingEmptySubsequences: false) guard parts.count == 3, parts[0].lowercased() == "agent" else { return true } return parts[1].lowercased() == normalizedAgent diff --git a/apps/shared/OpenClawKit/Sources/OpenClawChatUI/ChatTranscriptCache.swift b/apps/shared/OpenClawKit/Sources/OpenClawChatUI/ChatTranscriptCache.swift index ad5df747b619..81d44f0f8b06 100644 --- a/apps/shared/OpenClawKit/Sources/OpenClawChatUI/ChatTranscriptCache.swift +++ b/apps/shared/OpenClawKit/Sources/OpenClawChatUI/ChatTranscriptCache.swift @@ -85,6 +85,7 @@ public actor OpenClawChatSQLiteTranscriptCache: OpenClawChatTranscriptCache, OpenClawChatCommandOutbox { public static let maxCachedSessions = 50 + public static let maxCachedSessionOwners = 50 public static let maxCachedTranscripts = 50 public static let maxCachedMessagesPerSession = 200 public static let maxQueuedCommands = 50 @@ -137,17 +138,23 @@ public actor OpenClawChatSQLiteTranscriptCache: OpenClawChatTranscriptCache, // MARK: - Gateway cache public func loadSessions() async -> [OpenClawChatSessionEntry] { + await self.loadSessions(agentID: nil) + } + + public func loadSessions(agentID: String?) async -> [OpenClawChatSessionEntry] { guard !self.isRetired else { return [] } + let normalizedAgentID = Self.normalizedAgentID(agentID) let gatewayID = self.gatewayID do { return try await self.databases.cacheQueue.write { db in + try OpenClawClientDatabases.ensureAgentSessionCacheSchema(db) let rows = try Row.fetchAll( db, sql: """ - SELECT payload_json FROM cached_sessions - WHERE gateway_id = ? ORDER BY position + SELECT payload_json FROM cached_agent_sessions + WHERE gateway_id = ? AND agent_id = ? ORDER BY position """, - arguments: [gatewayID]) + arguments: [gatewayID, normalizedAgentID]) do { return try rows.map { row in let payload: String = row["payload_json"] @@ -156,11 +163,20 @@ public actor OpenClawChatSQLiteTranscriptCache: OpenClawChatTranscriptCache, } catch { cacheLogger.error( "gateway session cache decode failed: \(error.localizedDescription, privacy: .public)") - // Decode and cleanup share one transaction so a newer - // snapshot can never land between the failed read and delete. + // Decode and cleanup share one agent partition transaction; + // corruption in one roster must not erase another agent. try db.execute( - sql: "DELETE FROM cached_sessions WHERE gateway_id = ?", - arguments: [gatewayID]) + sql: """ + DELETE FROM cached_agent_sessions + WHERE gateway_id = ? AND agent_id = ? + """, + arguments: [gatewayID, normalizedAgentID]) + try db.execute( + sql: """ + DELETE FROM cached_session_rosters + WHERE gateway_id = ? AND agent_id = ? + """, + arguments: [gatewayID, normalizedAgentID]) return [] } } @@ -214,30 +230,80 @@ public actor OpenClawChatSQLiteTranscriptCache: OpenClawChatTranscriptCache, } public func storeSessions(_ sessions: [OpenClawChatSessionEntry]) async { + await self.storeSessions(sessions, agentID: nil) + } + + public func storeSessions(_ sessions: [OpenClawChatSessionEntry], agentID: String?) async { guard !self.isRetired else { return } - let bounded = Self.boundedSessions(sessions) + let normalizedAgentID = Self.normalizedAgentID(agentID) + guard let owned = Self.sessionsOwnedBy(sessions, agentID: normalizedAgentID) else { + cacheLogger.error("gateway session cache rejected a mixed-agent snapshot") + return + } + let bounded = Self.boundedSessions(owned) let gatewayID = self.gatewayID do { let encoded = try bounded.map(Self.encodeJSON) try await self.databases.cacheQueue.write { db in + try OpenClawClientDatabases.ensureAgentSessionCacheSchema(db) + // The legacy gateway-wide rows cannot represent agent ownership. + // Keep them empty so a downgraded client cannot paint a mixed roster. try db.execute( sql: "DELETE FROM cached_sessions WHERE gateway_id = ?", arguments: [gatewayID]) + try db.execute( + sql: """ + INSERT INTO cached_session_rosters(gateway_id, agent_id, last_used_at) + VALUES (?, ?, ?) + ON CONFLICT(gateway_id, agent_id) + DO UPDATE SET last_used_at = excluded.last_used_at + """, + arguments: [gatewayID, normalizedAgentID, Date().timeIntervalSince1970]) + try db.execute( + sql: """ + DELETE FROM cached_agent_sessions + WHERE gateway_id = ? AND agent_id = ? + """, + arguments: [gatewayID, normalizedAgentID]) for (position, pair) in zip(bounded, encoded).enumerated() { try db.execute( sql: """ - INSERT INTO cached_sessions( - gateway_id, session_key, position, updated_at, payload_json - ) VALUES (?, ?, ?, ?, ?) + INSERT INTO cached_agent_sessions( + gateway_id, agent_id, session_key, position, updated_at, payload_json + ) VALUES (?, ?, ?, ?, ?, ?) """, arguments: [ gatewayID, + normalizedAgentID, pair.0.key, position, pair.0.updatedAt ?? 0, pair.1, ]) } + let staleOwners = try String.fetchAll( + db, + sql: """ + SELECT agent_id FROM cached_session_rosters + WHERE gateway_id = ? + ORDER BY last_used_at DESC, agent_id + LIMIT -1 OFFSET ? + """, + arguments: [gatewayID, Self.maxCachedSessionOwners]) + for staleOwner in staleOwners { + try db.execute( + sql: """ + DELETE FROM cached_agent_sessions + WHERE gateway_id = ? AND agent_id = ? + """, + arguments: [gatewayID, staleOwner]) + try db.execute( + sql: """ + DELETE FROM cached_session_rosters + WHERE gateway_id = ? AND agent_id = ? + """, + arguments: [gatewayID, staleOwner]) + } } } catch { cacheLogger.error("gateway session cache write failed: \(error.localizedDescription, privacy: .public)") @@ -1489,6 +1555,26 @@ extension OpenClawChatSQLiteTranscriptCache { .prefix(self.maxCachedSessions)) } + private nonisolated static func sessionsOwnedBy( + _ sessions: [OpenClawChatSessionEntry], + agentID: String) -> [OpenClawChatSessionEntry]? + { + guard !agentID.isEmpty else { return sessions } + var owned: [OpenClawChatSessionEntry] = [] + owned.reserveCapacity(sessions.count) + for var session in sessions { + let rowAgentID = self.normalizedAgentID(session.agentId) + let keyAgentID = self.normalizedAgentID(OpenClawChatSessionKey.agentID(from: session.key)) + guard rowAgentID.isEmpty || rowAgentID == agentID, + keyAgentID.isEmpty || keyAgentID == agentID + else { return nil } + // The request owner is authoritative for bare and global keys. + session.agentId = agentID + owned.append(session) + } + return owned + } + private static func attachmentByteCount(_ attachments: [OpenClawChatOutboxAttachment]) -> Int? { var total = 0 for attachment in attachments { diff --git a/apps/shared/OpenClawKit/Sources/OpenClawChatUI/ChatTranscriptCacheContracts.swift b/apps/shared/OpenClawKit/Sources/OpenClawChatUI/ChatTranscriptCacheContracts.swift index f5db8914fe39..e28e2ca63635 100644 --- a/apps/shared/OpenClawKit/Sources/OpenClawChatUI/ChatTranscriptCacheContracts.swift +++ b/apps/shared/OpenClawKit/Sources/OpenClawChatUI/ChatTranscriptCacheContracts.swift @@ -4,13 +4,15 @@ import Foundation /// /// The cache only pre-paints cold opens and covers offline browsing; connected /// reads always come from the gateway and replace cached content wholesale. -/// Implementations must scope every row by gateway identity so one shared -/// installation database can safely serve all paired gateways. +/// Implementations must scope every row by gateway and agent identity so one +/// shared installation database can safely serve all paired gateways. public protocol OpenClawChatTranscriptCache: Sendable { func loadSessions() async -> [OpenClawChatSessionEntry] + func loadSessions(agentID: String?) async -> [OpenClawChatSessionEntry] func loadTranscript(sessionKey: String) async -> [OpenClawChatMessage] func loadTranscript(sessionKey: String, agentID: String?) async -> [OpenClawChatMessage] func storeSessions(_ sessions: [OpenClawChatSessionEntry]) async + func storeSessions(_ sessions: [OpenClawChatSessionEntry], agentID: String?) async /// Canonical gateway rows can prove that an ambiguously delivered local /// command landed after cancellation and must override local suppression. func storeCanonicalTranscript( @@ -24,6 +26,18 @@ public protocol OpenClawChatTranscriptCache: Sendable { } extension OpenClawChatTranscriptCache { + public func loadSessions(agentID: String?) async -> [OpenClawChatSessionEntry] { + // Legacy conformers have no agent partition. Scoped access must fail + // closed or an ownerless roster can cross an agent switch. + guard agentID == nil else { return [] } + return await self.loadSessions() + } + + public func storeSessions(_ sessions: [OpenClawChatSessionEntry], agentID: String?) async { + guard agentID == nil else { return } + await self.storeSessions(sessions) + } + public func loadTranscript(sessionKey: String, agentID: String?) async -> [OpenClawChatMessage] { guard agentID == nil else { return [] } return await self.loadTranscript(sessionKey: sessionKey) diff --git a/apps/shared/OpenClawKit/Sources/OpenClawChatUI/ChatTransport.swift b/apps/shared/OpenClawKit/Sources/OpenClawChatUI/ChatTransport.swift index aaa02ef757dd..b16599509398 100644 --- a/apps/shared/OpenClawKit/Sources/OpenClawChatUI/ChatTransport.swift +++ b/apps/shared/OpenClawKit/Sources/OpenClawChatUI/ChatTransport.swift @@ -1053,14 +1053,10 @@ extension OpenClawChatTransport { [] } - /// Conveniences for callers that only page a list. Transports must + /// Convenience for callers that only select archive state. Transports must /// implement the canonical `listSessions(limit:search:archived:)` - /// requirement; same-name methods on a conformer are shadowed by these - /// sugars and never called through the protocol. - public func listSessions(limit: Int?) async throws -> OpenClawChatSessionsListResponse { - try await self.listSessions(limit: limit, search: nil, archived: false) - } - + /// requirement; same-name methods on a conformer are shadowed by this + /// sugar and never called through the protocol. public func listSessions(limit: Int?, archived: Bool) async throws -> OpenClawChatSessionsListResponse { try await self.listSessions(limit: limit, search: nil, archived: archived) } diff --git a/apps/shared/OpenClawKit/Sources/OpenClawChatUI/ChatViewModel+TranscriptCache.swift b/apps/shared/OpenClawKit/Sources/OpenClawChatUI/ChatViewModel+TranscriptCache.swift index 95b209eccdfe..04ce982976ec 100644 --- a/apps/shared/OpenClawKit/Sources/OpenClawChatUI/ChatViewModel+TranscriptCache.swift +++ b/apps/shared/OpenClawKit/Sources/OpenClawChatUI/ChatViewModel+TranscriptCache.swift @@ -83,13 +83,13 @@ extension OpenClawChatViewModel { } } - func persistSessionsToCache(_ sessions: [OpenClawChatSessionEntry]) { + func persistSessionsToCache(_ sessions: [OpenClawChatSessionEntry], agentID: String?) { guard let transcriptCache else { return } let durableSessions = sessions.map(Self.durableSessionCacheProjection) let previous = pendingCacheWriteTask pendingCacheWriteTask = Task.detached { await previous?.value - await transcriptCache.storeSessions(durableSessions) + await transcriptCache.storeSessions(durableSessions, agentID: agentID) } } @@ -101,7 +101,7 @@ extension OpenClawChatViewModel { guard let transcriptCache else { return } if sessions.isEmpty, !hasAppliedLiveSessions { Task { [weak self] in - let cached = await transcriptCache.loadSessions() + let cached = await transcriptCache.loadSessions(agentID: session.deliveryAgentID) guard let self, !cached.isEmpty else { return } // A live sessions response (even an empty one) is authoritative; // a slow cache read must never repaint over it. @@ -116,8 +116,14 @@ extension OpenClawChatViewModel { } let durableSessions = cached.map(Self.durableSessionCacheProjection) let organized = OpenClawChatSessionListOrganizer.organize(durableSessions) + let agentScoped = organized.filter { + ChatSessionSidebarModel.isSessionInActiveAgentScope( + key: $0.key, + agentID: $0.agentId, + activeAgentID: self.activeAgentId) + } let scoped = ChatSessionSidebarModel.clearingForeignGlobalObserverDigest( - in: organized, + in: agentScoped, activeAgentId: self.activeAgentId) self.sessions = self.applyingLocalUnreadOverrides( to: scoped) diff --git a/apps/shared/OpenClawKit/Sources/OpenClawChatUI/ChatViewModel+TransportEvents.swift b/apps/shared/OpenClawKit/Sources/OpenClawChatUI/ChatViewModel+TransportEvents.swift index 7563816d4783..fb7fd65986bc 100644 --- a/apps/shared/OpenClawKit/Sources/OpenClawChatUI/ChatViewModel+TransportEvents.swift +++ b/apps/shared/OpenClawKit/Sources/OpenClawChatUI/ChatViewModel+TransportEvents.swift @@ -261,7 +261,9 @@ extension OpenClawChatViewModel { activeRunIDs: change.activeRunIds, activeRunIDsPresent: change.activeRunIdsPresent) self.sessions = OpenClawChatSessionListOrganizer.organize(updated) - self.persistSessionsToCache(self.sessions) + self.persistSessionsToCache( + self.sessions, + agentID: self.currentSessionSnapshot().deliveryAgentID) return .merged } diff --git a/apps/shared/OpenClawKit/Sources/OpenClawChatUI/ChatViewModel.swift b/apps/shared/OpenClawKit/Sources/OpenClawChatUI/ChatViewModel.swift index 786181788f98..365ab7fb0a48 100644 --- a/apps/shared/OpenClawKit/Sources/OpenClawChatUI/ChatViewModel.swift +++ b/apps/shared/OpenClawKit/Sources/OpenClawChatUI/ChatViewModel.swift @@ -1126,7 +1126,7 @@ extension OpenClawChatViewModel { self.hasAppliedLiveSessions = true self.syncSelectedModel() syncThinkingLevelOptions() - persistSessionsToCache(organized) + persistSessionsToCache(organized, agentID: session.deliveryAgentID) self.readySessionMetadataGeneration = metadataGeneration if self.healthOK { reconcilePendingOutboxBranchScopes() diff --git a/apps/shared/OpenClawKit/Sources/OpenClawChatUI/ClientDatabases.swift b/apps/shared/OpenClawKit/Sources/OpenClawChatUI/ClientDatabases.swift index 1c44c98892c9..21d449481d11 100644 --- a/apps/shared/OpenClawKit/Sources/OpenClawChatUI/ClientDatabases.swift +++ b/apps/shared/OpenClawKit/Sources/OpenClawChatUI/ClientDatabases.swift @@ -198,6 +198,12 @@ public final class OpenClawClientDatabases: @unchecked Sendable { arguments: [gatewayID]) } try self.cacheQueue.write { db in + try db.execute( + sql: "DELETE FROM cached_agent_sessions WHERE gateway_id = ?", + arguments: [gatewayID]) + try db.execute( + sql: "DELETE FROM cached_session_rosters WHERE gateway_id = ?", + arguments: [gatewayID]) try db.execute(sql: "DELETE FROM cached_sessions WHERE gateway_id = ?", arguments: [gatewayID]) try db.execute(sql: "DELETE FROM cached_transcripts WHERE gateway_id = ?", arguments: [gatewayID]) } @@ -641,8 +647,36 @@ extension OpenClawClientDatabases { INSERT OR REPLACE INTO cache_metadata(id, format_version) VALUES (1, \(self.gatewayCacheFormatVersion)); """) + try self.ensureAgentSessionCacheSchema(db) } } + + /// Session rosters are disposable cache state, so this additive surface is + /// lazily ensured without advancing the cache format or erasing transcripts. + static func ensureAgentSessionCacheSchema(_ db: Database) throws { + try db.execute(sql: """ + CREATE TABLE IF NOT EXISTS cached_session_rosters( + gateway_id TEXT NOT NULL, + agent_id TEXT NOT NULL, + last_used_at REAL NOT NULL, + PRIMARY KEY(gateway_id, agent_id) + ); + CREATE TABLE IF NOT EXISTS cached_agent_sessions( + gateway_id TEXT NOT NULL, + agent_id TEXT NOT NULL, + session_key TEXT NOT NULL, + position INTEGER NOT NULL, + updated_at REAL NOT NULL, + payload_json TEXT NOT NULL, + PRIMARY KEY(gateway_id, agent_id, session_key), + FOREIGN KEY(gateway_id, agent_id) + REFERENCES cached_session_rosters(gateway_id, agent_id) + ON DELETE CASCADE + ); + CREATE INDEX IF NOT EXISTS cached_agent_sessions_order + ON cached_agent_sessions(gateway_id, agent_id, position); + """) + } } extension OpenClawClientDatabases { diff --git a/apps/shared/OpenClawKit/Tests/OpenClawKitTests/ChatGatewayRequestTests.swift b/apps/shared/OpenClawKit/Tests/OpenClawKitTests/ChatGatewayRequestTests.swift index 8f29e89b9eaa..5ec41aca96aa 100644 --- a/apps/shared/OpenClawKit/Tests/OpenClawKitTests/ChatGatewayRequestTests.swift +++ b/apps/shared/OpenClawKit/Tests/OpenClawKitTests/ChatGatewayRequestTests.swift @@ -68,7 +68,8 @@ struct ChatGatewayRequestTests { let request = OpenClawChatGatewayRequests.sessionsList( limit: 12, search: " incident ", - archived: true) + archived: true, + agentID: " Reviewer ") #expect(request.method == "sessions.list") #expect(request.timeoutMs == 15000) @@ -77,6 +78,14 @@ struct ChatGatewayRequestTests { #expect(request.params["limit"]?.value as? Int == 12) #expect(request.params["search"]?.value as? String == "incident") #expect(request.params["archived"]?.value as? Bool == true) + #expect(request.params["agentId"]?.value as? String == "Reviewer") + + let unscoped = OpenClawChatGatewayRequests.sessionsList( + limit: nil, + search: nil, + archived: false, + agentID: " ") + #expect(unscoped.params["agentId"] == nil) } @Test func `child session request encodes focused pagination filters`() { diff --git a/apps/shared/OpenClawKit/Tests/OpenClawKitTests/ChatSessionSidebarModelTests.swift b/apps/shared/OpenClawKit/Tests/OpenClawKitTests/ChatSessionSidebarModelTests.swift index fced60672df8..328c5e39a6b5 100644 --- a/apps/shared/OpenClawKit/Tests/OpenClawKitTests/ChatSessionSidebarModelTests.swift +++ b/apps/shared/OpenClawKit/Tests/OpenClawKitTests/ChatSessionSidebarModelTests.swift @@ -682,6 +682,39 @@ struct ChatSessionSidebarModelTests { #expect(replayed[0].observerDigest?.revision == 4) } + @Test func `global observer events require the active agent owner`() { + let running = self.entry( + key: "global", + status: "running", + hasActiveRun: true, + activeRunIds: ["run-work"]) + let accepted = ChatSessionSidebarModel.applying( + observerDigest: SessionObserverDigest( + sessionkey: "global", + agentid: "work", + runid: "run-work", + revision: 1, + updatedat: 100, + headline: "Work status", + health: .onTrack), + to: [running], + activeAgentId: "work") + let rejected = ChatSessionSidebarModel.applying( + observerDigest: SessionObserverDigest( + sessionkey: "global", + agentid: "main", + runid: "run-work", + revision: 2, + updatedat: 200, + headline: "Foreign status", + health: .stuck), + to: accepted, + activeAgentId: "work") + + #expect(accepted[0].observerDigest?.headline == "Work status") + #expect(rejected[0].observerDigest?.headline == "Work status") + } + @Test func `run rollover clears a stale digest before the replacement event`() throws { let existing = self.entry( key: "agent:main:work", diff --git a/apps/shared/OpenClawKit/Tests/OpenClawKitTests/ChatTranscriptCacheStoreTests.swift b/apps/shared/OpenClawKit/Tests/OpenClawKitTests/ChatTranscriptCacheStoreTests.swift index 232b4befefab..9e8e4b369220 100644 --- a/apps/shared/OpenClawKit/Tests/OpenClawKitTests/ChatTranscriptCacheStoreTests.swift +++ b/apps/shared/OpenClawKit/Tests/OpenClawKitTests/ChatTranscriptCacheStoreTests.swift @@ -32,11 +32,16 @@ func cacheMessage( idempotencyKey: idempotencyKey) } -func cacheSessionEntry(key: String, updatedAt: Double) -> OpenClawChatSessionEntry { +func cacheSessionEntry( + key: String, + updatedAt: Double, + agentID: String? = nil) -> OpenClawChatSessionEntry +{ OpenClawChatSessionEntry( key: key, kind: nil, displayName: nil, + agentId: agentID, surface: nil, subject: nil, room: nil, @@ -277,6 +282,88 @@ final class ChatTranscriptCacheStoreTests: ClientDatabaseTestSuite, @unchecked S #expect(!messageRows[0].payloadJSON.hasPrefix("[")) } + @Test func `agent session snapshots preserve another agents offline roster`() async throws { + await store.storeSessions([ + cacheSessionEntry(key: "global", updatedAt: 1, agentID: "agent-a"), + ], agentID: "agent-a") + await store.storeSessions([ + cacheSessionEntry(key: "global", updatedAt: 2, agentID: "agent-b"), + ], agentID: "agent-b") + + #expect(await store.loadSessions(agentID: "agent-a").map(\.agentId) == ["agent-a"]) + #expect(await store.loadSessions(agentID: "agent-b").map(\.agentId) == ["agent-b"]) + try databases.close() + + let reopened = try OpenClawClientDatabases(directoryURL: directory) + #expect(await reopened.store(gatewayID: "gw-a").loadSessions(agentID: "agent-a").map(\.agentId) == [ + "agent-a", + ]) + #expect(await reopened.store(gatewayID: "gw-a").loadSessions(agentID: "agent-b").map(\.agentId) == [ + "agent-b", + ]) + } + + @Test func `empty or rejected agent snapshot cannot erase another roster`() async { + await store.storeSessions([ + cacheSessionEntry(key: "global", updatedAt: 1, agentID: "agent-a"), + ], agentID: "agent-a") + await store.storeSessions([ + cacheSessionEntry(key: "global", updatedAt: 2, agentID: "agent-b"), + ], agentID: "agent-b") + + await store.storeSessions([ + cacheSessionEntry(key: "agent:agent-a:main", updatedAt: 3, agentID: "agent-a"), + ], agentID: "agent-b") + + #expect(await store.loadSessions(agentID: "agent-a").map(\.agentId) == ["agent-a"]) + #expect(await store.loadSessions(agentID: "agent-b").map(\.agentId) == ["agent-b"]) + + await store.storeSessions([], agentID: "agent-b") + + #expect(await store.loadSessions(agentID: "agent-a").map(\.agentId) == ["agent-a"]) + #expect(await store.loadSessions(agentID: "agent-b").isEmpty) + } + + @Test func `agent session owner partitions remain bounded`() async throws { + for index in 0...OpenClawChatSQLiteTranscriptCache.maxCachedSessionOwners { + let agentID = "agent-\(index)" + await store.storeSessions([ + cacheSessionEntry(key: "global", updatedAt: Double(index), agentID: agentID), + ], agentID: agentID) + } + + let counts = try await databases.cacheQueue.read { db in + try ( + Int.fetchOne(db, sql: "SELECT COUNT(*) FROM cached_session_rosters WHERE gateway_id = 'gw-a'"), + Int.fetchOne(db, sql: "SELECT COUNT(*) FROM cached_agent_sessions WHERE gateway_id = 'gw-a'")) + } + #expect(counts.0 == OpenClawChatSQLiteTranscriptCache.maxCachedSessionOwners) + #expect(counts.1 == OpenClawChatSQLiteTranscriptCache.maxCachedSessionOwners) + } + + @Test func `agent snapshot discards legacy roster without touching transcripts`() async throws { + await store.storeTestTranscript( + sessionKey: "global", + agentID: "agent-a", + messages: [cacheMessage(role: "assistant", text: "preserved", timestamp: 1)]) + try await databases.cacheQueue.write { db in + try db.execute( + sql: """ + INSERT INTO cached_sessions(gateway_id, session_key, position, updated_at, payload_json) + VALUES ('gw-a', 'global', 0, 1, '{}') + """) + } + + await store.storeSessions([ + cacheSessionEntry(key: "global", updatedAt: 2, agentID: "agent-a"), + ], agentID: "agent-a") + + #expect(try await databases.cacheQueue.read { db in + try Int.fetchOne(db, sql: "SELECT COUNT(*) FROM cached_sessions WHERE gateway_id = 'gw-a'") + } == 0) + #expect(await messageTexts(store.loadTranscript(sessionKey: "global", agentID: "agent-a")) == ["preserved"]) + } + @Test func `cache format mismatch rebuilds without touching client state`() async throws { let stateIdentity = try #require(OpenClawChatSessionRoutingIdentity( scope: "per-sender", @@ -421,21 +508,33 @@ final class ChatTranscriptCacheStoreTests: ClientDatabaseTestSuite, @unchecked S } @Test func `malformed cache partitions are discarded atomically`() async throws { - await store.storeSessions([cacheSessionEntry(key: "main", updatedAt: 1)]) + await store.storeSessions([ + cacheSessionEntry(key: "global", updatedAt: 1, agentID: "agent-a"), + ], agentID: "agent-a") + await store.storeSessions([ + cacheSessionEntry(key: "global", updatedAt: 2, agentID: "agent-b"), + ], agentID: "agent-b") await store.storeTestTranscript( sessionKey: "main", messages: [cacheMessage(role: "assistant", text: "cached", timestamp: 1)]) try await databases.cacheQueue.write { db in try db.execute( - sql: "UPDATE cached_sessions SET payload_json = 'not-json' WHERE gateway_id = 'gw-a'") + sql: """ + UPDATE cached_agent_sessions SET payload_json = 'not-json' + WHERE gateway_id = 'gw-a' AND agent_id = 'agent-a' + """) try db.execute( sql: "UPDATE cached_messages SET payload_json = 'not-json' WHERE gateway_id = 'gw-a'") } - #expect(await store.loadSessions().isEmpty) + #expect(await store.loadSessions(agentID: "agent-a").isEmpty) + #expect(await store.loadSessions(agentID: "agent-b").map(\.agentId) == ["agent-b"]) #expect(await store.loadTranscript(sessionKey: "main").isEmpty) #expect(try await databases.cacheQueue.read { db in - try Int.fetchOne(db, sql: "SELECT COUNT(*) FROM cached_sessions WHERE gateway_id = 'gw-a'") + try Int.fetchOne(db, sql: """ + SELECT COUNT(*) FROM cached_agent_sessions + WHERE gateway_id = 'gw-a' AND agent_id = 'agent-a' + """) } == 0) #expect(try await databases.cacheQueue.read { db in try Int.fetchOne(db, sql: "SELECT COUNT(*) FROM cached_transcripts WHERE gateway_id = 'gw-a'") diff --git a/apps/shared/OpenClawKit/Tests/OpenClawKitTests/ChatViewModelTranscriptCacheTests.swift b/apps/shared/OpenClawKit/Tests/OpenClawKitTests/ChatViewModelTranscriptCacheTests.swift index 9383774f82ca..07e20cd9af9d 100644 --- a/apps/shared/OpenClawKit/Tests/OpenClawKitTests/ChatViewModelTranscriptCacheTests.swift +++ b/apps/shared/OpenClawKit/Tests/OpenClawKitTests/ChatViewModelTranscriptCacheTests.swift @@ -49,6 +49,10 @@ private actor TestTranscriptCache: OpenClawChatTranscriptCache { return self.sessions } + func loadSessions(agentID _: String?) async -> [OpenClawChatSessionEntry] { + await self.loadSessions() + } + func loadTranscript(sessionKey: String) async -> [OpenClawChatMessage] { self.transcripts[sessionKey] ?? [] } @@ -57,6 +61,10 @@ private actor TestTranscriptCache: OpenClawChatTranscriptCache { self.sessions = sessions } + func storeSessions(_ sessions: [OpenClawChatSessionEntry], agentID _: String?) async { + await self.storeSessions(sessions) + } + func storeCanonicalTranscript( sessionKey: String, agentID _: String?, @@ -68,6 +76,35 @@ private actor TestTranscriptCache: OpenClawChatTranscriptCache { } } +/// Old conformers compile through the protocol defaults but cannot partition +/// selected-agent rosters, so scoped access must never reach this storage. +private actor LegacyTranscriptCache: OpenClawChatTranscriptCache { + private var sessions: [OpenClawChatSessionEntry] + + init(sessions: [OpenClawChatSessionEntry] = []) { + self.sessions = sessions + } + + func loadSessions() async -> [OpenClawChatSessionEntry] { + self.sessions + } + + func loadTranscript(sessionKey _: String) async -> [OpenClawChatMessage] { + [] + } + + func storeSessions(_ sessions: [OpenClawChatSessionEntry]) async { + self.sessions = sessions + } + + func storeCanonicalTranscript( + sessionKey _: String, + agentID _: String?, + messages _: [OpenClawChatMessage], + canonicalMessageIdempotencyKeys _: Set) async + {} +} + /// Minimal FIFO transport whose history responses can be gated during cold open. private final class GatedHistoryChatTransport: @unchecked Sendable, OpenClawChatTransport { private let historyResult: @Sendable (String, Int) async throws -> OpenClawChatHistoryPayload @@ -127,7 +164,7 @@ private func makeViewModel( sessionKey: String = "main", transport: GatedHistoryChatTransport, activeAgentID: String? = nil, - cache: TestTranscriptCache, + cache: any OpenClawChatTranscriptCache, load: Bool = true) -> OpenClawChatViewModel { let vm = OpenClawChatViewModel( @@ -311,7 +348,7 @@ struct ChatViewModelTranscriptCacheTests { #expect(await MainActor.run { vm.sessions.isEmpty }) } - @Test func `ownerless shared cache cannot claim the selected global owner`() async throws { + @Test func `legacy cache cannot paint an ownerless roster for either selected agent`() async throws { var global = cacheSessionEntry(key: "global", updatedAt: 1000) global.observerDigest = OpenClawChatSessionObserverDigest( runId: "run-legacy", @@ -319,7 +356,7 @@ struct ChatViewModelTranscriptCacheTests { updatedAt: 1000, headline: "Ambiguous legacy owner", health: "on-track") - let cache = TestTranscriptCache(sessions: [global]) + let cache = LegacyTranscriptCache(sessions: [global]) let transport = GatedHistoryChatTransport { sessionKey, _ in historyPayload(sessionKey: sessionKey, sessionID: "unused-live-session") } @@ -332,12 +369,65 @@ struct ChatViewModelTranscriptCacheTests { let snapshot = await MainActor.run { vm.currentSessionSnapshot() } await MainActor.run { vm.paintFromCacheIfNeeded(session: snapshot) } - try await waitUntil("shared cache row painted without ambiguous digest") { - await MainActor.run { vm.sessions.count == 1 } + try await Task.sleep(nanoseconds: 100_000_000) + + await MainActor.run { vm.syncActiveAgentId("gadget") } + let switchedSnapshot = await MainActor.run { vm.currentSessionSnapshot() } + await MainActor.run { vm.paintFromCacheIfNeeded(session: switchedSnapshot) } + try await Task.sleep(nanoseconds: 100_000_000) + + #expect(await MainActor.run { vm.sessions.isEmpty }) + } + + @Test func `legacy cache rejects scoped roster writes`() async { + let original = cacheSessionEntry(key: "agent-a", updatedAt: 1000) + let cache = LegacyTranscriptCache(sessions: [original]) + + await cache.storeSessions( + [cacheSessionEntry(key: "agent-b", updatedAt: 2000)], + agentID: "agent-b") + + #expect(await cache.loadSessions().map(\.key) == ["agent-a"]) + #expect(await cache.loadSessions(agentID: "agent-a").isEmpty) + #expect(await cache.loadSessions(agentID: "agent-b").isEmpty) + } + + @Test func `cached session prepaint stays within the selected agent`() async throws { + var matchingBare = cacheSessionEntry(key: "shared-tool", updatedAt: 2000) + matchingBare.agentId = "main" + var foreignBare = cacheSessionEntry(key: "foreign-tool", updatedAt: 1750) + foreignBare.agentId = "gadget" + var foreignGlobal = cacheSessionEntry(key: "global", updatedAt: 1500) + foreignGlobal.agentId = "gadget" + let cache = TestTranscriptCache(sessions: [ + cacheSessionEntry(key: "agent:main:owned", updatedAt: 3000), + cacheSessionEntry(key: "agent:gadget:foreign", updatedAt: 2500), + matchingBare, + foreignBare, + foreignGlobal, + cacheSessionEntry(key: "legacy-ownerless", updatedAt: 1000), + ]) + let transport = GatedHistoryChatTransport { sessionKey, _ in + historyPayload(sessionKey: sessionKey, sessionID: "unused-live-session") + } + let vm = await makeViewModel( + sessionKey: "agent:main:main", + transport: transport, + activeAgentID: "main", + cache: cache, + load: false) + let snapshot = await MainActor.run { vm.currentSessionSnapshot() } + + await MainActor.run { vm.paintFromCacheIfNeeded(session: snapshot) } + try await waitUntil("cached sessions painted") { + await MainActor.run { !vm.sessions.isEmpty } } - #expect(await MainActor.run { vm.sessions[0].key == "global" }) - #expect(await MainActor.run { vm.sessions[0].observerDigest == nil }) + #expect(await MainActor.run { vm.sessions.map(\.key) } == [ + "agent:main:owned", + "shared-tool", + "legacy-ownerless", + ]) } @Test func `session cache strips active markers and preserves terminal recap`() {