Compare commits

...
Author SHA1 Message Date
MichaelandClaude Opus 5.5 775caad0ee Annule « Zap : ouvrir les chaînes par le HLS du panneau »
Mesure après déploiement : premier octet de turbo.ts à 4,5–4,8 s avec la
source HLS, contre 4,3–4,7 s en .ts — aucun gain. Les logs confirment
que le HLS était bien utilisé (aucun repli). La sonde qui annonçait
0,5 s ouvrait le .ts AVANT le .m3u8 : la chaîne était déjà démarrée chez
le panneau. Les ~4,5 s sont le temps de démarrage d'une chaîne chez le
fournisseur, quel que soit le format.

Retour au .ts, déjà validé en lecture, plutôt que de garder un chemin
plus complexe (segments de 10 s, repli) sans bénéfice mesuré.

This reverts commit 0ea60cc.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-10-10 23:30:51 +02:00
MichaelandClaude Opus 5.5 4d00409de8 Guide TV : régionales France 3, variantes de nom, chaînes +1, CH et BE
Mesuré en prod après le premier correctif : 69 chaînes sur 149 avec
guide. Les vides restantes relevaient surtout de formes de nom que le
dump écrit autrement, ou de pays non couverts :

- « F3 ALPES » → France 3 Alpes (France.3.-.Alpes.fr) : F2 à F5 développés
- clé souple en dernier recours, dans son propre espace de noms :
  sans code pays ni « channel/tv/hd » (« AB3 » → AB3.Channel.fr)
- chaîne « +1 » absente du dump (TF1 +1) : guide de la chaîne principale
  décalé d'une heure ; une +1 qui a son propre guide le garde
- sources par défaut : FR1 puis CH1 (RTS…) et BE2 (même fournisseur),
  le dump français reste prioritaire

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-10-10 23:17:28 +02:00
MichaelandClaude Opus 5.5 0ea60ccbfa Zap : ouvrir les chaînes par le HLS du panneau, ~4 s gagnées
Mesuré en prod avec la sonde admin : le panneau met 4,2 à 4,4 s à
répondre la redirection d'un .ts, contre 0,2 à 1,3 s pour un .m3u8 dont
le premier segment arrive ensuite en 50 ms.

turbo.ts et les sessions HLS live lisent désormais le .m3u8 du panneau
(-live_start_index -2 : deux segments d'avance d'un coup), avec repli
automatique sur le .ts si le HLS ne livre rien — dans la même réponse
pour turbo.ts, le lecteur ne voit qu'un démarrage plus long. Pas de
reconnect_at_eof en HLS : chaque segment finit par un EOF (leçon de
xtremobile). Le navigateur reçoit toujours le même flux MPEG-TS.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-10-10 23:08:08 +02:00
MichaelandClaude Opus 5.5 a3640a6b5f Diagnostic admin : chronométrer l'ouverture d'une chaîne chez le panneau
~4,3 s des ~4,6 s d'un zap se passent avant le premier octet du panneau.
Le compte autorise aussi le HLS (allowed_output_formats : m3u8, ts) ;
GET /api/admin/upstream-probe?stream=<id> mesure, pour .ts et .m3u8,
chaque redirection, les en-têtes, le premier octet et le premier segment,
pour décider du format à utiliser. Admin uniquement ; seuls l'hôte et les
durées sont rendus, jamais les identifiants ni les jetons.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-10-10 22:56:34 +02:00
MichaelandClaude Opus 5.5 584068b45b Guide TV : décoder la table des chaînes en UTF-8
Sans charset annoncé, Response.body pouvait décoder le JSON du panneau en
latin-1 et abîmer les noms accentués ou décorés (« Chérie », « ◉ »), donc
leur correspondance avec le guide XMLTV. Le test simule désormais des
octets UTF-8 sans charset, comme le panneau réel.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-10-10 22:54:46 +02:00
MichaelandClaude Opus 5.5 abe7d522c1 Guide TV : plus jamais de 500, et bien plus de chaînes couvertes
Constaté en prod : chaque /api/epg répondait 500 en 4 s. Un appel
player_api qui lève (connexion coupée, délai) n'était pas rattrapé, et le
dump XMLTV externe, qui couvre pourtant TF1, France 2, M6…, n'était jamais
consulté. Et 59 % des chaînes françaises n'ont pas d'epg_channel_id : leur
nom brut (« FR - TF1 FHD ◉ ») ne correspondait à aucune entrée XMLTV.

- échec player_api rattrapé : la chaîne passe à la source suivante
- dump externe (déjà en mémoire) avant l'appel lent chaîne par chaîne
- correspondance par nom nettoyé : préfixe pays, FHD/HD/UHD/4K/HEVC…,
  symboles retirés, forme « nom.pays » tentée ; « +1 » conservé
- tous les display-name XMLTV indexés, plus seulement le premier
- table des chaînes : un seul téléchargement partagé, décodé hors de la
  boucle principale, et un échec n'est plus gardé 6 h (2 min)
- guide vide gardé 5 min au lieu de 30

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-10-10 22:51:24 +02:00
MichaelandClaude Opus 5.5 8d79472ca8 Démarrage : la lecture du quota ne fait plus attendre le flux
AccountLimits interrogeait player_api.php avant d'ouvrir le flux quand
sa valeur n'était pas en cache (premier flux, puis toutes les 10 min).
Le panneau met jusqu'à 4 s à répondre (mesuré en prod) : autant de
retard au démarrage d'une chaîne ou d'un film.

La valeur connue, même périmée, ou le repli (1) est rendue tout de suite ;
la lecture chez le panneau se fait en arrière-plan, une seule à la fois.
Sans risque : le quota ne sert qu'à choisir quels flux orphelins couper.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-10-10 22:34:26 +02:00
MichaelandClaude Opus 5.5 6812a4964d Lecture : ne couper que les flux orphelins, jamais un flux regardé
La règle « le dernier flux gagne » coupait aussi un flux en cours de
lecture. Deux lecteurs ouverts sur un compte à une connexion se coupaient
alors en boucle (chacun relance aussitôt) : plus aucune image nulle part.
Constaté en prod : tant qu'un autre onglet lisait, chaque ouverture
échouait après 5 s sans un octet.

Une connexion n'est désormais coupée que si son spectateur ne donne plus
signe de vie depuis 15 s : aucune requête de playlist/segment (HLS, film)
ni octet remis au client (turbo.ts). Les orphelins partent toujours, un
second lecteur actif est refusé par le panneau comme avant.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-10-10 22:13:58 +02:00
MichaelandClaude Opus 5.5 fdf98f9659 Films : demander le fichier sous sa vraie extension
Le serveur réclamait toujours <id>.mkv au panneau. Un film stocké en .mp4
était refusé (HTTP 551 constaté en prod). Le client transmet désormais
container_extension (?ext=), validé côté serveur ; mkv reste le défaut
pour un client qui ne l'envoie pas. Même approche que xtremobile, qui lit
ces films sans souci.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-10-10 12:09:11 +02:00
MichaelandClaude Opus 5.5 6212948ea5 Déploiement : versionner les JS pour contourner le cache du proxy
main.dart.js, flutter_bootstrap.js et xf-player-core.js gardent le même
nom d'une build à l'autre. En prod, Nginx Proxy Manager les met en cache
(Cache Assets) jusqu'à 12 h, même sur rechargement forcé : après le
dernier déploiement, l'ancienne app et l'ancien lecteur tournaient avec
le nouveau serveur.

Au démarrage, le serveur réécrit index.html, flutter_bootstrap.js et
player*.html pour référencer ces fichiers avec ?v=<empreinte du contenu> :
chaque version a sa propre URL, aucun cache ne peut les confondre.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-10-10 12:09:11 +02:00
MichaelandClaude Opus 5.5 662e1034e5 Lecture : couper le flux orphelin avant d'en ouvrir un nouveau
L'abonnement n'autorise qu'une connexion simultanée (max_connections = 1,
vérifié en prod). Or un FFmpeg de lecture survivait au spectateur : 4 min
pour un direct HLS, jusqu'à 15 min ou la fin du téléchargement pour un
film, le temps de détecter la déconnexion pour turbo.ts. Le flux suivant
trouvait le slot pris : HTTP 551 ou timeout côté panneau (film qui ne
démarre pas, zap qui échoue).

Avant d'ouvrir une connexion amont, le serveur coupe désormais les plus
anciennes connexions de lecture du même compte, dans la limite du quota
max_connections lu chez le fournisseur (cache 10 min, 1 par défaut). Les
enregistrements ne sont jamais coupés.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-10-10 11:51:56 +02:00
MichaelandClaude Opus 5.5 a9b43ff581 Lecture TV : fin des coupures en rafale sur le direct
Les panneaux Xtream livrent le flux par rafales (~6 s de contenu puis
plusieurs secondes de silence). Le rattrapage mpegts.js réglé à 5 s max /
1 s de marge sautait en avant à chaque rafale puis tombait à sec : mesuré
en prod, 12 sauts et 14 coupures en 30 s sur une chaîne FHD.

Marge portée à 12 s, saut seulement au-delà de 30 s de retard, retard
résorbé par liveSync (x1,1) et purge du buffer déjà lu. Rejoué en prod
avec cette config : 0 coupure en continu.

Ajoute des tests Node du moteur web, lancés par la CI.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-10-10 11:37:55 +02:00
MichaelandClaude Opus 5.5 327873d254 Journaux : masquer les identifiants Xtream des URL et exceptions
Le repli XMLTV journalisait l'URL `xmltv.php?username=…&password=…` du
panneau en clair, et le texte des `ClientException` / `ProcessException`
recopiait aussi l'URL : les identifiants de l'abonné finissaient dans
`docker logs`. Passage par `LogRedactor.redactUrl` dans XmltvEpgService,
EpgApi, le relais live, FFmpegSessionManager et RecordingScheduler.

Test : capture des `print` via une Zone pour vérifier le masquage.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-10-10 09:25:00 +02:00
MichaelandClaude Opus 5.5 da67c479b3 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>
2026-10-10 09:11:54 +02:00
LogiFlow 9f4e98fcea Merge pull request #19 from R0m1k3/chore/lockfiles-3.47
pubspec.lock régénéré avec Flutter 3.47.6
2026-10-08 22:12:55 +02:00
23 changed files with 1732 additions and 75 deletions

No files matched your search

