From adf7f5ec2f3bc943ff4287c386c6fa955275d2b7 Mon Sep 17 00:00:00 2001 From: Peter Steinberger Date: Sat, 29 Aug 2026 05:57:00 -0700 Subject: [PATCH] fix(macos): stop cancelled share waits and cursor retries Integrate the cancellation repair from #114 on current main. Own mailbox expiry tasks with their continuations and close cancellation before waiter registration. Stop cancelled cursor operations while preserving retries for transient errors, with behavioral regressions at the capture owner. Co-authored-by: Sebastien Tardif --- CHANGELOG.md | 2 + docs/macos-native-client.md | 4 +- .../CrabfleetMac/MacScreenCapture.swift | 33 +++++-- .../Sources/CrabfleetMac/VideoMailbox.swift | 63 +++++++++--- .../VideoPipelineTests.swift | 96 +++++++++++++++++++ 5 files changed, 174 insertions(+), 24 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index aaa3f13..c8b51a4 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,8 @@ ## Unreleased +- Stop cancelled Share This Mac cursor reconciliation and release video mailbox waits and their timeout tasks promptly, including cancellation before waiter registration, thanks @SebTardif (#114). + ## 0.3.1 - 2026-08-28 ### Highlights diff --git a/docs/macos-native-client.md b/docs/macos-native-client.md index 8aa7b00..37ab7a7 100644 --- a/docs/macos-native-client.md +++ b/docs/macos-native-client.md @@ -379,7 +379,9 @@ VeNCrypt would bring an RFB-layer TLS boundary to TCP and other transports too. The relay never stores the registration token or RFB bytes. The direct listener is not reachable on Wi-Fi, Ethernet, loopback, or a public address. The host app must remain running, and stopping the share cancels the listener, relay -publisher, and capture stream. +publisher, and capture stream. Cancelled video mailbox waits return without +waiting for their frame timeout, and cursor configuration reconciliation stops +on cancellation while retaining bounded backoff for transient failures. ### Browser viewer diff --git a/macos/CrabfleetMac/Sources/CrabfleetMac/MacScreenCapture.swift b/macos/CrabfleetMac/Sources/CrabfleetMac/MacScreenCapture.swift index e4317cf..8900f1c 100644 --- a/macos/CrabfleetMac/Sources/CrabfleetMac/MacScreenCapture.swift +++ b/macos/CrabfleetMac/Sources/CrabfleetMac/MacScreenCapture.swift @@ -412,15 +412,8 @@ final class MacScreenCapture: NSObject, @unchecked Sendable { cursorReconcileTask?.cancel() let task = Task { [weak self] in guard let self else { return } - var delay = Duration.milliseconds(100) - while !Task.isCancelled { - do { - try await self.reconcileCursorConfiguration() - break - } catch { - try? await Task.sleep(for: delay) - delay = min(delay * 2, .seconds(5)) - } + await Self.reconcileCursorConfigurationWithRetry { + try await self.reconcileCursorConfiguration() } self.withFrameLock { if self.cursorReconcileGeneration == generation { @@ -434,6 +427,28 @@ final class MacScreenCapture: NSObject, @unchecked Sendable { _ = task } + static func reconcileCursorConfigurationWithRetry( + _ reconcile: () async throws -> Void + ) async { + var delay = Duration.milliseconds(100) + while !Task.isCancelled { + do { + try await reconcile() + return + } catch is CancellationError { + // The operation can be cancelled even when this retry task is not. + return + } catch { + do { + try await Task.sleep(for: delay) + } catch { + return + } + delay = min(delay * 2, .seconds(5)) + } + } + } + private func reconcileCursorConfiguration() async throws { try await configurationGate.run { [self] in guard let stream, let configuration else { return } diff --git a/macos/CrabfleetMac/Sources/CrabfleetMac/VideoMailbox.swift b/macos/CrabfleetMac/Sources/CrabfleetMac/VideoMailbox.swift index 3183d33..f16a144 100644 --- a/macos/CrabfleetMac/Sources/CrabfleetMac/VideoMailbox.swift +++ b/macos/CrabfleetMac/Sources/CrabfleetMac/VideoMailbox.swift @@ -12,6 +12,7 @@ final class VideoMailbox: @unchecked Sendable { private let lock = NSLock() private var latestElement: Element? private var waiter: Waiter? + private var timeoutTask: Task? private var finished = false var isFinished: Bool { @@ -30,27 +31,36 @@ final class VideoMailbox: @unchecked Sendable { } func offer(_ element: Element, onDrop: () -> Void = {}) { - let continuation = withLock { () -> CheckedContinuation? in - guard !finished else { return nil } + let (continuation, timeoutTask) = withLock { + () -> (CheckedContinuation?, Task?) in + guard !finished else { return (nil, nil) } guard let waiter else { if latestElement != nil { onDrop() } latestElement = element - return nil + return (nil, nil) } self.waiter = nil - return waiter.continuation + let timeoutTask = self.timeoutTask + self.timeoutTask = nil + return (waiter.continuation, timeoutTask) } + timeoutTask?.cancel() continuation?.resume(returning: element) } func finish() { - let continuation = withLock { () -> CheckedContinuation? in - guard !finished else { return nil } + let (continuation, timeoutTask) = withLock { + () -> (CheckedContinuation?, Task?) in + guard !finished else { return (nil, nil) } finished = true latestElement = nil - defer { waiter = nil } - return waiter?.continuation + let continuation = waiter?.continuation + let timeoutTask = self.timeoutTask + waiter = nil + self.timeoutTask = nil + return (continuation, timeoutTask) } + timeoutTask?.cancel() continuation?.resume(returning: nil) } @@ -64,6 +74,7 @@ final class VideoMailbox: @unchecked Sendable { var immediateElement: Element? var shouldResume = false var replacedWaiter: Waiter? + var replacedTimeout: Task? lock.lock() if let latestElement { immediateElement = latestElement @@ -73,16 +84,35 @@ final class VideoMailbox: @unchecked Sendable { shouldResume = true } else { replacedWaiter = waiter + replacedTimeout = timeoutTask + timeoutTask = nil waiter = (id, continuation) } lock.unlock() replacedWaiter?.continuation.resume(returning: nil) + replacedTimeout?.cancel() if shouldResume { continuation.resume(returning: immediateElement) } else { - Task { - try? await Task.sleep(for: timeout) + let task = Task { + do { + try await Task.sleep(for: timeout) + self.expire(id: id) + } catch is CancellationError { + // Do not expire a different waiter; expire already guards by id. + } catch { + self.expire(id: id) + } + } + let shouldCancelTimeout = withLock { () -> Bool in + guard waiter?.id == id else { return true } + timeoutTask = task + return false + } + if shouldCancelTimeout { + task.cancel() + } else if Task.isCancelled { self.expire(id: id) } } @@ -93,11 +123,16 @@ final class VideoMailbox: @unchecked Sendable { } private func expire(id: UUID) { - let continuation = withLock { () -> CheckedContinuation? in - guard waiter?.id == id else { return nil } - defer { waiter = nil } - return waiter?.continuation + let (continuation, timeoutTask) = withLock { + () -> (CheckedContinuation?, Task?) in + guard waiter?.id == id else { return (nil, nil) } + let continuation = waiter?.continuation + let timeoutTask = self.timeoutTask + waiter = nil + self.timeoutTask = nil + return (continuation, timeoutTask) } + timeoutTask?.cancel() continuation?.resume(returning: nil) } diff --git a/macos/CrabfleetMac/Tests/CrabfleetMacTests/VideoPipelineTests.swift b/macos/CrabfleetMac/Tests/CrabfleetMacTests/VideoPipelineTests.swift index 1e3ef0c..ed9c98d 100644 --- a/macos/CrabfleetMac/Tests/CrabfleetMacTests/VideoPipelineTests.swift +++ b/macos/CrabfleetMac/Tests/CrabfleetMacTests/VideoPipelineTests.swift @@ -377,6 +377,102 @@ struct VideoPipelineTests { #expect(await mailbox.next(timeout: .milliseconds(10)) == nil) } + @Test + func mailboxCancelledWaiterReturnsPromptlyWithoutResumingTwice() async { + let mailbox = VideoMailbox() + let startedAt = ContinuousClock().now + let task = Task { + withUnsafeCurrentTask { $0?.cancel() } + return await mailbox.next(timeout: .seconds(5)) + } + let result = await task.value + #expect(result == nil) + #expect(ContinuousClock().now - startedAt < .milliseconds(500)) + + mailbox.offer(4) + #expect(await mailbox.next(timeout: .milliseconds(50)) == 4) + } + + @Test + func mailboxCancellationRacingOfferLeavesMailboxUsable() async { + for value in 0..<100 { + let mailbox = VideoMailbox() + let waiter = Task { await mailbox.next(timeout: .seconds(5)) } + async let cancellation: Void = Task { waiter.cancel() }.value + async let offer: Void = Task { mailbox.offer(value) }.value + _ = await (cancellation, offer) + let received = await waiter.value + #expect(received == nil || received == value) + mailbox.offer(value + 1) + #expect(await mailbox.next(timeout: .milliseconds(50)) == value + 1) + mailbox.finish() + } + } + + @Test + func mailboxCancelledWaitDoesNotRetainTimeout() async { + weak var releasedMailbox: VideoMailbox? + let task = Task { + let mailbox = VideoMailbox() + releasedMailbox = mailbox + withUnsafeCurrentTask { $0?.cancel() } + #expect(await mailbox.next(timeout: .seconds(5)) == nil) + } + await task.value + let deadline = ContinuousClock.now.advanced(by: .milliseconds(500)) + while releasedMailbox != nil, ContinuousClock.now < deadline { + await Task.yield() + } + #expect(releasedMailbox == nil) + } + + @Test + func cursorCancellationErrorStopsReconciliation() async { + var attempts = 0 + await MacScreenCapture.reconcileCursorConfigurationWithRetry { + attempts += 1 + if attempts == 1 { throw CancellationError() } + } + #expect(attempts == 1) + } + + @Test + func cursorTransientFailureRetriesReconciliation() async { + var attempts = 0 + await MacScreenCapture.reconcileCursorConfigurationWithRetry { + attempts += 1 + if attempts == 1 { throw NSError(domain: "CrabfleetMacTests", code: 1) } + } + #expect(attempts == 2) + } + + @Test + func cursorCancelledTaskDoesNotStartReconciliation() async { + let task = Task { + withUnsafeCurrentTask { $0?.cancel() } + var attempts = 0 + await MacScreenCapture.reconcileCursorConfigurationWithRetry { attempts += 1 } + return attempts + } + #expect(await task.value == 0) + } + + @Test + func cursorCancellationStopsRetryBackoff() async { + let startedAt = ContinuousClock.now + let task = Task { + var attempts = 0 + await MacScreenCapture.reconcileCursorConfigurationWithRetry { + attempts += 1 + withUnsafeCurrentTask { $0?.cancel() } + throw NSError(domain: "CrabfleetMacTests", code: 1) + } + return attempts + } + #expect(await task.value == 1) + #expect(startedAt.duration(to: .now) < .milliseconds(500)) + } + @Test func videoNegotiationPrefersHEVCThenH264ThenTight() { let offered = [