diff --git a/bin/api/recordings_api.dart b/bin/api/recordings_api.dart index f9ee9ba..67f732e 100644 --- a/bin/api/recordings_api.dart +++ b/bin/api/recordings_api.dart @@ -36,8 +36,8 @@ class RecordingsApi { // (l'ancien replaceAll('.mp4', '.log') ne trouvait jamais le fichier). final logFilePath = p.setExtension(recording.filePath!, '.log'); - // Anti path-traversal : le log doit rester dans /app/recordings - final safeLogPath = SafePath.resolveWithin('/app/recordings', logFilePath); + // Anti path-traversal : le log doit rester dans le dossier des enregistrements + final safeLogPath = SafePath.resolveWithin(recordingsDirPath, logFilePath); if (safeLogPath == null) { return Response.forbidden( json.encode({'error': 'Chemin de log invalide'}), diff --git a/bin/database/database.dart b/bin/database/database.dart index 6c6f248..71dc5a5 100644 --- a/bin/database/database.dart +++ b/bin/database/database.dart @@ -529,17 +529,22 @@ class AppDatabase { } /// Mettre à jour le statut et éventuellement le chemin d'un enregistrement + /// + /// [clearError] efface le motif d'erreur au lieu de conserver l'ancien : + /// utile quand un enregistrement précédemment en échec est relancé. void updateRecordingStatus( String id, String status, { String? filePath, String? errorReason, + bool clearError = false, }) { final now = DateTime.now().toIso8601String(); + final errorExpr = clearError ? '?' : 'COALESCE(?, error_reason)'; _db.execute( ''' UPDATE tv_recordings - SET status = ?, file_path = COALESCE(?, file_path), error_reason = COALESCE(?, error_reason), updated_at = ? + SET status = ?, file_path = COALESCE(?, file_path), error_reason = $errorExpr, updated_at = ? WHERE id = ? ''', [status, filePath, errorReason, now, id], diff --git a/bin/services/recording_decision.dart b/bin/services/recording_decision.dart index 3e385e6..2a8f315 100644 --- a/bin/services/recording_decision.dart +++ b/bin/services/recording_decision.dart @@ -31,3 +31,96 @@ RecordingAction decideRecordingAction({ if (activeCount >= maxConcurrent) return RecordingAction.wait; return RecordingAction.start; } + +/// Ce que le planificateur doit faire d'un enregistrement resté au statut +/// « recording » en base alors qu'aucun processus FFmpeg ne lui correspond +/// (redémarrage du conteneur, crash de l'isolate…). +enum OrphanAction { + /// La fenêtre est encore ouverte : relancer la capture. + resume, + + /// La fenêtre est passée mais un fichier partiel existe : le conserver. + finish, + + /// Rien d'exploitable n'a été capturé. + fail, +} + +/// Décide du sort d'un enregistrement orphelin détecté au démarrage/à un tick. +OrphanAction decideOrphanAction({ + required DateTime now, + required DateTime endTime, + required bool hasFile, + Duration minRemaining = const Duration(seconds: 60), +}) { + final remaining = endTime.toUtc().difference(now.toUtc()); + if (remaining > minRemaining) return OrphanAction.resume; + return hasFile ? OrphanAction.finish : OrphanAction.fail; +} + +/// Ce que le planificateur doit faire quand un processus FFmpeg se termine. +enum PostExitAction { + /// La capture est allée jusqu'au bout (ou assez loin) : marquer terminé. + complete, + + /// FFmpeg s'est arrêté trop tôt : relancer la capture sur la fin de fenêtre. + retry, + + /// Arrêt prématuré sans rien d'exploitable. + fail, +} + +/// Nombre d'échecs FFmpeg consécutifs tolérés avant d'abandonner un +/// enregistrement dont la fenêtre est encore ouverte. +/// +/// Le compteur est remis à zéro dès qu'un lancement a duré assez longtemps +/// pour être considéré sain : une coupure toutes les dix minutes ne doit pas +/// finir par épuiser le quota. +const int maxFfmpegAttempts = 30; + +/// Décide de la suite à donner à la sortie d'un processus FFmpeg. +/// +/// FFmpeg rend la main dès que la source se tarit : sur un flux IPTV, une +/// coupure amont à mi-parcours ne doit pas condamner l'heure restante. +PostExitAction decidePostExitAction({ + required DateTime now, + required DateTime endTime, + required int exitCode, + required int consecutiveFailures, + required bool hasFile, + int maxAttempts = maxFfmpegAttempts, + Duration minRemaining = const Duration(seconds: 60), +}) { + final remaining = endTime.toUtc().difference(now.toUtc()); + if (remaining > minRemaining && consecutiveFailures + 1 < maxAttempts) { + return PostExitAction.retry; + } + // `255` est le code renvoyé par FFmpeg quand on l'interrompt proprement. + final cleanExit = exitCode == 0 || exitCode == 255; + return (cleanExit || hasFile) ? PostExitAction.complete : PostExitAction.fail; +} + +/// Durée de capture à demander à FFmpeg (`-t`). +/// +/// Se base sur le temps qu'il reste jusqu'à la fin programmée, et non sur la +/// durée théorique du programme : un démarrage tardif (ou une relance après +/// coupure) ne doit pas décaler la fin de l'enregistrement. +Duration captureDuration({ + required DateTime now, + required DateTime endTime, + Duration minimum = const Duration(seconds: 30), +}) { + final remaining = endTime.toUtc().difference(now.toUtc()); + return remaining < minimum ? minimum : remaining; +} + +/// Délai avant relance de FFmpeg, croissant avec les échecs consécutifs. +/// +/// Rouvrir la source dans la seconde se fait refuser sur les comptes limités +/// en connexions simultanées ; s'acharner à cette cadence pendant une panne +/// amont ne ferait qu'épuiser le quota de relances. +Duration ffmpegRetryDelay(int consecutiveFailures) { + final steps = (consecutiveFailures - 1).clamp(0, 4); + final seconds = 3 * (1 << steps); + return Duration(seconds: seconds > 30 ? 30 : seconds); +} diff --git a/bin/services/recording_scheduler.dart b/bin/services/recording_scheduler.dart index 395c8c8..27cabbf 100644 --- a/bin/services/recording_scheduler.dart +++ b/bin/services/recording_scheduler.dart @@ -8,10 +8,108 @@ import '../utils/log_redactor.dart'; import 'recording_decision.dart'; import 'package:path/path.dart' as p; +/// Dossier où sont écrits les enregistrements et leurs logs. +final String recordingsDirPath = + Platform.environment['RECORDINGS_DIR'] ?? '/app/recordings'; + +/// Binaire FFmpeg utilisé pour la capture et la fusion des parties. +Future resolveFfmpegPath() async { + final fromEnv = Platform.environment['FFMPEG_PATH']; + if (fromEnv != null && fromEnv.isNotEmpty) return fromEnv; + if (Platform.isLinux && await File('/usr/local/bin/ffmpeg').exists()) { + return '/usr/local/bin/ffmpeg'; + } + return 'ffmpeg'; // Depuis le PATH +} + +/// Écrivain de log tolérant aux fermetures. +/// +/// FFmpeg continue de livrer des morceaux de `stderr` après la résolution de +/// `exitCode` : écrire sur un `IOSink` déjà fermé lève une `StateError` depuis +/// un callback de stream, c'est-à-dire une erreur asynchrone non rattrapée qui +/// tue l'isolate — et donc le serveur entier, laissant l'enregistrement en +/// cours au statut « recording » (« Interruption inattendue du serveur »). +class _RecordingLog { + final IOSink _sink; + bool _closed = false; + + _RecordingLog(this._sink) { + // Une erreur d'écriture disque ne doit pas non plus remonter en erreur + // asynchrone non rattrapée. + _sink.done.catchError((_) {}); + } + + /// Ouvre le log, ou renvoie `null` si le fichier n'est pas créable : + /// l'enregistrement reste prioritaire sur sa journalisation. + static _RecordingLog? open(String path, {bool append = false}) { + try { + return _RecordingLog( + File(path).openWrite(mode: append ? FileMode.append : FileMode.write), + ); + } catch (e) { + print( + '[RecordingScheduler] AVERTISSEMENT: log indisponible ($path): $e', + ); + return null; + } + } + + void writeln([String line = '']) => _guard(() => _sink.writeln(line)); + + void add(List data) => _guard(() => _sink.add(data)); + + void _guard(void Function() action) { + if (_closed) return; + try { + action(); + } catch (_) { + // Sink cassé : on arrête d'écrire plutôt que de propager. + _closed = true; + } + } + + Future close() async { + if (_closed) return; + _closed = true; + try { + await _sink.close(); + } catch (_) {} + } +} + class _ActiveRecording { final Recording recording; - final Process process; - _ActiveRecording(this.recording, this.process); + + /// Fichier principal (celui exposé par l'API). + final String filePath; + + final _RecordingLog? log; + + /// Fichiers additionnels écrits après une relance, dans l'ordre. + final List extraParts; + + /// Index du lancement FFmpeg en cours (0 = fichier principal). + int attempt; + + /// Échecs consécutifs de FFmpeg, remis à zéro dès qu'une capture tient + /// assez longtemps pour être considérée saine. + int consecutiveFailures = 0; + + /// Heure du dernier lancement de FFmpeg. + DateTime launchedAt = DateTime.now(); + + Process? process; + + /// Arrêt volontaire : la sortie de FFmpeg ne doit pas déclencher de relance. + bool stopping = false; + + _ActiveRecording({ + required this.recording, + required this.filePath, + required this.log, + this.attempt = 0, + List? extraParts, + }) : extraParts = extraParts ?? []; } class RecordingScheduler { @@ -73,8 +171,6 @@ class RecordingScheduler { // Toujours comparer en UTC pour éviter les problèmes de fuseau horaire final now = DateTime.now().toUtc(); - final recordings = _db.getAllRecordings(); - // Arrêter les enregistrements actifs dont l'heure de fin est passée for (final id in _active.keys.toList()) { final active = _active[id]!; @@ -86,19 +182,18 @@ class RecordingScheduler { } } + // Lire la base APRÈS les arrêts : un instantané pris avant ferait passer + // l'enregistrement tout juste terminé pour un orphelin (il n'est plus + // dans `_active` alors que l'instantané le dit encore « recording »). + final recordings = _db.getAllRecordings(); + // Rechercher les enregistrements planifiés for (final recording in recordings) { - // Nettoyer les enregistrements bloqués "recording" suite à un crash serveur + // Enregistrement marqué "recording" sans processus associé : le serveur + // a redémarré (ou l'enregistrement n'a jamais démarré). if (recording.status == 'recording' && !_active.containsKey(recording.id)) { - print( - '[RecordingScheduler] Enregistrement orphelin détecté: ${recording.id}', - ); - _db.updateRecordingStatus( - recording.id, - 'failed', - errorReason: 'Interruption inattendue du serveur', - ); + await _recoverOrphan(recording, now); continue; } @@ -145,7 +240,52 @@ class RecordingScheduler { } } - /// Vérifier les Season Passes et créer des enregistrements pour les nouvelles diffusions + /// Reprend (ou clôture) un enregistrement laissé au statut "recording" par + /// une interruption du serveur. + Future _recoverOrphan(Recording recording, DateTime now) async { + // Le statut a pu changer depuis la lecture (clôture asynchrone en cours) : + // ne jamais requalifier un enregistrement sur une donnée périmée. + final current = _db.getRecordingById(recording.id); + if (current == null || current.status != 'recording') return; + + final action = decideOrphanAction( + now: now, + endTime: current.endTime, + hasFile: _fileHasData(current.filePath), + ); + + switch (action) { + case OrphanAction.resume: + if (_active.length >= maxConcurrent) { + // Pas de slot libre pour l'instant : on retentera au prochain tick. + return; + } + print( + '[RecordingScheduler] Reprise après interruption : ${current.title}', + ); + await _startRecording(current, resume: true); + case OrphanAction.finish: + print( + '[RecordingScheduler] Enregistrement interrompu conservé (partiel): ${current.title}', + ); + _db.updateRecordingStatus( + current.id, + 'completed', + errorReason: + 'Interruption du serveur pendant la capture : fichier partiel', + ); + case OrphanAction.fail: + print( + '[RecordingScheduler] Enregistrement orphelin sans fichier: ${current.id}', + ); + _db.updateRecordingStatus( + current.id, + 'failed', + errorReason: 'Interruption inattendue du serveur', + ); + } + } + Future _checkSeasonPasses() async { final passes = _db.getAllSeasonPasses(); if (passes.isEmpty) return; @@ -233,14 +373,18 @@ class RecordingScheduler { } } - Future _startRecording(Recording recording) async { - _db.updateRecordingStatus(recording.id, 'recording'); + Future _startRecording( + Recording recording, { + bool resume = false, + }) async { + // Un (re)démarrage rend obsolète un éventuel motif d'échec précédent. + _db.updateRecordingStatus(recording.id, 'recording', clearError: true); // Un seul bloc try/catch englobant TOUT pour éviter les Unhandled exceptions // qui tueraient le serveur entier try { // Préparation du dossier d'enregistrement - final recordingsDir = Directory('/app/recordings'); + final recordingsDir = Directory(recordingsDirPath); if (!await recordingsDir.exists()) { await recordingsDir.create(recursive: true); } @@ -248,106 +392,37 @@ class RecordingScheduler { // Nettoyer l'espace disque si nécessaire await _checkDiskSpaceAndRotate(recordingsDir); - // Génération d'un nom de fichier unique et sûr - final safeTitle = - recording.title.replaceAll(RegExp(r'[^a-zA-Z0-9_\-]'), '_'); - final dateStr = recording.startTime - .toUtc() - .toIso8601String() - .replaceAll(':', '') - .split('.')[0]; - final fileName = '${safeTitle}_$dateStr.mkv'; - final filePath = p.join(recordingsDir.path, fileName); - final logFilePath = filePath.replaceAll('.mkv', '.log'); + final filePath = p.join(recordingsDir.path, _fileNameFor(recording)); + final logPath = p.setExtension(filePath, '.log'); - // Résoudre l'URL relative en URL absolue pour FFmpeg - // Le serveur Xtremflow tourne sur le port 8089 en interne Docker - String streamUrl = recording.streamUrl; - if (streamUrl.startsWith('/')) { - streamUrl = 'http://localhost:8089$streamUrl'; - } + final active = _ActiveRecording( + recording: recording, + filePath: filePath, + // Une reprise complète le log existant au lieu de l'écraser. + log: _RecordingLog.open(logPath, append: resume), + ); - final args = [ - '-y', - '-i', - streamUrl, - '-c', - 'copy', - '-t', - '${recording.endTime.difference(recording.startTime).inSeconds}', - filePath, - ]; - - // Créer le fichier de log de façon sécurisée (openWrite peut lancer une exception) - // Le contenu du log est exposé via l'API : ne jamais y écrire de credentials. - IOSink? logSink; - try { - logSink = File(logFilePath).openWrite(); - logSink.writeln('[${DateTime.now()}] Démarrage: ${recording.title}'); - logSink.writeln('URL: ${LogRedactor.redactUrl(streamUrl)}'); - logSink.writeln('Destination: $filePath'); - } catch (logError) { - // Si on ne peut pas créer le log, on continue quand même (l'enregistrement reste prioritaire) - print( - '[RecordingScheduler] AVERTISSEMENT: Impossible de créer le fichier log ($logFilePath): $logError', - ); - // On note la filePath dans la DB quand même pour pouvoir retourner le statut - _db.updateRecordingStatus( - recording.id, - 'recording', - filePath: filePath, + if (resume) { + // Le fichier principal (et d'éventuelles parties) survivent au + // redémarrage : on repart sur une nouvelle partie pour ne rien écraser. + active.extraParts.addAll(_existingParts(filePath)); + active.attempt = _fileHasData(filePath) || active.extraParts.isNotEmpty + ? active.extraParts.length + 1 + : 0; + active.log?.writeln( + '\n[${DateTime.now()}] Reprise après interruption du serveur', ); + } else { + active.log?.writeln('[${DateTime.now()}] Démarrage: ${recording.title}'); + active.log?.writeln('Destination: $filePath'); } - // Trouver ffmpeg - String ffmpegPath = 'ffmpeg'; - if (Platform.isLinux && await File('/usr/local/bin/ffmpeg').exists()) { - ffmpegPath = '/usr/local/bin/ffmpeg'; - } - - final redactedArgs = - args.map(LogRedactor.redactUrl).join(' '); - logSink?.writeln('Commande: $ffmpegPath $redactedArgs\n'); - print('[RecordingScheduler] Exécution: $ffmpegPath $redactedArgs'); - - final process = await Process.start(ffmpegPath, args); - _active[recording.id] = _ActiveRecording(recording, process); - - // Rediriger stdout/stderr dans le log - process.stdout.listen((event) => logSink?.add(event)); - process.stderr.listen((event) => logSink?.add(event)); + _active[recording.id] = active; // Enregistrer le chemin du fichier dans la BDD _db.updateRecordingStatus(recording.id, 'recording', filePath: filePath); - // Écouter la fin du processus FFmpeg de manière asynchrone - process.exitCode.then((exitCode) async { - if (_active.containsKey(recording.id)) { - logSink?.writeln( - '\n[${DateTime.now()}] FFmpeg terminé avec le code $exitCode', - ); - await logSink?.close(); - - if (exitCode == 0 || exitCode == 255) { - print( - '[RecordingScheduler] Enregistrement terminé: ${recording.title}', - ); - _db.updateRecordingStatus(recording.id, 'completed'); - } else { - print( - '[RecordingScheduler] Erreur FFmpeg (code: $exitCode) pour ${recording.title}', - ); - _db.updateRecordingStatus( - recording.id, - 'failed', - errorReason: 'Erreur FFmpeg code $exitCode. Voir logs.', - ); - } - _active.remove(recording.id); - } else { - await logSink?.close(); - } - }); + await _launchFfmpeg(active); } catch (e, st) { // Attraper TOUTES les exceptions pour éviter de crasher le serveur print('[RecordingScheduler] ERREUR dans _startRecording: $e\n$st'); @@ -356,10 +431,236 @@ class RecordingScheduler { 'failed', errorReason: 'Erreur au lancement: $e', ); - _active.remove(recording.id); + final active = _active.remove(recording.id); + unawaited(active?.log?.close() ?? Future.value()); } } + /// Lance (ou relance) le processus FFmpeg d'un enregistrement actif. + Future _launchFfmpeg(_ActiveRecording active) async { + final recording = active.recording; + final output = active.attempt == 0 + ? active.filePath + : _partPath(active.filePath, active.attempt); + if (active.attempt > 0 && !active.extraParts.contains(output)) { + active.extraParts.add(output); + } + + // Résoudre l'URL relative en URL absolue pour FFmpeg + // Le serveur Xtremflow tourne sur le port 8089 en interne Docker + String streamUrl = recording.streamUrl; + if (streamUrl.startsWith('/')) { + streamUrl = 'http://localhost:8089$streamUrl'; + } + + final duration = captureDuration( + now: DateTime.now().toUtc(), + endTime: recording.endTime.toUtc(), + ); + + final args = [ + '-hide_banner', + // Les logs sont relus entièrement par l'API : sans -nostats, une heure de + // capture produit des mégaoctets de lignes de progression. + '-nostats', + '-user_agent', 'VLC/3.0.18 LibVLC/3.0.18', + // Une coupure amont (bascule de source, recyclage nginx du panneau) ne + // doit pas terminer la capture : FFmpeg rouvre le flux lui-même. + '-reconnect', '1', + '-reconnect_at_eof', '1', + '-reconnect_streamed', '1', + '-reconnect_delay_max', '30', + '-rw_timeout', '30000000', + '-y', + '-i', streamUrl, + '-c', 'copy', + '-t', '${duration.inSeconds}', + output, + ]; + + final ffmpegPath = await resolveFfmpegPath(); + + // Le contenu du log est exposé via l'API : ne jamais y écrire de credentials. + final redactedArgs = args.map(LogRedactor.redactUrl).join(' '); + active.log?.writeln('URL: ${LogRedactor.redactUrl(streamUrl)}'); + active.log?.writeln('Commande: $ffmpegPath $redactedArgs\n'); + print('[RecordingScheduler] Exécution: $ffmpegPath $redactedArgs'); + + final process = await Process.start(ffmpegPath, args); + active.process = process; + active.launchedAt = DateTime.now(); + + // Rediriger stdout/stderr dans le log ; on attend la fin des deux flux + // avant de fermer le sink, sinon les derniers morceaux arrivent sur un + // fichier déjà fermé. + final drained = Future.wait([ + _pipeToLog(process.stdout, active.log), + _pipeToLog(process.stderr, active.log), + ]); + + unawaited( + process.exitCode.then((exitCode) async { + await drained; + await _onFfmpegExit(active, exitCode); + }).catchError((Object e, StackTrace st) { + print('[RecordingScheduler] ERREUR à la sortie de FFmpeg: $e\n$st'); + }), + ); + } + + Future _pipeToLog(Stream> stream, _RecordingLog? log) { + return stream.forEach((chunk) => log?.add(chunk)).catchError((_) {}); + } + + /// Décide de la suite quand un processus FFmpeg se termine. + Future _onFfmpegExit(_ActiveRecording active, int exitCode) async { + final recording = active.recording; + + // Arrêt volontaire, ou enregistrement déjà repris ailleurs : la clôture est + // gérée par celui qui a demandé l'arrêt. + if (active.stopping || !identical(_active[recording.id], active)) return; + + // Une capture qui a tenu longtemps repart d'un quota de relances neuf. + final lasted = DateTime.now().difference(active.launchedAt); + if (lasted > const Duration(minutes: 2) && _activeHasData(active)) { + active.consecutiveFailures = 0; + } + + final action = decidePostExitAction( + now: DateTime.now().toUtc(), + endTime: recording.endTime.toUtc(), + exitCode: exitCode, + consecutiveFailures: active.consecutiveFailures, + hasFile: _activeHasData(active), + ); + + switch (action) { + case PostExitAction.retry: + active.attempt++; + active.consecutiveFailures++; + final delay = ffmpegRetryDelay(active.consecutiveFailures); + final message = + 'FFmpeg s\'est arrêté (code $exitCode) après ${lasted.inSeconds}s, ' + 'avant la fin prévue — relance dans ${delay.inSeconds}s ' + '(tentative ${active.consecutiveFailures}/${maxFfmpegAttempts - 1})'; + print('[RecordingScheduler] ${recording.title}: $message'); + active.log?.writeln('\n[${DateTime.now()}] $message'); + // Laisser la source respirer : sur les comptes limités en connexions + // simultanées, rouvrir immédiatement se fait refuser. + await Future.delayed(delay); + if (active.stopping || !identical(_active[recording.id], active)) { + return; + } + try { + await _launchFfmpeg(active); + } catch (e) { + print('[RecordingScheduler] Relance impossible: $e'); + _active.remove(recording.id); + await _finalize( + active, + status: 'failed', + errorReason: 'Relance de FFmpeg impossible: $e', + ); + } + case PostExitAction.complete: + print( + '[RecordingScheduler] Enregistrement terminé: ${recording.title}', + ); + _active.remove(recording.id); + await _finalize(active, status: 'completed'); + case PostExitAction.fail: + print( + '[RecordingScheduler] Erreur FFmpeg (code: $exitCode) pour ${recording.title}', + ); + _active.remove(recording.id); + await _finalize( + active, + status: 'failed', + errorReason: 'Erreur FFmpeg code $exitCode. Voir logs.', + ); + } + } + + /// Clôture un enregistrement : fusion des parties, fermeture du log, statut. + Future _finalize( + _ActiveRecording active, { + required String status, + String? errorReason, + }) async { + try { + await _mergeParts(active); + } catch (e) { + print('[RecordingScheduler] Fusion des parties impossible: $e'); + } + active.log?.writeln( + '\n[${DateTime.now()}] Enregistrement clôturé avec le statut "$status"', + ); + await active.log?.close(); + _db.updateRecordingStatus( + active.recording.id, + status, + filePath: active.filePath, + errorReason: errorReason, + ); + } + + /// Recolle les parties issues des relances dans le fichier principal. + /// + /// Les parties proviennent de la même source avec `-c copy` : le demuxer + /// `concat` de FFmpeg les rassemble sans réencodage. En cas d'échec, le + /// fichier principal et les parties sont conservés tels quels. + Future _mergeParts(_ActiveRecording active) async { + final parts = active.extraParts.where(_fileHasData).toList(); + if (parts.isEmpty) return; + + final segments = [ + if (_fileHasData(active.filePath)) active.filePath, + ...parts, + ]; + if (segments.length < 2) { + // Le fichier principal est vide : la première partie le remplace. + await File(segments.first).rename(active.filePath); + return; + } + + final listFile = File('${active.filePath}.parts.txt'); + await listFile.writeAsString( + segments.map((s) => 'file \'${s.replaceAll("'", r"'\''")}\'').join('\n'), + ); + final mergedPath = '${active.filePath}.merged.mkv'; + final ffmpegPath = await resolveFfmpegPath(); + + active.log?.writeln( + '\n[${DateTime.now()}] Fusion de ${segments.length} parties…', + ); + final result = await Process.run(ffmpegPath, [ + '-hide_banner', + '-nostats', + '-y', + '-f', 'concat', + '-safe', '0', + '-i', listFile.path, + '-c', 'copy', + mergedPath, + ]); + await _deleteQuietly(listFile.path); + + if (result.exitCode != 0 || !_fileHasData(mergedPath)) { + active.log?.writeln( + 'Fusion échouée (code ${result.exitCode}) : les parties sont ' + 'conservées séparément.', + ); + await _deleteQuietly(mergedPath); + return; + } + + await File(mergedPath).rename(active.filePath); + for (final part in parts) { + await _deleteQuietly(part); + } + active.log?.writeln('Fusion terminée.'); + } + /// Arrêter un enregistrement en cours (appelé depuis l'API) Future stopRecording(String id) async { if (_active.containsKey(id)) { @@ -375,11 +676,83 @@ class RecordingScheduler { void _stopActiveRecording(String id, {String? reason}) { final active = _active.remove(id); if (active == null) return; + active.stopping = true; if (reason != null) { print('[RecordingScheduler] Arrêt: $reason (${active.recording.title})'); } - active.process.kill(ProcessSignal.sigterm); + active.log?.writeln( + '\n[${DateTime.now()}] Arrêt: ${reason ?? 'fin de la fenêtre programmée'}', + ); + active.process?.kill(ProcessSignal.sigterm); _db.updateRecordingStatus(id, 'completed'); + unawaited(_finalizeStopped(active)); + } + + /// Attend la fin effective de FFmpeg puis clôture proprement (fusion + log). + Future _finalizeStopped(_ActiveRecording active) async { + try { + final process = active.process; + if (process != null) { + await process.exitCode.timeout( + const Duration(seconds: 15), + onTimeout: () { + process.kill(ProcessSignal.sigkill); + return -1; + }, + ); + } + await _finalize(active, status: 'completed'); + } catch (e, st) { + print('[RecordingScheduler] ERREUR à la clôture: $e\n$st'); + await active.log?.close(); + } + } + + String _fileNameFor(Recording recording) { + // Génération d'un nom de fichier unique et sûr + final safeTitle = + recording.title.replaceAll(RegExp(r'[^a-zA-Z0-9_\-]'), '_'); + final dateStr = recording.startTime + .toUtc() + .toIso8601String() + .replaceAll(':', '') + .split('.')[0]; + return '${safeTitle}_$dateStr.mkv'; + } + + /// Chemin de la n-ième partie (la partie 1 étant le fichier principal). + String _partPath(String filePath, int attempt) => + '${p.withoutExtension(filePath)}.part${attempt + 1}.mkv'; + + /// Parties déjà écrites sur disque pour cet enregistrement, dans l'ordre. + List _existingParts(String filePath) { + final parts = []; + for (var attempt = 1;; attempt++) { + final path = _partPath(filePath, attempt); + if (!_fileHasData(path)) break; + parts.add(path); + } + return parts; + } + + bool _activeHasData(_ActiveRecording active) => + _fileHasData(active.filePath) || active.extraParts.any(_fileHasData); + + bool _fileHasData(String? path) { + if (path == null) return false; + try { + final file = File(path); + return file.existsSync() && file.lengthSync() > 0; + } catch (_) { + return false; + } + } + + Future _deleteQuietly(String path) async { + try { + final file = File(path); + if (await file.exists()) await file.delete(); + } catch (_) {} } Future _checkDiskSpaceAndRotate(Directory dir) async { diff --git a/bin/test/recording_decision_test.dart b/bin/test/recording_decision_test.dart index 57d3a40..3ab5615 100644 --- a/bin/test/recording_decision_test.dart +++ b/bin/test/recording_decision_test.dart @@ -70,4 +70,176 @@ void main() { ); }); }); + + group('decideOrphanAction', () { + test('resumes while the window is still open', () { + expect( + decideOrphanAction( + now: now, + endTime: now.add(const Duration(minutes: 30)), + hasFile: true, + ), + OrphanAction.resume, + ); + }); + + test('resumes even without a file (crash before the first byte)', () { + expect( + decideOrphanAction( + now: now, + endTime: now.add(const Duration(minutes: 30)), + hasFile: false, + ), + OrphanAction.resume, + ); + }); + + test('keeps a partial file when the window has closed', () { + expect( + decideOrphanAction( + now: now, + endTime: now.subtract(const Duration(minutes: 5)), + hasFile: true, + ), + OrphanAction.finish, + ); + }); + + test('fails when the window has closed with nothing captured', () { + expect( + decideOrphanAction( + now: now, + endTime: now.subtract(const Duration(minutes: 5)), + hasFile: false, + ), + OrphanAction.fail, + ); + }); + + test('does not resume for the last seconds of the window', () { + expect( + decideOrphanAction( + now: now, + endTime: now.add(const Duration(seconds: 20)), + hasFile: true, + ), + OrphanAction.finish, + ); + }); + }); + + group('decidePostExitAction', () { + test('retries when ffmpeg dies mid-window', () { + expect( + decidePostExitAction( + now: now, + endTime: now.add(const Duration(minutes: 40)), + exitCode: 1, + consecutiveFailures: 0, + hasFile: true, + ), + PostExitAction.retry, + ); + }); + + test('retries on a clean exit too (upstream ended early)', () { + expect( + decidePostExitAction( + now: now, + endTime: now.add(const Duration(minutes: 40)), + exitCode: 0, + consecutiveFailures: 2, + hasFile: true, + ), + PostExitAction.retry, + ); + }); + + test('stops retrying once the attempt budget is spent', () { + expect( + decidePostExitAction( + now: now, + endTime: now.add(const Duration(minutes: 40)), + exitCode: 1, + consecutiveFailures: maxFfmpegAttempts - 1, + hasFile: true, + ), + PostExitAction.complete, + ); + }); + + test('completes at the end of the window', () { + expect( + decidePostExitAction( + now: now, + endTime: now, + exitCode: 0, + consecutiveFailures: 0, + hasFile: true, + ), + PostExitAction.complete, + ); + }); + + test('completes when ffmpeg errored but a file was captured', () { + expect( + decidePostExitAction( + now: now, + endTime: now, + exitCode: 1, + consecutiveFailures: 0, + hasFile: true, + ), + PostExitAction.complete, + ); + }); + + test('fails when nothing was captured and the window is over', () { + expect( + decidePostExitAction( + now: now, + endTime: now, + exitCode: 1, + consecutiveFailures: 0, + hasFile: false, + ), + PostExitAction.fail, + ); + }); + }); + + group('captureDuration', () { + test('uses the time left until the scheduled end, not the planned length', + () { + // Démarrage avec 2 minutes de retard sur une fenêtre d'une heure. + expect( + captureDuration( + now: now, + endTime: now.add(const Duration(minutes: 58)), + ), + const Duration(minutes: 58), + ); + }); + + test('never asks ffmpeg for a zero or negative duration', () { + expect( + captureDuration( + now: now, + endTime: now.subtract(const Duration(minutes: 5)), + ), + const Duration(seconds: 30), + ); + }); + }); + + group('ffmpegRetryDelay', () { + test('backs off on repeated failures and caps at 30s', () { + expect(ffmpegRetryDelay(1), const Duration(seconds: 3)); + expect(ffmpegRetryDelay(2), const Duration(seconds: 6)); + expect(ffmpegRetryDelay(3), const Duration(seconds: 12)); + expect(ffmpegRetryDelay(4), const Duration(seconds: 24)); + expect(ffmpegRetryDelay(5), const Duration(seconds: 30)); + expect(ffmpegRetryDelay(20), const Duration(seconds: 30)); + }); + }); }