import Foundation import OpenClawKit import Testing @testable import OpenClaw private actor SystemAgentGatewayConfig { private var token = "a" func snapshotToken() -> String { self.token } func setToken(_ token: String) { self.token = token } } private actor SystemAgentSessionRecorder { private var sessionIDs: [String] = [] func record(_ sessionID: String) { self.sessionIDs.append(sessionID) } func snapshot() -> [String] { self.sessionIDs } } private actor SystemAgentMessageRecorder { private var messages: [String] = [] func record(_ message: String) { self.messages.append(message) } func snapshot() -> [String] { self.messages } } private actor SystemAgentMethodRecorder { private var methods: [String] = [] func record(_ method: String) { self.methods.append(method) } func snapshot() -> [String] { self.methods } } private actor SystemAgentRequestGate { private var consumed = false private var released = false private var continuation: CheckedContinuation? func waitIfFirst() async -> Bool { guard !self.consumed else { return false } self.consumed = true if !self.released { await withCheckedContinuation { continuation in self.continuation = continuation } } return true } func release() { self.released = true self.continuation?.resume() self.continuation = nil } } private func systemAgentSessionID(from message: URLSessionWebSocketTask.Message) -> String? { let data: Data? = switch message { case let .data(data): data case let .string(string): string.data(using: .utf8) @unknown default: nil } guard let data, let object = try? JSONSerialization.jsonObject(with: data) as? [String: Any], object["method"] as? String == "openclaw.chat", let params = object["params"] as? [String: Any] else { return nil } return params["sessionId"] as? String } private func systemAgentRequestMethod(from message: URLSessionWebSocketTask.Message) -> String? { let data: Data? = switch message { case let .data(data): data case let .string(string): string.data(using: .utf8) @unknown default: nil } guard let data, let object = try? JSONSerialization.jsonObject(with: data) as? [String: Any] else { return nil } return object["method"] as? String } private func systemAgentChatMessage(from message: URLSessionWebSocketTask.Message) -> String? { let data: Data? = switch message { case let .data(data): data case let .string(string): string.data(using: .utf8) @unknown default: nil } guard let data, let object = try? JSONSerialization.jsonObject(with: data) as? [String: Any], object["method"] as? String == "openclaw.chat", let params = object["params"] as? [String: Any] else { return nil } return params["message"] as? String } private func respondToSystemAgentHealth( task: GatewayTestWebSocketTask, id: String, method: String?) -> Bool { guard method == "health" else { return false } task.emitReceiveSuccess(.data(GatewayWebSocketTestSupport.okResponseData(id: id))) return true } private func systemAgentResponse( id: String, action: String = "none", agentDraft: String? = nil, questionJSON: String? = nil) -> Data { let agentDraftField = agentDraft.map { ",\n \"agentDraft\": \"\($0)\"" } ?? "" let questionField = questionJSON.map { ",\n \"question\": \($0)" } ?? "" return Data( """ { "type": "res", "id": "\(id)", "ok": true, "payload": { "sessionId": "test-session", "reply": "ready", "action": "\(action)", "sensitive": false\(agentDraftField)\(questionField) } } """.utf8) } private func verifiedInferenceResponse(id: String) -> Data { Data( """ { "type": "res", "id": "\(id)", "ok": true, "payload": { "ok": true, "modelRef": "openai/gpt-5.5", "latencyMs": 42 } } """.utf8) } private func configuredAgentsResponse(id: String) -> Data { Data( """ { "type": "res", "id": "\(id)", "ok": true, "payload": { "defaultId": "main", "mainKey": "main", "scope": "per-sender", "agents": [{ "id": "main", "model": { "primary": "openai/gpt-5.5" } }] } } """.utf8) } private func transientVerificationErrorResponse(id: String) -> Data { Data( """ { "type": "res", "id": "\(id)", "ok": false, "error": { "code": "UNAVAILABLE", "message": "temporary disconnect" } } """.utf8) } @Suite(.serialized) @MainActor struct OnboardingSystemAgentChatTests { @Test func `fresh inference connection finishes onboarding once`() { let state = AppState(preview: true) state.connectionMode = .local var handoffs: [OnboardingDashboardHandoff] = [] let view = OnboardingView( state: state, dashboardHandoffOpener: { handoffs.append($0) }) view.prepareSystemAgentHandoff() view.aiSetup.onConnected?() #expect(view.finishState.didFinish) // Fresh connections own the remaining first-run steps via the custodian. #expect(handoffs == [.custodianOnboarding]) #expect(!view.finish()) #expect(handoffs == [.custodianOnboarding]) } @Test func `first run effective model is live verified before handoff`() async throws { let suiteName = "OnboardingFirstRunEffectiveModelTests-\(UUID().uuidString)" let defaults = try #require(UserDefaults(suiteName: suiteName)) defer { defaults.removePersistentDomain(forName: suiteName) } let methods = SystemAgentMethodRecorder() let session = GatewayTestWebSocketSession(taskFactory: { GatewayTestWebSocketTask(sendHook: { task, message, sendIndex in guard sendIndex > 0, let id = GatewayWebSocketTestSupport.requestID(from: message), let method = systemAgentRequestMethod(from: message) else { return } await methods.record(method) if respondToSystemAgentHealth(task: task, id: id, method: method) { return } switch method { case "agents.list": task.emitReceiveSuccess(.data(configuredAgentsResponse(id: id))) case "openclaw.setup.verify": task.emitReceiveSuccess(.data(verifiedInferenceResponse(id: id))) default: break } }) }) let url = try #require(URL(string: "ws://localhost:18789")) let gateway = GatewayConnection( configProvider: { (url: url, token: nil, password: nil) }, sessionBox: WebSocketSessionBox(session: session)) let appState = AppState(preview: true) appState.connectionMode = .local var handoffs: [OnboardingDashboardHandoff] = [] let view = OnboardingView( state: appState, aiSetupGateway: gateway, systemAgentDefaults: defaults, aiSetupRouteIdentityProvider: { "local" }, dashboardHandoffOpener: { handoffs.append($0) }) let initialProbe = try #require(view.onboardingDidAppear()) await initialProbe.value for _ in 0..<200 { if view.aiSetup.connected { break } try? await Task.sleep(nanoseconds: 5_000_000) } #expect(view.aiSetup.connected) #expect(view.finishState.didFinish) // A live-verified pre-existing setup reopens the normal dashboard. #expect(handoffs == [.dashboard]) #expect(await methods.snapshot() == [ "agents.list", "health", "openclaw.setup.verify", ]) } @Test func `relaunch with pending inference resumes OpenClaw`() async throws { let suiteName = "OnboardingPendingInferenceResumeTests-\(UUID().uuidString)" let defaults = try #require(UserDefaults(suiteName: suiteName)) defer { defaults.removePersistentDomain(forName: suiteName) } let methods = SystemAgentMethodRecorder() let session = GatewayTestWebSocketSession(taskFactory: { GatewayTestWebSocketTask(sendHook: { task, message, sendIndex in guard sendIndex > 0, let id = GatewayWebSocketTestSupport.requestID(from: message) else { return } let method = systemAgentRequestMethod(from: message) if let method { await methods.record(method) } if respondToSystemAgentHealth(task: task, id: id, method: method) { return } switch method { case "openclaw.setup.verify": task.emitReceiveSuccess(.data(verifiedInferenceResponse(id: id))) case "openclaw.chat": task.emitReceiveSuccess(.data(systemAgentResponse(id: id))) default: break } }) }) let url = try #require(URL(string: "ws://example.invalid")) let gateway = GatewayConnection( configProvider: { (url: url, token: nil, password: nil) }, sessionBox: WebSocketSessionBox(session: session)) let appState = AppState(preview: true) appState.connectionMode = .remote appState.remoteTransport = .direct appState.remoteUrl = "ws://example.invalid" var handoffs: [OnboardingDashboardHandoff] = [] let view = OnboardingView( state: appState, aiSetupGateway: gateway, systemAgentDefaults: defaults, aiSetupRouteIdentityProvider: { "remote:direct:example.invalid" }, dashboardHandoffOpener: { handoffs.append($0) }) let task = view.resumePendingSystemAgent(modelRef: "openai/gpt-5.5") await task.value #expect(view.aiSetup.connected) #expect(view.aiSetup.selectedKind == "existing-model") #expect(view.finishState.didFinish) #expect(handoffs == [.dashboard]) #expect(!view.finish()) let repeatedResume = view.resumePendingSystemAgent(modelRef: "openai/gpt-5.5") await repeatedResume.value #expect(view.aiSetup.connected) #expect(view.aiSetup.selectedKind == "existing-model") #expect(handoffs == [.dashboard]) #expect(await methods.snapshot() == [ "health", "openclaw.setup.verify", ]) } @Test func `pending verification retry schedules deadline and stays read only`() async throws { let suiteName = "OnboardingPendingVerificationRetryTests-\(UUID().uuidString)" let defaults = try #require(UserDefaults(suiteName: suiteName)) defer { defaults.removePersistentDomain(forName: suiteName) } let methods = SystemAgentMethodRecorder() let session = GatewayTestWebSocketSession(taskFactory: { GatewayTestWebSocketTask(sendHook: { task, message, sendIndex in guard sendIndex > 0, let id = GatewayWebSocketTestSupport.requestID(from: message), let method = systemAgentRequestMethod(from: message) else { return } await methods.record(method) if respondToSystemAgentHealth(task: task, id: id, method: method) { return } switch method { case "openclaw.setup.verify": let priorVerifications = await methods.snapshot().filter { $0 == "openclaw.setup.verify" }.count let response = priorVerifications == 1 ? transientVerificationErrorResponse(id: id) : verifiedInferenceResponse(id: id) task.emitReceiveSuccess(.data(response)) case "openclaw.chat": task.emitReceiveSuccess(.data(systemAgentResponse(id: id))) default: break } }) }) let url = try #require(URL(string: "ws://example.invalid")) let gateway = GatewayConnection( configProvider: { (url: url, token: nil, password: nil) }, sessionBox: WebSocketSessionBox(session: session)) let appState = AppState(preview: true) appState.connectionMode = .local OnboardingSystemAgentResumeStore.markPending( routeIdentity: "local", defaults: defaults) let view = OnboardingView( state: appState, aiSetupGateway: gateway, systemAgentDefaults: defaults, aiSetupRouteIdentityProvider: { "local" }) await view.resumePendingSystemAgent(modelRef: "openai/gpt-5.5").value var scheduledDeadlines: [(deadline: Date, routeIdentity: String)] = [] view.aiSetup.onPendingActivationDeadline = { deadline, routeIdentity in scheduledDeadlines.append((deadline, routeIdentity)) } view.aiSetup.retryFromScratch() for _ in 0..<200 { if case .verified = OnboardingSystemAgentResumeStore.pendingState( for: "local", defaults: defaults) { break } try? await Task.sleep(nanoseconds: 5_000_000) } #expect(!view.aiSetup.connected) #expect(view.aiSetup.waitingForPendingActivationDeadline) #expect(scheduledDeadlines.count == 1) #expect(scheduledDeadlines.first?.routeIdentity == "local") if case let .verified(deadline) = OnboardingSystemAgentResumeStore.pendingState( for: "local", defaults: defaults) { #expect(scheduledDeadlines.first?.deadline == deadline) } else { Issue.record("expected verified activation lease") } view.aiSetup.retryFromScratch() #expect(scheduledDeadlines.count == 1) #expect(await methods.snapshot() == [ "health", "openclaw.setup.verify", "health", "openclaw.setup.verify", ]) } @Test func `superseded resume cannot finish a replacement route handoff`() async throws { let suiteName = "OnboardingSupersededResumeTests-\(UUID().uuidString)" let defaults = try #require(UserDefaults(suiteName: suiteName)) defer { defaults.removePersistentDomain(forName: suiteName) } let gate = SystemAgentRequestGate() let methods = SystemAgentMethodRecorder() let session = GatewayTestWebSocketSession(taskFactory: { GatewayTestWebSocketTask(sendHook: { task, message, sendIndex in guard sendIndex > 0, let id = GatewayWebSocketTestSupport.requestID(from: message), let method = systemAgentRequestMethod(from: message) else { return } await methods.record(method) if respondToSystemAgentHealth(task: task, id: id, method: method) { return } guard method == "openclaw.setup.verify" else { return } _ = await gate.waitIfFirst() task.emitReceiveSuccess(.data(verifiedInferenceResponse(id: id))) }) }) let url = try #require(URL(string: "ws://example.invalid")) let gateway = GatewayConnection( configProvider: { (url: url, token: nil, password: nil) }, sessionBox: WebSocketSessionBox(session: session)) let appState = AppState(preview: true) appState.connectionMode = .remote appState.remoteTransport = .direct appState.remoteUrl = "ws://example.invalid" var dashboardOpenCount = 0 let view = OnboardingView( state: appState, aiSetupGateway: gateway, systemAgentDefaults: defaults, aiSetupRouteIdentityProvider: { "remote:direct:example.invalid" }, dashboardHandoffOpener: { _ in dashboardOpenCount += 1 }) let staleResume = view.resumePendingSystemAgent(modelRef: "openai/gpt-5.5") for _ in 0..<200 { if await methods.snapshot() == ["health", "openclaw.setup.verify"] { break } try? await Task.sleep(nanoseconds: 5_000_000) } view.resetGatewayBoundAIState() // Simulate a newer route reaching connected state without handing off. // The stale wrapper must not infer success from this state. view.aiSetup.onConnected = nil view.aiSetup.resumeConfiguredInference(modelRef: "openai/gpt-5.5") view.aiSetup.acceptVerifiedPendingInference(modelRef: "openai/gpt-5.5") await gate.release() await staleResume.value #expect(view.aiSetup.connected) #expect(!view.finishState.didFinish) #expect(dashboardOpenCount == 0) #expect(await methods.snapshot() == ["health", "openclaw.setup.verify"]) } @Test func `cold launch resumes a completed activation immediately`() async throws { let suiteName = "OnboardingColdPendingHandoffTests-\(UUID().uuidString)" let defaults = try #require(UserDefaults(suiteName: suiteName)) defer { defaults.removePersistentDomain(forName: suiteName) } let methods = SystemAgentMethodRecorder() let session = GatewayTestWebSocketSession(taskFactory: { GatewayTestWebSocketTask(sendHook: { task, message, sendIndex in guard sendIndex > 0, let id = GatewayWebSocketTestSupport.requestID(from: message), let method = systemAgentRequestMethod(from: message) else { return } await methods.record(method) if respondToSystemAgentHealth(task: task, id: id, method: method) { return } switch method { case "agents.list": task.emitReceiveSuccess(.data(configuredAgentsResponse(id: id))) case "openclaw.setup.verify": task.emitReceiveSuccess(.data(verifiedInferenceResponse(id: id))) case "openclaw.chat": task.emitReceiveSuccess(.data(systemAgentResponse(id: id))) default: break } }) }) let url = try #require(URL(string: "ws://localhost:18789")) let gateway = GatewayConnection( configProvider: { (url: url, token: nil, password: nil) }, sessionBox: WebSocketSessionBox(session: session)) let appState = AppState(preview: true) appState.connectionMode = .local let routeIdentity = OnboardingSystemAgentResumeStore.selectedRouteIdentity(state: appState) let route = try #require(await gateway.captureRoute()) let activationOwner = try OnboardingSystemAgentResumeStore.ActivationOwner( id: "completed-before-relaunch", routeFingerprint: #require(route.activationOwnershipFingerprint)) OnboardingSystemAgentResumeStore.markPending( routeIdentity: routeIdentity, activationOwner: activationOwner, defaults: defaults) OnboardingSystemAgentResumeStore.markCompleted( ifOwnedBy: routeIdentity, activationOwner: activationOwner, defaults: defaults) var dashboardOpenCount = 0 let view = OnboardingView( state: appState, aiSetupGateway: gateway, systemAgentDefaults: defaults, aiSetupRouteIdentityProvider: { routeIdentity }, dashboardHandoffOpener: { _ in dashboardOpenCount += 1 }) let aiSetup = view.aiSetup let initialProbe = try #require(view.onboardingDidAppear()) await initialProbe.value for _ in 0..<200 { if aiSetup.connected { break } try? await Task.sleep(nanoseconds: 5_000_000) } #expect(aiSetup.connected) #expect(view.finishState.didFinish) #expect(dashboardOpenCount == 1) #expect(OnboardingSystemAgentResumeStore.pendingState( for: routeIdentity, defaults: defaults) == .none) #expect(await methods.snapshot() == [ "agents.list", "health", "openclaw.setup.verify", ]) } @Test func `typed question sends reply while transcript shows label`() async throws { let recordedMessages = SystemAgentMessageRecorder() let questionJSON = """ {"id":"next","header":"Next step","question":"What now?","options":[ {"label":"Talk to my agent","reply":"talk to agent","recommended":true}, {"label":"Connect WhatsApp","reply":"connect whatsapp","description":"Chat there."} ],"isOther":true} """ let session = GatewayTestWebSocketSession(taskFactory: { GatewayTestWebSocketTask(sendHook: { task, message, sendIndex in guard sendIndex > 0, let id = GatewayWebSocketTestSupport.requestID(from: message) else { return } if let message = systemAgentChatMessage(from: message) { await recordedMessages.record(message) task.emitReceiveSuccess(.data(systemAgentResponse(id: id))) } else { task.emitReceiveSuccess(.data(systemAgentResponse( id: id, questionJSON: questionJSON))) } }) }) let url = try #require(URL(string: "ws://example.invalid")) let gateway = GatewayConnection( configProvider: { (url: url, token: nil, password: nil) }, sessionBox: WebSocketSessionBox(session: session)) let chat = SystemAgentOnboardingChatModel(gateway: gateway) await chat.startIfNeeded() let assistant = try #require(chat.messages.first) let question = try #require(assistant.question) #expect(question.options.first?.recommended == true) let task = try #require(chat.answerQuestion( messageID: assistant.id, optionLabel: "Connect WhatsApp")) await task.value #expect(await recordedMessages.snapshot() == ["connect whatsapp"]) #expect(chat.messages.map(\.text) == ["ready", "Connect WhatsApp", "ready"]) #expect(!chat.canAnswerQuestion(assistant)) } @Test func `typed question skip sends fixed reply and dismisses cards`() async throws { let recordedMessages = SystemAgentMessageRecorder() let questionJSON = #"{"id":"next","header":"Next step","question":"What now?","options":[{"label":"A"},{"label":"B"}]}"# let session = GatewayTestWebSocketSession(taskFactory: { GatewayTestWebSocketTask(sendHook: { task, message, sendIndex in guard sendIndex > 0, let id = GatewayWebSocketTestSupport.requestID(from: message) else { return } if let message = systemAgentChatMessage(from: message) { await recordedMessages.record(message) task.emitReceiveSuccess(.data(systemAgentResponse(id: id))) } else { task.emitReceiveSuccess(.data(systemAgentResponse( id: id, questionJSON: questionJSON))) } }) }) let url = try #require(URL(string: "ws://example.invalid")) let gateway = GatewayConnection( configProvider: { (url: url, token: nil, password: nil) }, sessionBox: WebSocketSessionBox(session: session)) let chat = SystemAgentOnboardingChatModel(gateway: gateway) await chat.startIfNeeded() let assistant = try #require(chat.messages.first) let task = try #require(chat.skipQuestion(messageID: assistant.id)) await task.value #expect(await recordedMessages.snapshot() == ["Skip for now"]) #expect(chat.messages.map(\.text) == ["ready", "Skip for now", "ready"]) #expect(!chat.isQuestionVisible(assistant)) } @Test(arguments: [ #"{"id":"dupes","header":"Next step","question":"What now?","options":[{"label":"Same"},{"label":"same"}]}"#, #""invalid""#, #"[]"#, ]) func `malformed typed question keeps prose reply only`(questionJSON: String) async throws { let session = GatewayTestWebSocketSession(taskFactory: { GatewayTestWebSocketTask(sendHook: { task, message, sendIndex in guard sendIndex > 0, let id = GatewayWebSocketTestSupport.requestID(from: message) else { return } task.emitReceiveSuccess(.data(systemAgentResponse( id: id, questionJSON: questionJSON))) }) }) let url = try #require(URL(string: "ws://example.invalid")) let gateway = GatewayConnection( configProvider: { (url: url, token: nil, password: nil) }, sessionBox: WebSocketSessionBox(session: session)) let chat = SystemAgentOnboardingChatModel(gateway: gateway) await chat.startIfNeeded() #expect(chat.messages.map(\.text) == ["ready"]) #expect(chat.messages.first?.question == nil) } @Test func `agent handoff carries the hatch draft intent`() async throws { let session = GatewayTestWebSocketSession(taskFactory: { GatewayTestWebSocketTask(sendHook: { task, message, sendIndex in guard sendIndex > 0, let id = GatewayWebSocketTestSupport.requestID(from: message) else { return } task.emitReceiveSuccess(.data(systemAgentResponse( id: id, action: "open-agent", agentDraft: "hatch"))) }) }) let url = try #require(URL(string: "ws://example.invalid")) let gateway = GatewayConnection( configProvider: { (url: url, token: nil, password: nil) }, sessionBox: WebSocketSessionBox(session: session)) let chat = SystemAgentOnboardingChatModel(gateway: gateway) var receivedDraft: SystemAgentDraft? chat.onAgentHandoff = { receivedDraft = $0 } await chat.startIfNeeded() #expect(receivedDraft == .hatch) #expect(receivedDraft?.composerValue == "Wake up, my friend!") } @Test func `settings callback refreshes inference after assistant reply`() async throws { let session = GatewayTestWebSocketSession(taskFactory: { GatewayTestWebSocketTask(sendHook: { task, message, sendIndex in guard sendIndex > 0, let id = GatewayWebSocketTestSupport.requestID(from: message) else { return } task.emitReceiveSuccess(.data(systemAgentResponse(id: id))) }) }) let url = try #require(URL(string: "ws://example.invalid")) let gateway = GatewayConnection( configProvider: { (url: url, token: nil, password: nil) }, sessionBox: WebSocketSessionBox(session: session)) let chat = SystemAgentOnboardingChatModel(gateway: gateway) var refreshCount = 0 SystemAgentSettings.configureChatCallbacks( for: chat, onReplyReceived: { refreshCount += 1 }) await chat.startIfNeeded() #expect(chat.messages.map(\.text) == ["ready"]) #expect(refreshCount == 1) } @Test func `model invalidation cancels queued send and restart tasks`() async throws { let session = GatewayTestWebSocketSession() let url = try #require(URL(string: "ws://example.invalid")) let gateway = GatewayConnection( configProvider: { (url: url, token: nil, password: nil) }, sessionBox: WebSocketSessionBox(session: session)) let chat = SystemAgentOnboardingChatModel(gateway: gateway) var replyCount = 0 var handoffCount = 0 chat.onReplyReceived = { replyCount += 1 } chat.onAgentHandoff = { _ in handoffCount += 1 } chat.input = "route-bound secret" let sendTask = try #require(chat.send()) let restartTask = try #require(chat.restartAfterError()) chat.invalidate() await sendTask.value await restartTask.value #expect(session.snapshotMakeCount() == 0) #expect(chat.messages.isEmpty) #expect(replyCount == 0) #expect(handoffCount == 0) #expect(chat.send() == nil) #expect(chat.restartAfterError() == nil) } @Test func `chat session stays bound to its original gateway route`() async throws { let config = SystemAgentGatewayConfig() let recorder = SystemAgentSessionRecorder() let session = GatewayTestWebSocketSession(taskFactory: { GatewayTestWebSocketTask(sendHook: { task, message, sendIndex in guard sendIndex > 0, let id = GatewayWebSocketTestSupport.requestID(from: message) else { return } if let sessionID = systemAgentSessionID(from: message) { await recorder.record(sessionID) } task.emitReceiveSuccess(.data(systemAgentResponse(id: id))) }) }) let url = try #require(URL(string: "ws://example.invalid")) let gateway = GatewayConnection( configProvider: { let token = await config.snapshotToken() return (url: url, token: token, password: nil) }, sessionBox: WebSocketSessionBox(session: session)) let chat = SystemAgentOnboardingChatModel(gateway: gateway) await chat.startIfNeeded() #expect(chat.messages.map(\.text) == ["ready"]) #expect(session.snapshotMakeCount() == 1) #expect(session.latestTask()?.snapshotSendCount() == 2) let routeASessionIDs = await recorder.snapshot() #expect(routeASessionIDs.count == 1) let routeASessionID = try #require(routeASessionIDs.first) await config.setToken("b") chat.input = "must stay on route a" let sendTask = try #require(chat.send()) await sendTask.value #expect(session.snapshotMakeCount() == 1) #expect(session.latestTask()?.snapshotSendCount() == 2) #expect(chat.messages.map(\.text) == ["ready", "must stay on route a"]) #expect(chat.errorMessage == "The Gateway connection changed. Restart OpenClaw to reconnect.") #expect(await recorder.snapshot() == [routeASessionID]) let restartTask = try #require(chat.restartAfterError()) await restartTask.value #expect(session.snapshotMakeCount() == 2) #expect(session.latestTask()?.snapshotSendCount() == 2) #expect(chat.messages.map(\.text) == ["ready"]) #expect(chat.errorMessage == nil) let sessionIDs = await recorder.snapshot() #expect(sessionIDs.count == 2) #expect(sessionIDs.first == routeASessionID) #expect(sessionIDs.last != routeASessionID) } @Test func `route change while reply is in flight discards reply and action`() async throws { let config = SystemAgentGatewayConfig() let requestGate = SystemAgentRequestGate() let session = GatewayTestWebSocketSession(taskFactory: { GatewayTestWebSocketTask(sendHook: { task, message, sendIndex in guard sendIndex > 0, let id = GatewayWebSocketTestSupport.requestID(from: message) else { return } _ = await requestGate.waitIfFirst() task.emitReceiveSuccess(.data(systemAgentResponse(id: id, action: "open-agent"))) }) }) let url = try #require(URL(string: "ws://example.invalid")) let gateway = GatewayConnection( configProvider: { let token = await config.snapshotToken() return (url: url, token: token, password: nil) }, sessionBox: WebSocketSessionBox(session: session)) let chat = SystemAgentOnboardingChatModel(gateway: gateway) var replyCount = 0 var handoffCount = 0 chat.onReplyReceived = { replyCount += 1 } chat.onAgentHandoff = { _ in handoffCount += 1 } let startTask = Task { await chat.startIfNeeded() } var requestStarted = false for _ in 0..<1000 { if session.latestTask()?.snapshotSendCount() == 2 { requestStarted = true break } await Task.yield() } try #require(requestStarted) await config.setToken("b") await requestGate.release() await startTask.value #expect(chat.messages.isEmpty) #expect(replyCount == 0) #expect(handoffCount == 0) #expect(chat.errorMessage == "The Gateway connection changed. Restart OpenClaw to reconnect.") } @Test func `cancelled initial request exposes restart and recovers`() async throws { let requestGate = SystemAgentRequestGate() let session = GatewayTestWebSocketSession(taskFactory: { GatewayTestWebSocketTask(sendHook: { task, message, sendIndex in guard sendIndex > 0, let id = GatewayWebSocketTestSupport.requestID(from: message) else { return } if sendIndex == 1, await requestGate.waitIfFirst() { throw CancellationError() } task.emitReceiveSuccess(.data(systemAgentResponse(id: id))) }) }) let url = try #require(URL(string: "ws://example.invalid")) let gateway = GatewayConnection( configProvider: { (url: url, token: nil, password: nil) }, sessionBox: WebSocketSessionBox(session: session)) let chat = SystemAgentOnboardingChatModel(gateway: gateway) let startTask = Task { await chat.startIfNeeded() } var requestStarted = false for _ in 0..<1000 { if session.latestTask()?.snapshotSendCount() == 2 { requestStarted = true break } await Task.yield() } try #require(requestStarted) startTask.cancel() await requestGate.release() await startTask.value #expect(chat.errorMessage == "OpenClaw was interrupted. Restart to try again.") #expect(!chat.isSending) #expect(chat.messages.isEmpty) let restartTask = try #require(chat.restartAfterError()) await restartTask.value #expect(chat.errorMessage == nil) #expect(chat.messages.map(\.text) == ["ready"]) #expect(session.snapshotMakeCount() == 1) #expect(session.latestTask()?.snapshotSendCount() == 3) } }