Add swarm mode transport doc

This commit is contained in:
Janez T
2026-03-07 14:13:54 +01:00
parent 4e76898c8d
commit 0e4f727e26
43 changed files with 3831 additions and 1078 deletions

View File

@@ -22,6 +22,7 @@ import '../utils/drawing_message_parser.dart';
import '../utils/raw_route_probe.dart';
import '../utils/voice_message_parser.dart';
import '../utils/image_message_parser.dart';
import '../utils/media_swarm_protocol.dart';
import '../utils/message_airtime_estimator.dart';
/// Main App Provider - coordinates all other providers
@@ -63,6 +64,7 @@ class AppProvider with ChangeNotifier {
bool get autoAddDiscoveredContacts => _autoAddDiscoveredContacts;
static const Duration _packetRetryDelay = Duration(milliseconds: 1200);
static const Duration _mediaSwarmResponseWindow = Duration(seconds: 10);
static const int _maxPacketRetryAttempts = 4;
final Map<String, String> _voiceSessionSenderKey6 = {};
final Map<String, String> _imageSessionSenderKey6 = {};
@@ -72,6 +74,9 @@ class AppProvider with ChangeNotifier {
final Map<String, int> _imageMissingRetryAttempts = {};
final FragmentAckWaitRegistry _rawProbeWaiters = FragmentAckWaitRegistry();
final Map<String, Future<bool>> _pendingRawRouteProbes = {};
final Map<String, Future<bool>> _pendingMediaSwarmFetches = {};
final Map<String, Map<String, MediaSwarmAvailability>>
_pendingMediaSwarmResponses = {};
Timer? _packetCaptureFlushTimer;
String? _lastPersistedPacketSignature;
bool _isPersistingPacketCapture = false;
@@ -183,12 +188,12 @@ class AppProvider with ChangeNotifier {
void _restoreSessionMetadataFromMessages() {
final restored = restoreSessionMetadataFromMessages(
messagesProvider.messages.map((message) => message.text),
messagesProvider.messages,
);
_voiceSessionSenderKey6.addAll(restored.voiceSenderKeyBySession);
_imageSessionSenderKey6.addAll(restored.imageSenderKeyBySession);
for (final entry in restored.imageEnvelopeBySession.entries) {
_imageSessionSenderKey6[entry.key] = entry.value.senderKey6.toLowerCase();
imageProvider.registerEnvelope(entry.value);
}
@@ -681,9 +686,15 @@ class AppProvider with ChangeNotifier {
// Voice envelope message (new public/direct on-demand format).
final voiceEnvelope = VoiceEnvelope.tryParseText(enrichedMessage.text);
if (voiceEnvelope != null) {
_voiceSessionSenderKey6[voiceEnvelope.sessionId] = voiceEnvelope
.senderKey6
.toLowerCase();
final senderPrefix = enrichedMessage.senderPublicKeyPrefix;
if (senderPrefix != null && senderPrefix.length >= 6) {
_voiceSessionSenderKey6[voiceEnvelope.sessionId] = senderPrefix
.sublist(0, 6)
.map((b) => b.toRadixString(16).padLeft(2, '0'))
.join()
.toLowerCase();
}
voiceProvider.registerEnvelope(voiceEnvelope);
enrichedMessage = enrichedMessage.copyWith(
isVoice: true,
voiceId: voiceEnvelope.sessionId,
@@ -713,9 +724,14 @@ class AppProvider with ChangeNotifier {
// Image envelope (IE1): announce image availability.
final imageEnvelope = ImageEnvelope.tryParse(enrichedMessage.text);
if (imageEnvelope != null) {
_imageSessionSenderKey6[imageEnvelope.sessionId] = imageEnvelope
.senderKey6
.toLowerCase();
final senderPrefix = enrichedMessage.senderPublicKeyPrefix;
if (senderPrefix != null && senderPrefix.length >= 6) {
_imageSessionSenderKey6[imageEnvelope.sessionId] = senderPrefix
.sublist(0, 6)
.map((b) => b.toRadixString(16).padLeft(2, '0'))
.join()
.toLowerCase();
}
imageProvider.registerEnvelope(imageEnvelope);
messagesProvider.addMessage(
enrichedMessage,
@@ -739,19 +755,6 @@ class AppProvider with ChangeNotifier {
return;
}
// If it's a text-format voice packet, feed it to VoiceProvider
if (VoicePacket.isVoiceText(enrichedMessage.text)) {
final pkt = VoicePacket.tryParseText(enrichedMessage.text);
if (pkt != null) {
voiceProvider.addPacket(pkt);
// Mark the message with voice metadata before adding to chat
enrichedMessage = enrichedMessage.copyWith(
isVoice: true,
voiceId: pkt.sessionId,
);
}
}
// Pass contact lookup function to link channel messages with contacts
messagesProvider.addMessage(
enrichedMessage,
@@ -803,8 +806,9 @@ class AppProvider with ChangeNotifier {
};
// When raw binary data is received (PUSH_CODE_RAW_DATA 0x84)
// Magic 0x72 'r' = voice fetch request; 0x69 'i' = image fetch request.
// Magic 0x56 'V' = voice packet; magic 0x49 'I' = image packet.
// Magic 0x6d 'm' = swarm control; 0x72 'r' = voice fetch request.
// Magic 0x69 'i' = image fetch request; 0x56 'V' = voice packet.
// Magic 0x49 'I' = image packet.
connectionProvider.onRawDataReceived = (payload, snrRaw, rssiDbm) {
final rawProbeRequest = RawRouteProbeRequest.tryParseBinary(payload);
if (rawProbeRequest != null) {
@@ -824,6 +828,20 @@ class AppProvider with ChangeNotifier {
return;
}
final mediaSwarmRequest = MediaSwarmRequest.tryParseBinary(payload);
if (mediaSwarmRequest != null) {
_handleIncomingMediaSwarmRequest(mediaSwarmRequest);
return;
}
final mediaSwarmAvailability = MediaSwarmAvailability.tryParseBinary(
payload,
);
if (mediaSwarmAvailability != null) {
_handleIncomingMediaSwarmAvailability(mediaSwarmAvailability);
return;
}
final voiceFetchRequest = VoiceFetchRequest.tryParseBinary(payload);
if (voiceFetchRequest != null) {
debugPrint(
@@ -841,26 +859,34 @@ class AppProvider with ChangeNotifier {
);
return;
}
if (requester.outPathLen > _maxDirectPayloadHops) {
if (requester.routeHopCount > _maxDirectPayloadHops) {
debugPrint(
'⚠️ [AppProvider] Voice fetch requester too far: ${requester.outPathLen} hops',
'⚠️ [AppProvider] Voice fetch requester too far: ${requester.routeHopCount} hops',
);
messagesProvider.logSystemMessage(
text:
'Cannot fetch voice for ${requester.advName}: message is too far (${requester.outPathLen} hops, max $_maxDirectPayloadHops).',
'Cannot fetch voice for ${requester.advName}: message is too far (${requester.routeHopCount} hops, max $_maxDirectPayloadHops).',
level: 'warning',
);
return;
}
unawaited(
voiceProvider.serveSessionTo(
unawaited(() async {
final served = await voiceProvider.serveSessionTo(
sessionId: voiceFetchRequest.sessionId,
requester: requester,
requestedIndices: voiceFetchRequest.want == 'missing'
? voiceFetchRequest.missingIndices.toSet()
: null,
),
);
);
if (served) {
messagesProvider.recordMediaTransfer(
sessionId: voiceFetchRequest.sessionId,
mediaType: 'voice',
requesterKey6: voiceFetchRequest.requesterKey6,
requesterName: requester.advName,
);
}
}());
return;
}
@@ -880,32 +906,40 @@ class AppProvider with ChangeNotifier {
);
return;
}
if (requester.outPathLen > _maxDirectPayloadHops) {
if (requester.routeHopCount > _maxDirectPayloadHops) {
debugPrint(
'⚠️ [AppProvider] Image fetch requester too far: '
'${requester.outPathLen} hops for session '
'${requester.routeHopCount} hops for session '
'${imageFetchRequest.sessionId}',
);
messagesProvider.logSystemMessage(
text:
'Cannot fetch image for ${requester.advName}: message is too far (${requester.outPathLen} hops, max $_maxDirectPayloadHops).',
'Cannot fetch image for ${requester.advName}: message is too far (${requester.routeHopCount} hops, max $_maxDirectPayloadHops).',
level: 'warning',
);
return;
}
debugPrint(
'📷 [AppProvider] Serving image session ${imageFetchRequest.sessionId} '
'to ${requester.advName} via ${requester.outPathLen} hop(s)',
'to ${requester.advName} via ${requester.routeHopCount} hop(s)',
);
unawaited(
imageProvider.serveSessionTo(
unawaited(() async {
final served = await imageProvider.serveSessionTo(
sessionId: imageFetchRequest.sessionId,
requester: requester,
requestedIndices: imageFetchRequest.want == 'missing'
? imageFetchRequest.missingIndices.toSet()
: null,
),
);
);
if (served) {
messagesProvider.recordMediaTransfer(
sessionId: imageFetchRequest.sessionId,
mediaType: 'image',
requesterKey6: imageFetchRequest.requesterKey6,
requesterName: requester.advName,
);
}
}());
return;
}
@@ -914,8 +948,23 @@ class AppProvider with ChangeNotifier {
if (frag == null) return;
debugPrint('📷 [AppProvider] Binary image fragment received: $frag');
final session = imageProvider.session(frag.sessionId);
if (session == null && frag.total < 1) {
debugPrint(
'⚠️ [AppProvider] Dropping compact image fragment without envelope '
'for session ${frag.sessionId}',
);
return;
}
imageProvider.addFragment(
frag,
session == null
? frag
: ImagePacket(
sessionId: frag.sessionId,
format: session.format,
index: frag.index,
total: session.total,
data: frag.data,
),
width: session?.width ?? 0,
height: session?.height ?? 0,
);
@@ -930,7 +979,25 @@ class AppProvider with ChangeNotifier {
final pkt = VoicePacket.tryParseBinary(payload);
if (pkt == null) return;
debugPrint('🎙️ [AppProvider] Binary voice packet received: $pkt');
final justComplete = voiceProvider.addPacket(pkt);
final session = voiceProvider.session(pkt.sessionId);
if (session == null && pkt.total < 1) {
debugPrint(
'⚠️ [AppProvider] Dropping compact voice packet without envelope '
'for session ${pkt.sessionId}',
);
return;
}
final justComplete = voiceProvider.addPacket(
session == null
? pkt
: VoicePacket(
sessionId: pkt.sessionId,
mode: session.mode,
index: pkt.index,
total: session.total,
codec2Data: pkt.codec2Data,
),
);
_scheduleVoiceMissingRetry(pkt.sessionId, justComplete: justComplete);
// Insert or update the placeholder message in the chat list
_handleIncomingVoicePacket(pkt, justComplete: justComplete);
@@ -1362,6 +1429,291 @@ class AppProvider with ChangeNotifier {
return null;
}
String _mediaSwarmKey(String mediaType, String sessionId) =>
'$mediaType:$sessionId';
String? _deviceKey6Hex() {
final deviceKey = connectionProvider.deviceInfo.publicKey;
if (deviceKey == null || deviceKey.length < 6) {
return null;
}
return deviceKey
.sublist(0, 6)
.map((b) => b.toRadixString(16).padLeft(2, '0'))
.join('');
}
List<int> _availableIndicesForSession(String mediaType, String sessionId) {
return switch (mediaType) {
'voice' => voiceProvider.availablePacketIndices(sessionId),
'image' => imageProvider.availableFragmentIndices(sessionId),
_ => const <int>[],
};
}
List<int> _matchingAvailableIndices(MediaSwarmRequest request) {
final available = _availableIndicesForSession(
request.mediaType,
request.sessionId,
);
if (available.isEmpty) return const [];
if (request.requestsAll) return available;
final requested = request.missingIndices.toSet();
return available.where(requested.contains).toList()..sort();
}
List<Contact> _eligibleSwarmPeers({String? excludeKey6}) {
final ownKey6 = _deviceKey6Hex();
return contactsProvider.contacts.where((contact) {
if (!contact.routeHasPath ||
contact.routeHopCount > _maxDirectPayloadHops ||
!contact.routeSupportsLegacyRawTransport ||
contact.outPath.isEmpty ||
contact.publicKey.length < 6) {
return false;
}
final key6 = contact.publicKey
.sublist(0, 6)
.map((b) => b.toRadixString(16).padLeft(2, '0'))
.join();
if (key6 == ownKey6 || key6 == excludeKey6) {
return false;
}
return true;
}).toList();
}
void _handleIncomingMediaSwarmRequest(MediaSwarmRequest request) {
final ownKey6 = _deviceKey6Hex();
if (ownKey6 == null || request.requesterKey6 == ownKey6) {
return;
}
final available = _matchingAvailableIndices(request);
if (available.isEmpty) {
return;
}
final availability = MediaSwarmAvailability(
mediaType: request.mediaType,
sessionId: request.sessionId,
requesterKey6: request.requesterKey6,
responderKey6: ownKey6,
availableIndices: available,
);
debugPrint(
'🌐 [AppProvider] Media swarm availability for ${request.mediaType} '
'${request.sessionId}: ${available.length} fragment(s)',
);
final requester = _resolveContactByPrefixHex(request.requesterKey6);
if (requester == null ||
!requester.routeHasPath ||
requester.routeHopCount > _maxDirectPayloadHops ||
!requester.routeSupportsLegacyRawTransport ||
requester.outPath.isEmpty) {
return;
}
unawaited(
connectionProvider.sendRawVoicePacket(
contactPath: requester.outPath,
contactPathLen: requester.routeSignedPathLen,
payload: availability.encodeBinary(),
),
);
}
void _handleIncomingMediaSwarmAvailability(
MediaSwarmAvailability availability,
) {
final ownKey6 = _deviceKey6Hex();
if (ownKey6 == null || availability.requesterKey6 != ownKey6) {
return;
}
final key = _mediaSwarmKey(availability.mediaType, availability.sessionId);
final responses = _pendingMediaSwarmResponses[key];
if (responses == null) {
return;
}
responses[availability.responderKey6] = availability;
debugPrint(
'🌐 [AppProvider] Media swarm response for ${availability.mediaType} '
'${availability.sessionId} from ${availability.responderKey6} '
'(${availability.servesAll ? 'all' : availability.availableIndices.length})',
);
}
Future<bool> _requestMissingMediaViaSwarm({
required String mediaType,
required String sessionId,
required List<int> missingIndices,
required String? originalSenderKey6,
}) async {
if (!connectionProvider.deviceInfo.isConnected || missingIndices.isEmpty) {
return false;
}
final key = _mediaSwarmKey(mediaType, sessionId);
final pending = _pendingMediaSwarmFetches[key];
if (pending != null) {
return pending;
}
final requesterKey6 = _deviceKey6Hex();
if (requesterKey6 == null) {
return false;
}
final future = () async {
final responses = <String, MediaSwarmAvailability>{};
_pendingMediaSwarmResponses[key] = responses;
try {
final request = MediaSwarmRequest(
mediaType: mediaType,
sessionId: sessionId,
requesterKey6: requesterKey6,
missingIndices: missingIndices,
);
final peers = _eligibleSwarmPeers(excludeKey6: originalSenderKey6);
if (peers.isEmpty) {
return false;
}
debugPrint(
'🌐 [AppProvider] Media swarm request for $mediaType $sessionId '
'(${missingIndices.length} needed fragment(s), ${peers.length} peer(s))',
);
for (final peer in peers) {
await connectionProvider.sendRawVoicePacket(
contactPath: peer.outPath,
contactPathLen: peer.routeSignedPathLen,
payload: request.encodeBinary(),
);
}
await Future<void>.delayed(_mediaSwarmResponseWindow);
final orderedResponses =
responses.values
.where(
(response) => response.responderKey6 != originalSenderKey6,
)
.toList()
..sort((a, b) {
final aScore = _swarmResponseScore(a, missingIndices);
final bScore = _swarmResponseScore(b, missingIndices);
return bScore.compareTo(aScore);
});
for (final response in orderedResponses) {
final responder = _resolveContactByPrefixHex(response.responderKey6);
if (responder == null ||
!responder.routeHasPath ||
responder.routeHopCount > _maxDirectPayloadHops ||
!responder.routeSupportsLegacyRawTransport ||
responder.outPath.isEmpty) {
continue;
}
final requestedSubset = response.servesAll
? missingIndices
: missingIndices
.where(response.availableIndices.toSet().contains)
.toList();
if (requestedSubset.isEmpty) {
continue;
}
final requestedSet = requestedSubset.toSet();
final sent = await _sendDirectMediaFetchRequest(
mediaType: mediaType,
sessionId: sessionId,
target: responder,
requesterKey6: requesterKey6,
missingIndices: requestedSet,
);
if (sent) {
debugPrint(
'🌐 [AppProvider] Requested $mediaType $sessionId '
'from swarm peer ${responder.advName} '
'(${requestedSubset.length} fragment(s))',
);
return true;
}
}
return false;
} catch (e) {
debugPrint(
'⚠️ [AppProvider] Media swarm request failed for $mediaType '
'$sessionId: $e',
);
return false;
} finally {
_pendingMediaSwarmResponses.remove(key);
}
}();
_pendingMediaSwarmFetches[key] = future;
try {
return await future;
} finally {
_pendingMediaSwarmFetches.remove(key);
}
}
int _swarmResponseScore(
MediaSwarmAvailability response,
List<int> missingIndices,
) {
if (response.servesAll) {
return missingIndices.length;
}
final needed = missingIndices.toSet();
return response.availableIndices.where(needed.contains).length;
}
Future<bool> _sendDirectMediaFetchRequest({
required String mediaType,
required String sessionId,
required Contact target,
required String requesterKey6,
required Set<int> missingIndices,
}) async {
try {
final payload = switch (mediaType) {
'voice' => VoiceFetchRequest(
sessionId: sessionId,
want: missingIndices.isEmpty ? 'all' : 'missing',
missingIndices: missingIndices.toList()..sort(),
requesterKey6: requesterKey6,
).encodeBinary(),
'image' => ImageFetchRequest(
sessionId: sessionId,
want: missingIndices.isEmpty ? 'all' : 'missing',
missingIndices: missingIndices.toList()..sort(),
requesterKey6: requesterKey6,
).encodeBinary(),
_ => null,
};
if (payload == null) {
return false;
}
await connectionProvider.sendRawVoicePacket(
contactPath: target.outPath,
contactPathLen: target.routeSignedPathLen,
payload: payload,
);
return true;
} catch (e) {
debugPrint(
'⚠️ [AppProvider] Direct $mediaType fetch via ${target.advName} failed: $e',
);
return false;
}
}
void _scheduleVoiceMissingRetry(
String sessionId, {
required bool justComplete,
@@ -1434,8 +1786,8 @@ class AppProvider with ChangeNotifier {
final senderKey6 = _voiceSessionSenderKey6[sessionId];
if (senderKey6 == null) return;
final sender = _resolveContactByPrefixHex(senderKey6);
final deviceKey = connectionProvider.deviceInfo.publicKey;
if (sender == null || deviceKey == null || deviceKey.length < 6) return;
final requesterKey6 = _deviceKey6Hex();
if (requesterKey6 == null) return;
final missing = voiceProvider.missingPacketIndices(sessionId);
if (missing.isEmpty) {
@@ -1443,27 +1795,28 @@ class AppProvider with ChangeNotifier {
return;
}
final requesterKey6 = deviceKey
.sublist(0, 6)
.map((b) => b.toRadixString(16).padLeft(2, '0'))
.join('');
final request = VoiceFetchRequest(
sessionId: sessionId,
want: 'missing',
missingIndices: missing,
requesterKey6: requesterKey6,
timestampSec: DateTime.now().millisecondsSinceEpoch ~/ 1000,
version: 2,
);
try {
await connectionProvider.sendRawVoicePacket(
contactPath: sender.outPath,
contactPathLen: sender.outPathLen,
payload: request.encodeBinary(),
var sent = false;
if (sender != null) {
final routeOk = await verifyRawTransportRoute(sender);
if (routeOk) {
sent = await _sendDirectMediaFetchRequest(
mediaType: 'voice',
sessionId: sessionId,
target: sender,
requesterKey6: requesterKey6,
missingIndices: missing.toSet(),
);
}
}
if (!sent) {
sent = await _requestMissingMediaViaSwarm(
mediaType: 'voice',
sessionId: sessionId,
missingIndices: missing,
originalSenderKey6: senderKey6,
);
} catch (_) {
}
if (!sent) {
return;
}
@@ -1496,8 +1849,8 @@ class AppProvider with ChangeNotifier {
final senderKey6 = _imageSessionSenderKey6[sessionId];
if (senderKey6 == null) return;
final sender = _resolveContactByPrefixHex(senderKey6);
final deviceKey = connectionProvider.deviceInfo.publicKey;
if (sender == null || deviceKey == null || deviceKey.length < 6) return;
final requesterKey6 = _deviceKey6Hex();
if (requesterKey6 == null) return;
final missing = imageProvider.missingFragmentIndices(sessionId);
if (missing.isEmpty) {
@@ -1505,26 +1858,28 @@ class AppProvider with ChangeNotifier {
return;
}
final requesterKey6 = deviceKey
.sublist(0, 6)
.map((b) => b.toRadixString(16).padLeft(2, '0'))
.join('');
final request = ImageFetchRequest(
sessionId: sessionId,
want: 'missing',
missingIndices: missing,
requesterKey6: requesterKey6,
timestampSec: DateTime.now().millisecondsSinceEpoch ~/ 1000,
);
try {
await connectionProvider.sendRawVoicePacket(
contactPath: sender.outPath,
contactPathLen: sender.outPathLen,
payload: request.encodeBinary(),
var sent = false;
if (sender != null) {
final routeOk = await verifyRawTransportRoute(sender);
if (routeOk) {
sent = await _sendDirectMediaFetchRequest(
mediaType: 'image',
sessionId: sessionId,
target: sender,
requesterKey6: requesterKey6,
missingIndices: missing.toSet(),
);
}
}
if (!sent) {
sent = await _requestMissingMediaViaSwarm(
mediaType: 'image',
sessionId: sessionId,
missingIndices: missing,
originalSenderKey6: senderKey6,
);
} catch (_) {
}
if (!sent) {
return;
}
@@ -1585,7 +1940,10 @@ class AppProvider with ChangeNotifier {
if (!connectionProvider.deviceInfo.isConnected) {
return false;
}
if (target.outPathLen < 0 || target.outPathLen > _maxDirectPayloadHops) {
if (!target.routeHasPath || target.routeHopCount > _maxDirectPayloadHops) {
return false;
}
if (!target.routeSupportsLegacyRawTransport) {
return false;
}
if (target.outPath.isEmpty) {
@@ -1616,11 +1974,11 @@ class AppProvider with ChangeNotifier {
try {
debugPrint(
'📡 [AppProvider] Outgoing raw route probe: target=${target.advName} hops=${target.outPathLen} nonce=${nonce.toRadixString(16)}',
'📡 [AppProvider] Outgoing raw route probe: target=${target.advName} hops=${target.routeHopCount} nonce=${nonce.toRadixString(16)}',
);
await connectionProvider.sendRawVoicePacket(
contactPath: target.outPath,
contactPathLen: target.outPathLen,
contactPathLen: target.routeSignedPathLen,
payload: RawRouteProbeRequest(
nonce: nonce,
requesterKey6: requesterKey6,
@@ -1648,7 +2006,7 @@ class AppProvider with ChangeNotifier {
if (target.publicKeyHex.isNotEmpty) {
return 'pk:${target.publicKeyHex}';
}
return 'name:${target.advName}:${target.outPathLen}:${target.outPath.map((b) => b.toRadixString(16).padLeft(2, '0')).join()}';
return 'name:${target.advName}:${target.routeSignedPathLen}:${target.outPath.map((b) => b.toRadixString(16).padLeft(2, '0')).join()}';
}
void _handleRawRouteProbeRequest(RawRouteProbeRequest request) {
@@ -1659,23 +2017,26 @@ class AppProvider with ChangeNotifier {
);
return;
}
if (requester.outPathLen < 0 ||
requester.outPathLen > _maxDirectPayloadHops) {
if (!requester.routeHasPath ||
requester.routeHopCount > _maxDirectPayloadHops) {
debugPrint(
'⚠️ [AppProvider] Raw route probe requester out of range: ${requester.outPathLen}',
'⚠️ [AppProvider] Raw route probe requester out of range: ${requester.routeHopCount}',
);
return;
}
if (!requester.routeSupportsLegacyRawTransport) {
return;
}
if (requester.outPath.isEmpty) {
return;
}
debugPrint(
'📡 [AppProvider] Outgoing raw route probe ACK: requester=${requester.advName} hops=${requester.outPathLen} nonce=${request.nonce.toRadixString(16)}',
'📡 [AppProvider] Outgoing raw route probe ACK: requester=${requester.advName} hops=${requester.routeHopCount} nonce=${request.nonce.toRadixString(16)}',
);
unawaited(
connectionProvider.sendRawVoicePacket(
contactPath: requester.outPath,
contactPathLen: requester.outPathLen,
contactPathLen: requester.routeSignedPathLen,
payload: RawRouteProbeAck(nonce: request.nonce).encodeBinary(),
),
);

View File

@@ -4,10 +4,11 @@ import 'package:flutter/foundation.dart';
import 'package:flutter/scheduler.dart';
import 'package:flutter_blue_plus/flutter_blue_plus.dart';
import 'package:crypto/crypto.dart';
import '../models/contact.dart';
import '../models/device_info.dart';
import '../models/room_login_state.dart';
import '../models/sse_server_config.dart';
import 'package:meshcore_client/meshcore_client.dart';
import 'package:meshcore_client/meshcore_client.dart' hide Contact;
import '../services/sse_server_service.dart';
import '../utils/sar_message_parser.dart';
import 'helpers/room_login_manager.dart';
@@ -1121,9 +1122,11 @@ class ConnectionProvider with ChangeNotifier {
);
}
debugPrint(' Type: ${contact.type.displayName}');
debugPrint(' Path status: ${contact.pathDescription}');
if (contact.hasPath) {
debugPrint(' ✅ Using learned path (${contact.outPathLen} bytes)');
debugPrint(' Path status: ${contact.routeSummary}');
if (contact.routeHasPath) {
debugPrint(
' ✅ Using learned path (${contact.routeHopCount} hop(s), ${contact.routeHashSize}-byte hashes)',
);
} else {
debugPrint(' ⚠️ No path available - will use flood mode');
}
@@ -1173,6 +1176,19 @@ class ConnectionProvider with ChangeNotifier {
attempt: retryAttempt,
);
if (messageId != null) {
Future.delayed(const Duration(milliseconds: 350), () {
if (_messageDeliveryTracker.hasAckForMessage(messageId)) {
return;
}
debugPrint(
' [ConnectionProvider] Missing RESP_CODE_SENT for $messageId; promoting to sent via fallback',
);
onMessageSent?.call(messageId, 0, 0);
});
}
// Clear pending operation after successful send (no error)
// If ERR_CODE_NOT_FOUND occurs, the operation will be recovered automatically
if (contact != null) {
@@ -2006,6 +2022,30 @@ class ConnectionProvider with ChangeNotifier {
}
}
Future<void> setContactRoute(
Contact contact, {
required int signedEncodedPathLen,
required Uint8List paddedPathBytes,
}) async {
if (!_activeService.isConnected) {
_error = 'Not connected to device';
notifyListeners();
return;
}
try {
final updatedContact = contact.copyWith(
outPathLen: signedEncodedPathLen,
outPath: Uint8List.fromList(paddedPathBytes),
);
await _activeService.addOrUpdateContact(updatedContact);
} catch (e) {
_error = 'Failed to set route: $e';
notifyListeners();
rethrow;
}
}
/// Remove a contact from the companion radio
///
/// Deletes the contact from the device's internal contact table.

View File

@@ -521,7 +521,39 @@ class ContactsProvider with ChangeNotifier {
/// prefer flood routing until the radio reports a fresh route.
void markPathUnhealthy(Uint8List publicKey) {
final contact = findContactByKey(publicKey);
if (contact == null || !contact.hasPath) {
if (contact == null || !contact.routeHasPath) {
return;
}
_contacts[contact.publicKeyHex] = contact.copyWith(
outPathLen: -1,
outPath: Uint8List(0),
);
_persistContacts();
notifyListeners();
}
void setContactRouteLocal(
Uint8List publicKey, {
required int signedEncodedPathLen,
required Uint8List paddedPathBytes,
}) {
final contact = findContactByKey(publicKey);
if (contact == null) {
return;
}
_contacts[contact.publicKeyHex] = contact.copyWith(
outPathLen: signedEncodedPathLen,
outPath: Uint8List.fromList(paddedPathBytes),
);
_persistContacts();
notifyListeners();
}
void resetContactRouteLocal(Uint8List publicKey) {
final contact = findContactByKey(publicKey);
if (contact == null) {
return;
}

View File

@@ -98,6 +98,11 @@ class MessageDeliveryTracker {
return _ackTagToMessageId[ackCode];
}
/// Returns true once a message has been matched to a concrete ACK tag.
bool hasAckForMessage(String messageId) {
return _messageIdToAckTag.containsKey(messageId);
}
/// Remove ACK tag mapping after delivery confirmed or timeout
///
/// Cleans up both forward and reverse mappings.

View File

@@ -1,5 +1,4 @@
import 'dart:convert';
import 'dart:math' as math;
import '../../models/message.dart';
import '../../models/contact.dart';
@@ -55,8 +54,8 @@ class MessageRetryManager {
final payloadBytes = utf8.encode(text).length;
final airtimeMs = _estimateLoRaAirtimeMs(payloadBytes);
final hopCount = contact?.hasPath == true
? math.max(contact!.outPathLen, 0)
final hopCount = contact?.routeHasPath == true
? contact!.routeHopCount
: -1;
if (hopCount < 0) {
@@ -87,7 +86,7 @@ class MessageRetryManager {
// Only retry if contact has a learned path
// If no path, the device uses flood mode automatically - retrying won't help
return contact.hasPath;
return contact.routeHasPath;
}
/// Check if should fall back to flood mode
@@ -101,7 +100,7 @@ class MessageRetryManager {
/// Contacts without paths already use flood mode automatically.
bool shouldUseFloodFallback(Message message, Contact contact) {
return message.retryAttempt >= 3 &&
contact.hasPath && // ✅ FIXED: Flood fallback for failed direct paths
contact.routeHasPath &&
!message.usedFloodFallback;
}

View File

@@ -28,13 +28,19 @@ Future<bool> serveCachedSessionFragments<T>({
debugPrint('⚠️ [$providerLabel] sendRawPacketCallback not set');
return false;
}
if (requester.outPathLen < 0) {
if (!requester.routeHasPath) {
debugPrint('⚠️ [$providerLabel] ${requester.advName} has no direct path');
return false;
}
if (requester.outPathLen > maxDirectPayloadHops) {
if (requester.routeHopCount > maxDirectPayloadHops) {
debugPrint(
'⚠️ [$providerLabel] ${requester.advName} is too far: ${requester.outPathLen} hops (max $maxDirectPayloadHops)',
'⚠️ [$providerLabel] ${requester.advName} is too far: ${requester.routeHopCount} hops (max $maxDirectPayloadHops)',
);
return false;
}
if (!requester.routeSupportsLegacyRawTransport) {
debugPrint(
'⚠️ [$providerLabel] ${requester.advName} route uses unsupported 3-byte raw transport on current client',
);
return false;
}
@@ -58,7 +64,7 @@ Future<bool> serveCachedSessionFragments<T>({
try {
await sendRawPacket(
contactPath: requester.outPath,
contactPathLen: requester.outPathLen,
contactPathLen: requester.routeSignedPathLen,
payload: encodeBinary(fragment),
);
servedCount++;

View File

@@ -1,38 +1,57 @@
import '../../models/message.dart';
import '../../utils/image_message_parser.dart';
import '../../utils/voice_message_parser.dart';
class RestoredSessionMetadata {
final Map<String, String> voiceSenderKeyBySession;
final Map<String, String> imageSenderKeyBySession;
final Map<String, ImageEnvelope> imageEnvelopeBySession;
const RestoredSessionMetadata({
required this.voiceSenderKeyBySession,
required this.imageSenderKeyBySession,
required this.imageEnvelopeBySession,
});
}
RestoredSessionMetadata restoreSessionMetadataFromMessages(
Iterable<String> messageTexts,
Iterable<Message> messages,
) {
final voiceSenderKeyBySession = <String, String>{};
final imageSenderKeyBySession = <String, String>{};
final imageEnvelopeBySession = <String, ImageEnvelope>{};
for (final text in messageTexts) {
for (final message in messages) {
final text = message.text;
final voiceEnvelope = VoiceEnvelope.tryParseText(text);
if (voiceEnvelope != null) {
voiceSenderKeyBySession[voiceEnvelope.sessionId] = voiceEnvelope
.senderKey6
.toLowerCase();
final senderPrefix = message.senderPublicKeyPrefix;
if (senderPrefix != null && senderPrefix.length >= 6) {
voiceSenderKeyBySession[voiceEnvelope.sessionId] = senderPrefix
.sublist(0, 6)
.map((b) => b.toRadixString(16).padLeft(2, '0'))
.join()
.toLowerCase();
}
}
final imageEnvelope = ImageEnvelope.tryParse(text);
if (imageEnvelope != null) {
imageEnvelopeBySession[imageEnvelope.sessionId] = imageEnvelope;
final senderPrefix = message.senderPublicKeyPrefix;
if (senderPrefix != null && senderPrefix.length >= 6) {
imageSenderKeyBySession[imageEnvelope.sessionId] = senderPrefix
.sublist(0, 6)
.map((b) => b.toRadixString(16).padLeft(2, '0'))
.join()
.toLowerCase();
}
}
}
return RestoredSessionMetadata(
voiceSenderKeyBySession: voiceSenderKeyBySession,
imageSenderKeyBySession: imageSenderKeyBySession,
imageEnvelopeBySession: imageEnvelopeBySession,
);
}

View File

@@ -96,11 +96,28 @@ class ImageProvider with ChangeNotifier {
return missing;
}
List<int> availableFragmentIndices(String sessionId) {
final outgoing = _outgoing[sessionId];
if (outgoing != null) {
return outgoing.fragments.map((fragment) => fragment.index).toList()
..sort();
}
final session = _sessions[sessionId];
if (session == null) return const [];
final indices = <int>[];
for (var i = 0; i < session.fragments.length; i++) {
if (session.fragments[i] != null) {
indices.add(i);
}
}
return indices;
}
// ── Incoming fragment reception ──────────────────────────────────────────
/// Add a received [fragment]. Creates the session on first fragment using
/// metadata from the fragment itself (requires envelope to have been
/// announced first; if not, defaults width/height to 0 — corrected on save).
/// Add a received [fragment]. New compact fragments rely on prior envelope
/// metadata for total/format, while legacy fragments can still self-describe.
///
/// Returns true when the session just became complete.
bool addFragment(ImagePacket fragment, {int width = 0, int height = 0}) {
@@ -110,16 +127,20 @@ class ImageProvider with ChangeNotifier {
);
return false;
}
_sessions.putIfAbsent(
fragment.sessionId,
() => ImageSession(
_sessions.putIfAbsent(fragment.sessionId, () {
if (fragment.total < 1) {
throw StateError(
'Image envelope missing for compact fragment ${fragment.sessionId}',
);
}
return ImageSession(
sessionId: fragment.sessionId,
format: fragment.format,
total: fragment.total,
width: width,
height: height,
),
);
);
});
final session = _sessions[fragment.sessionId]!;
if (fragment.index < session.total) {
@@ -244,16 +265,22 @@ class ImageProvider with ChangeNotifier {
required Contact requester,
Set<int>? requestedIndices,
}) async {
final cached = _outgoing[sessionId];
if (cached == null) {
debugPrint('⚠️ [ImageProvider] No cached session for $sessionId');
final outgoing = _outgoing[sessionId];
final fragments = outgoing != null
? List<ImagePacket>.from(outgoing.fragments)
: _sessions[sessionId]?.fragments.whereType<ImagePacket>().toList() ??
const <ImagePacket>[];
if (fragments.isEmpty) {
debugPrint(
'⚠️ [ImageProvider] No cached or received session for $sessionId',
);
return false;
}
return serveCachedSessionFragments<ImagePacket>(
providerLabel: 'ImageProvider',
sessionId: sessionId,
requester: requester,
fragments: cached.fragments,
fragments: fragments,
maxDirectPayloadHops: maxDirectPayloadHops,
indexOf: (fragment) => fragment.index,
encodeBinary: (fragment) => fragment.encodeBinary(),

View File

@@ -4,6 +4,7 @@ import '../models/message.dart';
import '../models/contact.dart';
import '../models/message_contact_location.dart';
import '../models/message_reception_details.dart';
import '../models/message_transfer_details.dart';
import '../models/sar_marker.dart';
import '../models/map_drawing.dart';
import '../services/message_storage_service.dart';
@@ -25,6 +26,7 @@ class MessagesProvider with ChangeNotifier {
AppLocalizations? _localizations;
final Map<String, MessageContactLocation> _messageContactLocations = {};
final Map<String, MessageReceptionDetails> _messageReceptionDetails = {};
final Map<String, MessageTransferDetails> _messageTransferDetails = {};
// Track pending sent messages by expected ACK/TAG
final Map<int, Message> _pendingSentMessages = {};
@@ -121,6 +123,9 @@ class MessagesProvider with ChangeNotifier {
MessageReceptionDetails? getMessageReceptionDetails(String messageId) =>
_messageReceptionDetails[messageId];
MessageTransferDetails? getMessageTransferDetails(String messageId) =>
_messageTransferDetails[messageId];
/// Set localizations for notifications
void setLocalizations(AppLocalizations localizations) {
_localizations = localizations;
@@ -153,12 +158,17 @@ class MessagesProvider with ChangeNotifier {
.loadMessageContactLocations();
final storedReceptionDetails = await _storageService
.loadMessageReceptionDetails();
final storedTransferDetails = await _storageService
.loadMessageTransferDetails();
_messageContactLocations
..clear()
..addAll(storedContactLocations);
_messageReceptionDetails
..clear()
..addAll(storedReceptionDetails);
_messageTransferDetails
..clear()
..addAll(storedTransferDetails);
// Add stored messages with enhancement to ensure SAR detection
for (final message in storedMessages) {
@@ -198,14 +208,6 @@ class MessagesProvider with ChangeNotifier {
isVoice: true,
voiceId: envelope.sessionId,
);
} else if (VoicePacket.isVoiceText(enhancedMessage.text)) {
final pkt = VoicePacket.tryParseText(enhancedMessage.text);
if (pkt != null) {
enhancedMessage = enhancedMessage.copyWith(
isVoice: true,
voiceId: pkt.sessionId,
);
}
}
}
@@ -353,14 +355,6 @@ class MessagesProvider with ChangeNotifier {
isVoice: true,
voiceId: envelope.sessionId,
);
} else if (VoicePacket.isVoiceText(enhancedMessage.text)) {
final pkt = VoicePacket.tryParseText(enhancedMessage.text);
if (pkt != null) {
enhancedMessage = enhancedMessage.copyWith(
isVoice: true,
voiceId: pkt.sessionId,
);
}
}
}
@@ -649,18 +643,16 @@ class MessagesProvider with ChangeNotifier {
final voiceEnvelope = VoiceEnvelope.tryParseText(message.text);
if (voiceEnvelope != null) {
final seconds = (voiceEnvelope.durationMs / 1000).ceil();
final route = isChannelMessage
? 'Channel: ${channelName ?? _resolveChannelName(message.channelIdx)}'
: 'From: $senderName';
return '$route\nVoice message - ${voiceEnvelope.mode.label} - ${seconds}s - ${voiceEnvelope.total} packets';
final summary =
'Voice message - ${voiceEnvelope.mode.label} - ${seconds}s - ${voiceEnvelope.total} packets';
return isChannelMessage ? '$senderName\n$summary' : summary;
}
final imageEnvelope = ImageEnvelope.tryParse(message.text);
if (imageEnvelope != null) {
final route = isChannelMessage
? 'Channel: ${channelName ?? _resolveChannelName(message.channelIdx)}'
: 'From: $senderName';
return '$route\nImage - ${imageEnvelope.format.label} - ${imageEnvelope.width}x${imageEnvelope.height} - ${_formatBytes(imageEnvelope.sizeBytes)}';
final summary =
'Image - ${imageEnvelope.format.label} - ${imageEnvelope.width}x${imageEnvelope.height} - ${_formatBytes(imageEnvelope.sizeBytes)}';
return isChannelMessage ? '$senderName\n$summary' : summary;
}
if (!isChannelMessage && message.recipientPublicKey != null) {
@@ -669,12 +661,12 @@ class MessagesProvider with ChangeNotifier {
fallback: null,
);
if (recipientName != 'Unknown') {
return 'From: $senderName\nTo: $recipientName\n${message.text}';
return 'To: $recipientName\n${message.text}';
}
}
if (isChannelMessage && channelName != null && channelName.isNotEmpty) {
return 'Channel: $channelName\n${message.text}';
if (isChannelMessage) {
return '$senderName\n${message.text}';
}
return message.text;
@@ -695,6 +687,7 @@ class MessagesProvider with ChangeNotifier {
_messages,
messageContactLocations: _messageContactLocations,
messageReceptionDetails: _messageReceptionDetails,
messageTransferDetails: _messageTransferDetails,
);
} catch (e) {
debugPrint('❌ [MessagesProvider] Error persisting messages: $e');
@@ -812,6 +805,7 @@ class MessagesProvider with ChangeNotifier {
_groupedMessageMapping.remove(messageId);
_messageContactLocations.remove(messageId);
_messageReceptionDetails.remove(messageId);
_messageTransferDetails.remove(messageId);
debugPrint('🗑️ [MessagesProvider] Message $messageId deleted');
@@ -842,6 +836,7 @@ class MessagesProvider with ChangeNotifier {
_sarMarkers.clear();
_messageContactLocations.clear();
_messageReceptionDetails.clear();
_messageTransferDetails.clear();
_persistMessages();
notifyListeners();
}
@@ -858,10 +853,69 @@ class MessagesProvider with ChangeNotifier {
_sarMarkers.clear();
_messageContactLocations.clear();
_messageReceptionDetails.clear();
_messageTransferDetails.clear();
_persistMessages();
notifyListeners();
}
int transferCountForSession({
String? voiceSessionId,
String? imageSessionId,
}) {
final messageId = _findMessageIdByMediaSession(
voiceSessionId: voiceSessionId,
imageSessionId: imageSessionId,
);
if (messageId == null) return 0;
return _messageTransferDetails[messageId]?.totalTransfers ?? 0;
}
void recordMediaTransfer({
required String sessionId,
required String mediaType,
required String requesterKey6,
String? requesterName,
}) {
final messageId = _findMessageIdByMediaSession(
voiceSessionId: mediaType == 'voice' ? sessionId : null,
imageSessionId: mediaType == 'image' ? sessionId : null,
);
if (messageId == null) {
debugPrint(
'⚠️ [MessagesProvider] No message found for $mediaType session $sessionId',
);
return;
}
final current =
_messageTransferDetails[messageId] ??
const MessageTransferDetails.empty();
_messageTransferDetails[messageId] = current.registerTransfer(
requesterKey6: requesterKey6,
requesterName: requesterName,
);
_persistMessages();
notifyListeners();
}
String? _findMessageIdByMediaSession({
String? voiceSessionId,
String? imageSessionId,
}) {
for (final message in _messages.reversed) {
if (voiceSessionId != null && message.voiceId == voiceSessionId) {
return message.id;
}
if (imageSessionId != null) {
final envelope = ImageEnvelope.tryParse(message.text);
if (envelope != null && envelope.sessionId == imageSessionId) {
return message.id;
}
}
}
return null;
}
/// Get storage statistics
Future<Map<String, dynamic>> getStorageStats() async {
return await _storageService.getStorageStats();
@@ -957,14 +1011,6 @@ class MessagesProvider with ChangeNotifier {
isVoice: true,
voiceId: envelope.sessionId,
);
} else if (VoicePacket.isVoiceText(enhancedMessage.text)) {
final pkt = VoicePacket.tryParseText(enhancedMessage.text);
if (pkt != null) {
enhancedMessage = enhancedMessage.copyWith(
isVoice: true,
voiceId: pkt.sessionId,
);
}
}
}
@@ -1628,7 +1674,7 @@ class MessagesProvider with ChangeNotifier {
debugPrint('❌ [MessagesProvider] Message $messageId timeout/failed');
debugPrint(' Retry attempt: ${message.retryAttempt}');
debugPrint(' Contact has path: ${contact?.hasPath ?? false}');
debugPrint(' Contact has path: ${contact?.routeHasPath ?? false}');
debugPrint(' Used flood fallback: ${message.usedFloodFallback}');
// Decision tree for retry/flood/fail
@@ -1784,7 +1830,7 @@ class MessagesProvider with ChangeNotifier {
_retryManager.clearRetry(messageId);
final failedContact = _messageContactMap[messageId];
if (failedContact != null && failedContact.hasPath) {
if (failedContact != null && failedContact.routeHasPath) {
final failureStreak = _retryManager.recordPathFailure(failedContact);
debugPrint(
' Path failure streak for ${failedContact.advName}: $failureStreak',
@@ -1804,8 +1850,69 @@ class MessagesProvider with ChangeNotifier {
}
}
/// Reset an existing failed message back into a sending state so a manual
/// retry can reuse the same record instead of appending a duplicate.
bool prepareMessageForRetry(String messageId) {
final index = _messages.indexWhere((m) => m.id == messageId);
if (index == -1) {
debugPrint(
'⚠️ [MessagesProvider] prepareMessageForRetry: Message not found: $messageId',
);
return false;
}
final message = _messages[index];
_timeoutTimers[message.id]?.cancel();
_timeoutTimers.remove(message.id);
if (message.expectedAckTag != null) {
_pendingSentMessages.remove(message.expectedAckTag);
}
_clearAckHistoryForMessage(messageId);
_retryManager.clearRetry(messageId);
_messages[index] = Message(
id: message.id,
messageType: message.messageType,
senderPublicKeyPrefix: message.senderPublicKeyPrefix,
channelIdx: message.channelIdx,
pathLen: message.pathLen,
textType: message.textType,
senderTimestamp: message.senderTimestamp,
text: message.text,
isSarMarker: message.isSarMarker,
sarGpsCoordinates: message.sarGpsCoordinates,
sarNotes: message.sarNotes,
sarCustomEmoji: message.sarCustomEmoji,
sarColorIndex: message.sarColorIndex,
receivedAt: message.receivedAt,
senderName: message.senderName,
deliveryStatus: MessageDeliveryStatus.sending,
recipientPublicKey: message.recipientPublicKey,
retryAttempt: 0,
lastRetryAt: DateTime.now(),
usedFloodFallback: false,
isRead: message.isRead,
echoCount: message.echoCount,
firstEchoAt: message.firstEchoAt,
lastEchoSnrRaw: message.lastEchoSnrRaw,
lastEchoRssiDbm: message.lastEchoRssiDbm,
lastEchoAt: message.lastEchoAt,
isDrawing: message.isDrawing,
drawingId: message.drawingId,
groupId: message.groupId,
recipients: message.recipients,
isVoice: message.isVoice,
voiceId: message.voiceId,
);
_persistMessages();
notifyListeners();
return true;
}
/// Resend a failed message
Future<void> resendMessage(String messageId) async {
Future<void> resendMessage(String messageId, {Contact? contact}) async {
final index = _messages.indexWhere((m) => m.id == messageId);
if (index == -1) {
debugPrint(
@@ -1815,9 +1922,9 @@ class MessagesProvider with ChangeNotifier {
}
final message = _messages[index];
final contact = _messageContactMap[messageId];
final resolvedContact = contact ?? _messageContactMap[messageId];
if (contact == null) {
if (resolvedContact == null) {
debugPrint(
'⚠️ [MessagesProvider] Cannot resend: Contact not found for message $messageId',
);
@@ -1826,26 +1933,19 @@ class MessagesProvider with ChangeNotifier {
debugPrint('🔁 [MessagesProvider] Resending message $messageId');
// Reset retry state
_messages[index] = message.copyWith(
retryAttempt: 0,
usedFloodFallback: false,
deliveryStatus: MessageDeliveryStatus.sending,
lastRetryAt: DateTime.now(),
);
// Clear retry tracking
_retryManager.clearRetry(messageId);
notifyListeners();
_messageContactMap[messageId] = resolvedContact;
final prepared = prepareMessageForRetry(messageId);
if (!prepared) {
return;
}
// Send again
if (sendMessageCallback != null) {
final queued = await sendMessageCallback!(
contactPublicKey: contact.publicKey,
contactPublicKey: resolvedContact.publicKey,
text: message.text,
messageId: messageId,
contact: contact,
contact: resolvedContact,
retryAttempt: 0,
);
if (!queued) {

View File

@@ -125,6 +125,23 @@ class VoiceProvider with ChangeNotifier {
return missing;
}
List<int> availablePacketIndices(String sessionId) {
final outgoing = _outgoingSessions[sessionId];
if (outgoing != null) {
return outgoing.packets.map((packet) => packet.index).toList()..sort();
}
final session = _sessions[sessionId];
if (session == null) return const [];
final indices = <int>[];
for (var i = 0; i < session.packets.length; i++) {
if (session.packets[i] != null) {
indices.add(i);
}
}
return indices;
}
// ── Packet reception ─────────────────────────────────────────────────────
/// Add an incoming [packet] to its session. Creates the session on first packet.
@@ -136,14 +153,18 @@ class VoiceProvider with ChangeNotifier {
);
return false;
}
_sessions.putIfAbsent(
packet.sessionId,
() => VoiceSession(
_sessions.putIfAbsent(packet.sessionId, () {
if (packet.total < 1) {
throw StateError(
'Voice envelope missing for compact packet ${packet.sessionId}',
);
}
return VoiceSession(
sessionId: packet.sessionId,
mode: packet.mode,
total: packet.total,
),
);
);
});
final session = _sessions[packet.sessionId]!;
if (packet.index < session.total) {
@@ -179,6 +200,47 @@ class VoiceProvider with ChangeNotifier {
}
}
void registerEnvelope(VoiceEnvelope envelope) {
if (_ignoredIncomingSessions.contains(envelope.sessionId)) {
return;
}
final existing = _sessions[envelope.sessionId];
if (existing == null) {
_sessions[envelope.sessionId] = VoiceSession(
sessionId: envelope.sessionId,
mode: envelope.mode,
total: envelope.total,
);
_persistVoiceData();
notifyListeners();
return;
}
final needsMerge =
existing.total != envelope.total || existing.mode != envelope.mode;
if (!needsMerge) {
notifyListeners();
return;
}
final merged = VoiceSession(
sessionId: envelope.sessionId,
mode: envelope.mode,
total: envelope.total,
);
merged.firstPacketAt = existing.firstPacketAt;
merged.lastPacketAt = existing.lastPacketAt;
for (final packet in existing.packets) {
if (packet == null) continue;
if (packet.index < merged.total) {
merged.packets[packet.index] = packet;
}
}
_sessions[envelope.sessionId] = merged;
_persistVoiceData();
notifyListeners();
}
/// Cache encoded packets for deferred voice serving.
void cacheOutgoingSession(String sessionId, List<VoicePacket> packets) {
if (packets.isEmpty) return;
@@ -195,10 +257,14 @@ class VoiceProvider with ChangeNotifier {
required Contact requester,
Set<int>? requestedIndices,
}) async {
final cached = _outgoingSessions[sessionId];
if (cached == null) {
final outgoing = _outgoingSessions[sessionId];
final packets = outgoing != null
? List<VoicePacket>.from(outgoing.packets)
: _sessions[sessionId]?.packets.whereType<VoicePacket>().toList() ??
const <VoicePacket>[];
if (packets.isEmpty) {
debugPrint(
'⚠️ [VoiceProvider] No cached outgoing session for $sessionId',
'⚠️ [VoiceProvider] No cached or received session for $sessionId',
);
return false;
}
@@ -206,7 +272,7 @@ class VoiceProvider with ChangeNotifier {
providerLabel: 'VoiceProvider',
sessionId: sessionId,
requester: requester,
fragments: cached.packets,
fragments: packets,
maxDirectPayloadHops: maxDirectPayloadHops,
indexOf: (packet) => packet.index,
encodeBinary: (packet) => packet.encodeBinary(),