mirror of
https://github.com/openclaw/openclaw.git
synced 2026-08-17 08:02:12 -06:00
828 lines
32 KiB
Swift
828 lines
32 KiB
Swift
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<Void, Never>?
|
|
|
|
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 dashboardOpenCount = 0
|
|
let view = OnboardingView(
|
|
state: state,
|
|
dashboardOnboardingOpener: { dashboardOpenCount += 1 })
|
|
|
|
view.prepareSystemAgentHandoff()
|
|
view.aiSetup.onConnected?()
|
|
|
|
#expect(view.finishState.didFinish)
|
|
#expect(dashboardOpenCount == 1)
|
|
#expect(!view.finish())
|
|
#expect(dashboardOpenCount == 1)
|
|
}
|
|
|
|
@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 dashboardOpenCount = 0
|
|
let view = OnboardingView(
|
|
state: appState,
|
|
aiSetupGateway: gateway,
|
|
systemAgentDefaults: defaults,
|
|
aiSetupRouteIdentityProvider: { "remote:direct:example.invalid" },
|
|
dashboardOnboardingOpener: { dashboardOpenCount += 1 })
|
|
|
|
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(dashboardOpenCount == 1)
|
|
#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(dashboardOpenCount == 1)
|
|
#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" },
|
|
dashboardOnboardingOpener: { 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 },
|
|
dashboardOnboardingOpener: { 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)
|
|
}
|
|
}
|