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
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
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
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
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
40 changes: 14 additions & 26 deletions Tests/OpenClawLinuxRuntimeTests/AgentRuntimeExtensionsTests.swift
Original file line number Diff line number Diff line change
Expand Up @@ -38,7 +38,7 @@ actor SessionRoutingProvider: ModelProvider {
}
}

@Suite("Agent runtime extensions")
@Suite("Agent runtime extensions", .timeLimit(.minutes(1)))
struct AgentRuntimeExtensionsTests {
// MARK: - Exec gate

Expand Down Expand Up @@ -92,13 +92,8 @@ struct AgentRuntimeExtensionsTests {
#expect(await broker.pending().isEmpty)
// The fourth command escalates to a human.
let escalated = Task { await gate.evaluate(command: "danger 4", permissionMode: .workspace, sessionKey: "w") }
var pending: [AgentApproval] = []
for _ in 0..<100 {
pending = await broker.pending()
if !pending.isEmpty { break }
try await Task.sleep(nanoseconds: 10_000_000)
}
let approval = try #require(pending.first)
try await waitUntil("escalated approval pending") { await !broker.pending().isEmpty }
let approval = try #require(await broker.pending().first)
#expect(approval.presentation.warningText?.contains("Escalated") == true)
_ = try await broker.resolve(id: approval.id, decision: .allowOnce)
#expect(await escalated.value == .allow(source: .human))
Expand All @@ -108,9 +103,7 @@ struct AgentRuntimeExtensionsTests {
throw ReviewDown()
})
let asking = Task { await failing.evaluate(command: "anything", permissionMode: .workspace, sessionKey: "f") }
for _ in 0..<100 where await broker.pending().isEmpty {
try await Task.sleep(nanoseconds: 10_000_000)
}
try await waitUntil("failed-review approval pending") { await !broker.pending().isEmpty }
let failed = try #require(await broker.pending().first)
#expect(failed.presentation.warningText?.contains("Automatic review failed") == true)
await broker.cancel(runID: "none")
Expand All @@ -132,8 +125,7 @@ struct AgentRuntimeExtensionsTests {

// MARK: - Sub-agents and ledger

// Bounded so a lost wake-up fails fast instead of hanging the suite.
@Test(.timeLimit(.minutes(1)))
@Test
func spawnAnnouncesCompletionAndWakesTheYieldedParent() async throws {
let provider = SessionRoutingProvider()
let store = SessionStore(fileURL: FileManager.default.temporaryDirectory.appendingPathComponent("sub-\(UUID().uuidString)/sessions.json"))
Expand All @@ -148,7 +140,8 @@ struct AgentRuntimeExtensionsTests {
let manager = SubagentManager(runtime: runtime, ledger: ledger)
await manager.registerTools()

let result = try await runtime.run(AgentRunRequest(sessionKey: "agent:main:main", prompt: "delegate"), timeoutMs: 10_000)
// The run ends on cancellation, so the hour-long timeout leaves the time limit as the only bound.
let result = try await runtime.run(AgentRunRequest(sessionKey: "agent:main:main", prompt: "delegate"), timeoutMs: 3_600_000)
#expect(result.toolResults.map(\.name) == ["sessions_spawn", "sessions_yield"])
let spawnDetails = try #require(result.toolResults.first?.output.details?.dictionaryValue)
#expect(spawnDetails["status"] == AnyCodable("accepted"))
Expand All @@ -157,20 +150,15 @@ struct AgentRuntimeExtensionsTests {
#expect(await store.recordForKey(childKey)?.spawnedBy == "agent:main:main")
#expect(await store.recordForKey(childKey)?.spawnDepth == 1)

var prompts: [String] = []
for _ in 0..<300 {
prompts = await provider.prompts()
if prompts.contains(where: { $0.contains("child result") }) { break }
try await Task.sleep(nanoseconds: 10_000_000)
try await waitUntil("parent woke with the child result") {
await provider.prompts().contains { $0.contains("child result") }
}
#expect(prompts.contains { $0.contains("[Subagent completion]") && $0.contains("child result") })
#expect(await provider.prompts().contains { $0.contains("[Subagent completion]") && $0.contains("child result") })

let tasks = await ledger.list().tasks
#expect(tasks.count == 1)
#expect(tasks.first?.kind == .subagent)
for _ in 0..<100 where await ledger.list().tasks.first?.status.isTerminal != true {
try await Task.sleep(nanoseconds: 10_000_000)
}
try await waitUntil("sub-agent task terminal") { await ledger.list().tasks.first?.status.isTerminal == true }
#expect(await ledger.list().tasks.first?.status == .completed)
#expect(await ledger.list(statuses: [.running]).tasks.isEmpty)
let children = await manager.children(of: "agent:main:main")
Expand All @@ -185,7 +173,7 @@ struct AgentRuntimeExtensionsTests {
}
}

@Test(.timeLimit(.minutes(1)))
@Test
func spawnParamsRejectUnsupportedOptionsAndKillWorks() async throws {
#expect(throws: SubagentError.self) { try SubagentSpawnParams.parse(["task": AnyCodable("x"), "runtime": AnyCodable("acp")]) }
#expect(throws: SubagentError.self) { try SubagentSpawnParams.parse(["task": AnyCodable("x"), "visible": AnyCodable(true)]) }
Expand All @@ -212,7 +200,7 @@ struct AgentRuntimeExtensionsTests {
#expect(await manager.resolve(target: "1", parentSessionKey: "p").count == 1)
let killed = try await manager.kill(target: "last", parentSessionKey: "p")
#expect(killed.map(\.status) == ["killed"])
#expect(await runtime.wait(runID: record.runID, timeoutMs: 2_000)?.status == "error")
#expect(try await awaitCancellable("killed run finished") { await runtime.wait(runID: record.runID) }?.status == "error")
}

@Test
Expand Down Expand Up @@ -344,7 +332,7 @@ struct AgentRuntimeExtensionsTests {
)
)
let runID = try #require(refresh.payload?.dictionaryValue?["runId"]?.stringValue)
#expect(await runtime.wait(runID: runID, timeoutMs: 5_000)?.status == "ok")
#expect(try await awaitCancellable("refresh run finished") { await runtime.wait(runID: runID) }?.status == "ok")
// The refresh instruction is hidden from chat history.
#expect(try await runtime.history(sessionKey: "pc").map(\.role) == ["assistant"])
}
Expand Down
Loading
Loading