Annule « Zap : ouvrir les chaînes par le HLS du panneau »

Mesure après déploiement : premier octet de turbo.ts à 4,5–4,8 s avec la
source HLS, contre 4,3–4,7 s en .ts — aucun gain. Les logs confirment
que le HLS était bien utilisé (aucun repli). La sonde qui annonçait
0,5 s ouvrait le .ts AVANT le .m3u8 : la chaîne était déjà démarrée chez
le panneau. Les ~4,5 s sont le temps de démarrage d'une chaîne chez le
fournisseur, quel que soit le format.

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

This reverts commit 0ea60cc.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
MichaelandClaude Opus 5.5 committed 2026-10-10 23:30:51 +02:00
1 parent 4d00409de8
commit 775caad0ee
2 files changed
+80 -164

No files matched your search

+80 -129
View File
@@ -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<List<int>> _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<String> 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<PlaylistConfig?> 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<Process?> 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<List<int>> 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,
-35
View File
@@ -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');
});
});
}