diff --git a/Tests/OpenClawKitTests/CoreAIModelRuntimeTests.swift b/Tests/OpenClawKitTests/CoreAIModelRuntimeTests.swift index ddcf386..f98005f 100644 --- a/Tests/OpenClawKitTests/CoreAIModelRuntimeTests.swift +++ b/Tests/OpenClawKitTests/CoreAIModelRuntimeTests.swift @@ -296,7 +296,7 @@ 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 }) @@ -304,9 +304,7 @@ struct CoreAIModelRuntimeTests { 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) { diff --git a/Tests/OpenClawKitTests/GatewayServerChatUICompatTests.swift b/Tests/OpenClawKitTests/GatewayServerChatUICompatTests.swift index 6d0186c..3191832 100644 --- a/Tests/OpenClawKitTests/GatewayServerChatUICompatTests.swift +++ b/Tests/OpenClawKitTests/GatewayServerChatUICompatTests.swift @@ -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)") @@ -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 } } @@ -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") } diff --git a/Tests/OpenClawKitTests/MediaPipelineTests.swift b/Tests/OpenClawKitTests/MediaPipelineTests.swift index 7e6f64d..315d30c 100644 --- a/Tests/OpenClawKitTests/MediaPipelineTests.swift +++ b/Tests/OpenClawKitTests/MediaPipelineTests.swift @@ -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 { @@ -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 { diff --git a/Tests/OpenClawKitTests/ModelRoutingTests.swift b/Tests/OpenClawKitTests/ModelRoutingTests.swift index 6b22ef8..c5bb624 100644 --- a/Tests/OpenClawKitTests/ModelRoutingTests.swift +++ b/Tests/OpenClawKitTests/ModelRoutingTests.swift @@ -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() @@ -355,9 +357,9 @@ 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")], @@ -365,22 +367,39 @@ struct ModelRoutingTests { 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 diff --git a/Tests/OpenClawKitTests/SystemStateReportingTests.swift b/Tests/OpenClawKitTests/SystemStateReportingTests.swift index a2f88f9..00f681b 100644 --- a/Tests/OpenClawKitTests/SystemStateReportingTests.swift +++ b/Tests/OpenClawKitTests/SystemStateReportingTests.swift @@ -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])]) } diff --git a/Tests/OpenClawKitTests/VoiceNoteRecorderTests.swift b/Tests/OpenClawKitTests/VoiceNoteRecorderTests.swift index 0e164ab..3cd7405 100644 --- a/Tests/OpenClawKitTests/VoiceNoteRecorderTests.swift +++ b/Tests/OpenClawKitTests/VoiceNoteRecorderTests.swift @@ -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) diff --git a/Tests/OpenClawLinuxRuntimeTests/AgentRuntimeExtensionsTests.swift b/Tests/OpenClawLinuxRuntimeTests/AgentRuntimeExtensionsTests.swift index abec741..571aba6 100644 --- a/Tests/OpenClawLinuxRuntimeTests/AgentRuntimeExtensionsTests.swift +++ b/Tests/OpenClawLinuxRuntimeTests/AgentRuntimeExtensionsTests.swift @@ -38,7 +38,7 @@ actor SessionRoutingProvider: ModelProvider { } } -@Suite("Agent runtime extensions") +@Suite("Agent runtime extensions", .timeLimit(.minutes(1))) struct AgentRuntimeExtensionsTests { // MARK: - Exec gate @@ -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)) @@ -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") @@ -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")) @@ -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")) @@ -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") @@ -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)]) } @@ -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 @@ -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"]) } diff --git a/Tests/OpenClawLinuxRuntimeTests/ApprovalQuestionBrokerTests.swift b/Tests/OpenClawLinuxRuntimeTests/ApprovalQuestionBrokerTests.swift index 348b43b..89b3974 100644 --- a/Tests/OpenClawLinuxRuntimeTests/ApprovalQuestionBrokerTests.swift +++ b/Tests/OpenClawLinuxRuntimeTests/ApprovalQuestionBrokerTests.swift @@ -5,7 +5,7 @@ import OpenClawModels import OpenClawProtocol @testable import OpenClawAgents -@Suite("Approval and question brokers") +@Suite("Approval and question brokers", .timeLimit(.minutes(1))) struct ApprovalQuestionBrokerTests { // MARK: - Approvals @@ -56,7 +56,7 @@ struct ApprovalQuestionBrokerTests { func expiryFailsClosed() async throws { let broker = ApprovalBroker() let pending = await broker.request(presentation: .exec(commandText: "rm -rf build"), timeoutMs: 30) - let result = try await broker.waitDecision(id: pending.id, timeoutMs: 2_000) + let result = try await awaitCancellable("approval expired") { try await broker.waitDecision(id: pending.id) } #expect(result.state == .expired) #expect(result.reason == .timeout) #expect(result.isAllowed == false) @@ -76,13 +76,8 @@ struct ApprovalQuestionBrokerTests { let waiter = Task { await broker.requestAndWait(presentation: .exec(commandText: "git status"), agentID: "main", grantKey: key) } - var pendingID: String? - for _ in 0..<50 { - pendingID = await broker.pending().first?.id - if pendingID != nil { break } - try await Task.sleep(nanoseconds: 10_000_000) - } - let id = try #require(pendingID) + try await waitUntil("grant approval pending") { await !broker.pending().isEmpty } + let id = try #require(await broker.pending().first?.id) _ = try await broker.resolve(id: id, decision: .allowAlways, grantExpiresInDays: 7) #expect(await waiter.value.isAllowed) @@ -135,13 +130,11 @@ struct ApprovalQuestionBrokerTests { hooks: AgentLoopHooks(beforeToolCall: { _ in .requireApproval(AgentToolApprovalRequest(title: "Echo", description: "Allow?")) }) ) let runID = await runtime.start(AgentRunRequest(runID: "wait-approval", sessionKey: "w", prompt: "go"), streaming: false) - for _ in 0..<100 where await runtime.approvals.pending().isEmpty { - try await Task.sleep(nanoseconds: 10_000_000) - } + try await waitUntil("tool approval pending") { await !runtime.approvals.pending().isEmpty } let pending = try #require(await runtime.approvals.pending().first) #expect(pending.runID == runID) await runtime.abort(runID: runID) - let result = try #require(await runtime.wait(runID: runID, timeoutMs: 2_000)) + let result = try #require(await awaitCancellable("aborted run finished") { await runtime.wait(runID: runID) }) #expect(result.status == "error") #expect(await runtime.approvals.get(id: pending.id)?.state == .cancelled) } @@ -208,13 +201,8 @@ struct ApprovalQuestionBrokerTests { let asking = Task { try await tool.invoke(AgentToolInvocation(arguments: arguments, context: AgentToolInvocationContext(runID: "r", sessionKey: "chat")), update: nil) } - var questionID: String? - for _ in 0..<100 { - questionID = await broker.pendingQuestion(sessionKey: "chat")?.id - if questionID != nil { break } - try await Task.sleep(nanoseconds: 10_000_000) - } - let id = try #require(questionID) + try await waitUntil("question pending") { await broker.pendingQuestion(sessionKey: "chat") != nil } + let id = try #require(await broker.pendingQuestion(sessionKey: "chat")?.id) // One pending question per session. let second = try await tool.invoke(AgentToolInvocation(arguments: arguments, context: AgentToolInvocationContext(sessionKey: "chat")), update: nil) @@ -243,7 +231,7 @@ struct ApprovalQuestionBrokerTests { options: [AgentQuestionOption(label: "Yes"), AgentQuestionOption(label: "No")] ) let expiring = try await broker.request(questions: [prompt], timeoutMs: 30) - #expect(try await broker.waitAnswer(id: expiring.id, timeoutMs: 2_000) == .expired) + #expect(try await awaitCancellable("question expired") { try await broker.waitAnswer(id: expiring.id) } == .expired) let cancelled = try await broker.request(questions: [prompt], runID: "run-9") #expect(try await broker.waitAnswer(id: cancelled.id, timeoutMs: 20) == .pending) diff --git a/Tests/OpenClawLinuxRuntimeTests/ChannelAdaptersLinuxSmokeTests.swift b/Tests/OpenClawLinuxRuntimeTests/ChannelAdaptersLinuxSmokeTests.swift index 53831f4..46f8ba8 100644 --- a/Tests/OpenClawLinuxRuntimeTests/ChannelAdaptersLinuxSmokeTests.swift +++ b/Tests/OpenClawLinuxRuntimeTests/ChannelAdaptersLinuxSmokeTests.swift @@ -69,14 +69,22 @@ struct ChannelAdaptersLinuxSmokeTests { @Test func a2aRoundTripCompletesTaskWithReply() async throws { - let config = A2AChannelConfig(enabled: true, replyTimeoutMs: 5_000, peers: ["peer": A2APeerConfig(token: "secret")]) + // The longest reply timeout (10 min) outlasts the time limit; the task-store wait ignores + // cancellation, so the request is awaited through awaitCancellable. + let config = A2AChannelConfig( + enabled: true, + replyTimeoutMs: A2AChannelConfig.replyTimeoutRangeMs.upperBound, + peers: ["peer": A2APeerConfig(token: "secret")] + ) let adapter = A2AChannelAdapter(config: config) await adapter.setInboundHandler { message in try? await adapter.send(OutboundMessage(channel: .a2a, peerID: message.peerID, text: "done")) } try await adapter.start() let request = #"{"jsonrpc":"2.0","id":1,"method":"SendMessage","params":{"message":{"role":"ROLE_USER","contextId":"c1","parts":[{"text":"task"}]}}}"# - let response = await adapter.handleHTTP(method: "POST", path: "/a2a/v1", headers: ["Authorization": "Bearer secret"], body: Data(request.utf8)) + let response = try await awaitCancellable("A2A task completed") { + await adapter.handleHTTP(method: "POST", path: "/a2a/v1", headers: ["Authorization": "Bearer secret"], body: Data(request.utf8)) + } let decoded = try JSONDecoder().decode(A2AJSONRPCResponse.self, from: response.body) #expect(decoded.result?.task?.status.state == .completed) #expect(decoded.result?.task?.replyText == "done") @@ -106,7 +114,8 @@ struct ChannelAdaptersLinuxSmokeTests { try FileManager.default.setAttributes([.posixPermissions: 0o755], ofItemAtPath: script.path) let client = IMsgRPCClient(pipe: IMsgProcessPipe(cliPath: script.path)) try await client.start() - let result = try await client.request("ping", timeoutMs: 10_000) + // No request timeout (0); the reply wait ignores cancellation, so it goes through awaitCancellable. + let result = try await awaitCancellable("imsg ping answered") { try await client.request("ping", timeoutMs: 0) } #expect(result.dictionaryValue?["ok"]?.boolValue == true) await client.stop() } diff --git a/Tests/OpenClawLinuxRuntimeTests/ExecApprovalHardeningTests.swift b/Tests/OpenClawLinuxRuntimeTests/ExecApprovalHardeningTests.swift index 733bcac..726ea7a 100644 --- a/Tests/OpenClawLinuxRuntimeTests/ExecApprovalHardeningTests.swift +++ b/Tests/OpenClawLinuxRuntimeTests/ExecApprovalHardeningTests.swift @@ -5,15 +5,10 @@ import OpenClawProtocol @testable import OpenClawAgents /// Exec approval gate and grant-key hardening (allow-always scope, raw allowlist text, closure rules). -@Suite("Exec approval hardening") +@Suite("Exec approval hardening", .timeLimit(.minutes(1))) struct ExecApprovalHardeningTests { private func firstPending(_ broker: ApprovalBroker) async throws -> AgentApproval { - for _ in 0..<300 { - if let approval = await broker.pending().first { - return approval - } - try await Task.sleep(nanoseconds: 10_000_000) - } + try await waitUntil("exec approval pending") { await !broker.pending().isEmpty } return try #require(await broker.pending().first) } diff --git a/Tests/OpenClawLinuxRuntimeTests/GatewayRunLifecycleTests.swift b/Tests/OpenClawLinuxRuntimeTests/GatewayRunLifecycleTests.swift index 6828211..cff332e 100644 --- a/Tests/OpenClawLinuxRuntimeTests/GatewayRunLifecycleTests.swift +++ b/Tests/OpenClawLinuxRuntimeTests/GatewayRunLifecycleTests.swift @@ -8,7 +8,7 @@ import OpenClawProtocol /// Built-in run tracking, client timeouts, pagination cursors, run-id collisions and /// `progressCard.put` validation (2026.3.0 FX1 review fixes). -@Suite("Gateway run lifecycle") +@Suite("Gateway run lifecycle", .timeLimit(.minutes(1))) struct GatewayRunLifecycleTests { typealias Harness = GatewayServerTestHarness @@ -33,22 +33,24 @@ struct GatewayRunLifecycleTests { @Test func builtinAgentWaitTimesOutWithoutWaitingForTheRun() async throws { + // The run only ends when aborted, so a wait that ignored its timeout would hang until the time limit. let server = Self.bareServer("lifecycle-wait-timeout") { request in - try await Task.sleep(nanoseconds: 10_000_000_000) + try await Task.sleep(nanoseconds: 3_600_000_000_000) return Self.ok(request) } let accepted = try Harness.payload(await Harness.call(server, "agent", ["message": AnyCodable("slow"), "idempotencyKey": AnyCodable("slow-1")])) #expect(accepted["runId"] == AnyCodable("slow-1")) - let started = Date() - let waited = try Harness.payload(await Harness.call(server, "agent.wait", ["runId": AnyCodable("slow-1"), "timeoutMs": AnyCodable(50)])) + let waited = try Harness.payload(try await awaitCancellable("timed-out agent.wait returned") { + await Harness.call(server, "agent.wait", ["runId": AnyCodable("slow-1"), "timeoutMs": AnyCodable(50)]) + }) #expect(waited["status"] == AnyCodable("timeout")) - // Well before the 10 s run ends (generous margin for loaded CI machines). - #expect(Date().timeIntervalSince(started) < 5) // The run keeps being tracked after a timed-out wait. #expect(await server.trackedRuns["slow-1"] != nil) let aborted = try Harness.payload(await Harness.call(server, "sessions.abort", ["runId": AnyCodable("slow-1")])) #expect(aborted["status"] == AnyCodable("aborted")) - let final = try Harness.payload(await Harness.call(server, "agent.wait", ["runId": AnyCodable("slow-1"), "timeoutMs": AnyCodable(5_000)])) + let final = try Harness.payload(try await awaitCancellable("aborted run reported") { + await Harness.call(server, "agent.wait", ["runId": AnyCodable("slow-1")]) + }) #expect(final["status"] == AnyCodable("error")) #expect(final["error"] == AnyCodable("aborted")) } @@ -62,7 +64,9 @@ struct GatewayRunLifecycleTests { return Self.ok(request) } _ = await Harness.call(server, "agent", ["message": AnyCodable("explode"), "idempotencyKey": AnyCodable("boom-1")]) - let failed = try Harness.payload(await Harness.call(server, "agent.wait", ["runId": AnyCodable("boom-1"), "timeoutMs": AnyCodable(5_000)])) + let failed = try Harness.payload(try await awaitCancellable("failed run reported") { + await Harness.call(server, "agent.wait", ["runId": AnyCodable("boom-1")]) + }) #expect(failed["status"] == AnyCodable("error")) #expect(failed["error"] == AnyCodable("exploded")) #expect(failed["endedAt"] != nil) @@ -71,9 +75,7 @@ struct GatewayRunLifecycleTests { _ = await Harness.call(server, "sessions.send", [ "key": AnyCodable("agent:main:quick"), "message": AnyCodable("quick"), "idempotencyKey": AnyCodable("quick-1"), ]) - for _ in 0..<200 where await server.completedRuns["quick-1"] == nil { - try await Task.sleep(nanoseconds: 5_000_000) - } + try await waitUntil("unwatched run completed") { await server.completedRuns["quick-1"] != nil } #expect(await server.agentRuns.isEmpty) #expect(await server.trackedRuns.isEmpty) let abortByKey = try Harness.payload(await Harness.call(server, "sessions.abort", ["key": AnyCodable("agent:main:quick")])) @@ -169,7 +171,9 @@ struct GatewayRunLifecycleTests { "key": AnyCodable("agent:main:main"), "message": AnyCodable("hi"), ])) let runID = try #require(sent["runId"]?.stringValue) - _ = await Harness.call(stack.server, "agent.wait", ["runId": AnyCodable(runID), "timeoutMs": AnyCodable(5_000)]) + _ = try await awaitCancellable("sessions.send run finished") { + await Harness.call(stack.server, "agent.wait", ["runId": AnyCodable(runID)]) + } let task = await ledger.create(kind: .tool, sessionKey: "agent:main:main") for method in ["tasks.list", "approval.history"] { @@ -205,9 +209,11 @@ struct GatewayRunLifecycleTests { @Test func reusedIdempotencyKeysDoNotStartASecondRunUnderTheSameID() async throws { + // The first run stays in flight until the duplicate sends are refused. + let inFlight = AsyncGate() let stack = await Harness.runtimeStack("lifecycle-run-ids", turns: [ { _ in - try await Task.sleep(nanoseconds: 3_000_000_000) + await inFlight.wait() return ModelGenerationResponse(text: "first", providerID: "scripted") }, ], fallback: ScriptedToolProvider.text("later")) @@ -225,7 +231,10 @@ struct GatewayRunLifecycleTests { ])) #expect(agent["status"] == AnyCodable("in_flight")) #expect(await stack.runtime.activeRunIDs() == ["msg-1"]) - let waited = try Harness.payload(await Harness.call(stack.server, "agent.wait", ["runId": AnyCodable("msg-1"), "timeoutMs": AnyCodable(5_000)])) + await inFlight.open() + let waited = try Harness.payload(try await awaitCancellable("first run finished") { + await Harness.call(stack.server, "agent.wait", ["runId": AnyCodable("msg-1")]) + }) #expect(waited["output"] == AnyCodable("first")) } diff --git a/Tests/OpenClawLinuxRuntimeTests/MCPHardeningTests.swift b/Tests/OpenClawLinuxRuntimeTests/MCPHardeningTests.swift index d092169..9fa9949 100644 --- a/Tests/OpenClawLinuxRuntimeTests/MCPHardeningTests.swift +++ b/Tests/OpenClawLinuxRuntimeTests/MCPHardeningTests.swift @@ -370,11 +370,10 @@ struct MCPHardeningTests { feed.yield(Data(": ping\n\n".utf8)) let transport = MCPLegacySSETransport(url: URL(string: "https://mcp.example.com/mcp")!, http: http) let client = MCPClient(serverName: "legacy", transport: transport, connectionTimeoutMs: 100) - let started = Date() + // The error names the 100 ms deadline; without one the stream never ends and the time limit fails the test. await #expect(throws: MCPTransportError.timeout(method: "connect", milliseconds: 100)) { try await client.connect() } - #expect(Date().timeIntervalSince(started) < 2) await #expect(throws: MCPTransportError.closed("transport closed")) { try await transport.send(.notification(method: "ping", params: nil)) } diff --git a/Tests/OpenClawLinuxRuntimeTests/MCPStdioTransportTests.swift b/Tests/OpenClawLinuxRuntimeTests/MCPStdioTransportTests.swift index 47132b5..57a86a2 100644 --- a/Tests/OpenClawLinuxRuntimeTests/MCPStdioTransportTests.swift +++ b/Tests/OpenClawLinuxRuntimeTests/MCPStdioTransportTests.swift @@ -89,9 +89,8 @@ struct MCPStdioTransportTests { ) try await transport.start() let pid = try #require(await transport.processIdentifier) - let startedAt = Date() await transport.close() - #expect(Date().timeIntervalSince(startedAt) < 5) + // The server traps TERM, so it is gone only if close escalated to KILL. #expect(kill(pid, 0) != 0) } @@ -111,20 +110,26 @@ struct MCPStdioTransportTests { @Test func closeDoesNotWaitForAWriteBlockedOnAFullPipe() async throws { - // The server never reads stdin, so a 200 KB message fills the pipe and the write blocks. + // The server never reads stdin, so a 200 KB message fills the pipe and the write blocks. The + // server outlives the time limit: a close that waited for the write would hang until the + // cancellation handler below kills the server. let transport = try MCPStdioTransport( serverName: "deaf", - config: self.config("exec sleep 5"), + config: self.config("exec sleep 3600"), allowlist: ExecCommandAllowlist(patterns: ["/bin/*", "/usr/bin/*"]), shutdownGraceSeconds: 0.2 ) try await transport.start() + let pid = try #require(await transport.processIdentifier) let big = MCPJSONRPCMessage.request(id: .int(1), method: "tools/call", params: AnyCodable(["blob": AnyCodable(String(repeating: "x", count: 200_000))])) let pending = Task { try await transport.send(big) } try await Task.sleep(nanoseconds: 200_000_000) - let startedAt = Date() - await transport.close() - #expect(Date().timeIntervalSince(startedAt) < 3, "close escalates without waiting for the blocked write") + await withTaskCancellationHandler { + await transport.close() + } onCancel: { + _ = kill(pid, SIGKILL) + } + #expect(kill(pid, 0) != 0, "close escalates without waiting for the blocked write") // Once the child is gone the blocked write fails with EPIPE (no SIGPIPE crash). await #expect(throws: MCPTransportError.self) { try await pending.value diff --git a/Tests/OpenClawLinuxRuntimeTests/RuntimeIntegrationWiringTests.swift b/Tests/OpenClawLinuxRuntimeTests/RuntimeIntegrationWiringTests.swift index 5892f24..d2f1b4f 100644 --- a/Tests/OpenClawLinuxRuntimeTests/RuntimeIntegrationWiringTests.swift +++ b/Tests/OpenClawLinuxRuntimeTests/RuntimeIntegrationWiringTests.swift @@ -8,7 +8,7 @@ import OpenClawProtocol import OpenClawSkills @testable import OpenClawAgents -@Suite("Runtime integration wiring") +@Suite("Runtime integration wiring", .timeLimit(.minutes(1))) struct RuntimeIntegrationWiringTests { private static func document(_ json: String) throws -> OpenClawConfigDocument { try OpenClawConfigDocument.decode(Data(json.utf8)) @@ -433,9 +433,7 @@ struct RuntimeIntegrationWiringTests { ) let manager = SubagentManager(runtime: runtime) let record = try await manager.spawn(SubagentSpawnParams(task: "summarize", label: "Summary"), parentSessionKey: "agent:main:main") - for _ in 0..<300 where await !log.names.contains("subagent_ended") { - try await Task.sleep(nanoseconds: 10_000_000) - } + try await waitUntil("subagent_ended hook ran") { await log.names.contains("subagent_ended") } #expect(await log.names == ["subagent_spawned", "subagent_ended"]) let spawned = try #require(await log.first(.subagentSpawned)) #expect(spawned["childSessionKey"]?.stringValue == record.childSessionKey) diff --git a/Tests/OpenClawLinuxRuntimeTests/SignInWithChatGPTSessionTests.swift b/Tests/OpenClawLinuxRuntimeTests/SignInWithChatGPTSessionTests.swift index a43aaf0..5091d87 100644 --- a/Tests/OpenClawLinuxRuntimeTests/SignInWithChatGPTSessionTests.swift +++ b/Tests/OpenClawLinuxRuntimeTests/SignInWithChatGPTSessionTests.swift @@ -392,7 +392,8 @@ struct SignInWithChatGPTSessionTests { let server = SIWCFakeServer() let session = SIWCTest.makeSession(server: server, clock: SIWCTestClock(), callbackPort: 0) let browser = LoopbackCallingBrowser() - let result = try await session.signIn(using: browser, reauthenticating: nil, consent: .automatic, timeout: 20, secrets: SIWCTest.secrets) + // Sign-in ends on cancellation, so an hour-long timeout leaves the suite's time limit as the bound. + let result = try await session.signIn(using: browser, reauthenticating: nil, consent: .automatic, timeout: 3_600, secrets: SIWCTest.secrets) #expect(result.account.subject == "user-abc") #expect(browser.presentedURL.map { SIWCTest.queryItems($0)["client_id"] } == "dynamic_agent_client") #expect(browser.dismissCount >= 1) @@ -408,7 +409,7 @@ struct SignInWithChatGPTSessionTests { let server = SIWCFakeServer() let session = SIWCTest.makeSession(server: server, clock: SIWCTestClock(), callbackPort: 0) let browser = LoopbackCallingBrowser(closesAfterCallback: true) - let result = try await session.signIn(using: browser, reauthenticating: nil, consent: .automatic, timeout: 20, secrets: SIWCTest.secrets) + let result = try await session.signIn(using: browser, reauthenticating: nil, consent: .automatic, timeout: 3_600, secrets: SIWCTest.secrets) #expect(result.account.subject == "user-abc") } @@ -417,13 +418,13 @@ struct SignInWithChatGPTSessionTests { let server = SIWCFakeServer() let session = SIWCTest.makeSession(server: server, clock: SIWCTestClock(), callbackPort: 0) await #expect(throws: SignInWithChatGPTError.cancelled) { - try await session.signIn(using: CancellingBrowser(), timeout: 20) + try await session.signIn(using: CancellingBrowser(), timeout: 3_600) } await #expect(throws: SignInWithChatGPTError.timedOut) { try await session.signIn(using: SignInWithChatGPTExternalBrowser { _ in true }, timeout: 1) } await #expect(throws: OpenClawCoreError.self) { - try await session.signIn(using: SignInWithChatGPTExternalBrowser { _ in false }, timeout: 20) + try await session.signIn(using: SignInWithChatGPTExternalBrowser { _ in false }, timeout: 3_600) } #expect(try await session.accounts().isEmpty) } diff --git a/Tests/OpenClawLinuxRuntimeTests/SpotlightMemoryTests.swift b/Tests/OpenClawLinuxRuntimeTests/SpotlightMemoryTests.swift index 64f8c86..b247c9e 100644 --- a/Tests/OpenClawLinuxRuntimeTests/SpotlightMemoryTests.swift +++ b/Tests/OpenClawLinuxRuntimeTests/SpotlightMemoryTests.swift @@ -10,12 +10,31 @@ import OpenClawProtocol final class FakeSpotlightStore: SpotlightItemIndexing, @unchecked Sendable { private let lock = NSLock() private var items: [String: String] = [:] - private let hangs: Bool + private var hangs: Bool + private var stalled: [@Sendable ((any Error)?) -> Void] = [] init(hangs: Bool = false) { self.hangs = hangs } + /// Answers the calls a hanging store left pending, and every later call. + func stopHanging() { + let stalled = self.lock.withLock { + self.hangs = false + defer { self.stalled.removeAll() } + return self.stalled + } + for completion in stalled { completion(nil) } + } + + /// Keeps `completion` pending when the store hangs. + private func stalls(_ completion: @escaping @Sendable ((any Error)?) -> Void) -> Bool { + self.lock.withLock { + if self.hangs { self.stalled.append(completion) } + return self.hangs + } + } + var isAvailable: Bool { true } @@ -28,7 +47,7 @@ final class FakeSpotlightStore: SpotlightItemIndexing, @unchecked Sendable { } func indexItems(_ items: [CSSearchableItem], completion: @escaping @Sendable ((any Error)?) -> Void) { - guard !self.hangs else { return } + guard !self.stalls(completion) else { return } self.lock.lock() for item in items { self.items[item.uniqueIdentifier] = item.domainIdentifier ?? "" } self.lock.unlock() @@ -36,7 +55,7 @@ final class FakeSpotlightStore: SpotlightItemIndexing, @unchecked Sendable { } func deleteItems(identifiers: [String], completion: @escaping @Sendable ((any Error)?) -> Void) { - guard !self.hangs else { return } + guard !self.stalls(completion) else { return } self.lock.lock() for id in identifiers { self.items.removeValue(forKey: id) } self.lock.unlock() @@ -44,7 +63,7 @@ final class FakeSpotlightStore: SpotlightItemIndexing, @unchecked Sendable { } func deleteItems(domains: [String], completion: @escaping @Sendable ((any Error)?) -> Void) { - guard !self.hangs else { return } + guard !self.stalls(completion) else { return } self.lock.lock() self.items = self.items.filter { _, domain in !domains.contains { domain == $0 || domain.hasPrefix($0 + ".") } @@ -54,6 +73,20 @@ final class FakeSpotlightStore: SpotlightItemIndexing, @unchecked Sendable { } } +/// Tells a stalled test operation to stop. +final class StallFlag: @unchecked Sendable { + private let lock = NSLock() + private var ended = false + + var isEnded: Bool { + self.lock.withLock { self.ended } + } + + func end() { + self.lock.withLock { self.ended = true } + } +} + @Suite("Spotlight memory backends") struct SpotlightMemoryTests { static var live: Bool { @@ -83,17 +116,22 @@ struct SpotlightMemoryTests { #expect(SpotlightTimeoutRace.deadlineNanoseconds(-5) == 50_000_000) } - @Test + @Test(.timeLimit(.minutes(1))) func indexWritesAreBoundedWhenSpotlightNeverAnswers() async throws { - let index = Self.index(store: FakeSpotlightStore(hangs: true), writeTimeoutSeconds: 0.2) - let started = Date() - await #expect(throws: SpotlightTimeoutError.self) { - try await index.upsert([MemoryDocument(id: "a", source: .systemNote, text: "stalled write")], sessionKey: nil) - } - await #expect(throws: SpotlightTimeoutError.self) { - try await index.deleteAll() + let store = FakeSpotlightStore(hangs: true) + let index = Self.index(store: store, writeTimeoutSeconds: 0.2) + // The errors name the 0.2 s deadline. A write without one would hang until the time limit, + // whose cancellation answers the stalled calls so the test can end. + await withTaskCancellationHandler { + await #expect(throws: SpotlightTimeoutError(operation: "indexSearchableItems", seconds: 0.2)) { + try await index.upsert([MemoryDocument(id: "a", source: .systemNote, text: "stalled write")], sessionKey: nil) + } + await #expect(throws: SpotlightTimeoutError(operation: "deleteSearchableItems(withDomainIdentifiers:)", seconds: 0.2)) { + try await index.deleteAll() + } + } onCancel: { + store.stopHanging() } - #expect(Date().timeIntervalSince(started) < 3) } @Test @@ -229,18 +267,26 @@ struct SpotlightMemoryTests { feed.finish() } - @Test - func timeoutRaceReturnsWithoutJoiningAStalledOperation() async { - let started = Date() + @Test(.timeLimit(.minutes(1))) + func timeoutRaceReturnsWithoutJoiningAStalledOperation() async throws { // The operation ignores cancellation and never finishes on its own, like a stalled CSUserQuery. - let value: Int? = await SpotlightTimeoutRace.first(timeoutSeconds: 0.2) { - while true { - try? await Task.sleep(nanoseconds: 50_000_000) + // It stops once the test ends, or when the time limit cancels a race that waits for it. + let stall = StallFlag() + defer { stall.end() } + let value: Int? = await withTaskCancellationHandler { + await SpotlightTimeoutRace.first(timeoutSeconds: 0.2) { + while !stall.isEnded { + try? await Task.sleep(nanoseconds: 50_000_000) + } + return 0 } + } onCancel: { + stall.end() } #expect(value == nil) - #expect(Date().timeIntervalSince(started) < 2) - let fast: Int? = await SpotlightTimeoutRace.first(timeoutSeconds: 5) { 42 } + let fast: Int? = try await awaitCancellable("fast operation won the race") { + await SpotlightTimeoutRace.first(timeoutSeconds: 3_600) { 42 } + } #expect(fast == 42) } diff --git a/Tests/OpenClawLinuxRuntimeTests/SubagentRuntimeHardeningTests.swift b/Tests/OpenClawLinuxRuntimeTests/SubagentRuntimeHardeningTests.swift index 1a855f1..189fded 100644 --- a/Tests/OpenClawLinuxRuntimeTests/SubagentRuntimeHardeningTests.swift +++ b/Tests/OpenClawLinuxRuntimeTests/SubagentRuntimeHardeningTests.swift @@ -108,7 +108,7 @@ struct SubagentRuntimeHardeningTests { ) let manager = SubagentManager(runtime: runtime) let child = try await manager.spawn(SubagentSpawnParams(task: "run rm -rf build and write X"), parentSessionKey: parent) - #expect(await runtime.wait(runID: child.runID, timeoutMs: 10_000)?.status == "ok") + #expect(try await awaitCancellable("child run finished") { await runtime.wait(runID: child.runID) }?.status == "ok") let childRecord = try #require(await store.recordForKey(child.childSessionKey)) #expect(childRecord.permissionMode == .readOnly) @@ -143,11 +143,11 @@ struct SubagentRuntimeHardeningTests { await manager.registerTools() let result = try await runtime.run( AgentRunRequest(sessionKey: "agent:main:main", prompt: "delegate", toolPolicy: ToolPolicy(deny: ["exec"])), - timeoutMs: 10_000 + timeoutMs: 3_600_000 ) #expect(result.output == "parent done") let child = try #require(await manager.children(of: "agent:main:main").first) - #expect(await runtime.wait(runID: child.runID, timeoutMs: 10_000)?.status == "ok") + #expect(try await awaitCancellable("child run finished") { await runtime.wait(runID: child.runID) }?.status == "ok") #expect(await log.count("exec") == 0) let results = toolResults(try await runtime.history(sessionKey: child.childSessionKey)) #expect(results.first?.text == "Tool exec is not allowed by the current tool policy") @@ -190,20 +190,28 @@ struct SubagentRuntimeHardeningTests { await runtime.start(AgentRunRequest(runID: "dup", sessionKey: "two", prompt: "b"), streaming: false) #expect(await runtime.activeRunIDs() == ["dup"]) await provider.release() - #expect(await runtime.wait(runID: "dup", timeoutMs: 10_000)?.output == "first") + #expect(try await awaitCancellable("joined run finished") { await runtime.wait(runID: "dup") }?.output == "first") #expect(await provider.count() == 1) } @Test(.timeLimit(.minutes(1))) func aStaleTimerNeverCancelsANewRunWithTheSameID() async throws { - let provider = GatedTextProvider(replies: ["old", "new"]) - await provider.release() - let runtime = EmbeddedAgentRuntime( - toolRegistry: AgentToolRegistry(), - modelRouter: ModelRouter(defaultProviderID: provider.id, providers: [provider]) - ) - await runtime.start(AgentRunRequest(runID: "reuse", sessionKey: "s", prompt: "a"), timeoutMs: 200, streaming: false) - #expect(await runtime.wait(runID: "reuse", timeoutMs: 10_000)?.output == "old") + // The first run has to finish before its 200 ms timer so the timer is stale during the second + // run. A stalled pool can let the timer win, so start over until the first run finishes first. + var provider: GatedTextProvider + var runtime: EmbeddedAgentRuntime + var first: AgentRunWaitResult? + repeat { + provider = GatedTextProvider(replies: ["old", "new"]) + await provider.release() + runtime = EmbeddedAgentRuntime( + toolRegistry: AgentToolRegistry(), + modelRouter: ModelRouter(defaultProviderID: provider.id, providers: [provider]) + ) + await runtime.start(AgentRunRequest(runID: "reuse", sessionKey: "s", prompt: "a"), timeoutMs: 200, streaming: false) + first = try await awaitCancellable("first run finished") { [runtime] in await runtime.wait(runID: "reuse") } + } while first?.status == "timeout" + #expect(first?.output == "old") await provider.block() await runtime.start(AgentRunRequest(runID: "reuse", sessionKey: "s", prompt: "b"), streaming: false) @@ -213,7 +221,8 @@ struct SubagentRuntimeHardeningTests { #expect(await runtime.activeRunIDs() == ["reuse"]) await provider.release() // The waiter sees the new run, not the previous run's retained result. - #expect(await runtime.wait(runID: "reuse", timeoutMs: 10_000)?.output == "new") + let second = try await awaitCancellable("second run finished") { [runtime] in await runtime.wait(runID: "reuse") } + #expect(second?.output == "new") _ = await runtime.approvals.cancel(id: approval.id) } diff --git a/Tests/OpenClawLinuxRuntimeTests/TestAsyncHelpers.swift b/Tests/OpenClawLinuxRuntimeTests/TestAsyncHelpers.swift index ac54535..27cb5d2 100644 --- a/Tests/OpenClawLinuxRuntimeTests/TestAsyncHelpers.swift +++ b/Tests/OpenClawLinuxRuntimeTests/TestAsyncHelpers.swift @@ -31,3 +31,27 @@ func recordedWaitTimeout(_ label: String) -> AsyncWaitTimeoutError { Issue.record(timeout) return timeout } + +/// Awaits `operation` in its own task, so a time-limit cancellation ends the wait even when the +/// operation ignores cancellation. +/// +/// Product waits such as `EmbeddedAgentRuntime.wait(runID:)`, the gateway's `agent.wait` and the +/// approval and question brokers park on a continuation that only their own timer or the awaited +/// event 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(_ label: String, _ operation: @escaping @Sendable () async throws -> T) async throws -> T { + let (results, continuation) = AsyncStream>.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) +} diff --git a/Tests/OpenClawLinuxRuntimeTests/ToolsGatewayMethodsTests.swift b/Tests/OpenClawLinuxRuntimeTests/ToolsGatewayMethodsTests.swift index ab2c7c8..057bd5a 100644 --- a/Tests/OpenClawLinuxRuntimeTests/ToolsGatewayMethodsTests.swift +++ b/Tests/OpenClawLinuxRuntimeTests/ToolsGatewayMethodsTests.swift @@ -7,7 +7,7 @@ import OpenClawProtocol @testable import OpenClawAgents /// `tools.catalog`, `tools.effective` and `tools.invoke` over the runtime tool registry and policy. -@Suite("Tools gateway methods") +@Suite("Tools gateway methods", .timeLimit(.minutes(1))) struct ToolsGatewayMethodsTests { private typealias Harness = GatewayServerTestHarness @@ -164,13 +164,8 @@ struct ToolsGatewayMethodsTests { "name": AnyCodable("lookup"), "args": AnyCodable(["text": AnyCodable("w")]), "confirm": AnyCodable(true), ])) } - var pendingID: String? - for _ in 0..<200 { - pendingID = await stack.runtime.approvals.pending().first?.id - if pendingID != nil { break } - try await Task.sleep(nanoseconds: 10_000_000) - } - let waitingID = try #require(pendingID) + try await waitUntil("confirmation approval pending") { await !stack.runtime.approvals.pending().isEmpty } + let waitingID = try #require(await stack.runtime.approvals.pending().first?.id) _ = await Harness.call(stack.server, "approval.resolve", ["id": AnyCodable(waitingID), "decision": AnyCodable("deny")]) let denied = try await waiting.value #expect(denied["ok"] == AnyCodable(false))