mirror of
https://github.com/permissionlesstech/bitchat.git
synced 2026-07-25 02:25:20 +00:00
Harden PTT audio and the 1.7.1 release (#1423)
Centralize PTT and voice audio-session ownership, harden courier/bridge/outbox delivery and recovery, correct location and delivery-state races, add privacy/release metadata, and ship reproducible universal Arti slices with Release CI coverage. Validated by the full iOS suite, repeated audio/fragment/performance regressions, BitFoundation tests, strict lint and dead-code analysis, universal iOS Release builds, and iOS/macOS archives.
This commit is contained in:
@@ -37,6 +37,12 @@ final class BLEService: NSObject {
|
||||
// 1. Consolidated BLE link tracking for both central and peripheral roles.
|
||||
private var linkStateStore = BLELinkStateStore()
|
||||
|
||||
// A peer ID can retain an established Noise session after its physical
|
||||
// link disappears. Courier handover therefore needs the stronger fact
|
||||
// that the session was established *on this current ingress link*, not
|
||||
// merely that some session exists for the claimed ID. bleQueue-owned.
|
||||
private var noiseAuthenticatedLinkOwners: [BLEIngressLinkID: PeerID] = [:]
|
||||
|
||||
// Rotation-rebind cooldown per link UUID (bleQueue-owned, like the link
|
||||
// store): entries older than the cooldown are pruned on insert.
|
||||
private var lastLinkRebindAt: [String: Date] = [:]
|
||||
@@ -680,6 +686,7 @@ final class BLEService: NSObject {
|
||||
// Clear peripheral references (synchronized access to avoid races with BLE callbacks)
|
||||
bleQueue.sync {
|
||||
linkStateStore.clearAll()
|
||||
noiseAuthenticatedLinkOwners.removeAll()
|
||||
connectionScheduler.reset()
|
||||
subscriptionAnnounceLimiter.removeAll()
|
||||
}
|
||||
@@ -1188,9 +1195,91 @@ final class BLEService: NSObject {
|
||||
}
|
||||
}
|
||||
|
||||
private func sendOnAllLinks(packet: BitchatPacket, data: Data, pad: Bool, directedOnlyPeer: PeerID?) {
|
||||
/// Synchronously admits a notification to the link-specific retry queue.
|
||||
/// Destructive courier handoff uses this result as its commit point, so a
|
||||
/// full process-local queue must be reported as rejection, not success.
|
||||
private func enqueuePendingNotificationIfAccepted(
|
||||
data: Data,
|
||||
centrals: [CBCentral],
|
||||
context: String
|
||||
) -> Bool {
|
||||
let result = collectionsQueue.sync(flags: .barrier) {
|
||||
pendingNotifications.enqueue(
|
||||
data: data,
|
||||
targets: centrals,
|
||||
capCount: TransportConfig.blePendingNotificationsCapCount
|
||||
)
|
||||
}
|
||||
switch result {
|
||||
case let .enqueued(count):
|
||||
SecureLogger.debug("📋 Queued \(context) packet for retry (pending=\(count))", category: .session)
|
||||
return true
|
||||
case let .full(count):
|
||||
SecureLogger.warning("⚠️ Rejecting \(context) packet: notification queue full (pending=\(count))", category: .session)
|
||||
return false
|
||||
}
|
||||
}
|
||||
|
||||
/// Serializes the final authenticated-link check with CoreBluetooth's
|
||||
/// notification admission on `bleQueue`, closing the rebind/disconnect
|
||||
/// race between fanout planning and the actual handoff.
|
||||
private func notifyOrEnqueueIfAccepted(
|
||||
data: Data,
|
||||
centrals: [CBCentral],
|
||||
characteristic: CBMutableCharacteristic,
|
||||
context: String,
|
||||
requiredAuthenticatedPeer: PeerID?
|
||||
) -> Bool {
|
||||
let accept = { [self] in
|
||||
let eligible: [CBCentral]
|
||||
if let peerID = requiredAuthenticatedPeer {
|
||||
eligible = centrals.filter { central in
|
||||
let link = BLEIngressLinkID.central(central.identifier.uuidString)
|
||||
return noiseAuthenticatedLinkOwners[link] == peerID
|
||||
&& linkStateStore.peerID(forCentralUUID: central.identifier.uuidString) == peerID
|
||||
}
|
||||
} else {
|
||||
eligible = centrals
|
||||
}
|
||||
guard !eligible.isEmpty else { return false }
|
||||
if peripheralManager?.updateValue(data, for: characteristic, onSubscribedCentrals: eligible) == true {
|
||||
return true
|
||||
}
|
||||
return enqueuePendingNotificationIfAccepted(
|
||||
data: data,
|
||||
centrals: eligible,
|
||||
context: context
|
||||
)
|
||||
}
|
||||
|
||||
if DispatchQueue.getSpecific(key: bleQueueKey) != nil {
|
||||
return accept()
|
||||
}
|
||||
return bleQueue.sync(execute: accept)
|
||||
}
|
||||
|
||||
/// Returns true only when the packet was accepted by at least one current
|
||||
/// physical link (including its link-specific backpressure queue). A
|
||||
/// process-local directed spool is deliberately not success: callers
|
||||
/// that own a durable upstream copy must keep it retryable.
|
||||
@discardableResult
|
||||
private func sendOnAllLinks(
|
||||
packet: BitchatPacket,
|
||||
data: Data,
|
||||
pad: Bool,
|
||||
directedOnlyPeer: PeerID?,
|
||||
requireDirectPeerLink: Bool = false,
|
||||
requireNoiseAuthenticatedPeerLink: Bool = false
|
||||
) -> Bool {
|
||||
let ingressRecord = collectionsQueue.sync { ingressLinks.record(for: packet) }
|
||||
let excludedPeerLinks = links(to: ingressRecord?.peerID)
|
||||
var excludedPeerLinks = links(to: ingressRecord?.peerID)
|
||||
if requireNoiseAuthenticatedPeerLink {
|
||||
guard let directedOnlyPeer else { return false }
|
||||
let boundLinks = links(to: directedOnlyPeer)
|
||||
let authenticatedLinks = currentNoiseAuthenticatedLinks(to: directedOnlyPeer)
|
||||
guard !authenticatedLinks.isEmpty else { return false }
|
||||
excludedPeerLinks.formUnion(boundLinks.subtracting(authenticatedLinks))
|
||||
}
|
||||
let outboundPriority = BLEOutboundPacketPolicy.priority(for: packet, data: data)
|
||||
|
||||
let states = snapshotPeripheralStates()
|
||||
@@ -1222,12 +1311,22 @@ final class BLEService: NSObject {
|
||||
// as a combined snapshot.
|
||||
preferredPeripheralPerPeer: readLinkState { $0.preferredPeripheralBindings },
|
||||
directAnnounceTTL: messageTTL,
|
||||
directedOnlyPeer: directedOnlyPeer
|
||||
directedOnlyPeer: directedOnlyPeer,
|
||||
requireDirectPeerLink: requireDirectPeerLink || requireNoiseAuthenticatedPeerLink
|
||||
)
|
||||
|
||||
if let chunk = plan.fragmentChunkSize {
|
||||
sendFragmentedPacket(packet, pad: pad, maxChunk: chunk, directedOnlyPeer: directedOnlyPeer)
|
||||
return
|
||||
guard !plan.selectedLinks.peripheralIDs.isEmpty || !plan.selectedLinks.centralIDs.isEmpty else {
|
||||
return false
|
||||
}
|
||||
return sendFragmentedPacket(
|
||||
packet,
|
||||
pad: pad,
|
||||
maxChunk: chunk,
|
||||
directedOnlyPeer: directedOnlyPeer,
|
||||
requireDirectPeerLink: requireDirectPeerLink || requireNoiseAuthenticatedPeerLink,
|
||||
requireNoiseAuthenticatedPeerLink: requireNoiseAuthenticatedPeerLink
|
||||
)
|
||||
}
|
||||
|
||||
// If directed and we currently have no links to forward on, spool for a short window
|
||||
@@ -1236,36 +1335,73 @@ final class BLEService: NSObject {
|
||||
spoolDirectedPacket(packet, recipientPeerID: only)
|
||||
}
|
||||
|
||||
var acceptedByPhysicalLink = false
|
||||
|
||||
// Writes to selected connected peripherals
|
||||
for s in connectedStates {
|
||||
let pid = s.peripheral.identifier.uuidString
|
||||
guard plan.selectedLinks.peripheralIDs.contains(pid) else { continue }
|
||||
if let ch = s.characteristic {
|
||||
writeOrEnqueue(data, to: s.peripheral, characteristic: ch, priority: outboundPriority)
|
||||
if requireDirectPeerLink || requireNoiseAuthenticatedPeerLink {
|
||||
acceptedByPhysicalLink = writeOrEnqueueIfAccepted(
|
||||
data,
|
||||
to: s.peripheral,
|
||||
characteristic: ch,
|
||||
priority: outboundPriority,
|
||||
requiredAuthenticatedPeer: requireNoiseAuthenticatedPeerLink ? directedOnlyPeer : nil
|
||||
) || acceptedByPhysicalLink
|
||||
} else {
|
||||
writeOrEnqueue(data, to: s.peripheral, characteristic: ch, priority: outboundPriority)
|
||||
}
|
||||
}
|
||||
}
|
||||
// Notify selected subscribed centrals
|
||||
if let ch = characteristic {
|
||||
let targets = subscribedCentrals.filter { plan.selectedLinks.centralIDs.contains($0.identifier.uuidString) }
|
||||
if !targets.isEmpty {
|
||||
let success = peripheralManager?.updateValue(data, for: ch, onSubscribedCentrals: targets) ?? false
|
||||
if !success {
|
||||
// Notification queue full - queue for retry to prevent silent packet loss
|
||||
// This is critical for fragment delivery reliability
|
||||
let context = packet.type == MessageType.fragment.rawValue ? "fragment" : "broadcast"
|
||||
enqueuePendingNotification(data: data, centrals: targets, context: context)
|
||||
if requireDirectPeerLink || requireNoiseAuthenticatedPeerLink {
|
||||
acceptedByPhysicalLink = notifyOrEnqueueIfAccepted(
|
||||
data: data,
|
||||
centrals: targets,
|
||||
characteristic: ch,
|
||||
context: "directed",
|
||||
requiredAuthenticatedPeer: requireNoiseAuthenticatedPeerLink ? directedOnlyPeer : nil
|
||||
) || acceptedByPhysicalLink
|
||||
} else {
|
||||
let success = peripheralManager?.updateValue(data, for: ch, onSubscribedCentrals: targets) ?? false
|
||||
if !success {
|
||||
// Notification queue full - queue for retry to prevent silent packet loss
|
||||
// This is critical for fragment delivery reliability
|
||||
let context = packet.type == MessageType.fragment.rawValue ? "fragment" : "broadcast"
|
||||
enqueuePendingNotification(data: data, centrals: targets, context: context)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
if requireDirectPeerLink || requireNoiseAuthenticatedPeerLink { return acceptedByPhysicalLink }
|
||||
return !plan.selectedLinks.peripheralIDs.isEmpty || !plan.selectedLinks.centralIDs.isEmpty
|
||||
}
|
||||
|
||||
// Directed send helper (unicast to a specific peerID) without altering packet contents
|
||||
private func sendPacketDirected(_ packet: BitchatPacket, to peerID: PeerID) {
|
||||
@discardableResult
|
||||
private func sendPacketDirected(
|
||||
_ packet: BitchatPacket,
|
||||
to peerID: PeerID,
|
||||
requireDirectPeerLink: Bool = false,
|
||||
requireNoiseAuthenticatedPeerLink: Bool = false
|
||||
) -> Bool {
|
||||
#if DEBUG
|
||||
_test_onOutboundPacket?(packet)
|
||||
#endif
|
||||
guard let data = packet.toBinaryData(padding: false) else { return }
|
||||
sendOnAllLinks(packet: packet, data: data, pad: false, directedOnlyPeer: peerID)
|
||||
guard let data = packet.toBinaryData(padding: false) else { return false }
|
||||
return sendOnAllLinks(
|
||||
packet: packet,
|
||||
data: data,
|
||||
pad: false,
|
||||
directedOnlyPeer: peerID,
|
||||
requireDirectPeerLink: requireDirectPeerLink,
|
||||
requireNoiseAuthenticatedPeerLink: requireNoiseAuthenticatedPeerLink
|
||||
)
|
||||
}
|
||||
|
||||
// MARK: - Directed store-and-forward
|
||||
@@ -1938,6 +2074,10 @@ extension BLEService: CBCentralManagerDelegate {
|
||||
#endif
|
||||
|
||||
// Clean up references and peer mappings
|
||||
collectionsQueue.sync(flags: .barrier) {
|
||||
pendingPeripheralWrites.discardAll(for: peripheralID)
|
||||
}
|
||||
noiseAuthenticatedLinkOwners.removeValue(forKey: .peripheral(peripheralID))
|
||||
_ = linkStateStore.removePeripheral(peripheralID)
|
||||
// A duplicate link can drop while the peer stays live on another
|
||||
// (the dual-role central link, or a second bound link after a
|
||||
@@ -1988,6 +2128,10 @@ extension BLEService: CBCentralManagerDelegate {
|
||||
let peripheralID = peripheral.identifier.uuidString
|
||||
|
||||
// Clean up the references
|
||||
collectionsQueue.sync(flags: .barrier) {
|
||||
pendingPeripheralWrites.discardAll(for: peripheralID)
|
||||
}
|
||||
noiseAuthenticatedLinkOwners.removeValue(forKey: .peripheral(peripheralID))
|
||||
_ = linkStateStore.removePeripheral(peripheralID)
|
||||
|
||||
SecureLogger.error("❌ Failed to connect to peripheral: \(peripheral.name ?? "Unknown") [\(peripheralID)] - Error: \(error?.localizedDescription ?? "Unknown")", category: .session)
|
||||
@@ -2088,6 +2232,10 @@ extension BLEService {
|
||||
|
||||
SecureLogger.debug("⏱️ Timeout: \(candidate.name)", category: .session)
|
||||
central.cancelPeripheralConnection(peripheral)
|
||||
self.collectionsQueue.sync(flags: .barrier) {
|
||||
self.pendingPeripheralWrites.discardAll(for: peripheralID)
|
||||
}
|
||||
self.noiseAuthenticatedLinkOwners.removeValue(forKey: .peripheral(peripheralID))
|
||||
_ = self.linkStateStore.removePeripheral(peripheralID)
|
||||
self.connectionScheduler.recordConnectionTimeout(peripheralID: peripheralID, at: Date())
|
||||
self.tryConnectFromQueue()
|
||||
@@ -2134,6 +2282,23 @@ extension BLEService {
|
||||
handleReceivedPacket(packet, from: fromPeerID)
|
||||
}
|
||||
|
||||
/// Waits until fragment ingress already submitted by a test has finished
|
||||
/// reassembly/reinjection and any resulting transport event has crossed
|
||||
/// the MainActor delivery hop. This is a deterministic pipeline fence,
|
||||
/// avoiding wall-clock sleeps that become flaky under a parallel suite.
|
||||
func _test_drainFragmentPipeline() async {
|
||||
await withCheckedContinuation { continuation in
|
||||
messageQueue.async(flags: .barrier) {
|
||||
// Reassembled packets are reinjected synchronously on
|
||||
// `messageQueue`; their UI delivery task is therefore already
|
||||
// enqueued before this later MainActor marker.
|
||||
Task { @MainActor in
|
||||
continuation.resume()
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func _test_hasGossipPrekeyBundle(for peerID: PeerID) -> Bool {
|
||||
gossipSyncManager?._hasPrekeyBundle(for: peerID) ?? false
|
||||
}
|
||||
@@ -2164,6 +2329,13 @@ extension BLEService {
|
||||
bleQueue.sync { linkStateStore.peerID(forCentralUUID: centralUUID) }
|
||||
}
|
||||
|
||||
func _test_markNoiseAuthenticatedCentral(_ centralUUID: String, to peerID: PeerID) {
|
||||
bleQueue.sync {
|
||||
guard linkStateStore.peerID(forCentralUUID: centralUUID) == peerID else { return }
|
||||
noiseAuthenticatedLinkOwners[.central(centralUUID)] = peerID
|
||||
}
|
||||
}
|
||||
|
||||
func _test_seedConnectedPeer(_ peerID: PeerID, nickname: String) {
|
||||
collectionsQueue.sync(flags: .barrier) {
|
||||
peerRegistry.upsert(BLEPeerInfo(
|
||||
@@ -2578,7 +2750,12 @@ extension BLEService: CBPeripheralManagerDelegate {
|
||||
}
|
||||
|
||||
func peripheralManager(_ peripheral: CBPeripheralManager, central: CBCentral, didUnsubscribeFrom characteristic: CBCharacteristic) {
|
||||
SecureLogger.debug("📤 Central unsubscribed: \(central.identifier.uuidString.prefix(8))…", category: .session)
|
||||
let centralID = central.identifier.uuidString
|
||||
SecureLogger.debug("📤 Central unsubscribed: \(centralID.prefix(8))…", category: .session)
|
||||
collectionsQueue.sync(flags: .barrier) {
|
||||
pendingNotifications.removeTarget { $0.identifier.uuidString == centralID }
|
||||
}
|
||||
noiseAuthenticatedLinkOwners.removeValue(forKey: .central(centralID))
|
||||
let removedPeerID = linkStateStore.removeSubscribedCentral(central)
|
||||
|
||||
// Ensure we're still advertising for other devices to find us
|
||||
@@ -3118,6 +3295,45 @@ extension BLEService {
|
||||
private func links(to peerID: PeerID?) -> Set<BLEIngressLinkID> {
|
||||
readLinkState { $0.links(to: peerID) }
|
||||
}
|
||||
|
||||
private func boundPeerID(for link: BLEIngressLinkID, in store: BLELinkStateStore) -> PeerID? {
|
||||
switch link {
|
||||
case .peripheral(let peripheralUUID):
|
||||
store.peerID(forPeripheralID: peripheralUUID)
|
||||
case .central(let centralUUID):
|
||||
store.peerID(forCentralUUID: centralUUID)
|
||||
}
|
||||
}
|
||||
|
||||
/// Marks the exact physical ingress link that completed a fresh Noise
|
||||
/// handshake. An old session keyed only by peer ID is insufficient: a
|
||||
/// replayed announce can rebind an attacker's link to that ID.
|
||||
private func markNoiseAuthenticatedIngressLink(for packet: BitchatPacket, peerID: PeerID) {
|
||||
guard let link = collectionsQueue.sync(execute: { ingressLinks.link(for: packet) }) else { return }
|
||||
readLinkState { store in
|
||||
guard boundPeerID(for: link, in: store) == peerID else { return }
|
||||
noiseAuthenticatedLinkOwners[link] = peerID
|
||||
}
|
||||
}
|
||||
|
||||
private func isNoiseAuthenticatedIngressLink(for packet: BitchatPacket, peerID: PeerID) -> Bool {
|
||||
guard let link = collectionsQueue.sync(execute: { ingressLinks.link(for: packet) }) else { return false }
|
||||
return readLinkState { store in
|
||||
noiseAuthenticatedLinkOwners[link] == peerID && boundPeerID(for: link, in: store) == peerID
|
||||
}
|
||||
}
|
||||
|
||||
private func hasCurrentNoiseAuthenticatedLink(to peerID: PeerID) -> Bool {
|
||||
!currentNoiseAuthenticatedLinks(to: peerID).isEmpty
|
||||
}
|
||||
|
||||
private func currentNoiseAuthenticatedLinks(to peerID: PeerID) -> Set<BLEIngressLinkID> {
|
||||
readLinkState { store in
|
||||
Set(noiseAuthenticatedLinkOwners.compactMap { link, owner in
|
||||
owner == peerID && boundPeerID(for: link, in: store) == peerID ? link : nil
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
private func configureNoiseServiceCallbacks(for service: NoiseEncryptionService) {
|
||||
service.onPeerAuthenticated = { [weak self] peerID, fingerprint in
|
||||
@@ -3299,22 +3515,32 @@ extension BLEService {
|
||||
guard CourierEnvelope.candidateTags(noiseStaticKey: myKey, around: Date()).contains(envelope.recipientTag) else {
|
||||
return false
|
||||
}
|
||||
openCourierEnvelope(envelope)
|
||||
return true
|
||||
return openCourierEnvelope(envelope)
|
||||
}
|
||||
|
||||
/// Hands a bridge-fetched envelope directly to the matching local peer
|
||||
/// as a directed courier packet. Delivery-only by design: the recipient's
|
||||
/// tag matched, so this never lands in a stranger's carry quota.
|
||||
/// Returns true only if a current Noise-authenticated physical link
|
||||
/// accepted the packet; a stale peer-level session, reachability record,
|
||||
/// replay-rebound link, or process-local spool is not delivery.
|
||||
@discardableResult
|
||||
func deliverBridgedEnvelope(_ envelope: CourierEnvelope, to peerID: PeerID) -> Bool {
|
||||
guard isPeerConnected(peerID) || isPeerReachable(peerID) else { return false }
|
||||
guard hasCurrentNoiseAuthenticatedLink(to: peerID) else { return false }
|
||||
guard let payload = envelope.encode() else { return false }
|
||||
let packet = makeCourierPacket(payload, to: peerID)
|
||||
messageQueue.async { [weak self] in
|
||||
self?.sendPacketDirected(packet, to: peerID)
|
||||
let send = { [weak self] in
|
||||
self?.sendPacketDirected(
|
||||
packet,
|
||||
to: peerID,
|
||||
requireDirectPeerLink: true,
|
||||
requireNoiseAuthenticatedPeerLink: true
|
||||
) ?? false
|
||||
}
|
||||
return true
|
||||
if DispatchQueue.getSpecific(key: messageQueueKey) != nil {
|
||||
return send()
|
||||
}
|
||||
return messageQueue.sync(execute: send)
|
||||
}
|
||||
|
||||
/// Our own Noise static public key (for computing our courier tags).
|
||||
@@ -3385,7 +3611,8 @@ extension BLEService {
|
||||
}
|
||||
}
|
||||
|
||||
private func openCourierEnvelope(_ envelope: CourierEnvelope) {
|
||||
@discardableResult
|
||||
private func openCourierEnvelope(_ envelope: CourierEnvelope) -> Bool {
|
||||
do {
|
||||
let typedPayload: Data
|
||||
let senderStaticKey: Data
|
||||
@@ -3410,12 +3637,12 @@ extension BLEService {
|
||||
let payloadType = NoisePayloadType(rawValue: typeRaw),
|
||||
payloadType == .privateMessage else {
|
||||
SecureLogger.warning("⚠️ Courier envelope carried unsupported payload type", category: .session)
|
||||
return
|
||||
return true // decrypted but deterministically unsupported
|
||||
}
|
||||
let payload = Data(typedPayload.dropFirst())
|
||||
guard let innerMessageID = PrivateMessagePacket.decode(from: payload)?.messageID else {
|
||||
SecureLogger.warning("⚠️ Courier envelope carried undecodable private message", category: .session)
|
||||
return
|
||||
return true // decrypted but deterministically malformed
|
||||
}
|
||||
// Redundant copies of one message arrive as distinct envelopes
|
||||
// (fresh seal each: mesh couriers, bridge drops across relays),
|
||||
@@ -3427,14 +3654,14 @@ extension BLEService {
|
||||
}
|
||||
guard firstOpen else {
|
||||
SecureLogger.debug("📦 Dropping duplicate courier envelope for message \(innerMessageID.prefix(8))…", category: .session)
|
||||
return
|
||||
return true
|
||||
}
|
||||
// Couriered mail arrives while the sender is absent, so the UI's
|
||||
// block check can't resolve their fingerprint from a live session.
|
||||
// Gate here, where the full static key is in hand.
|
||||
guard !identityManager.isBlocked(fingerprint: senderStaticKey.sha256Fingerprint()) else {
|
||||
SecureLogger.debug("🚫 Dropping courier envelope from blocked sender", category: .security)
|
||||
return
|
||||
return true
|
||||
}
|
||||
// A present sender resolves to their live mesh thread via the
|
||||
// derived short ID. An absent sender — the usual courier case —
|
||||
@@ -3454,9 +3681,11 @@ extension BLEService {
|
||||
timestamp: Date()
|
||||
))
|
||||
}
|
||||
return true
|
||||
} catch {
|
||||
// Tag collision or stale key: not addressed to us after all.
|
||||
SecureLogger.debug("📦 Courier envelope failed to open: \(error)", category: .encryption)
|
||||
return false
|
||||
}
|
||||
}
|
||||
|
||||
@@ -3498,13 +3727,23 @@ extension BLEService {
|
||||
|
||||
/// Hand over any carried envelopes addressed to a peer we just heard from.
|
||||
private func deliverCourierMail(to peerID: PeerID, noiseKey: Data) {
|
||||
let envelopes = courierStore.takeEnvelopes(for: noiseKey)
|
||||
guard !envelopes.isEmpty else { return }
|
||||
SecureLogger.debug("📦 Handing over \(envelopes.count) courier envelope(s) to \(peerID.id.prefix(8))…", category: .session)
|
||||
for envelope in envelopes {
|
||||
guard let payload = envelope.encode() else { continue }
|
||||
sendPacketDirected(makeCourierPacket(payload, to: peerID), to: peerID)
|
||||
sfMetrics?.record(.courierHandedOver)
|
||||
let metrics = sfMetrics
|
||||
let accepted = courierStore.handoverEnvelopes(for: noiseKey) { [weak self] envelope in
|
||||
guard let self,
|
||||
let payload = envelope.encode(),
|
||||
self.sendPacketDirected(
|
||||
self.makeCourierPacket(payload, to: peerID),
|
||||
to: peerID,
|
||||
requireDirectPeerLink: true,
|
||||
requireNoiseAuthenticatedPeerLink: true
|
||||
) else {
|
||||
return false
|
||||
}
|
||||
metrics?.record(.courierHandedOver)
|
||||
return true
|
||||
}
|
||||
if accepted > 0 {
|
||||
SecureLogger.debug("📦 Handed over \(accepted) courier envelope(s) to \(peerID.id.prefix(8))…", category: .session)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -3534,13 +3773,23 @@ extension BLEService {
|
||||
private func sprayCourierMail(to peerID: PeerID, noiseKey: Data, isVerifiedPeer: Bool) {
|
||||
let store = courierStore
|
||||
let metrics = sfMetrics
|
||||
let sendSpray: ([CourierEnvelope]) -> Void = { [weak self] envelopes in
|
||||
guard let self, !envelopes.isEmpty else { return }
|
||||
SecureLogger.debug("📦 Spraying \(envelopes.count) envelope copy(ies) to courier \(peerID.id.prefix(8))…", category: .session)
|
||||
for envelope in envelopes {
|
||||
guard let payload = envelope.encode() else { continue }
|
||||
self.sendPacketDirected(self.makeCourierPacket(payload, to: peerID), to: peerID)
|
||||
let sendSpray: () -> Void = { [weak self] in
|
||||
guard let self else { return }
|
||||
let accepted = store.transferSprayCopies(to: noiseKey) { envelope in
|
||||
guard let payload = envelope.encode(),
|
||||
self.sendPacketDirected(
|
||||
self.makeCourierPacket(payload, to: peerID),
|
||||
to: peerID,
|
||||
requireDirectPeerLink: true,
|
||||
requireNoiseAuthenticatedPeerLink: true
|
||||
) else {
|
||||
return false
|
||||
}
|
||||
metrics?.record(.courierSprayed)
|
||||
return true
|
||||
}
|
||||
if accepted > 0 {
|
||||
SecureLogger.debug("📦 Sprayed \(accepted) envelope copy(ies) to courier \(peerID.id.prefix(8))…", category: .session)
|
||||
}
|
||||
}
|
||||
let policy = courierDepositPolicy
|
||||
@@ -3548,7 +3797,7 @@ extension BLEService {
|
||||
// Same trust gate as deposits: don't hand mail to a peer who
|
||||
// would reject it from us.
|
||||
guard policy(noiseKey, isVerifiedPeer) != nil else { return }
|
||||
sendSpray(store.takeSprayCopies(for: noiseKey))
|
||||
sendSpray()
|
||||
}
|
||||
}
|
||||
|
||||
@@ -3829,6 +4078,62 @@ extension BLEService {
|
||||
}
|
||||
}
|
||||
|
||||
/// Writes immediately or synchronously admits the packet to this
|
||||
/// peripheral's bounded retry queue. Unlike `writeOrEnqueue`, the return
|
||||
/// value distinguishes a retained queue item from one rejected or trimmed
|
||||
/// immediately, which lets durable courier state commit truthfully.
|
||||
private func writeOrEnqueueIfAccepted(
|
||||
_ data: Data,
|
||||
to peripheral: CBPeripheral,
|
||||
characteristic: CBCharacteristic,
|
||||
priority: BLEOutboundWritePriority,
|
||||
requiredAuthenticatedPeer: PeerID?
|
||||
) -> Bool {
|
||||
let accept = { [self] in
|
||||
let uuid = peripheral.identifier.uuidString
|
||||
guard let state = linkStateStore.state(forPeripheralID: uuid),
|
||||
state.isConnected,
|
||||
state.characteristic?.uuid == characteristic.uuid else {
|
||||
return false
|
||||
}
|
||||
if let peerID = requiredAuthenticatedPeer {
|
||||
let link = BLEIngressLinkID.peripheral(uuid)
|
||||
guard state.peerID == peerID,
|
||||
noiseAuthenticatedLinkOwners[link] == peerID else {
|
||||
return false
|
||||
}
|
||||
}
|
||||
|
||||
if peripheral.canSendWriteWithoutResponse {
|
||||
peripheral.writeValue(data, for: characteristic, type: .withoutResponse)
|
||||
return true
|
||||
}
|
||||
|
||||
let attempt = collectionsQueue.sync(flags: .barrier) {
|
||||
pendingPeripheralWrites.enqueueReportingAcceptance(
|
||||
data: data,
|
||||
for: uuid,
|
||||
priority: priority,
|
||||
capBytes: TransportConfig.blePendingWriteBufferCapBytes
|
||||
)
|
||||
}
|
||||
switch attempt.result {
|
||||
case .oversized(let bytes):
|
||||
SecureLogger.warning("⚠️ Rejecting oversized write chunk (\(bytes)B) for peripheral \(uuid)", category: .session)
|
||||
case let .enqueued(trimmedBytes, remainingBytes) where trimmedBytes > 0:
|
||||
SecureLogger.warning("📉 Trimmed pending write buffer for \(uuid) by \(trimmedBytes)B to \(remainingBytes)B", category: .session)
|
||||
case .enqueued:
|
||||
break
|
||||
}
|
||||
return attempt.accepted
|
||||
}
|
||||
|
||||
if DispatchQueue.getSpecific(key: bleQueueKey) != nil {
|
||||
return accept()
|
||||
}
|
||||
return bleQueue.sync(execute: accept)
|
||||
}
|
||||
|
||||
private func drainPendingWrites(for peripheral: CBPeripheral) {
|
||||
let uuid = peripheral.identifier.uuidString
|
||||
bleQueue.async { [weak self] in
|
||||
@@ -3976,6 +4281,10 @@ extension BLEService {
|
||||
guard age > TransportConfig.bleConnectTimeoutSeconds else { continue }
|
||||
let peripheralID = state.peripheral.identifier.uuidString
|
||||
central.cancelPeripheralConnection(state.peripheral)
|
||||
self.collectionsQueue.sync(flags: .barrier) {
|
||||
self.pendingPeripheralWrites.discardAll(for: peripheralID)
|
||||
}
|
||||
self.noiseAuthenticatedLinkOwners.removeValue(forKey: .peripheral(peripheralID))
|
||||
_ = self.linkStateStore.removePeripheral(peripheralID)
|
||||
cancelled += 1
|
||||
}
|
||||
@@ -4119,25 +4428,37 @@ extension BLEService {
|
||||
|
||||
// MARK: Fragmentation (Required for messages > BLE MTU)
|
||||
|
||||
private func sendFragmentedPacket(_ packet: BitchatPacket, pad: Bool, maxChunk: Int? = nil, directedOnlyPeer: PeerID? = nil, transferId: String? = nil) {
|
||||
@discardableResult
|
||||
private func sendFragmentedPacket(
|
||||
_ packet: BitchatPacket,
|
||||
pad: Bool,
|
||||
maxChunk: Int? = nil,
|
||||
directedOnlyPeer: PeerID? = nil,
|
||||
transferId: String? = nil,
|
||||
requireDirectPeerLink: Bool = false,
|
||||
requireNoiseAuthenticatedPeerLink: Bool = false
|
||||
) -> Bool {
|
||||
let request = BLEOutboundFragmentTransferRequest(
|
||||
packet: packet,
|
||||
pad: pad,
|
||||
maxChunk: maxChunk,
|
||||
directedPeer: directedOnlyPeer,
|
||||
transferId: transferId
|
||||
transferId: transferId,
|
||||
requireDirectPeerLink: requireDirectPeerLink,
|
||||
requireNoiseAuthenticatedPeerLink: requireNoiseAuthenticatedPeerLink
|
||||
)
|
||||
|
||||
let result = collectionsQueue.sync(flags: .barrier) {
|
||||
outboundFragmentTransfers.submit(request, maxConcurrentTransfers: TransportConfig.bleMaxConcurrentTransfers)
|
||||
}
|
||||
handleFragmentTransferSubmitResult(result)
|
||||
return handleFragmentTransferSubmitResult(result)
|
||||
}
|
||||
|
||||
private func handleFragmentTransferSubmitResult(_ result: BLEOutboundFragmentTransferScheduler.SubmitResult) {
|
||||
@discardableResult
|
||||
private func handleFragmentTransferSubmitResult(_ result: BLEOutboundFragmentTransferScheduler.SubmitResult) -> Bool {
|
||||
switch result {
|
||||
case let .start(request, reservedTransferId):
|
||||
startFragmentedPacket(request, reservedTransferId: reservedTransferId)
|
||||
return startFragmentedPacket(request, reservedTransferId: reservedTransferId)
|
||||
|
||||
case let .queued(_, transferId, _):
|
||||
if let transferId {
|
||||
@@ -4145,16 +4466,29 @@ extension BLEService {
|
||||
} else {
|
||||
SecureLogger.debug("🚦 Queued fragment transfer waiting for slot", category: .session)
|
||||
}
|
||||
return false
|
||||
|
||||
case let .rejectedStrict(_, transferId):
|
||||
SecureLogger.debug(
|
||||
"🚫 Strict directed fragment transfer \(transferId?.prefix(8) ?? "?")… rejected while scheduler busy",
|
||||
category: .session
|
||||
)
|
||||
return false
|
||||
|
||||
case let .droppedDuplicate(_, activeTransferId):
|
||||
SecureLogger.debug(
|
||||
"🔁 Skipping duplicate outbound transfer — same content already in flight as \(activeTransferId?.prefix(8) ?? "?")…",
|
||||
category: .session
|
||||
)
|
||||
return false
|
||||
}
|
||||
}
|
||||
|
||||
private func startFragmentedPacket(_ request: BLEOutboundFragmentTransferRequest, reservedTransferId: String?) {
|
||||
@discardableResult
|
||||
private func startFragmentedPacket(
|
||||
_ request: BLEOutboundFragmentTransferRequest,
|
||||
reservedTransferId: String?
|
||||
) -> Bool {
|
||||
let releaseReservedSlot: (String) -> Void = { [weak self] id in
|
||||
guard let self = self else { return }
|
||||
TransferProgressManager.shared.cancel(id: id)
|
||||
@@ -4174,7 +4508,7 @@ extension BLEService {
|
||||
if let id = reservedTransferId {
|
||||
releaseReservedSlot(id)
|
||||
}
|
||||
return
|
||||
return false
|
||||
}
|
||||
|
||||
// Lightweight pacing to reduce floods and allow BLE buffers to drain
|
||||
@@ -4200,6 +4534,42 @@ extension BLEService {
|
||||
return id
|
||||
}()
|
||||
|
||||
let sendFragment: (BitchatPacket) -> Bool = { [weak self] fragmentPacket in
|
||||
guard let self else { return false }
|
||||
if request.requireDirectPeerLink, let directedPeer = request.directedPeer {
|
||||
return self.sendPacketDirected(
|
||||
fragmentPacket,
|
||||
to: directedPeer,
|
||||
requireDirectPeerLink: true,
|
||||
requireNoiseAuthenticatedPeerLink: request.requireNoiseAuthenticatedPeerLink
|
||||
)
|
||||
}
|
||||
self.broadcastPacket(fragmentPacket)
|
||||
return true
|
||||
}
|
||||
|
||||
// Strict courier handoff is transactional at the fragment-admission
|
||||
// boundary: every fragment must enter the intended authenticated
|
||||
// link or its bounded retry queue before the durable owner may commit.
|
||||
// A partial train is harmlessly abandoned and the envelope stays
|
||||
// retryable with a fresh fragment ID on the next encounter.
|
||||
if request.requireDirectPeerLink {
|
||||
let admitted = BLEStrictFragmentAdmission.admitAll(plan.fragmentPackets) { fragmentPacket in
|
||||
guard sendFragment(fragmentPacket) else { return false }
|
||||
if let transferId = transferIdentifier {
|
||||
markFragmentSent(transferId: transferId)
|
||||
}
|
||||
return true
|
||||
}
|
||||
guard admitted else {
|
||||
if let id = reservedTransferId {
|
||||
releaseReservedSlot(id)
|
||||
}
|
||||
return false
|
||||
}
|
||||
return true
|
||||
}
|
||||
|
||||
var scheduledItems: [(item: DispatchWorkItem, index: Int)] = []
|
||||
|
||||
for (index, fragmentPacket) in plan.fragmentPackets.enumerated() {
|
||||
@@ -4212,7 +4582,7 @@ extension BLEService {
|
||||
if fragmentPacket.recipientID == nil || fragmentPacket.recipientID?.allSatisfy({ $0 == 0xFF }) == true {
|
||||
self.gossipSyncManager?.onPublicPacketSeen(fragmentPacket)
|
||||
}
|
||||
self.broadcastPacket(fragmentPacket)
|
||||
_ = sendFragment(fragmentPacket)
|
||||
if let transferId = transferIdentifier {
|
||||
self.markFragmentSent(transferId: transferId)
|
||||
}
|
||||
@@ -4232,6 +4602,7 @@ extension BLEService {
|
||||
let delayMs = index * plan.spacingMs
|
||||
messageQueue.asyncAfter(deadline: .now() + .milliseconds(delayMs), execute: workItem)
|
||||
}
|
||||
return true
|
||||
}
|
||||
|
||||
// MARK: - Fragmentation (Required for messages > BLE MTU)
|
||||
@@ -4503,15 +4874,32 @@ extension BLEService {
|
||||
let result,
|
||||
result.isVerified else { return }
|
||||
let noiseKey = result.announcement.noisePublicKey
|
||||
if result.isDirectAnnounce {
|
||||
// Established link: destructive handover is safe, and the peer is
|
||||
// close enough to become a courier for other carried mail.
|
||||
let authenticatedIngress = result.isDirectAnnounce
|
||||
&& canDeliverSecurely(to: result.peerID)
|
||||
&& isNoiseAuthenticatedIngressLink(for: packet, peerID: result.peerID)
|
||||
if authenticatedIngress {
|
||||
// The session was established on this still-bound ingress link.
|
||||
// A peer-level Noise session alone is not enough: it can outlive
|
||||
// its physical link while a replay rebinds an attacker's link to
|
||||
// the victim's ID.
|
||||
deliverCourierMail(to: result.peerID, noiseKey: noiseKey)
|
||||
sprayCourierMail(to: result.peerID, noiseKey: noiseKey, isVerifiedPeer: true)
|
||||
} else {
|
||||
// Relayed announce: recipient is multi-hop away. Push a copy
|
||||
// toward them speculatively; the carried copy stays put.
|
||||
// Relayed announce, or a direct-looking announce that has not yet
|
||||
// proved link ownership with Noise: push a speculative copy while
|
||||
// retaining the durable carried original.
|
||||
deliverCourierMailRemotely(to: result.peerID, noiseKey: noiseKey)
|
||||
if result.isDirectAnnounce,
|
||||
!hasCurrentNoiseAuthenticatedLink(to: result.peerID) {
|
||||
if noiseService.hasEstablishedSession(with: result.peerID) {
|
||||
// A session with no surviving authenticated link is stale;
|
||||
// force the current link to prove possession again.
|
||||
noiseService.clearSession(for: result.peerID)
|
||||
}
|
||||
if !noiseService.hasSession(with: result.peerID) {
|
||||
initiateNoiseHandshake(with: result.peerID)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -4559,6 +4947,10 @@ extension BLEService {
|
||||
}
|
||||
self.lastLinkRebindAt[linkUUID] = now
|
||||
|
||||
// A Noise proof belongs to the old physical binding. Never carry
|
||||
// it across an announce-driven rebind, whose direct TTL is
|
||||
// replayable; the new owner must complete a fresh handshake.
|
||||
self.noiseAuthenticatedLinkOwners.removeValue(forKey: link)
|
||||
switch link {
|
||||
case .peripheral(let peripheralUUID):
|
||||
self.linkStateStore.bindPeripheral(peripheralUUID, to: peerID)
|
||||
@@ -4652,6 +5044,10 @@ extension BLEService {
|
||||
)
|
||||
for uuid in retiring {
|
||||
guard let state = linkStateStore.state(forPeripheralID: uuid) else { continue }
|
||||
collectionsQueue.sync(flags: .barrier) {
|
||||
pendingPeripheralWrites.discardAll(for: uuid)
|
||||
}
|
||||
noiseAuthenticatedLinkOwners.removeValue(forKey: .peripheral(uuid))
|
||||
_ = linkStateStore.removePeripheral(uuid)
|
||||
SecureLogger.info(
|
||||
"🔗 Retiring redundant link \(uuid.prefix(8))… bound to \(peerID.id.prefix(8))…\(keptUUID.map { " (keeping \($0.prefix(8))…)" } ?? "")",
|
||||
@@ -5029,7 +5425,11 @@ extension BLEService {
|
||||
}
|
||||
|
||||
private func handleNoiseHandshake(_ packet: BitchatPacket, from peerID: PeerID) {
|
||||
let wasEstablished = noiseService.hasEstablishedSession(with: peerID)
|
||||
noisePacketHandler.handleHandshake(packet, from: peerID)
|
||||
if !wasEstablished, noiseService.hasEstablishedSession(with: peerID) {
|
||||
markNoiseAuthenticatedIngressLink(for: packet, peerID: peerID)
|
||||
}
|
||||
}
|
||||
|
||||
private func handleNoiseEncrypted(_ packet: BitchatPacket, from peerID: PeerID) {
|
||||
|
||||
Reference in New Issue
Block a user