Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
16 commits
Select commit Hold shift + click to select a range
3935801
test: wait for events instead of a 5 s wall-clock deadline
MarcoDotIO Sep 30, 2026
f8eb5fb
test: bound lifecycle timeouts with a time limit, not a 5 s wall clock
MarcoDotIO Sep 30, 2026
79c50c2
test: show the device-auth write leaves the actor free by ordering, n…
MarcoDotIO Sep 30, 2026
bccaaac
test: bound the FIFO, AsyncTimeout and link-preview tests with a time…
MarcoDotIO Sep 30, 2026
ffeef5c
test: drop wall-clock deadlines from the gateway and session-action w…
MarcoDotIO Sep 30, 2026
470d25e
Merge remote-tracking branch 'origin/main' into test/wait-until-no-de…
MarcoDotIO Sep 30, 2026
f5bee2a
test: drop the wall-clock deadline from waitUntil
MarcoDotIO Sep 30, 2026
d5620d1
test(e2e): wait without a wall-clock deadline in the channel and reco…
MarcoDotIO Sep 30, 2026
45849eb
test: drop the 1 s deadline from the chat UI view-model waits
MarcoDotIO Sep 30, 2026
7661599
test: drop the wall-clock waits left in the provider and realtime rel…
MarcoDotIO Sep 30, 2026
3b3466e
test(linux): wait without a wall-clock deadline in the runtime tests
MarcoDotIO Sep 30, 2026
ba41908
test: drop the bounded polls, elapsed bounds and positive-path timeouts
MarcoDotIO Sep 30, 2026
07e589b
test: wait for agent runs without a positive-path timeout
MarcoDotIO Sep 30, 2026
5912d21
Merge remote-tracking branch 'origin/main' into test/agent-wait-no-ti…
MarcoDotIO Sep 30, 2026
e17d424
Merge remote-tracking branch 'origin/main' into test/bounded-polls-an…
MarcoDotIO Oct 1, 2026
f90e771
Merge branch 'test/bounded-polls-and-elapsed-asserts' into test/agent…
MarcoDotIO Oct 1, 2026
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
6 changes: 2 additions & 4 deletions Tests/OpenClawKitTests/CoreAIModelRuntimeTests.swift
Original file line number Diff line number Diff line change
Expand Up @@ -296,17 +296,15 @@ struct CoreAIModelRuntimeTests {
#expect(texts == ["aaaaa", "bbbbb"])
}

