Files
xtremflow/bin/api/streaming_handler.dart
T
MichaelandClaude Opus 5.5 0ea60ccbfa 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>
2026-10-10 23:08:08 +02:00

1065 lines
39 KiB
Dart
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
import 'dart:io';
import 'package:shelf/shelf.dart';
import 'package:shelf_router/shelf_router.dart';
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 'recording_playlist.dart';
/// Directory for temporary HLS segments
final Directory _hlsTempDir =
Directory('${Directory.systemTemp.path}/xtremflow_streams');
/// 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));
}
/// Délai sans consommation au-delà duquel un flux est tenu pour orphelin.
///
/// Un lecteur HLS redemande sa playlist toutes les 2 à 4 s ; un flux
/// `turbo.ts` reçoit au moins une rafale toutes les ~7 s même sur un panneau
/// qui livre par à-coups. 15 s laissent une marge à ces deux rythmes.
const _viewerGrace = Duration(seconds: 15);
/// Inscrit [process] comme connexion amont de [playlist] tant qu'il vit.
///
/// [lastActivity] : dernier signe de vie du spectateur ; la connexion n'est
/// coupée pour en ouvrir une autre que s'il remonte à plus de
/// [_viewerGrace].
void _registerUpstream(
PlaylistConfig playlist,
String id,
Process process, {
required void Function() release,
required DateTime Function() lastActivity,
}) {
final token = _upstreamSlots.register(
id: id,
account: accountKeyOf(playlist),
release: release,
isActive: () =>
DateTime.now().difference(lastActivity()) < _viewerGrace,
);
process.exitCode.then((_) => _upstreamSlots.unregister(id, token: token));
}
/// Helper to resolve FFmpeg path (SYSTEM PATH vs Portable)
String _getFFmpegPath() {
if (Platform.isWindows) {
if (File('ffmpeg.exe').existsSync()) return 'ffmpeg.exe';
if (File('bin/ffmpeg.exe').existsSync()) return 'bin/ffmpeg.exe';
if (File('ffmpeg/bin/ffmpeg.exe').existsSync()) {
return 'ffmpeg/bin/ffmpeg.exe';
}
}
return 'ffmpeg'; // Default to PATH
}
/// Check if NVIDIA GPU acceleration is enabled via environment variable
bool _isNvidiaGpuEnabled() {
final envValue = Platform.environment['NVIDIA_GPU'] ??
Platform.environment['nvidia_gpu'] ??
'false';
return envValue.toLowerCase() == 'true' || envValue == '1';
}
/// Initialize streaming subsystem
Future<void> initStreaming() async {
await sessionManager.init();
}
/// Supported quality presets. `source` skips video transcoding entirely
/// (`-c:v copy`) — a huge CPU win since most Xtream streams are already
/// H.264. Audio is always normalized to AAC for HLS compatibility.
const supportedQualities = {'source', 'high', 'medium', 'low'};
String _sanitizeQuality(String? quality) {
return supportedQualities.contains(quality) ? quality! : 'high';
}
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) {
case 'source':
return ['-c:v', 'copy'];
case 'medium':
return gpu
? [
'-c:v', 'h264_nvenc', '-preset', 'p4', '-tune', 'hq',
'-b:v', '3000k', '-maxrate', '4500k', '-bufsize', '6000k',
'-profile:v', 'high', '-level', '4.0', '-pix_fmt', 'yuv420p',
'-g', '50',
]
: [
'-c:v', 'libx264', '-preset', 'veryfast', '-tune', 'zerolatency',
'-profile:v', 'high', '-level', '4.0',
'-b:v', '3000k', '-maxrate', '4500k', '-bufsize', '6000k',
'-pix_fmt', 'yuv420p', '-g', '50',
];
case 'low':
return gpu
? [
'-c:v', 'h264_nvenc', '-preset', 'p4',
'-b:v', '1500k', '-maxrate', '2000k', '-bufsize', '3000k',
'-vf', 'scale=-2:720',
'-pix_fmt', 'yuv420p', '-g', '50',
]
: [
'-c:v', 'libx264', '-preset', 'veryfast', '-tune', 'zerolatency',
'-b:v', '1500k', '-maxrate', '2000k', '-bufsize', '3000k',
'-vf', 'scale=-2:720',
'-pix_fmt', 'yuv420p', '-g', '50',
];
case 'high':
default:
return gpu
? [
'-c:v', 'h264_nvenc', '-preset', 'p4', '-tune', 'hq',
'-b:v', '8000k', '-maxrate', '12000k', '-bufsize', '16000k',
'-profile:v', 'high', '-level', '4.0', '-pix_fmt', 'yuv420p',
'-g', '50',
]
: [
// veryfast au lieu de medium : en live, l'encodeur doit tenir
// le temps réel ET produire le premier segment vite ; medium
// ajoutait plusieurs secondes au zap sans gain visible à ce
// débit.
'-c:v', 'libx264', '-preset', 'veryfast', '-tune', 'zerolatency',
'-profile:v', 'high', '-level', '4.0',
'-b:v', '6000k', '-maxrate', '8000k', '-bufsize', '12000k',
'-pix_fmt', 'yuv420p', '-g', '50',
];
}
}
/// Video encoding args for the requested quality (VOD profile).
List<String> _vodVideoArgs(String quality, bool gpu) {
switch (quality) {
case 'source':
return ['-c:v', 'copy'];
case 'medium':
return gpu
? [
'-c:v', 'h264_nvenc', '-preset', 'p4', '-tune', 'hq',
'-rc', 'cbr', '-b:v', '2500k', '-maxrate', '3000k',
'-bufsize', '5000k', '-g', '48', '-bf', '2',
'-pix_fmt', 'yuv420p',
]
: [
'-c:v', 'libx264', '-preset', 'fast', '-crf', '23',
'-maxrate', '6000k', '-bufsize', '12000k',
'-pix_fmt', 'yuv420p', '-g', '48', '-threads', '0',
];
case 'low':
return gpu
? [
'-c:v', 'h264_nvenc', '-preset', 'p4',
'-rc', 'cbr', '-b:v', '1500k', '-maxrate', '2000k',
'-bufsize', '3000k', '-vf', 'scale=-2:720',
'-g', '48', '-pix_fmt', 'yuv420p',
]
: [
'-c:v', 'libx264', '-preset', 'fast', '-crf', '26',
'-maxrate', '3000k', '-bufsize', '6000k',
'-vf', 'scale=-2:720',
'-pix_fmt', 'yuv420p', '-g', '48', '-threads', '0',
];
case 'high':
default:
return gpu
? [
'-c:v', 'h264_nvenc', '-preset', 'p4', '-tune', 'hq',
'-rc', 'cbr', '-b:v', '4000k', '-maxrate', '4500k',
'-bufsize', '8000k', '-g', '48', '-bf', '2',
'-pix_fmt', 'yuv420p',
]
: [
'-c:v', 'libx264', '-preset', 'medium', '-crf', '18',
'-maxrate', '12000k', '-bufsize', '24000k',
'-pix_fmt', 'yuv420p', '-g', '48', '-threads', '0',
];
}
}
/// Audio args. Source mode keeps it simple (plain AAC); transcoded modes
/// keep the downmix + loudness normalization filter.
List<String> _audioArgs({required bool withFilters}) {
return [
'-c:a', 'aac', '-b:a', '192k', '-ac', '2', '-ar', '48000',
if (withFilters)
...['-af',
'pan=stereo|FL=1.0*FL+0.707*FC+0.5*BL+0.5*SL+0.5*LFE|FR=1.0*FR+0.707*FC+0.5*BR+0.5*SR+0.5*LFE,dynaudnorm=f=150:g=15']
else
...['-af', 'aresample=async=1'],
];
}
Map<String, String> _hlsHeaders() => {
'Content-Type': 'application/vnd.apple.mpegurl',
'Cache-Control': 'no-cache',
'X-Content-Type-Options': 'nosniff',
};
Map<String, String> _segmentHeaders({int maxAge = 60}) => {
'Content-Type': 'video/mp2t',
'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),
// Touchée à chaque requête de playlist ou de segment.
lastActivity: () => session.lastAccess,
);
}
final outcome =
await sessionManager.waitForPlaylist(session, minSegments: minSegments);
return (session: session, ready: outcome.ready, error: outcome.error);
}
/// 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)
// ==========================================
/// Connexion amont ouverte, avec son client pour pouvoir le refermer.
class _LiveUpstream {
final http.Client client;
final http.StreamedResponse response;
_LiveUpstream(this.client, this.response);
}
Future<_LiveUpstream?> _openLiveUpstream(String targetUrl) async {
final client = http.Client();
try {
final request = http.Request('GET', Uri.parse(targetUrl))
..headers['User-Agent'] = 'VLC/3.0.18 LibVLC/3.0.18'
..headers['Accept'] = '*/*';
final response = await client.send(request);
return _LiveUpstream(client, response);
} catch (_) {
client.close();
return null;
}
}
/// Nombre de reconnexions consécutives tentées avant d'abandonner le flux.
const _liveProxyMaxRetries = 4;
/// Corps du proxy `.ts` direct, avec reconnexion amont.
///
/// Les panneaux Xtream ferment régulièrement la connexion en cours de route
/// (bascule de source, limite de connexions simultanées, recyclage nginx).
/// Sans reprise, le `MediaSource` du navigateur reçoit un `sourceEnded` et la
/// lecture s'arrête net — le symptôme « ça coupe au bout de 30 s ».
/// On rouvre donc la source tant que le client, lui, écoute toujours.
///
/// Le compteur de tentatives est remis à zéro dès qu'une connexion a duré
/// assez longtemps pour être considérée saine : une coupure toutes les
/// 10 minutes ne doit pas finir par épuiser le quota.
Stream<List<int>> _resilientLiveBody(
String streamId,
String targetUrl,
_LiveUpstream first,
) async* {
var upstream = first;
var attempt = 0;
while (true) {
final startedAt = DateTime.now();
var bytes = 0;
try {
await for (final chunk in upstream.response.stream) {
bytes += chunk.length;
yield chunk;
}
} catch (e) {
// Un `ClientException` recopie l'URL amont, identifiants compris.
print(
'[Live Proxy] $streamId : coupure amont '
'(${LogRedactor.redactUrl('$e')})',
);
} finally {
upstream.client.close();
}
final lasted = DateTime.now().difference(startedAt);
if (lasted > const Duration(seconds: 60) && bytes > 0) {
attempt = 0;
}
print(
'[Live Proxy] $streamId : flux terminé après ${lasted.inSeconds}s '
'($bytes octets)',
);
// Rouvrir la source, en réessayant tant qu'il reste du quota.
_LiveUpstream? next;
while (next == null && attempt < _liveProxyMaxRetries) {
attempt++;
// Laisser le panneau libérer le slot de connexion précédent : sur les
// comptes limités à une connexion simultanée, rouvrir immédiatement se
// fait refuser.
await Future<void>.delayed(const Duration(seconds: 1));
final candidate = await _openLiveUpstream(targetUrl);
if (candidate != null && candidate.response.statusCode == 200) {
next = candidate;
} else {
candidate?.client.close();
print(
'[Live Proxy] $streamId : reconnexion '
'$attempt/$_liveProxyMaxRetries échouée',
);
}
}
if (next == null) {
print('[Live Proxy] $streamId : abandon, source injoignable');
return;
}
upstream = next;
}
}
/// 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,
}) {
final router = Router();
Future<Response> servePlaylist(
Request request,
String streamId,
String quality,
) async {
if (!_isValidStreamId(streamId)) {
return Response.badRequest(body: 'Invalid stream ID');
}
final playlist = await getPlaylist(request);
if (playlist == null) return Response.forbidden('No playlist');
final useNvidiaGpu = isGpuEnabled?.call() ?? _isNvidiaGpuEnabled();
final sessionId = 'live_${streamId}_$quality';
if (!sessionManager.contains(sessionId)) {
print('[Live HLS] Starting $sessionId');
}
Future<_SessionAttempt> run({required bool hls}) => _runWithGpuFallback(
id: sessionId,
isLive: true,
upstream: playlist,
wantGpu: useNvidiaGpu && quality != 'source',
buildArgs: (gpu) => [
'-hide_banner', '-loglevel', 'warning',
if (gpu) ...['-hwaccel', 'cuda'],
...liveInputArgs(
_liveSourceUrl(playlist, streamId, hls: hls),
hls: hls,
),
..._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).
// hls_init_time 1: first segment closes after ~1s so playback can
// start sooner; subsequent segments use hls_time.
'-f', 'hls',
'-hls_init_time', '1',
'-hls_time', '2',
'-hls_list_size', '10',
'-hls_flags', 'delete_segments+independent_segments',
'-hls_segment_type', 'mpegts',
'-hls_segment_filename', 'seg_%03d.ts',
'playlist.m3u8',
],
);
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);
return Response(502, body: 'Live transcoder failed: ${result.error}');
}
session.touch();
return Response.ok(
File('${session.dir.path}/playlist.m3u8').openRead(),
headers: _hlsHeaders(),
);
}
// Route: /api/live/{streamId}/{quality}/playlist.m3u8
router.get('/<streamId>/<quality>/playlist.m3u8',
(Request request, String streamId, String quality) {
return servePlaylist(request, streamId, _sanitizeQuality(quality));
});
// Back-compat route: /api/live/{streamId}/playlist.m3u8 (?quality=...)
// Redirects so relative segment URLs resolve inside the quality path.
router.get('/<streamId>/playlist.m3u8',
(Request request, String streamId) async {
final quality =
_sanitizeQuality(request.url.queryParameters['quality']);
return Response.found('/api/live/$streamId/$quality/playlist.m3u8');
});
// Route: /api/live/{streamId}/turbo.ts — flux MPEG-TS continu pour le
// player web (mpegts.js) : vidéo copiée telle quelle, audio TOUJOURS
// réencodé en AAC.
//
// POURQUOI : mpegts.js ne démuxe que l'AAC et le MP3. Les chaînes qui
// diffusent en AC-3/E-AC-3 (0x81/0x87) ou en MPEG-1 Layer II passaient par
// le proxy brut → image sans aucun son. Le réencodage audio seul coûte
// quelques % de CPU (pas de transcodage vidéo) et garde le zapping
// instantané : un seul flux continu, pas de segmentation HLS à amorcer.
//
// Latence : -fflags nobuffer + probesize réduit → première image en
// ~0,5-1,5 s au lieu des 2-6 s d'analyse par défaut de FFmpeg.
router.get('/<streamId>/turbo.ts', (Request request, String streamId) async {
if (!_isValidStreamId(streamId)) {
return Response.badRequest(body: 'Invalid stream ID');
}
final playlist = await getPlaylist(request);
if (playlist == null) return Response.forbidden('No playlist');
print('[Live Turbo] $streamId');
// 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);
// 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();
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;
}
_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();
return Response(
200,
body: body,
headers: {
'Content-Type': 'video/mp2t',
'Cache-Control': 'no-store',
'Connection': 'keep-alive',
},
);
});
// Route: /api/live/{streamId}.ts (Direct Proxy for recordings/raw playback)
router.get('/<streamId>.ts', (Request request, String streamId) async {
final playlist = await getPlaylist(request);
if (playlist == null) return Response.forbidden('No playlist');
final targetUrl =
'${playlist.dns}/live/${playlist.username}/${playlist.password}/$streamId.ts';
print(
'[Live Proxy] Forwarding $streamId: ${LogRedactor.redactUrl(targetUrl)}');
// Première connexion faite hors du générateur : elle seule décide du code
// de réponse renvoyé au client.
final first = await _openLiveUpstream(targetUrl);
if (first == null) {
return Response(502, body: 'Upstream connection failed');
}
if (first.response.statusCode != 200) {
first.client.close();
return Response(first.response.statusCode, body: 'Upstream refused');
}
return Response(
200,
body: _resilientLiveBody(streamId, targetUrl, first),
headers: {
'Content-Type': 'video/mp2t',
'Connection': 'keep-alive',
},
);
});
// Route: /api/live/{streamId}/{quality}/{segment}
router.get('/<streamId>/<quality>/<segment>',
(Request request, String streamId, String quality, String segment) {
if (!_isValidStreamId(streamId) || segment.contains('..')) {
return Response.badRequest(body: 'Invalid request');
}
final sessionId = 'live_${streamId}_${_sanitizeQuality(quality)}';
sessionManager.touch(sessionId);
final file = File('${_hlsTempDir.path}/$sessionId/$segment');
if (!file.existsSync()) return Response.notFound('Segment not found');
return Response.ok(file.openRead(), headers: _segmentHeaders());
});
return router.call;
}
// ==========================================
// 2. VOD HANDLER (FFmpeg HLS Transcoding)
// ==========================================
Handler createVodStreamHandler(
Future<PlaylistConfig?> Function(Request) getPlaylist, {
bool Function()? isGpuEnabled,
}) {
final router = Router();
Future<Response> servePlaylist(
Request request,
String streamId,
String quality,
) async {
if (!_isValidStreamId(streamId)) {
return Response.badRequest(body: 'Invalid stream ID');
}
final playlist = await getPlaylist(request);
if (playlist == null) {
return Response.internalServerError(body: 'No playlist');
}
// ?type=series or ?type=movie selects the upstream path
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.${vodContainerExtension(request.url.queryParameters['ext'])}';
final useNvidiaGpu = isGpuEnabled?.call() ?? _isNvidiaGpuEnabled();
final sessionId = 'vod_${streamId}_$quality';
if (!sessionManager.contains(sessionId)) {
print(
'[VOD] Starting $sessionId ($contentType): ${LogRedactor.redactUrl(targetUrl)}',
);
}
final result = await _runWithGpuFallback(
id: sessionId,
isLive: false,
upstream: playlist,
wantGpu: useNvidiaGpu && quality != 'source',
buildArgs: (gpu) => [
'-hide_banner', '-loglevel', 'warning',
if (gpu) ...['-hwaccel', 'cuda'],
'-headers', 'User-Agent: VLC/3.0.18 LibVLC/3.0.18\r\n',
'-reconnect', '1',
'-reconnect_at_eof', '1',
'-reconnect_streamed', '1',
'-reconnect_delay_max', '10',
'-rw_timeout', '30000000',
'-timeout', '30000000',
// Sonde réduite (10 Mo → 5 Mo) : sur un upstream lent, télécharger
// 10 Mo avant la première image ajoutait plusieurs secondes au
// démarrage de chaque film.
'-analyzeduration', '2000000',
'-probesize', '5000000',
'-i', targetUrl,
..._vodVideoArgs(quality, gpu),
..._audioArgs(withFilters: quality != 'source'),
'-f', 'hls',
'-hls_time', '4',
'-hls_list_size', '0',
'-hls_playlist_type', 'event',
'-hls_allow_cache', '1',
'-hls_flags', 'independent_segments',
'-hls_segment_type', 'mpegts',
'-hls_segment_filename', 'segment_%03d.ts',
'-start_number', '0',
'playlist.m3u8',
],
);
final session = result.session;
if (!result.ready || session == null) {
sessionManager.killSession(sessionId);
return Response(502, body: 'VOD transcoder failed: ${result.error}');
}
session.touch();
return Response.ok(
File('${session.dir.path}/playlist.m3u8').openRead(),
headers: _hlsHeaders(),
);
}
// Route: /api/vod/{streamId}/{quality}/playlist.m3u8
router.get('/<streamId>/<quality>/playlist.m3u8',
(Request request, String streamId, String quality) {
return servePlaylist(request, streamId, _sanitizeQuality(quality));
});
// Back-compat route: /api/vod/{streamId}/playlist.m3u8 (?quality=...)
router.get('/<streamId>/playlist.m3u8',
(Request request, String streamId) async {
final quality =
_sanitizeQuality(request.url.queryParameters['quality']);
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}
router.get('/<streamId>/<quality>/<segment>',
(Request request, String streamId, String quality, String segment) {
if (!_isValidStreamId(streamId) || segment.contains('..')) {
return Response.badRequest(body: 'Invalid request');
}
final sessionId = 'vod_${streamId}_${_sanitizeQuality(quality)}';
sessionManager.touch(sessionId);
final file = File('${_hlsTempDir.path}/$sessionId/$segment');
if (!file.existsSync()) return Response.notFound('Segment not found');
return Response.ok(file.openRead(), headers: _segmentHeaders(maxAge: 3600));
});
return router.call;
}
// ==========================================
// 3. RECORDING HANDLER (FFmpeg HLS Transcoding)
// ==========================================
Handler createRecordingStreamHandler(
AppDatabase db, {
bool Function()? isGpuEnabled,
}) {
final router = Router();
router.get('/<streamId>/playlist.m3u8',
(Request request, String streamId) async {
if (!_isValidStreamId(streamId)) {
return Response.badRequest(body: 'Invalid request');
}
final recording = db.getRecordingById(streamId);
if (recording == null) {
return Response.notFound('Recording not found');
}
// Check recording status - only allow streaming if completed or recording
if (recording.status == 'scheduled') {
return Response(202, body: 'Recording not yet started');
}
if (recording.status == 'failed') {
return Response(422,
body: 'Recording failed: ${recording.errorReason ?? 'unknown error'}');
}
if (recording.status != 'completed' && recording.status != 'recording') {
return Response.internalServerError(
body: 'Invalid recording status: ${recording.status}');
}
if (recording.filePath == null) {
return Response.internalServerError(body: 'Recording file path not set');
}
final targetUrl = recording.filePath!;
if (!File(targetUrl).existsSync()) {
return Response.internalServerError(
body: 'Physical recording file not found');
}
final start = parseRecordingStart(request.url.queryParameters['start']);
final offsetKey = recordingOffsetKey(start);
final sessionPrefix = 'rec_${streamId}_t';
final sessionId = 'rec_${streamId}_$offsetKey';
// La capture est faite en `-c copy` : le fichier contient le codec
// d'origine de la chaîne, presque toujours du H.264 — directement
// lisible en HLS. Le ré-encoder tenait à peine le temps réel en 1080p,
// si bien que la lecture démarrait collée au front d'encodage et se
// coupait à la moindre hésitation. En copie, la segmentation va à la
// vitesse du disque : l'enregistrement devient navigable en quelques
// secondes. Seuls les codecs que le navigateur ne sait pas lire
// (HEVC, MPEG-2…) justifient encore un ré-encodage.
final videoCodec = await MediaProbe.videoCodec(targetUrl);
final canCopyVideo = videoCodec == 'h264';
List<String> buildArgs({required bool copyVideo, required bool gpu}) {
final useNvidiaGpu = !copyVideo && gpu;
return [
'-hide_banner', '-loglevel', 'warning',
if (useNvidiaGpu) ...['-hwaccel', 'cuda'],
// -ss AVANT -i : recherche rapide par index, sans décoder ce qui
// précède. La sortie repart de 0, le lecteur ré-ajoute l'offset.
// En copie, la position atteinte est l'image-clé qui précède.
if (start > 0) ...['-ss', '$start'],
'-i', targetUrl,
if (copyVideo) ...[
'-c:v', 'copy',
] else if (useNvidiaGpu) ...[
'-c:v', 'h264_nvenc', '-preset', 'p4', '-tune', 'hq',
'-rc', 'cbr', '-b:v', '3000k', '-maxrate', '3500k',
'-bufsize', '6000k',
'-g', '48', '-bf', '2', '-pix_fmt', 'yuv420p',
] else ...[
// veryfast/crf 20 au lieu de medium/crf 18 : à ce niveau le rendu
// est visuellement identique, mais medium ne tenait pas le temps
// réel en 1080p sur un CPU modeste → lecture d'enregistrement qui
// démarre lentement puis bufferise.
'-c:v', 'libx264', '-preset', 'veryfast', '-crf', '20',
'-maxrate', '12000k', '-bufsize', '24000k', '-pix_fmt', 'yuv420p',
'-g', '48', '-threads', '0',
],
..._audioArgs(withFilters: true),
// `event` et non `vod` : en VOD, hls.js considère la playlist comme
// définitive et ne la relit jamais. Il ne voyait donc que les
// quelques secondes déjà transcodées à l'ouverture — d'où une durée
// absurde et une barre de progression inutilisable. En `event`, la
// playlist est rechargée et la zone navigable grandit avec l'encodage.
'-f', 'hls', '-hls_time', '4', '-hls_list_size', '0',
'-hls_playlist_type', 'event', '-hls_allow_cache', '1',
'-hls_flags', 'independent_segments', '-hls_segment_type', 'mpegts',
'-hls_segment_filename', 'segment_%03d.ts', '-start_number', '0',
'playlist.m3u8',
];
}
/// Démarre (ou récupère) la session et attend qu'elle ait de l'avance.
///
/// Trois segments avant de rendre la main : c'est la marge qui manquait
/// 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<_SessionAttempt> run(bool copyVideo) {
if (!sessionManager.contains(sessionId)) {
print(
'[Recording] $sessionId : '
'${copyVideo ? 'copie vidéo' : 'ré-encodage'} '
'(codec source : ${videoCodec ?? 'inconnu'})',
);
}
if (copyVideo) {
return _runSession(
id: sessionId,
isLive: false,
argsBuilder: (_) => buildArgs(copyVideo: true, gpu: false),
minSegments: 3,
);
}
return _runWithGpuFallback(
id: sessionId,
isLive: false,
wantGpu: isGpuEnabled?.call() ?? _isNvidiaGpuEnabled(),
buildArgs: (gpu) => buildArgs(copyVideo: false, gpu: gpu),
minSegments: 3,
);
}
var attempt = await run(canCopyVideo);
// La copie a échoué : certains conteneurs ou horodatages ne s'y prêtent
// pas. Plutôt qu'un écran d'erreur, on retombe sur le ré-encodage, qui
// était le seul chemin jusqu'ici.
if (!attempt.ready && canCopyVideo) {
print(
'[Recording] $sessionId : copie vidéo refusée (${attempt.error}) '
'— repli sur le ré-encodage',
);
sessionManager.killSession(sessionId);
attempt = await run(false);
}
final session = attempt.session;
if (!attempt.ready || session == null) {
sessionManager.killSession(sessionId);
return Response(502,
body: 'Recording transcoder failed: ${attempt.error}');
}
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.
sessionManager.killIdleSiblings(sessionPrefix, keep: sessionId);
// Les segments sont préfixés par l'offset pour rester rattachés à LEUR
// session : `segment_000.ts` d'une playlist à 0 s et d'une playlist à
// 45 min sont deux fichiers différents.
final playlist = await File('${session.dir.path}/playlist.m3u8')
.readAsString();
return Response.ok(
rewriteRecordingPlaylist(playlist, offsetKey),
headers: _hlsHeaders(),
);
});
Response serveSegment(String streamId, String offsetKey, String segment) {
if (!_isValidStreamId(streamId) ||
parseRecordingOffsetKey(offsetKey) == null ||
segment.contains('..') ||
segment.contains('/')) {
return Response.badRequest(body: 'Invalid request');
}
final sessionId = 'rec_${streamId}_$offsetKey';
sessionManager.touch(sessionId);
final file = File('${_hlsTempDir.path}/$sessionId/$segment');
if (!file.existsSync()) return Response.notFound('Segment not found');
return Response.ok(file.openRead(), headers: _segmentHeaders(maxAge: 3600));
}
router.get('/<streamId>/<offsetKey>/<segment>',
(Request request, String streamId, String offsetKey, String segment) {
return serveSegment(streamId, offsetKey, segment);
});
// Compat : playlist servie avant la mise en place des offsets, encore en
// cache dans un onglet ouvert.
router.get('/<streamId>/<segment>',
(Request request, String streamId, String segment) {
return serveSegment(streamId, recordingOffsetKey(0), segment);
});
return router.call;
}