diff --git a/lib/connector/meshcore_connector.dart b/lib/connector/meshcore_connector.dart index 73e3dc3..52048b0 100644 --- a/lib/connector/meshcore_connector.dart +++ b/lib/connector/meshcore_connector.dart @@ -7110,8 +7110,15 @@ class MeshCoreConnector extends ChangeNotifier { } final retryService = _retryService; + // Channel sends draw this same frame. While one is awaiting its own + // RESP_CODE_SENT the frame is ambiguous, so the retry service must not + // adopt an unpredicted hash into a direct message (#581). if (retryService != null && - retryService.updateMessageFromSent(ackHash, timeoutMs)) { + retryService.updateMessageFromSent( + ackHash, + timeoutMs, + allowUnpredictedAdoption: _pendingChannelSentQueue.isEmpty, + )) { return; } diff --git a/lib/services/message_retry_service.dart b/lib/services/message_retry_service.dart index 8d048ed..6475834 100644 --- a/lib/services/message_retry_service.dart +++ b/lib/services/message_retry_service.dart @@ -391,7 +391,74 @@ class MessageRetryService extends ChangeNotifier { _onMessageResolved(messageId, contact.publicKeyHex); } - bool updateMessageFromSent(int ackHash, int timeoutMs) { + /// Message ids whose send reached the radio but whose RESP_CODE_SENT has not + /// arrived yet. + List get _sendsAwaitingConfirmation => _sentConfirmationTimers.keys + .where((id) => _pendingMessages[id]?.status == MessageStatus.pending) + .toList(); + + /// The radio is authoritative for the expected-ACK hash it reports in + /// RESP_CODE_SENT. Firmware forks build the message payload differently, so + /// the hash recomputed locally can disagree while the send itself is fine. + /// When that happened the message stayed pending, the 8s watchdog from #395 + /// marked an already-delivered DM failed, and the genuine ACK later matched + /// nothing because every downstream map is only populated on the match path. + /// Captured on a Wadamesh radio in #449. + /// + /// This frame is the reply to our own CMD_SEND_TXT_MSG, so adopt the radio's + /// value when exactly one send is awaiting confirmation. With zero or several + /// candidates the correlation would be a guess, so keep the old behaviour. + /// + /// [allowed] is false when a channel message is also awaiting its + /// RESP_CODE_SENT. Channel sends share this frame, so adopting there could + /// attach a channel message's confirmation to an unrelated direct message. + String? _adoptUnpredictedSentHash( + String ackHashHex, + RetryServiceConfig config, { + required bool allowed, + }) { + if (!allowed) { + config.debugLogService?.warn( + 'RESP_CODE_SENT: ACK hash $ackHashHex matches no pending message and a ' + 'channel send is also awaiting confirmation, not adopting', + tag: 'AckHash', + ); + return null; + } + + final awaiting = _sendsAwaitingConfirmation; + if (awaiting.length != 1) { + config.debugLogService?.warn( + 'RESP_CODE_SENT: ACK hash $ackHashHex matches no pending message and ' + '${awaiting.length} sends are awaiting confirmation, ignoring', + tag: 'AckHash', + ); + return null; + } + + final messageId = awaiting.first; + final message = _pendingMessages[messageId]; + final text = message?.text ?? ''; + final shortText = text.length > 20 ? '${text.substring(0, 20)}...' : text; + config.debugLogService?.warn( + 'RESP_CODE_SENT: ACK hash $ackHashHex is not the hash we predicted, ' + 'adopting the radio value for "$shortText" (the radio is authoritative, ' + 'see #449)', + tag: 'AckHash', + ); + + // Drop the stale prediction so it cannot mis-match a later reply. + _expectedHashToMessageId.removeWhere((_, id) => id == messageId); + return messageId; + } + + /// [allowUnpredictedAdoption] must be false when a channel message is also + /// awaiting its RESP_CODE_SENT, because that frame could belong to either. + bool updateMessageFromSent( + int ackHash, + int timeoutMs, { + bool allowUnpredictedAdoption = true, + }) { final config = _config; if (config == null) return false; @@ -424,8 +491,15 @@ class MessageRetryService extends ChangeNotifier { } if (messageId == null || contact == null) { - debugPrint('No pending message found for ACK hash: $ackHashHex'); - return false; + final adopted = _adoptUnpredictedSentHash( + ackHashHex, + config, + allowed: allowUnpredictedAdoption, + ); + if (adopted == null) return false; + messageId = adopted; + contact = _pendingContacts[adopted]; + if (contact == null) return false; } final message = _pendingMessages[messageId]!; diff --git a/test/services/retry_and_protocol_test.dart b/test/services/retry_and_protocol_test.dart index c957ff8..0860f2d 100644 --- a/test/services/retry_and_protocol_test.dart +++ b/test/services/retry_and_protocol_test.dart @@ -709,4 +709,238 @@ void main() { }, ); }); + + group('RESP_CODE_SENT correlation uses the radio hash (#449/#581)', () { + // The radio is authoritative for the expected-ACK hash. Firmware forks + // (Wadamesh on the HV4 TFT, in the #449 capture) build the payload + // differently, so the client's locally recomputed hash can disagree while + // delivery works perfectly. Correlation must not depend on that guess. + const int radioHash = 0xDEADBEEF; // deliberately not the client's value + + test('a RESP_CODE_SENT hash the client did not predict still marks the ' + 'message sent, and its ACK still marks it delivered', () async { + final retryService = MessageRetryService(); + final contact = _makeContact( + publicKey: recipientKey, + pathLength: 2, + path: const [0x10, 0x20], + ); + final updates = []; + + retryService.initialize( + RetryServiceConfig( + sendMessage: (_, _, _, _) async {}, + addMessage: (_, _) {}, + updateMessage: updates.add, + getSelfPublicKey: () => fixedKey, + ), + ); + + await retryService.sendMessageWithRetry( + contact: contact, + text: 'Weird that I am getting errors though.', + ); + await Future.delayed(const Duration(milliseconds: 20)); + + final matched = retryService.updateMessageFromSent(radioHash, 4884); + + expect( + matched, + isTrue, + reason: + 'the radio replied to our own CMD_SEND_TXT_MSG, so it must be ' + 'correlated even though the hash is not the one we predicted', + ); + expect( + updates.last.status, + equals(MessageStatus.sent), + reason: + 'an unpredicted hash must not leave the message pending, ' + 'because the 8s watchdog would then mark a delivered DM failed', + ); + + // The real ACK carries the radio's hash, not ours. + retryService.handleAckReceived(radioHash, 7354); + + expect( + updates.last.status, + equals(MessageStatus.delivered), + reason: 'the ACK must resolve against the radio-supplied hash', + ); + + retryService.dispose(); + }); + + test('an unpredicted hash is not adopted while two sends are awaiting ' + 'confirmation, because the correlation would be a guess', () async { + final retryService = MessageRetryService(); + final contactA = _makeContact(publicKey: recipientKey, pathLength: 2); + final contactB = _makeContact(publicKey: _makeKey(0x55), pathLength: 2); + final updates = []; + + retryService.initialize( + RetryServiceConfig( + sendMessage: (_, _, _, _) async {}, + addMessage: (_, _) {}, + updateMessage: updates.add, + getSelfPublicKey: () => fixedKey, + ), + ); + + // Per-contact queues, so two contacts means two in flight at once. + await retryService.sendMessageWithRetry(contact: contactA, text: 'a'); + await retryService.sendMessageWithRetry(contact: contactB, text: 'b'); + await Future.delayed(const Duration(milliseconds: 20)); + + expect( + retryService.updateMessageFromSent(radioHash, 4884), + isFalse, + reason: + 'with two candidates the owner would rather keep the old ' + 'behaviour than attach the confirmation to the wrong message', + ); + + retryService.dispose(); + }); + + test('an unpredicted hash is not adopted while a channel send is also ' + 'awaiting confirmation, because that frame could be either', () async { + final retryService = MessageRetryService(); + final contact = _makeContact(publicKey: recipientKey, pathLength: 2); + final updates = []; + + retryService.initialize( + RetryServiceConfig( + sendMessage: (_, _, _, _) async {}, + addMessage: (_, _) {}, + updateMessage: updates.add, + getSelfPublicKey: () => fixedKey, + ), + ); + + await retryService.sendMessageWithRetry(contact: contact, text: 'dm'); + await Future.delayed(const Duration(milliseconds: 20)); + + // The connector passes false while _pendingChannelSentQueue is not + // empty. Adopting here would attach a channel message's confirmation + // to this DM and lose both. + expect( + retryService.updateMessageFromSent( + radioHash, + 4884, + allowUnpredictedAdoption: false, + ), + isFalse, + ); + expect( + updates.every((m) => m.status != MessageStatus.sent), + isTrue, + reason: 'the DM must not be marked sent off an ambiguous frame', + ); + + retryService.dispose(); + }); + + test( + 'an unpredicted hash with nothing awaiting confirmation is ignored', + () { + final retryService = MessageRetryService(); + + retryService.initialize( + RetryServiceConfig( + sendMessage: (_, _, _, _) async {}, + addMessage: (_, _) {}, + updateMessage: (_) {}, + getSelfPublicKey: () => fixedKey, + ), + ); + + expect( + retryService.updateMessageFromSent(radioHash, 4884), + isFalse, + reason: + 'with no send in flight the frame must fall through to the ' + 'channel handler, not be swallowed', + ); + + retryService.dispose(); + }, + ); + + test( + 'a late stray frame is not adopted once the message has resolved', + () async { + final retryService = MessageRetryService(); + final contact = _makeContact(publicKey: recipientKey, pathLength: 2); + final updates = []; + + retryService.initialize( + RetryServiceConfig( + sendMessage: (_, _, _, _) async {}, + addMessage: (_, _) {}, + updateMessage: updates.add, + getSelfPublicKey: () => fixedKey, + ), + ); + + await retryService.sendMessageWithRetry(contact: contact, text: 'one'); + await Future.delayed(const Duration(milliseconds: 20)); + + expect(retryService.updateMessageFromSent(radioHash, 4884), isTrue); + retryService.handleAckReceived(radioHash, 5000); + expect(updates.last.status, equals(MessageStatus.delivered)); + + // A second, unrelated unpredicted frame arrives afterwards. Nothing is + // awaiting confirmation now, so it must not be attached to anything. + expect( + retryService.updateMessageFromSent(0xFEEDFACE, 4884), + isFalse, + reason: 'a stray frame must not be adopted by a resolved message', + ); + + retryService.dispose(); + }, + ); + + test( + 'a hash the client did predict still matches on the fast path', + () async { + final retryService = MessageRetryService(); + final contact = _makeContact(publicKey: recipientKey, pathLength: 2); + final updates = []; + int? sentTs; + int? sentAttempt; + + retryService.initialize( + RetryServiceConfig( + sendMessage: (_, _, attempt, ts) async { + sentAttempt = attempt; + sentTs = ts; + }, + addMessage: (_, _) {}, + updateMessage: updates.add, + getSelfPublicKey: () => fixedKey, + ), + ); + + await retryService.sendMessageWithRetry(contact: contact, text: 'Yep.'); + await Future.delayed(const Duration(milliseconds: 20)); + + final predicted = _manualAckHash( + sentTs!, + sentAttempt! & 0x03, + 'Yep.', + fixedKey, + ); + + expect(retryService.updateMessageFromSent(predicted, 3252), isTrue); + expect(updates.last.status, equals(MessageStatus.sent)); + + retryService.handleAckReceived(predicted, 1200); + expect(updates.last.status, equals(MessageStatus.delivered)); + + retryService.dispose(); + }, + ); + }); }