+3
View File
@@ -31,6 +31,9 @@ jobs:
# fatal; errors and warnings still fail the build.
- run: flutter analyze --no-fatal-infos
- run: flutter test
# Moteur de lecture web (web/xf-player-core.js) : Node est préinstallé
# sur les runners ubuntu, aucune dépendance npm.
- run: node --test "test/web/*_test.mjs"
- run: flutter build web --release
backend:
+1 -1
View File
@@ -37,7 +37,7 @@ Au premier démarrage, un compte `admin` est créé avec un **mot de passe aléa
| `RECORDINGS_PATH` | `./data/recordings` | Dossier hôte des enregistrements |
| `TZ` | `Europe/Paris` | Fuseau du conteneur (les enregistrements sont stockés en UTC) |
| `MAX_CONCURRENT_RECORDINGS` | `2` | Enregistrements simultanés |
| `EPG_XMLTV_URLS` | dump FR | Sources XMLTV de repli (vide = aucun appel sortant) |
| `EPG_XMLTV_URLS` | dumps FR, CH, BE | Sources XMLTV de repli (vide = aucun appel sortant) |
| `NVIDIA_GPU` | `false` | Transcodage NVENC |
| `ADMIN_INITIAL_PASSWORD` | *(généré)* | Mot de passe initial du compte admin |
| `MIN_FREE_DISK_MB` | `500` | Espace libre minimal pour démarrer une capture |
+97 -25
View File
@@ -1,4 +1,5 @@
import 'dart:convert';
import 'dart:isolate';
import 'package:shelf/shelf.dart';
import 'package:http/http.dart' as http;
import '../models/playlist_config.dart';
@@ -97,10 +98,15 @@ class EpgApi {
// 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.
//
// Le dump externe passe ensuite AVANT l'interrogation chaîne par
// chaîne : il est déjà indexé en mémoire (réponse instantanée), alors
// que `player_api` coûte plusieurs secondes par chaîne et échouait en
// prod — ce qui laissait toute la grille sans guide.
final sources = <(String, Future<List<Map<String, dynamic>>> Function())>[
('panel-xmltv', () => _panelProgrammes(playlist, channelId)),
('xtream', () => _xtreamProgrammes(playlist, channelId)),
('xmltv', () => _xmltvProgrammes(playlist, channelId)),
('xtream', () => _xtreamProgrammes(playlist, channelId)),
];
var epgData = <String, dynamic>{
@@ -133,10 +139,13 @@ class EpgApi {
final jsonStr = json.encode(epgData);
// Mettre en cache 30 minutes
// 30 min pour un guide trouvé ; 5 min seulement pour un guide vide,
// qui vient souvent d'une source momentanément indisponible et ne
// doit pas priver la chaîne de guide pendant une demi-heure.
_cache[cacheKey] = _CacheEntry(
data: jsonStr,
expiresAt: DateTime.now().add(const Duration(minutes: 30)),
expiresAt: DateTime.now()
.add(Duration(minutes: found ? 30 : 5)),
);
return Response.ok(
@@ -191,8 +200,19 @@ class EpgApi {
'&password=${playlist.password}'
'&action=$action&stream_id=$channelId';
final response =
await _http.get(Uri.parse(url)).timeout(const Duration(seconds: 60));
final http.Response response;
try {
response = await _http
.get(Uri.parse(url))
.timeout(const Duration(seconds: 20));
} catch (e) {
// Délai dépassé, connexion coupée… : une source parmi d'autres, pas
// une raison de répondre 500 (ce qui privait la chaîne des sources
// suivantes). L'exception recopie l'URL, identifiants compris.
print('[EpgApi] player_api $action indisponible pour $channelId : '
'${LogRedactor.redactUrl('$e')}');
continue;
}
if (response.statusCode != 200) continue;
Map<String, dynamic> parsed;
@@ -257,46 +277,98 @@ class EpgApi {
.map((p) => p.toJson(channelId))
.toList();
} catch (e) {
print('[EpgApi] source XMLTV indisponible pour $channelId : $e');
// L'exception peut recopier l'URL `player_api`/`xmltv.php`, identifiants
// compris.
print(
'[EpgApi] source XMLTV indisponible pour $channelId : '
'${LogRedactor.redactUrl('$e')}',
);
return const [];
}
}
/// `stream_id` → (`epg_channel_id`, nom) depuis la réponse
/// `get_live_streams`.
static (Map<String, String>, Map<String, String>) _parseChannelTable(
String body,
) {
final epgIds = <String, String>{};
final names = <String, String>{};
final decoded = json.decode(body);
if (decoded is List) {
for (final item in decoded) {
if (item is! Map) continue;
final streamId = item['stream_id']?.toString();
if (streamId == null || streamId.isEmpty) continue;
final epgId = item['epg_channel_id']?.toString();
if (epgId != null && epgId.isNotEmpty) epgIds[streamId] = epgId;
final name = item['name']?.toString();
if (name != null && name.isNotEmpty) names[streamId] = name;
}
}
return (epgIds, names);
}
/// Téléchargements de la table des chaînes en cours, par compte.
final Map<String, Future<_ChannelMap>> _channelMapsInFlight = {};
/// Table des chaînes du compte, rafraîchie toutes les 6 heures.
Future<_ChannelMap> _channelMapFor(PlaylistConfig playlist) async {
///
/// Les requêtes simultanées partagent un seul téléchargement : la grille
/// demande le guide de 30 chaînes d'un coup, et chacune déclenchait sinon
/// son propre `get_live_streams` (7 Mo en prod).
Future<_ChannelMap> _channelMapFor(PlaylistConfig playlist) {
final key = '${playlist.dns}|${playlist.username}';
final cached = _channelMaps[key];
if (cached != null && DateTime.now().isBefore(cached.expiresAt)) {
return cached;
return Future.value(cached);
}
return _channelMapsInFlight[key] ??=
_loadChannelMap(playlist, key).whenComplete(() {
_channelMapsInFlight.remove(key);
});
}
Future<_ChannelMap> _loadChannelMap(
PlaylistConfig playlist,
String key,
) async {
final url = '${playlist.dns}/player_api.php'
'?username=${playlist.username}&password=${playlist.password}'
'&action=get_live_streams';
final response =
await _http.get(Uri.parse(url)).timeout(const Duration(seconds: 90));
http.Response? response;
try {
response =
await _http.get(Uri.parse(url)).timeout(const Duration(seconds: 90));
} catch (e) {
print('[EpgApi] table des chaînes indisponible : '
'${LogRedactor.redactUrl('$e')}');
}
final epgIds = <String, String>{};
final names = <String, String>{};
if (response.statusCode == 200) {
final decoded = json.decode(response.body);
if (decoded is List) {
for (final item in decoded) {
if (item is! Map) continue;
final streamId = item['stream_id']?.toString();
if (streamId == null || streamId.isEmpty) continue;
final epgId = item['epg_channel_id']?.toString();
if (epgId != null && epgId.isNotEmpty) epgIds[streamId] = epgId;
final name = item['name']?.toString();
if (name != null && name.isNotEmpty) names[streamId] = name;
}
var epgIds = <String, String>{};
var names = <String, String>{};
if (response != null && response.statusCode == 200) {
// Décodage dans un isolate : 7 Mo de JSON sur la boucle principale la
// gelaient, et le relais turbo.ts cessait d'émettre pendant ce temps.
// UTF-8 explicite : sans charset annoncé, `response.body` peut décoder
// en latin-1 et abîmer les noms (« Chérie », « ◉ »), donc leur
// correspondance avec le guide.
final body = utf8.decode(response.bodyBytes, allowMalformed: true);
try {
(epgIds, names) = await Isolate.run(() => _parseChannelTable(body));
} catch (e) {
print('[EpgApi] table des chaînes illisible : ${e.runtimeType}');
}
}
// Une table vide (panneau en échec) n'est gardée que 2 minutes : la
// garder 6 h privait TOUTES les chaînes du repli XMLTV pendant 6 h.
final map = _ChannelMap(
epgIds: epgIds,
names: names,
expiresAt: DateTime.now().add(const Duration(hours: 6)),
expiresAt: DateTime.now().add(
names.isEmpty ? const Duration(minutes: 2) : const Duration(hours: 6),
),
);
_channelMaps[key] = map;
return map;
+130 -8
View File
@@ -6,6 +6,7 @@ import '../database/database.dart';
import '../models/playlist_config.dart';
import '../services/ffmpeg_session_manager.dart';
import '../services/gpu_fallback.dart';
import '../services/upstream_slots.dart';
import '../utils/log_redactor.dart';
import '../utils/media_probe.dart';
import '../utils/stream_pipe.dart';
@@ -21,6 +22,61 @@ final FfmpegSessionManager sessionManager = FfmpegSessionManager(_hlsTempDir);
/// Pannes GPU récentes, partagées par toutes les routes de transcodage.
final GpuHealth _gpuHealth = GpuHealth();
/// Connexions de lecture ouvertes chez le fournisseur, par compte.
final UpstreamSlots _upstreamSlots = UpstreamSlots();
/// Numérote les flux `turbo.ts` (un processus FFmpeg par requête).
int _turboCounter = 0;
/// Quota `max_connections` de chaque compte Xtream.
final AccountLimits _accountLimits = AccountLimits();
/// Fait de la place chez le fournisseur avant d'ouvrir la connexion [id].
///
/// Voir [UpstreamSlots] : sur un compte à une connexion, le flux précédent
/// (souvent un FFmpeg orphelin que plus personne ne regarde) bloquait le
/// suivant.
Future<void> _claimUpstream(PlaylistConfig playlist, String id) async {
final max = await _accountLimits.maxFor(playlist);
final freed =
_upstreamSlots.makeRoom(accountKeyOf(playlist), max: max, keep: id);
if (freed.isEmpty) return;
print('[Upstream] $id : ${freed.join(', ')} coupé(s) '
'(quota fournisseur : $max connexion(s))');
// Laisser le panneau constater la fermeture : rouvrir aussitôt se fait
// refuser sur les comptes à une seule connexion.
await Future<void>.delayed(const Duration(seconds: 1));
}
/// Délai sans consommation au-delà duquel un flux est tenu pour orphelin.
///
/// Un lecteur HLS redemande sa playlist toutes les 2 à 4 s ; un flux
/// `turbo.ts` reçoit au moins une rafale toutes les ~7 s même sur un panneau
/// qui livre par à-coups. 15 s laissent une marge à ces deux rythmes.
const _viewerGrace = Duration(seconds: 15);
/// Inscrit [process] comme connexion amont de [playlist] tant qu'il vit.
///
/// [lastActivity] : dernier signe de vie du spectateur ; la connexion n'est
/// coupée pour en ouvrir une autre que s'il remonte à plus de
/// [_viewerGrace].
void _registerUpstream(
PlaylistConfig playlist,
String id,
Process process, {
required void Function() release,
required DateTime Function() lastActivity,
}) {
final token = _upstreamSlots.register(
id: id,
account: accountKeyOf(playlist),
release: release,
isActive: () =>
DateTime.now().difference(lastActivity()) < _viewerGrace,
);
process.exitCode.then((_) => _upstreamSlots.unregister(id, token: token));
}
/// Helper to resolve FFmpeg path (SYSTEM PATH vs Portable)
String _getFFmpegPath() {
if (Platform.isWindows) {
@@ -57,6 +113,16 @@ String _sanitizeQuality(String? quality) {
bool _isValidStreamId(String id) => RegExp(r'^[a-zA-Z0-9_-]+$').hasMatch(id);
/// Extension de conteneur d'un film/épisode, validée (elle finit dans l'URL
/// du panneau). `mkv` par défaut : c'était l'unique valeur jusqu'ici, les
/// anciens clients qui n'envoient pas `ext` gardent leur comportement.
String vodContainerExtension(String? ext) {
final value = ext?.toLowerCase();
return value != null && RegExp(r'^[a-z0-9]{2,5}$').hasMatch(value)
? value
: 'mkv';
}
/// Video encoding args for the requested quality (live profile).
List<String> _liveVideoArgs(String quality, bool gpu) {
switch (quality) {
@@ -192,12 +258,19 @@ typedef _SessionAttempt = ({FfmpegSession? session, bool ready, String? error});
/// Un `Process.start` qui échoue (serveur à court de processus ou de
/// mémoire) devient un échec ordinaire au lieu d'une exception non
/// rattrapée qui remontait en 500.
///
/// [upstream] : playlist dont la session ouvre une connexion chez le
/// fournisseur (direct, film). Null pour un enregistrement, lu sur disque.
Future<_SessionAttempt> _runSession({
required String id,
required bool isLive,
required List<String> Function(Directory dir) argsBuilder,
int minSegments = 1,
PlaylistConfig? upstream,
}) async {
final starting = sessionManager.needsStart(id, isLive: isLive);
if (upstream != null && starting) await _claimUpstream(upstream, id);
final FfmpegSession session;
try {
session = await sessionManager.getOrStart(
@@ -209,6 +282,16 @@ Future<_SessionAttempt> _runSession({
} on ProcessException catch (e) {
return (session: null, ready: false, error: 'FFmpeg start failed: $e');
}
if (upstream != null && starting) {
_registerUpstream(
upstream,
id,
session.process,
release: () => sessionManager.killSession(id),
// Touchée à chaque requête de playlist ou de segment.
lastActivity: () => session.lastAccess,
);
}
final outcome =
await sessionManager.waitForPlaylist(session, minSegments: minSegments);
return (session: session, ready: outcome.ready, error: outcome.error);
@@ -226,6 +309,7 @@ Future<_SessionAttempt> _runWithGpuFallback({
required bool wantGpu,
required List<String> Function(bool gpu) buildArgs,
int minSegments = 1,
PlaylistConfig? upstream,
}) async {
final gpu = wantGpu && _gpuHealth.available;
var attempt = await _runSession(
@@ -233,6 +317,7 @@ Future<_SessionAttempt> _runWithGpuFallback({
isLive: isLive,
argsBuilder: (_) => buildArgs(gpu),
minSegments: minSegments,
upstream: upstream,
);
if (!attempt.ready && gpu && isGpuFailure(attempt.error)) {
@@ -245,6 +330,7 @@ Future<_SessionAttempt> _runWithGpuFallback({
isLive: isLive,
argsBuilder: (_) => buildArgs(false),
minSegments: minSegments,
upstream: upstream,
);
}
return attempt;
@@ -307,7 +393,11 @@ Stream<List<int>> _resilientLiveBody(
yield chunk;
}
} catch (e) {
print('[Live Proxy] $streamId : coupure amont ($e)');
// Un `ClientException` recopie l'URL amont, identifiants compris.
print(
'[Live Proxy] $streamId : coupure amont '
'(${LogRedactor.redactUrl('$e')})',
);
} finally {
upstream.client.close();
}
@@ -382,6 +472,7 @@ Handler createLiveStreamHandler(
final result = await _runWithGpuFallback(
id: sessionId,
isLive: true,
upstream: playlist,
wantGpu: useNvidiaGpu && quality != 'source',
buildArgs: (gpu) => [
'-hide_banner', '-loglevel', 'warning',
@@ -469,6 +560,11 @@ Handler createLiveStreamHandler(
'[Live Turbo] $streamId: ${LogRedactor.redactUrl(targetUrl)}',
);
// Un identifiant par requête : deux onglets sur la même chaîne sont
// deux connexions distinctes chez le fournisseur.
final upstreamId = 'turbo_${streamId}_${++_turboCounter}';
await _claimUpstream(playlist, upstreamId);
final Process process;
try {
process = await Process.start(_getFFmpegPath(), [
@@ -494,10 +590,24 @@ Handler createLiveStreamHandler(
} on ProcessException catch (e) {
// Serveur saturé (plus de processus ou de mémoire) : un 503 que le
// lecteur sait traiter, plutôt qu'une exception qui remonte en 500.
print('[Live Turbo] $streamId : FFmpeg n\'a pas démarré ($e)');
// `ProcessException` liste les arguments, donc l'URL `-i` du panneau.
print('[Live Turbo] $streamId : FFmpeg n\'a pas démarré '
'(${LogRedactor.redactUrl('$e')})');
return Response(503, body: 'FFmpeg start failed');
}
// Dernier octet remis au client : un lecteur en vie en reçoit à chaque
// rafale ; un client parti (zap mal détecté derrière le proxy) n'en
// reçoit plus, la contre-pression bloquant FFmpeg.
var lastDelivered = DateTime.now();
_registerUpstream(
playlist,
upstreamId,
process,
release: () => process.kill(ProcessSignal.sigterm),
lastActivity: () => lastDelivered,
);
// Journaliser les erreurs FFmpeg (redactées) sans bloquer le flux.
process.stderr.transform(const SystemEncoding().decoder).listen((line) {
final trimmed = line.trim();
@@ -513,7 +623,10 @@ Handler createLiveStreamHandler(
final body = pipeWithBackpressure(
process.stdout,
onStop: () => process.kill(ProcessSignal.sigterm),
);
).map((chunk) {
lastDelivered = DateTime.now();
return chunk;
});
return Response(
200,
@@ -601,8 +714,11 @@ Handler createVodStreamHandler(
final contentType = request.url.queryParameters['type'] ?? 'movie';
final basePath = contentType == 'series' ? 'series' : 'movie';
// Extension réelle du fichier (`container_extension` du catalogue) : le
// panneau refuse un `.mkv` imposé à un film stocké en `.mp4`.
final targetUrl =
'${playlist.dns}/$basePath/${playlist.username}/${playlist.password}/$streamId.mkv';
'${playlist.dns}/$basePath/${playlist.username}/${playlist.password}/'
'$streamId.${vodContainerExtension(request.url.queryParameters['ext'])}';
final useNvidiaGpu = isGpuEnabled?.call() ?? _isNvidiaGpuEnabled();
final sessionId = 'vod_${streamId}_$quality';
@@ -615,6 +731,7 @@ Handler createVodStreamHandler(
final result = await _runWithGpuFallback(
id: sessionId,
isLive: false,
upstream: playlist,
wantGpu: useNvidiaGpu && quality != 'source',
buildArgs: (gpu) => [
'-hide_banner', '-loglevel', 'warning',
@@ -671,10 +788,15 @@ Handler createVodStreamHandler(
(Request request, String streamId) async {
final quality =
_sanitizeQuality(request.url.queryParameters['quality']);
final query = request.url.queryParameters['type'] != null
? '?type=${request.url.queryParameters['type']}'
: '';
return Response.found('/api/vod/$streamId/$quality/playlist.m3u8$query');
final params = request.url.queryParameters;
final query = Uri(queryParameters: {
if (params['type'] != null) 'type': params['type']!,
if (params['ext'] != null) 'ext': params['ext']!,
}).query;
return Response.found(
'/api/vod/$streamId/$quality/playlist.m3u8'
'${query.isEmpty ? '' : '?$query'}',
);
});
// Route: /api/vod/{streamId}/{quality}/{segment}
+133
View File
@@ -0,0 +1,133 @@
import 'dart:async';
import 'dart:convert';
import 'package:http/http.dart' as http;
import '../models/playlist_config.dart';
/// Diagnostic admin : chronomètre l'ouverture d'une chaîne chez le
/// fournisseur, en `.ts` et en `.m3u8`.
///
/// POURQUOI : ~4,3 s des ~4,6 s d'un zap se passent avant le premier octet
/// du panneau (mesuré en prod). Le panneau autorise aussi le HLS ; reste à
/// savoir lequel des deux formats démarre le plus vite, et où part le temps
/// (redirection, en-têtes, premier octet). Les identifiants n'apparaissent
/// jamais dans la réponse : seuls l'hôte et les durées sont rendus.
Future<Map<String, Object?>> probeUpstream(
PlaylistConfig playlist,
String streamId, {
http.Client? client,
}) async {
final c = client ?? http.Client();
try {
final base = '${playlist.dns}/live/${playlist.username}/'
'${playlist.password}/$streamId';
final ts = await _probeStream(c, Uri.parse('$base.ts'));
// Laisser le panneau libérer la connexion (comptes à une connexion).
await Future<void>.delayed(const Duration(seconds: 2));
final hls = await _probeHls(c, Uri.parse('$base.m3u8'));
return {'stream': streamId, 'ts': ts, 'm3u8': hls};
} finally {
if (client == null) c.close();
}
}
/// Suit les redirections à la main pour chronométrer chaque saut.
Future<({http.StreamedResponse? response, List<Map<String, Object?>> hops, Uri url})>
_open(http.Client c, Uri url, Stopwatch sw) async {
final hops = <Map<String, Object?>>[];
var current = url;
for (var i = 0; i < 5; i++) {
final request = http.Request('GET', current)
..followRedirects = false
..headers['User-Agent'] = 'VLC/3.0.18 LibVLC/3.0.18';
final response =
await c.send(request).timeout(const Duration(seconds: 20));
hops.add({
'host': current.host,
'status': response.statusCode,
'headersMs': sw.elapsedMilliseconds,
});
final location = response.headers['location'];
if (const {301, 302, 303, 307, 308}.contains(response.statusCode) &&
location != null) {
await response.stream.drain<void>().catchError((_) {});
current = current.resolve(location);
continue;
}
return (response: response, hops: hops, url: current);
}
return (response: null, hops: hops, url: current);
}
Future<Map<String, Object?>> _probeStream(http.Client c, Uri url) async {
final sw = Stopwatch()..start();
try {
final opened = await _open(c, url, sw);
final response = opened.response;
if (response == null || response.statusCode != 200) {
return {'hops': opened.hops, 'error': 'status ${response?.statusCode}'};
}
var bytes = 0;
int? firstByteMs;
final sub = response.stream.listen(null);
final done = Completer<void>();
sub.onData((chunk) {
firstByteMs ??= sw.elapsedMilliseconds;
bytes += chunk.length;
// 1 Mo suffit à voir le premier envoi (préchargement du panneau).
if (bytes > 1000000 && !done.isCompleted) done.complete();
});
sub.onDone(() {
if (!done.isCompleted) done.complete();
});
await done.future.timeout(const Duration(seconds: 15), onTimeout: () {});
await sub.cancel();
return {
'hops': opened.hops,
'firstByteMs': firstByteMs,
'firstMegabyteMs': bytes > 1000000 ? sw.elapsedMilliseconds : null,
};
} catch (e) {
return {'error': e.runtimeType.toString(), 'afterMs': sw.elapsedMilliseconds};
}
}
Future<Map<String, Object?>> _probeHls(http.Client c, Uri url) async {
final sw = Stopwatch()..start();
try {
final opened = await _open(c, url, sw);
final response = opened.response;
if (response == null || response.statusCode != 200) {
return {'hops': opened.hops, 'error': 'status ${response?.statusCode}'};
}
final body = await response.stream.bytesToString();
final playlistMs = sw.elapsedMilliseconds;
final lines = const LineSplitter().convert(body);
final segments =
lines.where((l) => l.isNotEmpty && !l.startsWith('#')).toList();
final target = lines
.firstWhere((l) => l.startsWith('#EXT-X-TARGETDURATION'),
orElse: () => '')
.split(':')
.last;
if (segments.isEmpty) {
return {'hops': opened.hops, 'playlistMs': playlistMs, 'segments': 0};
}
// Un lecteur démarre près du direct : chronométrer l'avant-dernier.
final pick = segments.length >= 2
? segments[segments.length - 2]
: segments.last;
final segment = await _probeStream(c, opened.url.resolve(pick));
return {
'hops': opened.hops,
'playlistMs': playlistMs,
'segments': segments.length,
'targetDuration': target,
'segment': segment,
'segmentStartedAtMs': playlistMs,
};
} catch (e) {
return {'error': e.runtimeType.toString(), 'afterMs': sw.elapsedMilliseconds};
}
}
+52 -2
View File
@@ -27,6 +27,8 @@ import 'middleware/auth_middleware.dart';
import 'middleware/security_middleware.dart';
import 'services/cleanup_service.dart';
import 'services/recording_scheduler.dart';
import 'utils/asset_versioning.dart';
import 'api/upstream_probe.dart';
void main(List<String> args) async {
// Parse command line arguments
@@ -130,7 +132,9 @@ void main(List<String> args) async {
// Source XMLTV de repli. Vider EPG_XMLTV_URLS désactive tout appel sortant :
// l'EPG se limite alors au panneau de l'abonné.
final xmltvUrls = (Platform.environment['EPG_XMLTV_URLS'] ??
'https://epgshare01.online/epgshare01/epg_ripper_FR1.xml.gz')
'https://epgshare01.online/epgshare01/epg_ripper_FR1.xml.gz,'
'https://epgshare01.online/epgshare01/epg_ripper_CH1.xml.gz,'
'https://epgshare01.online/epgshare01/epg_ripper_BE2.xml.gz')
.split(',')
.map((u) => u.trim())
.where((u) => u.isNotEmpty)
@@ -268,6 +272,29 @@ void main(List<String> args) async {
);
});
// GET /api/admin/upstream-probe?stream=<id> : chronomètre l'ouverture
// d'une chaîne chez le fournisseur en .ts et en .m3u8 (diagnostic du
// délai de zap ; aucun identifiant dans la réponse).
router.get('/upstream-probe', (Request req) async {
final user = req.context['user'] as User?;
if (user == null || !user.isAdmin) {
return Response.forbidden(
jsonEncode({'error': 'Admin access required'}),
);
}
final streamId = req.url.queryParameters['stream'] ?? '';
if (!RegExp(r'^[0-9]+$').hasMatch(streamId)) {
return Response.badRequest(body: 'stream=<id numérique> requis');
}
final playlist = await getPlaylist(req);
if (playlist == null) return Response.forbidden('No playlist');
final result = await probeUpstream(playlist, streamId);
return Response.ok(
jsonEncode(result),
headers: {'content-type': 'application/json'},
);
});
return router(request);
}),
);
@@ -293,14 +320,37 @@ void main(List<String> args) async {
listDirectories: false,
);
// Points d'entrée réécrits avec `?v=<empreinte>` sur leurs JS : voir
// asset_versioning.dart (un proxy qui cache les .js servait l'ancienne
// app et l'ancien lecteur après chaque déploiement).
final versionedEntryPoints = buildVersionedEntryPoints(webPath);
print('[Server] Assets versionnés : ${versionedEntryPoints.keys.join(', ')}');
// Wrap static handler to enforce cache policies
FutureOr<Response> staticHandler(Request request) async {
final path = request.url.path;
final versioned =
versionedEntryPoints[path.isEmpty ? 'index.html' : path];
if (versioned != null && request.method == 'GET') {
return Response.ok(
versioned,
headers: {
'Content-Type': path.endsWith('.js')
? 'text/javascript; charset=utf-8'
: 'text/html; charset=utf-8',
'Cache-Control': 'no-store, no-cache, must-revalidate, max-age=0',
'Pragma': 'no-cache',
'Expires': '0',
},
);
}
final response = await baseStaticHandler(request);
// Disable cache for entry points to ensure updates are seen immediately.
// Vendored player libraries (hls.js/mpegts.js, ~750 KB) are pinned
// versions: keep them cacheable or every player open re-downloads them.
final path = request.url.path;
final isVendored = path.startsWith('vendor/');
if (!isVendored &&
(path.isEmpty ||
+17 -8
View File
@@ -2,6 +2,8 @@ import 'dart:async';
import 'dart:convert';
import 'dart:io';
import '../utils/log_redactor.dart';
/// One running FFmpeg transcoding session (live, VOD or recording playback).
class FfmpegSession {
final String id;
@@ -72,6 +74,17 @@ class FfmpegSessionManager {
bool contains(String id) => _sessions.containsKey(id);
/// Vrai si [getOrStart] lancerait un nouveau processus FFmpeg pour [id] —
/// donc une nouvelle connexion au fournisseur.
bool needsStart(String id, {required bool isLive}) {
final existing = _sessions[id];
if (existing == null) return true;
if (!existing.exited) return false;
// A VOD/recording transcode that finished cleanly is still fully
// playable from its segments — reuse it instead of re-transcoding.
return isLive || existing.exitCode != 0 || !_playlistComplete(existing);
}
/// Returns the existing healthy session or starts a new FFmpeg process.
///
/// [argsBuilder] receives the session working directory and returns the
@@ -84,17 +97,11 @@ class FfmpegSessionManager {
required List<String> Function(Directory dir) argsBuilder,
}) async {
final existing = _sessions[id];
if (existing != null && !existing.exited) {
if (existing != null && !needsStart(id, isLive: isLive)) {
existing.touch();
return existing;
}
if (existing != null) {
// A VOD/recording transcode that finished cleanly is still fully
// playable from its segments — reuse it instead of re-transcoding.
if (!isLive && existing.exitCode == 0 && _playlistComplete(existing)) {
existing.touch();
return existing;
}
// Process died: clean up before restarting
killSession(id);
}
@@ -121,7 +128,9 @@ class FfmpegSessionManager {
process.stderr.transform(utf8.decoder).listen((data) {
session.recentStderr.add(data);
if (session.recentStderr.length > 20) session.recentStderr.removeAt(0);
print('[FFmpeg $id] $data');
// FFmpeg rappelle l'URL d'entrée (« Input #0 … from 'http://…' ») :
// masquer les identifiants Xtream avant de journaliser.
print('[FFmpeg $id] ${LogRedactor.redactUrl(data)}');
});
process.exitCode.then((code) {
+10 -2
View File
@@ -533,7 +533,12 @@ class RecordingScheduler {
await _launchFfmpeg(active);
} catch (e, st) {
// Attraper TOUTES les exceptions pour éviter de crasher le serveur
print('[RecordingScheduler] ERREUR dans _startRecording: $e\n$st');
// `ProcessException` liste les arguments FFmpeg, URL de la source
// comprise : masquer les identifiants avant de journaliser.
print(
'[RecordingScheduler] ERREUR dans _startRecording: '
'${LogRedactor.redactUrl('$e')}\n$st',
);
_db.updateRecordingStatus(
recording.id,
'failed',
@@ -662,7 +667,10 @@ class RecordingScheduler {
try {
await _launchFfmpeg(active);
} catch (e) {
print('[RecordingScheduler] Relance impossible: $e');
print(
'[RecordingScheduler] Relance impossible: '
'${LogRedactor.redactUrl('$e')}',
);
_active.remove(recording.id);
await _finalize(
active,
+183
View File
@@ -0,0 +1,183 @@
import 'dart:async';
import 'dart:convert';
import 'package:http/http.dart' as http;
import '../models/playlist_config.dart';
/// Connexions ouvertes vers le fournisseur Xtream, par compte.
///
/// POURQUOI : la plupart des abonnements n'autorisent qu'UNE connexion
/// simultanée (`max_connections: "1"`, vérifié en prod). Or un FFmpeg de
/// lecture survit au spectateur : 4 min pour un direct HLS, jusqu'à 15 min
/// ou la fin du téléchargement pour un film, le temps de détecter la
/// déconnexion pour `turbo.ts`. Le flux suivant tombait sur un slot occupé :
/// le panneau répondait `HTTP 551` ou ne répondait pas (timeout), d'où les
/// films « indisponibles » et les zaps qui échouent.
///
/// Règle : avant d'ouvrir une nouvelle connexion amont, on coupe les
/// connexions ORPHELINES du même compte (plus personne ne les consomme),
/// les plus anciennes d'abord, jusqu'à lui faire de la place.
///
/// Un flux encore regardé n'est jamais coupé. Une première version coupait
/// sans distinction (« le dernier gagne ») : deux lecteurs ouverts sur un
/// compte à une connexion se coupaient alors en boucle, chacun relançant
/// aussitôt — plus aucune image nulle part (constaté en prod). Le second
/// lecteur est désormais refusé par le panneau, comme avant, et le premier
/// continue.
///
/// Les enregistrements n'y sont pas inscrits : ils ne sont jamais coupés.
class UpstreamSlots {
final DateTime Function() _now;
final Map<String, _Slot> _slots = {};
int _nextToken = 0;
UpstreamSlots({DateTime Function()? now}) : _now = now ?? DateTime.now;
/// Inscrit la connexion [id] du compte [account]. [release] la coupe.
///
/// Renvoie un jeton à repasser à [unregister] : une session relancée sous
/// le même identifiant ne doit pas être désinscrite par la fin de
/// l'ancien processus.
///
/// [isActive] : vrai tant qu'un spectateur consomme le flux. Une
/// connexion active n'est jamais coupée par [makeRoom].
int register({
required String id,
required String account,
required void Function() release,
bool Function()? isActive,
}) {
final token = ++_nextToken;
_slots[id] = _Slot(token, account, _now(), release, isActive);
return token;
}
void unregister(String id, {int? token}) {
final slot = _slots[id];
if (slot == null) return;
if (token != null && slot.token != token) return;
_slots.remove(id);
}
int countFor(String account) =>
_slots.values.where((s) => s.account == account).length;
/// Libère assez de connexions de [account] pour qu'une nouvelle tienne
/// sous [max]. [keep] (la session demandée) n'est jamais coupée.
///
/// [max] <= 0 = quota inconnu : on ne coupe rien plutôt que de risquer
/// d'interrompre un autre spectateur.
List<String> makeRoom(String account, {required int max, required String keep}) {
if (max <= 0) return const [];
final others = _slots.entries
.where((e) => e.value.account == account && e.key != keep)
.toList();
final excess = others.length - (max - 1);
if (excess <= 0) return const [];
final orphans = others.where((e) => !e.value.active).toList()
..sort((a, b) => a.value.since.compareTo(b.value.since));
final freed = <String>[];
for (final entry in orphans.take(excess)) {
_slots.remove(entry.key);
try {
entry.value.release();
} catch (_) {}
freed.add(entry.key);
}
return freed;
}
}
class _Slot {
final int token;
final String account;
final DateTime since;
final void Function() release;
final bool Function()? isActive;
_Slot(this.token, this.account, this.since, this.release, this.isActive);
bool get active {
try {
return isActive?.call() ?? false;
} catch (_) {
return false;
}
}
}
/// Clé de compte : deux playlists sur les mêmes identifiants partagent le
/// même quota chez le fournisseur.
String accountKeyOf(PlaylistConfig p) => '${p.dns}|${p.username}';
/// `user_info.max_connections` d'une réponse `player_api.php`, ou null.
int? parseMaxConnections(Object? json) {
if (json is! Map) return null;
final info = json['user_info'];
if (info is! Map) return null;
final raw = info['max_connections'];
final value = raw is int ? raw : int.tryParse('$raw');
return (value != null && value > 0) ? value : null;
}
/// Quota de connexions par compte, lu chez le fournisseur et mis en cache.
///
/// Ne fait JAMAIS attendre un flux : `player_api.php` met jusqu'à 4 s à
/// répondre (mesuré en prod), délai qui s'ajoutait au démarrage du premier
/// flux puis de chaque flux suivant l'expiration du cache. La valeur connue
/// (même périmée), sinon [fallback], est rendue immédiatement ; la lecture
/// chez le panneau se fait en arrière-plan. Sans risque : le quota ne sert
/// qu'à décider quels flux ORPHELINS couper.
class AccountLimits {
final Duration ttl;
final Future<int?> Function(PlaylistConfig) _fetch;
final Map<String, ({int max, DateTime at})> _cache = {};
final Set<String> _refreshing = {};
/// Valeur retenue tant que le panneau n'a pas répondu : un seul flux, le
/// cas de loin le plus courant chez les fournisseurs Xtream.
static const fallback = 1;
AccountLimits({
this.ttl = const Duration(minutes: 10),
Future<int?> Function(PlaylistConfig)? fetch,
}) : _fetch = fetch ?? _fetchFromPanel;
Future<int> maxFor(PlaylistConfig p) async {
final key = accountKeyOf(p);
final cached = _cache[key];
if (cached == null || DateTime.now().difference(cached.at) >= ttl) {
_refresh(key, p);
}
return cached?.max ?? fallback;
}
void _refresh(String key, PlaylistConfig p) {
if (!_refreshing.add(key)) return; // une seule lecture à la fois
Future<int?>.sync(() => _fetch(p))
.then<int?>((v) => v, onError: (Object _) => null)
.then((value) {
// Un échec n'est pas mis en cache : on réessaiera au prochain flux.
if (value != null) _cache[key] = (max: value, at: DateTime.now());
}).whenComplete(() => _refreshing.remove(key));
}
static Future<int?> _fetchFromPanel(PlaylistConfig p) async {
final uri = Uri.parse('${p.dns}/player_api.php').replace(
queryParameters: {'username': p.username, 'password': p.password},
);
final client = http.Client();
try {
final response = await client.get(uri, headers: {
'User-Agent': 'VLC/3.0.18 LibVLC/3.0.18',
}).timeout(const Duration(seconds: 5));
if (response.statusCode != 200) return null;
return parseMaxConnections(jsonDecode(response.body));
} finally {
client.close();
}
}
}
+212 -18
View File
@@ -1,9 +1,13 @@
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';
import '../utils/log_redactor.dart';
/// Guide TV construit à partir de dumps XMLTV publics.
///
/// Beaucoup de panneaux Xtream servent un EPG figé depuis plusieurs jours, ou
@@ -78,8 +82,39 @@ class XmltvEpgService {
String? displayName,
}) async {
await ensureFresh();
for (final candidate in [channelId, displayName]) {
final key = normalizeKey(candidate);
final direct = _lookup(channelId, displayName);
if (direct.isNotEmpty) return direct;
// Chaîne décalée sans guide propre (« TF1 +1 » absent des dumps) : le
// guide de la chaîne principale, une heure plus tard, est exact.
if (displayName != null && _plusOne.hasMatch(displayName)) {
final base = displayName.replaceAll(_plusOne, ' ');
final hits = _lookup(null, base);
return [
for (final p in hits)
XmltvProgramme(
title: p.title,
description: p.description,
start: p.start.add(const Duration(hours: 1)),
stop: p.stop.add(const Duration(hours: 1)),
),
];
}
return const [];
}
static final _plusOne = RegExp(r'\+\s*1(?![0-9])');
List<XmltvProgramme> _lookup(String? channelId, String? displayName) {
final keys = [
normalizeKey(channelId),
if (displayName != null) ...nameKeys(displayName),
// Clés souples en dernier recours, dans leur propre espace de noms.
_looseIndexKey(looseKey(channelId)),
if (displayName != null)
_looseIndexKey(looseKey(cleanChannelName(displayName).name)),
];
for (final key in keys) {
if (key.isEmpty) continue;
final hit = _index[key];
if (hit != null && hit.isNotEmpty) return hit;
@@ -87,6 +122,110 @@ class XmltvEpgService {
return const [];
}
static String _looseIndexKey(String loose) => loose.isEmpty ? '' : '~$loose';
/// Jetons sans valeur distinctive pour reconnaître une chaîne.
static const _looseNoise = {
'channel', 'tv', 'hd', 'fhd', 'uhd', 'sd', '4k', 'hevc',
};
/// Codes pays en fin d'identifiant XMLTV (`TF1.fr`, `RTS1.ch`).
static const _countryCodes = {
'fr', 'be', 'ch', 'lu', 'mc', 'ca', 'uk', 'us', 'es', 'it', 'de', 'pt',
'nl',
};
/// Clé « souple » : sans code pays final ni mots génériques
/// (« AB3.Channel.fr » et « AB3 » → `ab3`).
///
/// Servie dans un espace de noms à part, consultée en tout dernier : elle
/// rattrape les variantes d'écriture sans fausser les correspondances
/// exactes.
static String looseKey(String? raw) {
if (raw == null) return '';
final folded = StringBuffer();
for (final rune in raw.toLowerCase().runes) {
final char = String.fromCharCode(rune);
folded.write(_accents[char] ?? char);
}
final tokens = folded
.toString()
.split(RegExp(r'[^a-z0-9]+'))
.where((t) => t.isNotEmpty)
.toList();
if (tokens.length > 1 && _countryCodes.contains(tokens.last)) {
tokens.removeLast();
}
tokens.removeWhere(_looseNoise.contains);
return tokens.join();
}
/// Préfixe pays des noms de panneau : `FR - `, `FR: `, `|FR| `, `[FR] `…
static final _countryPrefix = RegExp(
r'^\s*[\|\[\(]?\s*([A-Za-z]{2,3})\s*[\|\]\)]?\s*[-:|]\s*|^\s*[\|\[\(]\s*([A-Za-z]{2,3})\s*[\|\]\)]\s*',
);
/// Marqueurs de qualité ou de source en fin de nom, sans rapport avec la
/// chaîne elle-même.
static final _qualitySuffix = RegExp(
r'(\s+|^)(fhd|uhd|hd|sd|hq|lq|4k|8k|hevc|h\.?265|h\.?264|1080p?|720p?|'
r'50fps|60fps|backup|raw|vip|multi|\(backup\)|\(multi\))\s*$',
caseSensitive: false,
);
/// Clés candidates pour retrouver une chaîne du panneau par son nom.
///
/// POURQUOI : en prod, 59 % des chaînes françaises n'ont pas
/// d'`epg_channel_id`. Leur nom brut (« FR - TF1 FHD ◉ » → `frtf1fhd`) ne
/// correspond à aucune entrée XMLTV (`TF1.fr`, « TF1 »). On retire le
/// préfixe pays, les marqueurs de qualité et les symboles, puis on tente
/// aussi la forme `nom.pays` qu'utilisent les dumps publics.
///
/// Le décalage horaire (« +1 ») est conservé : TF1 +1 n'a pas le guide
/// de TF1.
static List<String> nameKeys(String name) {
final keys = <String>[normalizeKey(name)];
final cleaned = cleanChannelName(name);
final country = cleaned.country;
final base = normalizeKey(cleaned.name);
if (base.isNotEmpty) {
keys.add(base);
keys.add('$base${country ?? 'fr'}');
}
return keys.where((k) => k.isNotEmpty).toSet().toList();
}
/// Abréviations du panneau pour les chaînes publiques (« F3 ALPES »), là
/// où les dumps écrivent `France.3.-.Alpes.fr`.
static final _franceAbbrev = RegExp(r'^F\s?([2-5])\b', caseSensitive: false);
/// Nom de chaîne débarrassé du préfixe pays, des symboles et des
/// marqueurs de qualité ; [country] = code pays du préfixe, s'il y en a.
static ({String name, String? country}) cleanChannelName(String raw) {
var cleaned = raw;
String? country;
final prefix = _countryPrefix.firstMatch(cleaned);
if (prefix != null) {
country = (prefix.group(1) ?? prefix.group(2))?.toLowerCase();
cleaned = cleaned.substring(prefix.end);
}
// Symboles décoratifs (◉, ᴴᴰ, ★…) : tout ce qui n'est ni lettre, ni
// chiffre, ni ponctuation utile.
cleaned = cleaned
.replaceAll(RegExp(r'[^\p{L}\p{N}\s+.&\-()]', unicode: true), ' ')
.trim();
String previous;
do {
previous = cleaned;
cleaned = cleaned.replaceFirst(_qualitySuffix, '').trim();
} while (cleaned != previous && cleaned.isNotEmpty);
cleaned = cleaned.replaceFirstMapped(
_franceAbbrev,
(m) => 'France ${m[1]}',
);
return (name: cleaned, country: country);
}
/// Recharge l'index s'il est absent ou périmé. Les appels concurrents
/// partagent le même téléchargement.
Future<void> ensureFresh() {
@@ -114,9 +253,13 @@ class XmltvEpgService {
var ok = 0;
for (final url in sourceUrls) {
// La source peut être le `xmltv.php` du panneau, dont l'URL porte les
// identifiants de l'abonné en clair : ne jamais la journaliser telle
// quelle (elle finirait dans `docker logs`). Le texte de l'exception
// est masqué aussi, un `ClientException` recopiant l'URL demandée.
final safeUrl = LogRedactor.redactUrl(url);
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) {
@@ -124,10 +267,12 @@ class XmltvEpgService {
}
ok++;
print(
'[XmltvEpg] $url : ${parsed.length} chaînes indexées',
'[XmltvEpg] $safeUrl : ${parsed.length} chaînes indexées',
);
} catch (e) {
print('[XmltvEpg] $url : échec ($e)');
print(
'[XmltvEpg] $safeUrl : échec (${LogRedactor.redactUrl('$e')})',
);
}
}
@@ -144,34 +289,70 @@ 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>>{};
// Alias : plusieurs dumps déclarent <channel id="X"> avec un
// <display-name> différent de l'identifiant. On indexe les deux pour
// maximiser les correspondances.
// Alias : plusieurs dumps déclarent <channel id="X"> avec un ou plusieurs
// <display-name> différents de l'identifiant (« TF1 HD », « TF1 »). On
// les indexe TOUS : seul le premier l'était, et c'était souvent la
// variante la moins proche du nom annoncé par le panneau.
final aliases = <String, String>{};
final displayNames = <String>[];
String? channelId;
String? programmeChannel;
@@ -188,6 +369,7 @@ class XmltvEpgService {
case 'channel':
channelId = _attr(event, 'id');
displayName.clear();
displayNames.clear();
case 'programme':
programmeChannel = _attr(event, 'channel');
start = parseXmltvDate(_attr(event, 'start'));
@@ -210,15 +392,27 @@ class XmltvEpgService {
case 'desc':
desc.write(text);
case 'display-name':
if (displayName.isEmpty) displayName.write(text);
displayName.write(text);
}
} else if (event is XmlEndElementEvent) {
switch (event.name) {
case 'display-name':
displayNames.add(displayName.toString());
displayName.clear();
case 'channel':
final id = normalizeKey(channelId);
final name = normalizeKey(displayName.toString());
if (id.isNotEmpty && name.isNotEmpty && id != name) {
aliases[name] = id;
for (final raw in displayNames) {
final name = normalizeKey(raw);
if (id.isNotEmpty && name.isNotEmpty && id != name) {
aliases.putIfAbsent(name, () => id);
}
}
// Clés souples (espace de noms « ~ ») : premier arrivé gagne.
if (id.isNotEmpty) {
for (final raw in [channelId, ...displayNames]) {
final loose = _looseIndexKey(looseKey(raw));
if (loose.isNotEmpty) aliases.putIfAbsent(loose, () => id);
}
}
channelId = null;
case 'programme':
+47
View File
@@ -0,0 +1,47 @@
import 'package:test/test.dart';
import '../utils/asset_versioning.dart';
void main() {
group('versionAssetRefs', () {
test('versionne la balise script du lecteur', () {
expect(
versionAssetRefs(
'<script src="xf-player-core.js"></script>',
{'xf-player-core.js': 'abc123'},
),
'<script src="xf-player-core.js?v=abc123"></script>',
);
});
test('versionne mainJsPath dans la config de build Flutter', () {
expect(
versionAssetRefs(
'{"compileTarget":"dart2js","mainJsPath":"main.dart.js"}',
{'main.dart.js': 'f00'},
),
'{"compileTarget":"dart2js","mainJsPath":"main.dart.js?v=f00"}',
);
});
test('ne touche ni un nom qui ne fait que contenir l\'asset ni une ref déjà versionnée', () {
const html = '<script src="vendor/xf-player-core.js"></script>'
'<script src="xf-player-core.js?v=old"></script>';
expect(versionAssetRefs(html, {'xf-player-core.js': 'new'}), html);
});
test('guillemets simples acceptés', () {
expect(
versionAssetRefs("<link href='xf-player.css'>", {'xf-player.css': '1'}),
"<link href='xf-player.css?v=1'>",
);
});
});
group('contentVersion', () {
test('stable pour un même contenu, différente sinon', () {
expect(contentVersion([1, 2, 3]), contentVersion([1, 2, 3]));
expect(contentVersion([1, 2, 3]), isNot(contentVersion([1, 2, 4])));
expect(contentVersion([1, 2, 3]), hasLength(12));
});
});
}
+108
View File
@@ -217,5 +217,113 @@ void main() {
expect(response.headers['X-Epg-Source'], 'xtream');
expect(actions, contains('get_simple_data_table'));
});
// Dump externe couvrant TF1, et panneau dont le xmltv.php est vide.
String externalDump() {
final now = DateTime.now().toUtc();
return xmltvDump(now, now.add(const Duration(hours: 1)))
.replaceAll('France2.fr', 'TF1.fr')
.replaceAll('FRANCE 2', 'TF1');
}
EpgApi apiWith(http.Client client) => EpgApi(
(_) async => playlist,
httpClient: client,
xmltv: XmltvEpgService(
sourceUrls: const ['http://ext.example/epg.xml'],
client: client,
),
panelXmltvBuilder: (config) => XmltvEpgService(
sourceUrls: ['${config.dns}/xmltv.php'],
client: client,
),
);
test('un panneau qui plante ne prive plus la chaîne du dump externe',
() async {
// Constaté en prod : chaque /api/epg répondait 500 en 4 s, le guide
// externe (qui couvrait pourtant la chaîne) n'étant jamais consulté.
final client = MockClient((request) async {
final url = request.url;
if (url.host == 'ext.example') return http.Response(externalDump(), 200);
if (url.path.endsWith('/xmltv.php')) return http.Response('<tv></tv>', 200);
final action = url.queryParameters['action'] ?? '';
if (action == 'get_live_streams') {
// Octets UTF-8 et JSON sans charset annoncé, comme le panneau réel.
return http.Response.bytes(
utf8.encode(jsonEncode([
{'stream_id': 1, 'name': 'FR - TF1 FHD ◉', 'epg_channel_id': ''},
])),
200,
headers: {'content-type': 'application/json'},
);
}
throw http.ClientException('Connection closed', url);
});
final response = await apiWith(client).handleGetEpg(requestFor('1'), '1');
final body = jsonDecode(await response.readAsString()) as Map;
expect(response.statusCode, 200);
expect(response.headers['X-Epg-Source'], 'xmltv');
expect((body['programmes'] as List).single['title'], 'Journal de 20h');
});
test('le dump externe passe avant l\'appel lent chaîne par chaîne',
() async {
final actions = <String>[];
final client = MockClient((request) async {
final url = request.url;
if (url.host == 'ext.example') return http.Response(externalDump(), 200);
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([
{'stream_id': 1, 'name': 'TF1', 'epg_channel_id': 'TF1.fr'},
]),
200,
);
}
return http.Response('{}', 200);
});
final response = await apiWith(client).handleGetEpg(requestFor('1'), '1');
expect(response.headers['X-Epg-Source'], 'xmltv');
expect(actions, isNot(contains('get_simple_data_table')));
});
test('la grille entière partage un seul téléchargement de la table',
() async {
// 30 tuiles demandent leur guide d'un coup : sans mutualisation, autant
// de `get_live_streams` (7 Mo en prod) partaient en parallèle.
var liveStreamsCalls = 0;
final client = MockClient((request) async {
final url = request.url;
if (url.host == 'ext.example') return http.Response(externalDump(), 200);
if (url.path.endsWith('/xmltv.php')) return http.Response('<tv></tv>', 200);
final action = url.queryParameters['action'] ?? '';
if (action == 'get_live_streams') {
liveStreamsCalls++;
await Future<void>.delayed(const Duration(milliseconds: 50));
return http.Response(
jsonEncode([
for (var i = 1; i <= 5; i++)
{'stream_id': i, 'name': 'TF1', 'epg_channel_id': 'TF1.fr'},
]),
200,
);
}
return http.Response('{}', 200);
});
final api = apiWith(client);
final responses = await Future.wait([
for (var i = 1; i <= 5; i++) api.handleGetEpg(requestFor('$i'), '$i'),
]);
expect(responses.every((r) => r.statusCode == 200), isTrue);
expect(liveStreamsCalls, 1);
});
});
}
+62
View File
@@ -0,0 +1,62 @@
import 'dart:convert';
import 'package:http/http.dart' as http;
import 'package:http/testing.dart';
import 'package:test/test.dart';
import '../api/upstream_probe.dart';
import '../models/playlist_config.dart';
void main() {
final playlist = PlaylistConfig(
id: 'p',
name: 'n',
dns: 'http://panel.test',
username: 'secretuser',
password: 'secretpass',
createdAt: DateTime(2026),
);
// Panneau simulé : redirige vers un serveur de diffusion avec jeton.
final client = MockClient.streaming((request, _) async {
final path = request.url.path;
if (request.url.host == 'panel.test') {
return http.StreamedResponse(
const Stream.empty(),
302,
headers: {
'location': 'http://lb.test$path?token=secrettoken',
},
);
}
if (path.endsWith('.m3u8')) {
return http.StreamedResponse(
Stream.value(utf8.encode(
'#EXTM3U\n#EXT-X-TARGETDURATION:6\n#EXTINF:6,\nseg1.ts\n'
'#EXTINF:6,\nseg2.ts\n',
)),
200,
);
}
return http.StreamedResponse(Stream.value(List.filled(2000, 0x47)), 200);
});
test('chronomètre .ts et .m3u8 en suivant les redirections', () async {
final result = await probeUpstream(playlist, '123', client: client);
final ts = result['ts'] as Map;
expect((ts['hops'] as List).map((h) => (h as Map)['status']), [302, 200]);
expect(ts['firstByteMs'], isNotNull);
final hls = result['m3u8'] as Map;
expect(hls['segments'], 2);
expect(hls['targetDuration'], '6');
expect((hls['segment'] as Map)['firstByteMs'], isNotNull);
});
test('aucun identifiant ni jeton dans la réponse', () async {
final json = jsonEncode(await probeUpstream(playlist, '123', client: client));
expect(json, isNot(contains('secretuser')));
expect(json, isNot(contains('secretpass')));
expect(json, isNot(contains('secrettoken')));
});
}
+164
View File
@@ -0,0 +1,164 @@
import 'dart:async';
import 'package:test/test.dart';
import '../models/playlist_config.dart';
import '../services/upstream_slots.dart';
void main() {
group('UpstreamSlots.makeRoom', () {
late UpstreamSlots slots;
late List<String> released;
var clock = DateTime(2026, 1, 1);
void add(String id, String account) {
clock = clock.add(const Duration(seconds: 1));
slots.register(id: id, account: account, release: () => released.add(id));
}
setUp(() {
slots = UpstreamSlots(now: () => clock);
released = [];
});
test('compte à 1 connexion : le flux précédent est coupé', () {
add('live_1_source', 'acc');
final freed = slots.makeRoom('acc', max: 1, keep: 'live_2_source');
expect(freed, ['live_1_source']);
expect(released, ['live_1_source']);
expect(slots.countFor('acc'), 0);
});
test('la session demandée elle-même n\'est jamais coupée', () {
add('live_1_source', 'acc');
expect(slots.makeRoom('acc', max: 1, keep: 'live_1_source'), isEmpty);
expect(released, isEmpty);
});
test('les plus anciennes partent d\'abord, dans la limite du quota', () {
add('a', 'acc');
add('b', 'acc');
add('c', 'acc');
// 3 connexions max : il faut une place pour la nouvelle → 1 seule coupée.
expect(slots.makeRoom('acc', max: 3, keep: 'new'), ['a']);
expect(slots.countFor('acc'), 2);
});
test('un autre compte n\'est jamais touché', () {
add('x', 'other');
expect(slots.makeRoom('acc', max: 1, keep: 'new'), isEmpty);
expect(slots.countFor('other'), 1);
});
test('unregister avec un jeton périmé ne libère pas la session relancée', () {
final old = slots.register(id: 's', account: 'acc', release: () {});
slots.register(id: 's', account: 'acc', release: () {});
slots.unregister('s', token: old);
expect(slots.countFor('acc'), 1);
});
test('un flux encore regardé n\'est jamais coupé (pas de ping-pong entre deux lecteurs)', () {
slots.register(
id: 'watched',
account: 'acc',
release: () => released.add('watched'),
isActive: () => true,
);
expect(slots.makeRoom('acc', max: 1, keep: 'new'), isEmpty);
expect(released, isEmpty);
expect(slots.countFor('acc'), 1);
});
test('seuls les orphelins partent, même plus récents qu\'un flux regardé', () {
slots.register(
id: 'watched',
account: 'acc',
release: () => released.add('watched'),
isActive: () => true,
);
add('orphan', 'acc');
expect(slots.makeRoom('acc', max: 1, keep: 'new'), ['orphan']);
expect(released, ['orphan']);
});
test('quota inconnu ou nul : rien n\'est coupé', () {
add('a', 'acc');
expect(slots.makeRoom('acc', max: 0, keep: 'new'), isEmpty);
});
});
group('AccountLimits', () {
PlaylistConfig playlist() => PlaylistConfig(
id: 'p',
name: 'n',
dns: 'http://panel.test',
username: 'u',
password: 'x',
createdAt: DateTime(2026),
);
test('ne fait jamais attendre un flux : repli immédiat, valeur réelle ensuite', () async {
final pending = Completer<int?>();
final limits = AccountLimits(fetch: (_) => pending.future);
// Le panneau répond lentement (4 s mesurées en prod) : le premier flux
// part tout de suite sur la valeur de repli.
expect(await limits.maxFor(playlist()), AccountLimits.fallback);
pending.complete(3);
await Future<void>.delayed(Duration.zero);
expect(await limits.maxFor(playlist()), 3);
});
test('valeur périmée servie pendant le rafraîchissement', () async {
var calls = 0;
final limits = AccountLimits(
ttl: Duration.zero,
fetch: (_) async => ++calls == 1 ? 2 : 5,
);
await limits.maxFor(playlist());
await Future<void>.delayed(Duration.zero);
// Périmée (ttl nul) : la valeur connue sort sans attendre la nouvelle.
expect(await limits.maxFor(playlist()), 2);
await Future<void>.delayed(Duration.zero);
expect(await limits.maxFor(playlist()), 5);
});
test('panneau en échec : repli, sans exception', () async {
final limits = AccountLimits(fetch: (_) async => throw Exception('down'));
expect(await limits.maxFor(playlist()), AccountLimits.fallback);
await Future<void>.delayed(Duration.zero);
expect(await limits.maxFor(playlist()), AccountLimits.fallback);
});
});
group('parseMaxConnections', () {
test('lit la valeur texte renvoyée par player_api', () {
expect(
parseMaxConnections({
'user_info': {'max_connections': '1'},
}),
1,
);
});
test('accepte un entier', () {
expect(
parseMaxConnections({
'user_info': {'max_connections': 3},
}),
3,
);
});
test('réponse inexploitable → null', () {
expect(parseMaxConnections({'user_info': {}}), isNull);
expect(parseMaxConnections('oops'), isNull);
expect(
parseMaxConnections({
'user_info': {'max_connections': '0'},
}),
isNull,
);
});
});
}
+23
View File
@@ -0,0 +1,23 @@
import 'package:test/test.dart';
import '../api/streaming_handler.dart';
void main() {
group('vodContainerExtension', () {
test('garde l\'extension réelle du catalogue', () {
expect(vodContainerExtension('mp4'), 'mp4');
expect(vodContainerExtension('MKV'), 'mkv');
expect(vodContainerExtension('avi'), 'avi');
});
test('ancien client sans ext : mkv, comme avant', () {
expect(vodContainerExtension(null), 'mkv');
expect(vodContainerExtension(''), 'mkv');
});
test('refuse tout ce qui pourrait sortir du nom de fichier', () {
expect(vodContainerExtension('mp4/../x'), 'mkv');
expect(vodContainerExtension('mp4?a=b'), 'mkv');
expect(vodContainerExtension('toolongext'), 'mkv');
});
});
}
@@ -0,0 +1,75 @@
import 'dart:async';
import 'package:http/http.dart' as http;
import 'package:http/testing.dart';
import 'package:test/test.dart';
import '../services/xmltv_epg_service.dart';
/// URL du dump `xmltv.php` d'un panneau Xtream : les identifiants de
/// l'abonné y figurent en clair, comme en production.
const _panelUrl =
'http://panel.example:8080/xmltv.php?username=john&password=hunter2';
/// Exécute [body] en capturant tout ce qui passe par `print`.
Future<List<String>> _capturePrints(Future<void> Function() body) async {
final lines = <String>[];
await runZoned(
body,
zoneSpecification: ZoneSpecification(
print: (self, parent, zone, line) => lines.add(line),
),
);
return lines;
}
void main() {
group('XmltvEpgService — journalisation', () {
test('masque les identifiants de l\'URL quand la source répond', () async {
final service = XmltvEpgService(
sourceUrls: const [_panelUrl],
client: MockClient(
(_) async => http.Response('<?xml version="1.0"?><tv></tv>', 200),
),
);
final lines = await _capturePrints(service.ensureFresh);
final xmltvLines = lines.where((l) => l.startsWith('[XmltvEpg]'));
expect(xmltvLines, isNotEmpty);
for (final line in xmltvLines) {
expect(line, isNot(contains('john')));
expect(line, isNot(contains('hunter2')));
}
expect(
lines,
contains(contains('username=***&password=***')),
);
});
test('masque les identifiants de l\'URL et de l\'exception en cas d\'échec',
() async {
// `ClientException` recopie l'URL demandée dans son message : c'est
// par là que les identifiants fuyaient aussi.
final service = XmltvEpgService(
sourceUrls: const [_panelUrl],
client: MockClient(
(request) async => throw http.ClientException(
'Connection refused',
request.url,
),
),
);
final lines = await _capturePrints(service.ensureFresh);
final failure = lines.firstWhere(
(l) => l.contains('échec'),
orElse: () => fail('aucune ligne d\'échec journalisée : $lines'),
);
expect(failure, isNot(contains('john')));
expect(failure, isNot(contains('hunter2')));
expect(failure, contains('username=***&password=***'));
expect(failure, contains('Connection refused'));
});
});
}
+185
View File
@@ -1,3 +1,6 @@
import 'dart:async';
import 'dart:convert';
import 'package:http/http.dart' as http;
import 'package:http/testing.dart';
import 'package:test/test.dart';
@@ -26,6 +29,13 @@ String _fixture(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';
}
void main() {
final service = XmltvEpgService(sourceUrls: const []);
@@ -49,6 +59,117 @@ void main() {
});
});
group('nameKeys', () {
// Noms réels du panneau (prod) : 59 % des chaînes françaises n'ont pas
// d'epg_channel_id, seul leur nom permet de retrouver le guide.
test('retire le préfixe pays et les marqueurs de qualité', () {
expect(XmltvEpgService.nameKeys('FR - TF1 FHD ◉'), contains('tf1'));
expect(XmltvEpgService.nameKeys('FR - FRANCE 4 FHD'), contains('france4'));
expect(XmltvEpgService.nameKeys('|FR| M6 HD'), contains('m6'));
expect(XmltvEpgService.nameKeys('[FR] Arte UHD 4K'), contains('arte'));
expect(XmltvEpgService.nameKeys('FR: RMC Story HEVC'), contains('rmcstory'));
});
test('propose aussi la forme identifiant « nom.pays »', () {
// Les dumps publics identifient souvent la chaîne par `TF1.fr`.
expect(XmltvEpgService.nameKeys('FR - TF1 FHD'), contains('tf1fr'));
});
test('garde la variante décalée distincte de la chaîne principale', () {
expect(XmltvEpgService.nameKeys('FR - TF1 +1 FHD'), isNot(contains('tf1')));
});
test('ne réduit pas un nom court à rien', () {
expect(XmltvEpgService.nameKeys('M6'), contains('m6'));
expect(XmltvEpgService.nameKeys('HD'), isNot(contains('')));
});
});
group('correspondances élargies (chaînes vides mesurées en prod)', () {
XmltvEpgService serviceWith(String dump) => XmltvEpgService(
sourceUrls: const ['http://ext/epg.xml'],
client: MockClient((_) async => http.Response.bytes(
// Octets UTF-8 : les ids réels portent des accents.
const Utf8Encoder().convert(dump),
200,
)),
);
String channel(String id, String name, String title, DateTime start) => '''
<channel id="$id"><display-name>$name</display-name></channel>
<programme start="${_stamp(start)}" stop="${_stamp(start.add(const Duration(hours: 1)))}" channel="$id"><title>$title</title></programme>''';
test('« F3 ALPES » du panneau trouve France.3.-.Alpes.fr', () async {
final now = DateTime.now().toUtc();
final svc = serviceWith('<tv>${channel('France.3.-.Alpes.fr', 'France 3 Alpes', 'JT Alpes', now)}</tv>');
final hits = await svc.programmesFor(null, displayName: 'FR - F3 ALPES HD');
expect(hits.single.title, 'JT Alpes');
});
test('« AB3 » trouve AB3.Channel.fr (clé souple)', () async {
final now = DateTime.now().toUtc();
final svc = serviceWith('<tv>${channel('AB3.Channel.fr', 'AB3 Channel', 'Série', now)}</tv>');
final hits = await svc.programmesFor('AB3.fr', displayName: 'FR - AB3 FHD');
expect(hits.single.title, 'Série');
});
test('« TF1 +1 » reprend le guide de TF1 décalé d\'une heure', () async {
final now = DateTime.now().toUtc();
final svc = serviceWith('<tv>${channel('TF1.fr', 'TF1', 'Le 13h', now)}</tv>');
final hits = await svc.programmesFor(null, displayName: 'FR - TF1 +1 FHD');
expect(hits.single.title, 'Le 13h');
expect(hits.single.start, now.add(const Duration(hours: 1)).copyWith(microsecond: 0, millisecond: 0));
});
test('une chaîne +1 qui a son propre guide le garde', () async {
final now = DateTime.now().toUtc();
final svc = serviceWith('<tv>'
'${channel('Boomerang.fr', 'Boomerang', 'Normal', now)}'
'${channel('Boomerang.+1.fr', 'Boomerang +1', 'Décalé', now)}</tv>');
final hits = await svc.programmesFor(null, displayName: 'FR - BOOMERANG +1');
expect(hits.single.title, 'Décalé');
});
});
group('programmesFor par nom', () {
test('trouve le guide d\'une chaîne sans epg_channel_id par son nom nettoyé',
() async {
final now = DateTime.now().toUtc();
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';
}
final dump = '''
<tv>
<channel id="TF1.fr"><display-name>TF1 HD</display-name><display-name>TF1</display-name></channel>
<programme start="${stamp(now)}" stop="${stamp(now.add(const Duration(hours: 1)))}" channel="TF1.fr">
<title>Le 13h</title>
</programme>
</tv>
''';
final svc = XmltvEpgService(
sourceUrls: const ['http://ext/epg.xml'],
client: MockClient((_) async => http.Response(dump, 200)),
);
final hits = await svc.programmesFor(null, displayName: 'FR - TF1 FHD ◉');
expect(hits.single.title, 'Le 13h');
});
test('indexe chaque display-name, pas seulement le premier', () {
final now = DateTime.now().toUtc();
final index = service.parseForTest('''
<tv>
<channel id="x.fr"><display-name>Premier Nom</display-name><display-name>Second Nom</display-name></channel>
<programme start="${_stamp(now)}" stop="${_stamp(now.add(const Duration(hours: 1)))}" channel="x.fr"><title>T</title></programme>
</tv>
''');
expect(index[XmltvEpgService.normalizeKey('Second Nom')], isNotNull);
});
});
group('parseXmltvDate', () {
test('applique le décalage horaire annoncé', () {
expect(
@@ -148,6 +269,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
+84
View File
@@ -0,0 +1,84 @@
import 'dart:convert';
import 'dart:io';
import 'package:crypto/crypto.dart';
/// Versionnage des assets non hachés de la build web.
///
/// POURQUOI : `main.dart.js`, `flutter_bootstrap.js` et `xf-player-core.js`
/// gardent le même nom d'une build à l'autre. Un reverse proxy qui met les
/// `.js` en cache (option « Cache Assets » de Nginx Proxy Manager, constatée
/// en prod : cache jusqu'à 12 h, même sur un rechargement forcé) servait donc
/// l'ANCIEN lecteur et l'ANCIENNE app avec le NOUVEAU serveur après chaque
/// déploiement. En ajoutant `?v=<empreinte du contenu>` aux références,
/// chaque version a sa propre URL : aucun cache intermédiaire ne peut la
/// confondre avec la précédente.
/// Empreinte courte et stable de [bytes].
String contentVersion(List<int> bytes) =>
md5.convert(bytes).toString().substring(0, 12);
/// Ajoute `?v=<version>` à chaque référence entre guillemets à un asset de
/// [versions] (`"main.dart.js"` → `"main.dart.js?v=…"`).
///
/// Seules les références exactes sont réécrites : `vendor/x.js` ou une
/// référence déjà versionnée restent intactes.
String versionAssetRefs(String content, Map<String, String> versions) {
var out = content;
versions.forEach((asset, version) {
final pattern = RegExp('(["\'])${RegExp.escape(asset)}\\1');
out = out.replaceAllMapped(
pattern,
(m) => '${m[1]}$asset?v=$version${m[1]}',
);
});
return out;
}
/// Contenus réécrits, indexés par chemin relatif (`index.html`…), prêts à
/// être servis à la place des fichiers d'origine.
///
/// Calculé une fois au démarrage : la build web est figée dans l'image.
Map<String, String> buildVersionedEntryPoints(String webPath) {
String? read(String name) {
final file = File('$webPath/$name');
return file.existsSync() ? file.readAsStringSync() : null;
}
String? versionOf(String name) {
final file = File('$webPath/$name');
return file.existsSync() ? contentVersion(file.readAsBytesSync()) : null;
}
final rewritten = <String, String>{};
// Le bootstrap référence main.dart.js : sa version dépend donc de celle
// de main.dart.js, d'où une version calculée APRÈS réécriture.
final bootstrap = read('flutter_bootstrap.js');
final mainVersion = versionOf('main.dart.js');
if (bootstrap != null && mainVersion != null) {
final versioned =
versionAssetRefs(bootstrap, {'main.dart.js': mainVersion});
rewritten['flutter_bootstrap.js'] = versioned;
final index = read('index.html');
if (index != null) {
rewritten['index.html'] = versionAssetRefs(index, {
'flutter_bootstrap.js': contentVersion(utf8.encode(versioned)),
});
}
}
final playerAssets = <String, String>{
for (final name in const ['xf-player-core.js', 'xf-player.css'])
if (versionOf(name) case final v?) name: v,
};
for (final page in const [
'player.html',
'player_lite.html',
'player_mobile.html',
]) {
final html = read(page);
if (html != null) rewritten[page] = versionAssetRefs(html, playerAssets);
}
return rewritten;
}
+1 -1
View File
@@ -9,7 +9,7 @@ services:
- MAX_CONCURRENT_RECORDINGS=2
# Source XMLTV de repli, utilisée seulement quand le panneau ne rend
# aucun programme actuel. Vider la variable coupe tout appel sortant.
- EPG_XMLTV_URLS=https://epgshare01.online/epgshare01/epg_ripper_FR1.xml.gz
- EPG_XMLTV_URLS=https://epgshare01.online/epgshare01/epg_ripper_FR1.xml.gz,https://epgshare01.online/epgshare01/epg_ripper_CH1.xml.gz,https://epgshare01.online/epgshare01/epg_ripper_BE2.xml.gz
- NVIDIA_GPU=true
volumes:
- xtremflow-data:/app/data
+1 -1
View File
@@ -27,7 +27,7 @@ services:
# aucun programme couvrant l'instant présent. Vider la variable coupe
# tout appel sortant : l'EPG se limite alors au panneau de l'abonné.
# Plusieurs dumps se séparent par des virgules.
- EPG_XMLTV_URLS=${EPG_XMLTV_URLS-https://epgshare01.online/epgshare01/epg_ripper_FR1.xml.gz}
- EPG_XMLTV_URLS=${EPG_XMLTV_URLS-https://epgshare01.online/epgshare01/epg_ripper_FR1.xml.gz,https://epgshare01.online/epgshare01/epg_ripper_CH1.xml.gz,https://epgshare01.online/epgshare01/epg_ripper_BE2.xml.gz}
# Accélération GPU NVIDIA pour le transcodage (NVENC) : true/false
- NVIDIA_GPU=${NVIDIA_GPU:-false}
restart: unless-stopped
+12 -2
View File
@@ -191,7 +191,16 @@ class XtreamService {
String quality = 'high',
}) {
if (_currentPlaylist == null) throw Exception('No playlist configured');
return '$_backendBaseUrl/api/vod/$streamId/$quality/playlist.m3u8';
// `ext` : le serveur demande le fichier au panneau sous son extension
// réelle — un `.mp4` réclamé en `.mkv` est refusé (HTTP 551).
return '$_backendBaseUrl/api/vod/$streamId/$quality/playlist.m3u8'
'${_extQuery(containerExtension, first: true)}';
}
String _extQuery(String ext, {required bool first}) {
final value = ext.trim().toLowerCase();
if (value.isEmpty) return '';
return '${first ? '?' : '&'}ext=${Uri.encodeQueryComponent(value)}';
}
/// Generate stream URL for series episodes
@@ -201,7 +210,8 @@ class XtreamService {
String quality = 'high',
}) {
if (_currentPlaylist == null) throw Exception('No playlist configured');
return '$_backendBaseUrl/api/vod/$streamId/$quality/playlist.m3u8?type=series';
return '$_backendBaseUrl/api/vod/$streamId/$quality/playlist.m3u8?type=series'
'${_extQuery(containerExtension, first: false)}';
}
/// Authenticate and get server info
+110
View File
@@ -0,0 +1,110 @@
// Tests du moteur de lecture web (web/xf-player-core.js).
//
// Lancement : node --test "test/web/*_test.mjs"
//
// Le moteur est un script navigateur (IIFE sur `window`) : on l'exécute dans
// un contexte vm avec un DOM minimal et un faux mpegts.js qui capture la
// configuration reçue.
import { test } from 'node:test';
import assert from 'node:assert/strict';
import { readFileSync } from 'node:fs';
import vm from 'node:vm';
const source = readFileSync(
new URL('../../web/xf-player-core.js', import.meta.url),
'utf8',
);
function fakeEventTarget() {
return { addEventListener() {}, removeEventListener() {} };
}
/** Charge le moteur et renvoie la config passée à mpegts.createPlayer. */
function liveMpegtsConfig({ profile } = {}) {
const created = [];
const video = {
...fakeEventTarget(),
canPlayType: () => '',
play: () => Promise.resolve(),
buffered: { length: 0 },
seekable: { length: 0 },
};
const window = {
...fakeEventTarget(),
console: { log() {} },
location: { origin: 'https://xf.test', href: 'https://xf.test/', search: '' },
parent: { postMessage() {} },
mpegts: {
isSupported: () => true,
Events: { MEDIA_INFO: 'media_info', ERROR: 'error' },
createPlayer(mediaDataSource, config) {
created.push({ mediaDataSource, config });
return { on() {}, attachMediaElement() {}, load() {} };
},
},
};
const context = vm.createContext({
window,
location: window.location,
document: fakeEventTarget(),
URLSearchParams,
setInterval: () => 0,
clearInterval() {},
setTimeout,
});
vm.runInContext(source, context);
const XFPlayer = window.XFPlayer;
const player = new XFPlayer({
video,
url: '/api/live/1/turbo.ts',
type: 'live',
});
if (profile) player.profile = window.XFPlayerProfiles[profile];
player._createMpegts(true);
assert.equal(created.length, 1);
return created[0].config;
}
// Les panneaux Xtream relaient souvent une source HLS : le flux arrive par
// rafales d'environ 6 s séparées de silences complets (mesuré en prod).
const UPSTREAM_BURST_GAP_S = 7;
for (const profile of ['fast', 'balanced', 'safe']) {
test(`live (${profile}) : la marge gardée après rattrapage couvre un silence amont`, () => {
const c = liveMpegtsConfig({ profile });
if (!c.liveBufferLatencyChasing) return; // pas de saut = rien à vérifier
assert.ok(
c.liveBufferLatencyMinRemain >= UPSTREAM_BURST_GAP_S,
`MinRemain ${c.liveBufferLatencyMinRemain}s < silence amont ${UPSTREAM_BURST_GAP_S}s : ` +
'chaque rafale déclenche un saut puis une coupure',
);
});
test(`live (${profile}) : une rafale normale ne déclenche pas de saut`, () => {
const c = liveMpegtsConfig({ profile });
if (!c.liveBufferLatencyChasing) return;
// Après une rafale, le buffer vaut la marge + une rafale entière.
assert.ok(
c.liveBufferLatencyMaxLatency >=
c.liveBufferLatencyMinRemain + 2 * UPSTREAM_BURST_GAP_S,
'seuil de rattrapage trop proche de la marge : sauts à chaque rafale',
);
});
}
test('live : le retard est résorbé en douceur (liveSync) plutôt que par sauts', () => {
const c = liveMpegtsConfig();
assert.equal(c.liveSync, true);
assert.ok(c.liveSyncPlaybackRate > 1 && c.liveSyncPlaybackRate <= 1.2);
assert.ok(c.liveSyncTargetLatency < c.liveSyncMaxLatency);
assert.ok(c.liveSyncMaxLatency <= c.liveBufferLatencyMaxLatency);
});
test('live : le buffer déjà lu est purgé (sinon QuotaExceeded en séance longue)', () => {
const c = liveMpegtsConfig();
assert.equal(c.autoCleanupSourceBuffer, true);
assert.ok(
c.autoCleanupMinBackwardDuration < c.autoCleanupMaxBackwardDuration,
);
});
+22 -7
View File
@@ -419,13 +419,28 @@
// 64 Ko en profil « fast » contre 512 Ko auparavant.
stashInitialSize: this.profile.mpegtsStash,
maxStashSize: 30 * 1024 * 1024,
// Rattrapage du direct : sans lui, chaque micro-coupure fait dériver
// la lecture derrière le direct (latence qui s'accumule). Activé
// uniquement en profil « fast » — les profils élargis privilégient
// la stabilité du buffer.
liveBufferLatencyChasing: isLive && this.profile.name === 'fast',
liveBufferLatencyMaxLatency: 5,
liveBufferLatencyMinRemain: 1,
// POURQUOI ces seuils : les panneaux Xtream relaient le plus souvent
// une source HLS et livrent le flux par RAFALES (~6 s de contenu d'un
// coup, puis plusieurs secondes de silence total — mesuré en prod).
// L'ancien rattrapage (5 s max, 1 s de marge) sautait en avant à
// chaque rafale puis tombait à sec pendant le silence suivant :
// coupure toutes les 3-4 s et contenu sauté (12 sauts et 14 coupures
// en 30 s sur une chaîne FHD). La marge de 12 s couvre deux silences ;
// le saut ne sert plus qu'à borner un retard réellement accumulé.
liveBufferLatencyChasing: isLive,
liveBufferLatencyMaxLatency: 30,
liveBufferLatencyMinRemain: 12,
// Au-delà de 20 s de retard, accélérer légèrement (x1,1) jusqu'à
// revenir à 12 s : rattrapage invisible, sans saut ni coupure.
liveSync: isLive,
liveSyncMaxLatency: 20,
liveSyncTargetLatency: 12,
liveSyncPlaybackRate: 1.1,
// Purge du buffer déjà lu : sans elle, une longue séance finit par
// saturer le SourceBuffer (QuotaExceededError → erreur MPEG-TS).
autoCleanupSourceBuffer: true,
autoCleanupMaxBackwardDuration: 60,
autoCleanupMinBackwardDuration: 30,
lazyLoad: false
}
);