mirror of
https://github.com/dz0ny/meshcore-sar.git
synced 2026-10-09 19:55:05 +00:00
chore: Initial commit
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01XPovRHpejR4zRxuhjY4Qs3
This commit is contained in:
43
lib/providers/helpers/fragment_ack_wait_registry.dart
Normal file
43
lib/providers/helpers/fragment_ack_wait_registry.dart
Normal 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;
|
||||
}
|
||||
}
|
||||
234
lib/providers/helpers/message_delivery_tracker.dart
Normal file
234
lib/providers/helpers/message_delivery_tracker.dart
Normal file
@@ -0,0 +1,234 @@
|
||||
import 'dart:async';
|
||||
import 'dart:typed_data';
|
||||
|
||||
/// Message delivery tracking helper
|
||||
///
|
||||
/// Manages message delivery tracking for sent messages, including:
|
||||
/// - FIFO queue for matching RESP_CODE_SENT with message IDs
|
||||
/// - ACK tag to message ID mapping
|
||||
/// - Timeout tracking for stale ACK mappings
|
||||
/// - Message sent/delivered coordination
|
||||
///
|
||||
/// IMPORTANT: Based on MeshCore firmware analysis:
|
||||
/// - Firmware tracks max 8 pending ACKs in circular buffer
|
||||
/// - ACK entries overwritten after 8 messages → need rate limiting
|
||||
/// - Duplicate ACKs suppressed after first match
|
||||
/// - No automatic retry → app must implement
|
||||
class MessageDeliveryTracker {
|
||||
/// FIFO queue of pending message IDs
|
||||
/// Messages tracked here before sending, popped when RESP_CODE_SENT arrives
|
||||
final List<String> _pendingMessageIds = [];
|
||||
|
||||
/// Contact-scoped FIFOs for matching direct-message SENT responses.
|
||||
final Map<String, List<String>> _pendingMessageIdsByContact = {};
|
||||
|
||||
/// Map of ACK tag to message ID for delivery confirmation
|
||||
final Map<int, String> _ackTagToMessageId = {};
|
||||
|
||||
/// Map of message ID to ACK tag (reverse mapping for cleanup)
|
||||
final Map<String, int> _messageIdToAckTag = {};
|
||||
|
||||
/// Map of ACK tag to timestamp for timeout cleanup
|
||||
final Map<int, DateTime> _ackTagTimestamps = {};
|
||||
|
||||
/// Completer signalled when a pending ACK slot is freed (delivery or removal).
|
||||
Completer<void>? _slotFreedCompleter;
|
||||
|
||||
/// Track a pending message ID before sending
|
||||
///
|
||||
/// This is called BEFORE sending the message. When RESP_CODE_SENT
|
||||
/// arrives, we pop from this FIFO queue to match with the ACK tag.
|
||||
void trackPendingMessage(String messageId) {
|
||||
_pendingMessageIds.add(messageId);
|
||||
}
|
||||
|
||||
/// Track a pending direct message ID for a specific contact.
|
||||
void trackPendingDirectMessage(String messageId, Uint8List contactPublicKey) {
|
||||
trackPendingMessage(messageId);
|
||||
final contactKey = _contactKey(contactPublicKey);
|
||||
_pendingMessageIdsByContact
|
||||
.putIfAbsent(contactKey, () => [])
|
||||
.add(messageId);
|
||||
}
|
||||
|
||||
/// Pop the next pending message ID from FIFO queue
|
||||
///
|
||||
/// Called when RESP_CODE_SENT arrives. Returns null if queue empty.
|
||||
String? popPendingMessageId() {
|
||||
if (_pendingMessageIds.isEmpty) {
|
||||
return null;
|
||||
}
|
||||
return _pendingMessageIds.removeAt(0);
|
||||
}
|
||||
|
||||
/// Pop the next pending direct message ID for a specific contact.
|
||||
///
|
||||
/// Falls back to the legacy global FIFO if the contact queue is empty.
|
||||
String? popPendingDirectMessageId(Uint8List contactPublicKey) {
|
||||
final contactKey = _contactKey(contactPublicKey);
|
||||
final queue = _pendingMessageIdsByContact[contactKey];
|
||||
if (queue == null || queue.isEmpty) {
|
||||
return popPendingMessageId();
|
||||
}
|
||||
|
||||
final messageId = queue.removeAt(0);
|
||||
if (queue.isEmpty) {
|
||||
_pendingMessageIdsByContact.remove(contactKey);
|
||||
}
|
||||
_pendingMessageIds.remove(messageId);
|
||||
return messageId;
|
||||
}
|
||||
|
||||
/// Map ACK tag to message ID after RESP_CODE_SENT received
|
||||
///
|
||||
/// Creates bidirectional mapping for efficient cleanup and tracking.
|
||||
///
|
||||
/// WARNING: Firmware only tracks 8 pending ACKs! Caller should
|
||||
/// enforce rate limiting before calling this.
|
||||
void mapAckTagToMessageId(int ackTag, String messageId) {
|
||||
// Store bidirectional mapping
|
||||
_ackTagToMessageId[ackTag] = messageId;
|
||||
_messageIdToAckTag[messageId] = ackTag;
|
||||
_ackTagTimestamps[ackTag] = DateTime.now();
|
||||
}
|
||||
|
||||
/// Get message ID for ACK code
|
||||
///
|
||||
/// Called when SEND_CONFIRMED arrives. Returns the message ID
|
||||
/// that corresponds to this ACK code.
|
||||
///
|
||||
/// Returns null if ACK tag not found.
|
||||
String? getMessageIdForAck(int ackCode) {
|
||||
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.
|
||||
void removeAckTag(int ackCode) {
|
||||
final messageId = _ackTagToMessageId.remove(ackCode);
|
||||
if (messageId != null) {
|
||||
_messageIdToAckTag.remove(messageId);
|
||||
}
|
||||
_ackTagTimestamps.remove(ackCode);
|
||||
_notifySlotFreed();
|
||||
}
|
||||
|
||||
/// Remove ACK tag mapping by message ID
|
||||
///
|
||||
/// Used when message times out or is cancelled.
|
||||
void removeByMessageId(String messageId) {
|
||||
final ackTag = _messageIdToAckTag.remove(messageId);
|
||||
if (ackTag != null) {
|
||||
_ackTagToMessageId.remove(ackTag);
|
||||
_ackTagTimestamps.remove(ackTag);
|
||||
_notifySlotFreed();
|
||||
}
|
||||
_pendingMessageIds.remove(messageId);
|
||||
final emptyKeys = <String>[];
|
||||
for (final entry in _pendingMessageIdsByContact.entries) {
|
||||
entry.value.remove(messageId);
|
||||
if (entry.value.isEmpty) {
|
||||
emptyKeys.add(entry.key);
|
||||
}
|
||||
}
|
||||
for (final key in emptyKeys) {
|
||||
_pendingMessageIdsByContact.remove(key);
|
||||
}
|
||||
}
|
||||
|
||||
/// Clean up stale ACK mappings
|
||||
///
|
||||
/// Removes ACK tags that haven't received delivery confirmation
|
||||
/// within the specified timeout (default: 5 minutes).
|
||||
///
|
||||
/// Returns count of cleaned up entries.
|
||||
int cleanupStaleAcks({Duration timeout = const Duration(minutes: 5)}) {
|
||||
final now = DateTime.now();
|
||||
final staleAcks = <int>[];
|
||||
|
||||
for (final entry in _ackTagTimestamps.entries) {
|
||||
if (now.difference(entry.value) > timeout) {
|
||||
staleAcks.add(entry.key);
|
||||
}
|
||||
}
|
||||
|
||||
for (final ackTag in staleAcks) {
|
||||
removeAckTag(ackTag);
|
||||
}
|
||||
|
||||
return staleAcks.length;
|
||||
}
|
||||
|
||||
/// Clear all tracking state
|
||||
void clearTracking() {
|
||||
_pendingMessageIds.clear();
|
||||
_pendingMessageIdsByContact.clear();
|
||||
_ackTagToMessageId.clear();
|
||||
_messageIdToAckTag.clear();
|
||||
_ackTagTimestamps.clear();
|
||||
_notifySlotFreed();
|
||||
}
|
||||
|
||||
/// Get count of pending ACK tags
|
||||
///
|
||||
/// WARNING: Firmware only tracks 8 pending ACKs in circular buffer.
|
||||
/// If this exceeds 7, message sending should be rate limited.
|
||||
int get pendingCount => _ackTagToMessageId.length;
|
||||
|
||||
/// Check if should rate limit message sending
|
||||
///
|
||||
/// Returns true if >= 7 pending ACKs (stay under firmware limit of 8)
|
||||
bool get shouldRateLimit => pendingCount >= 7;
|
||||
|
||||
/// Wait until a pending ACK slot is freed, or [timeout] elapses.
|
||||
///
|
||||
/// Returns immediately if not at the rate limit.
|
||||
Future<void> waitForSlot({
|
||||
Duration timeout = const Duration(milliseconds: 500),
|
||||
}) async {
|
||||
if (!shouldRateLimit) return;
|
||||
_slotFreedCompleter ??= Completer<void>();
|
||||
await _slotFreedCompleter!.future.timeout(
|
||||
timeout,
|
||||
onTimeout: () {},
|
||||
);
|
||||
}
|
||||
|
||||
void _notifySlotFreed() {
|
||||
if (_slotFreedCompleter != null && !_slotFreedCompleter!.isCompleted) {
|
||||
_slotFreedCompleter!.complete();
|
||||
}
|
||||
_slotFreedCompleter = null;
|
||||
}
|
||||
|
||||
/// Get oldest pending ACK timestamp (for debugging)
|
||||
DateTime? get oldestPendingTimestamp {
|
||||
if (_ackTagTimestamps.isEmpty) return null;
|
||||
return _ackTagTimestamps.values.reduce((a, b) => a.isBefore(b) ? a : b);
|
||||
}
|
||||
|
||||
/// Get diagnostic info for debugging
|
||||
Map<String, dynamic> getDiagnostics() {
|
||||
return {
|
||||
'pendingCount': pendingCount,
|
||||
'shouldRateLimit': shouldRateLimit,
|
||||
'oldestPending': oldestPendingTimestamp?.toIso8601String(),
|
||||
'ackTags': _ackTagToMessageId.keys.toList(),
|
||||
'pendingByContact': _pendingMessageIdsByContact.map(
|
||||
(key, value) => MapEntry(key, value.length),
|
||||
),
|
||||
};
|
||||
}
|
||||
|
||||
String _contactKey(Uint8List contactPublicKey) {
|
||||
return contactPublicKey
|
||||
.map((b) => b.toRadixString(16).padLeft(2, '0'))
|
||||
.join();
|
||||
}
|
||||
}
|
||||
178
lib/providers/helpers/message_retry_manager.dart
Normal file
178
lib/providers/helpers/message_retry_manager.dart
Normal file
@@ -0,0 +1,178 @@
|
||||
import 'dart:convert';
|
||||
import 'dart:math' as math;
|
||||
|
||||
import '../../models/message.dart';
|
||||
import '../../models/contact.dart';
|
||||
|
||||
/// Manages message retry state and logic
|
||||
///
|
||||
/// This helper class centralizes retry logic for direct messages.
|
||||
///
|
||||
/// IMPORTANT: Based on MeshCore firmware analysis:
|
||||
/// - Firmware calculates timeout based on path length and airtime
|
||||
/// - Direct mode: ~(path_len * airtime * 2) + margin
|
||||
/// - Flood mode: ~10-30 seconds for multi-hop
|
||||
/// - Our retry delays (1s, 2s, 4s, 8s) are app-level backoff timers
|
||||
/// - Firmware does NOT automatically retry - app must implement
|
||||
class MessageRetryManager {
|
||||
// Track retry state for each message ID
|
||||
final Map<String, int> _retryAttempts = {};
|
||||
final Map<String, DateTime> _lastRetryTimes = {};
|
||||
final Map<String, int> _pathFailureStreaks = {};
|
||||
|
||||
/// Max retry attempts when the contact has a known path.
|
||||
/// Sequence: 2 direct attempts, flood, then last successful route.
|
||||
static const int maxRetryAttemptsWithPath = 3;
|
||||
|
||||
/// No retries for flood-only contacts (no known path).
|
||||
/// Value 0 means: don't retry at all, go straight to fallback/fail.
|
||||
static const int maxRetryAttemptsFloodOnly = 0;
|
||||
|
||||
@Deprecated('Use maxRetryAttemptsForContact instead')
|
||||
static const int maxRetryAttempts = maxRetryAttemptsWithPath;
|
||||
|
||||
// Retry backoff values in milliseconds.
|
||||
static const List<int> _retryDelays = [1000, 2000, 4000, 8000, 8000];
|
||||
static const int _defaultLoRaSf = 10;
|
||||
static const int _defaultLoRaCr = 5;
|
||||
static const int _defaultLoRaBwHz = 250000;
|
||||
static const int _defaultLoRaPreambleSymbols = 8;
|
||||
static const int _defaultLoRaCrcEnabled = 1;
|
||||
static const int _defaultLoRaExplicitHeader = 1;
|
||||
|
||||
/// Get backoff delay for the next retry attempt.
|
||||
int getDelayForAttempt(int attempt) {
|
||||
if (attempt < 0 || attempt >= _retryDelays.length) {
|
||||
return _retryDelays.last;
|
||||
}
|
||||
return _retryDelays[attempt];
|
||||
}
|
||||
|
||||
static final math.Random _rng = math.Random();
|
||||
|
||||
/// Calculate a delivery-ACK timeout with random jitter.
|
||||
///
|
||||
/// Matches the official MeshCore app: `suggestedTimeout + random(1-8s)`.
|
||||
/// The jitter prevents collision when multiple messages are in flight.
|
||||
int calculateAckTimeoutMs({
|
||||
required String text,
|
||||
required Contact? contact,
|
||||
int? suggestedTimeoutMs,
|
||||
}) {
|
||||
int baseTimeout;
|
||||
if (suggestedTimeoutMs != null && suggestedTimeoutMs > 0) {
|
||||
baseTimeout = suggestedTimeoutMs;
|
||||
} else {
|
||||
final payloadBytes = utf8.encode(text).length;
|
||||
final airtimeMs = _estimateLoRaAirtimeMs(payloadBytes);
|
||||
final hopCount = contact?.routeHasPath == true
|
||||
? contact!.routeHopCount
|
||||
: -1;
|
||||
|
||||
if (hopCount < 0) {
|
||||
baseTimeout = ((airtimeMs * 10) + 4000).clamp(10000, 30000);
|
||||
} else {
|
||||
baseTimeout = ((airtimeMs * (hopCount + 1) * 2) + 1500).clamp(4000, 20000);
|
||||
}
|
||||
}
|
||||
|
||||
// Add random jitter: 500-2000ms to avoid collision
|
||||
final jitterMs = 500 + _rng.nextInt(1501);
|
||||
return baseTimeout + jitterMs;
|
||||
}
|
||||
|
||||
/// Max attempts for a given contact based on whether it has a known path.
|
||||
static int maxRetryAttemptsForContact(Contact? contact) {
|
||||
final hasPath = contact?.routeHasPath ?? false;
|
||||
return hasPath ? maxRetryAttemptsWithPath : maxRetryAttemptsFloodOnly;
|
||||
}
|
||||
|
||||
bool canRetry(Message message, Contact contact) {
|
||||
return message.retryAttempt < maxRetryAttemptsForContact(contact);
|
||||
}
|
||||
|
||||
/// Whether this is the last retry attempt.
|
||||
/// When true, the caller should reset the path to force flood mode.
|
||||
///
|
||||
/// Note: called AFTER retryAttempt has been incremented, so we compare
|
||||
/// directly against maxAttempts (not +1).
|
||||
bool isLastAttempt(Message message, Contact contact) {
|
||||
final maxAttempts = maxRetryAttemptsForContact(contact);
|
||||
return maxAttempts > 1 && message.retryAttempt >= maxAttempts;
|
||||
}
|
||||
|
||||
/// Track a retry attempt for a message
|
||||
void trackRetry(String messageId, int attempt) {
|
||||
_retryAttempts[messageId] = attempt;
|
||||
_lastRetryTimes[messageId] = DateTime.now();
|
||||
}
|
||||
|
||||
/// Clear retry tracking for a message (on success or permanent failure)
|
||||
void clearRetry(String messageId) {
|
||||
_retryAttempts.remove(messageId);
|
||||
_lastRetryTimes.remove(messageId);
|
||||
}
|
||||
|
||||
/// Clear all retry tracking (on disconnect)
|
||||
void clearAll() {
|
||||
_retryAttempts.clear();
|
||||
_lastRetryTimes.clear();
|
||||
_pathFailureStreaks.clear();
|
||||
}
|
||||
|
||||
/// Get current retry attempt for a message (for debugging)
|
||||
int? getRetryAttempt(String messageId) {
|
||||
return _retryAttempts[messageId];
|
||||
}
|
||||
|
||||
/// Get last retry time for a message (for debugging)
|
||||
DateTime? getLastRetryTime(String messageId) {
|
||||
return _lastRetryTimes[messageId];
|
||||
}
|
||||
|
||||
/// Record a successful delivery for a contact and clear any accumulated
|
||||
/// route failure streak for future sends.
|
||||
void recordDeliverySuccess(Contact contact) {
|
||||
_pathFailureStreaks.remove(contact.publicKeyHex);
|
||||
}
|
||||
|
||||
/// Record a permanent route failure for a contact.
|
||||
///
|
||||
/// Returns the updated failure streak so callers can decide when to reset
|
||||
/// the learned path on the radio and in local state.
|
||||
int recordPathFailure(Contact contact) {
|
||||
final contactKey = contact.publicKeyHex;
|
||||
final next = (_pathFailureStreaks[contactKey] ?? 0) + 1;
|
||||
_pathFailureStreaks[contactKey] = next;
|
||||
return next;
|
||||
}
|
||||
|
||||
int? getPathFailureStreak(Contact contact) {
|
||||
return _pathFailureStreaks[contact.publicKeyHex];
|
||||
}
|
||||
|
||||
int _estimateLoRaAirtimeMs(int payloadLenBytes) {
|
||||
final sf = _defaultLoRaSf;
|
||||
final bw = _defaultLoRaBwHz.toDouble();
|
||||
final cr = (_defaultLoRaCr - 4).clamp(1, 4);
|
||||
final ih = _defaultLoRaExplicitHeader == 1 ? 0 : 1;
|
||||
final de = (sf >= 11 && _defaultLoRaBwHz <= 125000) ? 1 : 0;
|
||||
|
||||
final symbolMs = ((1 << sf) / bw) * 1000.0;
|
||||
final preambleMs = (_defaultLoRaPreambleSymbols + 4.25) * symbolMs;
|
||||
|
||||
final num =
|
||||
(8 * payloadLenBytes) -
|
||||
(4 * sf) +
|
||||
28 +
|
||||
(16 * _defaultLoRaCrcEnabled) -
|
||||
(20 * ih);
|
||||
final den = 4 * (sf - (2 * de));
|
||||
final payloadSymCoeff = den <= 0 ? 0 : (num / den).ceil();
|
||||
final payloadSymbols =
|
||||
8 + (payloadSymCoeff < 0 ? 0 : payloadSymCoeff) * (cr + 4);
|
||||
final payloadMs = payloadSymbols * symbolMs;
|
||||
|
||||
return (preambleMs + payloadMs).ceil();
|
||||
}
|
||||
}
|
||||
120
lib/providers/helpers/ping_tracker.dart
Normal file
120
lib/providers/helpers/ping_tracker.dart
Normal file
@@ -0,0 +1,120 @@
|
||||
import 'dart:async';
|
||||
import 'package:flutter/foundation.dart';
|
||||
|
||||
/// Helper class to track pending ping (telemetry) requests
|
||||
/// and implement automatic fallback to flooding if no response received
|
||||
class PingTracker {
|
||||
// Map of public key hex string to ping request state
|
||||
final Map<String, _PingRequest> _pendingPings = {};
|
||||
|
||||
// Timeout duration for ping responses (seconds)
|
||||
static const int _pingTimeoutSeconds = 5;
|
||||
|
||||
/// Track a new ping request
|
||||
/// Returns a Future that completes when either:
|
||||
/// - A response is received (completes with true)
|
||||
/// - Timeout occurs (completes with false)
|
||||
Future<bool> trackPing({
|
||||
required Uint8List publicKey,
|
||||
required bool wasDirectAttempt,
|
||||
}) {
|
||||
final String keyHex = _publicKeyToHex(publicKey);
|
||||
|
||||
// Cancel any existing pending ping for this contact
|
||||
_pendingPings[keyHex]?.cancel();
|
||||
|
||||
// Create new ping request tracker
|
||||
final completer = Completer<bool>();
|
||||
final timer = Timer(const Duration(seconds: _pingTimeoutSeconds), () {
|
||||
// Timeout occurred - mark as failed
|
||||
_pendingPings.remove(keyHex);
|
||||
if (!completer.isCompleted) {
|
||||
completer.complete(false);
|
||||
}
|
||||
});
|
||||
|
||||
_pendingPings[keyHex] = _PingRequest(
|
||||
publicKey: publicKey,
|
||||
wasDirectAttempt: wasDirectAttempt,
|
||||
timer: timer,
|
||||
completer: completer,
|
||||
);
|
||||
|
||||
return completer.future;
|
||||
}
|
||||
|
||||
/// Mark a ping as successful (response received)
|
||||
/// Should be called when telemetry response arrives
|
||||
void markPingSuccessful(Uint8List publicKey) {
|
||||
final requestKey = _findMatchingPendingPingKey(publicKey);
|
||||
final request = requestKey != null ? _pendingPings.remove(requestKey) : null;
|
||||
|
||||
if (request != null) {
|
||||
request.cancel();
|
||||
if (!request.completer.isCompleted) {
|
||||
request.completer.complete(true);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Check if there's a pending ping for this contact
|
||||
bool hasPendingPing(Uint8List publicKey) {
|
||||
final String keyHex = _publicKeyToHex(publicKey);
|
||||
return _pendingPings.containsKey(keyHex);
|
||||
}
|
||||
|
||||
/// Get pending ping info (was it a direct attempt?)
|
||||
bool? wasPingDirect(Uint8List publicKey) {
|
||||
final String keyHex = _publicKeyToHex(publicKey);
|
||||
return _pendingPings[keyHex]?.wasDirectAttempt;
|
||||
}
|
||||
|
||||
/// Clear all pending pings (useful on disconnect)
|
||||
void clearAll() {
|
||||
for (final request in _pendingPings.values) {
|
||||
request.cancel();
|
||||
}
|
||||
_pendingPings.clear();
|
||||
}
|
||||
|
||||
/// Convert public key to hex string for map key
|
||||
String _publicKeyToHex(Uint8List publicKey) {
|
||||
return publicKey.map((b) => b.toRadixString(16).padLeft(2, '0')).join('');
|
||||
}
|
||||
|
||||
String? _findMatchingPendingPingKey(Uint8List responseKey) {
|
||||
final responseHex = _publicKeyToHex(responseKey);
|
||||
|
||||
if (_pendingPings.containsKey(responseHex)) {
|
||||
return responseHex;
|
||||
}
|
||||
|
||||
for (final entry in _pendingPings.entries) {
|
||||
final pendingHex = entry.key;
|
||||
if (pendingHex.startsWith(responseHex) || responseHex.startsWith(pendingHex)) {
|
||||
return pendingHex;
|
||||
}
|
||||
}
|
||||
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
/// Internal class to track a single ping request
|
||||
class _PingRequest {
|
||||
final Uint8List publicKey;
|
||||
final bool wasDirectAttempt;
|
||||
final Timer timer;
|
||||
final Completer<bool> completer;
|
||||
|
||||
_PingRequest({
|
||||
required this.publicKey,
|
||||
required this.wasDirectAttempt,
|
||||
required this.timer,
|
||||
required this.completer,
|
||||
});
|
||||
|
||||
void cancel() {
|
||||
timer.cancel();
|
||||
}
|
||||
}
|
||||
90
lib/providers/helpers/raw_session_retransmit.dart
Normal file
90
lib/providers/helpers/raw_session_retransmit.dart
Normal file
@@ -0,0 +1,90 @@
|
||||
import 'package:flutter/foundation.dart';
|
||||
|
||||
import '../../models/contact.dart';
|
||||
|
||||
typedef RawPacketSender =
|
||||
Future<void> Function({
|
||||
required Uint8List contactPath,
|
||||
required int contactPathLen,
|
||||
required Uint8List payload,
|
||||
});
|
||||
|
||||
/// Pacing between consecutive raw fragments. Firmware v1.15 rejects bursts of
|
||||
/// raw packets with `ERROR: Table full` when fragments are blasted too quickly,
|
||||
/// so we space them out to keep the radio's raw transmit queue from overflowing.
|
||||
const Duration _interFragmentDelay = Duration(milliseconds: 350);
|
||||
|
||||
Future<bool> serveCachedSessionFragments<T>({
|
||||
required String providerLabel,
|
||||
required String sessionId,
|
||||
required Contact requester,
|
||||
required List<T> fragments,
|
||||
required int maxDirectPayloadHops,
|
||||
required int Function(T fragment) indexOf,
|
||||
required Uint8List Function(T fragment) encodeBinary,
|
||||
required RawPacketSender? sendRawPacket,
|
||||
Set<int>? requestedIndices,
|
||||
Duration interFragmentDelay = _interFragmentDelay,
|
||||
}) async {
|
||||
if (fragments.isEmpty) {
|
||||
debugPrint('⚠️ [$providerLabel] No cached fragments for $sessionId');
|
||||
return false;
|
||||
}
|
||||
if (sendRawPacket == null) {
|
||||
debugPrint('⚠️ [$providerLabel] sendRawPacketCallback not set');
|
||||
return false;
|
||||
}
|
||||
if (!requester.routeHasPath) {
|
||||
debugPrint('⚠️ [$providerLabel] ${requester.advName} has no direct path');
|
||||
return false;
|
||||
}
|
||||
if (requester.routeHopCount > maxDirectPayloadHops) {
|
||||
debugPrint(
|
||||
'⚠️ [$providerLabel] ${requester.advName} is too far: ${requester.routeHopCount} hops (max $maxDirectPayloadHops)',
|
||||
);
|
||||
return false;
|
||||
}
|
||||
if (requester.outPath.isEmpty) {
|
||||
debugPrint(
|
||||
'⚠️ [$providerLabel] ${requester.advName} has empty outPath payload',
|
||||
);
|
||||
return false;
|
||||
}
|
||||
|
||||
var servedCount = 0;
|
||||
for (final fragment in fragments) {
|
||||
final index = indexOf(fragment);
|
||||
if (index < 0) {
|
||||
debugPrint('⚠️ [$providerLabel] Invalid fragment index $index');
|
||||
continue;
|
||||
}
|
||||
if (requestedIndices != null && !requestedIndices.contains(index)) {
|
||||
continue;
|
||||
}
|
||||
try {
|
||||
if (servedCount > 0 && interFragmentDelay > Duration.zero) {
|
||||
await Future<void>.delayed(interFragmentDelay);
|
||||
}
|
||||
await sendRawPacket(
|
||||
contactPath: requester.outPath,
|
||||
contactPathLen: requester.routeEncodedPathLen,
|
||||
payload: encodeBinary(fragment),
|
||||
);
|
||||
servedCount++;
|
||||
} catch (e, st) {
|
||||
debugPrint(
|
||||
'❌ [$providerLabel] Serve error for $sessionId#$index: $e\n$st',
|
||||
);
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
if (servedCount == 0) {
|
||||
debugPrint(
|
||||
'⚠️ [$providerLabel] No fragments matched request for $sessionId',
|
||||
);
|
||||
return false;
|
||||
}
|
||||
debugPrint('✅ [$providerLabel] Served $servedCount fragments for $sessionId');
|
||||
return true;
|
||||
}
|
||||
84
lib/providers/helpers/room_login_manager.dart
Normal file
84
lib/providers/helpers/room_login_manager.dart
Normal file
@@ -0,0 +1,84 @@
|
||||
import 'package:flutter/foundation.dart';
|
||||
import 'package:shared_preferences/shared_preferences.dart';
|
||||
import '../../models/room_login_state.dart';
|
||||
|
||||
/// Room login state management helper
|
||||
///
|
||||
/// Manages login state tracking for room contacts, including:
|
||||
/// - Room login state per contact (Map of String to RoomLoginState)
|
||||
/// - Password checking logic
|
||||
/// - Login success/fail state updates
|
||||
class RoomLoginManager {
|
||||
/// Map of room public key prefix (hex string) to login state
|
||||
final Map<String, RoomLoginState> _roomLoginStates = {};
|
||||
|
||||
/// Get all room login states (unmodifiable view)
|
||||
Map<String, RoomLoginState> get roomLoginStates => Map.unmodifiable(_roomLoginStates);
|
||||
|
||||
/// Get login state for a room by public key prefix
|
||||
RoomLoginState? getRoomLoginState(Uint8List publicKeyPrefix) {
|
||||
final prefixHex = _publicKeyPrefixToHex(publicKeyPrefix);
|
||||
return _roomLoginStates[prefixHex];
|
||||
}
|
||||
|
||||
/// Check if logged into a specific room
|
||||
bool isLoggedIntoRoom(Uint8List publicKeyPrefix) {
|
||||
final state = getRoomLoginState(publicKeyPrefix);
|
||||
return state?.isLoggedIn ?? false;
|
||||
}
|
||||
|
||||
/// Update room login state after successful login
|
||||
Future<void> handleLoginSuccess({
|
||||
required Uint8List publicKeyPrefix,
|
||||
required int permissions,
|
||||
required bool isAdmin,
|
||||
required int tag,
|
||||
}) async {
|
||||
final prefixHex = _publicKeyPrefixToHex(publicKeyPrefix);
|
||||
final hasPassword = await _hasPasswordForRoom(publicKeyPrefix);
|
||||
|
||||
_roomLoginStates[prefixHex] = RoomLoginState.loggedIn(
|
||||
publicKeyPrefix: publicKeyPrefix,
|
||||
permissions: permissions,
|
||||
isAdmin: isAdmin,
|
||||
tag: tag,
|
||||
hasPassword: hasPassword,
|
||||
);
|
||||
}
|
||||
|
||||
/// Update room login state after failed login
|
||||
void handleLoginFail({
|
||||
required Uint8List publicKeyPrefix,
|
||||
}) {
|
||||
final prefixHex = _publicKeyPrefixToHex(publicKeyPrefix);
|
||||
|
||||
_roomLoginStates[prefixHex] = RoomLoginState.loggedOut(
|
||||
publicKeyPrefix: publicKeyPrefix,
|
||||
hasPassword: false, // Password was incorrect
|
||||
);
|
||||
}
|
||||
|
||||
/// Clear all room login states (call on disconnect)
|
||||
void clearRoomLoginStates() {
|
||||
_roomLoginStates.clear();
|
||||
}
|
||||
|
||||
/// Check if a password exists for a room (by public key prefix)
|
||||
Future<bool> _hasPasswordForRoom(Uint8List publicKeyPrefix) async {
|
||||
try {
|
||||
final prefs = await SharedPreferences.getInstance();
|
||||
// Convert prefix to hex string for storage key
|
||||
final prefixHex = _publicKeyPrefixToHex(publicKeyPrefix);
|
||||
final roomKey = 'room_password_$prefixHex';
|
||||
return prefs.getString(roomKey) != null;
|
||||
} catch (e) {
|
||||
debugPrint('Error checking password for room: $e');
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
/// Convert public key prefix to hex string (colon-separated)
|
||||
String _publicKeyPrefixToHex(Uint8List publicKeyPrefix) {
|
||||
return publicKeyPrefix.map((b) => b.toRadixString(16).padLeft(2, '0')).join(':');
|
||||
}
|
||||
}
|
||||
57
lib/providers/helpers/session_metadata_restore.dart
Normal file
57
lib/providers/helpers/session_metadata_restore.dart
Normal file
@@ -0,0 +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<Message> messages,
|
||||
) {
|
||||
final voiceSenderKeyBySession = <String, String>{};
|
||||
final imageSenderKeyBySession = <String, String>{};
|
||||
final imageEnvelopeBySession = <String, ImageEnvelope>{};
|
||||
|
||||
for (final message in messages) {
|
||||
final text = message.text;
|
||||
final voiceEnvelope = VoiceEnvelope.tryParseText(text);
|
||||
if (voiceEnvelope != null) {
|
||||
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,
|
||||
);
|
||||
}
|
||||
Reference in New Issue
Block a user