Enhance message handling and synchronization in BluetoothMeshService

- Introduce thread safety for recentlySentMessages using NSLock to prevent race conditions.
- Adjust message delay parameters for faster synchronization: reduce min/max message delays and adjust initial announce timing.
- Update sendMessage and sendEncryptedRoomMessage methods to accept optional messageID and timestamp parameters for better tracking.
- Improve message deduplication logic in ChatViewModel to prevent duplicates from being added when echoing back sent messages.
- Modify MessageRetryService to include original message ID and timestamp for better retry handling.
This commit is contained in:
Nelson Campos
2025-07-07 23:19:24 -03:00
parent d3c1b77015
commit 602657fec2
3 changed files with 151 additions and 55 deletions
+74 -43
View File
@@ -74,6 +74,7 @@ class BluetoothMeshService: NSObject {
private var cachedMessagesSentToPeer: Set<String> = [] // Track which peers have already received cached messages private var cachedMessagesSentToPeer: Set<String> = [] // Track which peers have already received cached messages
private var receivedMessageTimestamps: [String: Date] = [:] // Track timestamps of received messages for debugging private var receivedMessageTimestamps: [String: Date] = [:] // Track timestamps of received messages for debugging
private var recentlySentMessages: Set<String> = [] // Short-term cache to prevent any duplicate sends private var recentlySentMessages: Set<String> = [] // Short-term cache to prevent any duplicate sends
private let recentlySentMessagesLock = NSLock() // Thread safety for recentlySentMessages
private var lastMessageFromPeer: [String: Date] = [:] // Track last message time from each peer for connection prioritization private var lastMessageFromPeer: [String: Date] = [:] // Track last message time from each peer for connection prioritization
// Battery and range optimizations // Battery and range optimizations
@@ -100,8 +101,8 @@ class BluetoothMeshService: NSObject {
private var advertisingTimer: Timer? // Timer for interval-based advertising private var advertisingTimer: Timer? // Timer for interval-based advertising
// Timing randomization for privacy // Timing randomization for privacy
private let minMessageDelay: TimeInterval = 0.05 // 50ms minimum private let minMessageDelay: TimeInterval = 0.01 // 10ms minimum for faster sync
private let maxMessageDelay: TimeInterval = 0.5 // 500ms maximum private let maxMessageDelay: TimeInterval = 0.1 // 100ms maximum for faster sync
// Fragment handling // Fragment handling
private var incomingFragments: [String: [Int: Data]] = [:] // fragmentID -> [index: data] private var incomingFragments: [String: [Int: Data]] = [:] // fragmentID -> [index: data]
@@ -366,7 +367,7 @@ class BluetoothMeshService: NSObject {
} }
// Send initial announces after services are ready // Send initial announces after services are ready
DispatchQueue.main.asyncAfter(deadline: .now() + 1.0) { [weak self] in DispatchQueue.main.asyncAfter(deadline: .now() + 0.2) { [weak self] in
self?.sendBroadcastAnnounce() self?.sendBroadcastAnnounce()
} }
@@ -395,7 +396,7 @@ class BluetoothMeshService: NSObject {
} }
// Send multiple times for reliability with jittered delays // Send multiple times for reliability with jittered delays
for baseDelay in [0.5, 1.0, 2.0] { for baseDelay in [0.2, 0.5, 1.0] {
let jitteredDelay = baseDelay + self.randomDelay() let jitteredDelay = baseDelay + self.randomDelay()
DispatchQueue.main.asyncAfter(deadline: .now() + jitteredDelay) { [weak self] in DispatchQueue.main.asyncAfter(deadline: .now() + jitteredDelay) { [weak self] in
guard let self = self else { return } guard let self = self else { return }
@@ -489,7 +490,7 @@ class BluetoothMeshService: NSObject {
self.characteristic = characteristic self.characteristic = characteristic
} }
func sendMessage(_ content: String, mentions: [String] = [], room: String? = nil, to recipientID: String? = nil) { func sendMessage(_ content: String, mentions: [String] = [], room: String? = nil, to recipientID: String? = nil, messageID: String? = nil, timestamp: Date? = nil) {
messageQueue.async { [weak self] in messageQueue.async { [weak self] in
guard let self = self else { return } guard let self = self else { return }
@@ -497,9 +498,10 @@ class BluetoothMeshService: NSObject {
let senderNick = nickname?.nickname ?? self.myPeerID let senderNick = nickname?.nickname ?? self.myPeerID
let message = BitchatMessage( let message = BitchatMessage(
id: messageID,
sender: senderNick, sender: senderNick,
content: content, content: content,
timestamp: Date(), timestamp: timestamp ?? Date(),
isRelay: false, isRelay: false,
originalSender: nil, originalSender: nil,
isPrivate: false, isPrivate: false,
@@ -532,12 +534,21 @@ class BluetoothMeshService: NSObject {
// Track this message to prevent duplicate sends // Track this message to prevent duplicate sends
let msgID = "\(packet.timestamp)-\(self.myPeerID)-\(packet.payload.prefix(32).hashValue)" let msgID = "\(packet.timestamp)-\(self.myPeerID)-\(packet.payload.prefix(32).hashValue)"
if !self.recentlySentMessages.contains(msgID) {
self.recentlySentMessagesLock.lock()
let shouldSend = !self.recentlySentMessages.contains(msgID)
if shouldSend {
self.recentlySentMessages.insert(msgID) self.recentlySentMessages.insert(msgID)
}
self.recentlySentMessagesLock.unlock()
if shouldSend {
// Clean up old entries after 10 seconds // Clean up old entries after 10 seconds
self.messageQueue.asyncAfter(deadline: .now() + 10.0) { [weak self] in self.messageQueue.asyncAfter(deadline: .now() + 10.0) { [weak self] in
self?.recentlySentMessages.remove(msgID) guard let self = self else { return }
self.recentlySentMessagesLock.lock()
self.recentlySentMessages.remove(msgID)
self.recentlySentMessagesLock.unlock()
} }
// Add random delay before initial send // Add random delay before initial send
@@ -630,12 +641,21 @@ class BluetoothMeshService: NSObject {
// Track to prevent duplicate sends // Track to prevent duplicate sends
let msgID = "\(packet.timestamp)-\(self.myPeerID)-\(packet.payload.prefix(32).hashValue)" let msgID = "\(packet.timestamp)-\(self.myPeerID)-\(packet.payload.prefix(32).hashValue)"
if !self.recentlySentMessages.contains(msgID) {
self.recentlySentMessagesLock.lock()
let shouldSend = !self.recentlySentMessages.contains(msgID)
if shouldSend {
self.recentlySentMessages.insert(msgID) self.recentlySentMessages.insert(msgID)
}
self.recentlySentMessagesLock.unlock()
if shouldSend {
// Clean up after 10 seconds // Clean up after 10 seconds
self.messageQueue.asyncAfter(deadline: .now() + 10.0) { [weak self] in self.messageQueue.asyncAfter(deadline: .now() + 10.0) { [weak self] in
self?.recentlySentMessages.remove(msgID) guard let self = self else { return }
self.recentlySentMessagesLock.lock()
self.recentlySentMessages.remove(msgID)
self.recentlySentMessagesLock.unlock()
} }
// Message tracking is now done in ChatViewModel to ensure consistent message IDs // Message tracking is now done in ChatViewModel to ensure consistent message IDs
@@ -784,7 +804,7 @@ class BluetoothMeshService: NSObject {
} }
} }
func sendEncryptedRoomMessage(_ content: String, mentions: [String], room: String, roomKey: SymmetricKey) { func sendEncryptedRoomMessage(_ content: String, mentions: [String], room: String, roomKey: SymmetricKey, messageID: String? = nil, timestamp: Date? = nil) {
messageQueue.async { [weak self] in messageQueue.async { [weak self] in
guard let self = self else { return } guard let self = self else { return }
@@ -802,9 +822,10 @@ class BluetoothMeshService: NSObject {
// Create message with encrypted content // Create message with encrypted content
let message = BitchatMessage( let message = BitchatMessage(
id: messageID,
sender: senderNick, sender: senderNick,
content: "", // Empty placeholder since actual content is encrypted content: "", // Empty placeholder since actual content is encrypted
timestamp: Date(), timestamp: timestamp ?? Date(),
isRelay: false, isRelay: false,
originalSender: nil, originalSender: nil,
isPrivate: false, isPrivate: false,
@@ -1184,7 +1205,7 @@ class BluetoothMeshService: NSObject {
// Send cached messages with slight delay between each // Send cached messages with slight delay between each
for (index, storedMessage) in messagesToSend.enumerated() { for (index, storedMessage) in messagesToSend.enumerated() {
let delay = Double(index) * 0.1 // 100ms between messages let delay = Double(index) * 0.02 // 20ms between messages for faster sync
DispatchQueue.main.asyncAfter(deadline: .now() + delay) { [weak peripheral] in DispatchQueue.main.asyncAfter(deadline: .now() + delay) { [weak peripheral] in
guard let peripheral = peripheral, guard let peripheral = peripheral,
@@ -1238,8 +1259,10 @@ class BluetoothMeshService: NSObject {
private func broadcastPacket(_ packet: BitchatPacket) { private func broadcastPacket(_ packet: BitchatPacket) {
guard let data = packet.toBinaryData() else { guard let data = packet.toBinaryData() else {
// print("[ERROR] Failed to convert packet to binary data") // print("[ERROR] Failed to convert packet to binary data")
// Add to retry queue if this is a message packet // Add to retry queue if this is a message packet AND it's our own message
if packet.type == MessageType.message.rawValue, if let senderID = String(data: packet.senderID.trimmingNullBytes(), encoding: .utf8),
senderID == self.myPeerID,
packet.type == MessageType.message.rawValue,
let message = BitchatMessage.fromBinaryPayload(packet.payload) { let message = BitchatMessage.fromBinaryPayload(packet.payload) {
MessageRetryService.shared.addMessageForRetry( MessageRetryService.shared.addMessageForRetry(
content: message.content, content: message.content,
@@ -1247,7 +1270,10 @@ class BluetoothMeshService: NSObject {
room: message.room, room: message.room,
isPrivate: message.isPrivate, isPrivate: message.isPrivate,
recipientPeerID: nil, recipientPeerID: nil,
recipientNickname: message.recipientNickname recipientNickname: message.recipientNickname,
roomKey: nil,
originalMessageID: message.id,
originalTimestamp: message.timestamp
) )
} }
return return
@@ -1294,29 +1320,36 @@ class BluetoothMeshService: NSObject {
} }
} }
// If no peers received the message, add to retry queue // If no peers received the message, add to retry queue ONLY if it's our own message
if sentToPeripherals == 0 && sentToCentrals == 0 { if sentToPeripherals == 0 && sentToCentrals == 0 {
if packet.type == MessageType.message.rawValue, // Check if this packet originated from us
let message = BitchatMessage.fromBinaryPayload(packet.payload) { if let senderID = String(data: packet.senderID.trimmingNullBytes(), encoding: .utf8),
// For encrypted room messages, we need to preserve the room key senderID == self.myPeerID {
var roomKeyData: Data? = nil // This is our own message that failed to send
if let room = message.room, message.isEncrypted { if packet.type == MessageType.message.rawValue,
// This is an encrypted room message let message = BitchatMessage.fromBinaryPayload(packet.payload) {
if let viewModel = delegate as? ChatViewModel, // For encrypted room messages, we need to preserve the room key
let roomKey = viewModel.roomKeys[room] { var roomKeyData: Data? = nil
roomKeyData = roomKey.withUnsafeBytes { Data($0) } if let room = message.room, message.isEncrypted {
// This is an encrypted room message
if let viewModel = delegate as? ChatViewModel,
let roomKey = viewModel.roomKeys[room] {
roomKeyData = roomKey.withUnsafeBytes { Data($0) }
}
} }
MessageRetryService.shared.addMessageForRetry(
content: message.content,
mentions: message.mentions,
room: message.room,
isPrivate: message.isPrivate,
recipientPeerID: nil,
recipientNickname: message.recipientNickname,
roomKey: roomKeyData,
originalMessageID: message.id,
originalTimestamp: message.timestamp
)
} }
MessageRetryService.shared.addMessageForRetry(
content: message.content,
mentions: message.mentions,
room: message.room,
isPrivate: message.isPrivate,
recipientPeerID: nil,
recipientNickname: message.recipientNickname,
roomKey: roomKeyData
)
} }
} }
} }
@@ -1722,11 +1755,9 @@ class BluetoothMeshService: NSObject {
// Send announce with our nickname immediately // Send announce with our nickname immediately
self.sendAnnouncementToPeer(senderID) self.sendAnnouncementToPeer(senderID)
// Delay sending cached messages to ensure connection is fully established // Send cached messages immediately for faster sync
DispatchQueue.main.asyncAfter(deadline: .now() + 0.5) { [weak self] in // Check if this peer has cached messages (especially for favorites)
// Check if this peer has cached messages (especially for favorites) self.sendCachedMessages(to: senderID)
self?.sendCachedMessages(to: senderID)
}
} }
} }
+23 -5
View File
@@ -12,6 +12,8 @@ import CryptoKit
struct RetryableMessage { struct RetryableMessage {
let id: String let id: String
let originalMessageID: String?
let originalTimestamp: Date?
let content: String let content: String
let mentions: [String]? let mentions: [String]?
let room: String? let room: String?
@@ -29,7 +31,7 @@ class MessageRetryService {
private var retryQueue: [RetryableMessage] = [] private var retryQueue: [RetryableMessage] = []
private var retryTimer: Timer? private var retryTimer: Timer?
private let retryInterval: TimeInterval = 5.0 // Retry every 5 seconds private let retryInterval: TimeInterval = 2.0 // Retry every 2 seconds for faster sync
private let maxQueueSize = 50 private let maxQueueSize = 50
weak var meshService: BluetoothMeshService? weak var meshService: BluetoothMeshService?
@@ -55,7 +57,9 @@ class MessageRetryService {
isPrivate: Bool = false, isPrivate: Bool = false,
recipientPeerID: String? = nil, recipientPeerID: String? = nil,
recipientNickname: String? = nil, recipientNickname: String? = nil,
roomKey: Data? = nil roomKey: Data? = nil,
originalMessageID: String? = nil,
originalTimestamp: Date? = nil
) { ) {
// Don't queue if we're at capacity // Don't queue if we're at capacity
guard retryQueue.count < maxQueueSize else { guard retryQueue.count < maxQueueSize else {
@@ -64,6 +68,8 @@ class MessageRetryService {
let retryMessage = RetryableMessage( let retryMessage = RetryableMessage(
id: UUID().uuidString, id: UUID().uuidString,
originalMessageID: originalMessageID,
originalTimestamp: originalTimestamp,
content: content, content: content,
mentions: mentions, mentions: mentions,
room: room, room: room,
@@ -113,13 +119,16 @@ class MessageRetryService {
meshService.sendPrivateMessage( meshService.sendPrivateMessage(
message.content, message.content,
to: recipientID, to: recipientID,
recipientNickname: message.recipientNickname ?? "unknown" recipientNickname: message.recipientNickname ?? "unknown",
messageID: message.originalMessageID
) )
} else { } else {
// Recipient not connected, keep in queue with updated retry time // Recipient not connected, keep in queue with updated retry time
var updatedMessage = message var updatedMessage = message
updatedMessage = RetryableMessage( updatedMessage = RetryableMessage(
id: message.id, id: message.id,
originalMessageID: message.originalMessageID,
originalTimestamp: message.originalTimestamp,
content: message.content, content: message.content,
mentions: message.mentions, mentions: message.mentions,
room: message.room, room: message.room,
@@ -141,13 +150,17 @@ class MessageRetryService {
message.content, message.content,
mentions: message.mentions ?? [], mentions: message.mentions ?? [],
room: room, room: room,
roomKey: roomKey roomKey: roomKey,
messageID: message.originalMessageID,
timestamp: message.originalTimestamp
) )
} else { } else {
// No peers connected, keep in queue // No peers connected, keep in queue
var updatedMessage = message var updatedMessage = message
updatedMessage = RetryableMessage( updatedMessage = RetryableMessage(
id: message.id, id: message.id,
originalMessageID: message.originalMessageID,
originalTimestamp: message.originalTimestamp,
content: message.content, content: message.content,
mentions: message.mentions, mentions: message.mentions,
room: message.room, room: message.room,
@@ -166,13 +179,18 @@ class MessageRetryService {
meshService.sendMessage( meshService.sendMessage(
message.content, message.content,
mentions: message.mentions ?? [], mentions: message.mentions ?? [],
room: message.room room: message.room,
to: nil,
messageID: message.originalMessageID,
timestamp: message.originalTimestamp
) )
} else { } else {
// No peers connected, keep in queue // No peers connected, keep in queue
var updatedMessage = message var updatedMessage = message
updatedMessage = RetryableMessage( updatedMessage = RetryableMessage(
id: message.id, id: message.id,
originalMessageID: message.originalMessageID,
originalTimestamp: message.originalTimestamp,
content: message.content, content: message.content,
mentions: message.mentions, mentions: message.mentions,
room: message.room, room: message.room,
+54 -7
View File
@@ -2118,8 +2118,32 @@ extension ChatViewModel: BitchatDelegate {
if roomMessages[room] == nil { if roomMessages[room] == nil {
roomMessages[room] = [] roomMessages[room] = []
} }
roomMessages[room]?.append(messageToAdd)
roomMessages[room]?.sort { $0.timestamp < $1.timestamp } // Check if this is our own message being echoed back
if messageToAdd.sender != nickname {
roomMessages[room]?.append(messageToAdd)
roomMessages[room]?.sort { $0.timestamp < $1.timestamp }
} else {
// Our own message - check if we already have it (by ID and content)
let messageExists = roomMessages[room]?.contains { existingMsg in
// Check by ID first
if existingMsg.id == messageToAdd.id {
return true
}
// Check by content and sender with time window (within 1 second)
if existingMsg.content == messageToAdd.content &&
existingMsg.sender == messageToAdd.sender {
let timeDiff = abs(existingMsg.timestamp.timeIntervalSince(messageToAdd.timestamp))
return timeDiff < 1.0
}
return false
} ?? false
if !messageExists {
// This is a message we sent from another device or it's missing locally
roomMessages[room]?.append(messageToAdd)
roomMessages[room]?.sort { $0.timestamp < $1.timestamp }
}
}
// Save message if room has retention enabled // Save message if room has retention enabled
if retentionEnabledRooms.contains(room) { if retentionEnabledRooms.contains(room) {
@@ -2135,8 +2159,8 @@ extension ChatViewModel: BitchatDelegate {
} else { } else {
} }
// Update unread count if not currently viewing this room // Update unread count if not currently viewing this room and it's not our own message
if currentRoom != room { if currentRoom != room && messageToAdd.sender != nickname {
unreadRoomMessages[room] = (unreadRoomMessages[room] ?? 0) + 1 unreadRoomMessages[room] = (unreadRoomMessages[room] ?? 0) + 1
} }
} else { } else {
@@ -2144,9 +2168,32 @@ extension ChatViewModel: BitchatDelegate {
} }
} else { } else {
// Regular public message (main chat) // Regular public message (main chat)
messages.append(message) // Check if this is our own message being echoed back
// Sort messages by timestamp to ensure proper ordering if message.sender != nickname {
messages.sort { $0.timestamp < $1.timestamp } messages.append(message)
// Sort messages by timestamp to ensure proper ordering
messages.sort { $0.timestamp < $1.timestamp }
} else {
// Our own message - check if we already have it (by ID and content)
let messageExists = messages.contains { existingMsg in
// Check by ID first
if existingMsg.id == message.id {
return true
}
// Check by content and sender with time window (within 1 second)
if existingMsg.content == message.content &&
existingMsg.sender == message.sender {
let timeDiff = abs(existingMsg.timestamp.timeIntervalSince(message.timestamp))
return timeDiff < 1.0
}
return false
}
if !messageExists {
// This is a message we sent from another device or it's missing locally
messages.append(message)
messages.sort { $0.timestamp < $1.timestamp }
}
}
} }
// Check if we're mentioned // Check if we're mentioned