From f798b43a1ec8f5450bc300086ccb04bda048ae66 Mon Sep 17 00:00:00 2001 From: Claude Date: Wed, 26 Aug 2026 19:09:31 +0000 Subject: [PATCH] =?UTF-8?q?fix(enregistrements):=20un=20enregistrement=20m?= =?UTF-8?q?en=C3=A9=20=C3=A0=20terme=20n'est=20plus=20marqu=C3=A9=20=C2=AB?= =?UTF-8?q?=20=C3=A9chou=C3=A9=20=C2=BB?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Chaque enregistrement arrivé au bout de sa fenêtre était marqué « Échoué / Interruption inattendue du serveur » alors que le fichier était bien sur le disque : le tick lisait la base AVANT d'arrêter les captures terminées, si bien que l'instantané annonçait encore « recording » pour un enregistrement déjà retiré de la table des processus actifs — la détection d'orphelin le requalifiait aussitôt en échec. La base est désormais lue après les arrêts, et la requalification revérifie le statut courant. Deux autres façons de perdre un enregistrement sont corrigées au passage : - FFmpeg livre encore des morceaux de stderr après la résolution de `exitCode` ; écrire sur l'`IOSink` déjà fermé levait une `StateError` depuis un callback de stream, donc une erreur asynchrone non rattrapée qui tue l'isolate — et avec lui le serveur et toutes les captures en cours. Le log passe par un écrivain tolérant et n'est fermé qu'une fois stdout et stderr drainés. - Une coupure amont terminait la capture définitivement. FFmpeg reçoit maintenant les options de reconnexion (comme le proxy live), et le planificateur relance la capture sur la fin de fenêtre quand le process sort trop tôt (backoff 3→30 s, quota remis à zéro après une capture saine). Les parties issues des relances sont recollées dans le fichier principal via le demuxer `concat`, la lecture reste donc un seul fichier. Également : - reprise des captures interrompues par un redémarrage du conteneur tant que la fenêtre est ouverte, au lieu d'un échec sec ; fichier partiel conservé (statut « terminé ») quand la fenêtre est passée - `-t` calculé sur le temps restant jusqu'à la fin programmée : un démarrage tardif ne rogne plus la fin du programme - `-hide_banner -nostats` : le log d'enregistrement redevient lisible (et ne pèse plus des mégaoctets, il est relu en entier par l'API) - `RECORDINGS_DIR` et `FFMPEG_PATH` surchargeables, et le motif d'erreur est effacé au (re)démarrage d'une capture Vérifié avec un faux ffmpeg sur un planificateur réel : avant correctif, un enregistrement mené jusqu'à la fin de fenêtre ressort « failed / Interruption inattendue du serveur » ; après, « completed » — de même que la reprise après interruption, la relance après coupure et la fusion des parties. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_01E5xWYTbtZwJuYB3E51K243 --- bin/api/recordings_api.dart | 4 +- bin/database/database.dart | 7 +- bin/services/recording_decision.dart | 93 ++++ bin/services/recording_scheduler.dart | 593 +++++++++++++++++++++----- bin/test/recording_decision_test.dart | 172 ++++++++ 5 files changed, 756 insertions(+), 113 deletions(-) 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)); + }); + }); }