diff --git a/PeerConnectivity.xcodeproj/project.pbxproj b/PeerConnectivity.xcodeproj/project.pbxproj index 0404a82..bc2a407 100644 --- a/PeerConnectivity.xcodeproj/project.pbxproj +++ b/PeerConnectivity.xcodeproj/project.pbxproj @@ -42,6 +42,10 @@ 30IDENTTEST2606230002 /* PeerIdentityTests.swift in Sources */ = {isa = PBXBuildFile; fileRef = 30IDENTTEST2606230001 /* PeerIdentityTests.swift */; }; 30NETPROTOTEST2606232 /* PeerNetworkProtocolTests.swift in Sources */ = {isa = PBXBuildFile; fileRef = 30NETPROTOTEST2606231 /* PeerNetworkProtocolTests.swift */; }; 30OBSTEST2607210002 /* ObservableThreadSafetyTests.swift in Sources */ = {isa = PBXBuildFile; fileRef = 30OBSTEST2607210001 /* ObservableThreadSafetyTests.swift */; }; + 30ASYNCOBS260721002 /* AsyncObservable.swift in Sources */ = {isa = PBXBuildFile; fileRef = 30ASYNCOBS260721001 /* AsyncObservable.swift */; }; + 30ASYNCMULTI2607212 /* AsyncMultiObservable.swift in Sources */ = {isa = PBXBuildFile; fileRef = 30ASYNCMULTI2607211 /* AsyncMultiObservable.swift */; }; + 30SYNCBRIDGE2607212 /* SyncObservableBridge.swift in Sources */ = {isa = PBXBuildFile; fileRef = 30SYNCBRIDGE2607211 /* SyncObservableBridge.swift */; }; + 30ASYNCTEST26072102 /* AsyncObservableTests.swift in Sources */ = {isa = PBXBuildFile; fileRef = 30ASYNCTEST26072101 /* AsyncObservableTests.swift */; }; 30SECURITY26060300000001 /* PeerSecurityConfiguration.swift in Sources */ = {isa = PBXBuildFile; fileRef = 30SECURITY26060300000002 /* PeerSecurityConfiguration.swift */; }; 30SECURITY26060300000003 /* PeerSecurityConfigurationTests.swift in Sources */ = {isa = PBXBuildFile; fileRef = 30SECURITY26060300000004 /* PeerSecurityConfigurationTests.swift */; }; B20000022F30600000000001 /* ObservableTests.swift in Sources */ = {isa = PBXBuildFile; fileRef = B20000022F30600000000002 /* ObservableTests.swift */; }; @@ -76,6 +80,10 @@ 30IDENTTEST2606230001 /* PeerIdentityTests.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = PeerIdentityTests.swift; sourceTree = ""; }; 30NETPROTOTEST2606231 /* PeerNetworkProtocolTests.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = PeerNetworkProtocolTests.swift; sourceTree = ""; }; 30OBSTEST2607210001 /* ObservableThreadSafetyTests.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = ObservableThreadSafetyTests.swift; sourceTree = ""; }; + 30ASYNCOBS260721001 /* AsyncObservable.swift */ = {isa = PBXFileReference; fileEncoding = 4; lastKnownFileType = sourcecode.swift; name = AsyncObservable.swift; path = Sources/AsyncObservable.swift; sourceTree = ""; }; + 30ASYNCMULTI2607211 /* AsyncMultiObservable.swift */ = {isa = PBXFileReference; fileEncoding = 4; lastKnownFileType = sourcecode.swift; name = AsyncMultiObservable.swift; path = Sources/AsyncMultiObservable.swift; sourceTree = ""; }; + 30SYNCBRIDGE2607211 /* SyncObservableBridge.swift */ = {isa = PBXFileReference; fileEncoding = 4; lastKnownFileType = sourcecode.swift; name = SyncObservableBridge.swift; path = Sources/SyncObservableBridge.swift; sourceTree = ""; }; + 30ASYNCTEST26072101 /* AsyncObservableTests.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = AsyncObservableTests.swift; sourceTree = ""; }; B20000022F30600000000002 /* ObservableTests.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = ObservableTests.swift; sourceTree = ""; }; B20000022F30600000000004 /* PeerTests.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = PeerTests.swift; sourceTree = ""; }; 3080C7DB1D80A1D600AF9EA3 /* Info.plist */ = {isa = PBXFileReference; fileEncoding = 4; lastKnownFileType = text.plist.xml; name = Info.plist; path = Sources/Info.plist; sourceTree = ""; }; @@ -128,6 +136,9 @@ 3080C7DB1D80A1D600AF9EA3 /* Info.plist */, 3080C7DC1D80A1D600AF9EA3 /* MultiObservable.swift */, 3080C7DD1D80A1D600AF9EA3 /* Observable.swift */, + 30ASYNCOBS260721001 /* AsyncObservable.swift */, + 30ASYNCMULTI2607211 /* AsyncMultiObservable.swift */, + 30SYNCBRIDGE2607211 /* SyncObservableBridge.swift */, 30NETREG260624000001 /* NetworkPeerConnectionRegistry.swift */, 30NETDATA26062400001 /* NetworkPeerDataSender.swift */, 30NETCOORD2606250001 /* NetworkPeerCoordinator.swift */, @@ -185,6 +196,7 @@ 30NETSTATETEST2606241 /* NetworkPeerStateMappingTests.swift */, 30NETPROTOTEST2606231 /* PeerNetworkProtocolTests.swift */, 30OBSTEST2607210001 /* ObservableThreadSafetyTests.swift */, + 30ASYNCTEST26072101 /* AsyncObservableTests.swift */, 30PEERMSG2602020000000001 /* PeerMessageTests.swift */, 30SECURITY26060300000004 /* PeerSecurityConfigurationTests.swift */, B20000022F30600000000002 /* ObservableTests.swift */, @@ -329,6 +341,9 @@ 30SECURITY26060300000001 /* PeerSecurityConfiguration.swift in Sources */, 3080C7FB1D80A1D700AF9EA3 /* PeerSession.swift in Sources */, 3080C7EE1D80A1D700AF9EA3 /* Observable.swift in Sources */, + 30ASYNCOBS260721002 /* AsyncObservable.swift in Sources */, + 30ASYNCMULTI2607212 /* AsyncMultiObservable.swift in Sources */, + 30SYNCBRIDGE2607212 /* SyncObservableBridge.swift in Sources */, 3080C7ED1D80A1D700AF9EA3 /* MultiObservable.swift in Sources */, 3080C7F61D80A1D700AF9EA3 /* PeerBrowserEventProducer.swift in Sources */, ); @@ -347,6 +362,7 @@ 30NETSTATETEST2606242 /* NetworkPeerStateMappingTests.swift in Sources */, 30NETPROTOTEST2606232 /* PeerNetworkProtocolTests.swift in Sources */, 30OBSTEST2607210002 /* ObservableThreadSafetyTests.swift in Sources */, + 30ASYNCTEST26072102 /* AsyncObservableTests.swift in Sources */, 30PEERMSG2602020000000002 /* PeerMessageTests.swift in Sources */, 30SECURITY26060300000003 /* PeerSecurityConfigurationTests.swift in Sources */, B20000022F30600000000001 /* ObservableTests.swift in Sources */, diff --git a/PeerConnectivityTests/AsyncObservableTests.swift b/PeerConnectivityTests/AsyncObservableTests.swift new file mode 100644 index 0000000..355741a --- /dev/null +++ b/PeerConnectivityTests/AsyncObservableTests.swift @@ -0,0 +1,162 @@ +// +// AsyncObservableTests.swift +// PeerConnectivityTests +// +// Created by Reid Chatham on 7/21/26. +// Copyright © 2026 Reid Chatham. All rights reserved. +// + +import XCTest +@testable import PeerConnectivity + +@available(iOS 13.0, macOS 10.15, *) +final class AsyncObservableTests: XCTestCase { + + // MARK: - AsyncObservable + + func testAsyncObservableReplaysCurrentValueWhenObserverIsAdded() async { + let observable = AsyncObservable(42) + let lock = NSLock() + var deliveredValues : [Int] = [] + + await observable.addObserver { value in + lock.lock() + deliveredValues.append(value) + lock.unlock() + } + + XCTAssertEqual(deliveredValues, [42]) + } + + func testAsyncObservableDeliversUpdatesInRegistrationOrder() async { + let observable = AsyncObservable(0) + let lock = NSLock() + var deliveredValues : [String] = [] + + await observable.addObserver { value in + lock.lock() + deliveredValues.append("first-\(value)") + lock.unlock() + } + await observable.addObserver { value in + lock.lock() + deliveredValues.append("second-\(value)") + lock.unlock() + } + await observable.update(1) + + let currentValue = await observable.value + XCTAssertEqual(deliveredValues, ["first-0", "second-0", "first-1", "second-1"]) + XCTAssertEqual(currentValue, 1) + } + + func testAsyncObservableHandlesConcurrentUpdates() async { + let observable = AsyncObservable(0) + let lock = NSLock() + var deliveryCount = 0 + + await observable.addObserver { _ in + lock.lock() + deliveryCount += 1 + lock.unlock() + } + + await withTaskGroup(of: Void.self) { group in + for index in 1...100 { + group.addTask { + await observable.update(index) + } + } + } + + XCTAssertGreaterThanOrEqual(deliveryCount, 101) + } + + // MARK: - AsyncMultiObservable + + func testAsyncMultiObservableCanRemoveObserverDuringDelivery() async { + let observable = AsyncMultiObservable(0) + let lock = NSLock() + var deliveredValues : [Int] = [] + + await observable.addObserver({ value in + lock.lock() + deliveredValues.append(value) + lock.unlock() + Task { + await observable.removeObserverForkey("self-removing") + } + }, key: "self-removing") + + await observable.update(1) + await waitForAsyncRemoval(from: observable) + let observerCount = await observable.observerCount + await observable.update(2) + + XCTAssertEqual(observerCount, 0) + XCTAssertEqual(deliveredValues, [0, 1]) + } + + private func waitForAsyncRemoval(from observable: AsyncMultiObservable) async { + for _ in 0..<10 { + if await observable.observerCount == 0 { return } + await Task.yield() + } + } + + // MARK: - SyncObservableBridge + + func testSyncObservableBridgeWaitsForAsyncObservableDelivery() { + let observable = AsyncObservable(0) + let queue = DispatchQueue(label: "PeerConnectivity.AsyncObservableTests.syncBridge") + let completed = expectation(description: "sync bridge completed") + let lock = NSLock() + var deliveredValues : [Int] = [] + + queue.async { + SyncObservableBridge.addObserverAndWait({ value in + lock.lock() + deliveredValues.append(value) + lock.unlock() + }, to: observable) + SyncObservableBridge.updateAndWait(observable, value: 1) + + lock.lock() + let values = deliveredValues + lock.unlock() + + XCTAssertEqual(values, [0, 1]) + completed.fulfill() + } + + wait(for: [completed], timeout: 5) + } + + func testSyncObservableBridgeWaitsForAsyncMultiObservableDelivery() { + let observable = AsyncMultiObservable(0) + let queue = DispatchQueue(label: "PeerConnectivity.AsyncObservableTests.syncMultiBridge") + let completed = expectation(description: "sync multi bridge completed") + let lock = NSLock() + var deliveredValues : [Int] = [] + + queue.async { + SyncObservableBridge.addObserverAndWait({ value in + lock.lock() + deliveredValues.append(value) + lock.unlock() + }, to: observable, key: "listener") + SyncObservableBridge.updateAndWait(observable, value: 1) + SyncObservableBridge.removeObserverAndWait(forKey: "listener", from: observable) + SyncObservableBridge.updateAndWait(observable, value: 2) + + lock.lock() + let values = deliveredValues + lock.unlock() + + XCTAssertEqual(values, [0, 1]) + completed.fulfill() + } + + wait(for: [completed], timeout: 5) + } +} diff --git a/Sources/AsyncMultiObservable.swift b/Sources/AsyncMultiObservable.swift new file mode 100644 index 0000000..cf34e3a --- /dev/null +++ b/Sources/AsyncMultiObservable.swift @@ -0,0 +1,52 @@ +// +// AsyncMultiObservable.swift +// PeerConnectivity +// +// Created by Reid Chatham on 7/21/26. +// Copyright © 2026 Reid Chatham. All rights reserved. +// + +import Foundation + +@available(iOS 13.0, macOS 10.15, *) +internal actor AsyncMultiObservable { + internal typealias Observer = (T) -> Void + + fileprivate var storedValue : T + fileprivate var observers : [String:Observer] = [:] + + internal var value : T { + return storedValue + } + + internal var observerCount : Int { + return observers.count + } + + internal init(_ v: T) { + storedValue = v + } + + internal func addObserver(_ observer: @escaping Observer, key: String) { + let currentValue = storedValue + observers[key] = observer + observer(currentValue) + } + + internal func removeObserverForkey(_ key: String) { + observers.removeValue(forKey: key) + } + + internal func removeAllObservers() { + observers.removeAll() + } + + internal func update(_ newValue: T) { + storedValue = newValue + let currentObservers = Array(observers.values) + + currentObservers.forEach { observer in + observer(newValue) + } + } +} diff --git a/Sources/AsyncObservable.swift b/Sources/AsyncObservable.swift new file mode 100644 index 0000000..86c660b --- /dev/null +++ b/Sources/AsyncObservable.swift @@ -0,0 +1,44 @@ +// +// AsyncObservable.swift +// PeerConnectivity +// +// Created by Reid Chatham on 7/21/26. +// Copyright © 2026 Reid Chatham. All rights reserved. +// + +import Foundation + +@available(iOS 13.0, macOS 10.15, *) +internal actor AsyncObservable { + internal typealias Observer = (T) -> Void + + fileprivate var storedValue : T + fileprivate var observers : [Observer] = [] + + internal var value : T { + return storedValue + } + + internal init(_ v: T) { + storedValue = v + } + + internal func addObserver(_ observer: @escaping Observer) { + let currentValue = storedValue + observers.append(observer) + observer(currentValue) + } + + internal func removeAllObservers() { + observers.removeAll() + } + + internal func update(_ newValue: T) { + storedValue = newValue + let currentObservers = observers + + currentObservers.forEach { observer in + observer(newValue) + } + } +} diff --git a/Sources/SyncObservableBridge.swift b/Sources/SyncObservableBridge.swift new file mode 100644 index 0000000..abd352a --- /dev/null +++ b/Sources/SyncObservableBridge.swift @@ -0,0 +1,73 @@ +// +// SyncObservableBridge.swift +// PeerConnectivity +// +// Created by Reid Chatham on 7/21/26. +// Copyright © 2026 Reid Chatham. All rights reserved. +// + +import Foundation + +@available(iOS 13.0, macOS 10.15, *) +internal enum SyncObservableBridge { + /// Runs an async observable operation and blocks the current thread until it completes. + /// + /// Use this bridge only from synchronous delegate/GCD entry points that must preserve + /// current in-line delivery semantics. Prefer `await` from async contexts. Do not call + /// this from actor-isolated code or cooperative Swift-concurrency executors, because + /// blocking those executors can deadlock or starve unrelated tasks. + internal static func waitForAsync(_ operation: @escaping () async -> Void) { + let semaphore = DispatchSemaphore(value: 0) + + Task { + await operation() + semaphore.signal() + } + + semaphore.wait() + } + + internal static func addObserverAndWait(_ observer: @escaping AsyncObservable.Observer, + to observable: AsyncObservable) { + waitForAsync { + await observable.addObserver(observer) + } + } + + internal static func updateAndWait(_ observable: AsyncObservable, value: T) { + waitForAsync { + await observable.update(value) + } + } + + internal static func removeAllObserversAndWait(from observable: AsyncObservable) { + waitForAsync { + await observable.removeAllObservers() + } + } + + internal static func addObserverAndWait(_ observer: @escaping AsyncMultiObservable.Observer, + to observable: AsyncMultiObservable, key: String) { + waitForAsync { + await observable.addObserver(observer, key: key) + } + } + + internal static func updateAndWait(_ observable: AsyncMultiObservable, value: T) { + waitForAsync { + await observable.update(value) + } + } + + internal static func removeObserverAndWait(forKey key: String, from observable: AsyncMultiObservable) { + waitForAsync { + await observable.removeObserverForkey(key) + } + } + + internal static func removeAllObserversAndWait(from observable: AsyncMultiObservable) { + waitForAsync { + await observable.removeAllObservers() + } + } +}