fix: Dedupe replayed messages #123

This commit is contained in:
Janez T
2026-04-06 21:00:07 +02:00
parent daa88fb471
commit d92e52c5c2
2 changed files with 197 additions and 27 deletions

View File

@@ -459,23 +459,12 @@ class MessagesProvider with ChangeNotifier {
.loadMessageRouteMetadata(namespace: _storageNamespace); .loadMessageRouteMetadata(namespace: _storageNamespace);
final storedRemovedSarMarkerIds = await _storageService final storedRemovedSarMarkerIds = await _storageService
.loadRemovedSarMarkerIds(namespace: _storageNamespace); .loadRemovedSarMarkerIds(namespace: _storageNamespace);
_messageContactLocations
..clear()
..addAll(storedContactLocations);
_messageReceptionDetails
..clear()
..addAll(storedReceptionDetails);
_messageTransferDetails
..clear()
..addAll(storedTransferDetails);
_messageRouteMetadata
..clear()
..addAll(storedRouteMetadata);
_removedSarMarkerIds _removedSarMarkerIds
..clear() ..clear()
..addAll(storedRemovedSarMarkerIds); ..addAll(storedRemovedSarMarkerIds);
// Add stored messages with enhancement to ensure SAR detection // Rebuild the in-memory list through the normal duplicate rules so
// stale duplicate rows already persisted on-device are collapsed.
for (final message in storedMessages) { for (final message in storedMessages) {
// Re-enhance each message to ensure SAR markers are properly detected // Re-enhance each message to ensure SAR markers are properly detected
// This handles cases where messages were stored before enhancement logic // This handles cases where messages were stored before enhancement logic
@@ -516,7 +505,66 @@ class MessagesProvider with ChangeNotifier {
} }
} }
enhancedMessage = _resolveSenderNameIfNeeded(enhancedMessage);
final duplicateIndex = _findDuplicateMessageIndex(enhancedMessage);
if (duplicateIndex != -1) {
final existingId = _messages[duplicateIndex].id;
final incomingReceptionDetails = storedReceptionDetails[message.id];
final existingMessage = _messages[duplicateIndex];
final duplicatePathBytes = incomingReceptionDetails?.pathBytes;
final storedLocation = storedContactLocations[message.id];
if (storedLocation != null &&
!_messageContactLocations.containsKey(existingId)) {
_messageContactLocations[existingId] = storedLocation;
}
_messageReceptionDetails[existingId] =
MessageReceptionDetails.mergeDuplicate(
existing: _messageReceptionDetails[existingId],
incoming: incomingReceptionDetails,
);
final storedTransferDetailsForMessage = storedTransferDetails[message.id];
if (storedTransferDetailsForMessage != null &&
!_messageTransferDetails.containsKey(existingId)) {
_messageTransferDetails[existingId] = storedTransferDetailsForMessage;
}
final storedRouteMetadataForMessage = storedRouteMetadata[message.id];
if (storedRouteMetadataForMessage != null &&
!_messageRouteMetadata.containsKey(existingId)) {
_messageRouteMetadata[existingId] = storedRouteMetadataForMessage;
}
if (existingMessage.pathBytes == null &&
duplicatePathBytes != null &&
duplicatePathBytes.isNotEmpty) {
_messages[duplicateIndex] = existingMessage.copyWith(
pathBytes: Uint8List.fromList(duplicatePathBytes),
);
}
continue;
}
_messages.add(enhancedMessage); _messages.add(enhancedMessage);
final storedLocation = storedContactLocations[message.id];
if (storedLocation != null) {
_messageContactLocations[message.id] = storedLocation;
}
final storedReception = storedReceptionDetails[message.id];
if (storedReception != null) {
_messageReceptionDetails[message.id] = storedReception;
}
final storedTransfer = storedTransferDetails[message.id];
if (storedTransfer != null) {
_messageTransferDetails[message.id] = storedTransfer;
}
final storedRoute = storedRouteMetadata[message.id];
if (storedRoute != null) {
_messageRouteMetadata[message.id] = storedRoute;
}
// Extract SAR markers // Extract SAR markers
if (enhancedMessage.isSarMarker) { if (enhancedMessage.isSarMarker) {
@@ -830,11 +878,12 @@ class MessagesProvider with ChangeNotifier {
} }
int _findExactDuplicateMessageIndex(Message message) { int _findExactDuplicateMessageIndex(Message message) {
final normalizedIncomingText = _normalizedDuplicateText(message.text);
for (int index = 0; index < _messages.length; index++) { for (int index = 0; index < _messages.length; index++) {
final existing = _messages[index]; final existing = _messages[index];
if (existing.isSentMessage || if (existing.isSentMessage ||
!_matchesExactDuplicateScope(existing, message) || !_matchesExactDuplicateScope(existing, message) ||
existing.text != message.text) { _normalizedDuplicateText(existing.text) != normalizedIncomingText) {
continue; continue;
} }
if (message.isChannelMessage) { if (message.isChannelMessage) {
@@ -851,21 +900,25 @@ class MessagesProvider with ChangeNotifier {
} }
int _findLastConversationDuplicateMessageIndex(Message message) { int _findLastConversationDuplicateMessageIndex(Message message) {
final normalizedIncomingText = _normalizedDuplicateText(message.text);
for (int index = _messages.length - 1; index >= 0; index--) { for (int index = _messages.length - 1; index >= 0; index--) {
final existing = _messages[index]; final existing = _messages[index];
if (!_isSameConversation(existing, message)) { if (!_isSameConversation(existing, message)) {
continue; continue;
} }
if (existing.isSentMessage ||
existing.isSystemMessage ||
existing.text != message.text) {
return -1;
}
final receivedDelta = existing.receivedAt.difference(message.receivedAt).abs(); final receivedDelta = existing.receivedAt.difference(message.receivedAt).abs();
if (receivedDelta > _receivedDuplicateWindow) { if (receivedDelta > _receivedDuplicateWindow) {
return -1; continue;
}
if (existing.isSentMessage || existing.isSystemMessage) {
continue;
}
if (_normalizedDuplicateText(existing.text) != normalizedIncomingText) {
continue;
}
if (_matchesDuplicateSenderIdentity(existing, message)) {
return index;
} }
return _matchesDuplicateSenderIdentity(existing, message) ? index : -1;
} }
return -1; return -1;
@@ -944,11 +997,12 @@ class MessagesProvider with ChangeNotifier {
return -1; return -1;
} }
final normalizedIncomingText = _normalizedDuplicateText(message.text);
for (int index = 0; index < _messages.length; index++) { for (int index = 0; index < _messages.length; index++) {
final existing = _messages[index]; final existing = _messages[index];
if (!existing.isSentMessage || if (!existing.isSentMessage ||
!existing.isChannelMessage || !existing.isChannelMessage ||
existing.text != message.text) { _normalizedDuplicateText(existing.text) != normalizedIncomingText) {
continue; continue;
} }
@@ -1091,6 +1145,10 @@ class MessagesProvider with ChangeNotifier {
return _normalizeSenderName(resolved); return _normalizeSenderName(resolved);
} }
String _normalizedDuplicateText(String value) {
return value.replaceAll('\r\n', '\n').trim();
}
void _scheduleChannelEchoWarning(String messageId) { void _scheduleChannelEchoWarning(String messageId) {
_channelEchoWarningTimers[messageId]?.cancel(); _channelEchoWarningTimers[messageId]?.cancel();
_channelEchoWarningMessageIds.remove(messageId); _channelEchoWarningMessageIds.remove(messageId);
@@ -1398,9 +1456,12 @@ class MessagesProvider with ChangeNotifier {
return false; return false;
} }
return existing.text == message.text && final receivedDelta = existing.receivedAt.difference(message.receivedAt).abs();
return _normalizedDuplicateText(existing.text) ==
_normalizedDuplicateText(message.text) &&
(_matchesExactDuplicateScope(existing, message) || (_matchesExactDuplicateScope(existing, message) ||
(_isSameConversation(existing, message) && (receivedDelta <= _receivedDuplicateWindow &&
_isSameConversation(existing, message) &&
_matchesDuplicateSenderIdentity(existing, message))); _matchesDuplicateSenderIdentity(existing, message)));
} }

View File

@@ -8,6 +8,7 @@ import 'package:meshcore_sar_app/models/message_reception_details.dart';
import 'package:meshcore_sar_app/models/path_selection.dart'; import 'package:meshcore_sar_app/models/path_selection.dart';
import 'package:meshcore_sar_app/providers/helpers/message_retry_manager.dart'; import 'package:meshcore_sar_app/providers/helpers/message_retry_manager.dart';
import 'package:meshcore_sar_app/providers/messages_provider.dart'; import 'package:meshcore_sar_app/providers/messages_provider.dart';
import 'package:meshcore_sar_app/services/message_storage_service.dart';
import 'package:shared_preferences/shared_preferences.dart'; import 'package:shared_preferences/shared_preferences.dart';
Contact _buildContact({ Contact _buildContact({
@@ -462,7 +463,7 @@ void main() {
expect(provider.messages.single.id, equals('handle-1')); expect(provider.messages.single.id, equals('handle-1'));
}); });
test('channel duplicates only dedupe against the latest channel message', () { test('channel duplicates dedupe across recent interleaved channel traffic', () {
final provider = MessagesProvider(); final provider = MessagesProvider();
final sender = Uint8List.fromList([9, 8, 7, 6, 5, 4]); final sender = Uint8List.fromList([9, 8, 7, 6, 5, 4]);
@@ -509,8 +510,9 @@ void main() {
), ),
); );
expect(provider.messages, hasLength(3)); expect(provider.messages, hasLength(2));
expect(provider.messages.last.id, equals('channel-repeat')); expect(provider.messages.first.id, equals('channel-first'));
expect(provider.messages.last.id, equals('channel-middle'));
}); });
test('adjacent channel duplicates still dedupe within 5 seconds', () { test('adjacent channel duplicates still dedupe within 5 seconds', () {
@@ -551,6 +553,62 @@ void main() {
expect(provider.messages.single.id, equals('adjacent-1')); expect(provider.messages.single.id, equals('adjacent-1'));
}); });
test('channel duplicates still dedupe when traffic lands between them', () {
final provider = MessagesProvider();
final sender = Uint8List.fromList([4, 5, 6, 7, 8, 9]);
final baseTime = DateTime.fromMillisecondsSinceEpoch(1700000840000);
provider.addMessage(
Message(
id: 'interleaved-1',
messageType: MessageType.channel,
channelIdx: 3,
pathLen: 4,
textType: MessageTextType.plain,
senderTimestamp: 1700000840,
text: 'Tudi vsm',
senderName: 'Tady (SI)',
receivedAt: baseTime,
senderPublicKeyPrefix: sender,
),
);
provider.addMessage(
Message(
id: 'interleaved-other',
messageType: MessageType.channel,
channelIdx: 3,
pathLen: 1,
textType: MessageTextType.plain,
senderTimestamp: 1700000841,
text: 'Different payload',
senderName: 'MHQ-1 [SI]',
receivedAt: baseTime.add(const Duration(seconds: 1)),
senderPublicKeyPrefix: Uint8List.fromList([1, 2, 3, 4, 5, 6]),
),
);
provider.addMessage(
Message(
id: 'interleaved-2',
messageType: MessageType.channel,
channelIdx: 3,
pathLen: 4,
textType: MessageTextType.plain,
senderTimestamp: 1700000842,
text: 'Tudi vsm',
senderName: 'Tady (SI)',
receivedAt: baseTime.add(const Duration(seconds: 2)),
senderPublicKeyPrefix: sender,
),
);
expect(provider.messages, hasLength(2));
expect(provider.messages.first.id, equals('interleaved-1'));
expect(
provider.getMessageReceptionDetails('interleaved-1')?.receivedCopies,
equals(2),
);
});
test('channel duplicates outside 5 second window are kept', () { test('channel duplicates outside 5 second window are kept', () {
final provider = MessagesProvider(); final provider = MessagesProvider();
final sender = Uint8List.fromList([6, 7, 8, 9, 0, 1]); final sender = Uint8List.fromList([6, 7, 8, 9, 0, 1]);
@@ -653,6 +711,57 @@ void main() {
expect(display.single.message.id, equals('display-sent')); expect(display.single.message.id, equals('display-sent'));
}); });
test('initialize collapses persisted duplicates from storage', () async {
final storage = MessageStorageService();
final sender = Uint8List.fromList([5, 5, 5, 5, 5, 5]);
final firstReceivedAt = DateTime.fromMillisecondsSinceEpoch(1700001000000);
await storage.saveMessages(
[
Message(
id: 'persist-a',
messageType: MessageType.channel,
channelIdx: 3,
pathLen: 4,
textType: MessageTextType.plain,
senderTimestamp: 1700001000,
text: 'Tudi vsm',
senderName: 'Tady (SI)',
receivedAt: firstReceivedAt,
senderPublicKeyPrefix: sender,
),
Message(
id: 'persist-b',
messageType: MessageType.channel,
channelIdx: 3,
pathLen: 4,
textType: MessageTextType.plain,
senderTimestamp: 1700001001,
text: 'Tudi vsm',
senderName: 'Tady (SI)',
receivedAt: firstReceivedAt.add(const Duration(seconds: 2)),
senderPublicKeyPrefix: sender,
),
],
messageReceptionDetails: {
'persist-a': MessageReceptionDetails(capturedAt: firstReceivedAt),
'persist-b': MessageReceptionDetails(
capturedAt: firstReceivedAt.add(const Duration(seconds: 2)),
),
},
);
final provider = MessagesProvider();
await provider.initialize();
expect(provider.messages, hasLength(1));
expect(provider.messages.single.id, equals('persist-a'));
expect(
provider.getMessageReceptionDetails('persist-a')?.receivedCopies,
equals(2),
);
});
test('missing ACK schedules a delayed retransmission', () { test('missing ACK schedules a delayed retransmission', () {
fakeAsync((async) { fakeAsync((async) {
final provider = MessagesProvider(); final provider = MessagesProvider();