Harden private media migration compatibility

This commit is contained in:
jack
2026-07-25 20:39:29 +02:00
committed by jack
parent 3f18768112
commit 86d8bb70f7
31 changed files with 2715 additions and 121 deletions
@@ -121,7 +121,7 @@ final class BLENoisePacketHandler {
let payloadType = decrypted[0]
let payloadData = decrypted.dropFirst()
guard let noisePayloadType = NoisePayloadType(rawValue: payloadType) else {
guard let noisePayloadType = NoisePayloadType.decoded(rawValue: payloadType) else {
SecureLogger.warning("⚠️ Unknown noise payload type: \(payloadType)")
return
}
@@ -53,6 +53,12 @@ struct BLENoiseSessionQueues {
return payloads
}
func containsTypedPayload(transferId: String) -> Bool {
typedPayloadsByPeerID.values.contains { payloads in
payloads.contains { $0.transferId == transferId }
}
}
@discardableResult
mutating func removeTypedPayload(transferId: String) -> Bool {
for peerID in Array(typedPayloadsByPeerID.keys) {
@@ -17,6 +17,9 @@ struct BLEOutboundFragmentPlan {
}
enum BLEOutboundFragmentPlanner {
/// Current Android receivers reject fragment sets above 256. Private
/// media v1 treats that deployed ceiling as a cross-platform contract.
static let privateMediaV1MaxFragments = 256
private static let minimumChunkSize = 64
private static let fragmentIDLength = 8
@@ -71,6 +74,10 @@ enum BLEOutboundFragmentPlanner {
)
}
static func isPrivateMediaV1Compatible(_ plan: BLEOutboundFragmentPlan) -> Bool {
plan.totalFragments <= privateMediaV1MaxFragments
}
private static func sizingPolicy(
for packet: BitchatPacket,
requestedMaxChunk: Int?,
+10 -2
View File
@@ -10,6 +10,9 @@ struct BLEPeerInfo: Equatable {
var isVerifiedNickname: Bool
var lastSeen: Date
var capabilities: PeerCapabilities = []
/// Distinguishes an old client that omitted the capabilities TLV from a
/// modern client that explicitly advertised a set without a given bit.
var capabilitiesWereExplicitlyAdvertised: Bool = false
/// Rendezvous cell from the peer's announce when it advertises `.bridge`.
var bridgeGeohash: String?
}
@@ -114,6 +117,10 @@ struct BLEPeerRegistry {
peers[peerID.toShort()]?.capabilities ?? []
}
func capabilitiesWereExplicitlyAdvertised(for peerID: PeerID) -> Bool {
peers[peerID.toShort()]?.capabilitiesWereExplicitlyAdvertised == true
}
/// Peers whose last verified announce advertised the given capability.
func peers(advertising capability: PeerCapabilities) -> [PeerID] {
peers.values.filter { $0.capabilities.contains(capability) }.map(\.peerID)
@@ -181,7 +188,7 @@ struct BLEPeerRegistry {
signingPublicKey: Data?,
isConnected: Bool,
now: Date,
capabilities: PeerCapabilities = [],
capabilities: PeerCapabilities? = nil,
bridgeGeohash: String? = nil
) -> BLEPeerAnnounceUpdate {
let existing = peers[peerID]
@@ -199,7 +206,8 @@ struct BLEPeerRegistry {
signingPublicKey: signingPublicKey,
isVerifiedNickname: true,
lastSeen: now,
capabilities: capabilities,
capabilities: capabilities ?? [],
capabilitiesWereExplicitlyAdvertised: capabilities != nil,
bridgeGeohash: bridgeGeohash
)
+752 -39
View File
@@ -7,6 +7,175 @@ import Combine
import UIKit
#endif
/// Linearizes app-private-media admission against cancellation before work is
/// handed to the fragment scheduler. A transfer starts here synchronously,
/// before its `messageQueue` work item is enqueued; cancel/delete can therefore
/// leave a tombstone that the deferred work must observe.
///
/// Active admissions and cancellation tombstones have independent count
/// bounds. Tombstones may age out or evict older tombstones; active entries
/// are never evicted under pressure. A one-hour active timeout is reported as
/// an explicit transfer failure and removes any handshake-queued payload.
private final class BLEPrivateMediaTransferAdmissionRegistry {
enum BeginResult: Equatable {
case admitted
case alreadyKnown
case capacityExhausted
}
private enum State: Equatable {
case active
case cancelled
}
private struct Entry {
var state: State
var updatedAt: Date
}
private let lock = NSLock()
private let maxActiveEntries = 512
private let maxCancelledTombstones = 512
private let lifetime: TimeInterval = 60 * 60
private let onActiveExpired: (String) -> Void
private var entries: [String: Entry] = [:]
init(onActiveExpired: @escaping (String) -> Void) {
self.onActiveExpired = onActiveExpired
}
func begin(_ transferId: String, now: Date = Date()) -> BeginResult {
guard !transferId.isEmpty else { return .alreadyKnown }
lock.lock()
let expiredActive = pruneLocked(now: now)
// Transfer IDs are invocation-unique. Never revive a cancellation or
// admit a duplicate invocation that reused an in-flight identifier.
let result: BeginResult
if entries[transferId] != nil {
result = .alreadyKnown
} else if activeCountLocked >= maxActiveEntries {
// Never evict an admitted transfer: doing so strands its UI
// placeholder with no completion event. Reject the newcomer and
// let the caller surface the bounded-pressure failure instead.
result = .capacityExhausted
} else {
entries[transferId] = Entry(state: .active, updatedAt: now)
result = .admitted
}
lock.unlock()
notifyExpired(expiredActive)
return result
}
func cancel(_ transferId: String, now: Date = Date()) {
guard !transferId.isEmpty else { return }
lock.lock()
// Cancel the requested active entry before expiry pruning so a user
// cancellation wins over a simultaneous timeout notification.
entries[transferId] = Entry(state: .cancelled, updatedAt: now)
let expiredActive = pruneLocked(now: now)
trimCancelledTombstonesLocked()
lock.unlock()
notifyExpired(expiredActive)
}
func isActive(_ transferId: String, now: Date = Date()) -> Bool {
lock.lock()
let expiredActive = pruneLocked(now: now)
let active = entries[transferId]?.state == .active
if active {
entries[transferId]?.updatedAt = now
}
lock.unlock()
notifyExpired(expiredActive)
return active
}
/// Runs `body` while holding the admission lock. Callers use this at the
/// collections-queue append/submit boundary so cancellation and admission
/// have one deterministic order: whichever acquires this lock first wins.
func withActive<Result>(
_ transferId: String,
now: Date = Date(),
_ body: () -> Result
) -> Result? {
lock.lock()
let expiredActive = pruneLocked(now: now)
guard entries[transferId]?.state == .active else {
lock.unlock()
notifyExpired(expiredActive)
return nil
}
entries[transferId]?.updatedAt = now
let result = body()
lock.unlock()
notifyExpired(expiredActive)
return result
}
func finish(_ transferId: String) {
lock.lock()
entries.removeValue(forKey: transferId)
lock.unlock()
}
var count: Int {
lock.lock()
let expiredActive = pruneLocked(now: Date())
let result = entries.count
lock.unlock()
notifyExpired(expiredActive)
return result
}
func prune(now: Date = Date()) {
lock.lock()
let expiredActive = pruneLocked(now: now)
lock.unlock()
notifyExpired(expiredActive)
}
private var activeCountLocked: Int {
entries.values.reduce(into: 0) { count, entry in
if entry.state == .active { count += 1 }
}
}
/// Removes stale tombstones silently and stale active admissions with a
/// caller-visible timeout notification. Must be called with `lock` held;
/// notifications are delivered only after the lock is released.
private func pruneLocked(now: Date) -> [String] {
var expiredActive: [String] = []
let expiredEntries = entries.filter {
now.timeIntervalSince($0.value.updatedAt) > lifetime
}
for (transferId, entry) in expiredEntries {
if entry.state == .active {
expiredActive.append(transferId)
}
entries.removeValue(forKey: transferId)
}
trimCancelledTombstonesLocked()
return expiredActive
}
private func trimCancelledTombstonesLocked() {
let cancelled = entries
.filter { $0.value.state == .cancelled }
.sorted { $0.value.updatedAt < $1.value.updatedAt }
let overflow = max(0, cancelled.count - maxCancelledTombstones)
for victim in cancelled.prefix(overflow) {
entries.removeValue(forKey: victim.key)
}
}
private func notifyExpired(_ transferIds: [String]) {
for transferId in transferIds {
onActiveExpired(transferId)
}
}
}
/// BLEService Bluetooth Mesh Transport
/// - Emits events exclusively via `BitchatDelegate` for UI.
/// - ChatViewModel must consume delegate callbacks (`didReceivePublicMessage`, `didReceiveNoisePayload`).
@@ -108,6 +277,10 @@ final class BLEService: NSObject {
/// before it hands a packet to `messageQueue`.
var _test_beforeReceivePacketHandoff: (() -> Void)?
var _test_onReceivePacketHandoff: (() -> Void)?
var _test_onPrivateMediaSessionReconciled: ((PeerID) -> Void)?
/// May block in tests to hold the serial message queue immediately before
/// the deferred private-media admission check.
var _test_beforePrivateMediaDeferredSend: ((String) -> Void)?
#endif
private var selfBroadcastTracker = BLESelfBroadcastTracker()
private let meshTopology = MeshTopologyTracker()
@@ -136,6 +309,9 @@ final class BLEService: NSObject {
// 5. Fragment Reassembly (necessary for messages > MTU)
private var fragmentAssemblyBuffer = BLEFragmentAssemblyBuffer()
private var outboundFragmentTransfers = BLEOutboundFragmentTransferScheduler()
private lazy var privateMediaTransferAdmissions = BLEPrivateMediaTransferAdmissionRegistry { [weak self] transferId in
self?.handlePrivateMediaAdmissionExpiry(transferId)
}
private let incomingFileStore: BLEIncomingFileStore
// Simple announce throttling
@@ -883,6 +1059,69 @@ final class BLEService: NSObject {
collectionsQueue.sync { peerRegistry.capabilities(for: peerID) }
}
func privateMediaSendPolicy(to peerID: PeerID) -> PrivateMediaSendPolicy {
let state: (
capabilities: PeerCapabilities,
fingerprint: String?
) = collectionsQueue.sync {
let info = peerRegistry.info(for: peerID.toShort())
return (
info?.capabilities ?? [],
info?.noisePublicKey?.sha256Fingerprint()
)
}
if state.capabilities.contains(.privateMedia) {
return .encrypted
}
guard let fingerprint = state.fingerprint else {
// A raw fallback must be bound to the stable Noise key from a
// verified registry entry; a routing ID alone can rotate or be
// spoofed. Session authentication is required for pinning, below,
// but old clients may need the consented fallback before a session.
return .blockedDowngrade
}
// During the mixed-version migration, an unpinned peer is legacy
// eligible whether the capabilities TLV was absent or explicitly
// omitted this bit. Only a capability observation bound to a matching
// authenticated Noise session may make the no-bit state a downgrade.
return identityManager.hasObservedPrivateMediaCapability(fingerprint: fingerprint)
? .blockedDowngrade
: .legacyRequiresConsent
}
/// Pins private-media support only when the capability-bearing registry
/// identity and a completed Noise session authenticate the same static
/// key. A signed announce alone does not prove possession of that Noise
/// key and must never create a downgrade pin.
private func reconcilePrivateMediaCapabilityPin(
for peerID: PeerID,
authenticatedFingerprint: String? = nil
) {
let normalizedPeerID = peerID.toShort()
let advertisedFingerprint: String? = collectionsQueue.sync {
guard let info = peerRegistry.info(for: normalizedPeerID),
info.capabilities.contains(.privateMedia) else { return nil }
return info.noisePublicKey?.sha256Fingerprint()
}
guard let advertisedFingerprint else { return }
let sessionFingerprint = authenticatedFingerprint
?? noiseService.getPeerFingerprint(normalizedPeerID)
guard let sessionFingerprint,
sessionFingerprint.caseInsensitiveCompare(advertisedFingerprint) == .orderedSame else {
if authenticatedFingerprint != nil {
SecureLogger.warning(
"Refusing private-media capability pin for \(normalizedPeerID.id.prefix(8))…: authenticated Noise key does not match verified announce",
category: .security
)
}
return
}
identityManager.markPrivateMediaCapable(fingerprint: advertisedFingerprint)
}
/// Enables or disables a runtime-advertised capability bit (e.g. the
/// internet-gateway toggle) and re-announces so peers learn promptly.
/// Build-time bits stay in `PeerCapabilities.localSupported`.
@@ -1010,7 +1249,28 @@ final class BLEService: NSObject {
// MARK: Messaging
private func handlePrivateMediaAdmissionExpiry(_ transferId: String) {
// Expiry can be discovered from the BLE maintenance queue or while a
// caller already owns collectionsQueue. Cleanup is therefore
// fire-and-forget; never synchronously re-enter the collections lock.
collectionsQueue.async(flags: .barrier) { [weak self] in
_ = self?.pendingNoiseSessionQueues.removeTypedPayload(transferId: transferId)
}
TransferProgressManager.shared.rejectBeforeStart(
id: transferId,
reason: String(
localized: "content.delivery.reason.private_media_admission_expired",
defaultValue: "Media transfer timed out before it could start",
comment: "Failure reason when private-media admission expires before fragment scheduling"
)
)
}
func cancelTransfer(_ transferId: String) {
// Cancellation must become visible synchronously. Scheduler/pending-
// Noise cleanup remains asynchronous, but deferred private-media work
// cannot pass another admission boundary after this returns.
privateMediaTransferAdmissions.cancel(transferId)
collectionsQueue.async(flags: .barrier) { [weak self] in
guard let self = self else { return }
@@ -1087,50 +1347,236 @@ final class BLEService: NSObject {
}
func sendFilePrivate(_ filePacket: BitchatFilePacket, to peerID: PeerID, transferId: String) {
sendFilePrivate(
filePacket,
to: peerID,
transferId: transferId,
allowLegacyFallback: false
)
}
func sendFilePrivate(
_ filePacket: BitchatFilePacket,
to peerID: PeerID,
transferId: String,
allowLegacyFallback: Bool
) {
// Register before enqueueing onto messageQueue. This closes the window
// where cancel/delete could run first, observe no scheduler state, and
// then be followed by a deferred clear-media send.
switch privateMediaTransferAdmissions.begin(transferId) {
case .admitted:
break
case .alreadyKnown:
SecureLogger.debug(
"Private media admission already cancelled or duplicated for \(transferId.prefix(8))",
category: .security
)
return
case .capacityExhausted:
SecureLogger.warning(
"Private media admission capacity exhausted for \(transferId.prefix(8))",
category: .security
)
TransferProgressManager.shared.rejectBeforeStart(
id: transferId,
reason: String(
localized: "content.delivery.reason.private_media_admission_full",
defaultValue: "Too many media transfers are waiting; try again shortly",
comment: "Failure reason when too many private-media transfers are awaiting admission"
)
)
return
}
messageQueue.async { [weak self] in
guard let self = self else { return }
guard !self.isPanicSuspended else { return }
let targetID = peerID.toShort()
let supportsPrivateMedia = self.collectionsQueue.sync {
self.peerRegistry.capabilities(for: targetID).contains(.privateMedia)
#if DEBUG
self._test_beforePrivateMediaDeferredSend?(transferId)
#endif
guard !self.isPanicSuspended else {
self.privateMediaTransferAdmissions.finish(transferId)
return
}
guard supportsPrivateMedia else {
guard self.privateMediaTransferAdmissions.isActive(transferId) else {
self.privateMediaTransferAdmissions.finish(transferId)
return
}
let targetID = peerID.toShort()
switch self.privateMediaSendPolicy(to: targetID) {
case .encrypted:
break
case .legacyRequiresConsent:
guard allowLegacyFallback else {
SecureLogger.warning(
"Private media blocked pending explicit legacy-clear consent for \(targetID.id.prefix(8))",
category: .security
)
TransferProgressManager.shared.rejectBeforeStart(
id: transferId,
reason: String(
localized: "content.delivery.reason.legacy_media_consent_required",
defaultValue: "Confirmation required before sending without end-to-end encryption",
comment: "Failure reason when a legacy private-media send lacks per-send consent"
)
)
self.privateMediaTransferAdmissions.finish(transferId)
return
}
// Migration path accepted by current Android and used by older
// iOS releases: preserve the directed raw file-transfer wire
// shape, but require the signature the receive path verifies.
// The allow flag belongs to this invocation only and is
// consumed here; a retry must obtain fresh user consent.
self.sendSignedLegacyPrivateFile(
filePacket,
to: targetID,
transferId: transferId
)
return
case .blockedDowngrade:
SecureLogger.warning(
"Private media not sent: \(targetID.id.prefix(8)) did not advertise encrypted-media support",
"Private media downgrade blocked for \(targetID.id.prefix(8))",
category: .security
)
TransferProgressManager.shared.rejectBeforeStart(id: transferId)
TransferProgressManager.shared.rejectBeforeStart(
id: transferId,
reason: String(
localized: "content.delivery.reason.private_media_downgrade_blocked",
defaultValue: "Encrypted media required; ask this contact to upgrade",
comment: "Failure reason when a peer that previously supported encrypted media appears to downgrade"
)
)
self.privateMediaTransferAdmissions.finish(transferId)
return
}
guard let typedPayload = BLENoisePayloadFactory.privateFile(filePacket) else {
SecureLogger.error("❌ Failed to encode file packet for private send", category: .session)
TransferProgressManager.shared.rejectBeforeStart(id: transferId)
TransferProgressManager.shared.rejectBeforeStart(
id: transferId,
reason: String(localized: "content.delivery.reason.media_encoding_failed", defaultValue: "Failed to prepare media", comment: "Failure reason when private media cannot be encoded")
)
self.privateMediaTransferAdmissions.finish(transferId)
return
}
guard self.noiseService.hasEstablishedSession(with: targetID) else {
self.collectionsQueue.sync(flags: .barrier) {
self.pendingNoiseSessionQueues.appendTypedPayload(
typedPayload,
transferId: transferId,
for: targetID
)
let queued = self.collectionsQueue.sync(flags: .barrier) {
self.privateMediaTransferAdmissions.withActive(transferId) {
self.pendingNoiseSessionQueues.appendTypedPayload(
typedPayload,
transferId: transferId,
for: targetID
)
return true
} ?? false
}
guard queued else {
self.privateMediaTransferAdmissions.finish(transferId)
return
}
SecureLogger.debug("📥 Queued private file for \(targetID.id.prefix(8))… pending handshake", category: .session)
guard self.privateMediaTransferAdmissions.isActive(transferId) else {
self.collectionsQueue.sync(flags: .barrier) {
_ = self.pendingNoiseSessionQueues.removeTypedPayload(transferId: transferId)
}
self.privateMediaTransferAdmissions.finish(transferId)
return
}
self.initiateNoiseHandshake(with: targetID)
return
}
do {
guard self.privateMediaTransferAdmissions.isActive(transferId) else {
self.privateMediaTransferAdmissions.finish(transferId)
return
}
let packet = try self.makeEncryptedNoisePacket(typedPayload, to: targetID)
guard self.privateMediaTransferAdmissions.isActive(transferId) else {
self.privateMediaTransferAdmissions.finish(transferId)
return
}
SecureLogger.debug("📁 Sending encrypted private file to \(targetID.id.prefix(8))… plaintextBytes=\(typedPayload.count)", category: .session)
self.broadcastPacket(packet, transferId: transferId)
self.broadcastPacket(
packet,
transferId: transferId,
requiresPrivateMediaAdmission: true
)
} catch {
SecureLogger.error("❌ Failed to encrypt private file for \(targetID.id.prefix(8))…: \(error)", category: .security)
TransferProgressManager.shared.rejectBeforeStart(id: transferId)
TransferProgressManager.shared.rejectBeforeStart(
id: transferId,
reason: String(localized: "content.delivery.reason.encryption_failed", comment: "Failure reason shown when a message could not be encrypted for the peer")
)
self.privateMediaTransferAdmissions.finish(transferId)
}
}
}
/// Compatibility-only fallback for peers that have not advertised
/// encrypted private media. The payload is authenticated but visible to
/// relays, matching the pre-migration behavior until those clients upgrade.
private func sendSignedLegacyPrivateFile(
_ filePacket: BitchatFilePacket,
to targetID: PeerID,
transferId: String
) {
guard privateMediaTransferAdmissions.isActive(transferId) else {
privateMediaTransferAdmissions.finish(transferId)
return
}
guard let payload = filePacket.encode(),
let recipientData = Data(hexString: targetID.id) else {
SecureLogger.error("❌ Failed to encode legacy private file transfer", category: .session)
TransferProgressManager.shared.rejectBeforeStart(
id: transferId,
reason: String(localized: "content.delivery.reason.media_encoding_failed", defaultValue: "Failed to prepare media", comment: "Failure reason when private media cannot be encoded")
)
privateMediaTransferAdmissions.finish(transferId)
return
}
let unsigned = BitchatPacket(
type: MessageType.fileTransfer.rawValue,
senderID: myPeerIDData,
recipientID: recipientData,
timestamp: UInt64(Date().timeIntervalSince1970 * 1000),
payload: payload,
signature: nil,
ttl: messageTTL,
version: 2
)
guard let signed = noiseService.signPacket(unsigned) else {
SecureLogger.error("❌ Failed to sign legacy private file transfer", category: .security)
TransferProgressManager.shared.rejectBeforeStart(
id: transferId,
reason: String(localized: "content.delivery.reason.media_signing_failed", defaultValue: "Failed to authenticate media", comment: "Failure reason when a legacy private-media packet cannot be signed")
)
privateMediaTransferAdmissions.finish(transferId)
return
}
// Signing can be non-trivial; cancellation that won while it ran must
// still prevent the clear payload from reaching the broadcast path.
guard privateMediaTransferAdmissions.isActive(transferId) else {
privateMediaTransferAdmissions.finish(transferId)
return
}
SecureLogger.warning(
"📁 Sending signed legacy private file to \(targetID.id.prefix(8))…; peer has not advertised E2E media",
category: .security
)
broadcastPacket(
signed,
transferId: transferId,
requiresPrivateMediaAdmission: true
)
}
func sendReadReceipt(_ receipt: ReadReceipt, to peerID: PeerID) {
// Hop like sendMessage: callers are often on the main actor, and the
@@ -1251,8 +1697,26 @@ final class BLEService: NSObject {
// MARK: - Packet Broadcasting
private func broadcastPacket(_ packet: BitchatPacket, transferId: String? = nil) {
guard !isPanicSuspended else { return }
private func broadcastPacket(
_ packet: BitchatPacket,
transferId: String? = nil,
requiresPrivateMediaAdmission: Bool = false
) {
guard !isPanicSuspended else {
if requiresPrivateMediaAdmission, let transferId {
privateMediaTransferAdmissions.finish(transferId)
}
return
}
if requiresPrivateMediaAdmission {
guard let transferId,
privateMediaTransferAdmissions.isActive(transferId) else {
if let transferId {
privateMediaTransferAdmissions.finish(transferId)
}
return
}
}
// Apply route if recipient exists (centralized route application)
let packetToSend: BitchatPacket
if let recipientPeerID = PeerID(hexData: packet.recipientID) {
@@ -1261,14 +1725,75 @@ final class BLEService: NSObject {
packetToSend = packet
}
// Encode once using a small per-type padding policy, then delegate by type
let padForBLE = BLEOutboundPacketPolicy.padsBLEFrame(for: packetToSend.type)
// Cross-platform private-media v1 is bounded by Android's deployed
// 256-fragment receive cap. Run the same planner the scheduler will
// use, after route application, for both encrypted and consented raw
// migration sends. Reject before reserving a transfer slot or writing
// any fragment; public media is intentionally unaffected.
if let transferId,
let recipientPeerID = PeerID(hexData: packetToSend.recipientID),
packetToSend.type == MessageType.noiseEncrypted.rawValue
|| packetToSend.type == MessageType.fileTransfer.rawValue {
let compatibilityRequest = BLEOutboundFragmentTransferRequest(
packet: packetToSend,
pad: padForBLE,
maxChunk: nil,
directedPeer: recipientPeerID,
transferId: transferId
)
guard let plan = BLEOutboundFragmentPlanner.makePlan(
for: compatibilityRequest,
defaultChunkSize: defaultFragmentSize,
bleMaxMTU: bleMaxMTU
), BLEOutboundFragmentPlanner.isPrivateMediaV1Compatible(plan) else {
SecureLogger.warning(
"Private media rejected: exceeds cross-platform 256-fragment limit",
category: .security
)
TransferProgressManager.shared.rejectBeforeStart(
id: transferId,
reason: String(
localized: "content.delivery.reason.private_media_too_many_fragments",
defaultValue: "File is too large for this contact's client (more than 256 mesh fragments)",
comment: "Failure reason when private media exceeds the Android-compatible fragment limit"
)
)
if requiresPrivateMediaAdmission {
privateMediaTransferAdmissions.finish(transferId)
}
return
}
}
// Route planning and fragment preflight can take enough time for a
// user cancellation to win. Recheck before exposing even the test tap,
// then check atomically with scheduler admission below.
if requiresPrivateMediaAdmission {
guard let transferId,
privateMediaTransferAdmissions.isActive(transferId) else {
if let transferId {
privateMediaTransferAdmissions.finish(transferId)
}
return
}
}
#if DEBUG
_test_onOutboundPacket?(packetToSend)
#endif
// Encode once using a small per-type padding policy, then delegate by type
let padForBLE = BLEOutboundPacketPolicy.padsBLEFrame(for: packetToSend.type)
if packetToSend.type == MessageType.fileTransfer.rawValue {
sendFragmentedPacket(packetToSend, pad: padForBLE, maxChunk: nil, directedOnlyPeer: nil, transferId: transferId)
sendFragmentedPacket(
packetToSend,
pad: padForBLE,
maxChunk: nil,
directedOnlyPeer: nil,
transferId: transferId,
requiresPrivateMediaAdmission: requiresPrivateMediaAdmission
)
return
}
// App-initiated private media is already one opaque Noise ciphertext.
@@ -1282,7 +1807,18 @@ final class BLEService: NSObject {
pad: padForBLE,
maxChunk: nil,
directedOnlyPeer: recipientPeerID,
transferId: transferId
transferId: transferId,
requiresPrivateMediaAdmission: requiresPrivateMediaAdmission
)
return
}
if requiresPrivateMediaAdmission {
if let transferId {
privateMediaTransferAdmissions.finish(transferId)
}
SecureLogger.error(
"Private media admission reached an unsupported non-directed packet shape",
category: .security
)
return
}
@@ -2631,18 +3167,20 @@ extension BLEService {
func _test_seedConnectedPeer(
_ peerID: PeerID,
nickname: String,
capabilities: PeerCapabilities = []
capabilities: PeerCapabilities? = nil,
noisePublicKey: Data? = nil
) {
collectionsQueue.sync(flags: .barrier) {
peerRegistry.upsert(BLEPeerInfo(
peerID: peerID,
nickname: nickname,
isConnected: true,
noisePublicKey: nil,
noisePublicKey: noisePublicKey,
signingPublicKey: nil,
isVerifiedNickname: true,
lastSeen: Date(),
capabilities: capabilities
capabilities: capabilities ?? [],
capabilitiesWereExplicitlyAdvertised: capabilities != nil
))
}
}
@@ -2657,6 +3195,90 @@ extension BLEService {
try noiseService.processHandshakeMessage(from: peerID, message: message)
}
func _test_enqueuePendingNoisePayload(
_ payload: Data,
transferId: String,
for peerID: PeerID
) {
guard privateMediaTransferAdmissions.begin(transferId) == .admitted else { return }
collectionsQueue.sync(flags: .barrier) {
pendingNoiseSessionQueues.appendTypedPayload(
payload,
transferId: transferId,
for: peerID
)
}
}
func _test_sendPendingNoisePayloadsAfterHandshake(for peerID: PeerID) {
sendPendingNoisePayloadsAfterHandshake(for: peerID)
}
func _test_privateMediaTransferState(
transferId: String
) -> (admissionActive: Bool, pendingNoise: Bool, activeScheduler: Int, pendingScheduler: Int) {
let scheduler = collectionsQueue.sync {
(
pendingNoiseSessionQueues.containsTypedPayload(transferId: transferId),
outboundFragmentTransfers.activeCount,
outboundFragmentTransfers.pendingCount
)
}
return (
privateMediaTransferAdmissions.isActive(transferId),
scheduler.0,
scheduler.1,
scheduler.2
)
}
func _test_privateMediaAdmissionEntryCount() -> Int {
privateMediaTransferAdmissions.count
}
@discardableResult
func _test_beginPrivateMediaAdmission(_ transferId: String, now: Date) -> Bool {
privateMediaTransferAdmissions.begin(transferId, now: now) == .admitted
}
func _test_isPrivateMediaAdmissionActive(_ transferId: String, now: Date) -> Bool {
privateMediaTransferAdmissions.isActive(transferId, now: now)
}
func _test_finishPrivateMediaAdmission(_ transferId: String) {
privateMediaTransferAdmissions.finish(transferId)
}
func _test_drainPrivateMediaSendPipeline() async {
let collectionsQueue = self.collectionsQueue
await withCheckedContinuation { continuation in
messageQueue.async {
collectionsQueue.async(flags: .barrier) {
continuation.resume()
}
}
}
}
func _test_broadcastPrivateMediaPacket(
_ packet: BitchatPacket,
transferId: String
) {
broadcastPacket(
packet,
transferId: transferId,
requiresPrivateMediaAdmission: true
)
}
/// Builds an authenticated-session packet from an exact typed plaintext.
/// Compatibility tests use this to model Android's deployed 0x20 file
/// payload and the short-lived 0x09 prerelease payload without exposing a
/// production API that can emit the old value.
func _test_makeEncryptedNoisePacket(_ typedPayload: Data, to peerID: PeerID) throws -> BitchatPacket {
try makeEncryptedNoisePacket(typedPayload, to: peerID)
}
static func _test_shouldRediscoverBitChatService(
invalidatedServiceUUIDs: [CBUUID],
cachedServiceUUIDs: [CBUUID]?
@@ -3698,6 +4320,13 @@ extension BLEService {
service.onPeerAuthenticated = { [weak self] peerID, fingerprint in
SecureLogger.debug("🔐 Noise session authenticated with \(peerID.id.prefix(8))…, fingerprint: \(fingerprint.prefix(16))")
self?.messageQueue.async { [weak self] in
self?.reconcilePrivateMediaCapabilityPin(
for: peerID,
authenticatedFingerprint: fingerprint
)
#if DEBUG
self?._test_onPrivateMediaSessionReconciled?(peerID)
#endif
self?.sendPendingMessagesAfterHandshake(for: peerID)
self?.sendPendingNoisePayloadsAfterHandshake(for: peerID)
}
@@ -3759,7 +4388,7 @@ extension BLEService {
private func makeEncryptedNoisePacket(_ typedPayload: Data, to peerID: PeerID) throws -> BitchatPacket {
let encrypted: Data
let isPrivateFile = typedPayload.first == NoisePayloadType.privateFile.rawValue
let isPrivateFile = NoisePayloadType.isPrivateFile(rawValue: typedPayload.first)
if isPrivateFile {
encrypted = try noiseService.encryptPrivateFilePayload(typedPayload, for: peerID)
} else {
@@ -4807,7 +5436,8 @@ extension BLEService {
directedOnlyPeer: PeerID? = nil,
transferId: String? = nil,
requireDirectPeerLink: Bool = false,
requireNoiseAuthenticatedPeerLink: Bool = false
requireNoiseAuthenticatedPeerLink: Bool = false,
requiresPrivateMediaAdmission: Bool = false
) -> Bool {
let request = BLEOutboundFragmentTransferRequest(
packet: packet,
@@ -4819,8 +5449,34 @@ extension BLEService {
requireNoiseAuthenticatedPeerLink: requireNoiseAuthenticatedPeerLink
)
let result = collectionsQueue.sync(flags: .barrier) {
outboundFragmentTransfers.submit(request, maxConcurrentTransfers: TransportConfig.bleMaxConcurrentTransfers)
let result: BLEOutboundFragmentTransferScheduler.SubmitResult? = collectionsQueue.sync(flags: .barrier) {
if requiresPrivateMediaAdmission {
guard let transferId else { return nil }
// This lock is taken while the scheduler is already protected
// by collectionsQueue. Cancellation takes the admission lock
// synchronously but never waits on collectionsQueue, avoiding
// lock inversion while giving submit/cancel one linear order.
return privateMediaTransferAdmissions.withActive(transferId) {
outboundFragmentTransfers.submit(
request,
maxConcurrentTransfers: TransportConfig.bleMaxConcurrentTransfers
)
}
}
return outboundFragmentTransfers.submit(
request,
maxConcurrentTransfers: TransportConfig.bleMaxConcurrentTransfers
)
}
guard let result else {
if let transferId, requiresPrivateMediaAdmission {
privateMediaTransferAdmissions.finish(transferId)
}
return false
}
if let transferId, requiresPrivateMediaAdmission {
// The scheduler now owns normal cancellation (active or pending).
privateMediaTransferAdmissions.finish(transferId)
}
return handleFragmentTransferSubmitResult(result)
}
@@ -4896,14 +5552,24 @@ extension BLEService {
}
}
let transferIdentifier: String? = {
guard let id = reservedTransferId else { return nil }
collectionsQueue.sync(flags: .barrier) {
_ = self.outboundFragmentTransfers.activateReservedTransfer(id: id, totalFragments: plan.totalFragments, workItems: [])
let transferIdentifier: String?
if let id = reservedTransferId {
let activated = collectionsQueue.sync(flags: .barrier) {
self.outboundFragmentTransfers.activateReservedTransfer(
id: id,
totalFragments: plan.totalFragments,
workItems: []
)
}
// Cancellation may remove the reservation between submit and plan
// construction. Treat that as cancellation, not as permission to
// schedule an untracked fragment train.
guard activated else { return false }
TransferProgressManager.shared.start(id: id, totalFragments: plan.totalFragments)
return id
}()
transferIdentifier = id
} else {
transferIdentifier = nil
}
let sendFragment: (BitchatPacket) -> Bool = { [weak self] fragmentPacket in
guard let self else { return false }
@@ -5236,6 +5902,16 @@ extension BLEService {
private func handleAnnounce(_ packet: BitchatPacket, from peerID: PeerID) {
let result = announceHandler.handle(packet, from: peerID)
// The verified announce and Noise handshake can arrive in either
// order. If authentication already completed, compare its static-key
// fingerprint with the newly upserted registry key before pinning;
// otherwise the session callback performs this reconciliation later.
if let result,
result.isVerified,
result.announcement.capabilities?.contains(.privateMedia) == true {
reconcilePrivateMediaCapabilityPin(for: result.peerID)
}
// A verified announce is the moment a signing key becomes bound to this
// owner's noise key: retry any prekey bundle that raced ahead of it.
if let result, result.isVerified {
@@ -5546,7 +6222,7 @@ extension BLEService {
signingPublicKey: announcement.signingPublicKey,
isConnected: isConnected,
now: now,
capabilities: announcement.capabilities ?? [],
capabilities: announcement.capabilities,
bridgeGeohash: announcement.bridgeGeohash
) ?? BLEPeerAnnounceUpdate(isNewPeer: false, wasDisconnected: false, previousNickname: nil)
},
@@ -5897,13 +6573,45 @@ extension BLEService {
guard !payloads.isEmpty else { return }
SecureLogger.debug("📤 Sending \(payloads.count) pending noise payloads to \(peerID.id.prefix(8))… after handshake", category: .session)
for pending in payloads {
let privateMediaTransferId: String? = {
guard NoisePayloadType.isPrivateFile(rawValue: pending.payload.first) else { return nil }
return pending.transferId
}()
if let transferId = privateMediaTransferId,
!privateMediaTransferAdmissions.isActive(transferId) {
privateMediaTransferAdmissions.finish(transferId)
continue
}
do {
if let transferId = privateMediaTransferId,
!privateMediaTransferAdmissions.isActive(transferId) {
privateMediaTransferAdmissions.finish(transferId)
continue
}
let packet = try makeEncryptedNoisePacket(pending.payload, to: peerID)
if let transferId = privateMediaTransferId,
!privateMediaTransferAdmissions.isActive(transferId) {
privateMediaTransferAdmissions.finish(transferId)
continue
}
broadcastPacket(
try makeEncryptedNoisePacket(pending.payload, to: peerID),
transferId: pending.transferId
packet,
transferId: pending.transferId,
requiresPrivateMediaAdmission: privateMediaTransferId != nil
)
} catch {
SecureLogger.error("❌ Failed to send pending noise payload to \(peerID.id.prefix(8))…: \(error)")
if let transferId = pending.transferId {
TransferProgressManager.shared.rejectBeforeStart(
id: transferId,
reason: String(
localized: "content.delivery.reason.encryption_failed",
defaultValue: "Failed to encrypt media",
comment: "Failure reason shown when queued private media cannot be encrypted after handshake"
)
)
privateMediaTransferAdmissions.finish(transferId)
}
}
}
}
@@ -6064,6 +6772,11 @@ extension BLEService {
private func performCleanup() {
let now = Date()
// Admission expiry is a visible transfer failure, never a silent
// eviction. The registry delivers notifications after releasing its
// lock, so this maintenance pass cannot deadlock a concurrent cancel.
privateMediaTransferAdmissions.prune(now: now)
// Clean old processed messages efficiently
messageDeduplicator.cleanup()