@Test
@Test(.timeLimit(.minutes(1)))
func cancelStopsOnlyTheRunningGeneration() async throws {
let executor = SharedCacheExecutor(limit: 40)
let engine = CoreAILocalModelEngine(tokenizer: ScalarTokenizer(), executorFactory: { _, _ in executor })
var configuration = LocalModelConfig(enabled: true, runtime: CoreAILocalModelEngine.runtimeID, modelPath: "/m")
configuration.temperature = 0
try await engine.loadModel(path: "/m", configuration: configuration)
let running = Task { try await engine.generate(prompt: "a", systemPrompt: nil, configuration: configuration, onToken: nil) }
for _ in 0..<500 where await executor.runs == 0 {
try await Task.sleep(nanoseconds: 1_000_000)
}
try await waitUntil("first generation running") { await executor.runs > 0 }
let queued = Task { try await engine.generate(prompt: "b", systemPrompt: nil, configuration: configuration, onToken: nil) }
await engine.cancelGeneration(token: nil)
await #expect(throws: CoreAIRuntimeError.cancelled) {
Expand Down
4 changes: 2 additions & 2 deletions Tests/OpenClawKitTests/EmbeddedAgentStackTests.swift
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@ import Foundation
import Testing
@testable import OpenClawKit

@Suite("Embedded agent stack")
@Suite("Embedded agent stack", .timeLimit(.minutes(1)))
struct EmbeddedAgentStackTests {
struct ToolCallingProvider: ModelProvider {
let id = "stack"
Expand Down Expand Up @@ -32,7 +32,7 @@ struct EmbeddedAgentStackTests {
RequestFrame(type: "req", id: "1", method: "sessions.send", params: AnyCodable(["key": AnyCodable("agent:main:main"), "message": AnyCodable("hi")]))
)
let runID = try #require(sent.payload?.dictionaryValue?["runId"]?.stringValue)
#expect(await stack.runtime.wait(runID: runID, timeoutMs: 5_000)?.output == "stack says hi")
#expect(try await awaitCancellable("sent run finished") { await stack.runtime.wait(runID: runID) }?.output == "stack says hi")

let history = await stack.server.handle(
RequestFrame(type: "req", id: "2", method: "chat.history", params: AnyCodable(["sessionKey": AnyCodable("agent:main:main")]))
Expand Down
11 changes: 4 additions & 7 deletions Tests/OpenClawKitTests/GatewayServerChatUICompatTests.swift
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,7 @@ import Testing

/// The in-process server's session rows, `sessions.changed`, `session.message` and protocol-v4
/// `chat` payloads decode through the ChatUI transport models.
@Suite("Gateway server ChatUI compatibility")
@Suite("Gateway server ChatUI compatibility", .timeLimit(.minutes(1)))
struct GatewayServerChatUICompatTests {
private func makeStack(_ name: String, turns: [String]) async throws -> (GatewayServer, EmbeddedAgentRuntime, URL) {
let root = FileManager.default.temporaryDirectory.appendingPathComponent("\(name)-\(UUID().uuidString)")
Expand Down Expand Up @@ -43,11 +43,8 @@ struct GatewayServerChatUICompatTests {
private(set) var frames: [EventFrame] = []
func record(_ frame: EventFrame) { self.frames.append(frame) }

func wait(_ predicate: @Sendable ([EventFrame]) -> Bool) async -> [EventFrame] {
for _ in 0..<500 {
if predicate(self.frames) { return self.frames }
try? await Task.sleep(nanoseconds: 10_000_000)
}
func wait(_ label: String, _ predicate: @escaping @Sendable ([EventFrame]) -> Bool) async throws -> [EventFrame] {
try await waitUntil(label) { predicate(await self.frames) }
return self.frames
}
}
Expand Down Expand Up @@ -101,7 +98,7 @@ struct GatewayServerChatUICompatTests {
let send = try GatewayPayloadCodec.decode(OpenClawChatSendResponse.self, from: ack.payload)
#expect(send.runId == "ui-run-1")

let frames = await recorder.wait { frames in
let frames = try await recorder.wait("final chat, both session messages and the end change") { frames in
frames.contains { $0.event == "chat" && $0.payload?.dictionaryValue?["state"] == AnyCodable("final") }
&& frames.filter { $0.event == "session.message" }.count >= 2
&& frames.contains { $0.event == "sessions.changed" && $0.payload?.dictionaryValue?["phase"] == AnyCodable("end") }
Expand Down
10 changes: 7 additions & 3 deletions Tests/OpenClawKitTests/GatewayServerRegistryTests.swift
Original file line number Diff line number Diff line change
Expand Up @@ -404,7 +404,7 @@ struct GatewayServerRegistryTests {
#expect(await reloaded.listSecretKeys() == ["TOKEN", "legacy"])
}

@Test
@Test(.timeLimit(.minutes(1)))
func builtinHandlersAcceptUpstreamWireShapes() async throws {
let root = try Self.makeTempDirectory(named: "upstream-wire")
defer { try? FileManager.default.removeItem(at: root) }
Expand Down Expand Up @@ -461,10 +461,14 @@ struct GatewayServerRegistryTests {
#expect(missing.error?.errorCode == .invalidRequest)

// agent.wait accepts the upstream `runId` key and the legacy `runID` key.
let upstreamWait = await server.handle(Self.frame("agent.wait", params: ["runId": AnyCodable("run-main"), "timeoutMs": AnyCodable(1_000)]))
let upstreamWait = try await awaitCancellable("upstream-keyed run finished") {
await server.handle(Self.frame("agent.wait", params: ["runId": AnyCodable("run-main")]))
}
#expect(upstreamWait.payload?.dictionaryValue?["status"] == AnyCodable("ok"))
#expect(upstreamWait.payload?.dictionaryValue?["runId"] == AnyCodable("run-main"))
let legacyWait = await server.handle(Self.frame("agent.wait", params: ["runID": AnyCodable("run-legacy")]))
let legacyWait = try await awaitCancellable("legacy-keyed run finished") {
await server.handle(Self.frame("agent.wait", params: ["runID": AnyCodable("run-legacy")]))
}
#expect(try GatewayPayloadCodec.decode(GatewayAgentWaitResult.self, from: legacyWait.payload).output == "hi")

// sessions.patch accepts upstream `agentId`/`model`/`fastMode` and null clears.
Expand Down
46 changes: 28 additions & 18 deletions Tests/OpenClawKitTests/GatewayServerTests.swift
Original file line number Diff line number Diff line change
Expand Up @@ -149,7 +149,7 @@ struct GatewayServerTests {
#expect(deletedSession.deleted == true)
}

@Test
@Test(.timeLimit(.minutes(1)))
func gatewayServerSupportsAgentRunWaitTimeoutAndCleanup() async throws {
let root = try self.makeTempDirectory(named: "gateway-server-agent")
defer { try? FileManager.default.removeItem(at: root) }
Expand Down Expand Up @@ -191,12 +191,14 @@ struct GatewayServerTests {
)
#expect(!accepted.runID.isEmpty)

let waited = try await self.request(
server,
method: "agent.wait",
params: GatewayAgentWaitParams(runID: accepted.runID, timeoutMs: 5_000),
as: GatewayAgentWaitResult.self
)
let waited = try await awaitCancellable("first run finished") {
try await self.request(
server,
method: "agent.wait",
params: GatewayAgentWaitParams(runID: accepted.runID),
as: GatewayAgentWaitResult.self
)
}
#expect(waited.status == "ok")
#expect(waited.output == "first")

Expand Down Expand Up @@ -227,12 +229,14 @@ struct GatewayServerTests {
)
#expect(timedOut.status == "timeout")

let eventuallyCompleted = try await self.request(
server,
method: "agent.wait",
params: GatewayAgentWaitParams(runID: slow.runID, timeoutMs: 1_000),
as: GatewayAgentWaitResult.self
)
let eventuallyCompleted = try await awaitCancellable("slow run finished") {
try await self.request(
server,
method: "agent.wait",
params: GatewayAgentWaitParams(runID: slow.runID),
as: GatewayAgentWaitResult.self
)
}
#expect(eventuallyCompleted.status == "ok")
#expect(eventuallyCompleted.output == "slow")
}
Expand Down Expand Up @@ -741,7 +745,7 @@ struct GatewayServerTests {
}
}

@Test
@Test(.timeLimit(.minutes(1)))
func sdkGatewayServerSupportsCatalogSkillsAndRuntimeHandlers() async throws {
let root = try self.makeTempDirectory(named: "gateway-server-sdk")
defer { try? FileManager.default.removeItem(at: root) }
Expand Down Expand Up @@ -812,10 +816,16 @@ struct GatewayServerTests {
modelProviderID: "openai"
)
)
let completed: GatewayAgentWaitResult = try await client.request(
"agent.wait",
params: GatewayAgentWaitParams(runID: accepted.runID, timeoutMs: 5_000)
)
// The client's own request timeout is pushed past the time limit too, so neither side can
// expire while the pool is stalled.
let completed = try await awaitCancellable("runtime run finished") {
try await client.request(
"agent.wait",
params: GatewayAgentWaitParams(runID: accepted.runID),
timeoutMs: 3_600_000,
as: GatewayAgentWaitResult.self
)
}
#expect(completed.status == "ok")
#expect(completed.output == "runtime-output")
}
Expand Down
26 changes: 15 additions & 11 deletions Tests/OpenClawKitTests/MediaPipelineTests.swift
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,7 @@ import Glibc
import Darwin
#endif

@Suite("Media pipeline", .serialized)
@Suite("Media pipeline", .serialized, .timeLimit(.minutes(1)))
struct MediaPipelineTests {
@Test
func normalizeRejectsEmptyMimeTypeAndOversizedPayloads() async throws {
Expand Down Expand Up @@ -415,22 +415,26 @@ struct MediaPipelineTests {
}

let baseURL = URL(string: "http://127.0.0.1:\(port)/")!
try await self.waitForHTTPServer(baseURL: baseURL)
try await self.waitForHTTPServer(baseURL: baseURL, process: process)
try await body(baseURL)
}

private func waitForHTTPServer(baseURL: URL) async throws {
var lastError: (any Error)?
for _ in 0..<20 {
do {
_ = try Data(contentsOf: baseURL)
/// Polls until the server answers. There is no deadline: a loaded runner only delays Python's
/// startup, and the suite's time limit bounds a server that never answers.
private func waitForHTTPServer(baseURL: URL, process: Process) async throws {
while !Task.isCancelled {
if (try? Data(contentsOf: baseURL)) != nil {
return
} catch {
lastError = error
try await Task.sleep(nanoseconds: 100_000_000)
}
guard process.isRunning else {
throw OpenClawCoreError.unavailable("Local HTTP server exited with status \(process.terminationStatus)")
}
try? await Task.sleep(nanoseconds: 100_000_000)
}
throw lastError ?? OpenClawCoreError.unavailable("Local HTTP server did not start")
// Swift Testing drops errors thrown after a time-limit cancellation, so record which wait hung.
let timeout = AsyncWaitTimeoutError(label: "local HTTP server answering")
Issue.record(timeout)
throw timeout
}

private func reserveTCPPort() throws -> Int {
Expand Down
49 changes: 34 additions & 15 deletions Tests/OpenClawKitTests/ModelRoutingTests.swift
Original file line number Diff line number Diff line change
Expand Up @@ -332,9 +332,11 @@ struct ModelRoutingTests {
StaticProvider(id: "primary", text: "primary-output"),
StaticProvider(id: "secondary", text: "secondary-output"),
],
// An hour-long window: the second request is inside it however slow the runner is, and a
// drop never waits for the window to pass.
throttlePolicy: ModelProviderThrottlePolicy(
maxRequestsPerWindow: 1,
windowMs: 1_000,
windowMs: 3_600_000,
strategy: .drop
),
diagnosticsSink: await pipeline.sink()
Expand All @@ -355,32 +357,49 @@ struct ModelRoutingTests {
#expect(events.contains(where: { $0.name == "model.request.retry" }))
}

@Test
@Test(.timeLimit(.minutes(1)))
func routerThrottleDelayAppliesCooldownAndEmitsDiagnostics() async throws {
let pipeline = RuntimeDiagnosticsPipeline(eventLimit: 50)
let request = ModelGenerationRequest(sessionKey: "main", prompt: "hello", providerID: "primary")
let router = ModelRouter(
defaultProviderID: "primary",
providers: [StaticProvider(id: "primary", text: "primary-output")],
throttlePolicy: ModelProviderThrottlePolicy(
maxRequestsPerWindow: 1,
windowMs: 70,
strategy: .delay
),
diagnosticsSink: await pipeline.sink()
)
)

// The second request proceeds once the window has passed. A lower bound: a slow runner only
// adds to it (and may let the window pass before the second request arrives).
let startedAt = Date()
_ = try await router.generate(
ModelGenerationRequest(sessionKey: "main", prompt: "hello", providerID: "primary")
)
_ = try await router.generate(
ModelGenerationRequest(sessionKey: "main", prompt: "hello", providerID: "primary")
)
let elapsed = Date().timeIntervalSince(startedAt)
_ = try await router.generate(request)
_ = try await router.generate(request)
#expect(Date().timeIntervalSince(startedAt) >= 0.06)

#expect(elapsed >= 0.06)
let events = await pipeline.recentEvents(limit: 50)
#expect(events.contains(where: { $0.name == "model.throttle.delay" }))
// With an hour-long window the second request always lands inside it and waits.
let pipeline = RuntimeDiagnosticsPipeline(eventLimit: 50)
let slowRouter = ModelRouter(
defaultProviderID: "primary",
providers: [StaticProvider(id: "primary", text: "primary-output")],
throttlePolicy: ModelProviderThrottlePolicy(
maxRequestsPerWindow: 1,
windowMs: 3_600_000,
strategy: .delay
),
diagnosticsSink: await pipeline.sink()
)
_ = try await slowRouter.generate(request)
let delayed = Task { try await slowRouter.generate(request) }
try await waitUntil("model.throttle.delay emitted") {
await pipeline.recentEvents(limit: 50).contains { $0.name == "model.throttle.delay" }
}
let delay = try #require(await pipeline.recentEvents(limit: 50).first { $0.name == "model.throttle.delay" })
#expect((Int(delay.metadata["delayMs"] ?? "") ?? 0) > 3_000_000)
delayed.cancel()
await #expect(throws: CancellationError.self) {
_ = try await delayed.value
}
}

@Test
Expand Down
11 changes: 6 additions & 5 deletions Tests/OpenClawKitTests/SystemStateReportingTests.swift
Original file line number Diff line number Diff line change
Expand Up @@ -191,16 +191,17 @@ struct SystemStateReportingTests {
#expect(recorder.reports.last == .volatile(.gateway, ["backoffMs": 4000]))
}

@Test
@Test(.timeLimit(.minutes(1)))
func coalescingReporterSchedulesTrailingFlush() async throws {
let recorder = RecordingStateReporter()
let reporter = CoalescingSystemStateReporter(wrapping: recorder, minimumInterval: 0.05)
// The clock stands still, so the second update always lands inside the interval; only the
// trailing flush can forward it.
let clock = TestClock()
let reporter = CoalescingSystemStateReporter(wrapping: recorder, minimumInterval: 0.05, now: { clock.now })
reporter.reportVolatileUpdate(.talk, ["level": 1])
reporter.reportVolatileUpdate(.talk, ["level": 2])
#expect(recorder.reports == [.volatile(.talk, ["level": 1])])
for _ in 0..<50 where recorder.reports.count < 2 {
try await Task.sleep(nanoseconds: 20_000_000)
}
try await waitUntil("trailing flush forwarded the held update") { recorder.reports.count >= 2 }
#expect(recorder.reports == [.volatile(.talk, ["level": 1]), .volatile(.talk, ["level": 2])])
}

Expand Down
33 changes: 31 additions & 2 deletions Tests/OpenClawKitTests/TestAsyncHelpers.swift
Original file line number Diff line number Diff line change
Expand Up @@ -21,8 +21,37 @@ func waitUntil(
if await condition() { return }
try? await Task.sleep(nanoseconds: pollMs * 1_000_000)
}
// Swift Testing drops errors thrown after a time-limit cancellation, so record which wait hung.
throw recordedWaitTimeout(label)
}

/// Records which wait hung and returns the error to throw. Swift Testing drops errors thrown after a
/// time-limit cancellation, so the recorded issue is what names the wait.
func recordedWaitTimeout(_ label: String) -> AsyncWaitTimeoutError {
let timeout = AsyncWaitTimeoutError(label: label)
Issue.record(timeout)
throw timeout
return timeout
}

/// Awaits `operation` in its own task, so a time-limit cancellation ends the wait even when the
/// operation ignores cancellation.
///
/// `EmbeddedAgentRuntime.wait(runID:)` and the gateway's `agent.wait` park on a continuation that
/// only their own timer or the finished run resumes. Wrapping them here lets a test wait without a
/// timeout: a regression shows up as a hang, and the time limit turns it into an
/// ``AsyncWaitTimeoutError`` naming `label`.
func awaitCancellable<T: Sendable>(_ label: String, _ operation: @escaping @Sendable () async throws -> T) async throws -> T {
let (results, continuation) = AsyncStream<Result<T, any Error>>.makeStream(bufferingPolicy: .bufferingNewest(1))
let task = Task {
do {
continuation.yield(.success(try await operation()))
} catch {
continuation.yield(.failure(error))
}
continuation.finish()
}
defer { task.cancel() }
for await result in results {
return try result.get()
}
throw recordedWaitTimeout(label)
}
5 changes: 2 additions & 3 deletions Tests/OpenClawKitTests/VoiceNoteRecorderTests.swift
Original file line number Diff line number Diff line change
Expand Up @@ -359,10 +359,9 @@ struct VoiceNoteRecorderTests {
#expect(started)
#expect(recorder.level == 0)

for _ in 0..<200 where recorder.level == 0 {
try await Task.sleep(nanoseconds: 3_000_000)
try await waitUntil("capture level published") {
await MainActor.run { recorder.level > 0 }
}
#expect(recorder.level > 0)

_ = try #require(recorder.finish())
#expect(recorder.level == 0)
Expand Down
Loading
Loading