diff --git a/apps/macos/Tests/OpenClawIPCTests/AppStateVoiceEarsTests.swift b/apps/macos/Tests/OpenClawIPCTests/AppStateVoiceEarsTests.swift index 5b8b2ee97a4b..113e403f054c 100644 --- a/apps/macos/Tests/OpenClawIPCTests/AppStateVoiceEarsTests.swift +++ b/apps/macos/Tests/OpenClawIPCTests/AppStateVoiceEarsTests.swift @@ -23,7 +23,10 @@ struct AppStateVoiceEarsTests { #expect(state.earBoostActive) - try await Task.sleep(for: .milliseconds(60)) + let deadline = ContinuousClock.now + .seconds(1) + while state.earBoostActive, ContinuousClock.now < deadline { + try await Task.sleep(for: .milliseconds(10)) + } #expect(!state.earBoostActive) } } diff --git a/apps/shared/OpenClawKit/Sources/OpenClawKit/GatewayNodeSession.swift b/apps/shared/OpenClawKit/Sources/OpenClawKit/GatewayNodeSession.swift index 40e7a324b75c..feb8bcde692d 100644 --- a/apps/shared/OpenClawKit/Sources/OpenClawKit/GatewayNodeSession.swift +++ b/apps/shared/OpenClawKit/Sources/OpenClawKit/GatewayNodeSession.swift @@ -160,7 +160,7 @@ public actor GatewayNodeSession { private var serverMethods: Set? private var serverCapabilities: Set? private var mainSessionKey: String? - private var snapshotWaiters: [CheckedContinuation] = [] + private var snapshotWaiters: [UUID: CheckedContinuation] = [:] private var snapshotReadyWaiters: [CheckedContinuation] = [] // `computer.act` is not safe to repeat after a response is lost. Keep recent // in-flight/results on the long-lived node session so a channel reconnect can @@ -979,12 +979,13 @@ extension GatewayNodeSession { return true } let clamped = max(0, timeoutMs) + let waiterID = UUID() return await withCheckedContinuation { cont in - self.snapshotWaiters.append(cont) + self.snapshotWaiters[waiterID] = cont Task { [weak self] in guard let self else { return } try? await Task.sleep(nanoseconds: UInt64(clamped) * 1_000_000) - await self.timeoutSnapshotWaiters() + await self.timeoutSnapshotWaiter(id: waiterID) } } } @@ -998,14 +999,16 @@ extension GatewayNodeSession { } } - private func timeoutSnapshotWaiters() { - guard !self.snapshotReceived else { return } - self.drainSnapshotWaiters(returning: false) + private func timeoutSnapshotWaiter(id: UUID) { + guard !self.snapshotReceived, + let waiter = self.snapshotWaiters.removeValue(forKey: id) + else { return } + waiter.resume(returning: false) } private func drainSnapshotWaiters(returning value: Bool) { if !self.snapshotWaiters.isEmpty { - let waiters = self.snapshotWaiters + let waiters = self.snapshotWaiters.values self.snapshotWaiters.removeAll() for waiter in waiters { waiter.resume(returning: value) @@ -1332,6 +1335,26 @@ extension GatewayNodeSession { self.broadcastServerEvent(event) } + // periphery:ignore - package tests reproduce a stale timeout across route reset. + func _test_waitForSnapshot(timeoutMs: Int) async -> Bool { + await self.waitForSnapshot(timeoutMs: timeoutMs) + } + + // periphery:ignore - package tests complete snapshot waits without a live socket. + func _test_markSnapshotReceived() { + self.markSnapshotReceived() + } + + // periphery:ignore - package tests reproduce route replacement snapshot state. + func _test_resetConnectionState() { + self.resetConnectionState() + } + + // periphery:ignore - package tests wait until a snapshot continuation is registered. + func _test_snapshotWaiterCount() -> Int { + self.snapshotWaiters.count + } + #endif private func cancelActiveInvokes( diff --git a/apps/shared/OpenClawKit/Tests/OpenClawKitTests/GatewayNodeSessionTests.swift b/apps/shared/OpenClawKit/Tests/OpenClawKitTests/GatewayNodeSessionTests.swift index 7a8d58847e47..646888f9f0b9 100644 --- a/apps/shared/OpenClawKit/Tests/OpenClawKitTests/GatewayNodeSessionTests.swift +++ b/apps/shared/OpenClawKit/Tests/OpenClawKitTests/GatewayNodeSessionTests.swift @@ -1320,6 +1320,30 @@ struct GatewayNodeSessionTests { await gateway.disconnect() } + @Test + func `completed snapshot timeout cannot release a later route waiter`() async throws { + let gateway = GatewayNodeSession() + let firstWait = Task { + await gateway._test_waitForSnapshot(timeoutMs: 1000) + } + try await waitUntil("initial snapshot waiter registered") { + await gateway._test_snapshotWaiterCount() == 1 + } + await gateway._test_markSnapshotReceived() + #expect(await firstWait.value) + + await gateway._test_resetConnectionState() + let replacementWait = Task { + await gateway._test_waitForSnapshot(timeoutMs: 3000) + } + try await waitUntil("replacement snapshot waiter registered") { + await gateway._test_snapshotWaiterCount() == 1 + } + try await Task.sleep(nanoseconds: 1_200_000_000) + await gateway._test_markSnapshotReceived() + #expect(await replacementWait.value) + } + @Test func `concurrent replacements wait for route invalidation before installing a channel`() async throws { let session = FakeGatewayWebSocketSession()