Fix restore-path main↔bleQueue deadlock and courier drop amplification (#1425)

* Fix restore-path main↔bleQueue deadlock and courier drop amplification

A device froze permanently in a two-phone test. Debugger stacks showed an
ABBA deadlock: the main actor was in bleQueue.sync (delivery-ack send →
broadcastPacket → readLinkState) while bleQueue was in main.sync
(captureBluetoothStatus reading backgroundTimeRemaining). The load that
lined the two edges up came from a courier-drop amplification storm:
drop dedup was in-memory only while the outbox driving 120s re-deposits
is persisted, so every relaunch republished the same undelivered DM as a
fresh 24h relay drop and every gateway relaunch re-fetched the whole
backlog — ~20 copies of one DM delivered in 40ms, each triggering
decrypt + delivery + ack + handshake work.

Fixes, in rank order:
- Edge B (P0): captureBluetoothStatus no longer main.syncs from bleQueue;
  backgroundTimeRemaining is sampled on main and cached behind a lock.
  Invariant documented: bleQueue must NEVER sync-dispatch to main.
- Edge A (P0, defense in depth): sendDeliveryAck / sendReadReceipt /
  sendPrivateMessage / sendNoisePayload / triggerHandshake hop to
  messageQueue like sendMessage, so no main-actor call path reaches
  readLinkState's bleQueue.sync.
- Drop dedup (P1): publishedDropKeys and seenDropEventIDs persist across
  relaunches (new BridgeDropDedupStore, entries expire with the 24h
  NIP-40 drop window; wiped on panic) — one drop per message ID per 24h
  regardless of relaunch count.
- Receiver dedup (P1): openCourierEnvelope dedups on the inner private
  message ID before delivery, so a duplicate copy costs one decrypt and
  never re-delivers, re-acks, or re-triggers a handshake.
- Handshake gating (P2): queued acks initiate a Noise handshake only for
  reachable peers; mail from absent/rotated identities no longer turns
  each copy into a mesh-wide handshake flood (the ack stays queued and
  flushes when a session eventually establishes).
- Outbox (P3): re-enqueueing a queued message ID carries over its
  depositedCourierKeys so resends stop re-burning the same courier slots.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

* Review fixes: offline-drop durability, gateway handoff retry, restore-log freshness, coalesced persist

Adversarial review of the storm/deadlock PR surfaced four issues:

- Offline blackhole (must-fix): a deposit made while relays were down
  persisted its dedup key even though the drop only sat in the in-memory
  pending queue — app killed before reconnect meant the relaunch lost the
  drop but the persisted key blocked every re-deposit for 24h. The
  persisted snapshot now excludes keys still pending; they become durable
  only when flushPendingDrops actually publishes them.
- Gateway handoff: seen-event IDs were consumed before the deliverToPeer
  handoff; a failed handoff (peer walked away) permanently dropped the
  event for a single-gateway island. deliverToPeer now reports whether
  the handoff was attempted, and a failure releases the seen slot so a
  relaunch or backlog redelivery retries.
- Restore-path logs: central/peripheral-restore captures logged the init
  sentinel bgRemaining=∞. The cache is now seeded in init's main-thread
  branch and restore captures route through the sampler, which refreshes
  the cached budget before logging.
- Persist cost: the dedup record was a full JSON encode + atomic write on
  the main actor per mutation (once per event during a backlog re-fetch).
  Writes now coalesce behind a 1s window, flushed immediately on
  background/terminate; panic wipe stays immediate.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

* Remove BoundedIDSet.remove, orphaned by the ExpiringIDSet migration

The drop-dedup sets that needed slot release moved to ExpiringIDSet;
remaining BoundedIDSet users only insert and check. Periphery caught it.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

---------

Co-authored-by: jack <jackjackbits@users.noreply.github.com>
Co-authored-by: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
jack
2026-07-09 18:23:46 +02:00
committed by GitHub
co-authored by jack Claude Fable 5
parent 304460ee83
commit f8ab0a7dc0
13 changed files with 754 additions and 55 deletions
+138 -21
View File
@@ -129,6 +129,15 @@ final class BLEService: NSObject {
// Application state tracking (thread-safe)
#if os(iOS)
private var isAppActive: Bool = true // Assume active initially
/// Last `UIApplication.shared.backgroundTimeRemaining` sampled on the
/// main thread, cached so bleQueue status logs can read it without ever
/// dispatching to main (see `captureBluetoothStatus` for the invariant).
private let backgroundTimeLock = NSLock()
private var _cachedBackgroundTimeRemaining: TimeInterval = .greatestFiniteMagnitude
private var cachedBackgroundTimeRemaining: TimeInterval {
backgroundTimeLock.lock(); defer { backgroundTimeLock.unlock() }
return _cachedBackgroundTimeRemaining
}
#endif
// MARK: - Core BLE Objects
@@ -176,6 +185,13 @@ final class BLEService: NSObject {
// Ingress link tracking for duplicate and last-hop suppression
private var ingressLinks = BLEIngressLinkRegistry()
// Inner message IDs of recently opened courier envelopes. Redundant
// copies of one message ride different envelopes (each seal uses a fresh
// ephemeral key, and bridge drops multiply across relays/couriers), so
// envelope-level dedup can't catch them; dedup on the inner ID before
// delivery so a duplicate costs one decrypt instead of a delivery + ack
// + handshake each. Owned by collectionsQueue barriers.
private var openedCourierMessageIDs = BoundedIDSet(capacity: TransportConfig.courierOpenedMessageIDCap)
private let logRateLimiter = BLELogRateLimiter(defaultMinimumInterval: 5)
private var pendingPeripheralWrites = BLEOutboundWriteBuffer()
@@ -265,12 +281,18 @@ final class BLEService: NSObject {
// Set up application state tracking (iOS only)
#if os(iOS)
// Check initial state on main thread
// Check initial state on main thread. The background-budget cache is
// seeded here too: a background-restore launch captures Bluetooth
// status before any lifecycle notification fires, and the init-time
// sentinel would log a meaningless bgRemaining= for exactly the
// wake window that matters.
if Thread.isMainThread {
isAppActive = UIApplication.shared.applicationState == .active
refreshCachedBackgroundTimeRemaining()
} else {
DispatchQueue.main.sync {
isAppActive = UIApplication.shared.applicationState == .active
refreshCachedBackgroundTimeRemaining()
}
}
@@ -777,7 +799,11 @@ final class BLEService: NSObject {
}
func triggerHandshake(with peerID: PeerID) {
initiateNoiseHandshake(with: peerID)
// Callers are on the main actor; the handshake broadcast sync-waits
// on bleQueue for link state, so hop off main first.
messageQueue.async { [weak self] in
self?.initiateNoiseHandshake(with: peerID)
}
}
// MARK: Noise identity/session access (narrow Transport wrappers)
@@ -932,6 +958,15 @@ final class BLEService: NSObject {
func sendReadReceipt(_ receipt: ReadReceipt, to peerID: PeerID) {
// Hop like sendMessage: callers are often on the main actor, and the
// send path sync-waits on bleQueue for link state the main thread
// must never block on bleQueue (see captureBluetoothStatus).
if DispatchQueue.getSpecific(key: messageQueueKey) == nil {
messageQueue.async { [weak self] in
self?.sendReadReceipt(receipt, to: peerID)
}
return
}
let payload = BLENoisePayloadFactory.readReceipt(originalMessageID: receipt.originalMessageID)
if noiseService.hasEstablishedSession(with: peerID) {
@@ -942,11 +977,15 @@ final class BLEService: NSObject {
SecureLogger.error("Failed to send read receipt: \(error)")
}
} else {
// Queue for after handshake and initiate if needed
// Queue for after handshake; initiate only while the peer is
// around to answer (see sendDeliveryAck absent senders must
// not turn queued acks into handshake floods).
collectionsQueue.sync(flags: .barrier) {
pendingNoiseSessionQueues.appendTypedPayload(payload, for: peerID)
}
if !noiseService.hasSession(with: peerID) { initiateNoiseHandshake(with: peerID) }
if !noiseService.hasSession(with: peerID), isPeerReachable(peerID) {
initiateNoiseHandshake(with: peerID)
}
SecureLogger.debug("🕒 Queued READ receipt for \(peerID.id.prefix(8))… until handshake completes", category: .session)
}
}
@@ -1395,6 +1434,15 @@ final class BLEService: NSObject {
}
func sendDeliveryAck(for messageID: String, to peerID: PeerID) {
// Hop like sendMessage: callers are often on the main actor, and the
// send path sync-waits on bleQueue for link state the main thread
// must never block on bleQueue (see captureBluetoothStatus).
if DispatchQueue.getSpecific(key: messageQueueKey) == nil {
messageQueue.async { [weak self] in
self?.sendDeliveryAck(for: messageID, to: peerID)
}
return
}
let payload = BLENoisePayloadFactory.delivered(messageID: messageID)
if noiseService.hasEstablishedSession(with: peerID) {
@@ -1404,11 +1452,18 @@ final class BLEService: NSObject {
SecureLogger.error("Failed to send delivery ACK: \(error)")
}
} else {
// Queue for after handshake and initiate if needed
// Queue for after handshake; initiate only while the peer is
// around to answer couriered/bridged mail routinely arrives
// from absent (or rotated) identities, and every duplicate copy
// initiating a handshake broadcast turns one undeliverable ack
// into a mesh-wide flood. The queued ack flushes whenever a
// session eventually establishes.
collectionsQueue.sync(flags: .barrier) {
pendingNoiseSessionQueues.appendTypedPayload(payload, for: peerID)
}
if !noiseService.hasSession(with: peerID) { initiateNoiseHandshake(with: peerID) }
if !noiseService.hasSession(with: peerID), isPeerReachable(peerID) {
initiateNoiseHandshake(with: peerID)
}
SecureLogger.debug("🕒 Queued DELIVERED ack for \(peerID.id.prefix(8))… until handshake completes", category: .session)
}
}
@@ -1680,7 +1735,10 @@ extension BLEService: CBCentralManagerDelegate {
recentPeripheralCache.record(peripheral, peripheralID: identifier, at: Date())
}
captureBluetoothStatus(context: "central-restore")
// Via the sampler (not a direct capture): it refreshes the cached
// background budget on main first, so the restore log shows the real
// wake window instead of the init sentinel.
logBluetoothStatus("central-restore")
if central.state == .poweredOn {
startScanning()
@@ -2439,7 +2497,8 @@ extension BLEService: CBPeripheralManagerDelegate {
}
}
captureBluetoothStatus(context: "peripheral-restore")
// Via the sampler for a fresh background budget (see central-restore).
logBluetoothStatus("peripheral-restore")
if peripheral.state == .poweredOn && !peripheral.isAdvertising {
peripheral.startAdvertising(buildAdvertisementData())
@@ -2720,19 +2779,37 @@ extension BLEService {
}
private func logBluetoothStatus(_ context: String) {
bleQueue.async { [weak self] in
guard let self = self else { return }
self.captureBluetoothStatus(context: context)
}
scheduleBluetoothStatusSample(after: 0, context: context)
}
private func scheduleBluetoothStatusSample(after delay: TimeInterval, context: String) {
bleQueue.asyncAfter(deadline: .now() + delay) { [weak self] in
guard let self = self else { return }
self.captureBluetoothStatus(context: context)
#if os(iOS)
// Sample the main-actor background budget first (async hop, never a
// sync wait), then log from bleQueue off the cache bleQueue must
// never block on main (see captureBluetoothStatus).
DispatchQueue.main.asyncAfter(deadline: .now() + delay) { [weak self] in
guard let self else { return }
self.refreshCachedBackgroundTimeRemaining()
self.bleQueue.async { self.captureBluetoothStatus(context: context) }
}
#else
bleQueue.asyncAfter(deadline: .now() + delay) { [weak self] in
self?.captureBluetoothStatus(context: context)
}
#endif
}
#if os(iOS)
/// Main thread only (reads main-actor UIApplication state).
private func refreshCachedBackgroundTimeRemaining() {
dispatchPrecondition(condition: .onQueue(.main))
let seconds = UIApplication.shared.backgroundTimeRemaining
backgroundTimeLock.lock()
_cachedBackgroundTimeRemaining = seconds
backgroundTimeLock.unlock()
}
#endif
private func captureBluetoothStatus(context: String) {
assert(DispatchQueue.getSpecific(key: bleQueueKey) != nil, "captureBluetoothStatus must run on bleQueue")
@@ -2750,11 +2827,15 @@ extension BLEService {
}
#if os(iOS)
var backgroundDescriptor = ""
var backgroundSeconds: TimeInterval = 0
DispatchQueue.main.sync {
backgroundSeconds = UIApplication.shared.backgroundTimeRemaining
}
// INVARIANT: bleQueue must NEVER sync-dispatch to the main thread.
// The main actor sync-waits on bleQueue along the send paths
// (readLinkState), so a main.sync here completes an ABBA deadlock
// field-verified as a permanent freeze when a courier-drop storm put
// an ack send (main bleQueue.sync) up against a status capture
// (bleQueue main.sync). backgroundTimeRemaining is main-actor
// state, so it is sampled on main and cached.
let backgroundSeconds = cachedBackgroundTimeRemaining
let backgroundDescriptor: String
if backgroundSeconds == .greatestFiniteMagnitude {
backgroundDescriptor = " bgRemaining=∞"
} else {
@@ -3044,6 +3125,15 @@ extension BLEService {
private func sendNoisePayload(_ typedPayload: Data, to peerID: PeerID) {
// Hop like sendMessage: the Transport-facing wrappers (verify/vouch/
// group payloads) call this from the main actor, and the send path
// sync-waits on bleQueue for link state.
if DispatchQueue.getSpecific(key: messageQueueKey) == nil {
messageQueue.async { [weak self] in
self?.sendNoisePayload(typedPayload, to: peerID)
}
return
}
guard noiseService.hasSession(with: peerID) else {
// No session yet - queue the payload SYNCHRONOUSLY before initiating handshake
// to prevent race where fast handshake completion drains empty queue
@@ -3291,6 +3381,23 @@ extension BLEService {
SecureLogger.warning("⚠️ Courier envelope carried unsupported payload type", category: .session)
return
}
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
}
// Redundant copies of one message arrive as distinct envelopes
// (fresh seal each: mesh couriers, bridge drops across relays),
// so dedup here on the inner message ID before delivery, ack,
// and handshake work. A duplicate costs only the decrypt above
// and at most one ack ever goes out per message ID.
let firstOpen = collectionsQueue.sync(flags: .barrier) {
openedCourierMessageIDs.insert(innerMessageID)
}
guard firstOpen else {
SecureLogger.debug("📦 Dropping duplicate courier envelope for message \(innerMessageID.prefix(8))", category: .session)
return
}
// 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.
@@ -3306,7 +3413,6 @@ extension BLEService {
let shortID = PeerID(publicKey: senderStaticKey)
let isKnownOnMesh = collectionsQueue.sync { peerRegistry.info(for: shortID) != nil }
let senderPeerID = isKnownOnMesh ? shortID : PeerID(hexData: senderStaticKey)
let payload = Data(typedPayload.dropFirst())
SecureLogger.debug("📦 Opened courier envelope from \(senderPeerID.id.prefix(8))", category: .session)
sfMetrics?.record(.courierOpened)
notifyUI { [weak self] in
@@ -3745,6 +3851,7 @@ extension BLEService {
#if os(iOS)
@objc private func appDidBecomeActive() {
isAppActive = true
refreshCachedBackgroundTimeRemaining()
// Restart scanning with allow duplicates when app becomes active
if centralManager?.state == .poweredOn {
centralManager?.stopScan()
@@ -3758,6 +3865,7 @@ extension BLEService {
@objc private func appDidEnterBackground() {
isAppActive = false
refreshCachedBackgroundTimeRemaining()
// Restart scanning without allow duplicates in background
if centralManager?.state == .poweredOn {
centralManager?.stopScan()
@@ -3851,6 +3959,15 @@ extension BLEService {
// MARK: Private Message Handling
private func sendPrivateMessage(_ content: String, to recipientID: PeerID, messageID: String) {
// Hop like sendMessage: the Transport-facing wrappers call this from
// the main actor (router sends, favorite notifications), and the send
// path sync-waits on bleQueue for link state.
if DispatchQueue.getSpecific(key: messageQueueKey) == nil {
messageQueue.async { [weak self] in
self?.sendPrivateMessage(content, to: recipientID, messageID: messageID)
}
return
}
// Sessions and wire recipient IDs are keyed by the short 16-hex form;
// callers may pass the full 64-hex noise key (mirrors sendFilePrivate).
let recipientID = recipientID.toShort()
@@ -32,14 +32,4 @@ struct BoundedIDSet {
}
return true
}
/// Releases a previously inserted ID so it can be re-added later. Used
/// when a tentatively-recorded item is later abandoned (e.g. a queued
/// drop evicted before it was ever published) and must become retryable.
mutating func remove(_ id: String) {
guard members.remove(id) != nil else { return }
if let index = order.firstIndex(of: id) {
order.remove(at: index)
}
}
}
@@ -9,6 +9,11 @@
import BitFoundation
import BitLogger
import Foundation
#if os(iOS)
import UIKit
#elseif os(macOS)
import AppKit
#endif
/// Courier delivery over the internet bridge: sealed courier envelopes are
/// parked on relays as kind-1401 "drops" tagged with their rotating
@@ -48,6 +53,10 @@ final class BridgeCourierService: ObservableObject {
/// Encoded envelope cap for a drop (16 KiB ciphertext + TLV slack).
static let maxDropEnvelopeBytes = 20 * 1024
static let maxTrackedIDs = 512
/// Coalescing window for dedup-record writes: a backlog re-fetch
/// mutates the seen set once per event, and each snapshot save is a
/// full JSON encode + atomic write on the main actor.
static let dedupPersistCoalesceSeconds: TimeInterval = 1.0
}
static let shared = BridgeCourierService()
@@ -70,7 +79,9 @@ final class BridgeCourierService: ObservableObject {
/// Opens a drop addressed to us (tag verified inside).
var openEnvelope: (@MainActor (CourierEnvelope) -> Void)?
/// Hands a drop to a matching local peer as a directed courier packet.
var deliverToPeer: (@MainActor (CourierEnvelope, PeerID) -> Void)?
/// Returns false when the handoff could not even be attempted (peer no
/// longer reachable), so the drop event stays retryable.
var deliverToPeer: (@MainActor (CourierEnvelope, PeerID) -> Bool)?
/// Held envelopes eligible for (re)publish, honoring the cooldown.
var heldEnvelopes: (@MainActor (TimeInterval) -> [CourierEnvelope])?
/// Timer injection for tests; nil arms a real `Task`.
@@ -81,21 +92,90 @@ final class BridgeCourierService: ObservableObject {
private(set) var myTagsHex: Set<String> = []
private(set) var watchedPeerTags: [(peerID: PeerID, tagsHex: Set<String>)] = []
private(set) var pendingDrops: [(envelope: CourierEnvelope, dedupKey: String?)] = []
/// Message IDs already published as drops (sender-side dedup).
private var publishedDropKeys: BoundedIDSet
/// Drop event IDs already handled (multi-relay dedup).
private var seenDropEventIDs: BoundedIDSet
/// Message IDs already published as drops (sender-side dedup) and drop
/// event IDs already handled (multi-relay dedup). Both persist across
/// relaunches: relays hold drops for the full 24h NIP-40 window and the
/// persisted outbox keeps re-depositing, so in-memory-only dedup meant
/// every relaunch republished the same message as a fresh drop and every
/// gateway relaunch re-delivered the whole backlog (field-verified
/// amplification storm). Entries age out with the 24h drop window.
private var publishedDropKeys: ExpiringIDSet
private var seenDropEventIDs: ExpiringIDSet
private var subscriptionOpen = false
private var lastSubscribedTags: Set<String> = []
private var refreshTimerArmed = false
private var lastAnnounceRefresh = Date.distantPast
private let now: () -> Date
private let dedupStore: BridgeDropDedupStore
init(now: @escaping () -> Date = Date.init) {
private var dedupPersistScheduled = false
init(now: @escaping () -> Date = Date.init, dedupStore: BridgeDropDedupStore? = nil) {
self.now = now
self.publishedDropKeys = BoundedIDSet(capacity: Limits.maxTrackedIDs)
self.seenDropEventIDs = BoundedIDSet(capacity: Limits.maxTrackedIDs)
self.dedupStore = dedupStore ?? BridgeDropDedupStore(persistsToDisk: !TestEnvironment.isRunningTests)
let snapshot = self.dedupStore.load()
let date = now()
self.publishedDropKeys = ExpiringIDSet(
capacity: Limits.maxTrackedIDs,
lifetime: CourierEnvelope.maxLifetimeSeconds,
entries: snapshot.publishedDropKeys,
now: date
)
self.seenDropEventIDs = ExpiringIDSet(
capacity: Limits.maxTrackedIDs,
lifetime: CourierEnvelope.maxLifetimeSeconds,
entries: snapshot.seenDropEventIDs,
now: date
)
// A coalesced dedup write scheduled just before a background kill
// would be lost; flush when the app backgrounds or terminates.
#if os(iOS)
let flushNotifications = [UIApplication.didEnterBackgroundNotification, UIApplication.willTerminateNotification]
#else
let flushNotifications = [NSApplication.willTerminateNotification]
#endif
for name in flushNotifications {
NotificationCenter.default.addObserver(forName: name, object: nil, queue: .main) { [weak self] _ in
MainActor.assumeIsolated { self?.flushDedupSnapshot() }
}
}
}
/// Schedules a coalesced write of the dedup record (see
/// `Limits.dedupPersistCoalesceSeconds`); lifecycle notifications flush
/// any scheduled write before a background kill could drop it.
private func persistDedup() {
guard !dedupPersistScheduled else { return }
dedupPersistScheduled = true
Task { @MainActor [weak self] in
try? await Task.sleep(nanoseconds: UInt64(Limits.dedupPersistCoalesceSeconds * 1_000_000_000))
guard let self else { return }
self.dedupPersistScheduled = false
self.flushDedupSnapshot()
}
}
/// Writes the dedup record now. Sender keys still sitting in the
/// in-memory `pendingDrops` queue are excluded: their drop is not durable
/// until it actually reaches a relay, and persisting the key early would
/// turn "app killed before relays connected" into a silent 24h blackhole
/// (the relaunch loses the queued drop but the persisted key blocks every
/// re-deposit). `flushPendingDrops` re-persists once they publish.
func flushDedupSnapshot() {
let pendingKeys = Set(pendingDrops.compactMap(\.dedupKey))
dedupStore.save(BridgeDropDedupStore.Snapshot(
publishedDropKeys: publishedDropKeys.entries.filter { !pendingKeys.contains($0.key) },
seenDropEventIDs: seenDropEventIDs.entries
))
}
/// Panic wipe: forget queued drops and the persisted dedup record.
func wipe() {
pendingDrops.removeAll()
publishedDropKeys = ExpiringIDSet(capacity: Limits.maxTrackedIDs, lifetime: CourierEnvelope.maxLifetimeSeconds)
seenDropEventIDs = ExpiringIDSet(capacity: Limits.maxTrackedIDs, lifetime: CourierEnvelope.maxLifetimeSeconds)
dedupStore.wipe()
}
// MARK: - Sender role
@@ -108,14 +188,15 @@ final class BridgeCourierService: ObservableObject {
@discardableResult
func depositDrop(content: String, messageID: String, recipientNoiseKey: Data) -> Bool {
guard bridgeEnabled?() ?? false else { return false }
guard !publishedDropKeys.contains(messageID) else { return false }
guard !publishedDropKeys.contains(messageID, now: now()) else { return false }
guard let envelope = sealEnvelope?(content, messageID, recipientNoiseKey) else { return false }
// An envelope that can't encode within the drop size caps fails the
// same way on every attempt (size is a function of the content, not
// of the sealing); consume the dedup slot so the retry sweep stops
// re-running Noise sealing on a drop that can never ship.
guard let encoded = envelope.encode(), encoded.count <= Limits.maxDropEnvelopeBytes else {
publishedDropKeys.insert(messageID)
publishedDropKeys.insert(messageID, now: now())
persistDedup()
return false
}
// Only consume the sender-side dedup slot once the drop is durably
@@ -124,7 +205,8 @@ final class BridgeCourierService: ObservableObject {
// router's retry sweep can attempt a fresh deposit rather than
// marking the message "carried" and blocking retries forever.
guard publishDrop(envelope, messageID: messageID) else { return false }
publishedDropKeys.insert(messageID)
publishedDropKeys.insert(messageID, now: now())
persistDedup()
return true
}
@@ -153,7 +235,10 @@ final class BridgeCourierService: ObservableObject {
let evicted = pendingDrops.removeFirst()
// The oldest queued drop is being dropped before it ever
// published; release its dedup slot so it stays retryable.
if let key = evicted.dedupKey { publishedDropKeys.remove(key) }
if let key = evicted.dedupKey {
publishedDropKeys.remove(key)
persistDedup()
}
}
return true
}
@@ -180,8 +265,13 @@ final class BridgeCourierService: ObservableObject {
for item in queued where !publishDrop(item.envelope, messageID: item.dedupKey) {
// Compose failed with relays up: release the slot so the router's
// retry sweep can attempt a fresh deposit.
if let key = item.dedupKey { publishedDropKeys.remove(key) }
if let key = item.dedupKey {
publishedDropKeys.remove(key)
}
}
// Flushed keys just became durable (published, so no longer excluded
// as pending) or were released above; either way the record changed.
persistDedup()
}
// MARK: - Subscription (recipient + gateway watch)
@@ -265,7 +355,8 @@ final class BridgeCourierService: ObservableObject {
func handleDropEvent(_ event: NostrEvent) {
guard bridgeEnabled?() ?? false else { return }
guard event.kind == NostrProtocol.EventKind.courierDrop.rawValue else { return }
guard seenDropEventIDs.insert(event.id) else { return }
guard seenDropEventIDs.insert(event.id, now: now()) else { return }
persistDedup()
guard let data = Data(base64Encoded: event.content),
data.count <= Limits.maxDropEnvelopeBytes,
let envelope = CourierEnvelope.decode(data),
@@ -285,7 +376,14 @@ final class BridgeCourierService: ObservableObject {
}
if let match = watchedPeerTags.first(where: { $0.tagsHex.contains(tagHex) }) {
SecureLogger.info("📦🌉 Courier drop fetched for local peer \(match.peerID.id.prefix(8))", category: .session)
deliverToPeer?(envelope, match.peerID)
if deliverToPeer?(envelope, match.peerID) != true {
// The best-effort handoff never left this device (the peer
// walked away between the relay fetch and the mesh send).
// Release the seen slot so a relaunch or backlog redelivery
// retries a single-gateway island has no other carrier.
seenDropEventIDs.remove(event.id)
persistDedup()
}
}
}
@@ -0,0 +1,137 @@
//
// BridgeDropDedupStore.swift
// bitchat
//
// This is free and unencumbered software released into the public domain.
// For more information, see <https://unlicense.org>
//
import BitLogger
import Foundation
/// ID set with per-entry timestamps: entries expire after `lifetime` and the
/// oldest are evicted past `capacity`. The bridge's dedup caches use this
/// instead of `BoundedIDSet` because their contents persist across relaunches
/// and must age out with the 24h drop window they guard.
struct ExpiringIDSet {
private(set) var entries: [String: Date]
let capacity: Int
let lifetime: TimeInterval
init(capacity: Int, lifetime: TimeInterval, entries: [String: Date] = [:], now: Date = Date()) {
self.capacity = capacity
self.lifetime = lifetime
self.entries = entries
prune(now: now)
}
func contains(_ id: String, now: Date) -> Bool {
guard let recorded = entries[id] else { return false }
return now.timeIntervalSince(recorded) <= lifetime
}
/// Returns false when the ID was already present (and unexpired).
@discardableResult
mutating func insert(_ id: String, now: Date) -> Bool {
guard !contains(id, now: now) else { return false }
entries[id] = now
prune(now: now)
return true
}
/// Releases a previously inserted ID so it can be re-added later (e.g. a
/// queued drop evicted before it ever published must become retryable).
mutating func remove(_ id: String) {
entries.removeValue(forKey: id)
}
private mutating func prune(now: Date) {
if entries.contains(where: { now.timeIntervalSince($0.value) > lifetime }) {
entries = entries.filter { now.timeIntervalSince($0.value) <= lifetime }
}
let overflow = entries.count - capacity
guard overflow > 0 else { return }
for (id, _) in entries.sorted(by: { $0.value < $1.value }).prefix(overflow) {
entries.removeValue(forKey: id)
}
}
}
/// Disk persistence for the bridge courier's drop-dedup record. Relays hold
/// drops for the full 24h NIP-40 window and redeliver them on every launch,
/// and the 120s outbox sweep re-deposits anything undelivered so with
/// in-memory-only dedup every relaunch republished the same message as a
/// fresh drop (fresh throwaway seal, undeduplicatable downstream) and every
/// gateway relaunch re-delivered the whole backlog. Field-verified: ~20
/// copies of one DM delivered in 40ms fed the storm behind a permanent
/// device freeze. Persisting both sides caps this at one drop per message
/// ID per 24h regardless of relaunch count.
///
/// Contents are opaque IDs (message UUIDs, relay event IDs) no plaintext,
/// no peer identities so until-first-unlock protection matches
/// `NostrProcessedEventStore`, and the file must load during a
/// locked-background restoration relaunch. Wiped on panic with the rest of
/// the courier state.
final class BridgeDropDedupStore {
struct Snapshot: Codable {
var publishedDropKeys: [String: Date]
var seenDropEventIDs: [String: Date]
}
private let fileURL: URL?
/// - Parameter fileURL: Overrides the on-disk location (tests). Ignored
/// when `persistsToDisk` is false.
init(persistsToDisk: Bool = true, fileURL: URL? = nil) {
self.fileURL = persistsToDisk ? (fileURL ?? Self.defaultFileURL()) : nil
}
func load() -> Snapshot {
guard let fileURL,
let data = try? Data(contentsOf: fileURL),
let snapshot = try? JSONDecoder().decode(Snapshot.self, from: data) else {
return Snapshot(publishedDropKeys: [:], seenDropEventIDs: [:])
}
return snapshot
}
func save(_ snapshot: Snapshot) {
guard let fileURL else { return }
guard !(snapshot.publishedDropKeys.isEmpty && snapshot.seenDropEventIDs.isEmpty) else {
try? FileManager.default.removeItem(at: fileURL)
return
}
do {
try FileManager.default.createDirectory(
at: fileURL.deletingLastPathComponent(),
withIntermediateDirectories: true
)
let data = try JSONEncoder().encode(snapshot)
var options: Data.WritingOptions = [.atomic]
#if os(iOS)
options.insert(.completeFileProtectionUntilFirstUserAuthentication)
#endif
try data.write(to: fileURL, options: options)
} catch {
SecureLogger.error("Failed to persist bridge drop dedup record: \(error)", category: .session)
}
}
/// Panic wipe: forget which drops we published or handled.
func wipe() {
guard let fileURL else { return }
try? FileManager.default.removeItem(at: fileURL)
}
private static func defaultFileURL() -> URL? {
guard let base = try? FileManager.default.url(
for: .applicationSupportDirectory,
in: .userDomainMask,
appropriateFor: nil,
create: true
) else { return nil }
return base
.appendingPathComponent("courier", isDirectory: true)
.appendingPathComponent("bridge-drop-dedup.json")
}
}
+8 -2
View File
@@ -351,9 +351,15 @@ final class MessageRouter {
}
private func enqueue(_ message: QueuedMessage, for peerID: PeerID) {
var message = message
var queue = outbox[peerID] ?? []
// Re-sending an already-queued ID replaces the entry (keeps attempt count fresh)
queue.removeAll { $0.messageID == message.messageID }
// Re-sending an already-queued ID replaces the entry (keeps attempt
// count fresh) but must not forget which couriers already carry it,
// or the replacement re-burns the same courier slots.
if let existing = queue.firstIndex(where: { $0.messageID == message.messageID }) {
message.depositedCourierKeys.formUnion(queue[existing].depositedCourierKeys)
queue.remove(at: existing)
}
queue.append(message)
// Enforce per-peer size limit with FIFO eviction
+3
View File
@@ -345,6 +345,9 @@ enum TransportConfig {
// Cooldown between speculative multi-hop handovers of the same envelope
// toward a recipient heard only via relayed announces.
static let courierRemoteHandoverCooldownSeconds: TimeInterval = 10 * 60
// Recently opened courier inner message IDs kept for receiver-side dedup
// (redundant copies ride distinct seals, so only the inner ID matches).
static let courierOpenedMessageIDCap: Int = 512
// One-time prekey bundles (forward-secret courier sealing)
// Own gossip-sync round for bundles: modest cadence, bounded peer count,
+1
View File
@@ -1173,6 +1173,7 @@ final class ChatViewModel: ObservableObject, BitchatDelegate, TransportEventDele
// our own queued outbox, the carried public history, and the
// counters describing all of it
CourierStore.shared.wipe()
BridgeCourierService.shared.wipe()
messageRouter.wipeOutbox()
GossipMessageArchive.wipeDefault()
StoreAndForwardMetrics.shared.reset()
@@ -567,7 +567,7 @@ private extension ChatViewModelBootstrapper {
bleService?.openBridgedCourierEnvelope(envelope)
}
courier.deliverToPeer = { [weak bleService] envelope, peerID in
bleService?.deliverBridgedEnvelope(envelope, to: peerID)
bleService?.deliverBridgedEnvelope(envelope, to: peerID) ?? false
}
courier.heldEnvelopes = { cooldown in
CourierStore.shared.envelopesForBridgePublish(cooldown: cooldown)