Skip to content
Merged
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
28 changes: 20 additions & 8 deletions Tests/OpenClawKitE2ETests/ChannelAdaptersE2ETests.swift
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,7 @@ import FoundationNetworking
import Testing
@testable import OpenClawKit

@Suite("Channel adapters E2E")
@Suite("Channel adapters E2E", .timeLimit(.minutes(1)))
struct ChannelAdaptersE2ETests {
actor TelegramInboundCollector {
private(set) var messages: [InboundMessage] = []
Expand Down Expand Up @@ -115,14 +115,26 @@ struct ChannelAdaptersE2ETests {
#expect(sent.first?.text == "pong")
}

struct WaitTimeout: Error, CustomStringConvertible {
let label: String
var description: String {
"Timeout waiting for: \(self.label)"
}
}

/// Polls a condition instead of sleeping a fixed interval (fixed sleeps flaked under load).
static func waitFor(timeoutSeconds: Double = 15, _ condition: @escaping @Sendable () async -> Bool) async throws {
let deadline = Date().addingTimeInterval(timeoutSeconds)
while Date() < deadline {
///
/// There is no wall-clock deadline either: a saturated test pool must only slow the wait down. The
/// suite's `.timeLimit` cancels a real hang, which ends the wait with ``WaitTimeout``.
static func waitFor(_ label: String, _ condition: @escaping @Sendable () async -> Bool) async throws {
while !Task.isCancelled {
if await condition() { return }
try await Task.sleep(nanoseconds: 10_000_000)
try? await Task.sleep(nanoseconds: 10_000_000)
}
#expect(await condition(), "timed out waiting for condition")
// Swift Testing drops errors thrown after a time-limit cancellation, so record which wait hung.
let timeout = WaitTimeout(label: label)
Issue.record(timeout)
throw timeout
}

@Test
Expand Down Expand Up @@ -150,7 +162,7 @@ struct ChannelAdaptersE2ETests {
await collector1.append(inbound)
}
try await adapter1.start()
try await Self.waitFor { await !collector1.snapshot().isEmpty }
try await Self.waitFor("first run delivered") { await !collector1.snapshot().isEmpty }
await adapter1.stop()

let collector2 = TelegramInboundCollector()
Expand All @@ -164,7 +176,7 @@ struct ChannelAdaptersE2ETests {
await collector2.append(inbound)
}
try await adapter2.start()
try await Self.waitFor { await !collector2.snapshot().isEmpty }
try await Self.waitFor("second run delivered") { await !collector2.snapshot().isEmpty }
await adapter2.stop()

let firstRun = await collector1.snapshot()
Expand Down
9 changes: 5 additions & 4 deletions Tests/OpenClawKitE2ETests/GatewayTransportE2ETests.swift
Original file line number Diff line number Diff line change
Expand Up @@ -250,7 +250,9 @@ struct GatewayTransportE2ETests {
#expect(counter.get() == baseline)
}

@Test
/// No wall-clock bound on the reconnect wait: a saturated test pool can stall the run for seconds. The
/// time limit ends the wait if the client stops reconnecting.
@Test(.timeLimit(.minutes(1)))
func reconnectFailureSchedulesAnotherAttempt() async throws {
let counter = Counter()
let client = GatewayClient(
Expand All @@ -272,9 +274,8 @@ struct GatewayTransportE2ETests {
try await client.connect(to: GatewayEndpoint(url: URL(string: "ws://127.0.0.1:18789")!))
// Wait for the third socket (stale tick -> reconnect -> connect failure -> another attempt)
// instead of a fixed 220 ms, which slower CI runners can overrun.
let deadline = Date().addingTimeInterval(10)
while counter.get() <= 2, Date() < deadline {
try await Task.sleep(nanoseconds: 10_000_000)
while counter.get() <= 2, !Task.isCancelled {
try? await Task.sleep(nanoseconds: 10_000_000)
}
#expect(counter.get() > 2)
await client.disconnect()
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
2 changes: 1 addition & 1 deletion Tests/OpenClawKitTests/ChannelDeliveryHardeningTests.swift
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,7 @@ import Testing

/// 2026.3.0 FX6 regressions: partial multi-part delivery, ambiguous 5xx, remote numeric input,
/// iMessage post-write failures and media budgets.
@Suite("Channel delivery hardening")
@Suite("Channel delivery hardening", .timeLimit(.minutes(1)))
struct ChannelDeliveryHardeningTests {
actor OffsetStore: TelegramUpdateOffsetStore {
func readLastUpdateID() async -> Int64? { nil }
Expand Down
2 changes: 1 addition & 1 deletion Tests/OpenClawKitTests/ChannelErrorRedactionTests.swift
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,7 @@ import OpenClawCore
import Testing

/// Credential redaction for channel health, `channels.status`, diagnostics and delivery failures.
@Suite("Channel error redaction")
@Suite("Channel error redaction", .timeLimit(.minutes(1)))
struct ChannelErrorRedactionTests {
actor OffsetStore: TelegramUpdateOffsetStore {
func readLastUpdateID() async -> Int64? { nil }
Expand Down
2 changes: 1 addition & 1 deletion Tests/OpenClawKitTests/ChannelPairingStoreTests.swift
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@ import Foundation
import Testing
@testable import OpenClawChannels

@Suite("Channel DM pairing store")
@Suite("Channel DM pairing store", .timeLimit(.minutes(1)))
struct ChannelPairingStoreTests {
final class Clock: @unchecked Sendable {
private let lock = NSLock()
Expand Down
1 change: 1 addition & 0 deletions Tests/OpenClawKitTests/ChatComposerParityTests.swift
Original file line number Diff line number Diff line change
Expand Up @@ -80,6 +80,7 @@ struct ChatReplyQuoteTests {
}

@MainActor
@Suite(.timeLimit(.minutes(1)))
struct ChatComposerStateTests {
@Test func `model selection target describes only gateway owned values`() {
let viewModel = OpenClawChatViewModel(sessionKey: "main", transport: ComposerParityTransport())
Expand Down
1 change: 1 addition & 0 deletions Tests/OpenClawKitTests/ChatComposerShellTests.swift
Original file line number Diff line number Diff line change
Expand Up @@ -264,6 +264,7 @@ struct ChatPrivateCloudQuotaNoticeTests {
}

@MainActor
@Suite(.timeLimit(.minutes(1)))
struct ChatSessionPagingTests {
@Test func `paging status summarizes truncated lists only`() {
#expect(ChatSessionPagingStatus(loadedCount: 10, totalCount: 10, isTruncated: false) == nil)
Expand Down
2 changes: 1 addition & 1 deletion Tests/OpenClawKitTests/ChatCoreCompatibilityTests.swift
Original file line number Diff line number Diff line change
Expand Up @@ -140,7 +140,7 @@ private func bootstrappedViewModel(_ transport: ScriptedEventTransport) async th
return viewModel
}

@Suite("Chat core 2026.3.0 compatibility")
@Suite("Chat core 2026.3.0 compatibility", .timeLimit(.minutes(1)))
struct ChatCoreCompatibilityTests {
// MARK: Legacy transport bridge

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
1 change: 1 addition & 0 deletions Tests/OpenClawKitTests/ChatHapticsTests.swift
Original file line number Diff line number Diff line change
Expand Up @@ -92,6 +92,7 @@ private func sendHapticsTestMessage(_ viewModel: OpenClawChatViewModel) async {
}
}

@Suite(.timeLimit(.minutes(1)))
struct ChatHapticsTests {
@Test func `send acceptance fires message sent exactly once`() async throws {
let (_, viewModel, recorder) = await makeHapticsViewModel(status: "started")
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
Original file line number Diff line number Diff line change
Expand Up @@ -63,6 +63,7 @@ private actor SidebarPreviewCache: OpenClawChatTranscriptCache {
}

@MainActor
@Suite(.timeLimit(.minutes(1)))
struct ChatSessionSidebarPreviewsTests {
@Test(arguments: [false, true])
func `changing Gateway owners never reuses an identical session preview`(oldLoadPending: Bool) async throws {
Expand Down
1 change: 1 addition & 0 deletions Tests/OpenClawKitTests/ChatStreamReplayTests.swift
Original file line number Diff line number Diff line change
Expand Up @@ -372,6 +372,7 @@ Closing paragraph with unicode — dashes, émojis 🦀🚀, and a trailing line
/// Covers streaming accumulation, provisional-final reconciliation against durable
/// `session.message` rows, duplicate delivery, out-of-order arrival, and reconnect
/// convergence. Tracking: #100196.
@Suite(.timeLimit(.minutes(1)))
struct ChatStreamReplayTests {
@Test func `live session message marker produces a visible transcript row`() async throws {
let harness = try await StreamReplayHarness.bootstrapped()
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -250,6 +250,7 @@ private final class AgentNavigationFixture {
}

@MainActor
@Suite(.timeLimit(.minutes(1)))
struct ChatViewModelAgentNavigationTests {
private func globalSession(owner: String) -> OpenClawChatSessionEntry {
var entry = OpenClawChatSessionEntry.placeholder(key: "global")
Expand Down
1 change: 1 addition & 0 deletions Tests/OpenClawKitTests/ChatViewModelAttachmentTests.swift
Original file line number Diff line number Diff line change
Expand Up @@ -261,6 +261,7 @@ private func chatAttachmentDimensions(for data: Data) -> (width: Int, height: In
return (width.intValue, height.intValue)
}

@Suite(.timeLimit(.minutes(1)))
struct ChatViewModelAttachmentTests {
@Test func imageAttachmentsAreProcessedBeforeStaging() async throws {
let imageData = try makeChatAttachmentJPEG(width: 3000, height: 4000)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@ private actor SettingsPatchCounter {
}
}

@Suite(.timeLimit(.minutes(1)))
struct ChatViewModelOutboxSettingsTests {
@Test func `background replay uses its command owned session settings`() async throws {
let (store, _, databaseDirectory) = try makeOutboxStore()
Expand Down
6 changes: 3 additions & 3 deletions Tests/OpenClawKitTests/ChatViewModelOutboxTests.swift
Original file line number Diff line number Diff line change
Expand Up @@ -832,7 +832,7 @@ actor ScriptedOutbox: OpenClawChatCommandOutbox {
}

// Serialized: every case opens a real SQLite outbox (see ChatTranscriptCacheStoreTests).
@Suite(.serialized)
@Suite(.serialized, .timeLimit(.minutes(1)))
struct ChatViewModelOutboxTests {
@Test func `offline send queues durably and renders queued row`() async throws {
let (store, _, databaseDirectory) = try makeOutboxStore()
Expand Down Expand Up @@ -1714,7 +1714,7 @@ struct ChatViewModelOutboxTests {
vm.messages.first { vm.outboxState(for: $0.id)?.isFailed == true }?.id
})
await MainActor.run { vm.retryOutboxMessage(failedMessageID) }
try await waitUntil("retried command drained", timeoutSeconds: 30) {
try await waitUntil("retried command drained") {
await store.loadCommands().isEmpty
}
#expect(await transport.state.sentIdempotencyKeys.count == 1)
Expand Down Expand Up @@ -1856,7 +1856,7 @@ struct ChatViewModelOutboxTests {

let messageID = try #require(await MainActor.run { vm.messages.last?.id })
await MainActor.run { vm.retryOutboxMessage(messageID) }
try await waitUntil("explicit retry drained", timeoutSeconds: 10) {
try await waitUntil("explicit retry drained") {
await store.loadCommands().isEmpty
}
#expect(await transport.state.sentIdempotencyKeys == [preserved.id])
Expand Down
Loading
Loading