Merge pull request #15 from R0m1k3/fix/picons-epg-flood

Logos, guide TV et fluidité : couper le flot de requêtes mortes
This commit is contained in:
LogiFlow authored and GitHub committed 2026-09-19 12:00:53 +02:00
commit bde32adda9
11 files changed
+716 -78

No files matched your search

+9
View File
@@ -2,6 +2,14 @@
## Non publié
### 🖼️ Logos des chaînes
- **Plus d'acharnement sur un hébergeur de picons hors service** : quand la machine qui héberge les logos répond en erreur (maintenance, quota), le proxy sert désormais un pixel transparent mis en cache dix minutes au lieu de relayer le 503 au navigateur — qui redemandait chaque logo à chaque affichage de la grille. Au bout de trois échecs consécutifs, l'hôte est considéré hors service pendant cinq minutes et plus aucune requête ne sort. Le serveur tournant sur un seul isolate, ce flot de centaines de requêtes mortes passait devant les paquets vidéo et faisait saccader la lecture
### 📺 Guide TV (EPG)
- **Guide servi par le dump `xmltv.php` du panneau** : `player_api` réclamait un appel par chaîne, à plusieurs secondes pièce chez la plupart des revendeurs — une grille de trente chaînes mettait plus d'une minute à se remplir. Le dump du panneau couvre toutes les chaînes en une requête, réutilisée trois heures. L'interrogation chaîne par chaîne reste le repli quand le dump est absent ou ne couvre pas la chaîne, suivie de la source XMLTV externe
- **Une source XMLTV morte n'est plus retéléchargée à chaque consultation** : après un échec, la source est laissée de côté quinze minutes
- **Index XMLTV borné à 48 h** : un dump national couvre sept jours pour un millier de chaînes ; tout garder coûtait des centaines de mégaoctets au conteneur pour un guide qui n'affiche que le programme courant et les suivants
### ▶️ Lecture vidéo
- **Fix du démarrage des enregistrements** : la lecture partait au milieu du programme et se coupait aussitôt, obligeant à relancer une deuxième fois. Une playlist encore en cours de transcodage n'a pas d'`EXT-X-ENDLIST` : hls.js la traite comme du direct et démarrait donc au « bord du direct », collé au front d'encodage, sans aucune avance de segments. `startPosition: 0` hors direct, plus la neutralisation du rattrapage de latence (qui accélérait la lecture puis forçait un saut en avant pour rejoindre un direct inexistant)
- **Plus de ré-encodage inutile à la lecture d'un enregistrement** : la capture étant faite en copie, le fichier contient déjà du H.264 dans la quasi-totalité des cas. Il est désormais servi tel quel (`-c:v copy`, seul l'audio est converti en AAC) au lieu d'être ré-encodé à peine plus vite que le temps réel — la segmentation va maintenant à la vitesse du disque, l'enregistrement devient navigable en quelques secondes. Les codecs illisibles par le navigateur (HEVC, MPEG-2…) restent ré-encodés
@@ -15,6 +23,7 @@
- **Latence live maîtrisée** : rattrapage du direct activé dans mpegts.js (profil rapide) — les micro-coupures ne font plus dériver la lecture derrière le direct
- **Fix « Échec du chargement : vendor/mpegts.min.js »** : un échec de chargement d'une lib de lecture n'affiche plus un écran d'erreur définitif. Le chargement est retenté une fois en contournant le cache HTTP (une entrée tronquée condamnait le lecteur jusqu'au vidage manuel du cache), le résultat n'est mémorisé qu'en cas de succès (une promesse rejetée en cache rendait tout réessai impossible) et, si la lib reste introuvable, le live bascule automatiquement sur la route HLS équivalente. Le message affiché précise désormais la cause (HTTP 404, 429, réseau injoignable)
- **Préchargement de la bonne lib** : les trois players préchargeaient hls.js en dur, soit 618 Ko téléchargés pour rien à chaque zap TV — où c'est mpegts.js qui sert — au détriment du flux et du chargement de mpegts.js. Le préchargement suit maintenant le flux réellement demandé (et ne charge rien sur Safari/iOS, qui lit le HLS nativement)
- **Contre-pression sur le flux turbo** : le relais de FFmpeg vers le navigateur ne transmettait pas la pause du client à la source. Un lecteur plus lent que le flux — réseau domestique, onglet en arrière-plan — laissait FFmpeg produire à pleine vitesse pendant que le serveur empilait les paquets en mémoire : la lecture dérivait derrière le direct et saccadait
- **Alerte au démarrage** : le serveur signale explicitement l'absence de `web/vendor/*.min.js` au lancement, au lieu de laisser le navigateur échouer sans explication
### 🧰 Qualité / Infra
+61
View File
@@ -0,0 +1,61 @@
/// Mémoire courte des hôtes d'images injoignables.
///
/// POURQUOI : les URL de picons pointent chez l'hébergeur du revendeur, pas
/// sur le panneau. Quand cette machine tombe (maintenance, DNS, quota), elle
/// répond 5xx en quelques millisecondes, le navigateur réessaie à chaque
/// rendu de la grille et le proxy relaie des centaines de requêtes mortes par
/// seconde. Le serveur tournant sur un unique isolate Dart, ce flot passe
/// devant les paquets vidéo dans la boucle d'événements : l'image saccade
/// pendant qu'on s'acharne sur un hôte qu'on sait hors service.
///
/// On coupe donc court : au bout de [threshold] échecs consécutifs, l'hôte
/// est considéré mort pendant [ttl] et les requêtes suivantes sont servies
/// localement, sans appel sortant.
class AssetFailureCache {
AssetFailureCache({
this.threshold = 3,
this.ttl = const Duration(minutes: 5),
DateTime Function()? clock,
}) : _clock = clock ?? DateTime.now;
/// Nombre d'échecs consécutifs avant de couper un hôte. Un 5xx passager ne
/// doit pas priver l'utilisateur de ses logos.
final int threshold;
/// Durée de la coupure. Assez courte pour qu'un hôte réparé revienne seul.
final Duration ttl;
final DateTime Function() _clock;
final Map<String, int> _failures = {};
final Map<String, DateTime> _downUntil = {};
/// L'hôte est-il réputé hors service en ce moment ?
bool isDown(String host) {
final key = host.toLowerCase();
final until = _downUntil[key];
if (until == null) return false;
if (_clock().isBefore(until)) return true;
// Coupure expirée : on repart d'une ardoise vierge pour laisser une
// vraie chance au prochain appel.
_downUntil.remove(key);
_failures.remove(key);
return false;
}
void recordFailure(String host) {
final key = host.toLowerCase();
final count = (_failures[key] ?? 0) + 1;
_failures[key] = count;
if (count >= threshold) {
_downUntil[key] = _clock().add(ttl);
}
}
void recordSuccess(String host) {
final key = host.toLowerCase();
_failures.remove(key);
_downUntil.remove(key);
}
}
+139 -45
View File
@@ -8,10 +8,12 @@ import '../utils/log_redactor.dart';
/// API EPG — proxy vers Xtream avec cache 30 minutes
/// GET /api/epg/<channel_id>?days=1
///
/// Le panneau de l'abonné reste la source de référence. Une source XMLTV
/// externe n'est interrogée qu'en dernier recours, quand le panneau ne rend
/// aucun programme couvrant l'instant présent ou à venir — cas fréquent des
/// revendeurs dont le guide est figé depuis plusieurs jours.
/// Le panneau de l'abonné reste la source de référence, mais par son dump
/// `xmltv.php` d'abord : une requête couvre toutes les chaînes, là où
/// `player_api` en réclame une par chaîne à plusieurs secondes pièce. Vient
/// ensuite l'interrogation chaîne par chaîne, puis une source XMLTV externe
/// en dernier recours — cas des revendeurs dont le guide est figé depuis
/// plusieurs jours.
class EpgApi {
final Future<PlaylistConfig?> Function(Request) _getPlaylist;
@@ -31,7 +33,43 @@ class EpgApi {
/// pas par `stream_id` propre au panneau.
final Map<String, _ChannelMap> _channelMaps = {};
EpgApi(this._getPlaylist, {XmltvEpgService? xmltv}) : _xmltv = xmltv;
/// Client HTTP des appels `player_api`. Injectable pour les tests.
final http.Client _http;
/// Fabrique de la source XMLTV du panneau, une par compte.
final XmltvEpgService Function(PlaylistConfig) _panelXmltvBuilder;
/// Index XMLTV du panneau, par compte.
final Map<String, XmltvEpgService> _panelXmltv = {};
EpgApi(
this._getPlaylist, {
XmltvEpgService? xmltv,
http.Client? httpClient,
XmltvEpgService Function(PlaylistConfig)? panelXmltvBuilder,
}) : _xmltv = xmltv,
_http = httpClient ?? http.Client(),
_panelXmltvBuilder = panelXmltvBuilder ?? _defaultPanelXmltv;
/// Dump XMLTV servi par le panneau lui-même.
///
/// POURQUOI : `player_api.php?action=get_simple_data_table` répond en
/// plusieurs secondes chez beaucoup de revendeurs, et il faut un appel par
/// chaîne — afficher une grille de trente chaînes demandait donc plus d'une
/// minute. `xmltv.php` renvoie le guide de toutes les chaînes en une seule
/// requête, réutilisée ensuite pendant des heures. C'est ce que font les
/// clients IPTV rapides.
static XmltvEpgService _defaultPanelXmltv(PlaylistConfig playlist) {
final user = Uri.encodeQueryComponent(playlist.username);
final password = Uri.encodeQueryComponent(playlist.password);
return XmltvEpgService(
sourceUrls: ['${playlist.dns}/xmltv.php?username=$user&password=$password'],
refreshInterval: const Duration(hours: 3),
// Le guide est facultatif : on ne fait pas patienter l'utilisateur
// plusieurs minutes sur un panneau qui traîne.
downloadTimeout: const Duration(seconds: 90),
);
}
Future<Response> handleGetEpg(Request request, String channelId) async {
final playlist = await _getPlaylist(request);
@@ -52,47 +90,44 @@ class EpgApi {
}
try {
final dns = playlist.dns;
Map<String, dynamic> epgData = {
'channel_id': channelId,
'programmes': [],
};
// Ordre d'essai des actions Xtream.
// Sources par ordre de préférence. On s'arrête à la première qui
// contient un programme en cours ou à venir ; à défaut, la première
// non vide sert de repli.
//
// `get_simple_data_table` d'abord : c'est la seule action réellement
// universelle. `get_epg` n'existe pas sur beaucoup de panneaux — au lieu
// d'une erreur, ils renvoient poliment le payload d'authentification,
// sans champ `epg_listings`, ce qui produisait un guide vide impossible
// à distinguer d'une chaîne sans programme.
const actions = [
'get_simple_data_table',
'get_short_epg',
// Le dump du panneau passe devant : une requête couvre toutes les
// chaînes, là où `player_api` en demande une par chaîne, à plusieurs
// secondes pièce.
final sources = <(String, Future<List<Map<String, dynamic>>> Function())>[
('panel-xmltv', () => _panelProgrammes(playlist, channelId)),
('xtream', () => _xtreamProgrammes(playlist, channelId)),
('xmltv', () => _xmltvProgrammes(playlist, channelId)),
];
for (final action in actions) {
final url =
'$dns/player_api.php?username=${playlist.username}&password=${playlist.password}'
'&action=$action&stream_id=$channelId';
final response =
await http.get(Uri.parse(url)).timeout(const Duration(seconds: 60));
if (response.statusCode != 200) continue;
try {
epgData = transformEpgData(json.decode(response.body), channelId);
} catch (_) {
continue;
}
if ((epgData['programmes'] as List).isNotEmpty) break;
}
var epgData = <String, dynamic>{
'channel_id': channelId,
'programmes': <Map<String, dynamic>>[],
};
var source = 'xtream';
if (!_hasCurrentProgramme(epgData)) {
final fallback = await _xmltvProgrammes(playlist, channelId);
if (fallback.isNotEmpty) {
epgData = {'channel_id': channelId, 'programmes': fallback};
source = 'xmltv';
var found = false;
for (final (name, fetch) in sources) {
final programmes = await fetch();
if (programmes.isEmpty) continue;
final candidate = <String, dynamic>{
'channel_id': channelId,
'programmes': programmes,
};
if (!found) {
// Meilleur repli connu à ce stade, même si le guide est périmé.
epgData = candidate;
source = name;
found = true;
}
if (_hasCurrentProgramme(candidate)) {
epgData = candidate;
source = name;
break;
}
}
@@ -123,6 +158,57 @@ class EpgApi {
}
}
/// Programmes tirés du dump XMLTV du panneau, ou liste vide.
Future<List<Map<String, dynamic>>> _panelProgrammes(
PlaylistConfig playlist,
String channelId,
) async {
final key = '${playlist.dns}|${playlist.username}';
final service =
_panelXmltv.putIfAbsent(key, () => _panelXmltvBuilder(playlist));
return _fromXmltv(service, playlist, channelId);
}
/// Programmes obtenus en interrogeant le panneau chaîne par chaîne.
///
/// Chemin lent : deux actions possibles, plusieurs secondes chacune chez la
/// plupart des revendeurs. Il ne sert que lorsque le dump XMLTV du panneau
/// est absent ou ne couvre pas la chaîne.
Future<List<Map<String, dynamic>>> _xtreamProgrammes(
PlaylistConfig playlist,
String channelId,
) async {
// `get_simple_data_table` d'abord : c'est la seule action réellement
// universelle. `get_epg` n'existe pas sur beaucoup de panneaux — au lieu
// d'une erreur, ils renvoient poliment le payload d'authentification,
// sans champ `epg_listings`, ce qui produisait un guide vide impossible
// à distinguer d'une chaîne sans programme.
const actions = ['get_simple_data_table', 'get_short_epg'];
for (final action in actions) {
final url =
'${playlist.dns}/player_api.php?username=${playlist.username}'
'&password=${playlist.password}'
'&action=$action&stream_id=$channelId';
final response =
await _http.get(Uri.parse(url)).timeout(const Duration(seconds: 60));
if (response.statusCode != 200) continue;
Map<String, dynamic> parsed;
try {
parsed = transformEpgData(json.decode(response.body), channelId);
} catch (_) {
continue;
}
final programmes = (parsed['programmes'] as List)
.whereType<Map<String, dynamic>>()
.toList();
if (programmes.isNotEmpty) return programmes;
}
return const [];
}
/// Le guide contient-il un programme en cours ou à venir ?
///
/// Un panneau dont l'EPG est figé répond avec des centaines de programmes,
@@ -142,8 +228,16 @@ class EpgApi {
Future<List<Map<String, dynamic>>> _xmltvProgrammes(
PlaylistConfig playlist,
String channelId,
) =>
_fromXmltv(_xmltv, playlist, channelId);
/// Programmes d'une chaîne dans un index XMLTV quelconque (panneau ou
/// source externe), ou liste vide si l'index ne la connaît pas.
Future<List<Map<String, dynamic>>> _fromXmltv(
XmltvEpgService? xmltv,
PlaylistConfig playlist,
String channelId,
) async {
final xmltv = _xmltv;
if (xmltv == null) return const [];
try {
@@ -163,7 +257,7 @@ class EpgApi {
.map((p) => p.toJson(channelId))
.toList();
} catch (e) {
print('[EpgApi] repli XMLTV indisponible pour $channelId : $e');
print('[EpgApi] source XMLTV indisponible pour $channelId : $e');
return const [];
}
}
@@ -180,7 +274,7 @@ class EpgApi {
'?username=${playlist.username}&password=${playlist.password}'
'&action=get_live_streams';
final response =
await http.get(Uri.parse(url)).timeout(const Duration(seconds: 90));
await _http.get(Uri.parse(url)).timeout(const Duration(seconds: 90));
final epgIds = <String, String>{};
final names = <String, String>{};
+67 -16
View File
@@ -1,12 +1,48 @@
import 'dart:async';
import 'dart:convert';
import 'dart:io';
import 'dart:typed_data';
import 'package:shelf/shelf.dart';
import 'package:http/http.dart' as http;
import '../database/database.dart';
import '../middleware/auth_middleware.dart';
import '../models/playlist_config.dart';
import '../utils/log_redactor.dart';
import 'asset_failure_cache.dart';
/// PNG transparent 1×1, servi à la place d'un logo introuvable.
final Uint8List _placeholderPixel = base64Decode(
'iVBORw0KGgoAAAANSUhEUgAAAAEAAAABCAQAAAC1HAwCAAAAC0lEQVR42mNkYAAAAAYAAjCB0C8AAAAASUVORK5CYII=',
);
/// L'URL désigne-t-elle une image ou un logo ?
///
/// Ces URL échappent à l'allowlist de domaine (les revendeurs hébergent leurs
/// picons ailleurs que le panneau) et méritent un repli visuel plutôt qu'une
/// erreur propagée au navigateur.
bool _isStaticAssetUrl(Uri url) =>
url.path.endsWith('.png') ||
url.path.endsWith('.jpg') ||
url.path.endsWith('.jpeg') ||
url.path.endsWith('.gif') ||
url.path.endsWith('.webp') ||
url.path.endsWith('.ico') ||
url.path.contains('/picons/') ||
url.path.contains('/logos/');
/// Réponse de repli pour une image indisponible.
///
/// Le `max-age` est essentiel : sans lui le navigateur redemande le logo à
/// chaque affichage de la grille et le flot de requêtes mortes reprend
/// immédiatement, coupure serveur ou pas.
Response _placeholderImageResponse() => Response.ok(
_placeholderPixel,
headers: {
'content-type': 'image/png',
'cache-control': 'public, max-age=600',
'access-control-allow-origin': '*',
},
);
/// Returns true when [host] must never be proxied (loopback, private LAN,
/// link-local/cloud-metadata ranges) — SSRF protection for asset URLs that
@@ -39,6 +75,9 @@ class ProxyHandler {
final Map<String, (PlaylistConfig, DateTime)> _playlistCache = {};
static const _cacheDuration = Duration(minutes: 5);
/// Hôtes d'images réputés hors service (voir [AssetFailureCache]).
final AssetFailureCache _assetFailures = AssetFailureCache();
static const _allowedHeaders = [
'content-type',
'content-range',
@@ -141,14 +180,14 @@ class ProxyHandler {
// SSRF Protection - but allow images/static assets from any host
// Xtream providers often use separate CDN servers for picons/images
final isStaticAsset = targetUrl.path.endsWith('.png') ||
targetUrl.path.endsWith('.jpg') ||
targetUrl.path.endsWith('.jpeg') ||
targetUrl.path.endsWith('.gif') ||
targetUrl.path.endsWith('.webp') ||
targetUrl.path.endsWith('.ico') ||
targetUrl.path.contains('/picons/') ||
targetUrl.path.contains('/logos/');
final isStaticAsset = _isStaticAssetUrl(targetUrl);
// Hôte d'images déjà constaté mort : on sert le pixel tout de suite.
// Aucun appel sortant, aucune entrée de log — c'est précisément le
// flot qu'on cherche à éteindre.
if (isStaticAsset && _assetFailures.isDown(targetUrl.host)) {
return _placeholderImageResponse();
}
String? allowedHost;
if (!isStaticAsset) {
@@ -251,6 +290,23 @@ class ProxyHandler {
}
}
// Une image en erreur n'a rien à transmettre au navigateur : le
// relais du 503 déclenchait un nouveau cycle de requêtes à chaque
// rendu. On sert le pixel, mis en cache, et on compte l'échec.
if (isStaticAsset) {
if (response.statusCode >= 400) {
// Compté sur l'hôte demandé, pas sur la cible finale d'une
// redirection : c'est cette clé-là que consulte le garde-fou
// en tête de requête.
_assetFailures.recordFailure(targetUrl.host);
// Le corps d'erreur doit être consommé, sinon la connexion
// reste ouverte jusqu'au timeout.
unawaited(response.stream.drain<void>().catchError((_) {}));
return _placeholderImageResponse();
}
_assetFailures.recordSuccess(targetUrl.host);
}
// Build response headers from source response
final responseHeaders = <String, String>{
'access-control-allow-origin': '*',
@@ -283,14 +339,9 @@ class ProxyHandler {
);
// Return transparent 1x1 pixel image fallback for images
if (targetUrl?.path.endsWith('.png') == true ||
targetUrl?.path.endsWith('.jpg') == true) {
return Response.ok(
base64Decode(
'iVBORw0KGgoAAAANSUhEUgAAAAEAAAABCAQAAAC1HAwCAAAAC0lEQVR42mNkYAAAAAYAAjCB0C8AAAAASUVORK5CYII=',
),
headers: {'content-type': 'image/png'},
);
if (targetUrl != null && _isStaticAssetUrl(targetUrl)) {
_assetFailures.recordFailure(targetUrl.host);
return _placeholderImageResponse();
}
return Response.internalServerError(
+8 -13
View File
@@ -1,4 +1,3 @@
import 'dart:async';
import 'dart:io';
import 'package:shelf/shelf.dart';
import 'package:shelf_router/shelf_router.dart';
@@ -8,6 +7,7 @@ import '../models/playlist_config.dart';
import '../services/ffmpeg_session_manager.dart';
import '../utils/log_redactor.dart';
import '../utils/media_probe.dart';
import '../utils/stream_pipe.dart';
import 'recording_playlist.dart';
/// Directory for temporary HLS segments
@@ -430,22 +430,17 @@ Handler createLiveStreamHandler(
}
});
// Relayer stdout vers le client ; tuer FFmpeg dès que le client zappe
// ou ferme l'onglet (sinon les processus s'accumulent à chaque zap).
final controller = StreamController<List<int>>();
final subscription = process.stdout.listen(
controller.add,
onError: controller.addError,
onDone: controller.close,
// Relayer stdout vers le client en gardant la contre-pression, et tuer
// FFmpeg dès que le client zappe ou ferme l'onglet (sinon les processus
// s'accumulent à chaque zap).
final body = pipeWithBackpressure(
process.stdout,
onStop: () => process.kill(ProcessSignal.sigterm),
);
controller.onCancel = () {
subscription.cancel();
process.kill(ProcessSignal.sigterm);
};
return Response(
200,
body: controller.stream,
body: body,
headers: {
'Content-Type': 'video/mp2t',
'Cache-Control': 'no-store',
+39 -4
View File
@@ -19,6 +19,9 @@ class XmltvEpgService {
required this.sourceUrls,
this.refreshInterval = const Duration(hours: 6),
this.retention = const Duration(hours: 6),
this.horizon = const Duration(hours: 48),
this.retryBackoff = const Duration(minutes: 15),
this.downloadTimeout = const Duration(minutes: 5),
http.Client? client,
}) : _client = client ?? http.Client();
@@ -33,11 +36,32 @@ class XmltvEpgService {
/// garder double la taille de l'index.
final Duration retention;
/// Les programmes qui commencent au-delà de cet horizon sont écartés à
/// l'indexation.
///
/// Un dump national couvre souvent sept jours pour un millier de chaînes ;
/// tout garder en mémoire dans le conteneur coûte des centaines de
/// mégaoctets pour un guide qui n'affiche que le programme courant et les
/// suivants.
final Duration horizon;
/// Délai minimal entre deux tentatives quand aucune source n'a répondu.
///
/// Sans ce recul, une source morte était retéléchargée à chaque consultation
/// du guide — une requête sortante par chaîne affichée, soit exactement le
/// flot qu'un index est censé éviter.
final Duration retryBackoff;
/// Plafond de téléchargement d'un dump. Le guide est facultatif : mieux vaut
/// abandonner et servir le panneau que faire attendre l'utilisateur.
final Duration downloadTimeout;
final http.Client _client;
/// clé de chaîne normalisée → programmes triés par heure de début.
Map<String, List<XmltvProgramme>> _index = {};
DateTime? _indexedAt;
DateTime? _failedAt;
Future<void>? _refreshInFlight;
bool get hasData => _index.isNotEmpty;
@@ -70,6 +94,14 @@ class XmltvEpgService {
? null
: DateTime.now().difference(_indexedAt!);
if (age != null && age < refreshInterval) return Future.value();
final sinceFailure = _failedAt == null
? null
: DateTime.now().difference(_failedAt!);
if (sinceFailure != null && sinceFailure < retryBackoff) {
return Future.value();
}
return _refreshInFlight ??= _refresh().whenComplete(() {
_refreshInFlight = null;
});
@@ -101,19 +133,20 @@ class XmltvEpgService {
if (ok == 0) {
// Garder l'index précédent plutôt que de servir un guide vide.
_failedAt = DateTime.now();
print('[XmltvEpg] aucune source disponible, index précédent conservé');
return;
}
_index = merged;
_indexedAt = DateTime.now();
_failedAt = null;
print('[XmltvEpg] index prêt : ${merged.length} chaînes');
}
Future<String> _download(String url) async {
final response = await _client
.get(Uri.parse(url))
.timeout(const Duration(minutes: 5));
final response =
await _client.get(Uri.parse(url)).timeout(downloadTimeout);
if (response.statusCode != 200) {
throw HttpException('HTTP ${response.statusCode}');
}
@@ -132,6 +165,7 @@ class XmltvEpgService {
Map<String, List<XmltvProgramme>> _parse(String xml) {
final cutoff = DateTime.now().toUtc().subtract(retention);
final limit = DateTime.now().toUtc().add(horizon);
final byChannel = <String, List<XmltvProgramme>>{};
// Alias : plusieurs dumps déclarent <channel id="X"> avec un
@@ -191,7 +225,8 @@ class XmltvEpgService {
if (programmeChannel != null &&
start != null &&
stop != null &&
stop.isAfter(cutoff)) {
stop.isAfter(cutoff) &&
start.isBefore(limit)) {
final key = normalizeKey(programmeChannel);
if (key.isNotEmpty) {
byChannel.putIfAbsent(key, () => []).add(
+57
View File
@@ -0,0 +1,57 @@
import 'package:test/test.dart';
import '../api/asset_failure_cache.dart';
void main() {
group('AssetFailureCache', () {
test('un hôte inconnu est considéré joignable', () {
final cache = AssetFailureCache();
expect(cache.isDown('cdn.example'), isFalse);
});
test('un échec isolé ne coupe pas l’hôte', () {
// Un 5xx passager ne doit pas priver l'utilisateur de ses logos
// pendant plusieurs minutes.
final cache = AssetFailureCache(threshold: 3);
cache.recordFailure('cdn.example');
expect(cache.isDown('cdn.example'), isFalse);
});
test('coupe l’hôte après le seuil d’échecs consécutifs', () {
final cache = AssetFailureCache(threshold: 3);
for (var i = 0; i < 3; i++) {
cache.recordFailure('cdn.example');
}
expect(cache.isDown('cdn.example'), isTrue);
});
test('un succès remet le compteur à zéro', () {
final cache = AssetFailureCache(threshold: 3);
cache.recordFailure('cdn.example');
cache.recordFailure('cdn.example');
cache.recordSuccess('cdn.example');
cache.recordFailure('cdn.example');
expect(cache.isDown('cdn.example'), isFalse);
});
test('la coupure expire après le TTL', () {
var now = DateTime(2026, 9, 19, 9);
final cache = AssetFailureCache(
threshold: 1,
ttl: const Duration(minutes: 5),
clock: () => now,
);
cache.recordFailure('cdn.example');
expect(cache.isDown('cdn.example'), isTrue);
now = now.add(const Duration(minutes: 6));
expect(cache.isDown('cdn.example'), isFalse);
});
test('les hôtes sont indépendants', () {
final cache = AssetFailureCache(threshold: 1);
cache.recordFailure('mort.example');
expect(cache.isDown('mort.example'), isTrue);
expect(cache.isDown('vivant.example'), isFalse);
});
});
}
+147
View File
@@ -1,5 +1,12 @@
import 'dart:convert';
import 'package:http/http.dart' as http;
import 'package:http/testing.dart';
import 'package:shelf/shelf.dart';
import 'package:test/test.dart';
import '../api/epg_api.dart';
import '../models/playlist_config.dart';
import '../services/xmltv_epg_service.dart';
/// Extrait réel d'une réponse `get_simple_data_table` (panneau Xtream) :
/// titres/descriptions en base64, horodatages epoch UTC doublés d'une chaîne
@@ -71,4 +78,144 @@ void main() {
expect(result['programmes'], isEmpty);
});
});
group('handleGetEpg', () {
final playlist = PlaylistConfig(
id: 'p1',
name: 'Test',
dns: 'http://panel.example',
username: 'u',
password: 'p',
createdAt: DateTime(2026, 1, 1),
);
String xmltvDump(DateTime start, DateTime stop) {
String stamp(DateTime d) {
final u = d.toUtc();
String two(int v) => v.toString().padLeft(2, '0');
return '${u.year}${two(u.month)}${two(u.day)}'
'${two(u.hour)}${two(u.minute)}${two(u.second)} +0000';
}
return '''
<?xml version="1.0" encoding="utf-8" ?>
<tv>
<channel id="France2.fr"><display-name>FRANCE 2</display-name></channel>
<programme start="${stamp(start)}" stop="${stamp(stop)}" channel="France2.fr">
<title lang="fr">Journal de 20h</title>
<desc lang="fr">L'info du soir</desc>
</programme>
</tv>
''';
}
Request requestFor(String channelId) =>
Request('GET', Uri.parse('http://localhost/api/epg/$channelId'));
test('sert le guide depuis le dump XMLTV du panneau, sans appel par '
'chaîne', () async {
// Un appel `player_api` par chaîne coûte plusieurs secondes chez le
// fournisseur : le dump du panneau couvre toutes les chaînes d'un coup.
final now = DateTime.now().toUtc();
final actions = <String>[];
final client = MockClient((request) async {
final url = request.url;
if (url.path.endsWith('/xmltv.php')) {
return http.Response(
xmltvDump(now.subtract(const Duration(minutes: 5)),
now.add(const Duration(minutes: 25))),
200,
);
}
final action = url.queryParameters['action'] ?? '';
actions.add(action);
if (action == 'get_live_streams') {
return http.Response(
jsonEncode([
{
'stream_id': 845452,
'name': 'FRANCE 2',
'epg_channel_id': 'France2.fr',
},
]),
200,
);
}
return http.Response('{}', 200);
});
final api = EpgApi(
(_) async => playlist,
httpClient: client,
panelXmltvBuilder: (config) => XmltvEpgService(
sourceUrls: ['${config.dns}/xmltv.php'],
client: client,
),
);
final response = await api.handleGetEpg(requestFor('845452'), '845452');
final body = jsonDecode(await response.readAsString()) as Map;
expect(response.statusCode, 200);
expect((body['programmes'] as List).single['title'], 'Journal de 20h');
expect(response.headers['X-Epg-Source'], 'panel-xmltv');
// Seule la table des chaînes est interrogée : aucune action EPG.
expect(actions, ['get_live_streams']);
});
test('retombe sur player_api quand le dump ne couvre pas la chaîne',
() async {
final actions = <String>[];
// Heure locale volontairement : le panneau émet des horodatages sans
// fuseau, que le serveur relit comme de l'heure locale.
final future = DateTime.now().add(const Duration(minutes: 10));
String panelStamp(DateTime d) =>
d.toIso8601String().substring(0, 19).replaceFirst('T', ' ');
final client = MockClient((request) async {
final url = request.url;
if (url.path.endsWith('/xmltv.php')) {
return http.Response('<tv></tv>', 200);
}
final action = url.queryParameters['action'] ?? '';
actions.add(action);
if (action == 'get_live_streams') {
return http.Response(jsonEncode(const []), 200);
}
if (action == 'get_simple_data_table') {
return http.Response(
jsonEncode({
'epg_listings': [
{
'title': base64Encode(utf8.encode('Match')),
'description': '',
'start': panelStamp(future),
'end': panelStamp(future.add(const Duration(hours: 2))),
},
],
}),
200,
);
}
return http.Response('{}', 200);
});
final api = EpgApi(
(_) async => playlist,
httpClient: client,
panelXmltvBuilder: (config) => XmltvEpgService(
sourceUrls: ['${config.dns}/xmltv.php'],
client: client,
),
);
final response = await api.handleGetEpg(requestFor('999'), '999');
final body = jsonDecode(await response.readAsString()) as Map;
expect((body['programmes'] as List).single['title'], 'Match');
expect(response.headers['X-Epg-Source'], 'xtream');
expect(actions, contains('get_simple_data_table'));
});
});
}
+77
View File
@@ -0,0 +1,77 @@
import 'dart:async';
import 'package:test/test.dart';
import '../utils/stream_pipe.dart';
void main() {
group('pipeWithBackpressure', () {
test('relaie les données au client', () async {
final source = StreamController<List<int>>();
final piped = pipeWithBackpressure(source.stream, onStop: () {});
final received = <List<int>>[];
final done = piped.listen(received.add).asFuture<void>();
source.add([1, 2, 3]);
await source.close();
await done;
expect(received, [
[1, 2, 3]
]);
});
test('propage la pause du client vers la source', () async {
// Sans cette propagation, un client plus lent que le flux laissait
// FFmpeg produire à pleine vitesse et le serveur empilait les paquets
// en mémoire : latence qui dérive et lecture qui saccade.
var paused = false;
var resumed = false;
final source = StreamController<List<int>>(
onPause: () => paused = true,
onResume: () => resumed = true,
);
final subscription = pipeWithBackpressure(
source.stream,
onStop: () {},
).listen((_) {});
subscription.pause();
await Future<void>.delayed(Duration.zero);
expect(paused, isTrue);
subscription.resume();
await Future<void>.delayed(Duration.zero);
expect(resumed, isTrue);
await subscription.cancel();
await source.close();
});
test('arrête le producteur quand le client se déconnecte', () async {
var stopped = 0;
final source = StreamController<List<int>>();
final subscription =
pipeWithBackpressure(source.stream, onStop: () => stopped++)
.listen((_) {});
await subscription.cancel();
expect(stopped, 1);
await source.close();
});
test('arrête le producteur quand la source se termine', () async {
var stopped = 0;
final source = StreamController<List<int>>();
final piped = pipeWithBackpressure(source.stream, onStop: () => stopped++);
final done = piped.listen((_) {}).asFuture<void>();
await source.close();
await done;
expect(stopped, 1);
});
});
}
+63
View File
@@ -1,3 +1,5 @@
import 'package:http/http.dart' as http;
import 'package:http/testing.dart';
import 'package:test/test.dart';
import '../services/xmltv_epg_service.dart';
@@ -101,4 +103,65 @@ void main() {
expect(index, isEmpty);
});
});
group('ensureFresh', () {
test('ne retélécharge pas tant que l’index est frais', () async {
final now = DateTime.now().toUtc();
var calls = 0;
final service = XmltvEpgService(
sourceUrls: const ['http://dump.example/epg.xml'],
client: MockClient((_) async {
calls++;
return http.Response(
_fixture(now, now.add(const Duration(minutes: 30))),
200,
);
}),
);
await service.ensureFresh();
await service.ensureFresh();
expect(calls, 1);
expect(service.hasData, isTrue);
});
test('espace les tentatives après un échec', () async {
// Sans ce recul, une source morte était retéléchargée à chaque
// consultation du guide : une requête sortante par chaîne affichée.
var calls = 0;
final service = XmltvEpgService(
sourceUrls: const ['http://dump.example/epg.xml'],
retryBackoff: const Duration(minutes: 15),
client: MockClient((_) async {
calls++;
return http.Response('nope', 500);
}),
);
await service.ensureFresh();
await service.ensureFresh();
await service.ensureFresh();
expect(calls, 1);
expect(service.hasData, isFalse);
});
});
group('horizon', () {
test('écarte les programmes au-delà de l’horizon d’indexation', () {
// Un dump national couvre sept jours : tout garder ferait grossir
// l'index pour un guide qui n'affiche que le programme courant.
final service = XmltvEpgService(
sourceUrls: const [],
horizon: const Duration(hours: 48),
);
final far = DateTime.now().toUtc().add(const Duration(days: 5));
expect(
service.parseForTest(_fixture(far, far.add(const Duration(hours: 1)))),
isEmpty,
);
});
});
}
+49
View File
@@ -0,0 +1,49 @@
import 'dart:async';
/// Relaie [source] vers le client en conservant la contre-pression.
///
/// POURQUOI : un `StreamController` nu ne transmet pas la pause de son
/// abonné à la source. Quand le navigateur lit plus lentement que FFmpeg ne
/// produit — réseau domestique, onglet en arrière-plan —, le serveur empilait
/// donc les paquets en mémoire sans jamais ralentir le producteur : la
/// consommation grimpe et la lecture dérive derrière le direct, ce qui se voit
/// à l'écran comme des saccades.
///
/// [onStop] est appelé une seule fois, quand le client se déconnecte ou que la
/// source se termine : c'est là qu'on tue le processus FFmpeg, sinon il en
/// reste un par zapping.
Stream<List<int>> pipeWithBackpressure(
Stream<List<int>> source, {
required void Function() onStop,
}) {
late final StreamController<List<int>> controller;
late final StreamSubscription<List<int>> subscription;
var stopped = false;
void stop() {
if (stopped) return;
stopped = true;
onStop();
}
controller = StreamController<List<int>>(
onPause: () => subscription.pause(),
onResume: () => subscription.resume(),
onCancel: () {
final cancelled = subscription.cancel();
stop();
return cancelled;
},
);
subscription = source.listen(
controller.add,
onError: controller.addError,
onDone: () {
controller.close();
stop();
},
);
return controller.stream;
}