mirror of
https://github.com/R0m1k3/xtremflow.git
synced 2026-10-11 17:30:00 +02:00
Compare commits
18
Commits
9f4e98fcea
...
main
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
f82c70ae05 | ||
|
|
6733f80ffd | ||
|
|
0539ad857f | ||
|
|
a739d6451c | ||
|
|
775caad0ee | ||
|
|
4d00409de8 | ||
|
|
0ea60ccbfa | ||
|
|
a3640a6b5f | ||
|
|
584068b45b | ||
|
|
abe7d522c1 | ||
|
|
8d79472ca8 | ||
|
|
6812a4964d | ||
|
|
fdf98f9659 | ||
|
|
6212948ea5 | ||
|
|
662e1034e5 | ||
|
|
a9b43ff581 | ||
|
|
327873d254 | ||
|
|
da67c479b3 |
No files matched your search
@@ -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:
|
||||
|
||||
@@ -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
@@ -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;
|
||||
|
||||
@@ -0,0 +1,70 @@
|
||||
import 'dart:convert';
|
||||
|
||||
import 'package:shelf/shelf.dart';
|
||||
|
||||
/// Journal des incidents de lecture remontés par le lecteur web.
|
||||
///
|
||||
/// POURQUOI : les freezes (image et son figés puis reprise) se produisent
|
||||
/// chez l'utilisateur, dans son navigateur ; le serveur n'en voit rien. Le
|
||||
/// lecteur envoie ici chaque coupure, erreur, recréation et saut, et le
|
||||
/// serveur les écrit dans `docker logs` (`[Player] …`) : on sait alors si un
|
||||
/// freeze vient d'un buffer à sec, d'une erreur du démuxeur ou d'une
|
||||
/// relance du flux.
|
||||
///
|
||||
/// Entrée non fiable : seuls des champs connus, courts et typés passent.
|
||||
|
||||
const _events = {
|
||||
'start', 'first_frame', 'stall', 'seek', 'mpegts_error', 'recreate',
|
||||
'fallback_hls', 'hls_error', 'profile', 'error_shown', 'heartbeat',
|
||||
};
|
||||
|
||||
const _numericFields = {
|
||||
't', 'ms', 'buffer', 'pos', 'from', 'to', 'stalls', 'stallMs', 'seeks',
|
||||
'errors', 'rate',
|
||||
};
|
||||
|
||||
const _textFields = {'session', 'ch', 'mode', 'profile', 'detail'};
|
||||
|
||||
/// Ligne de journal pour l'événement [raw], ou null s'il est invalide.
|
||||
String? formatPlayerEvent(Object? raw) {
|
||||
if (raw is! Map) return null;
|
||||
final event = raw['event'];
|
||||
if (event is! String || !_events.contains(event)) return null;
|
||||
|
||||
final parts = <String>[event];
|
||||
for (final key in _textFields) {
|
||||
final value = raw[key];
|
||||
if (value is String && value.isNotEmpty) {
|
||||
// Ni espace, ni retour ligne, ni URL : jamais d'injection de faux
|
||||
// journaux, jamais d'identifiants recopiés depuis une URL d'erreur.
|
||||
final safe = value
|
||||
.replaceAll(RegExp(r'https?://\S+'), '<url>')
|
||||
.replaceAll(RegExp(r'[^\w.\-:/<>]'), '_');
|
||||
parts.add('$key=${safe.length > 60 ? safe.substring(0, 60) : safe}');
|
||||
}
|
||||
}
|
||||
for (final key in _numericFields) {
|
||||
final value = raw[key];
|
||||
if (value is num && value.isFinite) {
|
||||
parts.add('$key=${value is int ? value : value.toStringAsFixed(1)}');
|
||||
}
|
||||
}
|
||||
return '[Player] ${parts.join(' ')}';
|
||||
}
|
||||
|
||||
Future<Response> handlePlayerLog(Request request) async {
|
||||
final body = await request.readAsString();
|
||||
if (body.length > 4096) return Response(413);
|
||||
Object? decoded;
|
||||
try {
|
||||
decoded = jsonDecode(body);
|
||||
} catch (_) {
|
||||
return Response.badRequest();
|
||||
}
|
||||
final events = decoded is List ? decoded.take(20) : [decoded];
|
||||
for (final event in events) {
|
||||
final line = formatPlayerEvent(event);
|
||||
if (line != null) print(line);
|
||||
}
|
||||
return Response(204);
|
||||
}
|
||||
@@ -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}
|
||||
|
||||
@@ -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};
|
||||
}
|
||||
}
|
||||
+61
-2
@@ -27,6 +27,9 @@ 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';
|
||||
import 'api/player_log.dart';
|
||||
|
||||
void main(List<String> args) async {
|
||||
// Parse command line arguments
|
||||
@@ -130,7 +133,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)
|
||||
@@ -217,6 +222,14 @@ void main(List<String> args) async {
|
||||
.addMiddleware(authMiddleware(db))
|
||||
.addHandler(logoApi.handle),
|
||||
)
|
||||
// Incidents de lecture remontés par le lecteur web (freezes, erreurs),
|
||||
// écrits dans les logs pour le diagnostic.
|
||||
..post(
|
||||
'/api/player-log',
|
||||
const Pipeline()
|
||||
.addMiddleware(authMiddleware(db))
|
||||
.addHandler(handlePlayerLog),
|
||||
)
|
||||
// EPG - guide TV (auth required)
|
||||
..mount(
|
||||
'/api/epg',
|
||||
@@ -268,6 +281,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 +329,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 ||
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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();
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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':
|
||||
|
||||
@@ -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));
|
||||
});
|
||||
});
|
||||
}
|
||||
@@ -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);
|
||||
});
|
||||
});
|
||||
}
|
||||
@@ -0,0 +1,47 @@
|
||||
import 'package:test/test.dart';
|
||||
import '../api/player_log.dart';
|
||||
|
||||
void main() {
|
||||
group('formatPlayerEvent', () {
|
||||
test('écrit une coupure avec sa durée et le buffer', () {
|
||||
expect(
|
||||
formatPlayerEvent({
|
||||
'event': 'stall',
|
||||
'session': 'ab12',
|
||||
'ch': '554034',
|
||||
'ms': 5230,
|
||||
'buffer': 0.04,
|
||||
}),
|
||||
'[Player] stall session=ab12 ch=554034 ms=5230 buffer=0.0',
|
||||
);
|
||||
});
|
||||
|
||||
test('refuse un événement inconnu ou un payload invalide', () {
|
||||
expect(formatPlayerEvent({'event': 'rm -rf'}), isNull);
|
||||
expect(formatPlayerEvent('stall'), isNull);
|
||||
expect(formatPlayerEvent(null), isNull);
|
||||
});
|
||||
|
||||
test('pas d\'injection de lignes ni d\'URL (identifiants Xtream)', () {
|
||||
final line = formatPlayerEvent({
|
||||
'event': 'mpegts_error',
|
||||
'detail': 'NetworkError http://panel/live/user/pass/1.ts\n[Player] faux',
|
||||
})!;
|
||||
expect(line, isNot(contains('\n')));
|
||||
expect(line, isNot(contains('user/pass')));
|
||||
expect(line, contains('<url>'));
|
||||
});
|
||||
|
||||
test('ignore les champs inattendus et les nombres non finis', () {
|
||||
final line = formatPlayerEvent({
|
||||
'event': 'seek',
|
||||
'from': 10,
|
||||
'to': double.infinity,
|
||||
'cookie': 'secret',
|
||||
})!;
|
||||
expect(line, contains('from=10'));
|
||||
expect(line, isNot(contains('to=')));
|
||||
expect(line, isNot(contains('secret')));
|
||||
});
|
||||
});
|
||||
}
|
||||
@@ -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')));
|
||||
});
|
||||
}
|
||||
@@ -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,
|
||||
);
|
||||
});
|
||||
});
|
||||
}
|
||||
@@ -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'));
|
||||
});
|
||||
});
|
||||
}
|
||||
@@ -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
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
@@ -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
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -0,0 +1,303 @@
|
||||
// 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',
|
||||
);
|
||||
});
|
||||
}
|
||||
|
||||
/** Réglages du régulateur de vitesse (exposés par le moteur). */
|
||||
function liveTuning() {
|
||||
const window = {
|
||||
addEventListener() {}, removeEventListener() {},
|
||||
console: { log() {} },
|
||||
location: { origin: 'https://xf.test', href: 'https://xf.test/', search: '' },
|
||||
parent: { postMessage() {} },
|
||||
};
|
||||
vm.runInContext(source, vm.createContext({
|
||||
window, location: window.location,
|
||||
document: { addEventListener() {}, removeEventListener() {} },
|
||||
URLSearchParams, setInterval: () => 0, clearInterval() {}, setTimeout,
|
||||
}));
|
||||
return window.XFPlayerLiveTuning;
|
||||
}
|
||||
|
||||
test('live : un seul régulateur de vitesse (liveSync de mpegts.js désactivé)', () => {
|
||||
// liveSync ne sait qu'accélérer et remet la vitesse à 1 sous sa cible :
|
||||
// il annulerait le ralenti qui reconstitue la réserve.
|
||||
const c = liveMpegtsConfig();
|
||||
assert.notEqual(c.liveSync, true);
|
||||
const t = liveTuning();
|
||||
assert.ok(t.slowRate < 1 && t.slowRate >= 0.95, 'ralenti audible au-delà');
|
||||
assert.ok(t.fastRate > 1 && t.fastRate <= 1.15);
|
||||
assert.ok(t.lowAhead < t.targetAhead && t.targetAhead < t.highAhead);
|
||||
assert.ok(t.highAhead < c.liveBufferLatencyMaxLatency);
|
||||
});
|
||||
|
||||
test('télémétrie : une coupure (waiting → playing) est remontée avec sa durée', async () => {
|
||||
const sent = [];
|
||||
const listeners = {};
|
||||
const video = {
|
||||
addEventListener(type, fn) { (listeners[type] ||= []).push(fn); },
|
||||
removeEventListener() {},
|
||||
canPlayType: () => '',
|
||||
play: () => Promise.resolve(),
|
||||
buffered: { length: 0 },
|
||||
seekable: { length: 0 },
|
||||
currentTime: 12,
|
||||
playbackRate: 1,
|
||||
};
|
||||
const fire = (type) => (listeners[type] || []).forEach((fn) => fn());
|
||||
const window = {
|
||||
addEventListener() {}, removeEventListener() {},
|
||||
console: { log() {} },
|
||||
location: { origin: 'https://xf.test', href: 'https://xf.test/', search: '' },
|
||||
parent: { postMessage() {} },
|
||||
localStorage: { getItem: () => 'tok' },
|
||||
};
|
||||
let now = 1000;
|
||||
const context = vm.createContext({
|
||||
window,
|
||||
location: window.location,
|
||||
document: { addEventListener() {}, removeEventListener() {} },
|
||||
URLSearchParams,
|
||||
setInterval: () => 0,
|
||||
clearInterval() {},
|
||||
setTimeout,
|
||||
Date: { now: () => now },
|
||||
fetch: (url, opts) => { sent.push({ url, body: JSON.parse(opts.body) }); return Promise.resolve(); },
|
||||
});
|
||||
vm.runInContext(source, context);
|
||||
const player = new window.XFPlayer({ video, url: '/api/live/554034/turbo.ts', type: 'live' });
|
||||
player._wireVideoEvents();
|
||||
|
||||
fire('playing'); // premier affichage
|
||||
now += 1000; fire('waiting');
|
||||
now += 4200; fire('playing'); // reprise 4,2 s plus tard
|
||||
|
||||
const stall = sent.map((s) => s.body).find((b) => b.event === 'stall');
|
||||
assert.ok(stall, 'aucun événement stall envoyé');
|
||||
assert.equal(stall.ms, 4200);
|
||||
assert.equal(stall.ch, '554034');
|
||||
assert.equal(sent[0].url, '/api/player-log');
|
||||
});
|
||||
|
||||
// Freeze mesuré en prod (télémétrie) : le panneau s'est tu ~15-20 s sur une
|
||||
// connexion restée ouverte, alors que le lecteur n'avait gardé que 11-16 s.
|
||||
const PROVIDER_OUTAGE_S = 18;
|
||||
|
||||
test('live : la réserve visée couvre un silence du panneau de ~18 s', () => {
|
||||
const c = liveMpegtsConfig();
|
||||
const t = liveTuning();
|
||||
// Le régulateur ramène l'avance vers la cible : c'est la réserve en
|
||||
// régime établi. Elle doit dépasser le silence mesuré — et le ralenti
|
||||
// doit se déclencher AVANT d'être passé sous ce seuil de sécurité.
|
||||
assert.ok(t.targetAhead > PROVIDER_OUTAGE_S,
|
||||
`cible ${t.targetAhead}s < silence mesuré ${PROVIDER_OUTAGE_S}s`);
|
||||
assert.ok(t.lowAhead > PROVIDER_OUTAGE_S);
|
||||
assert.ok(c.liveBufferLatencyMinRemain > PROVIDER_OUTAGE_S);
|
||||
});
|
||||
|
||||
test('régulateur : ralentit sous la réserve, accélère au-dessus, sans osciller', () => {
|
||||
const listeners = {};
|
||||
let ahead = 0;
|
||||
const video = {
|
||||
addEventListener(type, fn) { (listeners[type] ||= []).push(fn); },
|
||||
removeEventListener() {},
|
||||
canPlayType: () => '',
|
||||
play: () => Promise.resolve(),
|
||||
pause() {},
|
||||
get buffered() { return { length: 1, start: () => 0, end: () => 100 + ahead }; },
|
||||
seekable: { length: 0 },
|
||||
currentTime: 100,
|
||||
playbackRate: 1,
|
||||
paused: false,
|
||||
seeking: false,
|
||||
};
|
||||
const window = {
|
||||
addEventListener() {}, removeEventListener() {},
|
||||
console: { log() {} },
|
||||
location: { origin: 'https://xf.test', href: 'https://xf.test/', search: '' },
|
||||
parent: { postMessage() {} },
|
||||
};
|
||||
vm.runInContext(source, vm.createContext({
|
||||
window, location: window.location,
|
||||
document: { addEventListener() {}, removeEventListener() {} },
|
||||
URLSearchParams, setInterval: () => 0, clearInterval() {}, setTimeout, Date,
|
||||
}));
|
||||
const t = window.XFPlayerLiveTuning;
|
||||
const p = new window.XFPlayer({ video, url: '/api/live/1/turbo.ts', type: 'live' });
|
||||
p.mpegts = {};
|
||||
p._started = true;
|
||||
|
||||
ahead = 13; p._tuneRate(); // session réelle : 13 s
|
||||
assert.equal(video.playbackRate, t.slowRate);
|
||||
ahead = t.lowAhead + 2; p._tuneRate(); // remonte, pas encore la cible
|
||||
assert.equal(video.playbackRate, t.slowRate, 'repasse à 1 trop tôt : oscille');
|
||||
ahead = t.targetAhead; p._tuneRate();
|
||||
assert.equal(video.playbackRate, 1);
|
||||
ahead = t.highAhead + 5; p._tuneRate(); // rafale de rattrapage
|
||||
assert.equal(video.playbackRate, t.fastRate);
|
||||
ahead = t.highAhead - 2; p._tuneRate();
|
||||
assert.equal(video.playbackRate, t.fastRate, 'freine trop tôt : oscille');
|
||||
ahead = t.targetAhead + 2; p._tuneRate();
|
||||
assert.equal(video.playbackRate, 1);
|
||||
});
|
||||
|
||||
test('live : la rafale de rattrapage après un silence ne fait pas sauter le programme', () => {
|
||||
// Mesuré en prod : après un silence, le panneau renvoie d'un coup tout le
|
||||
// retard ; l'avance a dépassé 45 s et le lecteur a sauté 2 fois de 20 s
|
||||
// (programme perdu). Le retard doit se résorber en douceur (liveSync).
|
||||
const observedPeakAhead = 45;
|
||||
const c = liveMpegtsConfig();
|
||||
assert.ok(c.liveBufferLatencyMaxLatency >= 2 * observedPeakAhead - 10,
|
||||
`saut dès ${c.liveBufferLatencyMaxLatency}s d'avance : la rafale de rattrapage le déclenche`);
|
||||
assert.ok(liveTuning().highAhead < observedPeakAhead,
|
||||
'le régulateur doit déjà accélérer avant ce niveau');
|
||||
});
|
||||
|
||||
test('live : la vidéo déjà vue est purgée plus tôt, l\'avance étant plus grande', () => {
|
||||
const c = liveMpegtsConfig();
|
||||
assert.ok(c.autoCleanupMaxBackwardDuration <= 30);
|
||||
});
|
||||
|
||||
test('reprise après coupure : attendre une vraie réserve avant de relire', () => {
|
||||
const calls = [];
|
||||
const listeners = {};
|
||||
let ahead = 0;
|
||||
const video = {
|
||||
addEventListener(type, fn) { (listeners[type] ||= []).push(fn); },
|
||||
removeEventListener() {},
|
||||
canPlayType: () => '',
|
||||
play: () => { calls.push('play'); return Promise.resolve(); },
|
||||
pause: () => { calls.push('pause'); },
|
||||
get buffered() { return { length: 1, start: () => 0, end: () => 100 + ahead }; },
|
||||
seekable: { length: 0 },
|
||||
currentTime: 100,
|
||||
playbackRate: 1,
|
||||
seeking: false,
|
||||
};
|
||||
const fire = (type) => (listeners[type] || []).forEach((fn) => fn());
|
||||
const window = {
|
||||
addEventListener() {}, removeEventListener() {},
|
||||
console: { log() {} },
|
||||
location: { origin: 'https://xf.test', href: 'https://xf.test/', search: '' },
|
||||
parent: { postMessage() {} },
|
||||
};
|
||||
const context = vm.createContext({
|
||||
window, location: window.location,
|
||||
document: { addEventListener() {}, removeEventListener() {} },
|
||||
URLSearchParams, setInterval: () => 0, clearInterval() {}, setTimeout, Date,
|
||||
});
|
||||
vm.runInContext(source, context);
|
||||
const player = new window.XFPlayer({ video, url: '/api/live/1/turbo.ts', type: 'live' });
|
||||
player.mpegts = {}; // flux MPEG-TS en cours
|
||||
player._wireVideoEvents();
|
||||
|
||||
fire('playing');
|
||||
fire('waiting'); // le panneau se tait
|
||||
assert.deepEqual(calls.slice(-1), ['pause'], 'le lecteur doit se mettre en attente');
|
||||
|
||||
ahead = 1; // 1 s arrive : trop peu, ne pas repartir
|
||||
player._checkHold();
|
||||
assert.ok(!calls.slice(1).includes('play'), 'reparti avec 1 s : recoupera aussitôt');
|
||||
|
||||
ahead = 5; // vraie réserve reconstituée
|
||||
player._checkHold();
|
||||
assert.equal(calls.at(-1), 'play');
|
||||
});
|
||||
|
||||
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,
|
||||
);
|
||||
});
|
||||
@@ -141,6 +141,12 @@
|
||||
|
||||
// Contrat historique du lite : mise en pause manuelle → overlay lecture.
|
||||
video.addEventListener('pause', function () {
|
||||
// Pause technique : le lecteur attend que la réserve se
|
||||
// reconstitue après une coupure, il repartira seul.
|
||||
if (player && player._holding) {
|
||||
showLoading('Reconstitution de la réserve');
|
||||
return;
|
||||
}
|
||||
if (video.currentTime === 0 || video.paused) {
|
||||
loading.style.display = 'none';
|
||||
playOverlay.style.display = 'flex';
|
||||
|
||||
+251
-8
@@ -58,6 +58,32 @@
|
||||
// Erreurs MPEG-TS tolérées sur 60 s avant de basculer sur HLS.
|
||||
var MPEGTS_MAX_ERRORS = 3;
|
||||
|
||||
// Après une coupure en direct, réserve à reconstituer avant de relire.
|
||||
// Repartir dès la première seconde reçue transformait un silence du
|
||||
// panneau en cascade de micro-coupures (6,3 s puis 1,7 s puis 2,0 s,
|
||||
// mesuré) ; comme mpv (cache-pause), on attend une vraie réserve.
|
||||
var REBUFFER_SECONDS = 4;
|
||||
// Au-delà, on relance quand même avec ce qu'on a.
|
||||
var REBUFFER_MAX_WAIT_MS = 15000;
|
||||
|
||||
// Régulateur de vitesse du direct MPEG-TS (voir _tuneRate).
|
||||
//
|
||||
// La réserve ne venait que de l'avance envoyée par le panneau au
|
||||
// démarrage : 28 s une fois, 18 s une autre — et avec 13-20 s, un silence
|
||||
// du panneau a fini en coupure de 5,5 s (télémétrie). Sous lowAhead, on
|
||||
// lit à x0,97 (hauteur du son préservée) jusqu'à reconstituer targetAhead ;
|
||||
// au-delà de highAhead (rafale de rattrapage du panneau), x1,1 jusqu'à
|
||||
// settleHigh. Les seuils d'entrée et de sortie diffèrent : pas
|
||||
// d'oscillation de vitesse.
|
||||
var LIVE_TUNING = {
|
||||
lowAhead: 20,
|
||||
targetAhead: 25,
|
||||
settleHigh: 28,
|
||||
highAhead: 35,
|
||||
slowRate: 0.97,
|
||||
fastRate: 1.1
|
||||
};
|
||||
|
||||
// ------------------------------------------------------------------
|
||||
// Chargement paresseux des libs (aucune requête inutile)
|
||||
// ------------------------------------------------------------------
|
||||
@@ -190,8 +216,75 @@
|
||||
// ça, une relance après échec laisserait l'ancienne instance réagir aux
|
||||
// commandes du parent et republier des positions périmées.
|
||||
this._listeners = [];
|
||||
|
||||
// Télémétrie des incidents (voir _report) : chaque freeze vécu par
|
||||
// l'utilisateur finit dans les logs du serveur, avec sa cause.
|
||||
this._t0 = Date.now();
|
||||
this._session = Math.random().toString(36).slice(2, 8);
|
||||
var idMatch = /\/api\/(?:live|vod|recordings\/stream)\/([^/?#.]+)/.exec(this.url || '');
|
||||
this._channel = idMatch ? idMatch[1] : '';
|
||||
this._stallSince = null;
|
||||
this._stats = { stalls: 0, stallMs: 0, seeks: 0, errors: 0 };
|
||||
this._lastPos = 0;
|
||||
var self = this;
|
||||
var userOnError = this.onError;
|
||||
this.onError = function (msg) {
|
||||
self._report('error_shown', { detail: String(msg) });
|
||||
userOnError(msg);
|
||||
};
|
||||
}
|
||||
|
||||
// ---------------- Télémétrie ----------------
|
||||
|
||||
/** Secondes d'avance dans le buffer à la position courante. */
|
||||
XFPlayer.prototype._bufferAhead = function () {
|
||||
var v = this.video;
|
||||
try {
|
||||
for (var i = 0; i < v.buffered.length; i++) {
|
||||
if (v.buffered.start(i) <= v.currentTime + 0.5 &&
|
||||
v.buffered.end(i) >= v.currentTime) {
|
||||
return v.buffered.end(i) - v.currentTime;
|
||||
}
|
||||
}
|
||||
} catch (e) {}
|
||||
return 0;
|
||||
};
|
||||
|
||||
/**
|
||||
* Remonte un incident au serveur (`POST /api/player-log`, écrit dans les
|
||||
* logs `[Player] …`). Jamais bloquant : erreurs ignorées, aucun envoi
|
||||
* sans session.
|
||||
*/
|
||||
XFPlayer.prototype._report = function (event, data) {
|
||||
try {
|
||||
var token = null;
|
||||
try { token = global.localStorage.getItem('auth_token'); } catch (e) {}
|
||||
if (!token || typeof fetch !== 'function') return;
|
||||
var payload = {
|
||||
event: event,
|
||||
session: this._session,
|
||||
ch: this._channel,
|
||||
t: (Date.now() - this._t0) / 1000,
|
||||
mode: this.hls ? 'hls' : (this.mpegts ? 'ts' : 'direct'),
|
||||
profile: this.profile.name,
|
||||
buffer: this._bufferAhead(),
|
||||
pos: this.video.currentTime || 0
|
||||
};
|
||||
for (var k in data) {
|
||||
if (Object.prototype.hasOwnProperty.call(data, k)) payload[k] = data[k];
|
||||
}
|
||||
fetch('/api/player-log', {
|
||||
method: 'POST',
|
||||
keepalive: true,
|
||||
headers: {
|
||||
'Content-Type': 'application/json',
|
||||
'Authorization': 'Bearer ' + token
|
||||
},
|
||||
body: JSON.stringify(payload)
|
||||
}).catch(function () {});
|
||||
} catch (e) {}
|
||||
};
|
||||
|
||||
/** addEventListener + mémorisation pour un retrait propre. */
|
||||
XFPlayer.prototype._on = function (target, type, handler, options) {
|
||||
target.addEventListener(type, handler, options);
|
||||
@@ -224,6 +317,7 @@
|
||||
var isLive = this.type === 'live';
|
||||
|
||||
this.onLoading('Ouverture du flux');
|
||||
this._report('start', {});
|
||||
|
||||
// Safari / iOS : HLS natif = le chemin le plus rapide, aucune lib à charger.
|
||||
if (isHls && this._canPlayNativeHls()) {
|
||||
@@ -345,6 +439,8 @@
|
||||
hls.on(Hls.Events.ERROR, function (e, d) {
|
||||
if (!d.fatal) return;
|
||||
self.log('Erreur HLS fatale : ' + d.details);
|
||||
self._stats.errors++;
|
||||
self._report('hls_error', { detail: d.type + '/' + d.details });
|
||||
if (d.type === Hls.ErrorTypes.NETWORK_ERROR) {
|
||||
// Playlist jamais obtenue (transcodeur en échec → 502, serveur trop
|
||||
// chargé → délai dépassé), après les relances internes de hls.js.
|
||||
@@ -419,13 +515,41 @@
|
||||
// 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é.
|
||||
//
|
||||
// 25 s et non plus 12 : la télémétrie a mesuré un silence du panneau
|
||||
// de ~15-20 s sur une connexion restée ouverte, alors que le lecteur
|
||||
// n'avait gardé que 11-16 s — il venait de « consommer » à x1,1
|
||||
// l'avance que le panneau envoie au démarrage. On la garde : ~25 s
|
||||
// de retard sur le direct, en échange d'une lecture qui traverse ces
|
||||
// silences sans figer.
|
||||
//
|
||||
// Saut seulement au-delà de 90 s : après un silence, le panneau
|
||||
// renvoie d'un coup tout le retard. Avec un seuil à 45 s, l'avance
|
||||
// le dépassait et le lecteur sautait 20 s de programme (2 fois en
|
||||
// 1 min, mesuré). Le régulateur de vitesse (_tuneRate) résorbe
|
||||
// désormais cet excédent en douceur.
|
||||
liveBufferLatencyChasing: isLive,
|
||||
liveBufferLatencyMaxLatency: 90,
|
||||
liveBufferLatencyMinRemain: 25,
|
||||
// liveSync de mpegts.js désactivé : il ne sait qu'accélérer et remet
|
||||
// la vitesse à 1 sous sa cible, ce qui annulerait le ralenti de
|
||||
// _tuneRate qui reconstitue la réserve.
|
||||
liveSync: false,
|
||||
// Purge du buffer déjà lu : sans elle, une longue séance finit par
|
||||
// saturer le SourceBuffer (QuotaExceededError → erreur MPEG-TS).
|
||||
// Avance jusqu'à 90 s : la vidéo déjà vue est purgée plus tôt pour
|
||||
// rester sous le quota du SourceBuffer (~150 Mo en vidéo sur Chrome).
|
||||
autoCleanupSourceBuffer: true,
|
||||
autoCleanupMaxBackwardDuration: 30,
|
||||
autoCleanupMinBackwardDuration: 15,
|
||||
lazyLoad: false
|
||||
}
|
||||
);
|
||||
@@ -447,6 +571,7 @@
|
||||
|
||||
self._audioFallbackDone = true;
|
||||
self.log('Aucune piste audio démuxée (codec non supporté) — bascule HLS');
|
||||
self._report('fallback_hls', { detail: 'no_audio' });
|
||||
self.onLoading('Piste audio incompatible — réencodage');
|
||||
|
||||
try {
|
||||
@@ -464,6 +589,8 @@
|
||||
player.on(mpegts.Events.ERROR, function (type, detail) {
|
||||
self.log('Erreur MPEG-TS : ' + type + ' / ' + detail);
|
||||
if (self._destroyed) return;
|
||||
self._stats.errors++;
|
||||
self._report('mpegts_error', { detail: type + '/' + detail });
|
||||
|
||||
// Chaque recréation relance un FFmpeg côté serveur : sur un serveur
|
||||
// déjà saturé, recréer sans fin toutes les 1,2 s ne faisait
|
||||
@@ -488,6 +615,7 @@
|
||||
return;
|
||||
}
|
||||
self.log('Erreurs MPEG-TS répétées — bascule HLS');
|
||||
self._report('fallback_hls', { detail: 'repeated_mpegts_errors' });
|
||||
self.onLoading('Ouverture du flux');
|
||||
self.url = hlsUrl;
|
||||
if (self._canPlayNativeHls()) self._playDirect(hlsUrl);
|
||||
@@ -504,6 +632,7 @@
|
||||
player.detachMediaElement();
|
||||
player.destroy();
|
||||
} catch (e) {}
|
||||
self._report('recreate', {});
|
||||
self._createMpegts(self._mpegtsIsLive);
|
||||
}, 1200);
|
||||
});
|
||||
@@ -570,6 +699,7 @@
|
||||
var next = PROFILES[ORDER[idx + 1]];
|
||||
this.profile = next;
|
||||
this.log('Escalade buffer → ' + next.name + ' (' + reason + ')');
|
||||
this._report('profile', { detail: reason });
|
||||
this.onProfileChange(next.name, reason);
|
||||
|
||||
// HLS : les fenêtres de buffer sont modifiables à chaud, pas besoin de
|
||||
@@ -587,6 +717,62 @@
|
||||
// recréation du lecteur (déclenchée par le handler d'erreur).
|
||||
};
|
||||
|
||||
/**
|
||||
* Direct MPEG-TS à sec : se mettre en pause jusqu'à [REBUFFER_SECONDS] de
|
||||
* réserve (voir la constante). HLS gère déjà sa propre reprise.
|
||||
*/
|
||||
XFPlayer.prototype._holdUntilBuffered = function () {
|
||||
if (this._holding || !this.mpegts || this.type !== 'live') return;
|
||||
var self = this;
|
||||
this._holding = true;
|
||||
this._holdSince = Date.now();
|
||||
try { this.video.pause(); } catch (e) {}
|
||||
this._holdTimer = setInterval(function () { self._checkHold(); }, 250);
|
||||
};
|
||||
|
||||
/**
|
||||
* Ajuste la vitesse du direct MPEG-TS selon l'avance (voir LIVE_TUNING).
|
||||
* Appelé chaque seconde ; sans effet en pause, en attente de réserve,
|
||||
* hors direct ou en HLS.
|
||||
*/
|
||||
XFPlayer.prototype._tuneRate = function () {
|
||||
var v = this.video;
|
||||
if (!this.mpegts || this.type !== 'live' || !this._started ||
|
||||
this._holding || v.paused) {
|
||||
return;
|
||||
}
|
||||
var ahead = this._bufferAhead();
|
||||
var t = LIVE_TUNING;
|
||||
var mode = this._rateMode || 'normal';
|
||||
if (mode === 'slow' && ahead >= t.targetAhead) mode = 'normal';
|
||||
else if (mode === 'fast' && ahead <= t.settleHigh) mode = 'normal';
|
||||
else if (mode === 'normal' && ahead < t.lowAhead) mode = 'slow';
|
||||
else if (mode === 'normal' && ahead > t.highAhead) mode = 'fast';
|
||||
// Un saut de rattrapage peut faire passer directement d'un extrême à
|
||||
// l'autre.
|
||||
if (mode === 'slow' && ahead > t.highAhead) mode = 'fast';
|
||||
if (mode === 'fast' && ahead < t.lowAhead) mode = 'slow';
|
||||
this._rateMode = mode;
|
||||
|
||||
var rate = mode === 'slow' ? t.slowRate : (mode === 'fast' ? t.fastRate : 1);
|
||||
if (v.playbackRate !== rate) {
|
||||
try { v.playbackRate = rate; } catch (e) {}
|
||||
}
|
||||
};
|
||||
|
||||
/** Relance la lecture dès que la réserve est reconstituée. */
|
||||
XFPlayer.prototype._checkHold = function () {
|
||||
if (!this._holding) return;
|
||||
if (this._destroyed) { clearInterval(this._holdTimer); return; }
|
||||
var waited = Date.now() - this._holdSince;
|
||||
if (this._bufferAhead() < REBUFFER_SECONDS && waited < REBUFFER_MAX_WAIT_MS) {
|
||||
return;
|
||||
}
|
||||
clearInterval(this._holdTimer);
|
||||
this._holding = false;
|
||||
this._attemptPlay();
|
||||
};
|
||||
|
||||
XFPlayer.prototype._noteStall = function () {
|
||||
if (!this._started) return; // les attentes d'amorçage ne comptent pas
|
||||
var now = Date.now();
|
||||
@@ -613,18 +799,65 @@
|
||||
});
|
||||
|
||||
this._on(v, 'playing', function () {
|
||||
if (!self._started) {
|
||||
self._report('first_frame', { ms: Date.now() - self._t0 });
|
||||
}
|
||||
if (self._stallSince !== null) {
|
||||
var stallMs = Date.now() - self._stallSince;
|
||||
self._stallSince = null;
|
||||
// < 250 ms : imperceptible, ne pas noyer les logs.
|
||||
if (stallMs >= 250) {
|
||||
self._stats.stalls++;
|
||||
self._stats.stallMs += stallMs;
|
||||
self._report('stall', { ms: stallMs });
|
||||
}
|
||||
}
|
||||
self._started = true;
|
||||
self.onReady();
|
||||
self.send({ type: 'playback_status', status: 'playing' });
|
||||
});
|
||||
|
||||
this._on(v, 'pause', function () {
|
||||
// Pause posée par l'attente de réserve : pas une pause utilisateur,
|
||||
// l'interface ne doit pas basculer sur « lecture ».
|
||||
if (self._holding) return;
|
||||
self.send({ type: 'playback_status', status: 'paused' });
|
||||
});
|
||||
|
||||
this._on(v, 'waiting', function () { self._noteStall(); });
|
||||
this._on(v, 'waiting', function () {
|
||||
if (self._started && self._stallSince === null) {
|
||||
self._stallSince = Date.now();
|
||||
}
|
||||
self._noteStall();
|
||||
if (self._started) self._holdUntilBuffered();
|
||||
});
|
||||
this._on(v, 'stalled', function () { self._noteStall(); });
|
||||
|
||||
// Sauts (rattrapage du direct, recherche) : d'où, vers où.
|
||||
this._on(v, 'timeupdate', function () {
|
||||
if (!v.seeking) self._lastPos = v.currentTime;
|
||||
});
|
||||
this._on(v, 'seeking', function () {
|
||||
if (!self._started) return;
|
||||
self._stats.seeks++;
|
||||
self._report('seek', { from: self._lastPos, to: v.currentTime });
|
||||
});
|
||||
|
||||
// Régulateur de vitesse du direct (voir LIVE_TUNING).
|
||||
this._rateTimer = setInterval(function () { self._tuneRate(); }, 1000);
|
||||
|
||||
// Bilan toutes les 60 s : nombre et durée cumulée des coupures.
|
||||
this._beatTimer = setInterval(function () {
|
||||
if (!self._started) return;
|
||||
self._report('heartbeat', {
|
||||
stalls: self._stats.stalls,
|
||||
stallMs: self._stats.stallMs,
|
||||
seeks: self._stats.seeks,
|
||||
errors: self._stats.errors,
|
||||
rate: v.playbackRate
|
||||
});
|
||||
}, 60000);
|
||||
|
||||
this._on(v, 'ended', function () {
|
||||
self.send({
|
||||
type: 'playback_ended',
|
||||
@@ -766,6 +999,12 @@
|
||||
self._attemptPlay();
|
||||
break;
|
||||
case 'pause':
|
||||
// Pause voulue par l'utilisateur : l'attente de réserve ne doit
|
||||
// pas relancer la lecture derrière son dos.
|
||||
if (self._holding) {
|
||||
clearInterval(self._holdTimer);
|
||||
self._holding = false;
|
||||
}
|
||||
v.pause();
|
||||
break;
|
||||
case 'seek':
|
||||
@@ -791,6 +1030,9 @@
|
||||
XFPlayer.prototype.destroy = function () {
|
||||
this._destroyed = true;
|
||||
clearInterval(this._posTimer);
|
||||
clearInterval(this._beatTimer);
|
||||
clearInterval(this._holdTimer);
|
||||
clearInterval(this._rateTimer);
|
||||
this._listeners.forEach(function (l) {
|
||||
try { l.target.removeEventListener(l.type, l.handler); } catch (e) {}
|
||||
});
|
||||
@@ -807,4 +1049,5 @@
|
||||
|
||||
global.XFPlayer = XFPlayer;
|
||||
global.XFPlayerProfiles = PROFILES;
|
||||
global.XFPlayerLiveTuning = LIVE_TUNING;
|
||||
})(window);
|
||||
Reference in new issue
Block a user