diff --git a/bitchat/Services/BLE/BLELocalIdentityStateStore.swift b/bitchat/Services/BLE/BLELocalIdentityStateStore.swift index b379b0cd..7b5be20b 100644 --- a/bitchat/Services/BLE/BLELocalIdentityStateStore.swift +++ b/bitchat/Services/BLE/BLELocalIdentityStateStore.swift @@ -5,6 +5,20 @@ struct BLELocalIdentitySnapshot: Equatable, Sendable { let peerID: PeerID let peerIDData: Data let nickname: String + /// Runtime-toggled capability bits (e.g. the internet-gateway toggle) + /// ORed into `PeerCapabilities.localSupported` for every announce. + let runtimeCapabilities: PeerCapabilities + /// Rendezvous cell advertised while bridging; rides announces only + /// while the `.bridge` capability is enabled. + let bridgeGeohash: String? + + var advertisedCapabilities: PeerCapabilities { + PeerCapabilities.localSupported.union(runtimeCapabilities) + } + + var advertisedBridgeGeohash: String? { + runtimeCapabilities.contains(.bridge) ? bridgeGeohash : nil + } } /// Lock-backed local identity state shared by the transport's message, @@ -12,8 +26,8 @@ struct BLELocalIdentitySnapshot: Equatable, Sendable { /// /// `peerID` and its binary wire representation must change as one unit during /// panic rotation. A snapshot also gives announce construction one consistent -/// view of the nickname and identity instead of reading three independently -/// mutable properties across queues. +/// view of the nickname, identity, and advertised capabilities instead of +/// reading independently mutable properties across queues. final class BLELocalIdentityStateStore: @unchecked Sendable { private let lock = NSLock() private var state: BLELocalIdentitySnapshot @@ -25,7 +39,9 @@ final class BLELocalIdentityStateStore: @unchecked Sendable { state = BLELocalIdentitySnapshot( peerID: peerID, peerIDData: Data(hexString: peerID.id) ?? Data(), - nickname: nickname + nickname: nickname, + runtimeCapabilities: [], + bridgeGeohash: nil ) } @@ -38,7 +54,9 @@ final class BLELocalIdentityStateStore: @unchecked Sendable { state = BLELocalIdentitySnapshot( peerID: state.peerID, peerIDData: state.peerIDData, - nickname: nickname + nickname: nickname, + runtimeCapabilities: state.runtimeCapabilities, + bridgeGeohash: state.bridgeGeohash ) } } @@ -48,8 +66,48 @@ final class BLELocalIdentityStateStore: @unchecked Sendable { state = BLELocalIdentitySnapshot( peerID: peerID, peerIDData: Data(hexString: peerID.id) ?? Data(), - nickname: state.nickname + nickname: state.nickname, + runtimeCapabilities: state.runtimeCapabilities, + bridgeGeohash: state.bridgeGeohash ) } } + + /// Flips a runtime capability bit. Returns whether anything changed. + @discardableResult + func setCapability(_ capability: PeerCapabilities, enabled: Bool) -> Bool { + lock.withLock { + var capabilities = state.runtimeCapabilities + if enabled { + capabilities.insert(capability) + } else { + capabilities.remove(capability) + } + guard capabilities != state.runtimeCapabilities else { return false } + state = BLELocalIdentitySnapshot( + peerID: state.peerID, + peerIDData: state.peerIDData, + nickname: state.nickname, + runtimeCapabilities: capabilities, + bridgeGeohash: state.bridgeGeohash + ) + return true + } + } + + /// Sets the bridged rendezvous cell. Returns whether anything changed. + @discardableResult + func setBridgeGeohash(_ cell: String?) -> Bool { + lock.withLock { + guard cell != state.bridgeGeohash else { return false } + state = BLELocalIdentitySnapshot( + peerID: state.peerID, + peerIDData: state.peerIDData, + nickname: state.nickname, + runtimeCapabilities: state.runtimeCapabilities, + bridgeGeohash: cell + ) + return true + } + } } diff --git a/bitchat/Services/BLE/BLEPeerRegistryStore.swift b/bitchat/Services/BLE/BLEPeerRegistryStore.swift new file mode 100644 index 00000000..bc39a765 --- /dev/null +++ b/bitchat/Services/BLE/BLEPeerRegistryStore.swift @@ -0,0 +1,87 @@ +import BitFoundation +import Foundation + +/// Lock-backed shared ownership of the peer registry, readable from any +/// queue or the main actor without hopping onto a transport queue. +/// +/// Mutations stay serialized by the transport (they only run on its +/// queues), so the lock's job is to let the main actor answer questions +/// like `isPeerConnected` without blocking behind in-flight transport +/// work. Every `BLEPeerRegistry` mutation is a single whole-transition +/// method, so a reader between two mutations always observes a valid +/// pre- or post-state, never a torn one. +/// +/// Closures passed to `read`/`mutate` run under the (non-recursive) lock +/// and must not call back into the store. +final class BLEPeerRegistryStore: @unchecked Sendable { + private let lock = NSLock() + private var registry = BLEPeerRegistry() + + /// One consistent view across multiple registry reads. + func read(_ body: (BLEPeerRegistry) -> T) -> T { + lock.withLock { body(registry) } + } + + func mutate(_ body: (inout BLEPeerRegistry) -> T) -> T { + lock.withLock { body(®istry) } + } + + // MARK: - Single-question reads + + var isEmpty: Bool { read { $0.isEmpty } } + var count: Int { read { $0.count } } + var peerIDs: [PeerID] { read { $0.peerIDs } } + var connectedCount: Int { read { $0.connectedCount } } + var connectedPeerIDs: [PeerID] { read { $0.connectedPeerIDs } } + var connectedRoutingData: [Data] { read { $0.connectedRoutingData } } + var snapshotByID: [PeerID: BLEPeerInfo] { read { $0.snapshotByID } } + + func info(for peerID: PeerID) -> BLEPeerInfo? { + read { $0.info(for: peerID) } + } + + func isConnected(_ peerID: PeerID) -> Bool { + read { $0.isConnected(peerID) } + } + + func isReachable(_ peerID: PeerID, now: Date) -> Bool { + read { $0.isReachable(peerID, now: now) } + } + + func nickname(for peerID: PeerID, connectedOnly: Bool) -> String? { + read { $0.nickname(for: peerID, connectedOnly: connectedOnly) } + } + + func fingerprint(for peerID: PeerID) -> String? { + read { $0.fingerprint(for: peerID) } + } + + func capabilities(for peerID: PeerID) -> PeerCapabilities { + read { $0.capabilities(for: peerID) } + } + + func capabilitiesWereExplicitlyAdvertised(for peerID: PeerID) -> Bool { + read { $0.capabilitiesWereExplicitlyAdvertised(for: peerID) } + } + + func advertisedBridgeGeohash() -> String? { + read { $0.advertisedBridgeGeohash() } + } + + func displayNicknames(selfNickname: String) -> [PeerID: String] { + read { $0.displayNicknames(selfNickname: selfNickname) } + } + + func transportSnapshots(selfNickname: String) -> [TransportPeerSnapshot] { + read { $0.transportSnapshots(selfNickname: selfNickname) } + } + + /// Peers advertising `capability` that are reachable now, in one + /// consistent view. + func reachablePeers(advertising capability: PeerCapabilities, now: Date) -> [PeerID] { + read { registry in + registry.peers(advertising: capability) + .filter { registry.isReachable($0, now: now) } + } + } +} diff --git a/bitchat/Services/BLE/BLEService.swift b/bitchat/Services/BLE/BLEService.swift index 919c3820..882958c2 100644 --- a/bitchat/Services/BLE/BLEService.swift +++ b/bitchat/Services/BLE/BLEService.swift @@ -254,8 +254,10 @@ final class BLEService: NSObject { // BCH-01-004: Rate-limiting for subscription-triggered announces. private var subscriptionAnnounceLimiter = BLESubscriptionAnnounceLimiter() - // 3. Peer Information (single source of truth) - private var peerRegistry = BLEPeerRegistry() + // 3. Peer Information (single source of truth). Lock-backed so the main + // actor reads it directly instead of blocking on collectionsQueue; + // mutations still only run on the transport's queues. + private let peerRegistry = BLEPeerRegistryStore() // 4. Efficient Message Deduplication private let messageDeduplicator = MessageDeduplicator() @@ -297,8 +299,6 @@ final class BLEService: NSObject { /// Fired (off-main) when a signature-verified announce is processed — /// the bridge courier watch refreshes its tag set on new arrivals. var onVerifiedPeerAnnounce: ((_ peerID: PeerID) -> Void)? - private var runtimeCapabilities: PeerCapabilities = [] // collectionsQueue - private var localBridgeGeohash: String? // collectionsQueue #if DEBUG // Test-only tap on the outbound pipeline so multi-node tests can ferry @@ -896,9 +896,7 @@ final class BLEService: NSObject { weak var peerEventsDelegate: TransportPeerEventsDelegate? func currentPeerSnapshots() -> [TransportPeerSnapshot] { - collectionsQueue.sync { - peerRegistry.transportSnapshots(selfNickname: myNickname) - } + peerRegistry.transportSnapshots(selfNickname: myNickname) } // MARK: Identity @@ -1077,7 +1075,7 @@ final class BLEService: NSObject { maxResponses: TransportConfig.meshPingInboundMaxPerLink, window: TransportConfig.meshPingInboundWindowSeconds ) - peerRegistry.removeAll() + peerRegistry.mutate { $0.removeAll() } fragmentAssemblyBuffer.removeAll() sourceRouteFailures = BLESourceRouteFailureCache() // Also clear pending message queues to avoid stale state across sessions @@ -1110,14 +1108,12 @@ final class BLEService: NSObject { func isPeerConnected(_ peerID: PeerID) -> Bool { // Accept both 16-hex short IDs and 64-hex Noise keys - return collectionsQueue.sync { peerRegistry.isConnected(peerID) } + return peerRegistry.isConnected(peerID) } func isPeerReachable(_ peerID: PeerID) -> Bool { // Accept both 16-hex short IDs and 64-hex Noise keys - return collectionsQueue.sync { - peerRegistry.isReachable(peerID, now: Date()) - } + peerRegistry.isReachable(peerID, now: Date()) } func canDeliverSecurely(to peerID: PeerID) -> Bool { @@ -1134,15 +1130,13 @@ final class BLEService: NSObject { } func peerNickname(peerID: PeerID) -> String? { - collectionsQueue.sync { - peerRegistry.nickname(for: peerID, connectedOnly: true) - } + peerRegistry.nickname(for: peerID, connectedOnly: true) } /// Capabilities the peer advertised in its last verified announce. /// Empty for peers that predate the capabilities TLV. func peerCapabilities(_ peerID: PeerID) -> PeerCapabilities { - collectionsQueue.sync { peerRegistry.capabilities(for: peerID) } + peerRegistry.capabilities(for: peerID) } func authenticatedPrivateMediaReceiptSessionGeneration( @@ -1183,11 +1177,9 @@ final class BLEService: NSObject { // registry entry populated by a public announce. return fingerprint } - return collectionsQueue.sync { - peerRegistry.info(for: normalizedPeerID)? - .noisePublicKey? - .sha256Fingerprint() - } + return peerRegistry.info(for: normalizedPeerID)? + .noisePublicKey? + .sha256Fingerprint() } func privateMediaSendPolicy(to peerID: PeerID) -> PrivateMediaSendPolicy { @@ -1418,67 +1410,41 @@ final class BLEService: NSObject { /// internet-gateway toggle) and re-announces so peers learn promptly. /// Build-time bits stay in `PeerCapabilities.localSupported`. func setLocalCapability(_ capability: PeerCapabilities, enabled: Bool) { - let changed: Bool = collectionsQueue.sync(flags: .barrier) { - let before = runtimeCapabilities - if enabled { - runtimeCapabilities.insert(capability) - } else { - runtimeCapabilities.remove(capability) - } - return runtimeCapabilities != before - } - guard changed else { return } + guard localIdentityState.setCapability(capability, enabled: enabled) else { return } sendAnnounce(forceSend: true) } /// Reachable peers currently advertising the `.gateway` capability. func reachableGatewayPeers() -> [PeerID] { - let now = Date() - return collectionsQueue.sync { - peerRegistry.peers(advertising: .gateway) - .filter { peerRegistry.isReachable($0, now: now) } - } + peerRegistry.reachablePeers(advertising: .gateway, now: Date()) } /// Reachable peers currently advertising the `.bridge` capability. func reachableBridgePeers() -> [PeerID] { - let now = Date() - return collectionsQueue.sync { - peerRegistry.peers(advertising: .bridge) - .filter { peerRegistry.isReachable($0, now: now) } - } + peerRegistry.reachablePeers(advertising: .bridge, now: Date()) } /// A rendezvous cell advertised by a bridge-capable peer's announce. func advertisedBridgeGeohash() -> String? { - collectionsQueue.sync { peerRegistry.advertisedBridgeGeohash() } + peerRegistry.advertisedBridgeGeohash() } /// The rendezvous cell this device advertises in its own announces while /// bridging with the gateway toggle on. Set from the main actor; the /// value rides the next (forced) announce. func setLocalBridgeGeohash(_ cell: String?) { - let changed: Bool = collectionsQueue.sync(flags: .barrier) { - guard localBridgeGeohash != cell else { return false } - localBridgeGeohash = cell - return true - } - guard changed else { return } + guard localIdentityState.setBridgeGeohash(cell) else { return } sendAnnounce(forceSend: true) } func getPeerNicknames() -> [PeerID: String] { - return collectionsQueue.sync { - peerRegistry.displayNicknames(selfNickname: myNickname) - } + peerRegistry.displayNicknames(selfNickname: myNickname) } // MARK: Protocol utilities func getFingerprint(for peerID: PeerID) -> String? { - return collectionsQueue.sync { - peerRegistry.fingerprint(for: peerID) - } + peerRegistry.fingerprint(for: peerID) } func getNoiseSessionState(for peerID: PeerID) -> LazyHandshakeState { @@ -2624,7 +2590,7 @@ final class BLEService: NSObject { let content = String(data: packet.payload, encoding: .utf8)?.trimmedOrNilIfEmpty else { return nil } let senderPeerID = PeerID(hexData: packet.senderID) - let peers = collectionsQueue.sync { peerRegistry.snapshotByID } + let peers = peerRegistry.snapshotByID // Archived senders are usually long gone, so the signature-derived // identity is the best shot at a name; a live registry entry is // next; anonymous fallback matches the live path. @@ -2663,7 +2629,7 @@ final class BLEService: NSObject { }, peersSnapshot: { [weak self] in guard let self = self else { return [:] } - return self.collectionsQueue.sync { self.peerRegistry.snapshotByID } + return self.peerRegistry.snapshotByID }, verifyPacketSignature: { [weak self] packet, signingPublicKey in self?.noiseService.verifyPacketSignature(packet, publicKey: signingPublicKey) ?? false @@ -2846,10 +2812,8 @@ final class BLEService: NSObject { noiseReconnectPolicy.endLinkEpoch(link) } } - _ = collectionsQueue.sync(flags: .barrier) { - // Remove the peer when they leave - peerRegistry.remove(peerID) - } + // Remove the peer when they leave + peerRegistry.mutate { _ = $0.remove(peerID) } // Remove any stored announcement for sync purposes gossipSyncManager?.removeAnnouncementForPeer(peerID) // Send on main thread @@ -2857,7 +2821,7 @@ final class BLEService: NSObject { guard let self = self else { return } // Get current peer list (after removal) - let currentPeerIDs = self.collectionsQueue.sync { self.peerRegistry.peerIDs } + let currentPeerIDs = self.peerRegistry.peerIDs self.deliverTransportEvent(.peerDisconnected(peerID)) self.deliverTransportEvent(.peerListUpdated(currentPeerIDs)) @@ -2890,15 +2854,10 @@ final class BLEService: NSObject { let noisePub = noiseService.getStaticPublicKeyData() // For noise handshakes and peer identification let signingPub = noiseService.getSigningPublicKeyData() // For signature verification - let (connectedPeerIDs, advertisedCapabilities, advertisedBridgeCell): ([Data], PeerCapabilities, String?) = collectionsQueue.sync { - ( - peerRegistry.connectedRoutingData, - PeerCapabilities.localSupported.union(runtimeCapabilities), - runtimeCapabilities.contains(.bridge) ? localBridgeGeohash : nil - ) - } - + let connectedPeerIDs = peerRegistry.connectedRoutingData let localIdentity = localIdentityState.snapshot() + let advertisedCapabilities = localIdentity.advertisedCapabilities + let advertisedBridgeCell = localIdentity.advertisedBridgeGeohash let announcement = AnnouncementPacket( nickname: localIdentity.nickname, noisePublicKey: noisePub, @@ -3351,9 +3310,7 @@ extension BLEService: CBCentralManagerDelegate { let peerStillLinked = (remainingLinks?.hasPeripheral ?? false) || (remainingLinks?.hasCentral ?? false) if let peerID, !peerStillLinked { // Do not remove peer; mark as not connected but retain for reachability - collectionsQueue.sync(flags: .barrier) { - peerRegistry.markDisconnected(peerID) - } + peerRegistry.mutate { $0.markDisconnected(peerID) } refreshLocalTopology() } @@ -3374,7 +3331,7 @@ extension BLEService: CBCentralManagerDelegate { guard let self = self else { return } // Get current peer list (after removal) - let currentPeerIDs = self.collectionsQueue.sync { self.peerRegistry.peerIDs } + let currentPeerIDs = self.peerRegistry.peerIDs if let peerID, !peerStillLinked { self.notifyPeerDisconnectedDebounced(peerID) @@ -3567,15 +3524,15 @@ extension BLEService { if preseedPeer { // Ensure the synthetic peer is known and marked verified for public-message tests let normalizedID = PeerID(hexData: packet.senderID) - collectionsQueue.sync(flags: .barrier) { - if var existing = peerRegistry.info(for: normalizedID) { + peerRegistry.mutate { registry in + if var existing = registry.info(for: normalizedID) { existing.isConnected = true existing.isVerifiedNickname = true if let signingPublicKey { existing.signingPublicKey = signingPublicKey } existing.lastSeen = Date() - peerRegistry.upsert(existing) + registry.upsert(existing) } else { - peerRegistry.upsert(BLEPeerInfo( + registry.upsert(BLEPeerInfo( peerID: normalizedID, nickname: "TestPeer_\(fromPeerID.id.prefix(4))", isConnected: true, @@ -3656,8 +3613,8 @@ extension BLEService { capabilities: PeerCapabilities? = nil, noisePublicKey: Data? = nil ) { - collectionsQueue.sync(flags: .barrier) { - peerRegistry.upsert(BLEPeerInfo( + peerRegistry.mutate { + $0.upsert(BLEPeerInfo( peerID: peerID, nickname: nickname, isConnected: true, @@ -4295,9 +4252,7 @@ extension BLEService: CBPeripheralManagerDelegate { // the bookkeeping. guard linkStateStore.links(to: peerID).isEmpty else { return } // Mark peer as not connected; retain for reachability - collectionsQueue.sync(flags: .barrier) { - peerRegistry.markDisconnected(peerID) - } + peerRegistry.mutate { $0.markDisconnected(peerID) } refreshLocalTopology() @@ -4306,7 +4261,7 @@ extension BLEService: CBPeripheralManagerDelegate { guard let self = self else { return } // Get current peer list (after removal) - let currentPeerIDs = self.collectionsQueue.sync { self.peerRegistry.peerIDs } + let currentPeerIDs = self.peerRegistry.peerIDs self.notifyPeerDisconnectedDebounced(peerID) // Publish snapshots so UnifiedPeerService can refresh icons promptly @@ -5332,9 +5287,7 @@ extension BLEService { } guard shouldSend else { return } - let capabilities = collectionsQueue.sync { - PeerCapabilities.localSupported.union(runtimeCapabilities) - } + let capabilities = localIdentityState.snapshot().advertisedCapabilities let state = AuthenticatedPeerStatePacket( capabilities: capabilities, signingPublicKey: noiseService.getSigningPublicKeyData() @@ -5400,10 +5353,12 @@ extension BLEService { guard privateMediaSessionGenerations[normalizedPeerID] == generation else { return [] } - peerRegistry.bindAuthenticatedSigningPublicKey( - state.signingPublicKey, - for: normalizedPeerID - ) + peerRegistry.mutate { + $0.bindAuthenticatedSigningPublicKey( + state.signingPublicKey, + for: normalizedPeerID + ) + } authenticatedPeerStates[normalizedPeerID] = BLEAuthenticatedPeerStateObservation( fingerprint: fingerprint, sessionGeneration: generation, @@ -5723,7 +5678,7 @@ extension BLEService { /// message consumes exactly one prekey ID regardless of courier count. private func assignRecipientPrekey(messageID: String, recipientNoiseKey: Data) -> PrekeyBundle.Prekey? { let shortID = PeerID(publicKey: recipientNoiseKey) - let knownOnMesh = collectionsQueue.sync { peerRegistry.info(for: shortID) != nil } + let knownOnMesh = peerRegistry.info(for: shortID) != nil if knownOnMesh, !peerCapabilities(shortID).contains(.prekeys) { return nil } @@ -5821,7 +5776,7 @@ extension BLEService { // favorite conversation instead of an unresolvable short-ID // thread labeled "Unknown". let shortID = PeerID(publicKey: senderStaticKey) - let isKnownOnMesh = collectionsQueue.sync { peerRegistry.info(for: shortID) != nil } + let isKnownOnMesh = peerRegistry.info(for: shortID) != nil let senderPeerID = isKnownOnMesh ? shortID : PeerID(hexData: senderStaticKey) SecureLogger.debug("📦 Opened courier envelope from \(senderPeerID.id.prefix(8))…", category: .session) sfMetrics?.record(.courierOpened) @@ -5851,7 +5806,7 @@ extension BLEService { SecureLogger.debug("📦 Courier deposit rejected: relayed envelope claims sender \(PeerID(hexData: packet.senderID).id.prefix(8))… but arrived from \(peerID.id.prefix(8))…", category: .security) return } - let depositorInfo = collectionsQueue.sync { peerRegistry.info(for: peerID) } + let depositorInfo = peerRegistry.info(for: peerID) guard let depositorKey = depositorInfo?.noisePublicKey else { SecureLogger.debug("📦 Courier deposit from unknown peer \(peerID.id.prefix(8))… rejected", category: .session) return @@ -6088,7 +6043,7 @@ extension BLEService { /// from identities persisted for offline verification. private func announceBoundSigningKey(forNoiseKey noiseKey: Data) -> Data? { let shortID = PeerID(publicKey: noiseKey) - if let info = collectionsQueue.sync(execute: { peerRegistry.info(for: shortID) }), + if let info = peerRegistry.info(for: shortID), info.noisePublicKey == noiseKey, let signingKey = info.signingPublicKey { return signingKey @@ -6161,7 +6116,7 @@ extension BLEService { // packet signature must verify against the sender's announced // signing key. Unlike courier deposits the depositor may be // multi-hop away, so ingress-link identity is not required. - let signingKey = collectionsQueue.sync { peerRegistry.info(for: senderID)?.signingPublicKey } + let signingKey = peerRegistry.info(for: senderID)?.signingPublicKey guard let signingKey, noiseService.verifyPacketSignature(packet, publicKey: signingKey) else { SecureLogger.debug("🌐 nostrCarrier uplink from \(senderID.id.prefix(8))… rejected (missing/invalid packet signature)", category: .security) @@ -7101,7 +7056,7 @@ extension BLEService { SecureLogger.debug("⚠️ Duplicate packet ignored: \(messageID.prefix(24))…", category: .session) } - let connectedCount = collectionsQueue.sync { peerRegistry.connectedCount } + let connectedCount = peerRegistry.connectedCount if BLEReceivePipeline.shouldCancelScheduledRelayForDuplicate(connectedPeerCount: connectedCount) { collectionsQueue.async(flags: .barrier) { [weak self] in self?.scheduledRelays.cancel(messageID: messageID) @@ -7112,7 +7067,7 @@ extension BLEService { } private func scheduleRelayIfNeeded(_ packet: BitchatPacket, senderID: PeerID, messageID: String) { - let degree = collectionsQueue.sync { peerRegistry.connectedCount } + let degree = peerRegistry.connectedCount let decision = BLEReceivePipeline.relayDecision( for: packet, senderID: senderID, @@ -7400,15 +7355,13 @@ extension BLEService { /// link. The `.peerConnected` UI event already fired from the announce /// path (new/reconnected + direct), so only list state needs refreshing. private func promoteReboundPeerToConnected(_ peerID: PeerID) { - let promoted = collectionsQueue.sync(flags: .barrier) { - peerRegistry.markConnected(peerID) - } + let promoted = peerRegistry.mutate { $0.markConnected(peerID) } guard promoted else { return } refreshLocalTopology() publishFullPeerData() notifyUI { [weak self] in guard let self else { return } - let currentPeerIDs = self.collectionsQueue.sync { self.peerRegistry.peerIDs } + let currentPeerIDs = self.peerRegistry.peerIDs self.deliverTransportEvent(.peerListUpdated(currentPeerIDs)) } } @@ -7417,15 +7370,13 @@ extension BLEService { /// instead of letting a ghost duplicate linger for the reachability /// retention window. private func retireRotatedPeer(_ peerID: PeerID) { - let removed = collectionsQueue.sync(flags: .barrier) { - peerRegistry.remove(peerID) != nil - } + let removed = peerRegistry.mutate { $0.remove(peerID) != nil } guard removed else { return } gossipSyncManager?.removeAnnouncementForPeer(peerID) refreshLocalTopology() notifyUI { [weak self] in guard let self else { return } - let currentPeerIDs = self.collectionsQueue.sync { self.peerRegistry.peerIDs } + let currentPeerIDs = self.peerRegistry.peerIDs self.deliverTransportEvent(.peerDisconnected(peerID)) self.deliverTransportEvent(.peerListUpdated(currentPeerIDs)) } @@ -7494,21 +7445,23 @@ extension BLEService { guard let self = self else { return BLEPeerAnnounceUpdate(isNewPeer: false, wasDisconnected: false, previousNickname: nil) } - return self.peerRegistry.upsertVerifiedAnnounce( - peerID: peerID, - nickname: announcement.nickname, - noisePublicKey: announcement.noisePublicKey, - signingPublicKey: announcement.signingPublicKey, - isConnected: isConnected, - // Propagate `nil` (registry refused the announce because it - // carries a signing key different from the pinned one) so - // the handler's guard rejects it instead of overwriting the - // pinned identity. Main's capabilities/bridgeGeohash are - // preserved. - now: now, - capabilities: announcement.capabilities, - bridgeGeohash: announcement.bridgeGeohash - ) + return self.peerRegistry.mutate { + $0.upsertVerifiedAnnounce( + peerID: peerID, + nickname: announcement.nickname, + noisePublicKey: announcement.noisePublicKey, + signingPublicKey: announcement.signingPublicKey, + isConnected: isConnected, + // Propagate `nil` (registry refused the announce because it + // carries a signing key different from the pinned one) so + // the handler's guard rejects it instead of overwriting the + // pinned identity. Main's capabilities/bridgeGeohash are + // preserved. + now: now, + capabilities: announcement.capabilities, + bridgeGeohash: announcement.bridgeGeohash + ) + } }, shouldEmitReconnectLog: { [weak self] peerID, now in // Called from inside withRegistryBarrier; access debouncer directly. @@ -7547,7 +7500,7 @@ extension BLEService { self.gossipSyncManager?.scheduleInitialSyncToPeer(peerID, delaySeconds: 1.0) } // Get current peer list (after addition) - let currentPeerIDs = self.collectionsQueue.sync { self.peerRegistry.peerIDs } + let currentPeerIDs = self.peerRegistry.peerIDs self.requestPeerDataPublish() self.deliverTransportEvent(.peerListUpdated(currentPeerIDs)) } @@ -7638,7 +7591,7 @@ extension BLEService { // A response can replay the entire gossip store, so require proof the // requester owns the claimed sender ID: the request must verify // against the signing key from that peer's announce. - let signingKey = collectionsQueue.sync { peerRegistry.info(for: peerID)?.signingPublicKey } + let signingKey = peerRegistry.info(for: peerID)?.signingPublicKey guard let signingKey, noiseService.verifyPacketSignature(packet, publicKey: signingKey) else { if logRateLimiter.shouldLog(key: "sync-sig:\(peerID.id)") { SecureLogger.warning("🚫 Dropping REQUEST_SYNC without verifiable signature from \(peerID.id.prefix(8))…", category: .security) @@ -7672,7 +7625,7 @@ extension BLEService { now: { Date() }, peersSnapshot: { [weak self] in guard let self = self else { return [:] } - return self.collectionsQueue.sync { self.peerRegistry.snapshotByID } + return self.peerRegistry.snapshotByID }, verifyPacketSignature: { [weak self] packet, signingPublicKey in self?.noiseService.verifyPacketSignature(packet, publicKey: signingPublicKey) ?? false @@ -7742,7 +7695,7 @@ extension BLEService { maxAgeSeconds: TransportConfig.pttPublicFrameMaxAgeSeconds ) else { return false } - let peersSnapshot = collectionsQueue.sync { peerRegistry.snapshotByID } + let peersSnapshot = peerRegistry.snapshotByID let registrySigningKey = peersSnapshot[peerID]?.signingPublicKey let verifiedViaRegistry = registrySigningKey.map { noiseService.verifyPacketSignature(packet, publicKey: $0) } ?? false let signedDisplayName = verifiedViaRegistry ? nil : signedSenderDisplayName(for: packet, from: peerID) @@ -7970,10 +7923,7 @@ extension BLEService { } private func updatePeerLastSeen(_ peerID: PeerID) { - // Use async to avoid deadlock - we don't need immediate consistency for last seen updates - collectionsQueue.async(flags: .barrier) { - self.peerRegistry.updateLastSeen(peerID, at: Date()) - } + peerRegistry.mutate { $0.updateLastSeen(peerID, at: Date()) } } // Debounced disconnect notifier to avoid duplicate disconnect callbacks within a short window @@ -8006,12 +7956,12 @@ extension BLEService { lastMaintenanceAt = Date() let now = Date() - let connectedCount = collectionsQueue.sync { peerRegistry.connectedCount } + let connectedCount = peerRegistry.connectedCount let elapsed = announceThrottle.elapsed(since: now) let recentSeen = collectionsQueue.sync { () -> Bool in recentTrafficTracker.hasTraffic(within: 5.0, now: now) } - let hasNoPeers = collectionsQueue.sync { peerRegistry.isEmpty } + let hasNoPeers = peerRegistry.isEmpty let plan = BLEMaintenancePolicy.plan( cycle: maintenanceCounter, connectedCount: connectedCount, @@ -8082,7 +8032,7 @@ extension BLEService { private func checkPeerConnectivity() { let now = Date() - let peerIDsForLinkState: [PeerID] = collectionsQueue.sync { peerRegistry.peerIDs } + let peerIDsForLinkState: [PeerID] = peerRegistry.peerIDs var cachedLinkStates: [PeerID: BLEPeerLinkPresence] = [:] for peerID in peerIDsForLinkState { let state = linkState(for: peerID) @@ -8092,8 +8042,8 @@ extension BLEService { ) } - let changes = collectionsQueue.sync(flags: .barrier) { - peerRegistry.reconcileConnectivity(now: now, linkStates: cachedLinkStates) + let changes = peerRegistry.mutate { + $0.reconcileConnectivity(now: now, linkStates: cachedLinkStates) } for removedPeer in changes.removedPeers { SecureLogger.debug("🗑️ Removing stale peer after reachability window: \(removedPeer.peerID.id.prefix(8))… (\(removedPeer.nickname))", category: .session) @@ -8106,7 +8056,7 @@ extension BLEService { guard let self else { return } // Get current peer list (after removal) - let currentPeerIDs = self.collectionsQueue.sync { self.peerRegistry.peerIDs } + let currentPeerIDs = self.peerRegistry.peerIDs for peerID in changes.disconnectedPeerIDs { self.deliverTransportEvent(.peerDisconnected(peerID))