mirror of
https://github.com/R0m1k3/xtremflow.git
synced 2026-10-12 01:36:20 +02:00
Compare commits
11
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
fdf98f9659 | ||
|
|
6212948ea5 | ||
|
|
662e1034e5 | ||
|
|
a9b43ff581 | ||
|
|
327873d254 | ||
|
|
da67c479b3 | ||
|
|
9f4e98fcea | ||
|
|
117523f153 | ||
|
|
08d841180f | ||
|
|
73469060ff | ||
|
|
db13c85c1d |
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:
|
||||
|
||||
@@ -257,7 +257,12 @@ 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 [];
|
||||
}
|
||||
}
|
||||
|
||||
+232
-59
@@ -5,6 +5,8 @@ import 'package:http/http.dart' as http;
|
||||
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';
|
||||
@@ -17,6 +19,50 @@ final Directory _hlsTempDir =
|
||||
/// Global FFmpeg session registry (initialized in [initStreaming]).
|
||||
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));
|
||||
}
|
||||
|
||||
/// Inscrit [process] comme connexion amont de [playlist] tant qu'il vit.
|
||||
void _registerUpstream(
|
||||
PlaylistConfig playlist,
|
||||
String id,
|
||||
Process process, {
|
||||
required void Function() release,
|
||||
}) {
|
||||
final token = _upstreamSlots.register(
|
||||
id: id,
|
||||
account: accountKeyOf(playlist),
|
||||
release: release,
|
||||
);
|
||||
process.exitCode.then((_) => _upstreamSlots.unregister(id, token: token));
|
||||
}
|
||||
|
||||
/// Helper to resolve FFmpeg path (SYSTEM PATH vs Portable)
|
||||
String _getFFmpegPath() {
|
||||
if (Platform.isWindows) {
|
||||
@@ -53,6 +99,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) {
|
||||
@@ -181,6 +237,89 @@ Map<String, String> _segmentHeaders({int maxAge = 60}) => {
|
||||
'Cache-Control': 'max-age=$maxAge',
|
||||
};
|
||||
|
||||
typedef _SessionAttempt = ({FfmpegSession? session, bool ready, String? error});
|
||||
|
||||
/// Démarre (ou récupère) la session [id] et attend sa playlist.
|
||||
///
|
||||
/// 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(
|
||||
id: id,
|
||||
isLive: isLive,
|
||||
ffmpegPath: _getFFmpegPath(),
|
||||
argsBuilder: argsBuilder,
|
||||
);
|
||||
} 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),
|
||||
);
|
||||
}
|
||||
final outcome =
|
||||
await sessionManager.waitForPlaylist(session, minSegments: minSegments);
|
||||
return (session: session, ready: outcome.ready, error: outcome.error);
|
||||
}
|
||||
|
||||
/// Comme [_runSession], avec repli sur le CPU si l'encodage GPU échoue.
|
||||
///
|
||||
/// Le repli garde le même identifiant de session : les requêtes de segments
|
||||
/// qui suivent tombent sur la session CPU sans que le lecteur s'en aperçoive.
|
||||
/// Une panne GPU fait aussi partir les sessions suivantes directement sur le
|
||||
/// CPU pendant le délai de [GpuHealth], au lieu d'échouer une à une.
|
||||
Future<_SessionAttempt> _runWithGpuFallback({
|
||||
required String id,
|
||||
required bool isLive,
|
||||
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(
|
||||
id: id,
|
||||
isLive: isLive,
|
||||
argsBuilder: (_) => buildArgs(gpu),
|
||||
minSegments: minSegments,
|
||||
upstream: upstream,
|
||||
);
|
||||
|
||||
if (!attempt.ready && gpu && isGpuFailure(attempt.error)) {
|
||||
_gpuHealth.markFailed();
|
||||
print('[GPU] $id : encodage GPU refusé '
|
||||
'(${LogRedactor.redactUrl(attempt.error ?? '')}) — repli sur le CPU');
|
||||
sessionManager.killSession(id);
|
||||
attempt = await _runSession(
|
||||
id: id,
|
||||
isLive: isLive,
|
||||
argsBuilder: (_) => buildArgs(false),
|
||||
minSegments: minSegments,
|
||||
upstream: upstream,
|
||||
);
|
||||
}
|
||||
return attempt;
|
||||
}
|
||||
|
||||
// ==========================================
|
||||
// 1. LIVE TV HANDLER (FFmpeg HLS + direct TS proxy)
|
||||
// ==========================================
|
||||
@@ -238,7 +377,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();
|
||||
}
|
||||
@@ -310,13 +453,14 @@ Handler createLiveStreamHandler(
|
||||
);
|
||||
}
|
||||
|
||||
final session = await sessionManager.getOrStart(
|
||||
final result = await _runWithGpuFallback(
|
||||
id: sessionId,
|
||||
isLive: true,
|
||||
ffmpegPath: _getFFmpegPath(),
|
||||
argsBuilder: (dir) => [
|
||||
upstream: playlist,
|
||||
wantGpu: useNvidiaGpu && quality != 'source',
|
||||
buildArgs: (gpu) => [
|
||||
'-hide_banner', '-loglevel', 'warning',
|
||||
if (useNvidiaGpu && quality != 'source') ...['-hwaccel', 'cuda'],
|
||||
if (gpu) ...['-hwaccel', 'cuda'],
|
||||
'-headers', 'User-Agent: VLC/3.0.18 LibVLC/3.0.18\r\n',
|
||||
'-reconnect', '1', '-reconnect_streamed', '1',
|
||||
'-reconnect_at_eof', '1',
|
||||
@@ -330,7 +474,7 @@ Handler createLiveStreamHandler(
|
||||
'-probesize', '1000000',
|
||||
'-analyzeduration', '1000000',
|
||||
'-i', targetUrl,
|
||||
..._liveVideoArgs(quality, useNvidiaGpu),
|
||||
..._liveVideoArgs(quality, gpu),
|
||||
..._audioArgs(withFilters: false),
|
||||
// HLS sliding window: 10 x 2s segments (lower live latency than the
|
||||
// previous 20-segment window while still safe for iOS).
|
||||
@@ -347,8 +491,8 @@ Handler createLiveStreamHandler(
|
||||
],
|
||||
);
|
||||
|
||||
final result = await sessionManager.waitForPlaylist(session);
|
||||
if (!result.ready) {
|
||||
final session = result.session;
|
||||
if (!result.ready || session == null) {
|
||||
sessionManager.killSession(sessionId);
|
||||
return Response(502, body: 'Live transcoder failed: ${result.error}');
|
||||
}
|
||||
@@ -400,26 +544,48 @@ Handler createLiveStreamHandler(
|
||||
'[Live Turbo] $streamId: ${LogRedactor.redactUrl(targetUrl)}',
|
||||
);
|
||||
|
||||
final process = await Process.start(_getFFmpegPath(), [
|
||||
'-hide_banner', '-loglevel', 'warning',
|
||||
'-headers', 'User-Agent: VLC/3.0.18 LibVLC/3.0.18\r\n',
|
||||
'-reconnect', '1', '-reconnect_streamed', '1',
|
||||
'-reconnect_at_eof', '1',
|
||||
'-reconnect_delay_max', '10',
|
||||
'-rw_timeout', '30000000',
|
||||
// Démarrage rapide : ne pas bufferiser l'analyse, sonde réduite.
|
||||
'-fflags', 'nobuffer',
|
||||
'-flags', 'low_delay',
|
||||
'-probesize', '1000000',
|
||||
'-analyzeduration', '1000000',
|
||||
'-i', targetUrl,
|
||||
'-c:v', 'copy',
|
||||
'-c:a', 'aac', '-b:a', '160k', '-ac', '2', '-ar', '48000',
|
||||
'-af', 'aresample=async=1',
|
||||
// Pas de délai de mux : les paquets partent dès qu'ils existent.
|
||||
'-muxdelay', '0', '-muxpreload', '0',
|
||||
'-f', 'mpegts', 'pipe:1',
|
||||
]);
|
||||
// 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(), [
|
||||
'-hide_banner', '-loglevel', 'warning',
|
||||
'-headers', 'User-Agent: VLC/3.0.18 LibVLC/3.0.18\r\n',
|
||||
'-reconnect', '1', '-reconnect_streamed', '1',
|
||||
'-reconnect_at_eof', '1',
|
||||
'-reconnect_delay_max', '10',
|
||||
'-rw_timeout', '30000000',
|
||||
// Démarrage rapide : ne pas bufferiser l'analyse, sonde réduite.
|
||||
'-fflags', 'nobuffer',
|
||||
'-flags', 'low_delay',
|
||||
'-probesize', '1000000',
|
||||
'-analyzeduration', '1000000',
|
||||
'-i', targetUrl,
|
||||
'-c:v', 'copy',
|
||||
'-c:a', 'aac', '-b:a', '160k', '-ac', '2', '-ar', '48000',
|
||||
'-af', 'aresample=async=1',
|
||||
// Pas de délai de mux : les paquets partent dès qu'ils existent.
|
||||
'-muxdelay', '0', '-muxpreload', '0',
|
||||
'-f', 'mpegts', 'pipe:1',
|
||||
]);
|
||||
} 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.
|
||||
// `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');
|
||||
}
|
||||
|
||||
_registerUpstream(
|
||||
playlist,
|
||||
upstreamId,
|
||||
process,
|
||||
release: () => process.kill(ProcessSignal.sigterm),
|
||||
);
|
||||
|
||||
// Journaliser les erreurs FFmpeg (redactées) sans bloquer le flux.
|
||||
process.stderr.transform(const SystemEncoding().decoder).listen((line) {
|
||||
@@ -524,8 +690,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';
|
||||
|
||||
@@ -535,13 +704,14 @@ Handler createVodStreamHandler(
|
||||
);
|
||||
}
|
||||
|
||||
final session = await sessionManager.getOrStart(
|
||||
final result = await _runWithGpuFallback(
|
||||
id: sessionId,
|
||||
isLive: false,
|
||||
ffmpegPath: _getFFmpegPath(),
|
||||
argsBuilder: (dir) => [
|
||||
upstream: playlist,
|
||||
wantGpu: useNvidiaGpu && quality != 'source',
|
||||
buildArgs: (gpu) => [
|
||||
'-hide_banner', '-loglevel', 'warning',
|
||||
if (useNvidiaGpu && quality != 'source') ...['-hwaccel', 'cuda'],
|
||||
if (gpu) ...['-hwaccel', 'cuda'],
|
||||
'-headers', 'User-Agent: VLC/3.0.18 LibVLC/3.0.18\r\n',
|
||||
'-reconnect', '1',
|
||||
'-reconnect_at_eof', '1',
|
||||
@@ -555,7 +725,7 @@ Handler createVodStreamHandler(
|
||||
'-analyzeduration', '2000000',
|
||||
'-probesize', '5000000',
|
||||
'-i', targetUrl,
|
||||
..._vodVideoArgs(quality, useNvidiaGpu),
|
||||
..._vodVideoArgs(quality, gpu),
|
||||
..._audioArgs(withFilters: quality != 'source'),
|
||||
'-f', 'hls',
|
||||
'-hls_time', '4',
|
||||
@@ -570,8 +740,8 @@ Handler createVodStreamHandler(
|
||||
],
|
||||
);
|
||||
|
||||
final result = await sessionManager.waitForPlaylist(session);
|
||||
if (!result.ready) {
|
||||
final session = result.session;
|
||||
if (!result.ready || session == null) {
|
||||
sessionManager.killSession(sessionId);
|
||||
return Response(502, body: 'VOD transcoder failed: ${result.error}');
|
||||
}
|
||||
@@ -594,10 +764,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}
|
||||
@@ -677,9 +852,8 @@ Handler createRecordingStreamHandler(
|
||||
final videoCodec = await MediaProbe.videoCodec(targetUrl);
|
||||
final canCopyVideo = videoCodec == 'h264';
|
||||
|
||||
List<String> buildArgs({required bool copyVideo}) {
|
||||
final useNvidiaGpu =
|
||||
!copyVideo && (isGpuEnabled?.call() ?? _isNvidiaGpuEnabled());
|
||||
List<String> buildArgs({required bool copyVideo, required bool gpu}) {
|
||||
final useNvidiaGpu = !copyVideo && gpu;
|
||||
return [
|
||||
'-hide_banner', '-loglevel', 'warning',
|
||||
if (useNvidiaGpu) ...['-hwaccel', 'cuda'],
|
||||
@@ -724,9 +898,7 @@ Handler createRecordingStreamHandler(
|
||||
/// au démarrage, la playlist étant servie dès le premier segment. En
|
||||
/// copie vidéo elle est produite en une fraction de seconde, elle ne
|
||||
/// coûte donc que sur un ré-encodage.
|
||||
Future<({FfmpegSession session, bool ready, String? error})> run(
|
||||
bool copyVideo,
|
||||
) async {
|
||||
Future<_SessionAttempt> run(bool copyVideo) {
|
||||
if (!sessionManager.contains(sessionId)) {
|
||||
print(
|
||||
'[Recording] $sessionId : '
|
||||
@@ -734,18 +906,20 @@ Handler createRecordingStreamHandler(
|
||||
'(codec source : ${videoCodec ?? 'inconnu'})',
|
||||
);
|
||||
}
|
||||
final session = await sessionManager.getOrStart(
|
||||
if (copyVideo) {
|
||||
return _runSession(
|
||||
id: sessionId,
|
||||
isLive: false,
|
||||
argsBuilder: (_) => buildArgs(copyVideo: true, gpu: false),
|
||||
minSegments: 3,
|
||||
);
|
||||
}
|
||||
return _runWithGpuFallback(
|
||||
id: sessionId,
|
||||
isLive: false,
|
||||
ffmpegPath: _getFFmpegPath(),
|
||||
argsBuilder: (dir) => buildArgs(copyVideo: copyVideo),
|
||||
);
|
||||
final outcome =
|
||||
await sessionManager.waitForPlaylist(session, minSegments: 3);
|
||||
return (
|
||||
session: session,
|
||||
ready: outcome.ready,
|
||||
error: outcome.error,
|
||||
wantGpu: isGpuEnabled?.call() ?? _isNvidiaGpuEnabled(),
|
||||
buildArgs: (gpu) => buildArgs(copyVideo: false, gpu: gpu),
|
||||
minSegments: 3,
|
||||
);
|
||||
}
|
||||
|
||||
@@ -763,14 +937,13 @@ Handler createRecordingStreamHandler(
|
||||
attempt = await run(false);
|
||||
}
|
||||
|
||||
if (!attempt.ready) {
|
||||
final session = attempt.session;
|
||||
if (!attempt.ready || session == null) {
|
||||
sessionManager.killSession(sessionId);
|
||||
return Response(502,
|
||||
body: 'Recording transcoder failed: ${attempt.error}');
|
||||
}
|
||||
|
||||
final session = attempt.session;
|
||||
|
||||
session.touch();
|
||||
// Ménage des offsets abandonnés (aller-retours dans la barre de
|
||||
// progression), une fois la nouvelle session confirmée démarrée.
|
||||
|
||||
+25
-1
@@ -27,6 +27,7 @@ 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';
|
||||
|
||||
void main(List<String> args) async {
|
||||
// Parse command line arguments
|
||||
@@ -293,14 +294,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) {
|
||||
|
||||
@@ -0,0 +1,43 @@
|
||||
/// Repli GPU → CPU des transcodages.
|
||||
///
|
||||
/// Quand la VRAM est saturée (plusieurs sessions NVENC, un autre processus
|
||||
/// sur la carte) ou que le GPU dépasse son quota de sessions d'encodage,
|
||||
/// FFmpeg échoue à l'ouverture de l'encodeur et meurt en une seconde. Sans
|
||||
/// repli, le lecteur recevait un 502, relançait… et retombait sur le même
|
||||
/// GPU saturé : le flux restait illisible tant que la carte était pleine.
|
||||
|
||||
/// Signatures FFmpeg d'un échec imputable au GPU NVIDIA (NVENC/NVDEC/CUDA).
|
||||
///
|
||||
/// N'est consulté que pour une session lancée sur le GPU : « opening
|
||||
/// encoder » y désigne donc forcément h264_nvenc.
|
||||
final _gpuFailurePattern = RegExp(
|
||||
r'nvenc|cuda|cuvid|nvdec|OpenEncodeSession|No capable devices|'
|
||||
r'out of memory|Device creation failed|hwaccel|opening encoder',
|
||||
caseSensitive: false,
|
||||
);
|
||||
|
||||
/// Vrai si [error] (stderr récent de FFmpeg) trahit une panne GPU, et non
|
||||
/// une source injoignable ou un délai dépassé — qu'un repli CPU ne
|
||||
/// réglerait pas, et qui doublerait inutilement la charge d'un serveur déjà
|
||||
/// occupé.
|
||||
bool isGpuFailure(String? error) =>
|
||||
error != null && _gpuFailurePattern.hasMatch(error);
|
||||
|
||||
/// Mémorise une panne GPU récente pour que les sessions suivantes partent
|
||||
/// directement sur le CPU, au lieu de payer chacune un échec NVENC avant
|
||||
/// leur repli.
|
||||
class GpuHealth {
|
||||
final Duration cooldown;
|
||||
final DateTime Function() _now;
|
||||
DateTime? _unavailableUntil;
|
||||
|
||||
GpuHealth({
|
||||
this.cooldown = const Duration(minutes: 2),
|
||||
DateTime Function()? now,
|
||||
}) : _now = now ?? DateTime.now;
|
||||
|
||||
bool get available =>
|
||||
_unavailableUntil == null || !_now().isBefore(_unavailableUntil!);
|
||||
|
||||
void markFailed() => _unavailableUntil = _now().add(cooldown);
|
||||
}
|
||||
@@ -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,151 @@
|
||||
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 : le dernier flux demandé gagne. Avant d'ouvrir une nouvelle
|
||||
/// connexion amont, on coupe les plus anciennes du même compte jusqu'à lui
|
||||
/// faire de la place. 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.
|
||||
int register({
|
||||
required String id,
|
||||
required String account,
|
||||
required void Function() release,
|
||||
}) {
|
||||
final token = ++_nextToken;
|
||||
_slots[id] = _Slot(token, account, _now(), release);
|
||||
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()
|
||||
..sort((a, b) => a.value.since.compareTo(b.value.since));
|
||||
|
||||
final excess = others.length - (max - 1);
|
||||
if (excess <= 0) return const [];
|
||||
|
||||
final freed = <String>[];
|
||||
for (final entry in others.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;
|
||||
_Slot(this.token, this.account, this.since, this.release);
|
||||
}
|
||||
|
||||
/// 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.
|
||||
class AccountLimits {
|
||||
final Duration ttl;
|
||||
final Future<int?> Function(PlaylistConfig) _fetch;
|
||||
final Map<String, ({int max, DateTime at})> _cache = {};
|
||||
|
||||
/// Valeur retenue quand le panneau ne répond pas : 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) {
|
||||
return cached.max;
|
||||
}
|
||||
int? value;
|
||||
try {
|
||||
value = await _fetch(p);
|
||||
} catch (_) {
|
||||
value = null;
|
||||
}
|
||||
// Un échec n'est pas mis en cache longtemps : on réessaiera au prochain
|
||||
// flux plutôt que de rester 10 min sur une valeur devinée.
|
||||
if (value != null) _cache[key] = (max: value, at: DateTime.now());
|
||||
return value ?? fallback;
|
||||
}
|
||||
|
||||
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
|
||||
@@ -114,9 +118,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 +132,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,26 +154,60 @@ class XmltvEpgService {
|
||||
print('[XmltvEpg] index prêt : ${merged.length} chaînes');
|
||||
}
|
||||
|
||||
Future<String> _download(String url) async {
|
||||
/// Télécharge [url] puis l'indexe dans un isolate dédié.
|
||||
///
|
||||
/// Le serveur n'a qu'une boucle d'événements : décompresser et parser un
|
||||
/// dump de plusieurs milliers de chaînes dessus la gelait pendant des
|
||||
/// secondes. Plus rien d'autre ne répondait — en particulier le relais
|
||||
/// `turbo.ts`, qui cessait d'envoyer des octets au lecteur : la lecture
|
||||
/// se coupait systématiquement au démarrage, au moment précis où l'écran
|
||||
/// des chaînes demandait le guide.
|
||||
Future<Map<String, List<XmltvProgramme>>> _downloadAndParse(
|
||||
String url,
|
||||
) async {
|
||||
final response =
|
||||
await _client.get(Uri.parse(url)).timeout(downloadTimeout);
|
||||
if (response.statusCode != 200) {
|
||||
throw HttpException('HTTP ${response.statusCode}');
|
||||
}
|
||||
|
||||
List<int> bytes = response.bodyBytes;
|
||||
// Variables locales : la fermeture envoyée à l'isolate ne doit pas
|
||||
// capturer `this` (le client HTTP n'est pas transférable).
|
||||
final bytes = response.bodyBytes;
|
||||
final retention = this.retention;
|
||||
final horizon = this.horizon;
|
||||
return Isolate.run(
|
||||
() => _decodeAndParse(bytes, retention: retention, horizon: horizon),
|
||||
);
|
||||
}
|
||||
|
||||
static Map<String, List<XmltvProgramme>> _decodeAndParse(
|
||||
Uint8List raw, {
|
||||
required Duration retention,
|
||||
required Duration horizon,
|
||||
}) {
|
||||
List<int> bytes = raw;
|
||||
// Beaucoup de miroirs servent du .gz sans en-tête Content-Encoding : on
|
||||
// regarde le nombre magique plutôt que de se fier aux en-têtes.
|
||||
if (bytes.length > 2 && bytes[0] == 0x1f && bytes[1] == 0x8b) {
|
||||
bytes = gzip.decode(bytes);
|
||||
}
|
||||
return utf8.decode(bytes, allowMalformed: true);
|
||||
return _parse(
|
||||
utf8.decode(bytes, allowMalformed: true),
|
||||
retention: retention,
|
||||
horizon: horizon,
|
||||
);
|
||||
}
|
||||
|
||||
/// Exposé pour les tests.
|
||||
Map<String, List<XmltvProgramme>> parseForTest(String xml) => _parse(xml);
|
||||
Map<String, List<XmltvProgramme>> parseForTest(String xml) =>
|
||||
_parse(xml, retention: retention, horizon: horizon);
|
||||
|
||||
Map<String, List<XmltvProgramme>> _parse(String xml) {
|
||||
static Map<String, List<XmltvProgramme>> _parse(
|
||||
String xml, {
|
||||
required Duration retention,
|
||||
required Duration horizon,
|
||||
}) {
|
||||
final cutoff = DateTime.now().toUtc().subtract(retention);
|
||||
final limit = DateTime.now().toUtc().add(horizon);
|
||||
final byChannel = <String, List<XmltvProgramme>>{};
|
||||
|
||||
@@ -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));
|
||||
});
|
||||
});
|
||||
}
|
||||
@@ -0,0 +1,76 @@
|
||||
import 'package:test/test.dart';
|
||||
import '../services/gpu_fallback.dart';
|
||||
|
||||
void main() {
|
||||
group('isGpuFailure', () {
|
||||
test('detects NVENC out-of-memory', () {
|
||||
expect(
|
||||
isGpuFailure(
|
||||
'FFmpeg exited (1): [h264_nvenc @ 0x55] OpenEncodeSessionEx failed: '
|
||||
'out of memory (10)',
|
||||
),
|
||||
isTrue,
|
||||
);
|
||||
});
|
||||
|
||||
test('detects missing NVENC device', () {
|
||||
expect(
|
||||
isGpuFailure('[h264_nvenc @ 0x1] No capable devices found'),
|
||||
isTrue,
|
||||
);
|
||||
});
|
||||
|
||||
test('detects CUDA errors', () {
|
||||
expect(
|
||||
isGpuFailure('cu->cuInit(0) failed -> CUDA_ERROR_OUT_OF_MEMORY'),
|
||||
isTrue,
|
||||
);
|
||||
});
|
||||
|
||||
test('detects encoder opening failure', () {
|
||||
expect(
|
||||
isGpuFailure(
|
||||
'Error while opening encoder for output stream #0:0 - maybe '
|
||||
'incorrect parameters such as bit_rate, rate, width or height',
|
||||
),
|
||||
isTrue,
|
||||
);
|
||||
});
|
||||
|
||||
test('ignores upstream failures', () {
|
||||
expect(
|
||||
isGpuFailure(
|
||||
'FFmpeg exited (1): http://x/live/u/p/1.ts: Server returned 404 '
|
||||
'Not Found',
|
||||
),
|
||||
isFalse,
|
||||
);
|
||||
});
|
||||
|
||||
test('ignores timeouts and null', () {
|
||||
expect(isGpuFailure('Timeout waiting for transcoder'), isFalse);
|
||||
expect(isGpuFailure(null), isFalse);
|
||||
});
|
||||
});
|
||||
|
||||
group('GpuHealth', () {
|
||||
test('available until a failure, then again after cooldown', () {
|
||||
var now = DateTime.utc(2026, 10, 8, 20, 0);
|
||||
final health = GpuHealth(
|
||||
cooldown: const Duration(minutes: 2),
|
||||
now: () => now,
|
||||
);
|
||||
|
||||
expect(health.available, isTrue);
|
||||
|
||||
health.markFailed();
|
||||
expect(health.available, isFalse);
|
||||
|
||||
now = now.add(const Duration(minutes: 1, seconds: 59));
|
||||
expect(health.available, isFalse);
|
||||
|
||||
now = now.add(const Duration(seconds: 1));
|
||||
expect(health.available, isTrue);
|
||||
});
|
||||
});
|
||||
}
|
||||
@@ -0,0 +1,92 @@
|
||||
import 'package:test/test.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('quota inconnu ou nul : rien n\'est coupé', () {
|
||||
add('a', 'acc');
|
||||
expect(slots.makeRoom('acc', max: 0, keep: 'new'), isEmpty);
|
||||
});
|
||||
});
|
||||
|
||||
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,5 @@
|
||||
import 'dart:async';
|
||||
|
||||
import 'package:http/http.dart' as http;
|
||||
import 'package:http/testing.dart';
|
||||
import 'package:test/test.dart';
|
||||
@@ -148,6 +150,70 @@ void main() {
|
||||
});
|
||||
});
|
||||
|
||||
group('réactivité du serveur', () {
|
||||
test('l’indexation ne gèle pas la boucle d’événements', () async {
|
||||
// Le serveur n'a qu'un isolate : décoder et parser un dump de
|
||||
// plusieurs milliers de chaînes dessus bloquait tout le reste pendant
|
||||
// des secondes — y compris le relais `turbo.ts` des flux en cours,
|
||||
// d'où des coupures systématiques au démarrage de la lecture.
|
||||
final now = DateTime.now().toUtc();
|
||||
String stamp(DateTime d) {
|
||||
String two(int v) => v.toString().padLeft(2, '0');
|
||||
return '${d.year}${two(d.month)}${two(d.day)}'
|
||||
'${two(d.hour)}${two(d.minute)}${two(d.second)} +0000';
|
||||
}
|
||||
|
||||
// Volume d'un vrai dump de panneau : ~7000 chaînes, résumés de
|
||||
// plusieurs phrases.
|
||||
final desc = 'Résumé du programme, avec assez de texte pour peser '
|
||||
'comme un vrai guide. ' * 4;
|
||||
final xml = StringBuffer('<?xml version="1.0"?><tv>');
|
||||
for (var c = 0; c < 7000; c++) {
|
||||
xml.write('<channel id="Chaine.$c.fr">'
|
||||
'<display-name>FR - CHAINE $c</display-name></channel>');
|
||||
for (var p = 0; p < 30; p++) {
|
||||
final start = now.add(Duration(minutes: 30 * p));
|
||||
final stop = start.add(const Duration(minutes: 30));
|
||||
xml.write('<programme start="${stamp(start)}" '
|
||||
'stop="${stamp(stop)}" channel="Chaine.$c.fr">'
|
||||
'<title>Programme $p</title><desc>$desc</desc>'
|
||||
'</programme>');
|
||||
}
|
||||
}
|
||||
xml.write('</tv>');
|
||||
final body = xml.toString();
|
||||
|
||||
final service = XmltvEpgService(
|
||||
sourceUrls: const ['http://dump.example/epg.xml'],
|
||||
client: MockClient((_) async => http.Response(body, 200)),
|
||||
);
|
||||
|
||||
// Un « battement » toutes les 10 ms : le plus long écart observé est
|
||||
// le temps pendant lequel la boucle d'événements est restée bloquée.
|
||||
var last = DateTime.now();
|
||||
var worstGap = Duration.zero;
|
||||
final heartbeat = Timer.periodic(const Duration(milliseconds: 10), (_) {
|
||||
final now = DateTime.now();
|
||||
final gap = now.difference(last);
|
||||
if (gap > worstGap) worstGap = gap;
|
||||
last = now;
|
||||
});
|
||||
|
||||
await service.ensureFresh();
|
||||
heartbeat.cancel();
|
||||
// Le dernier blocage n'est suivi d'aucun battement : le compter aussi.
|
||||
final tail = DateTime.now().difference(last);
|
||||
if (tail > worstGap) worstGap = tail;
|
||||
|
||||
expect(service.channelCount, greaterThanOrEqualTo(7000));
|
||||
expect(
|
||||
worstGap,
|
||||
lessThan(const Duration(milliseconds: 250)),
|
||||
reason: 'boucle d’événements gelée ${worstGap.inMilliseconds} ms',
|
||||
);
|
||||
});
|
||||
});
|
||||
|
||||
group('horizon', () {
|
||||
test('écarte les programmes au-delà de l’horizon d’indexation', () {
|
||||
// Un dump national couvre sept jours : tout garder ferait grossir
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
@@ -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
|
||||
|
||||
+12
-12
@@ -149,10 +149,10 @@ packages:
|
||||
dependency: transitive
|
||||
description:
|
||||
name: characters
|
||||
sha256: f71061c654a3380576a52b451dd5532377954cf9dbd272a78fc8479606670803
|
||||
sha256: faf38497bda5ead2a8c7615f4f7939df04333478bf32e4173fcb06d428b5716b
|
||||
url: "https://pub.dev"
|
||||
source: hosted
|
||||
version: "1.4.0"
|
||||
version: "1.4.1"
|
||||
checked_yaml:
|
||||
dependency: transitive
|
||||
description:
|
||||
@@ -588,26 +588,26 @@ packages:
|
||||
dependency: transitive
|
||||
description:
|
||||
name: matcher
|
||||
sha256: dc58c723c3c24bf8d3e2d3ad3f2f9d7bd9cf43ec6feaa64181775e60190153f2
|
||||
sha256: "31bd099b47c10cd1aeb55146a2d46ce0277630ecef3f7dae54ad7873f36696cd"
|
||||
url: "https://pub.dev"
|
||||
source: hosted
|
||||
version: "0.12.17"
|
||||
version: "0.12.20"
|
||||
material_color_utilities:
|
||||
dependency: transitive
|
||||
description:
|
||||
name: material_color_utilities
|
||||
sha256: f7142bb1154231d7ea5f96bc7bde4bda2a0945d2806bb11670e30b850d56bdec
|
||||
sha256: "9c337007e82b1889149c82ed242ed1cb24a66044e30979c44912381e9be4c48b"
|
||||
url: "https://pub.dev"
|
||||
source: hosted
|
||||
version: "0.11.1"
|
||||
version: "0.13.0"
|
||||
meta:
|
||||
dependency: transitive
|
||||
description:
|
||||
name: meta
|
||||
sha256: "23f08335362185a5ea2ad3a4e597f1375e78bce8a040df5c600c8d3552ef2394"
|
||||
sha256: "307249ce4ff29d58a18e97f6345f539382eb9c9c29ecda628900f31de0443dd9"
|
||||
url: "https://pub.dev"
|
||||
source: hosted
|
||||
version: "1.17.0"
|
||||
version: "1.19.0"
|
||||
mime:
|
||||
dependency: transitive
|
||||
description:
|
||||
@@ -1073,10 +1073,10 @@ packages:
|
||||
dependency: transitive
|
||||
description:
|
||||
name: test_api
|
||||
sha256: ab2726c1a94d3176a45960b6234466ec367179b87dd74f1611adb1f3b5fb9d55
|
||||
sha256: "2a122cbe059f8b610d3a5415f42e255b6c17b1f21eee1d960f31080237fb4f11"
|
||||
url: "https://pub.dev"
|
||||
source: hosted
|
||||
version: "0.7.7"
|
||||
version: "0.7.12"
|
||||
timing:
|
||||
dependency: transitive
|
||||
description:
|
||||
@@ -1113,10 +1113,10 @@ packages:
|
||||
dependency: transitive
|
||||
description:
|
||||
name: vector_math
|
||||
sha256: d530bd74fea330e6e364cda7a85019c434070188383e1cd8d9777ee586914c5b
|
||||
sha256: "92b9910f66ed1057fd4da7b040ae7c74cafacf885bdc81be496928d5049b032d"
|
||||
url: "https://pub.dev"
|
||||
source: hosted
|
||||
version: "2.2.0"
|
||||
version: "2.4.3"
|
||||
video_player:
|
||||
dependency: "direct main"
|
||||
description:
|
||||
|
||||
@@ -0,0 +1,110 @@
|
||||
// Tests du moteur de lecture web (web/xf-player-core.js).
|
||||
//
|
||||
// Lancement : node --test "test/web/*_test.mjs"
|
||||
//
|
||||
// Le moteur est un script navigateur (IIFE sur `window`) : on l'exécute dans
|
||||
// un contexte vm avec un DOM minimal et un faux mpegts.js qui capture la
|
||||
// configuration reçue.
|
||||
import { test } from 'node:test';
|
||||
import assert from 'node:assert/strict';
|
||||
import { readFileSync } from 'node:fs';
|
||||
import vm from 'node:vm';
|
||||
|
||||
const source = readFileSync(
|
||||
new URL('../../web/xf-player-core.js', import.meta.url),
|
||||
'utf8',
|
||||
);
|
||||
|
||||
function fakeEventTarget() {
|
||||
return { addEventListener() {}, removeEventListener() {} };
|
||||
}
|
||||
|
||||
/** Charge le moteur et renvoie la config passée à mpegts.createPlayer. */
|
||||
function liveMpegtsConfig({ profile } = {}) {
|
||||
const created = [];
|
||||
const video = {
|
||||
...fakeEventTarget(),
|
||||
canPlayType: () => '',
|
||||
play: () => Promise.resolve(),
|
||||
buffered: { length: 0 },
|
||||
seekable: { length: 0 },
|
||||
};
|
||||
const window = {
|
||||
...fakeEventTarget(),
|
||||
console: { log() {} },
|
||||
location: { origin: 'https://xf.test', href: 'https://xf.test/', search: '' },
|
||||
parent: { postMessage() {} },
|
||||
mpegts: {
|
||||
isSupported: () => true,
|
||||
Events: { MEDIA_INFO: 'media_info', ERROR: 'error' },
|
||||
createPlayer(mediaDataSource, config) {
|
||||
created.push({ mediaDataSource, config });
|
||||
return { on() {}, attachMediaElement() {}, load() {} };
|
||||
},
|
||||
},
|
||||
};
|
||||
const context = vm.createContext({
|
||||
window,
|
||||
location: window.location,
|
||||
document: fakeEventTarget(),
|
||||
URLSearchParams,
|
||||
setInterval: () => 0,
|
||||
clearInterval() {},
|
||||
setTimeout,
|
||||
});
|
||||
vm.runInContext(source, context);
|
||||
|
||||
const XFPlayer = window.XFPlayer;
|
||||
const player = new XFPlayer({
|
||||
video,
|
||||
url: '/api/live/1/turbo.ts',
|
||||
type: 'live',
|
||||
});
|
||||
if (profile) player.profile = window.XFPlayerProfiles[profile];
|
||||
player._createMpegts(true);
|
||||
assert.equal(created.length, 1);
|
||||
return created[0].config;
|
||||
}
|
||||
|
||||
// Les panneaux Xtream relaient souvent une source HLS : le flux arrive par
|
||||
// rafales d'environ 6 s séparées de silences complets (mesuré en prod).
|
||||
const UPSTREAM_BURST_GAP_S = 7;
|
||||
|
||||
for (const profile of ['fast', 'balanced', 'safe']) {
|
||||
test(`live (${profile}) : la marge gardée après rattrapage couvre un silence amont`, () => {
|
||||
const c = liveMpegtsConfig({ profile });
|
||||
if (!c.liveBufferLatencyChasing) return; // pas de saut = rien à vérifier
|
||||
assert.ok(
|
||||
c.liveBufferLatencyMinRemain >= UPSTREAM_BURST_GAP_S,
|
||||
`MinRemain ${c.liveBufferLatencyMinRemain}s < silence amont ${UPSTREAM_BURST_GAP_S}s : ` +
|
||||
'chaque rafale déclenche un saut puis une coupure',
|
||||
);
|
||||
});
|
||||
|
||||
test(`live (${profile}) : une rafale normale ne déclenche pas de saut`, () => {
|
||||
const c = liveMpegtsConfig({ profile });
|
||||
if (!c.liveBufferLatencyChasing) return;
|
||||
// Après une rafale, le buffer vaut la marge + une rafale entière.
|
||||
assert.ok(
|
||||
c.liveBufferLatencyMaxLatency >=
|
||||
c.liveBufferLatencyMinRemain + 2 * UPSTREAM_BURST_GAP_S,
|
||||
'seuil de rattrapage trop proche de la marge : sauts à chaque rafale',
|
||||
);
|
||||
});
|
||||
}
|
||||
|
||||
test('live : le retard est résorbé en douceur (liveSync) plutôt que par sauts', () => {
|
||||
const c = liveMpegtsConfig();
|
||||
assert.equal(c.liveSync, true);
|
||||
assert.ok(c.liveSyncPlaybackRate > 1 && c.liveSyncPlaybackRate <= 1.2);
|
||||
assert.ok(c.liveSyncTargetLatency < c.liveSyncMaxLatency);
|
||||
assert.ok(c.liveSyncMaxLatency <= c.liveBufferLatencyMaxLatency);
|
||||
});
|
||||
|
||||
test('live : le buffer déjà lu est purgé (sinon QuotaExceeded en séance longue)', () => {
|
||||
const c = liveMpegtsConfig();
|
||||
assert.equal(c.autoCleanupSourceBuffer, true);
|
||||
assert.ok(
|
||||
c.autoCleanupMinBackwardDuration < c.autoCleanupMaxBackwardDuration,
|
||||
);
|
||||
});
|
||||
+44
-21
@@ -90,25 +90,49 @@
|
||||
playOverlay.style.display = 'none';
|
||||
}
|
||||
|
||||
var player = new XFPlayer({
|
||||
video: video,
|
||||
logPrefix: '[LitePlayer]',
|
||||
onLoading: showLoading,
|
||||
onReady: function () {
|
||||
hideLoading();
|
||||
playOverlay.style.display = 'none';
|
||||
},
|
||||
onError: showError,
|
||||
onBlocked: function () {
|
||||
hideLoading();
|
||||
playOverlay.style.display = 'flex';
|
||||
},
|
||||
onProfileChange: function (name) {
|
||||
loadingHint.textContent =
|
||||
name === 'safe' ? 'Réseau instable — buffer maximal'
|
||||
: 'Réseau instable — buffer élargi';
|
||||
}
|
||||
});
|
||||
// Un serveur chargé peut rater le démarrage du transcodeur : on
|
||||
// relance un lecteur neuf avant d'afficher l'erreur, comme le
|
||||
// font déjà player.html et player_mobile.html.
|
||||
var retryCount = 0;
|
||||
var maxRetries = 3;
|
||||
var player = null;
|
||||
|
||||
function buildPlayer() {
|
||||
if (player) { try { player.destroy(); } catch (e) {} }
|
||||
|
||||
player = new XFPlayer({
|
||||
video: video,
|
||||
logPrefix: '[LitePlayer]',
|
||||
onLoading: showLoading,
|
||||
onReady: function () {
|
||||
hideLoading();
|
||||
retryCount = 0;
|
||||
playOverlay.style.display = 'none';
|
||||
},
|
||||
onError: function (msg) {
|
||||
if (retryCount < maxRetries) {
|
||||
retryCount++;
|
||||
showLoading(
|
||||
'Reconnexion ' + retryCount + '/' + maxRetries
|
||||
);
|
||||
setTimeout(buildPlayer, 1500);
|
||||
} else {
|
||||
showError(msg);
|
||||
}
|
||||
},
|
||||
onBlocked: function () {
|
||||
hideLoading();
|
||||
playOverlay.style.display = 'flex';
|
||||
},
|
||||
onProfileChange: function (name) {
|
||||
loadingHint.textContent =
|
||||
name === 'safe' ? 'Réseau instable — buffer maximal'
|
||||
: 'Réseau instable — buffer élargi';
|
||||
}
|
||||
});
|
||||
window.xfPlayer = player;
|
||||
player.start();
|
||||
}
|
||||
|
||||
playOverlay.addEventListener('click', function () {
|
||||
video.play();
|
||||
@@ -123,8 +147,7 @@
|
||||
}
|
||||
});
|
||||
|
||||
player.start();
|
||||
window.xfPlayer = player;
|
||||
buildPlayer();
|
||||
})();
|
||||
</script>
|
||||
</body>
|
||||
|
||||
+71
-7
@@ -55,6 +55,9 @@
|
||||
|
||||
var ORDER = ['fast', 'balanced', 'safe'];
|
||||
|
||||
// Erreurs MPEG-TS tolérées sur 60 s avant de basculer sur HLS.
|
||||
var MPEGTS_MAX_ERRORS = 3;
|
||||
|
||||
// ------------------------------------------------------------------
|
||||
// Chargement paresseux des libs (aucune requête inutile)
|
||||
// ------------------------------------------------------------------
|
||||
@@ -178,6 +181,7 @@
|
||||
this._started = false;
|
||||
this._stalls = 0;
|
||||
this._stallTimes = [];
|
||||
this._mpegtsErrorTimes = [];
|
||||
this._destroyed = false;
|
||||
this._reportedTime = 0;
|
||||
this._hlsDuration = 0;
|
||||
@@ -317,10 +321,15 @@
|
||||
self.hls = hls;
|
||||
global.hlsInstance = hls; // compat : code existant qui inspecte l'instance
|
||||
|
||||
// Faux tant que la première playlist n'a pas été obtenue : voir le
|
||||
// traitement des erreurs réseau plus bas.
|
||||
var manifestParsed = false;
|
||||
|
||||
hls.loadSource(self.url);
|
||||
hls.attachMedia(self.video);
|
||||
|
||||
hls.on(Hls.Events.MANIFEST_PARSED, function () {
|
||||
manifestParsed = true;
|
||||
// Démarrage immédiat : on ne sonde plus le buffer toutes les 200 ms.
|
||||
// Si le réseau ne suit pas, l'escalade de profil s'en chargera.
|
||||
self._attemptPlay();
|
||||
@@ -337,6 +346,15 @@
|
||||
if (!d.fatal) return;
|
||||
self.log('Erreur HLS fatale : ' + 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.
|
||||
// `startLoad()` ne relance pas ce chargement-là : le lecteur
|
||||
// restait muet sur un spinner infini. On remonte l'erreur pour que
|
||||
// la page relance un lecteur neuf ou l'affiche.
|
||||
if (!manifestParsed) {
|
||||
self.onError('Flux indisponible : le serveur n\'a pas pu le préparer');
|
||||
return;
|
||||
}
|
||||
self._escalate('erreur réseau');
|
||||
hls.startLoad();
|
||||
} else if (d.type === Hls.ErrorTypes.MEDIA_ERROR) {
|
||||
@@ -401,13 +419,28 @@
|
||||
// 64 Ko en profil « fast » contre 512 Ko auparavant.
|
||||
stashInitialSize: this.profile.mpegtsStash,
|
||||
maxStashSize: 30 * 1024 * 1024,
|
||||
// Rattrapage du direct : sans lui, chaque micro-coupure fait dériver
|
||||
// la lecture derrière le direct (latence qui s'accumule). Activé
|
||||
// uniquement en profil « fast » — les profils élargis privilégient
|
||||
// la stabilité du buffer.
|
||||
liveBufferLatencyChasing: isLive && this.profile.name === 'fast',
|
||||
liveBufferLatencyMaxLatency: 5,
|
||||
liveBufferLatencyMinRemain: 1,
|
||||
// POURQUOI ces seuils : les panneaux Xtream relaient le plus souvent
|
||||
// une source HLS et livrent le flux par RAFALES (~6 s de contenu d'un
|
||||
// coup, puis plusieurs secondes de silence total — mesuré en prod).
|
||||
// L'ancien rattrapage (5 s max, 1 s de marge) sautait en avant à
|
||||
// chaque rafale puis tombait à sec pendant le silence suivant :
|
||||
// coupure toutes les 3-4 s et contenu sauté (12 sauts et 14 coupures
|
||||
// en 30 s sur une chaîne FHD). La marge de 12 s couvre deux silences ;
|
||||
// le saut ne sert plus qu'à borner un retard réellement accumulé.
|
||||
liveBufferLatencyChasing: isLive,
|
||||
liveBufferLatencyMaxLatency: 30,
|
||||
liveBufferLatencyMinRemain: 12,
|
||||
// Au-delà de 20 s de retard, accélérer légèrement (x1,1) jusqu'à
|
||||
// revenir à 12 s : rattrapage invisible, sans saut ni coupure.
|
||||
liveSync: isLive,
|
||||
liveSyncMaxLatency: 20,
|
||||
liveSyncTargetLatency: 12,
|
||||
liveSyncPlaybackRate: 1.1,
|
||||
// Purge du buffer déjà lu : sans elle, une longue séance finit par
|
||||
// saturer le SourceBuffer (QuotaExceededError → erreur MPEG-TS).
|
||||
autoCleanupSourceBuffer: true,
|
||||
autoCleanupMaxBackwardDuration: 60,
|
||||
autoCleanupMinBackwardDuration: 30,
|
||||
lazyLoad: false
|
||||
}
|
||||
);
|
||||
@@ -446,6 +479,37 @@
|
||||
player.on(mpegts.Events.ERROR, function (type, detail) {
|
||||
self.log('Erreur MPEG-TS : ' + type + ' / ' + detail);
|
||||
if (self._destroyed) return;
|
||||
|
||||
// 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
|
||||
// qu'aggraver la charge. Au-delà de quelques échecs rapprochés, on
|
||||
// passe par la route HLS, qui renvoie un vrai code d'erreur.
|
||||
var now = Date.now();
|
||||
self._mpegtsErrorTimes = self._mpegtsErrorTimes.filter(function (t) {
|
||||
return now - t < 60000;
|
||||
});
|
||||
self._mpegtsErrorTimes.push(now);
|
||||
if (self._mpegtsErrorTimes.length > MPEGTS_MAX_ERRORS) {
|
||||
try {
|
||||
player.unload();
|
||||
player.detachMediaElement();
|
||||
player.destroy();
|
||||
} catch (e) {}
|
||||
self.mpegts = null;
|
||||
|
||||
var hlsUrl = self._hlsEquivalent(self.url);
|
||||
if (!hlsUrl) {
|
||||
self.onError('Flux indisponible');
|
||||
return;
|
||||
}
|
||||
self.log('Erreurs MPEG-TS répétées — bascule HLS');
|
||||
self.onLoading('Ouverture du flux');
|
||||
self.url = hlsUrl;
|
||||
if (self._canPlayNativeHls()) self._playDirect(hlsUrl);
|
||||
else self._startHls();
|
||||
return;
|
||||
}
|
||||
|
||||
self._escalate('erreur MPEG-TS');
|
||||
// Recréation propre avec le nouveau profil.
|
||||
setTimeout(function () {
|
||||
|
||||
Reference in new issue
Block a user