diff --git a/.changes/async-completer-cancel-deadlock b/.changes/async-completer-cancel-deadlock new file mode 100644 index 000000000..767cc2ae2 --- /dev/null +++ b/.changes/async-completer-cancel-deadlock @@ -0,0 +1 @@ +patch type="fixed" "Fix a deadlock when a `waitUntilActive` / `waitUntilAnyActive` timeout fired at the moment the wait was cancelled, which could hang the app's task and a timer thread" diff --git a/.github/stall-dump/stall-dump.sh b/.github/stall-dump/stall-dump.sh new file mode 100755 index 000000000..d87662578 --- /dev/null +++ b/.github/stall-dump/stall-dump.sh @@ -0,0 +1,22 @@ +#!/bin/bash +# Usage: stall-dump.sh [stall-seconds] +# +# Watches ; each time it stops growing for , dumps every test host into +# : a thread sample (a blocked thread) and the Swift concurrency runtime's task tree with +# async backtraces (a suspended task, which no thread sample can show). Up to three dumps, then exits. +log=$1; out=$2; stall=${3:-300} +last=-1; quiet=0; dumps=0 +while sleep 10; do + size=$(stat -f%z "$log" 2>/dev/null || echo 0) + if [ "$size" != "$last" ]; then last=$size; quiet=0; continue; fi + quiet=$((quiet + 10)) + [ "$quiet" -lt "$stall" ] && continue + quiet=0; dumps=$((dumps + 1)); mkdir -p "$out" + pids=$(pgrep -x xctest); [ -z "$pids" ] && pids=$(pgrep -x xcodebuild) + for pid in $pids; do + ps -o pid=,ppid=,etime=,command= -p "$pid" >> "$out/$dumps-processes.txt" + sample "$pid" 5 -file "$out/$dumps-sample-$pid.txt" >/dev/null 2>&1 + swift-inspect dump-concurrency "$pid" > "$out/$dumps-tasks-$pid.txt" 2>&1 + done + [ "$dumps" -ge 3 ] && exit 0 +done diff --git a/.github/workflows/ci.yaml b/.github/workflows/ci.yaml index 1d3359b3c..057327650 100644 --- a/.github/workflows/ci.yaml +++ b/.github/workflows/ci.yaml @@ -124,6 +124,8 @@ jobs: - name: Build & Test timeout-minutes: 30 run: | + # A hung job's log goes quiet; capture where the test host is stuck before the timeout kills it. + .github/stall-dump/stall-dump.sh "$RUNNER_TEMP/xcodebuild.log" "$RUNNER_TEMP/stall-dump" & set -o pipefail && xcodebuild test \ -scheme LiveKit \ -destination 'platform=${{ matrix.platform }}' \ @@ -134,7 +136,6 @@ jobs: -only-testing:LiveKitNanopbTests \ -only-testing:LiveKitObjCTests \ -parallel-testing-enabled NO \ - -retry-tests-on-failure -test-iterations 2 \ | tee "$RUNNER_TEMP/xcodebuild.log" \ | xcbeautify --renderer github-actions @@ -146,7 +147,9 @@ jobs: uses: actions/upload-artifact@v7 with: name: test-log-${{ strategy.job-index }} - path: ${{ runner.temp }}/xcodebuild.log + path: | + ${{ runner.temp }}/xcodebuild.log + ${{ runner.temp }}/stall-dump retention-days: 5 # Client-side logs can't distinguish "we gave up" from "the SFU hung up on diff --git a/Sources/LiveKit/Support/Async/AsyncCompleter.swift b/Sources/LiveKit/Support/Async/AsyncCompleter.swift index c2c5b493b..5c2770d45 100644 --- a/Sources/LiveKit/Support/Async/AsyncCompleter.swift +++ b/Sources/LiveKit/Support/Async/AsyncCompleter.swift @@ -63,6 +63,8 @@ actor CompleterMapActor { } } +/// Waiters are always resumed outside `_lock`: the runtime resumes a continuation under the task's +/// status-record lock, and a cancellation handler runs under that same lock while taking `_lock`. final class AsyncCompleter: @unchecked Sendable, Loggable { // struct WaitEntry { @@ -116,12 +118,14 @@ final class AsyncCompleter: @unchecked Sendable, Loggable { } func reset(throwing error: Error? = nil) { - _lock.sync { - for entry in _entries.values { - entry.cancel(throwing: LiveKitError.from(error: error)) - } + let entries = _lock.sync { + let entries = Array(_entries.values) _entries.removeAll() _result = nil + return entries + } + for entry in entries { + entry.cancel(throwing: LiveKitError.from(error: error)) } } @@ -134,16 +138,18 @@ final class AsyncCompleter: @unchecked Sendable, Loggable { } func resume(with result: Result) { - _lock.sync { + let entries = _lock.sync { if let _result { log("\(label) already resolved \(_entries) with \(_result)", .debug) } - for entry in _entries.values { - entry.resume(with: result) - } + let entries = Array(_entries.values) _entries.removeAll() _result = result + return entries + } + for entry in entries { + entry.resume(with: result) } } @@ -180,12 +186,7 @@ final class AsyncCompleter: @unchecked Sendable, Loggable { let timeoutBlock = DispatchWorkItem { [weak self] in guard let self else { return } log("\(label) id: \(entryId) timed out") - _lock.sync { - if let entry = self._entries[entryId] { - entry.timeout() - } - self._entries.removeValue(forKey: entryId) - } + _lock.sync { _entries.removeValue(forKey: entryId) }?.timeout() } _lock.sync { @@ -200,12 +201,7 @@ final class AsyncCompleter: @unchecked Sendable, Loggable { } } onCancel: { // Cancel only this completer when Task gets cancelled - _lock.sync { - if let entry = self._entries[entryId] { - entry.cancel() - } - self._entries.removeValue(forKey: entryId) - } + _lock.sync { _entries.removeValue(forKey: entryId) }?.cancel() } } } diff --git a/Tests/LiveKitCoreTests/CompleterTests.swift b/Tests/LiveKitCoreTests/CompleterTests.swift index 1fabd9e72..5eb5c749b 100644 --- a/Tests/LiveKitCoreTests/CompleterTests.swift +++ b/Tests/LiveKitCoreTests/CompleterTests.swift @@ -174,6 +174,32 @@ struct CompleterTests { completer.resume(returning: ()) try await secondTask.value } + + /// A waiter cancelled while its own timeout is firing must settle, not deadlock: the first child + /// to time out cancels the rest at the very moment their timers go off. Races for a fixed wall-clock + /// budget rather than a count, so a slow host cannot turn slowness into a timeout: only a deadlock + /// leaves the loop unfinished. + @Test func cancelRacingTimeoutSettles() async throws { + let races = Task.detached { + let deadline = Date().addingTimeInterval(3) + while Date() < deadline { + let first = AsyncCompleter(label: "first", defaultTimeout: 1) + let second = AsyncCompleter(label: "second", defaultTimeout: 1) + _ = try? await withThrowingTaskGroup(of: Void.self) { group in + group.addTask { try await first.wait(timeout: 0.001) } + group.addTask { try await second.wait(timeout: 0.001) } + for try await _ in group.prefix(1) { + group.cancelAll() + } + } + } + } + // Bounded by a completer of its own, so a regression fails the test instead of the job. Generous, + // because utility-QoS timers have been seen to fire tens of seconds late on loaded simulators. + let finished = AsyncCompleter(label: "races", defaultTimeout: 120) + Task.detached { await races.value; finished.resume(returning: ()) } + try await finished.wait() + } } @Suite(.tags(.concurrency)) diff --git a/Tests/LiveKitCoreTests/DataChannel/DataChannelDrainTests.swift b/Tests/LiveKitCoreTests/DataChannel/DataChannelDrainTests.swift index e20ff8c6b..e06ba9ba4 100644 --- a/Tests/LiveKitCoreTests/DataChannel/DataChannelDrainTests.swift +++ b/Tests/LiveKitCoreTests/DataChannel/DataChannelDrainTests.swift @@ -14,6 +14,8 @@ * limitations under the License. */ +// swiftlint:disable file_length + import Foundation @testable import LiveKit import Testing @@ -274,15 +276,25 @@ struct DropOldestContinuationTests { drain.attach(sendTarget: channel) } - private func sendAsync(_ tag: UInt8) -> Task { - Task { try await drain.send(DrainFixture.frame(tag)) } + /// Starts a send and returns once its submit has reached the drain's event stream, so a + /// `flushEvents()` that follows is a barrier behind it: `Task {}` alone may not have run yet. + private func sendAsync(_ tag: UInt8) async -> Task { + let (submitted, mark) = AsyncStream.makeStream(of: Void.self) + let task = Task { + try await withCheckedThrowingContinuation { continuation in + drain.submit(DrainFixture.frame(tag), continuation: continuation) + mark.finish() + } + } + for await _ in submitted {} + return task } @Test(.spec("https://github.com/livekit/client-sdk-js/blob/499c8420/src/room/RTCEngine.ts#L1458")) func evictionResolvesTheDisplacedWaiter() async throws { try await drain.fillBuffer(of: channel) - let displaced = sendAsync(1) + let displaced = await sendAsync(1) try await drain.flushEvents() // A newer group evicts the queued one; its waiter must not be left suspended. @@ -294,7 +306,7 @@ struct DropOldestContinuationTests { @Test func channelSwapResolvesQueuedWaiters() async throws { try await drain.fillBuffer(of: channel) - let queued = sendAsync(1) + let queued = await sendAsync(1) try await drain.flushEvents() drain.attach(sendTarget: FakeSendChannel()) @@ -309,7 +321,7 @@ struct DropOldestContinuationTests { @Test func rejectedSendFailsItsWaiterExactlyOnce() async throws { channel.acceptsSends = false - let waiter = sendAsync(1) + let waiter = await sendAsync(1) await #expect { try await waiter.value @@ -324,7 +336,7 @@ struct DropOldestContinuationTests { @Test func teardownFailsQueuedWaiters() async throws { try await drain.fillBuffer(of: channel) - let queued = sendAsync(1) + let queued = await sendAsync(1) try await drain.flushEvents() drain.reset(throwing: LiveKitError(.invalidState, message: "torn down")) diff --git a/Tests/LiveKitCoreTests/Participant/RemoteParticipantTests.swift b/Tests/LiveKitCoreTests/Participant/RemoteParticipantTests.swift index 8940840d2..64e121694 100644 --- a/Tests/LiveKitCoreTests/Participant/RemoteParticipantTests.swift +++ b/Tests/LiveKitCoreTests/Participant/RemoteParticipantTests.swift @@ -26,6 +26,13 @@ import LiveKitTestSupport struct RemoteParticipantTests { let timeout: TimeInterval = 0.1 + /// Makes `room` forget that `participant` became active, so `waitUntilActive` has to wait for a + /// transition that never comes rather than return the cached outcome. + private func forgetActive(_ participant: RemoteParticipant, in room: Room) async throws { + let identity = try #require(participant.identity) + await room.activeParticipantCompleters.completer(for: identity.stringValue).reset() + } + @Test func waitUntilActiveSuccess() async throws { try await TestEnvironment.withRooms(Array(repeating: RoomTestingOptions(), count: 2)) { rooms in let active = try #require(rooms[0].remoteParticipants.values.first) @@ -36,10 +43,10 @@ struct RemoteParticipantTests { @Test func waitUntilActiveTimeout() async throws { try await TestEnvironment.withRooms(Array(repeating: RoomTestingOptions(), count: 2)) { rooms in - let disconnected = try #require(rooms[0].remoteParticipants.values.first) - disconnected.set(info: .init(), connectionState: .disconnected) + let inactive = try #require(rooms[0].remoteParticipants.values.first) + try await forgetActive(inactive, in: rooms[0]) - await #expect(throws: (any Error).self) { try await disconnected.waitUntilActive(timeout: self.timeout) } + await #expect { try await inactive.waitUntilActive(timeout: self.timeout) } throws: { ($0 as? LiveKitError)?.type == .timedOut } } } @@ -53,10 +60,10 @@ struct RemoteParticipantTests { @Test func waitUntillAllActiveTimeout() async throws { try await TestEnvironment.withRooms(Array(repeating: RoomTestingOptions(), count: 3)) { rooms in - let oneDisconnected = try #require(rooms[0].remoteParticipants.values.first) - oneDisconnected.set(info: .init(), connectionState: .disconnected) + let oneInactive = try #require(rooms[0].remoteParticipants.values.first) + try await forgetActive(oneInactive, in: rooms[0]) - await #expect(throws: (any Error).self) { try await rooms[0].remoteParticipants.values.waitUntilAllActive(timeout: self.timeout) } + await #expect { try await rooms[0].remoteParticipants.values.waitUntilAllActive(timeout: self.timeout) } throws: { ($0 as? LiveKitError)?.type == .timedOut } try await rooms[1].remoteParticipants.values.waitUntilAllActive(timeout: timeout) try await rooms[2].remoteParticipants.values.waitUntilAllActive(timeout: timeout) } @@ -72,8 +79,8 @@ struct RemoteParticipantTests { @Test func waitUntillAnyActiveNoTimeout() async throws { try await TestEnvironment.withRooms(Array(repeating: RoomTestingOptions(), count: 3)) { rooms in - let oneDisconnected = try #require(rooms[0].remoteParticipants.values.first) - oneDisconnected.set(info: .init(), connectionState: .disconnected) + let oneInactive = try #require(rooms[0].remoteParticipants.values.first) + try await forgetActive(oneInactive, in: rooms[0]) try await rooms[0].remoteParticipants.values.waitUntilAnyActive(timeout: timeout) try await rooms[1].remoteParticipants.values.waitUntilAnyActive(timeout: timeout) @@ -83,10 +90,11 @@ struct RemoteParticipantTests { @Test func waitUntillAnyActiveTimeout() async throws { try await TestEnvironment.withRooms(Array(repeating: RoomTestingOptions(), count: 3)) { rooms in - let allDisconnected = rooms[0].remoteParticipants.values - allDisconnected.forEach { $0.set(info: .init(), connectionState: .disconnected) } + for participant in rooms[0].remoteParticipants.values { + try await forgetActive(participant, in: rooms[0]) + } - await #expect(throws: (any Error).self) { try await rooms[0].remoteParticipants.values.waitUntilAnyActive(timeout: self.timeout) } + await #expect { try await rooms[0].remoteParticipants.values.waitUntilAnyActive(timeout: self.timeout) } throws: { ($0 as? LiveKitError)?.type == .timedOut } try await rooms[1].remoteParticipants.values.waitUntilAnyActive(timeout: timeout) try await rooms[2].remoteParticipants.values.waitUntilAnyActive(timeout: timeout) } diff --git a/Tests/LiveKitTestSupport/Room+DataTrack.swift b/Tests/LiveKitTestSupport/Room+DataTrack.swift index c5d9c2d04..10ed31537 100644 --- a/Tests/LiveKitTestSupport/Room+DataTrack.swift +++ b/Tests/LiveKitTestSupport/Room+DataTrack.swift @@ -133,17 +133,12 @@ public final class DataTrackDelegateRecorder: NSObject, RoomDelegate, Participan /// waits forever if a frame is lost — on an unreliable channel that turns a failed assertion into /// a hung job. public extension DataTrackStream { - /// The next frame, or `nil` if none arrives in time. + /// The next frame, or `nil` if none arrives in time. The read is not cancelled on timeout — the + /// UniFFI future under `next()` cannot be — so it is left to finish when the stream ends. func next(within timeout: TimeInterval = 15) async -> DataTrackFrame? { - await withTaskGroup(of: DataTrackFrame?.self) { group in - group.addTask { await self.next() } - group.addTask { - try? await Task.sleep(nanoseconds: UInt64(timeout * 1_000_000_000)) - return nil - } - defer { group.cancelAll() } - return await group.next() ?? nil - } + let frame = AsyncCompleter(label: "data track frame", defaultTimeout: timeout) + Task { await frame.resume(returning: self.next()) } + return try? await frame.wait() } /// Up to `count` frames matching `predicate`, or fewer if the deadline passes first.