diff --git a/Sources/OpenClawKit/GatewayChannel+Testing.swift b/Sources/OpenClawKit/GatewayChannel+Testing.swift index 3e5c9f3..6d7f3ef 100644 --- a/Sources/OpenClawKit/GatewayChannel+Testing.swift +++ b/Sources/OpenClawKit/GatewayChannel+Testing.swift @@ -22,6 +22,11 @@ extension GatewayChannelActor { func _test_setRequestResumedHandler(_ handler: (@Sendable () async -> Void)?) { self.testRequestResumedHandler = handler } + + /// Called on the actor just before hello-ok's issued tokens are handed to the persistence hop. + func _test_setDeviceTokenPersistenceStartedHandler(_ handler: (@Sendable () -> Void)?) { + self.testDeviceTokenPersistenceStartedHandler = handler + } #endif func _test_pendingRequestCount() -> Int { diff --git a/Sources/OpenClawKit/GatewayChannel.swift b/Sources/OpenClawKit/GatewayChannel.swift index 283afad..ca38b24 100644 --- a/Sources/OpenClawKit/GatewayChannel.swift +++ b/Sources/OpenClawKit/GatewayChannel.swift @@ -94,6 +94,7 @@ public actor GatewayChannelActor { var testConnectRunFinishedHandler: (@Sendable () -> Void)? var testConnectFailureBackoffWaitHandler: (@Sendable () async throws -> Void)? var testRequestResumedHandler: (@Sendable () async -> Void)? + var testDeviceTokenPersistenceStartedHandler: (@Sendable () -> Void)? #endif private let connectChallengeTimeoutSeconds: Double = 6.0 // Some networks will silently drop idle TCP/TLS flows around ~30s. The gateway tick is server->client, @@ -1349,6 +1350,9 @@ extension GatewayChannelActor { } var persistedRoles = Set() if let identity, !writes.isEmpty { + #if DEBUG + self.testDeviceTokenPersistenceStartedHandler?() + #endif persistedRoles = await Self.persistDeviceTokens( writes, deviceId: identity.deviceId, diff --git a/Tests/OpenClawKitTests/AsyncTimeoutRaceTests.swift b/Tests/OpenClawKitTests/AsyncTimeoutRaceTests.swift index d8808e7..d090c77 100644 --- a/Tests/OpenClawKitTests/AsyncTimeoutRaceTests.swift +++ b/Tests/OpenClawKitTests/AsyncTimeoutRaceTests.swift @@ -21,21 +21,56 @@ private final class RaceCounter: @unchecked Sendable { } } +/// An operation that ignores cancellation entirely (the keepalive wedge): it parks on a checked +/// continuation until the test itself calls `open()`. +private final class RaceGate: @unchecked Sendable { + private let lock = NSLock() + private var isOpen = false + private var waiter: CheckedContinuation? + + func wait() async { + await withCheckedContinuation { continuation in + self.lock.lock() + guard !self.isOpen else { + self.lock.unlock() + continuation.resume() + return + } + self.waiter = continuation + self.lock.unlock() + } + } + + func open() { + self.lock.lock() + self.isOpen = true + let waiter = self.waiter + self.waiter = nil + self.lock.unlock() + waiter?.resume() + } +} + @Suite("AsyncTimeout race") struct AsyncTimeoutRaceTests { - @Test + /// Only the 50 ms deadline can end the race: the operation stays parked until the test opens its + /// gate. A race that joined its loser would hang and trip the time limit, whose cancellation opens + /// the gate so the test body can end. No wall-clock bound: a saturated test pool can stall the run + /// for seconds. + @Test(.timeLimit(.minutes(1))) func operationThatIgnoresCancellationStillTimesOut() async { - let start = ContinuousClock.now - await #expect(throws: RaceTimeoutError.self) { - try await AsyncTimeout.withTimeout( - seconds: 0.05, - onTimeout: { RaceTimeoutError() }, - operation: { - // A never-resumed continuation ignores cancellation entirely (the keepalive wedge). - await withCheckedContinuation { (_: CheckedContinuation) in } - }) + let gate = RaceGate() + defer { gate.open() } + await withTaskCancellationHandler { + await #expect(throws: RaceTimeoutError.self) { + try await AsyncTimeout.withTimeout( + seconds: 0.05, + onTimeout: { RaceTimeoutError() }, + operation: { await gate.wait() }) + } + } onCancel: { + gate.open() } - #expect(ContinuousClock.now - start < .seconds(5)) } @Test diff --git a/Tests/OpenClawKitTests/ChatGatewaySessionTransportTests.swift b/Tests/OpenClawKitTests/ChatGatewaySessionTransportTests.swift index 3249a1f..6b357b9 100644 --- a/Tests/OpenClawKitTests/ChatGatewaySessionTransportTests.swift +++ b/Tests/OpenClawKitTests/ChatGatewaySessionTransportTests.swift @@ -62,7 +62,7 @@ private func sentParams(_ socket: GatewayCoreFakeSocket, method: String) -> [[St socket.sentFrames(method: method).compactMap { $0["params"] as? [String: Any] } } -@Suite("Gateway session chat transport", .serialized) +@Suite("Gateway session chat transport", .serialized, .timeLimit(.minutes(1))) struct ChatGatewaySessionTransportTests { @Test func `init normalizes the agent and keeps exact gateway id bytes`() { let transport = OpenClawGatewaySessionChatTransport( @@ -384,7 +384,7 @@ struct ChatGatewaySessionTransportTests { // Same connection context, new socket: seqGap plus a fresh per-socket subscription. socket.emitReceiveFailure() - try await gatewayCoreWaitUntil("reconnect seqGap", timeoutSeconds: 15) { + try await gatewayCoreWaitUntil("reconnect seqGap") { recorder.values.contains("seqGap") } let reconnected = try #require(session.latestSocket) @@ -396,7 +396,7 @@ struct ChatGatewaySessionTransportTests { // A different endpoint is a different connection context: routeChanged. let replacement = transportTestSession(capabilities: []) try await gateway.connectForChatTransportTest("ws://replacement.example.invalid", session: replacement) - try await gatewayCoreWaitUntil("route change reported", timeoutSeconds: 15) { + try await gatewayCoreWaitUntil("route change reported") { recorder.values.contains("routeChanged") } let replacementSocket = try #require(replacement.latestSocket) diff --git a/Tests/OpenClawKitTests/ChatLinkPreviewTests.swift b/Tests/OpenClawKitTests/ChatLinkPreviewTests.swift index fae006b..cbf04ad 100644 --- a/Tests/OpenClawKitTests/ChatLinkPreviewTests.swift +++ b/Tests/OpenClawKitTests/ChatLinkPreviewTests.swift @@ -1,5 +1,6 @@ import Foundation import ImageIO +import OpenClawKit import Testing import UniformTypeIdentifiers @testable import OpenClawChatUI @@ -212,9 +213,17 @@ struct ChatLinkPreviewNetworkTests { #expect(ChatLinkPreviewStubURLProtocol.lastAcceptHeader == "text/html") } - @Test func `total deadline can fire before the session starts`() async throws { + /// The protocol never answers and URLSession's own timeouts sit past the time limit, so only the + /// fetcher's zero-second deadline can end the fetch in time. A deadline lost before `start` could + /// leave the fetch unresumed even on cancellation, so it is awaited through a no-deadline + /// `AsyncTimeout` race, which the time limit's cancellation always ends. No wall-clock bound: a + /// saturated test pool can stall the run for seconds. + @Test(.timeLimit(.minutes(1))) + func `total deadline can fire before the session starts`() async throws { let configuration = URLSessionConfiguration.ephemeral configuration.protocolClasses = [ChatLinkPreviewHangingURLProtocol.self] + configuration.timeoutIntervalForRequest = 3600 + configuration.timeoutIntervalForResource = 3600 let fetcher = ChatLinkPreviewFetcher( configuration: configuration, timeout: 0, @@ -222,11 +231,11 @@ struct ChatLinkPreviewNetworkTests { resolutionPolicy: { _ in true }, connectionPolicy: { _ in true }) let url = try #require(URL(string: "https://preview.test/slow")) - let clock = ContinuousClock() - let start = clock.now - #expect(await fetcher.fetch(url) == .failed) - #expect(start.duration(to: clock.now) < .seconds(1)) + let result = try await AsyncTimeout.withTimeout(seconds: 0, onTimeout: { CancellationError() }) { + await fetcher.fetch(url) + } + #expect(result == .failed) } @Test func `image fetch accepts only images and enforces its body cap`() async throws { diff --git a/Tests/OpenClawKitTests/ChatViewModelSessionActionTests.swift b/Tests/OpenClawKitTests/ChatViewModelSessionActionTests.swift index d0bcf9c..adbda1e 100644 --- a/Tests/OpenClawKitTests/ChatViewModelSessionActionTests.swift +++ b/Tests/OpenClawKitTests/ChatViewModelSessionActionTests.swift @@ -511,6 +511,7 @@ private struct BatchTestError: LocalizedError { } @MainActor +@Suite(.timeLimit(.minutes(1))) struct ChatViewModelSessionActionTests { @Test func `batch mutations continue after per-row failure with bounded fan-out`() async { let probe = BatchMutationProbe() @@ -1402,81 +1403,38 @@ struct ChatViewModelSessionActionTests { #expect(await transport.forkedParentKeys() == ["main"]) } - private func waitForForkStart( - _ gate: SessionActionCompletionGate, - timeout: Duration = .seconds(15)) async -> Bool - { - // The stream controls ordering; this deadline only bounds a broken fake or call path. - await withTaskGroup(of: Bool.self) { group in - group.addTask { await gate.waitUntilStarted() } - group.addTask { - try? await Task.sleep(for: timeout) - return false - } - let started = await group.next() ?? false - group.cancelAll() - return started - } + private func waitForForkStart(_ gate: SessionActionCompletionGate) async -> Bool { + // The stream controls ordering; the suite's time limit bounds a broken fake or call path, + // and its cancellation ends the stream wait with `false`. + await gate.waitUntilStarted() } - private func waitForBranchSwitchActivityToClear( - _ viewModel: OpenClawChatViewModel, - timeout: Duration = .seconds(15)) async -> Bool - { - let clock = ContinuousClock() - let deadline = clock.now + timeout - while clock.now < deadline { - if viewModel.hasBlockingRunActivity == false { - return true - } - await Task.yield() - } - return false + private func waitForBranchSwitchActivityToClear(_ viewModel: OpenClawChatViewModel) async -> Bool { + await self.eventually { viewModel.hasBlockingRunActivity == false } } - private func waitForOutboxRestore( - _ viewModel: OpenClawChatViewModel, - timeout: Duration = .seconds(15)) async -> Bool - { - let clock = ContinuousClock() - let deadline = clock.now + timeout - while clock.now < deadline { - if viewModel.hasRestoredOutboxMessages, - viewModel.hasPendingOutboxCommandsForCurrentSession - { - return true - } - await Task.yield() + private func waitForOutboxRestore(_ viewModel: OpenClawChatViewModel) async -> Bool { + await self.eventually { + viewModel.hasRestoredOutboxMessages && viewModel.hasPendingOutboxCommandsForCurrentSession } - return false } - private func waitForSend( - _ transport: SessionActionTransport, - timeout: Duration = .seconds(15)) async -> Bool - { - let clock = ContinuousClock() - let deadline = clock.now + timeout - while clock.now < deadline { - if await transport.sentSessionKeys().isEmpty == false { - return true - } - await Task.yield() - } - return false + private func waitForSend(_ transport: SessionActionTransport) async -> Bool { + await self.eventually { await transport.sentSessionKeys().isEmpty == false } } private func waitForBranchReload( _ viewModel: OpenClawChatViewModel, - branches: [OpenClawChatSessionBranch], - timeout: Duration = .seconds(15)) async -> Bool + branches: [OpenClawChatSessionBranch]) async -> Bool { - let clock = ContinuousClock() - let deadline = clock.now + timeout - while clock.now < deadline { - if viewModel.sessionBranches == branches, !viewModel.isLoading { - return true - } + await self.eventually { viewModel.sessionBranches == branches && !viewModel.isLoading } + } + + /// Yields until `condition` holds. There is no wall-clock deadline, so a saturated test pool only + /// slows the wait down; the suite's time limit cancels a real hang, which ends the wait with `false`. + private func eventually(_ condition: () async -> Bool) async -> Bool { + while !Task.isCancelled { + if await condition() { return true } await Task.yield() } return false diff --git a/Tests/OpenClawKitTests/FileTransferNodeCommandsTests.swift b/Tests/OpenClawKitTests/FileTransferNodeCommandsTests.swift index d33dc18..e84a465 100644 --- a/Tests/OpenClawKitTests/FileTransferNodeCommandsTests.swift +++ b/Tests/OpenClawKitTests/FileTransferNodeCommandsTests.swift @@ -217,15 +217,23 @@ struct FileTransferNodeCommandsTests { // MARK: - Drift between preflight and the final call - @Test func `file fetch refuses a FIFO promptly`() async throws { + /// A fetch that opened the FIFO without `O_NONBLOCK` would wait forever for a writer and trip the + /// time limit. That hang is a syscall, not an await, so the time limit's cancellation opens the write + /// end once to release the reader and let the test body end. No wall-clock bound: a saturated test + /// pool can stall the run for seconds. + @Test(.timeLimit(.minutes(1))) + func `file fetch refuses a FIFO promptly`() async throws { let sandbox = try FileTransferSandbox() defer { sandbox.cleanUp() } let fifo = sandbox.root + "/pipe" try #require(mkfifo(fifo, 0o600) == 0) let commands = sandbox.commands - let started = ContinuousClock.now - let result = try await self.invoke(commands, "file.fetch", ["path": fifo]) - #expect(ContinuousClock.now - started < .seconds(2)) + let result = try await withTaskCancellationHandler { + try await self.invoke(commands, "file.fetch", ["path": fifo]) + } onCancel: { + let writer = open(fifo, O_WRONLY | O_NONBLOCK) + if writer >= 0 { close(writer) } + } #expect(result["ok"] as? Bool == false) #expect(result["code"] as? String == "IS_DIRECTORY") } diff --git a/Tests/OpenClawKitTests/GatewayChannelLifecycleTests.swift b/Tests/OpenClawKitTests/GatewayChannelLifecycleTests.swift index ba546f2..2e05193 100644 --- a/Tests/OpenClawKitTests/GatewayChannelLifecycleTests.swift +++ b/Tests/OpenClawKitTests/GatewayChannelLifecycleTests.swift @@ -37,7 +37,7 @@ private func makeChannel( extraHeadersProvider: extraHeadersProvider) } -@Suite("Gateway channel lifecycle") +@Suite("Gateway channel lifecycle", .timeLimit(.minutes(1))) struct GatewayChannelLifecycleTests { @Test func operatorConnectOffersProtocolFourWithPlatformDefaults() async throws { @@ -496,15 +496,15 @@ struct GatewayChannelLifecycleTests { await channel.shutdown() } - @Test + /// The socket never pongs, so only the ping deadline can end the wait (with `URLError`); an unbounded + /// ping trips the time limit. No wall-clock bound: a saturated test pool can stall the run for seconds. + @Test(.timeLimit(.minutes(1))) func keepalivePingIsBoundedWhenNoPongArrives() async throws { let socket = GatewayCoreFakeSocket(script: GatewayCoreSocketScript(pingBehavior: .never)) let box = WebSocketTaskBox(task: socket) - let start = ContinuousClock.now await #expect(throws: URLError.self) { try await box.sendPing(timeout: .milliseconds(50)) } - #expect(ContinuousClock.now - start < .seconds(5)) let duplicate = WebSocketTaskBox(task: GatewayCoreFakeSocket(script: GatewayCoreSocketScript( pingBehavior: .duplicateSuccess))) @@ -545,14 +545,17 @@ struct GatewayChannelLifecycleTests { == .drop(.missingRecipientProfile)) } - @Test + /// The gateway never answers `connect`. Only the 100 ms option can time the handshake out within the + /// time limit: the fallback budget is an hour, so a channel that ignored the option would trip it. No + /// wall-clock bound: a saturated test pool can stall the run for seconds. + @Test(.timeLimit(.minutes(1))) func handshakeTimeoutOptionBoundsTheWholeHandshake() async throws { let session = GatewayCoreFakeSession(fixedScript: GatewayCoreSocketScript( connectReply: { _ in .none })) var options = gatewayCoreOptions() options.handshakeTimeoutMs = 100 let channel = try makeChannel(session: session, options: options) - let start = ContinuousClock.now + await channel._test_setConnectTimeoutSeconds(3600) do { try await channel.connect() Issue.record("expected a handshake timeout") @@ -560,7 +563,6 @@ struct GatewayChannelLifecycleTests { #expect((error as NSError).domain == NSURLErrorDomain) #expect((error as NSError).code == URLError.timedOut.rawValue) } - #expect(ContinuousClock.now - start < .seconds(5)) #expect(await channel.currentHandshakePhase() == .connectSent) #expect(session.latestSocket?.state != .running) await channel.shutdown() diff --git a/Tests/OpenClawKitTests/GatewayConnectRecoveryTests.swift b/Tests/OpenClawKitTests/GatewayConnectRecoveryTests.swift index ba7f8aa..7cc763a 100644 --- a/Tests/OpenClawKitTests/GatewayConnectRecoveryTests.swift +++ b/Tests/OpenClawKitTests/GatewayConnectRecoveryTests.swift @@ -37,7 +37,7 @@ private let pinMismatch = GatewayTLSValidationFailure( systemTrustOk: true, port: 443) -@Suite("Gateway connect recovery", .serialized) +@Suite("Gateway connect recovery", .serialized, .timeLimit(.minutes(1))) struct GatewayConnectRecoveryTests { // MARK: TLS pin mismatch @@ -232,7 +232,7 @@ struct GatewayConnectRecoveryTests { await channel.reconnectPauseReason() == .authFailure } #expect(session.makeCount == 2) - try await gatewayCoreWaitUntil("resumed after retryAfterMs", timeoutSeconds: 10) { + try await gatewayCoreWaitUntil("resumed after retryAfterMs") { guard session.makeCount == 3 else { return false } return await channel.currentConnectionGeneration() != nil } diff --git a/Tests/OpenClawKitTests/GatewayCoreTestSupport.swift b/Tests/OpenClawKitTests/GatewayCoreTestSupport.swift index e05807a..3c5c243 100644 --- a/Tests/OpenClawKitTests/GatewayCoreTestSupport.swift +++ b/Tests/OpenClawKitTests/GatewayCoreTestSupport.swift @@ -1,5 +1,6 @@ import Foundation import OpenClawProtocol +import Testing @testable import OpenClawKit // In-memory WebSocket transport for gateway channel and node session tests. It scripts @@ -392,18 +393,23 @@ struct GatewayCoreWaitTimeout: Error, CustomStringConvertible { } } +/// Polls `condition` until it holds. +/// +/// There is no wall-clock deadline: a saturated test pool must only slow the wait down. Every suite +/// that calls this carries a `.timeLimit`, whose cancellation ends the wait with +/// ``GatewayCoreWaitTimeout``. func gatewayCoreWaitUntil( _ label: String, - timeoutSeconds: Double = 10, _ condition: @escaping @Sendable () async -> Bool) async throws { - let deadline = ContinuousClock.now.advanced(by: .milliseconds(Int64(timeoutSeconds * 1000))) - while ContinuousClock.now < deadline { + while !Task.isCancelled { if await condition() { return } - try await Task.sleep(for: .milliseconds(5)) + try? await Task.sleep(for: .milliseconds(5)) } - if await condition() { return } - throw GatewayCoreWaitTimeout(label: label) + // Swift Testing drops errors thrown after a time-limit cancellation, so record which wait hung. + let timeout = GatewayCoreWaitTimeout(label: label) + Issue.record(timeout) + throw timeout } func gatewayCoreTemporaryStateDirectory() throws -> URL { diff --git a/Tests/OpenClawKitTests/GatewayDeviceAuthOffActorTests.swift b/Tests/OpenClawKitTests/GatewayDeviceAuthOffActorTests.swift index 47887c9..782a52a 100644 --- a/Tests/OpenClawKitTests/GatewayDeviceAuthOffActorTests.swift +++ b/Tests/OpenClawKitTests/GatewayDeviceAuthOffActorTests.swift @@ -35,8 +35,13 @@ private final class StateDatabaseWriteLock: @unchecked Sendable { } } -@Suite("Gateway device auth off the channel actor", .serialized) +@Suite("Gateway device auth off the channel actor", .serialized, .timeLimit(.minutes(1))) struct GatewayDeviceAuthOffActorTests { + /// Ordering, not elapsed time, proves the actor stayed responsive: the lock is released only after + /// the actor answers, so the token can reach disk only if the actor answered while the write was + /// still pending. A write on the actor would hold every call until SQLite's 30 s busy timeout fails + /// it; the token never lands and the final wait trips the time limit. The only timing left is that + /// the healthy write must see the release within the same 30 s busy timeout. @Test func issuedTokenPersistenceNeverBlocksTheChannelActor() async throws { let directory = try gatewayCoreTemporaryStateDirectory() @@ -59,15 +64,17 @@ struct GatewayDeviceAuthOffActorTests { token: "shared", session: WebSocketSessionBox(session: session), connectOptions: gatewayCoreOptions(includeDeviceIdentity: true)) + let (persistenceStarts, persistenceStarted) = AsyncStream.makeStream(of: Void.self) + await channel._test_setDeviceTokenPersistenceStartedHandler { persistenceStarted.yield() } let connect = Task { try await channel.connect() } - try await gatewayCoreWaitUntil("write lock held") { writeLock.isHeld } - try await Task.sleep(for: .milliseconds(200)) + // Wait until hello-ok's token is handed to the persistence hop, which shutdown cannot stop. + var starts = persistenceStarts.makeAsyncIterator() + guard await starts.next() != nil else { throw CancellationError() } + try #require(writeLock.isHeld) - // The token write waits on SQLite's 30 s busy timeout; the actor must stay responsive. - let started = ContinuousClock.now + // The write cannot finish while the lock is held; the actor must still answer. #expect(await channel.currentConnectionGeneration() == nil) await channel.shutdown() - #expect(ContinuousClock.now - started < .seconds(3)) writeLock.release() _ = try? await connect.value diff --git a/Tests/OpenClawKitTests/GatewayNetworkConnectionTransportTests.swift b/Tests/OpenClawKitTests/GatewayNetworkConnectionTransportTests.swift index 2a9ad51..8d4fc82 100644 --- a/Tests/OpenClawKitTests/GatewayNetworkConnectionTransportTests.swift +++ b/Tests/OpenClawKitTests/GatewayNetworkConnectionTransportTests.swift @@ -87,19 +87,28 @@ private final class NetworkConnectionLoopbackGateway: @unchecked Sendable { listener.newConnectionHandler = { [weak gateway] connection in gateway?.accept(connection) } + // Wait on the listener's own state updates rather than polling against a wall-clock deadline, + // which a saturated test pool can overrun even though the listener is long ready. The suite's + // time limit bounds a real hang: cancellation ends the stream iteration below. + let (states, stateContinuation) = AsyncStream.makeStream(of: NWListener.State.self) + defer { stateContinuation.finish() } + listener.stateUpdateHandler = { stateContinuation.yield($0) } listener.start(queue: gateway.queue) - let deadline = ContinuousClock.now.advanced(by: .seconds(5)) - while ContinuousClock.now < deadline { - if case .ready = listener.state, let port = listener.port, port.rawValue != 0 { + for await state in states { + switch state { + case .ready: return gateway - } - if case let .failed(error) = listener.state { + case let .failed(error): + listener.cancel() throw error + case .cancelled: + throw CancellationError() + default: + continue } - try await Task.sleep(for: .milliseconds(10)) } listener.cancel() - throw URLError(.timedOut) + throw CancellationError() } var url: URL { @@ -177,7 +186,7 @@ private final class NetworkConnectionLoopbackGateway: @unchecked Sendable { } } -@Suite("Network.framework gateway transport", .serialized, .gatewayTLSStoreIsolated) +@Suite("Network.framework gateway transport", .serialized, .gatewayTLSStoreIsolated, .timeLimit(.minutes(1))) struct GatewayNetworkConnectionTransportTests { @Test func handshakeAndRequestsRunOverNetworkConnection() async throws { diff --git a/Tests/OpenClawKitTests/GatewayNodeSessionRouteTests.swift b/Tests/OpenClawKitTests/GatewayNodeSessionRouteTests.swift index 8a36913..3603f8b 100644 --- a/Tests/OpenClawKitTests/GatewayNodeSessionRouteTests.swift +++ b/Tests/OpenClawKitTests/GatewayNodeSessionRouteTests.swift @@ -148,7 +148,7 @@ private func respondToSurfaceRefresh( reply: .ok(["surface": "canvas", "pluginSurfaceUrls": ["canvas": url]])) } -@Suite("Gateway node session routes", .serialized) +@Suite("Gateway node session routes", .serialized, .timeLimit(.minutes(1))) struct GatewayNodeSessionRouteTests { @Test func invokeMetadataReachesTheHandlerAndResultCarriesStructuredPayload() async throws { diff --git a/Tests/OpenClawKitTests/GatewayRequestBudgetTests.swift b/Tests/OpenClawKitTests/GatewayRequestBudgetTests.swift index d915160..6313210 100644 --- a/Tests/OpenClawKitTests/GatewayRequestBudgetTests.swift +++ b/Tests/OpenClawKitTests/GatewayRequestBudgetTests.swift @@ -10,7 +10,7 @@ private func budgetChannel(session: GatewayCoreFakeSession) throws -> GatewayCha connectOptions: gatewayCoreOptions()) } -@Suite("Gateway request budget", .serialized) +@Suite("Gateway request budget", .serialized, .timeLimit(.minutes(1))) struct GatewayRequestBudgetTests { @Test func slowConnectDoesNotConsumeTheRequestBudget() async throws { diff --git a/Tests/OpenClawKitTests/GatewayStateReportingWiringTests.swift b/Tests/OpenClawKitTests/GatewayStateReportingWiringTests.swift index dbca099..33f2394 100644 --- a/Tests/OpenClawKitTests/GatewayStateReportingWiringTests.swift +++ b/Tests/OpenClawKitTests/GatewayStateReportingWiringTests.swift @@ -21,7 +21,7 @@ private func gatewayLabels(_ reporter: RecordingStateReporter) -> [String?] { reporter.transitions.filter { $0.domain == .gateway }.map(\.label) } -@Suite("Gateway state reporting wiring", .serialized) +@Suite("Gateway state reporting wiring", .serialized, .timeLimit(.minutes(1))) struct GatewayStateReportingWiringTests { @Test func connectReportsTheHandshakeLifecycleWithStableContext() async throws { diff --git a/Tests/OpenClawKitTests/GatewayTLSPinRotationRecoveryTests.swift b/Tests/OpenClawKitTests/GatewayTLSPinRotationRecoveryTests.swift index 85dc577..36c6a3c 100644 --- a/Tests/OpenClawKitTests/GatewayTLSPinRotationRecoveryTests.swift +++ b/Tests/OpenClawKitTests/GatewayTLSPinRotationRecoveryTests.swift @@ -16,7 +16,7 @@ private func mismatch(storeKey: String) -> GatewayTLSValidationFailure { port: 443) } -@Suite("Gateway TLS pin rotation recovery", .serialized, .gatewayTLSStoreIsolated) +@Suite("Gateway TLS pin rotation recovery", .serialized, .gatewayTLSStoreIsolated, .timeLimit(.minutes(1))) struct GatewayTLSPinRotationRecoveryTests { @Test func pinningSessionAcceptsAReviewedRotationInPlace() throws { diff --git a/Tests/OpenClawKitTests/OpenClawAppIntentsRunMatchingTests.swift b/Tests/OpenClawKitTests/OpenClawAppIntentsRunMatchingTests.swift index 01f4814..c64b29f 100644 --- a/Tests/OpenClawKitTests/OpenClawAppIntentsRunMatchingTests.swift +++ b/Tests/OpenClawKitTests/OpenClawAppIntentsRunMatchingTests.swift @@ -5,14 +5,20 @@ import OpenClawKit /// Gateway transport whose `chat.send` stays pending until the test acknowledges it. private actor HeldChatSendRequester: OpenClawIntentGatewayRequesting { - private(set) var sendParams: [String: AnyCodable]? + /// Params of each `chat.send`, yielded the moment the request arrives. + nonisolated let sends: AsyncStream<[String: AnyCodable]> + private let sendsContinuation: AsyncStream<[String: AnyCodable]>.Continuation private(set) var abortParams: [[String: AnyCodable]] = [] private var pendingSend: CheckedContinuation? + init() { + (self.sends, self.sendsContinuation) = AsyncStream.makeStream() + } + func request(method: String, params: [String: AnyCodable]?, timeoutMs: Double?) async throws -> Data { switch method { case "chat.send": - self.sendParams = params ?? [:] + self.sendsContinuation.yield(params ?? [:]) return try await withCheckedThrowingContinuation { self.pendingSend = $0 } case "chat.abort": self.abortParams.append(params ?? [:]) @@ -26,26 +32,25 @@ private actor HeldChatSendRequester: OpenClawIntentGatewayRequesting { self.pendingSend?.resume(returning: Data(json.utf8)) self.pendingSend = nil } - - var idempotencyKey: String? { - self.sendParams?["idempotencyKey"]?.stringValue - } } private func chatEvent(_ fields: [String: String]) -> EventFrame { EventFrame(type: "event", event: "chat", payload: AnyCodable(fields.mapValues { AnyCodable($0) })) } +/// Returns the idempotency key of the first `chat.send` once it reaches the requester. +/// +/// Event-driven rather than deadline-polled: under a saturated test pool the send can take +/// seconds to arrive, which must only slow the test down. The suite's time limit bounds a real hang. private func waitForSend(_ requester: HeldChatSendRequester) async throws -> String { - let deadline = ContinuousClock.now.advanced(by: .seconds(5)) - while ContinuousClock.now < deadline { - if let key = await requester.idempotencyKey { return key } - try await Task.sleep(for: .milliseconds(5)) + for await params in requester.sends { + return try #require(params["idempotencyKey"]?.stringValue) } + // The stream is never finished, so iteration only ends when the time limit cancels the test. throw CancellationError() } -@Suite("App Intents run matching") +@Suite("App Intents run matching", .timeLimit(.minutes(1))) struct OpenClawAppIntentsRunMatchingTests { @Test func anotherAgentsRunBeforeTheAckNeverHijacksTheIntent() async throws { diff --git a/Tests/OpenClawKitTests/WatchNodeClientTests.swift b/Tests/OpenClawKitTests/WatchNodeClientTests.swift index f97f3a3..0d8776d 100644 --- a/Tests/OpenClawKitTests/WatchNodeClientTests.swift +++ b/Tests/OpenClawKitTests/WatchNodeClientTests.swift @@ -196,7 +196,7 @@ private final class InMemoryWatchNodeConfigurationStore: OpenClawWatchNodeConfig } } -@Suite(.serialized) +@Suite(.serialized, .timeLimit(.minutes(1))) struct WatchNodeClientTests { private static let nowMs = Int64(1_800_000_000_000) @@ -223,12 +223,10 @@ struct WatchNodeClientTests { sentAtMs: sentAtMs) } - private static func eventually( - timeout: Duration = .seconds(5), - _ condition: @Sendable () async -> Bool) async -> Bool - { - let deadline = ContinuousClock.now + timeout - while ContinuousClock.now < deadline { + /// Polls until `condition` holds. There is no wall-clock deadline, so a saturated test pool only + /// slows the wait down; the suite's time limit cancels a real hang, which ends the wait with `false`. + private static func eventually(_ condition: @Sendable () async -> Bool) async -> Bool { + while !Task.isCancelled { if await condition() { return true } try? await Task.sleep(for: .milliseconds(5)) }