diff --git a/bin/api/streaming_handler.dart b/bin/api/streaming_handler.dart index aafa2f5..20b5536 100644 --- a/bin/api/streaming_handler.dart +++ b/bin/api/streaming_handler.dart @@ -9,6 +9,7 @@ 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'; import 'recording_playlist.dart'; /// Directory for temporary HLS segments @@ -440,46 +441,6 @@ Stream> _resilientLiveBody( } } -/// URL d'une chaîne chez le panneau, en HLS (`.m3u8`) ou en flux continu -/// (`.ts`). -String _liveSourceUrl(PlaylistConfig p, String streamId, {required bool hls}) => - '${p.dns}/live/${p.username}/${p.password}/$streamId.${hls ? 'm3u8' : 'ts'}'; - -/// Arguments d'entrée FFmpeg pour une chaîne du panneau. -/// -/// POURQUOI le HLS d'abord : mesuré en prod (sonde admin), le panneau met -/// 4,2 à 4,4 s à répondre la redirection d'un `.ts`, contre 0,2 à 1,3 s pour -/// un `.m3u8` dont le premier segment arrive ensuite en 50 ms — un zap ~4× -/// plus rapide. La playlist du panneau garde en plus ~60 s de segments : -/// `-live_start_index -2` en récupère deux d'un coup, ce qui donne au -/// lecteur une avance immédiate pour absorber les à-coups. -/// -/// Pas de `-reconnect_at_eof` en HLS : chaque segment se termine par une -/// fin de fichier, l'option provoquerait une reconnexion à chaque segment -/// (leçon tirée de xtremobile, où elle causait gels et sauts). -List liveInputArgs(String url, {required bool hls}) => [ - '-headers', 'User-Agent: VLC/3.0.18 LibVLC/3.0.18\r\n', - if (hls) ...[ - '-reconnect', '1', - '-reconnect_on_network_error', '1', - '-reconnect_delay_max', '5', - '-live_start_index', '-2', - ] else ...[ - '-reconnect', '1', '-reconnect_streamed', '1', - '-reconnect_at_eof', '1', - '-reconnect_delay_max', '10', - ], - // Abort reads stuck for 30s so a stalled upstream triggers the - // reconnect logic instead of wedging the process forever. - '-rw_timeout', '30000000', - // Démarrage rapide : sans borne, FFmpeg peut passer plusieurs - // secondes à sonder le flux avant d'écrire la première sortie. - '-fflags', 'nobuffer', - '-probesize', '1000000', - '-analyzeduration', '1000000', - '-i', url, - ]; - Handler createLiveStreamHandler( Future Function(Request) getPlaylist, { bool Function()? isGpuEnabled, @@ -497,14 +458,18 @@ Handler createLiveStreamHandler( final playlist = await getPlaylist(request); if (playlist == null) return Response.forbidden('No playlist'); + final targetUrl = + '${playlist.dns}/live/${playlist.username}/${playlist.password}/$streamId.ts'; final useNvidiaGpu = isGpuEnabled?.call() ?? _isNvidiaGpuEnabled(); final sessionId = 'live_${streamId}_$quality'; if (!sessionManager.contains(sessionId)) { - print('[Live HLS] Starting $sessionId'); + print( + '[Live HLS] Starting $sessionId: ${LogRedactor.redactUrl(targetUrl)}', + ); } - Future<_SessionAttempt> run({required bool hls}) => _runWithGpuFallback( + final result = await _runWithGpuFallback( id: sessionId, isLive: true, upstream: playlist, @@ -512,10 +477,19 @@ Handler createLiveStreamHandler( buildArgs: (gpu) => [ '-hide_banner', '-loglevel', 'warning', if (gpu) ...['-hwaccel', 'cuda'], - ...liveInputArgs( - _liveSourceUrl(playlist, streamId, hls: hls), - hls: hls, - ), + '-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', + // Abort reads stuck for 30s so a stalled upstream triggers the + // reconnect logic instead of wedging the process forever. + '-rw_timeout', '30000000', + // Démarrage rapide : sans borne, FFmpeg peut passer plusieurs + // secondes à sonder le flux avant d'écrire le premier segment. + '-fflags', 'nobuffer', + '-probesize', '1000000', + '-analyzeduration', '1000000', + '-i', targetUrl, ..._liveVideoArgs(quality, gpu), ..._audioArgs(withFilters: false), // HLS sliding window: 10 x 2s segments (lower live latency than the @@ -533,15 +507,6 @@ Handler createLiveStreamHandler( ], ); - var result = await run(hls: true); - if (!result.ready) { - // Panneau sans HLS pour cette chaîne (ou HLS en panne) : repli sur le - // flux continu, plus lent à ouvrir mais universel. - print('[Live HLS] $sessionId : HLS du panneau indisponible — repli .ts'); - sessionManager.killSession(sessionId); - result = await run(hls: false); - } - final session = result.session; if (!result.ready || session == null) { sessionManager.killSession(sessionId); @@ -589,93 +554,79 @@ Handler createLiveStreamHandler( final playlist = await getPlaylist(request); if (playlist == null) return Response.forbidden('No playlist'); - print('[Live Turbo] $streamId'); + final targetUrl = + '${playlist.dns}/live/${playlist.username}/${playlist.password}/$streamId.ts'; + print( + '[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(), [ + '-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'); + } + // 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, + ); - Future start({required bool hls}) async { - final Process process; - try { - process = await Process.start(_getFFmpegPath(), [ - '-hide_banner', '-loglevel', 'warning', - ...liveInputArgs( - _liveSourceUrl(playlist, streamId, hls: hls), - hls: hls, - ), - '-flags', 'low_delay', - '-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) { - // `ProcessException` liste les arguments, donc l'URL `-i` du panneau. - print('[Live Turbo] $streamId : FFmpeg n\'a pas démarré ' - '(${LogRedactor.redactUrl('$e')})'); - return null; + // Journaliser les erreurs FFmpeg (redactées) sans bloquer le flux. + process.stderr.transform(const SystemEncoding().decoder).listen((line) { + final trimmed = line.trim(); + if (trimmed.isNotEmpty) { + print('[Live Turbo] $streamId ffmpeg: ' + '${LogRedactor.redactUrl(trimmed)}'); } - _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(); - if (trimmed.isNotEmpty) { - print('[Live Turbo] $streamId ffmpeg: ' - '${LogRedactor.redactUrl(trimmed)}'); - } - }); - return process; - } + }); - final first = await start(hls: true); - // 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. - if (first == null) return Response(503, body: 'FFmpeg start failed'); - - // Relais vers le client. Un générateur `async*` garde la contre-pression - // (`yield` attend tant que le client ne lit pas) et son `finally` tue - // FFmpeg dès que le client zappe ou ferme l'onglet. - // - // Si le HLS du panneau n'a livré aucun octet (format refusé pour cette - // chaîne, playlist vide), on rebascule sur le `.ts` dans la MÊME - // réponse : le lecteur ne voit qu'un démarrage un peu plus long. - Stream> relay() async* { - Process? process = first; - var hls = true; - try { - while (process != null) { - var delivered = false; - await for (final chunk in process.stdout) { - delivered = true; - lastDelivered = DateTime.now(); - yield chunk; - } - if (delivered || !hls) break; - print('[Live Turbo] $streamId : HLS du panneau muet — repli .ts'); - hls = false; - process = await start(hls: false); - } - } finally { - process?.kill(ProcessSignal.sigterm); - } - } - - final body = relay(); + // Relayer stdout vers le client en gardant la contre-pression, et tuer + // FFmpeg dès que le client zappe ou ferme l'onglet (sinon les processus + // s'accumulent à chaque zap). + final body = pipeWithBackpressure( + process.stdout, + onStop: () => process.kill(ProcessSignal.sigterm), + ).map((chunk) { + lastDelivered = DateTime.now(); + return chunk; + }); return Response( 200, diff --git a/bin/test/live_input_args_test.dart b/bin/test/live_input_args_test.dart deleted file mode 100644 index 1f1506c..0000000 --- a/bin/test/live_input_args_test.dart +++ /dev/null @@ -1,35 +0,0 @@ -import 'package:test/test.dart'; -import '../api/streaming_handler.dart'; - -void main() { - group('liveInputArgs', () { - test('HLS : jamais de reconnect_at_eof (chaque segment finit par un EOF)', - () { - final args = liveInputArgs('http://p/live/u/p/1.m3u8', hls: true); - expect(args, isNot(contains('-reconnect_at_eof'))); - expect(args, isNot(contains('-reconnect_streamed'))); - }); - - test('HLS : démarre deux segments avant le direct pour avoir de l\'avance', - () { - final args = liveInputArgs('http://p/live/u/p/1.m3u8', hls: true); - final i = args.indexOf('-live_start_index'); - expect(i, isNot(-1)); - expect(args[i + 1], '-2'); - // Option de démuxeur : doit précéder -i. - expect(i, lessThan(args.indexOf('-i'))); - }); - - test('.ts : garde la reconnexion en fin de flux (coupures du panneau)', () { - final args = liveInputArgs('http://p/live/u/p/1.ts', hls: false); - expect(args, contains('-reconnect_at_eof')); - expect(args, isNot(contains('-live_start_index'))); - }); - - test('l\'URL est la dernière entrée, juste après -i', () { - final args = liveInputArgs('http://p/x.m3u8', hls: true); - expect(args.last, 'http://p/x.m3u8'); - expect(args[args.length - 2], '-i'); - }); - }); -}