mirror of
https://github.com/R0m1k3/xtremflow.git
synced 2026-10-11 17:30:00 +02:00
Zap : ouvrir les chaînes par le HLS du panneau, ~4 s gagnées
Mesuré en prod avec la sonde admin : le panneau met 4,2 à 4,4 s à répondre la redirection d'un .ts, contre 0,2 à 1,3 s pour un .m3u8 dont le premier segment arrive ensuite en 50 ms. turbo.ts et les sessions HLS live lisent désormais le .m3u8 du panneau (-live_start_index -2 : deux segments d'avance d'un coup), avec repli automatique sur le .ts si le HLS ne livre rien — dans la même réponse pour turbo.ts, le lecteur ne voit qu'un démarrage plus long. Pas de reconnect_at_eof en HLS : chaque segment finit par un EOF (leçon de xtremobile). Le navigateur reçoit toujours le même flux MPEG-TS. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
1 parent
a3640a6b5f
commit
0ea60ccbfa
2 files changed
+164
-80
No files matched your search
+129
-80
@@ -9,7 +9,6 @@ import '../services/gpu_fallback.dart';
|
|||||||
import '../services/upstream_slots.dart';
|
import '../services/upstream_slots.dart';
|
||||||
import '../utils/log_redactor.dart';
|
import '../utils/log_redactor.dart';
|
||||||
import '../utils/media_probe.dart';
|
import '../utils/media_probe.dart';
|
||||||
import '../utils/stream_pipe.dart';
|
|
||||||
import 'recording_playlist.dart';
|
import 'recording_playlist.dart';
|
||||||
|
|
||||||
/// Directory for temporary HLS segments
|
/// Directory for temporary HLS segments
|
||||||
@@ -441,6 +440,46 @@ 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(
|
Handler createLiveStreamHandler(
|
||||||
Future<PlaylistConfig?> Function(Request) getPlaylist, {
|
Future<PlaylistConfig?> Function(Request) getPlaylist, {
|
||||||
bool Function()? isGpuEnabled,
|
bool Function()? isGpuEnabled,
|
||||||
@@ -458,18 +497,14 @@ Handler createLiveStreamHandler(
|
|||||||
final playlist = await getPlaylist(request);
|
final playlist = await getPlaylist(request);
|
||||||
if (playlist == null) return Response.forbidden('No playlist');
|
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 useNvidiaGpu = isGpuEnabled?.call() ?? _isNvidiaGpuEnabled();
|
||||||
final sessionId = 'live_${streamId}_$quality';
|
final sessionId = 'live_${streamId}_$quality';
|
||||||
|
|
||||||
if (!sessionManager.contains(sessionId)) {
|
if (!sessionManager.contains(sessionId)) {
|
||||||
print(
|
print('[Live HLS] Starting $sessionId');
|
||||||
'[Live HLS] Starting $sessionId: ${LogRedactor.redactUrl(targetUrl)}',
|
|
||||||
);
|
|
||||||
}
|
}
|
||||||
|
|
||||||
final result = await _runWithGpuFallback(
|
Future<_SessionAttempt> run({required bool hls}) => _runWithGpuFallback(
|
||||||
id: sessionId,
|
id: sessionId,
|
||||||
isLive: true,
|
isLive: true,
|
||||||
upstream: playlist,
|
upstream: playlist,
|
||||||
@@ -477,19 +512,10 @@ Handler createLiveStreamHandler(
|
|||||||
buildArgs: (gpu) => [
|
buildArgs: (gpu) => [
|
||||||
'-hide_banner', '-loglevel', 'warning',
|
'-hide_banner', '-loglevel', 'warning',
|
||||||
if (gpu) ...['-hwaccel', 'cuda'],
|
if (gpu) ...['-hwaccel', 'cuda'],
|
||||||
'-headers', 'User-Agent: VLC/3.0.18 LibVLC/3.0.18\r\n',
|
...liveInputArgs(
|
||||||
'-reconnect', '1', '-reconnect_streamed', '1',
|
_liveSourceUrl(playlist, streamId, hls: hls),
|
||||||
'-reconnect_at_eof', '1',
|
hls: hls,
|
||||||
'-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),
|
..._liveVideoArgs(quality, gpu),
|
||||||
..._audioArgs(withFilters: false),
|
..._audioArgs(withFilters: false),
|
||||||
// HLS sliding window: 10 x 2s segments (lower live latency than the
|
// HLS sliding window: 10 x 2s segments (lower live latency than the
|
||||||
@@ -507,6 +533,15 @@ 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;
|
final session = result.session;
|
||||||
if (!result.ready || session == null) {
|
if (!result.ready || session == null) {
|
||||||
sessionManager.killSession(sessionId);
|
sessionManager.killSession(sessionId);
|
||||||
@@ -554,79 +589,93 @@ Handler createLiveStreamHandler(
|
|||||||
final playlist = await getPlaylist(request);
|
final playlist = await getPlaylist(request);
|
||||||
if (playlist == null) return Response.forbidden('No playlist');
|
if (playlist == null) return Response.forbidden('No playlist');
|
||||||
|
|
||||||
final targetUrl =
|
print('[Live Turbo] $streamId');
|
||||||
'${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
|
// Un identifiant par requête : deux onglets sur la même chaîne sont
|
||||||
// deux connexions distinctes chez le fournisseur.
|
// deux connexions distinctes chez le fournisseur.
|
||||||
final upstreamId = 'turbo_${streamId}_${++_turboCounter}';
|
final upstreamId = 'turbo_${streamId}_${++_turboCounter}';
|
||||||
await _claimUpstream(playlist, upstreamId);
|
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
|
// 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
|
// rafale ; un client parti (zap mal détecté derrière le proxy) n'en
|
||||||
// reçoit plus, la contre-pression bloquant FFmpeg.
|
// reçoit plus, la contre-pression bloquant FFmpeg.
|
||||||
var lastDelivered = DateTime.now();
|
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.
|
Future<Process?> start({required bool hls}) async {
|
||||||
process.stderr.transform(const SystemEncoding().decoder).listen((line) {
|
final Process process;
|
||||||
final trimmed = line.trim();
|
try {
|
||||||
if (trimmed.isNotEmpty) {
|
process = await Process.start(_getFFmpegPath(), [
|
||||||
print('[Live Turbo] $streamId ffmpeg: '
|
'-hide_banner', '-loglevel', 'warning',
|
||||||
'${LogRedactor.redactUrl(trimmed)}');
|
...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;
|
||||||
}
|
}
|
||||||
});
|
_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;
|
||||||
|
}
|
||||||
|
|
||||||
// Relayer stdout vers le client en gardant la contre-pression, et tuer
|
final first = await start(hls: true);
|
||||||
// FFmpeg dès que le client zappe ou ferme l'onglet (sinon les processus
|
// Serveur saturé (plus de processus ou de mémoire) : un 503 que le
|
||||||
// s'accumulent à chaque zap).
|
// lecteur sait traiter, plutôt qu'une exception qui remonte en 500.
|
||||||
final body = pipeWithBackpressure(
|
if (first == null) return Response(503, body: 'FFmpeg start failed');
|
||||||
process.stdout,
|
|
||||||
onStop: () => process.kill(ProcessSignal.sigterm),
|
// Relais vers le client. Un générateur `async*` garde la contre-pression
|
||||||
).map((chunk) {
|
// (`yield` attend tant que le client ne lit pas) et son `finally` tue
|
||||||
lastDelivered = DateTime.now();
|
// FFmpeg dès que le client zappe ou ferme l'onglet.
|
||||||
return chunk;
|
//
|
||||||
});
|
// 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();
|
||||||
|
|
||||||
return Response(
|
return Response(
|
||||||
200,
|
200,
|
||||||
|
|||||||
@@ -0,0 +1,35 @@
|
|||||||
|
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');
|
||||||
|
});
|
||||||
|
});
|
||||||
|
}
|
||||||
Reference in new issue
Block a user