import 'dart:async'; import 'dart:typed_data'; import '../link_transport.dart'; import 'msp_frame.dart'; import 'msp_frame_decoder.dart'; /// Wird geworfen, wenn eine MSP-Anfrage nach allen Wiederholungen keine /// Antwort bekommen hat, oder wenn die Gegenstelle eine Fehlerantwort /// (`dir` = '!') gemeldet hat. class MspRequestException implements Exception { const MspRequestException(this.message); final String message; @override String toString() => 'MspRequestException: $message'; } /// MSP ist ein reines Frage-Antwort-Protokoll (Doku Kommunikationsschicht /// Abschnitt 4): die App fragt, der FC antwortet, von selbst kommt nichts. /// [MspClient] setzt darauf den in der Doku geforderten eigenen Takt auf: /// Anfragen werden serialisiert (immer nur eine offene Anfrage - sonst /// laufen bei schwacher Strecke die Antworten durcheinander), mit /// Zeitgrenze und Wiederholung. class MspClient { MspClient( this._transport, { this.requestTimeout = const Duration(milliseconds: 500), this.maxRetries = 2, }) { _incomingSub = _transport.incoming.listen(_onBytes); } final LinkTransport _transport; final Duration requestTimeout; final int maxRetries; final _decoder = MspFrameDecoder(); StreamSubscription? _incomingSub; // Einfache Mutex-Kette (siehe unten in request()): serialisiert // ueberlappende Aufrufe, ohne eine explizite Warteschlangen-Datenstruktur // zu brauchen. Future _lock = Future.value(); Completer? _pending; int? _pendingFunction; void _onBytes(Uint8List chunk) { for (final frame in _decoder.addBytes(chunk)) { final pending = _pending; if (pending != null && !pending.isCompleted && frame.function == _pendingFunction) { pending.complete(frame); } // Rahmen, die zu keiner offenen Anfrage passen (z.B. eine verspaetete // Antwort nach bereits abgelaufenem Timeout), werden verworfen - MSP // ist strikt Frage/Antwort, es gibt keine unaufgeforderten Pakete. } } /// Fragt [function] ab und liefert den rohen Antwort-Payload. Wartet ggf. /// auf eine bereits laufende Anfrage (Serialisierung), bevor die eigene /// gesendet wird. Future request(int function, {Uint8List? payload}) async { final previous = _lock; final release = Completer(); _lock = release.future; await previous; try { return await _requestOnce(function, payload: payload); } finally { release.complete(); } } Future _requestOnce(int function, {Uint8List? payload}) async { Object? lastError; for (var attempt = 0; attempt <= maxRetries; attempt++) { final completer = Completer(); _pending = completer; _pendingFunction = function; try { await _transport.send(encodeMspV2Request(function, payload)); final frame = await completer.future.timeout(requestTimeout); if (frame.direction == MspDirection.error) { throw MspRequestException( 'Flightcontroller meldet Fehler fuer MSP-Funktion $function', ); } return frame.payload; } on TimeoutException catch (e) { lastError = e; } finally { if (identical(_pending, completer)) { _pending = null; _pendingFunction = null; } } } throw MspRequestException( 'Keine Antwort auf MSP-Funktion $function nach ${maxRetries + 1} ' 'Versuchen${lastError != null ? ': $lastError' : ''}', ); } void dispose() { _incomingSub?.cancel(); } }