Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions Sources/OpenClawKit/GatewayChannel+Testing.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
4 changes: 4 additions & 0 deletions Sources/OpenClawKit/GatewayChannel.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -1349,6 +1350,9 @@ extension GatewayChannelActor {
}
var persistedRoles = Set<String>()
if let identity, !writes.isEmpty {
#if DEBUG
self.testDeviceTokenPersistenceStartedHandler?()
#endif
persistedRoles = await Self.persistDeviceTokens(
writes,
deviceId: identity.deviceId,
Expand Down
57 changes: 46 additions & 11 deletions Tests/OpenClawKitTests/AsyncTimeoutRaceTests.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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<Void, Never>?

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<Void, Never>) 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
Expand Down
6 changes: 3 additions & 3 deletions Tests/OpenClawKitTests/ChatGatewaySessionTransportTests.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand Down Expand Up @@ -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)
Expand All @@ -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)
Expand Down
19 changes: 14 additions & 5 deletions Tests/OpenClawKitTests/ChatLinkPreviewTests.swift
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
import Foundation
import ImageIO
import OpenClawKit
import Testing
import UniformTypeIdentifiers
@testable import OpenClawChatUI
Expand Down Expand Up @@ -212,21 +213,29 @@ 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,
hostPolicy: { _ in true },
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 {
Expand Down
84 changes: 21 additions & 63 deletions Tests/OpenClawKitTests/ChatViewModelSessionActionTests.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Expand Down Expand Up @@ -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
Expand Down
16 changes: 12 additions & 4 deletions Tests/OpenClawKitTests/FileTransferNodeCommandsTests.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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")
}
Expand Down
16 changes: 9 additions & 7 deletions Tests/OpenClawKitTests/GatewayChannelLifecycleTests.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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)))
Expand Down Expand Up @@ -545,22 +545,24 @@ 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")
} catch {
#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()
Expand Down
4 changes: 2 additions & 2 deletions Tests/OpenClawKitTests/GatewayConnectRecoveryTests.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down Expand Up @@ -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
}
Expand Down
Loading
Loading