Push meshcore_client and refresh pub

This commit is contained in:
Janez T
2026-03-05 20:58:29 +01:00
parent 2b34bc4162
commit 2482f94cd0
15 changed files with 885 additions and 291 deletions

View File

@@ -8,6 +8,8 @@ import 'drawing_provider.dart';
import 'channels_provider.dart';
import 'voice_provider.dart';
import 'image_provider.dart' as ip;
import 'helpers/fragment_ack_wait_registry.dart';
import 'helpers/session_metadata_restore.dart';
import '../services/tile_cache_service.dart';
import '../services/location_tracking_service.dart';
import '../services/packet_capture_storage_service.dart';
@@ -62,8 +64,10 @@ class AppProvider with ChangeNotifier {
final Map<String, String> _imageSessionSenderKey6 = {};
final Map<String, Timer> _voiceMissingRetryTimers = {};
final Map<String, int> _voiceMissingRetryAttempts = {};
final Map<String, Completer<void>> _voiceFragmentAckWaiters = {};
final Map<String, Completer<void>> _imageFragmentAckWaiters = {};
final FragmentAckWaitRegistry _voiceFragmentAckWaiters =
FragmentAckWaitRegistry();
final FragmentAckWaitRegistry _imageFragmentAckWaiters =
FragmentAckWaitRegistry();
Timer? _packetCaptureFlushTimer;
String? _lastPersistedPacketSignature;
bool _isPersistingPacketCapture = false;
@@ -165,12 +169,35 @@ class AppProvider with ChangeNotifier {
// Give DrawingProvider a moment to finish loading too
await Future.delayed(const Duration(milliseconds: 100));
_restoreSessionMetadataFromMessages();
debugPrint(
'🔄 [AppProvider] Early sync: syncing drawings from messages...',
);
messagesProvider.syncDrawingsWithProvider(drawingProvider);
}
void _restoreSessionMetadataFromMessages() {
final restored = restoreSessionMetadataFromMessages(
messagesProvider.messages.map((message) => message.text),
);
_voiceSessionSenderKey6.addAll(restored.voiceSenderKeyBySession);
for (final entry in restored.imageEnvelopeBySession.entries) {
_imageSessionSenderKey6[entry.key] = entry.value.senderKey6.toLowerCase();
imageProvider.registerEnvelope(entry.value);
}
final restoredVoice = restored.voiceSenderKeyBySession.length;
final restoredImage = restored.imageEnvelopeBySession.length;
if (restoredVoice > 0 || restoredImage > 0) {
debugPrint(
'🔄 [AppProvider] Restored session metadata from messages: '
'$restoredVoice voice, $restoredImage image',
);
}
}
/// Load simple mode setting from shared preferences
Future<void> _loadSimpleMode() async {
try {
@@ -788,12 +815,12 @@ class AppProvider with ChangeNotifier {
final imageFetchRequest = ImageFetchRequest.tryParseBinary(payload);
if (imageFetchRequest != null) {
final requester = contactsProvider.findContactByPrefixHex(
imageFetchRequest.requesterKey6,
);
final requester = _resolveImageFetchRequester(imageFetchRequest);
if (requester == null) {
debugPrint(
'⚠️ [AppProvider] Image fetch requester contact not found (binary)',
'⚠️ [AppProvider] Image fetch requester contact not found (binary) '
'for session ${imageFetchRequest.sessionId} / '
'${imageFetchRequest.requesterKey6}',
);
messagesProvider.logSystemMessage(
text:
@@ -804,7 +831,9 @@ class AppProvider with ChangeNotifier {
}
if (requester.outPathLen > _maxDirectPayloadHops) {
debugPrint(
'⚠️ [AppProvider] Image fetch requester too far: ${requester.outPathLen} hops',
'⚠️ [AppProvider] Image fetch requester too far: '
'${requester.outPathLen} hops for session '
'${imageFetchRequest.sessionId}',
);
messagesProvider.logSystemMessage(
text:
@@ -813,6 +842,10 @@ class AppProvider with ChangeNotifier {
);
return;
}
debugPrint(
'📷 [AppProvider] Serving image session ${imageFetchRequest.sessionId} '
'to ${requester.advName} via ${requester.outPathLen} hop(s)',
);
unawaited(
imageProvider.serveSessionTo(
sessionId: imageFetchRequest.sessionId,
@@ -1189,6 +1222,42 @@ class AppProvider with ChangeNotifier {
return contactsProvider.findContactByPrefixHex(prefixHex.toLowerCase());
}
Contact? _resolveImageFetchRequester(ImageFetchRequest request) {
final liveContact = _resolveContactByPrefixHex(request.requesterKey6);
if (liveContact != null) {
return liveContact;
}
for (final message in messagesProvider.messages.reversed) {
final envelope = ImageEnvelope.tryParse(message.text);
if (envelope == null || envelope.sessionId != request.sessionId) {
continue;
}
final recipientKey = message.recipientPublicKey;
if (recipientKey == null || recipientKey.isEmpty) {
continue;
}
final recipient = contactsProvider.findContactByKey(recipientKey);
if (recipient == null) {
continue;
}
final recipientKey6 = recipient.publicKey
.sublist(0, 6)
.map((b) => b.toRadixString(16).padLeft(2, '0'))
.join();
if (recipientKey6 != request.requesterKey6) {
continue;
}
debugPrint(
'📷 [AppProvider] Resolved image requester from sent message metadata '
'for session ${request.sessionId}: ${recipient.advName}',
);
return recipient;
}
return null;
}
void _scheduleVoiceMissingRetry(
String sessionId, {
required bool justComplete,
@@ -1312,54 +1381,48 @@ class AppProvider with ChangeNotifier {
required String sessionId,
required int index,
Duration timeout = const Duration(seconds: 8),
}) async {
final key = _fragmentAckKey(sessionId, index);
final completer = Completer<void>();
_voiceFragmentAckWaiters[key] = completer;
try {
await completer.future.timeout(timeout);
return true;
} catch (_) {
if (_voiceFragmentAckWaiters[key] == completer) {
_voiceFragmentAckWaiters.remove(key);
}
return false;
}
}
}) => _voiceFragmentAckWaiters.waitFor(
_fragmentAckKey(sessionId, index),
timeout: timeout,
);
void _completeVoiceFragmentAck(String sessionId, int index) {
final key = _fragmentAckKey(sessionId, index);
final completer = _voiceFragmentAckWaiters.remove(key);
if (completer != null && !completer.isCompleted) {
completer.complete();
final completed = _voiceFragmentAckWaiters.complete(
_fragmentAckKey(sessionId, index),
);
if (completed == 0) {
debugPrint(
' [AppProvider] Voice fragment ACK had no waiter: $sessionId#$index',
);
return;
}
debugPrint(
'✅ [AppProvider] Voice fragment ACK received for $sessionId#$index ($completed waiter(s))',
);
}
Future<bool> _waitForImageFragmentAck({
required String sessionId,
required int index,
Duration timeout = const Duration(seconds: 8),
}) async {
final key = _fragmentAckKey(sessionId, index);
final completer = Completer<void>();
_imageFragmentAckWaiters[key] = completer;
try {
await completer.future.timeout(timeout);
return true;
} catch (_) {
if (_imageFragmentAckWaiters[key] == completer) {
_imageFragmentAckWaiters.remove(key);
}
return false;
}
}
}) => _imageFragmentAckWaiters.waitFor(
_fragmentAckKey(sessionId, index),
timeout: timeout,
);
void _completeImageFragmentAck(String sessionId, int index) {
final key = _fragmentAckKey(sessionId, index);
final completer = _imageFragmentAckWaiters.remove(key);
if (completer != null && !completer.isCompleted) {
completer.complete();
final completed = _imageFragmentAckWaiters.complete(
_fragmentAckKey(sessionId, index),
);
if (completed == 0) {
debugPrint(
' [AppProvider] Image fragment ACK had no waiter: $sessionId#$index',
);
return;
}
debugPrint(
'✅ [AppProvider] Image fragment ACK received for $sessionId#$index ($completed waiter(s))',
);
}
void _sendVoiceFragmentAck(VoicePacket packet) {

View File

@@ -0,0 +1,43 @@
import 'dart:async';
/// Tracks one or more in-flight waiters for the same fragment ACK key.
///
/// Duplicate fetch requests can race and wait on the same fragment ACK at once.
/// Completing all registered waiters avoids losing the earlier completer when a
/// later request registers for the same key.
class FragmentAckWaitRegistry {
final Map<String, List<Completer<void>>> _waiters = {};
Future<bool> waitFor(
String key, {
Duration timeout = const Duration(seconds: 8),
}) async {
final completer = Completer<void>();
final waiters = _waiters.putIfAbsent(key, () => <Completer<void>>[]);
waiters.add(completer);
try {
await completer.future.timeout(timeout);
return true;
} catch (_) {
final pending = _waiters[key];
pending?.remove(completer);
if (pending != null && pending.isEmpty) {
_waiters.remove(key);
}
return false;
}
}
int complete(String key) {
final waiters = _waiters.remove(key);
if (waiters == null || waiters.isEmpty) {
return 0;
}
for (final completer in waiters) {
if (!completer.isCompleted) {
completer.complete();
}
}
return waiters.length;
}
}

View File

@@ -0,0 +1,38 @@
import '../../utils/image_message_parser.dart';
import '../../utils/voice_message_parser.dart';
class RestoredSessionMetadata {
final Map<String, String> voiceSenderKeyBySession;
final Map<String, ImageEnvelope> imageEnvelopeBySession;
const RestoredSessionMetadata({
required this.voiceSenderKeyBySession,
required this.imageEnvelopeBySession,
});
}
RestoredSessionMetadata restoreSessionMetadataFromMessages(
Iterable<String> messageTexts,
) {
final voiceSenderKeyBySession = <String, String>{};
final imageEnvelopeBySession = <String, ImageEnvelope>{};
for (final text in messageTexts) {
final voiceEnvelope = VoiceEnvelope.tryParseText(text);
if (voiceEnvelope != null) {
voiceSenderKeyBySession[voiceEnvelope.sessionId] = voiceEnvelope
.senderKey6
.toLowerCase();
}
final imageEnvelope = ImageEnvelope.tryParse(text);
if (imageEnvelope != null) {
imageEnvelopeBySession[imageEnvelope.sessionId] = imageEnvelope;
}
}
return RestoredSessionMetadata(
voiceSenderKeyBySession: voiceSenderKeyBySession,
imageEnvelopeBySession: imageEnvelopeBySession,
);
}

View File

@@ -58,6 +58,7 @@ class ImageProvider with ChangeNotifier {
/// Incoming sessions keyed by sessionId.
final Map<String, ImageSession> _sessions = {};
final Set<String> _ignoredIncomingSessions = {};
/// Outgoing sessions cached for deferred serving.
final Map<String, _OutgoingSession> _outgoing = {};
@@ -88,6 +89,8 @@ class ImageProvider with ChangeNotifier {
bool hasOutgoing(String sessionId) => _outgoing.containsKey(sessionId);
Duration? estimateRemainingTransferTime(String sessionId) =>
_sessions[sessionId]?.estimateRemaining();
bool isReceiveCanceled(String sessionId) =>
_ignoredIncomingSessions.contains(sessionId);
List<int> missingFragmentIndices(String sessionId) {
final session = _sessions[sessionId];
@@ -107,6 +110,12 @@ class ImageProvider with ChangeNotifier {
///
/// Returns true when the session just became complete.
bool addFragment(ImagePacket fragment, {int width = 0, int height = 0}) {
if (_ignoredIncomingSessions.contains(fragment.sessionId)) {
debugPrint(
'⏹️ [ImageProvider] Ignoring canceled incoming session ${fragment.sessionId}',
);
return false;
}
_sessions.putIfAbsent(
fragment.sessionId,
() => ImageSession(
@@ -135,9 +144,25 @@ class ImageProvider with ChangeNotifier {
return justComplete;
}
void cancelIncomingSession(String sessionId) {
_ignoredIncomingSessions.add(sessionId);
_sessions.remove(sessionId);
unawaited(_persist());
notifyListeners();
}
void resumeIncomingSession(String sessionId) {
if (_ignoredIncomingSessions.remove(sessionId)) {
notifyListeners();
}
}
/// Register envelope metadata for a session (called when IE1 is received
/// before any binary fragments arrive).
void registerEnvelope(ImageEnvelope envelope) {
if (_ignoredIncomingSessions.contains(envelope.sessionId)) {
return;
}
final existing = _sessions[envelope.sessionId];
if (existing == null) {
_sessions[envelope.sessionId] = ImageSession(
@@ -249,6 +274,7 @@ class ImageProvider with ChangeNotifier {
Future<void> clearAll() async {
_sessions.clear();
_outgoing.clear();
_ignoredIncomingSessions.clear();
notifyListeners();
try {
final prefs = await SharedPreferences.getInstance();

View File

@@ -59,6 +59,7 @@ class VoiceProvider with ChangeNotifier {
/// Active sessions keyed by sessionId.
final Map<String, VoiceSession> _sessions = {};
final Set<String> _ignoredIncomingSessions = {};
/// Currently playing session ID, or null.
String? _playingSessionId;
@@ -117,6 +118,8 @@ class VoiceProvider with ChangeNotifier {
_outgoingSessions.containsKey(sessionId);
Duration? estimateRemainingTransferTime(String sessionId) =>
_sessions[sessionId]?.estimateRemaining();
bool isReceiveCanceled(String sessionId) =>
_ignoredIncomingSessions.contains(sessionId);
List<int> missingPacketIndices(String sessionId) {
final session = _sessions[sessionId];
@@ -133,6 +136,12 @@ class VoiceProvider with ChangeNotifier {
/// Add an incoming [packet] to its session. Creates the session on first packet.
/// Returns true if the session just became complete.
bool addPacket(VoicePacket packet) {
if (_ignoredIncomingSessions.contains(packet.sessionId)) {
debugPrint(
'⏹️ [VoiceProvider] Ignoring canceled incoming session ${packet.sessionId}',
);
return false;
}
_sessions.putIfAbsent(
packet.sessionId,
() => VoiceSession(
@@ -159,6 +168,23 @@ class VoiceProvider with ChangeNotifier {
return justComplete;
}
void cancelIncomingSession(String sessionId) {
_ignoredIncomingSessions.add(sessionId);
_sessions.remove(sessionId);
if (_playingSessionId == sessionId) {
unawaited(_player.stop());
_playingSessionId = null;
}
_persistVoiceData();
notifyListeners();
}
void resumeIncomingSession(String sessionId) {
if (_ignoredIncomingSessions.remove(sessionId)) {
notifyListeners();
}
}
/// Cache encoded packets for deferred voice serving.
void cacheOutgoingSession(String sessionId, List<VoicePacket> packets) {
if (packets.isEmpty) return;
@@ -237,6 +263,7 @@ class VoiceProvider with ChangeNotifier {
Future<void> clearStoredVoiceData() async {
_sessions.clear();
_outgoingSessions.clear();
_ignoredIncomingSessions.clear();
_playingSessionId = null;
notifyListeners();
try {