mirror of
https://github.com/R0m1k3/xtremflow.git
synced 2026-10-11 17:30:00 +02:00
Guide TV : indexer le XMLTV hors de la boucle principale
Le serveur n'a qu'une boucle d'événements. La décompression et le parsing des dumps XMLTV (7000+ chaînes, ~90 Mo) la gelaient pendant plusieurs secondes : plus rien ne répondait, y compris le relais turbo.ts. La lecture se coupait donc au démarrage, au moment où l'écran des chaînes demandait le guide (EPG à 19-26 s dans les logs). - Téléchargement inchangé, puis gunzip + utf8 + parse dans Isolate.run - _parse devient statique (aucune capture de this, client HTTP non transférable) - Test : la boucle d'événements ne doit pas geler plus de 250 ms pendant l'indexation d'un dump de taille réelle (1,65 s avant) Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
1 parent
9f4e98fcea
commit
da67c479b3
2 files changed
+108
-7
No files matched your search
@@ -1,5 +1,7 @@
|
||||
import 'dart:convert';
|
||||
import 'dart:io';
|
||||
import 'dart:isolate';
|
||||
import 'dart:typed_data';
|
||||
|
||||
import 'package:http/http.dart' as http;
|
||||
import 'package:xml/xml_events.dart';
|
||||
@@ -115,8 +117,7 @@ class XmltvEpgService {
|
||||
|
||||
for (final url in sourceUrls) {
|
||||
try {
|
||||
final body = await _download(url);
|
||||
final parsed = _parse(body);
|
||||
final parsed = await _downloadAndParse(url);
|
||||
// Première source servie gagne : les suivantes ne comblent que les
|
||||
// chaînes encore absentes.
|
||||
for (final entry in parsed.entries) {
|
||||
@@ -144,26 +145,60 @@ class XmltvEpgService {
|
||||
print('[XmltvEpg] index prêt : ${merged.length} chaînes');
|
||||
}
|
||||
|
||||
Future<String> _download(String url) async {
|
||||
/// Télécharge [url] puis l'indexe dans un isolate dédié.
|
||||
///
|
||||
/// Le serveur n'a qu'une boucle d'événements : décompresser et parser un
|
||||
/// dump de plusieurs milliers de chaînes dessus la gelait pendant des
|
||||
/// secondes. Plus rien d'autre ne répondait — en particulier le relais
|
||||
/// `turbo.ts`, qui cessait d'envoyer des octets au lecteur : la lecture
|
||||
/// se coupait systématiquement au démarrage, au moment précis où l'écran
|
||||
/// des chaînes demandait le guide.
|
||||
Future<Map<String, List<XmltvProgramme>>> _downloadAndParse(
|
||||
String url,
|
||||
) async {
|
||||
final response =
|
||||
await _client.get(Uri.parse(url)).timeout(downloadTimeout);
|
||||
if (response.statusCode != 200) {
|
||||
throw HttpException('HTTP ${response.statusCode}');
|
||||
}
|
||||
|
||||
List<int> bytes = response.bodyBytes;
|
||||
// Variables locales : la fermeture envoyée à l'isolate ne doit pas
|
||||
// capturer `this` (le client HTTP n'est pas transférable).
|
||||
final bytes = response.bodyBytes;
|
||||
final retention = this.retention;
|
||||
final horizon = this.horizon;
|
||||
return Isolate.run(
|
||||
() => _decodeAndParse(bytes, retention: retention, horizon: horizon),
|
||||
);
|
||||
}
|
||||
|
||||
static Map<String, List<XmltvProgramme>> _decodeAndParse(
|
||||
Uint8List raw, {
|
||||
required Duration retention,
|
||||
required Duration horizon,
|
||||
}) {
|
||||
List<int> bytes = raw;
|
||||
// Beaucoup de miroirs servent du .gz sans en-tête Content-Encoding : on
|
||||
// regarde le nombre magique plutôt que de se fier aux en-têtes.
|
||||
if (bytes.length > 2 && bytes[0] == 0x1f && bytes[1] == 0x8b) {
|
||||
bytes = gzip.decode(bytes);
|
||||
}
|
||||
return utf8.decode(bytes, allowMalformed: true);
|
||||
return _parse(
|
||||
utf8.decode(bytes, allowMalformed: true),
|
||||
retention: retention,
|
||||
horizon: horizon,
|
||||
);
|
||||
}
|
||||
|
||||
/// Exposé pour les tests.
|
||||
Map<String, List<XmltvProgramme>> parseForTest(String xml) => _parse(xml);
|
||||
Map<String, List<XmltvProgramme>> parseForTest(String xml) =>
|
||||
_parse(xml, retention: retention, horizon: horizon);
|
||||
|
||||
Map<String, List<XmltvProgramme>> _parse(String xml) {
|
||||
static Map<String, List<XmltvProgramme>> _parse(
|
||||
String xml, {
|
||||
required Duration retention,
|
||||
required Duration horizon,
|
||||
}) {
|
||||
final cutoff = DateTime.now().toUtc().subtract(retention);
|
||||
final limit = DateTime.now().toUtc().add(horizon);
|
||||
final byChannel = <String, List<XmltvProgramme>>{};
|
||||
|
||||
@@ -1,3 +1,5 @@
|
||||
import 'dart:async';
|
||||
|
||||
import 'package:http/http.dart' as http;
|
||||
import 'package:http/testing.dart';
|
||||
import 'package:test/test.dart';
|
||||
@@ -148,6 +150,70 @@ void main() {
|
||||
});
|
||||
});
|
||||
|
||||
group('réactivité du serveur', () {
|
||||
test('l’indexation ne gèle pas la boucle d’événements', () async {
|
||||
// Le serveur n'a qu'un isolate : décoder et parser un dump de
|
||||
// plusieurs milliers de chaînes dessus bloquait tout le reste pendant
|
||||
// des secondes — y compris le relais `turbo.ts` des flux en cours,
|
||||
// d'où des coupures systématiques au démarrage de la lecture.
|
||||
final now = DateTime.now().toUtc();
|
||||
String stamp(DateTime d) {
|
||||
String two(int v) => v.toString().padLeft(2, '0');
|
||||
return '${d.year}${two(d.month)}${two(d.day)}'
|
||||
'${two(d.hour)}${two(d.minute)}${two(d.second)} +0000';
|
||||
}
|
||||
|
||||
// Volume d'un vrai dump de panneau : ~7000 chaînes, résumés de
|
||||
// plusieurs phrases.
|
||||
final desc = 'Résumé du programme, avec assez de texte pour peser '
|
||||
'comme un vrai guide. ' * 4;
|
||||
final xml = StringBuffer('<?xml version="1.0"?><tv>');
|
||||
for (var c = 0; c < 7000; c++) {
|
||||
xml.write('<channel id="Chaine.$c.fr">'
|
||||
'<display-name>FR - CHAINE $c</display-name></channel>');
|
||||
for (var p = 0; p < 30; p++) {
|
||||
final start = now.add(Duration(minutes: 30 * p));
|
||||
final stop = start.add(const Duration(minutes: 30));
|
||||
xml.write('<programme start="${stamp(start)}" '
|
||||
'stop="${stamp(stop)}" channel="Chaine.$c.fr">'
|
||||
'<title>Programme $p</title><desc>$desc</desc>'
|
||||
'</programme>');
|
||||
}
|
||||
}
|
||||
xml.write('</tv>');
|
||||
final body = xml.toString();
|
||||
|
||||
final service = XmltvEpgService(
|
||||
sourceUrls: const ['http://dump.example/epg.xml'],
|
||||
client: MockClient((_) async => http.Response(body, 200)),
|
||||
);
|
||||
|
||||
// Un « battement » toutes les 10 ms : le plus long écart observé est
|
||||
// le temps pendant lequel la boucle d'événements est restée bloquée.
|
||||
var last = DateTime.now();
|
||||
var worstGap = Duration.zero;
|
||||
final heartbeat = Timer.periodic(const Duration(milliseconds: 10), (_) {
|
||||
final now = DateTime.now();
|
||||
final gap = now.difference(last);
|
||||
if (gap > worstGap) worstGap = gap;
|
||||
last = now;
|
||||
});
|
||||
|
||||
await service.ensureFresh();
|
||||
heartbeat.cancel();
|
||||
// Le dernier blocage n'est suivi d'aucun battement : le compter aussi.
|
||||
final tail = DateTime.now().difference(last);
|
||||
if (tail > worstGap) worstGap = tail;
|
||||
|
||||
expect(service.channelCount, greaterThanOrEqualTo(7000));
|
||||
expect(
|
||||
worstGap,
|
||||
lessThan(const Duration(milliseconds: 250)),
|
||||
reason: 'boucle d’événements gelée ${worstGap.inMilliseconds} ms',
|
||||
);
|
||||
});
|
||||
});
|
||||
|
||||
group('horizon', () {
|
||||
test('écarte les programmes au-delà de l’horizon d’indexation', () {
|
||||
// Un dump national couvre sept jours : tout garder ferait grossir
|
||||
|
||||
Reference in new issue
Block a user