From bf641100e745cc676c56d6e77ea5ae0c94c28859 Mon Sep 17 00:00:00 2001 From: Arthur Liu Date: Wed, 22 Jul 2026 10:42:38 -0700 Subject: [PATCH] Fix concurrent share subscriptions --- RxSwift/Observables/ShareReplayScope.swift | 2 +- Sources/AllTestz/main.swift | 1 + Tests/RxSwiftTests/Anomalies.swift | 54 ++++++++++++++++++++++ 3 files changed, 56 insertions(+), 1 deletion(-) diff --git a/RxSwift/Observables/ShareReplayScope.swift b/RxSwift/Observables/ShareReplayScope.swift index dab3e82f2..e5732404e 100644 --- a/RxSwift/Observables/ShareReplayScope.swift +++ b/RxSwift/Observables/ShareReplayScope.swift @@ -408,8 +408,8 @@ private final class ShareWhileConnected: let connection = synchronized_subscribe(observer) let count = connection.observers.count - lock.unlock() let disposable = connection.synchronized_subscribe(observer) + lock.unlock() if count == 0 { connection.connect() diff --git a/Sources/AllTestz/main.swift b/Sources/AllTestz/main.swift index 511927e7e..5cf128a41 100644 --- a/Sources/AllTestz/main.swift +++ b/Sources/AllTestz/main.swift @@ -24,6 +24,7 @@ final class AnomaliesTest_ : AnomaliesTest, RxTestCase { ("test1323", AnomaliesTest.test1323), ("test1344", AnomaliesTest.test1344), ("testSeparationBetweenOnAndSubscriptionLocks", AnomaliesTest.testSeparationBetweenOnAndSubscriptionLocks), + ("test2718ShareWithoutReplayConnectsOnceForConcurrentFirstSubscriptions", AnomaliesTest.test2718ShareWithoutReplayConnectsOnceForConcurrentFirstSubscriptions), ("test2653ShareReplayOneInitialEmissionDeadlock", AnomaliesTest.test2653ShareReplayOneInitialEmissionDeadlock), ("test2653ShareReplayMoreInitialEmissionDeadlock", AnomaliesTest.test2653ShareReplayMoreInitialEmissionDeadlock), ("test2653ShareReplayOneForeverInitialEmissionDeadlock", AnomaliesTest.test2653ShareReplayOneForeverInitialEmissionDeadlock), diff --git a/Tests/RxSwiftTests/Anomalies.swift b/Tests/RxSwiftTests/Anomalies.swift index aa0a56bc4..b66da73ad 100644 --- a/Tests/RxSwiftTests/Anomalies.swift +++ b/Tests/RxSwiftTests/Anomalies.swift @@ -176,6 +176,60 @@ extension AnomaliesTest { } } + func test2718ShareWithoutReplayConnectsOnceForConcurrentFirstSubscriptions() { + let subscriberCount = 32 + + for _ in 0 ..< 100 { + let subscriptionsLock = NSLock() + var sourceSubscriptionCount = 0 + var subscriptions = [Disposable]() + + let shared = Observable.never() + .do(onSubscribe: { + subscriptionsLock.lock() + sourceSubscriptionCount += 1 + subscriptionsLock.unlock() + }) + .share(replay: 0, scope: .whileConnected) + + let ready = DispatchGroup() + let finished = DispatchGroup() + let start = DispatchSemaphore(value: 0) + + for _ in 0 ..< subscriberCount { + ready.enter() + finished.enter() + + DispatchQueue.global(qos: .userInitiated).async { + ready.leave() + start.wait() + + let subscription = shared.subscribe() + + subscriptionsLock.lock() + subscriptions.append(subscription) + subscriptionsLock.unlock() + + finished.leave() + } + } + + XCTAssertEqual(ready.wait(timeout: .now() + 5), .success) + for _ in 0 ..< subscriberCount { + start.signal() + } + XCTAssertEqual(finished.wait(timeout: .now() + 5), .success) + + subscriptionsLock.lock() + let subscriptionCount = sourceSubscriptionCount + let subscriptionsToDispose = subscriptions + subscriptionsLock.unlock() + + XCTAssertEqual(subscriptionCount, 1) + subscriptionsToDispose.forEach { $0.dispose() } + } + } + func test2653ShareReplayOneInitialEmissionDeadlock() { let immediatelyEmittingSource = Observable.create { observer in observer.on(.next(()))