mirror of
https://github.com/permissionlesstech/bitchat.git
synced 2026-07-27 05:25:22 +00:00
Make peer registry and local announce state lock-backed
The main actor answered isPeerConnected/peerNickname/currentPeerSnapshots and flipped runtime capability bits by blocking on collectionsQueue behind whatever transport work was in flight. Peer state now lives in a lock-backed BLEPeerRegistryStore (every registry mutation is a single whole-transition method, so readers never observe a torn state), and the runtime capability bits move into BLELocalIdentityStateStore next to the identity they ride announces with. No transport entry point called from the main actor blocks on a transport queue for peer state anymore. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
@@ -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
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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<T>(_ body: (BLEPeerRegistry) -> T) -> T {
|
||||
lock.withLock { body(registry) }
|
||||
}
|
||||
|
||||
func mutate<T>(_ 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) }
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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))
|
||||
|
||||
Reference in New Issue
Block a user