Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
16 changes: 16 additions & 0 deletions PeerConnectivity.xcodeproj/project.pbxproj
Original file line number Diff line number Diff line change
Expand Up @@ -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 */; };
Expand Down Expand Up @@ -76,6 +80,10 @@
30IDENTTEST2606230001 /* PeerIdentityTests.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = PeerIdentityTests.swift; sourceTree = "<group>"; };
30NETPROTOTEST2606231 /* PeerNetworkProtocolTests.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = PeerNetworkProtocolTests.swift; sourceTree = "<group>"; };
30OBSTEST2607210001 /* ObservableThreadSafetyTests.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = ObservableThreadSafetyTests.swift; sourceTree = "<group>"; };
30ASYNCOBS260721001 /* AsyncObservable.swift */ = {isa = PBXFileReference; fileEncoding = 4; lastKnownFileType = sourcecode.swift; name = AsyncObservable.swift; path = Sources/AsyncObservable.swift; sourceTree = "<group>"; };
30ASYNCMULTI2607211 /* AsyncMultiObservable.swift */ = {isa = PBXFileReference; fileEncoding = 4; lastKnownFileType = sourcecode.swift; name = AsyncMultiObservable.swift; path = Sources/AsyncMultiObservable.swift; sourceTree = "<group>"; };
30SYNCBRIDGE2607211 /* SyncObservableBridge.swift */ = {isa = PBXFileReference; fileEncoding = 4; lastKnownFileType = sourcecode.swift; name = SyncObservableBridge.swift; path = Sources/SyncObservableBridge.swift; sourceTree = "<group>"; };
30ASYNCTEST26072101 /* AsyncObservableTests.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = AsyncObservableTests.swift; sourceTree = "<group>"; };
B20000022F30600000000002 /* ObservableTests.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = ObservableTests.swift; sourceTree = "<group>"; };
B20000022F30600000000004 /* PeerTests.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = PeerTests.swift; sourceTree = "<group>"; };
3080C7DB1D80A1D600AF9EA3 /* Info.plist */ = {isa = PBXFileReference; fileEncoding = 4; lastKnownFileType = text.plist.xml; name = Info.plist; path = Sources/Info.plist; sourceTree = "<group>"; };
Expand Down Expand Up @@ -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 */,
Expand Down Expand Up @@ -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 */,
Expand Down Expand Up @@ -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 */,
);
Expand All @@ -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 */,
Expand Down
162 changes: 162 additions & 0 deletions PeerConnectivityTests/AsyncObservableTests.swift
Original file line number Diff line number Diff line change
@@ -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<Int>(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<Int>(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<Int>(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<Int>(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<Int>) async {
for _ in 0..<10 {
if await observable.observerCount == 0 { return }
await Task.yield()
}
}

// MARK: - SyncObservableBridge

func testSyncObservableBridgeWaitsForAsyncObservableDelivery() {
let observable = AsyncObservable<Int>(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<Int>(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)
}
}
52 changes: 52 additions & 0 deletions Sources/AsyncMultiObservable.swift
Original file line number Diff line number Diff line change
@@ -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<T> {
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)
}
}
}
44 changes: 44 additions & 0 deletions Sources/AsyncObservable.swift
Original file line number Diff line number Diff line change
@@ -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<T> {
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)
}
}
}
73 changes: 73 additions & 0 deletions Sources/SyncObservableBridge.swift
Original file line number Diff line number Diff line change
@@ -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<T>(_ observer: @escaping AsyncObservable<T>.Observer,
to observable: AsyncObservable<T>) {
waitForAsync {
await observable.addObserver(observer)
}
}

internal static func updateAndWait<T>(_ observable: AsyncObservable<T>, value: T) {
waitForAsync {
await observable.update(value)
}
}

internal static func removeAllObserversAndWait<T>(from observable: AsyncObservable<T>) {
waitForAsync {
await observable.removeAllObservers()
}
}

internal static func addObserverAndWait<T>(_ observer: @escaping AsyncMultiObservable<T>.Observer,
to observable: AsyncMultiObservable<T>, key: String) {
waitForAsync {
await observable.addObserver(observer, key: key)
}
}

internal static func updateAndWait<T>(_ observable: AsyncMultiObservable<T>, value: T) {
waitForAsync {
await observable.update(value)
}
}

internal static func removeObserverAndWait<T>(forKey key: String, from observable: AsyncMultiObservable<T>) {
waitForAsync {
await observable.removeObserverForkey(key)
}
}

internal static func removeAllObserversAndWait<T>(from observable: AsyncMultiObservable<T>) {
waitForAsync {
await observable.removeAllObservers()
}
}
}
Loading