112 lines
3.5 KiB
Dart
112 lines
3.5 KiB
Dart
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<Uint8List>? _incomingSub;
|
|
|
|
// Einfache Mutex-Kette (siehe unten in request()): serialisiert
|
|
// ueberlappende Aufrufe, ohne eine explizite Warteschlangen-Datenstruktur
|
|
// zu brauchen.
|
|
Future<void> _lock = Future.value();
|
|
|
|
Completer<MspFrame>? _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<Uint8List> request(int function, {Uint8List? payload}) async {
|
|
final previous = _lock;
|
|
final release = Completer<void>();
|
|
_lock = release.future;
|
|
await previous;
|
|
try {
|
|
return await _requestOnce(function, payload: payload);
|
|
} finally {
|
|
release.complete();
|
|
}
|
|
}
|
|
|
|
Future<Uint8List> _requestOnce(int function, {Uint8List? payload}) async {
|
|
Object? lastError;
|
|
for (var attempt = 0; attempt <= maxRetries; attempt++) {
|
|
final completer = Completer<MspFrame>();
|
|
_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();
|
|
}
|
|
}
|