Compare commits

...
5 changed files with 115 additions and 19 deletions
+57 -7
View File
@@ -44,7 +44,12 @@ class NostrRelayManager: ObservableObject {
private var cancellables = Set<AnyCancellable>() private var cancellables = Set<AnyCancellable>()
// Message queue for reliability // Message queue for reliability
private var messageQueue: [(event: NostrEvent, relayUrls: [String])] = [] // Pending sends held only for relays that are not yet connected.
private struct PendingSend {
var event: NostrEvent
var pendingRelays: Set<String>
}
private var messageQueue: [PendingSend] = []
private let messageQueueLock = NSLock() private let messageQueueLock = NSLock()
// Exponential backoff configuration // Exponential backoff configuration
@@ -96,15 +101,56 @@ class NostrRelayManager: ObservableObject {
let targetRelays = relayUrls ?? Self.defaultRelays let targetRelays = relayUrls ?? Self.defaultRelays
ensureConnections(to: targetRelays) ensureConnections(to: targetRelays)
// Add to queue for reliability // Attempt immediate send to relays with active connections; queue the rest
messageQueueLock.lock() var stillPending = Set<String>()
messageQueue.append((event, targetRelays))
messageQueueLock.unlock()
// Attempt immediate send
for relayUrl in targetRelays { for relayUrl in targetRelays {
if let connection = connections[relayUrl] { if let connection = connections[relayUrl] {
sendToRelay(event: event, connection: connection, relayUrl: relayUrl) sendToRelay(event: event, connection: connection, relayUrl: relayUrl)
} else {
stillPending.insert(relayUrl)
}
}
if !stillPending.isEmpty {
messageQueueLock.lock()
messageQueue.append(PendingSend(event: event, pendingRelays: stillPending))
messageQueueLock.unlock()
}
}
/// Try to flush any queued messages for relays that are now connected.
private func flushMessageQueue(for relayUrl: String? = nil) {
messageQueueLock.lock()
defer { messageQueueLock.unlock() }
guard !messageQueue.isEmpty else { return }
if let target = relayUrl {
// Flush only for a specific relay
for i in (0..<messageQueue.count).reversed() {
var item = messageQueue[i]
if item.pendingRelays.contains(target), let conn = connections[target] {
sendToRelay(event: item.event, connection: conn, relayUrl: target)
item.pendingRelays.remove(target)
if item.pendingRelays.isEmpty {
messageQueue.remove(at: i)
} else {
messageQueue[i] = item
}
}
}
} else {
// Flush for any relays that now have connections
for i in (0..<messageQueue.count).reversed() {
var item = messageQueue[i]
for url in item.pendingRelays {
if let conn = connections[url] {
sendToRelay(event: item.event, connection: conn, relayUrl: url)
item.pendingRelays.remove(url)
}
}
if item.pendingRelays.isEmpty {
messageQueue.remove(at: i)
} else {
messageQueue[i] = item
}
} }
} }
} }
@@ -389,6 +435,10 @@ class NostrRelayManager: ObservableObject {
} }
} }
updateConnectionStatus() updateConnectionStatus()
// If we just connected to this relay, flush any queued sends targeting it
if isConnected {
flushMessageQueue(for: url)
}
} }
private func updateConnectionStatus() { private func updateConnectionStatus() {
+24 -1
View File
@@ -223,7 +223,30 @@ final class BLEService: NSObject {
peripheral.writeValue(data, for: characteristic, type: .withoutResponse) peripheral.writeValue(data, for: characteristic, type: .withoutResponse)
} else { } else {
self.collectionsQueue.async(flags: .barrier) { self.collectionsQueue.async(flags: .barrier) {
self.pendingPeripheralWrites[uuid, default: []].append(data) var queue = self.pendingPeripheralWrites[uuid] ?? []
let capBytes = TransportConfig.blePendingWriteBufferCapBytes
let newSize = data.count
// If single chunk exceeds cap, drop it immediately
if newSize > capBytes {
SecureLogger.log("⚠️ Dropping oversized write chunk (\(newSize)B) for peripheral \(uuid)",
category: SecureLogger.session, level: .warning)
} else {
// Append and trim from the front to respect cap
var total = queue.reduce(0) { $0 + $1.count }
queue.append(data)
total += newSize
if total > capBytes {
var removedBytes = 0
while total > capBytes && !queue.isEmpty {
let removed = queue.removeFirst()
removedBytes += removed.count
total -= removed.count
}
SecureLogger.log("📉 Trimmed pending write buffer for \(uuid) by \(removedBytes)B to \(total)B",
category: SecureLogger.session, level: .warning)
}
self.pendingPeripheralWrites[uuid] = queue.isEmpty ? nil : queue
}
} }
} }
} }
+23 -8
View File
@@ -89,10 +89,19 @@ final class MessageRouter {
// MARK: - Outbox Management // MARK: - Outbox Management
private func canSendViaNostr(peerID: String) -> Bool { private func canSendViaNostr(peerID: String) -> Bool {
guard let noiseKey = Data(hexString: peerID) else { return false } // Two forms are supported:
if let fav = FavoritesPersistenceService.shared.getFavoriteStatus(for: noiseKey), // - 64-hex Noise public key (32 bytes)
fav.peerNostrPublicKey != nil { // - 16-hex short peer ID (derived from Noise pubkey)
return true if peerID.count == 64, let noiseKey = Data(hexString: peerID) {
if let fav = FavoritesPersistenceService.shared.getFavoriteStatus(for: noiseKey),
fav.peerNostrPublicKey != nil {
return true
}
} else if peerID.count == 16 {
if let fav = FavoritesPersistenceService.shared.getFavoriteStatus(forPeerID: peerID),
fav.peerNostrPublicKey != nil {
return true
}
} }
return false return false
} }
@@ -101,6 +110,7 @@ final class MessageRouter {
guard let queued = outbox[peerID], !queued.isEmpty else { return } guard let queued = outbox[peerID], !queued.isEmpty else { return }
SecureLogger.log("Flushing outbox for \(peerID.prefix(8))… count=\(queued.count)", SecureLogger.log("Flushing outbox for \(peerID.prefix(8))… count=\(queued.count)",
category: SecureLogger.session, level: .debug) category: SecureLogger.session, level: .debug)
var remaining: [(content: String, nickname: String, messageID: String)] = []
// Prefer mesh if connected; else try Nostr if mapping exists // Prefer mesh if connected; else try Nostr if mapping exists
for (content, nickname, messageID) in queued { for (content, nickname, messageID) in queued {
if mesh.isPeerReachable(peerID) { if mesh.isPeerReachable(peerID) {
@@ -112,14 +122,19 @@ final class MessageRouter {
category: SecureLogger.session, level: .debug) category: SecureLogger.session, level: .debug)
nostr.sendPrivateMessage(content, to: peerID, recipientNickname: nickname, messageID: messageID) nostr.sendPrivateMessage(content, to: peerID, recipientNickname: nickname, messageID: messageID)
} else { } else {
continue // Keep unsent items queued
remaining.append((content, nickname, messageID))
} }
} }
// Remove all flushed items (remaining ones, if any, will be re-queued on next call) // Persist only items we could not send
outbox[peerID]?.removeAll() if remaining.isEmpty {
outbox.removeValue(forKey: peerID)
} else {
outbox[peerID] = remaining
}
} }
func flushAllOutbox() { func flushAllOutbox() {
for key in outbox.keys { flushOutbox(for: key) } for key in Array(outbox.keys) { flushOutbox(for: key) }
} }
} }
+8 -2
View File
@@ -27,6 +27,7 @@ class UnifiedPeerService: ObservableObject, TransportPeerEventsDelegate {
private var peerIndex: [String: BitchatPeer] = [:] private var peerIndex: [String: BitchatPeer] = [:]
private var fingerprintCache: [String: String] = [:] // peerID -> fingerprint private var fingerprintCache: [String: String] = [:] // peerID -> fingerprint
private let meshService: Transport private let meshService: Transport
weak var messageRouter: MessageRouter?
private let favoritesService = FavoritesPersistenceService.shared private let favoritesService = FavoritesPersistenceService.shared
private var cancellables = Set<AnyCancellable>() private var cancellables = Set<AnyCancellable>()
@@ -330,8 +331,13 @@ class UnifiedPeerService: ObservableObject, TransportPeerEventsDelegate {
SecureLogger.log("⭐️ Toggled favorite for '\(finalNickname)' (peerID: \(peerID), was: \(wasFavorite), now: \(!wasFavorite))", SecureLogger.log("⭐️ Toggled favorite for '\(finalNickname)' (peerID: \(peerID), was: \(wasFavorite), now: \(!wasFavorite))",
category: SecureLogger.session, level: .debug) category: SecureLogger.session, level: .debug)
// Send favorite notification to the peer // Send favorite notification to the peer via router (mesh or Nostr)
meshService.sendFavoriteNotification(to: peerID, isFavorite: !wasFavorite) if let router = messageRouter {
router.sendFavoriteNotification(to: peerID, isFavorite: !wasFavorite)
} else {
// Fallback to mesh-only if router not yet wired
meshService.sendFavoriteNotification(to: peerID, isFavorite: !wasFavorite)
}
// Force update of peers to reflect the change // Force update of peers to reflect the change
updatePeers() updatePeers()
+2
View File
@@ -492,6 +492,8 @@ class ChatViewModel: ObservableObject, BitchatDelegate {
self.messageRouter = MessageRouter(mesh: meshService, nostr: nostrTransport) self.messageRouter = MessageRouter(mesh: meshService, nostr: nostrTransport)
// Route receipts from PrivateChatManager through MessageRouter // Route receipts from PrivateChatManager through MessageRouter
self.privateChatManager.messageRouter = self.messageRouter self.privateChatManager.messageRouter = self.messageRouter
// Allow UnifiedPeerService to route favorite notifications via mesh/Nostr
self.unifiedPeerService.messageRouter = self.messageRouter
self.autocompleteService = AutocompleteService() self.autocompleteService = AutocompleteService()
// Wire up dependencies // Wire up dependencies