Fiabilise le système d'enregistrement de bout en bout

Fuseaux horaires :
- POST /api/recordings rejette en 400 les dates sans fuseau et normalise
  tout en UTC à l'écriture (l'interprétation des dates naïves dans le TZ
  du conteneur décalait les enregistrements de 1-2 h)
- Helper frontend unique postRecording() : les 2 points de création
  (modal, guide EPG) envoient la même convention UTC

Contrôle d'accès :
- stop/delete/logs d'un enregistrement et delete d'un season pass
  vérifient la propriété (userId ou admin), comme playlists_handler
- Suppression des replis 'dev_user_id' et 'admin' (401 sans session)

SQLite :
- PRAGMA foreign_keys/WAL/busy_timeout (les ON DELETE CASCADE déclarés
  ne s'appliquaient pas : sessions et playlists orphelines)
- Migrations de schéma versionnées (schema_version) + index user_id,
  start_time, season_passes(user_id)

Gestion disque :
- Refus explicite de capture sous MIN_FREE_DISK_MB (défaut 500 Mo)
- Rotation par quota d'octets (RECORDINGS_QUOTA_GB, opt-in) qui ne touche
  jamais un enregistrement actif et supprime fichiers + ligne ensemble ;
  l'ancienne rotation « 50 fichiers » pouvait effacer une capture en cours
- DELETE /api/recordings/<id> supprime aussi .mkv/.log/parties (SafePath)

Scheduler :
- Arrêt gracieux orchestré par server.dart : clôture des enregistrements
  (fusion des parties, statut) avant killAll des sessions de streaming ;
  l'ancien handler SIGTERM de FfmpegSessionManager faisait exit(0) direct
- Noms de fichiers uniques par fragment d'id (deux enregistrements du
  même programme s'écrasaient mutuellement avec -y)
- Requête filtrée scheduled/recording au lieu de toute la table / 10 s
- Statut cancelled pour un scheduled arrêté (completed sans fichier
  cassait la lecture)

Season passes :
- Playlist du propriétaire du pass résolue à chaque scan (l'injection
  figée du 1er utilisateur rendait les passes muets après ajout de
  playlist, et mélangeait les credentials en multi-utilisateurs)
- Correspondance exacte par défaut (match_mode, migration en 'contains'
  pour l'existant), plafond de créations par scan, réalignement des
  horaires déplacés dans l'EPG, déduplication tolérante ±2 min
- Redaction des erreurs de scan (ClientException contient l'URL amont)

API de suivi (polling conservé) :
- GET /api/recordings enrichi : progress_pct, file_size_bytes,
  retry_count, is_active ; barre de progression + taille dans la liste

Validé : dart analyze (0 issue) et dart test (48/48) sur bin/.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_015oEu9QayWsw7hCKhenxgVa
This commit is contained in:
Claude committed 2026-08-27 12:11:56 +00:00
1 parent 60d3cc52d9
commit 1d8717bb40
10 files changed
+740 -187

No files matched your search

+12
View File
@@ -2,6 +2,18 @@
## Non publié ## Non publié
### 🔧 Fiabilité des enregistrements
- **Fuseaux horaires unifiés** : le backend exige des dates ISO-8601 avec fuseau (400 sinon) et stocke tout en UTC ; le frontend passe par un helper unique `postRecording()` — fini les enregistrements décalés de 1-2 h selon l'écran utilisé
- **Contrôle de propriété** : stop/suppression/logs d'un enregistrement et suppression d'un season pass ne sont plus possibles que par leur propriétaire (ou un admin)
- **SQLite durci** : `foreign_keys=ON` (les CASCADE déclarés s'appliquent enfin), WAL, `busy_timeout`, migrations de schéma versionnées, index sur `user_id`/`start_time`
- **Gestion disque** : refus explicite de démarrer une capture sous `MIN_FREE_DISK_MB` (défaut 500 Mo) ; nouvelle rotation par quota d'octets (`RECORDINGS_QUOTA_GB`, désactivée par défaut) qui ne touche jamais un enregistrement actif et supprime fichiers + ligne BDD ensemble (l'ancienne rotation « 50 fichiers » pouvait effacer une capture en cours) ; la suppression d'un enregistrement efface aussi ses fichiers (.mkv, .log, parties)
- **Arrêt gracieux** : `docker stop` clôture proprement les enregistrements (fusion des parties, statut en base) avant de tuer les sessions de streaming
- **Noms de fichiers uniques** (fragment d'id) : deux enregistrements du même programme ne s'écrasent plus
- **Statut `cancelled`** : arrêter un enregistrement planifié l'annule au lieu de le marquer « terminé » sans fichier (lecture cassée)
- **Season passes** : la playlist du propriétaire du pass est résolue à chaque scan (plus d'injection figée du premier utilisateur), correspondance de titre exacte par défaut (`match_mode`), plafond de créations par scan, réalignement automatique des horaires si le programme est déplacé dans l'EPG, déduplication tolérante (±2 min)
- **API de suivi** : `GET /api/recordings` renvoie désormais `progress_pct`, `file_size_bytes`, `retry_count`, `is_active` ; la liste affiche la barre de progression et la taille
- Le scheduler ne relit plus toute la table toutes les 10 s (requête filtrée sur `scheduled`/`recording`)
### 📺 Enregistrements ### 📺 Enregistrements
- La liste des enregistrements se met à jour automatiquement : rafraîchissement immédiat dès qu'un enregistrement est créé/arrêté n'importe où dans l'app (guide EPG, modal, widget rapide), et polling en arrière-plan (5 s quand un enregistrement est en cours ou planifié, 20 s sinon) pour suivre les statuts sans clic manuel - La liste des enregistrements se met à jour automatiquement : rafraîchissement immédiat dès qu'un enregistrement est créé/arrêté n'importe où dans l'app (guide EPG, modal, widget rapide), et polling en arrière-plan (5 s quand un enregistrement est en cours ou planifié, 20 s sinon) pour suivre les statuts sans clic manuel
- Indicateur « Suivi auto » avec heure de dernière actualisation dans l'onglet Enregistrements - Indicateur « Suivi auto » avec heure de dernière actualisation dans l'onglet Enregistrements
+168 -60
View File
@@ -3,6 +3,7 @@ import 'dart:convert';
import 'package:path/path.dart' as p; import 'package:path/path.dart' as p;
import 'package:shelf/shelf.dart'; import 'package:shelf/shelf.dart';
import '../database/database.dart'; import '../database/database.dart';
import '../models/recording.dart';
import '../models/user.dart'; import '../models/user.dart';
import '../services/recording_scheduler.dart'; import '../services/recording_scheduler.dart';
import '../utils/safe_path.dart'; import '../utils/safe_path.dart';
@@ -11,25 +12,43 @@ class RecordingsApi {
final AppDatabase _db; final AppDatabase _db;
final RecordingScheduler _scheduler; final RecordingScheduler _scheduler;
/// Durée maximale d'un enregistrement (env MAX_RECORDING_HOURS).
final int maxRecordingHours = int.tryParse(
Platform.environment['MAX_RECORDING_HOURS'] ?? '',
) ??
12;
RecordingsApi(this._db, this._scheduler); RecordingsApi(this._db, this._scheduler);
Response _json(int status, Map<String, dynamic> body) => Response(
status,
body: json.encode(body),
headers: {'Content-Type': 'application/json'},
);
/// L'utilisateur courant peut-il agir sur cet enregistrement ?
/// (même patron de contrôle de propriété que playlists_handler)
bool _canAccess(User? user, Recording recording) {
if (user == null) return false;
return user.isAdmin || recording.userId == user.id;
}
/// Handler pour GET /api/recordings/logs/<id> /// Handler pour GET /api/recordings/logs/<id>
/// Exposé séparément car shelf_router a un conflit entre DELETE /<id> et GET /logs/<id> /// Exposé séparément car shelf_router a un conflit entre DELETE /<id> et GET /logs/<id>
Future<Response> getLogHandler(Request request, String id) async { Future<Response> getLogHandler(Request request, String id) async {
final recording = _db.getRecordingById(id); final recording = _db.getRecordingById(id);
if (recording == null) { if (recording == null) {
return Response.notFound( return _json(404, {'error': 'Enregistrement non trouvé'});
json.encode({'error': 'Enregistrement non trouvé'}), }
headers: {'Content-Type': 'application/json'},
); final user = request.context['user'] as User?;
if (!_canAccess(user, recording)) {
return _json(403, {'error': 'Accès refusé'});
} }
if (recording.filePath == null) { if (recording.filePath == null) {
return Response.notFound( return _json(404, {'error': 'Aucun fichier ni log associé pour le moment.'});
json.encode({'error': 'Aucun fichier ni log associé pour le moment.'}),
headers: {'Content-Type': 'application/json'},
);
} }
// Les enregistrements sont écrits en .mkv avec un .log à côté // Les enregistrements sont écrits en .mkv avec un .log à côté
@@ -39,55 +58,115 @@ class RecordingsApi {
// Anti path-traversal : le log doit rester dans le dossier des enregistrements // Anti path-traversal : le log doit rester dans le dossier des enregistrements
final safeLogPath = SafePath.resolveWithin(recordingsDirPath, logFilePath); final safeLogPath = SafePath.resolveWithin(recordingsDirPath, logFilePath);
if (safeLogPath == null) { if (safeLogPath == null) {
return Response.forbidden( return _json(403, {'error': 'Chemin de log invalide'});
json.encode({'error': 'Chemin de log invalide'}),
headers: {'Content-Type': 'application/json'},
);
} }
final logFile = File(safeLogPath); final logFile = File(safeLogPath);
if (!await logFile.exists()) { if (!await logFile.exists()) {
return Response.notFound( return _json(404, {'error': 'Le fichier de log est introuvable.'});
json.encode({'error': 'Le fichier de log est introuvable.'}),
headers: {'Content-Type': 'application/json'},
);
} }
final logs = await logFile.readAsString(); final logs = await logFile.readAsString();
return Response.ok( return _json(200, {'logs': logs});
json.encode({'logs': logs}),
headers: {'Content-Type': 'application/json'},
);
} }
/// GET /api/recordings — Liste les enregistrements de l'utilisateur /// GET /api/recordings — Liste les enregistrements de l'utilisateur
/// (tous les enregistrements pour un admin) /// (tous les enregistrements pour un admin), enrichis des informations de
/// suivi : taille du fichier, progression, relances FFmpeg.
Response handleGetAll(Request request) { Response handleGetAll(Request request) {
final user = request.context['user'] as User?; final user = request.context['user'] as User?;
final recordings = (user != null && !user.isAdmin) final recordings = (user != null && !user.isAdmin)
? _db.getUserRecordings(user.id) ? _db.getUserRecordings(user.id)
: _db.getAllRecordings(); : _db.getAllRecordings();
final now = DateTime.now().toUtc();
return Response.ok( return Response.ok(
json.encode(recordings.map((r) => r.toMap()).toList()), json.encode(recordings.map((r) => _enrich(r, now)).toList()),
headers: {'Content-Type': 'application/json'}, headers: {'Content-Type': 'application/json'},
); );
} }
Map<String, dynamic> _enrich(Recording r, DateTime now) {
final map = r.toMap();
if (r.status == 'recording') {
final start = r.startTime.toUtc();
final end = r.endTime.toUtc();
final total = end.difference(start).inSeconds;
if (total > 0) {
final elapsed = now.difference(start).inSeconds;
map['progress_pct'] =
(elapsed * 100 / total).clamp(0, 100).round();
}
map['is_active'] = _scheduler.isCapturing(r.id);
final retries = _scheduler.retryCountOf(r.id);
if (retries != null) map['retry_count'] = retries;
}
final path = r.filePath;
if (path != null) {
try {
final file = File(path);
if (file.existsSync()) map['file_size_bytes'] = file.lengthSync();
} catch (_) {}
}
return map;
}
/// POST /api/recordings — Planifie un nouvel enregistrement /// POST /api/recordings — Planifie un nouvel enregistrement
Future<Response> handlePost(Request request) async { Future<Response> handlePost(Request request) async {
try { final user = request.context['user'] as User?;
final payload = await request.readAsString(); final userId = user?.id ?? request.context['userId'] as String?;
final data = json.decode(payload); if (userId == null) {
return _json(401, {'error': 'Authentification requise'});
}
final userId = request.context['userId'] as String? ?? 'dev_user_id'; Map<String, dynamic> data;
try {
data = json.decode(await request.readAsString()) as Map<String, dynamic>;
} catch (_) {
return _json(400, {'error': 'Corps JSON invalide'});
}
final channelId = data['channel_id']?.toString() ?? '';
final streamUrl = data['stream_url']?.toString() ?? '';
if (channelId.isEmpty) {
return _json(400, {'error': 'channel_id est requis'});
}
if (streamUrl.isEmpty) {
return _json(400, {'error': 'stream_url est requis'});
}
final startTime = _parseZonedDate(data['start_time']);
final endTime = _parseZonedDate(data['end_time']);
if (startTime == null || endTime == null) {
// Une date sans indicateur de fuseau ('Z' ou ±hh:mm) est ambiguë :
// l'interpréter dans le fuseau du serveur décale l'enregistrement
// de plusieurs heures selon le TZ du conteneur.
return _json(400, {
'error':
'start_time et end_time doivent être des dates ISO-8601 avec fuseau '
'(ex: 2026-08-27T21:00:00Z)',
});
}
if (!endTime.isAfter(startTime)) {
return _json(400, {'error': 'end_time doit être après start_time'});
}
if (endTime.difference(startTime) > Duration(hours: maxRecordingHours)) {
return _json(400, {
'error': 'Durée maximale dépassée ($maxRecordingHours h)',
});
}
try {
final recording = _db.createRecording( final recording = _db.createRecording(
userId: userId, userId: userId,
channelId: data['channel_id'], channelId: channelId,
streamUrl: data['stream_url'], streamUrl: streamUrl,
title: data['title'] ?? 'Sans Titre', title: data['title']?.toString() ?? 'Sans Titre',
startTime: DateTime.parse(data['start_time']), startTime: startTime,
endTime: DateTime.parse(data['end_time']), endTime: endTime,
); );
return Response.ok( return Response.ok(
@@ -95,58 +174,87 @@ class RecordingsApi {
headers: {'Content-Type': 'application/json'}, headers: {'Content-Type': 'application/json'},
); );
} catch (e) { } catch (e) {
return Response.internalServerError( print('[RecordingsApi] Erreur à la création: $e');
body: json.encode({'error': 'Erreur lors de la programmation: $e'}), return _json(500, {'error': 'Erreur lors de la programmation'});
headers: {'Content-Type': 'application/json'},
);
} }
} }
/// Parse une date ISO-8601 en exigeant un indicateur de fuseau, et la
/// normalise en UTC. Retourne null si absente, invalide ou naïve.
DateTime? _parseZonedDate(dynamic raw) {
final str = raw?.toString() ?? '';
if (str.isEmpty) return null;
// 'Z' final ou offset ±hh[:mm] après l'heure
final hasZone =
str.endsWith('Z') || RegExp(r'[+-]\d{2}:?\d{2}$').hasMatch(str);
if (!hasZone) return null;
return DateTime.tryParse(str)?.toUtc();
}
/// DELETE /api/recordings/<id> — Annule ou supprime un enregistrement /// DELETE /api/recordings/<id> — Annule ou supprime un enregistrement
/// Si un enregistrement FFmpeg est actif, il est arrêté avant la suppression /// Si un enregistrement FFmpeg est actif, il est arrêté avant la suppression.
/// Les fichiers associés (.mkv, .log, parties) sont supprimés avec la ligne.
Future<Response> handleDelete(Request request, String id) async { Future<Response> handleDelete(Request request, String id) async {
final recording = _db.getRecordingById(id); final recording = _db.getRecordingById(id);
if (recording == null) { if (recording == null) {
return Response.notFound( return _json(404, {'error': 'Enregistrement non trouvé'});
json.encode({'error': 'Enregistrement non trouvé'}), }
headers: {'Content-Type': 'application/json'},
); final user = request.context['user'] as User?;
if (!_canAccess(user, recording)) {
return _json(403, {'error': 'Accès refusé'});
} }
// Tuer FFmpeg si cet enregistrement est en cours AVANT de supprimer de la DB // Tuer FFmpeg si cet enregistrement est en cours AVANT de supprimer de la DB
await _scheduler.stopRecording(id); await _scheduler.stopRecording(id);
// Supprimer les fichiers pour ne pas laisser d'orphelins sur le volume,
// en restant confiné au dossier des enregistrements.
final path = recording.filePath;
if (path != null) {
final safePath = SafePath.resolveWithin(recordingsDirPath, path);
if (safePath != null) {
await _scheduler.deleteRecordingFiles(safePath);
}
}
_db.deleteRecording(id); _db.deleteRecording(id);
return Response.ok( return _json(200, {'message': 'Enregistrement supprimé avec succès'});
json.encode({'message': 'Enregistrement supprimé avec succès'}),
headers: {'Content-Type': 'application/json'},
);
} }
/// POST /api/recordings/stop/<id> — Arrête un enregistrement FFmpeg en cours /// POST /api/recordings/stop/<id> — Arrête un enregistrement FFmpeg en cours
Future<Response> handleStop(Request request, String id) async { Future<Response> handleStop(Request request, String id) async {
final recording = _db.getRecordingById(id); final recording = _db.getRecordingById(id);
if (recording == null) { if (recording == null) {
return Response.notFound( return _json(404, {'error': 'Enregistrement non trouvé'});
json.encode({'error': 'Enregistrement non trouvé'}),
headers: {'Content-Type': 'application/json'},
);
} }
final user = request.context['user'] as User?;
if (!_canAccess(user, recording)) {
return _json(403, {'error': 'Accès refusé'});
}
final stopped = await _scheduler.stopRecording(id); final stopped = await _scheduler.stopRecording(id);
if (stopped) { if (stopped) {
return Response.ok( return _json(200, {'message': 'Enregistrement arrêté'});
json.encode({'message': 'Enregistrement arrêté'}),
headers: {'Content-Type': 'application/json'},
);
} else {
// Pas de processus FFmpeg actif pour cet ID → marquer comme complété quand même
_db.updateRecordingStatus(id, 'completed');
return Response.ok(
json.encode({'message': 'Enregistrement marqué comme terminé'}),
headers: {'Content-Type': 'application/json'},
);
} }
if (recording.status == 'scheduled') {
// Rien n'a encore été capturé : annulé, pas « terminé ». Marquer
// completed sans fichier faisait ensuite échouer la lecture.
_db.updateRecordingStatus(id, 'cancelled');
return _json(200, {'message': 'Enregistrement annulé'});
}
if (recording.status == 'recording') {
// Statut « recording » sans processus actif (orphelin) : clôturer.
_db.updateRecordingStatus(id, 'completed');
return _json(200, {'message': 'Enregistrement marqué comme terminé'});
}
// Déjà completed/failed/cancelled : ne pas écraser le statut final.
return _json(200, {'message': 'Enregistrement déjà clôturé'});
} }
} }
+47 -8
View File
@@ -9,17 +9,29 @@ class SeasonPassesApi {
SeasonPassesApi(this._db); SeasonPassesApi(this._db);
/// GET /api/season-passes — liste tous les season passes /// GET /api/season-passes — liste les season passes de l'utilisateur
/// (tous les passes pour un admin)
Response handleGetAll(Request request) { Response handleGetAll(Request request) {
try { try {
final passes = _db.getAllSeasonPasses(); final user = request.context['user'] as User?;
if (user == null) {
return Response(
401,
body: json.encode({'error': 'Authentification requise'}),
headers: {'Content-Type': 'application/json'},
);
}
final passes = user.isAdmin
? _db.getAllSeasonPasses()
: _db.getSeasonPassesForUser(user.id);
return Response.ok( return Response.ok(
json.encode(passes), json.encode(passes),
headers: {'Content-Type': 'application/json'}, headers: {'Content-Type': 'application/json'},
); );
} catch (e) { } catch (e) {
print('[SeasonPass] Erreur au listage: $e');
return Response.internalServerError( return Response.internalServerError(
body: json.encode({'error': 'Erreur: $e'}), body: json.encode({'error': 'Erreur interne'}),
headers: {'Content-Type': 'application/json'}, headers: {'Content-Type': 'application/json'},
); );
} }
@@ -59,10 +71,20 @@ class SeasonPassesApi {
// Récupérer l'utilisateur depuis le contexte // Récupérer l'utilisateur depuis le contexte
final user = request.context['user'] as User?; final user = request.context['user'] as User?;
final userId = user?.id ?? 'admin'; // fallback if (user == null) {
return Response(
401,
body: json.encode({'error': 'Authentification requise'}),
headers: {'Content-Type': 'application/json'},
);
}
// Vérifier si un season pass identique existe déjà // 'exact' par défaut : « Journal » ne doit pas capturer tous les
final existing = _db.getAllSeasonPasses(); // programmes qui contiennent le mot. 'contains' reste disponible.
final matchMode = data['match_mode'] == 'contains' ? 'contains' : 'exact';
// Vérifier si un season pass identique existe déjà pour cet utilisateur
final existing = _db.getSeasonPassesForUser(user.id);
final duplicate = existing.any( final duplicate = existing.any(
(p) => (p) =>
(p['show_title'] as String).toLowerCase() == (p['show_title'] as String).toLowerCase() ==
@@ -78,10 +100,11 @@ class SeasonPassesApi {
} }
final pass = _db.createSeasonPass( final pass = _db.createSeasonPass(
userId: userId, userId: user.id,
showTitle: showTitle, showTitle: showTitle,
channelId: channelId, channelId: channelId,
streamUrl: streamUrl, streamUrl: streamUrl,
matchMode: matchMode,
); );
print('[SeasonPass] Créé: "$showTitle" sur chaîne $channelId'); print('[SeasonPass] Créé: "$showTitle" sur chaîne $channelId');
@@ -101,14 +124,30 @@ class SeasonPassesApi {
/// DELETE /api/season-passes/<id> — supprimer un season pass /// DELETE /api/season-passes/<id> — supprimer un season pass
Response handleDelete(Request request, String id) { Response handleDelete(Request request, String id) {
try { try {
final user = request.context['user'] as User?;
final pass = _db.getSeasonPassById(id);
if (pass == null) {
return Response.notFound(
json.encode({'error': 'Season Pass non trouvé'}),
headers: {'Content-Type': 'application/json'},
);
}
// Contrôle de propriété : seul le propriétaire ou un admin supprime.
if (user == null || (!user.isAdmin && pass['user_id'] != user.id)) {
return Response.forbidden(
json.encode({'error': 'Accès refusé'}),
headers: {'Content-Type': 'application/json'},
);
}
_db.deleteSeasonPass(id); _db.deleteSeasonPass(id);
return Response.ok( return Response.ok(
json.encode({'message': 'Season Pass supprimé'}), json.encode({'message': 'Season Pass supprimé'}),
headers: {'Content-Type': 'application/json'}, headers: {'Content-Type': 'application/json'},
); );
} catch (e) { } catch (e) {
print('[SeasonPass] Erreur à la suppression: $e');
return Response.internalServerError( return Response.internalServerError(
body: json.encode({'error': 'Erreur: $e'}), body: json.encode({'error': 'Erreur interne'}),
headers: {'Content-Type': 'application/json'}, headers: {'Content-Type': 'application/json'},
); );
} }
+165 -18
View File
@@ -23,10 +23,68 @@ class AppDatabase {
_db = sqlite3.open(dbPath); _db = sqlite3.open(dbPath);
// Sans foreign_keys, les ON DELETE CASCADE déclarés dans le schéma sont
// ignorés par SQLite : supprimer un utilisateur laissait ses sessions
// (donc des jetons valides), playlists et enregistrements orphelins.
_db.execute('PRAGMA foreign_keys = ON');
_db.execute('PRAGMA journal_mode = WAL');
_db.execute('PRAGMA busy_timeout = 5000');
await _createTables(); await _createTables();
_runMigrations();
print('Database initialized: $dbPath'); print('Database initialized: $dbPath');
} }
/// Migrations de schéma pour les bases créées par une version antérieure.
///
/// `CREATE TABLE IF NOT EXISTS` n'ajoute jamais de colonne à une table
/// existante : chaque colonne introduite après coup doit avoir sa migration.
/// Les migrations sont numérotées et rejouées uniquement si nécessaire.
void _runMigrations() {
_db.execute('''
CREATE TABLE IF NOT EXISTS schema_version (
version INTEGER PRIMARY KEY
)
''');
final result = _db.select(
'SELECT COALESCE(MAX(version), 0) AS v FROM schema_version',
);
var current = result.first['v'] as int;
final migrations = <int, void Function()>{
// v1 : colonne error_reason absente des bases d'avant son introduction.
1: () => _addColumnIfMissing('tv_recordings', 'error_reason', 'TEXT'),
// v2 : mode de correspondance des season passes. 'contains' pour les
// lignes existantes (comportement historique) ; les nouvelles créations
// passent par l'API qui choisit 'exact' par défaut.
2: () => _addColumnIfMissing(
'season_passes',
'match_mode',
"TEXT NOT NULL DEFAULT 'contains'",
),
};
for (final entry in migrations.entries) {
if (entry.key <= current) continue;
entry.value();
_db.execute(
'INSERT INTO schema_version (version) VALUES (?)',
[entry.key],
);
current = entry.key;
print('[DB] Migration v${entry.key} appliquée');
}
}
void _addColumnIfMissing(String table, String column, String definition) {
final columns = _db.select('PRAGMA table_info($table)');
final exists = columns.any((row) => row['name'] == column);
if (!exists) {
_db.execute('ALTER TABLE $table ADD COLUMN $column $definition');
}
}
/// Create database tables /// Create database tables
Future<void> _createTables() async { Future<void> _createTables() async {
// Users table // Users table
@@ -107,6 +165,7 @@ class AppDatabase {
channel_id TEXT NOT NULL, channel_id TEXT NOT NULL,
stream_url TEXT NOT NULL, stream_url TEXT NOT NULL,
enabled INTEGER DEFAULT 1, enabled INTEGER DEFAULT 1,
match_mode TEXT NOT NULL DEFAULT 'exact',
created_at TEXT DEFAULT CURRENT_TIMESTAMP, created_at TEXT DEFAULT CURRENT_TIMESTAMP,
FOREIGN KEY (user_id) REFERENCES users(id) ON DELETE CASCADE FOREIGN KEY (user_id) REFERENCES users(id) ON DELETE CASCADE
) )
@@ -125,6 +184,15 @@ class AppDatabase {
_db.execute( _db.execute(
'CREATE INDEX IF NOT EXISTS idx_recordings_status ON tv_recordings(status)', 'CREATE INDEX IF NOT EXISTS idx_recordings_status ON tv_recordings(status)',
); );
_db.execute(
'CREATE INDEX IF NOT EXISTS idx_recordings_user ON tv_recordings(user_id)',
);
_db.execute(
'CREATE INDEX IF NOT EXISTS idx_recordings_start ON tv_recordings(start_time)',
);
_db.execute(
'CREATE INDEX IF NOT EXISTS idx_season_passes_user ON season_passes(user_id)',
);
} }
/// Seed default admin user if no users exist /// Seed default admin user if no users exist
@@ -473,6 +541,11 @@ class AppDatabase {
final recordingId = _uuid.v4(); final recordingId = _uuid.v4();
final now = DateTime.now().toIso8601String(); final now = DateTime.now().toIso8601String();
// Toujours stocker en UTC avec suffixe 'Z' : les comparaisons du scheduler
// et la déduplication des season passes reposent sur ce format unique.
final startUtc = startTime.toUtc();
final endUtc = endTime.toUtc();
_db.execute( _db.execute(
''' '''
INSERT INTO tv_recordings (id, user_id, channel_id, stream_url, title, start_time, end_time, created_at, updated_at) INSERT INTO tv_recordings (id, user_id, channel_id, stream_url, title, start_time, end_time, created_at, updated_at)
@@ -484,8 +557,8 @@ class AppDatabase {
channelId, channelId,
streamUrl, streamUrl,
title, title,
startTime.toIso8601String(), startUtc.toIso8601String(),
endTime.toIso8601String(), endUtc.toIso8601String(),
now, now,
now, now,
], ],
@@ -497,21 +570,41 @@ class AppDatabase {
channelId: channelId, channelId: channelId,
streamUrl: streamUrl, streamUrl: streamUrl,
title: title, title: title,
startTime: startTime, startTime: startUtc,
endTime: endTime, endTime: endUtc,
status: 'scheduled', status: 'scheduled',
createdAt: DateTime.parse(now), createdAt: DateTime.parse(now),
updatedAt: DateTime.parse(now), updatedAt: DateTime.parse(now),
); );
} }
/// Lister tous les enregistrements (pour le Scheduler et l'admin) /// Lister tous les enregistrements (pour l'admin)
List<Recording> getAllRecordings() { List<Recording> getAllRecordings() {
final result = final result =
_db.select('SELECT * FROM tv_recordings ORDER BY start_time ASC'); _db.select('SELECT * FROM tv_recordings ORDER BY start_time ASC');
return result.map((row) => Recording.fromMap(row)).toList(); return result.map((row) => Recording.fromMap(row)).toList();
} }
/// Enregistrements qui intéressent le scheduler : à lancer ou en cours.
/// Évite de désérialiser tout l'historique toutes les 10 secondes.
List<Recording> getPendingRecordings() {
final result = _db.select(
"SELECT * FROM tv_recordings WHERE status IN ('scheduled', 'recording') "
'ORDER BY start_time ASC',
);
return result.map((row) => Recording.fromMap(row)).toList();
}
/// Enregistrements terminés (completed/failed/cancelled) du plus ancien au
/// plus récent — utilisé par la rotation disque.
List<Recording> getFinishedRecordingsOldestFirst() {
final result = _db.select(
"SELECT * FROM tv_recordings WHERE status IN ('completed', 'failed', 'cancelled') "
'ORDER BY start_time ASC',
);
return result.map((row) => Recording.fromMap(row)).toList();
}
/// Lister les enregistrements d'un utilisateur spécifique /// Lister les enregistrements d'un utilisateur spécifique
List<Recording> getUserRecordings(String userId) { List<Recording> getUserRecordings(String userId) {
final result = _db.select( final result = _db.select(
@@ -551,6 +644,24 @@ class AppDatabase {
); );
} }
/// Réaligner la fenêtre d'un enregistrement planifié (programme déplacé
/// dans l'EPG depuis sa création par un season pass).
void updateRecordingWindow(String id, DateTime start, DateTime end) {
_db.execute(
'''
UPDATE tv_recordings
SET start_time = ?, end_time = ?, updated_at = ?
WHERE id = ?
''',
[
start.toUtc().toIso8601String(),
end.toUtc().toIso8601String(),
DateTime.now().toIso8601String(),
id,
],
);
}
/// Supprimer un enregistrement depuis la BDD (ne supprime pas le fichier) /// Supprimer un enregistrement depuis la BDD (ne supprime pas le fichier)
void deleteRecording(String id) { void deleteRecording(String id) {
_db.execute('DELETE FROM tv_recordings WHERE id = ?', [id]); _db.execute('DELETE FROM tv_recordings WHERE id = ?', [id]);
@@ -564,12 +675,13 @@ class AppDatabase {
required String showTitle, required String showTitle,
required String channelId, required String channelId,
required String streamUrl, required String streamUrl,
String matchMode = 'exact',
}) { }) {
final id = _uuid.v4(); final id = _uuid.v4();
final now = DateTime.now().toIso8601String(); final now = DateTime.now().toIso8601String();
_db.execute( _db.execute(
'INSERT INTO season_passes (id, user_id, show_title, channel_id, stream_url, created_at) VALUES (?, ?, ?, ?, ?, ?)', 'INSERT INTO season_passes (id, user_id, show_title, channel_id, stream_url, match_mode, created_at) VALUES (?, ?, ?, ?, ?, ?, ?)',
[id, userId, showTitle, channelId, streamUrl, now], [id, userId, showTitle, channelId, streamUrl, matchMode, now],
); );
return { return {
'id': id, 'id': id,
@@ -577,12 +689,13 @@ class AppDatabase {
'show_title': showTitle, 'show_title': showTitle,
'channel_id': channelId, 'channel_id': channelId,
'stream_url': streamUrl, 'stream_url': streamUrl,
'match_mode': matchMode,
'enabled': 1, 'enabled': 1,
'created_at': now, 'created_at': now,
}; };
} }
/// Lister tous les Season Passes /// Lister les Season Passes actifs (pour le scheduler)
List<Map<String, dynamic>> getAllSeasonPasses() { List<Map<String, dynamic>> getAllSeasonPasses() {
final result = _db.select( final result = _db.select(
'SELECT * FROM season_passes WHERE enabled = 1 ORDER BY created_at DESC', 'SELECT * FROM season_passes WHERE enabled = 1 ORDER BY created_at DESC',
@@ -590,23 +703,57 @@ class AppDatabase {
return result.map((r) => Map<String, dynamic>.from(r)).toList(); return result.map((r) => Map<String, dynamic>.from(r)).toList();
} }
/// Lister les Season Passes d'un utilisateur (actifs ou non, pour l'API)
List<Map<String, dynamic>> getSeasonPassesForUser(String userId) {
final result = _db.select(
'SELECT * FROM season_passes WHERE user_id = ? ORDER BY created_at DESC',
[userId],
);
return result.map((r) => Map<String, dynamic>.from(r)).toList();
}
/// Récupérer un Season Pass par id (contrôle de propriété côté API)
Map<String, dynamic>? getSeasonPassById(String id) {
final result =
_db.select('SELECT * FROM season_passes WHERE id = ?', [id]);
if (result.isEmpty) return null;
return Map<String, dynamic>.from(result.first);
}
/// Supprimer un Season Pass /// Supprimer un Season Pass
void deleteSeasonPass(String id) { void deleteSeasonPass(String id) {
_db.execute('DELETE FROM season_passes WHERE id = ?', [id]); _db.execute('DELETE FROM season_passes WHERE id = ?', [id]);
} }
/// Vérifier si un enregistrement existe déjà pour ce titre (déduplication) /// Cherche un enregistrement existant pour cet épisode (déduplication des
/// Retourne true si un enregistrement non-échoué avec ce titre existe pour cetteémission programméeà la même heure /// season passes) : même titre, début à ±[tolerance] près.
bool existsRecordingForEpisode(String title, DateTime startTime) { ///
// Normaliser le titre pour la comparaison (insensible casse, sans espaces doubles) /// La comparaison par plage remplace l'ancienne égalité de chaîne, qui
/// échouait dès que le format stocké différait (avec/sans 'Z') ou que le
/// panneau décalait le programme de quelques secondes — l'épisode était
/// alors réenregistré en double.
Recording? findRecordingForEpisode(
String title,
DateTime startTime, {
Duration tolerance = const Duration(minutes: 2),
}) {
final startUtc = startTime.toUtc();
// Les dates sont stockées en ISO-8601 UTC : l'ordre lexicographique
// correspond à l'ordre chronologique, un BETWEEN sur chaînes suffit.
final result = _db.select( final result = _db.select(
'''SELECT COUNT(*) as cnt FROM tv_recordings '''SELECT * FROM tv_recordings
WHERE LOWER(title) = LOWER(?) WHERE LOWER(title) = LOWER(?)
AND start_time = ? AND start_time BETWEEN ? AND ?
AND status NOT IN ('failed')''', AND status NOT IN ('failed', 'cancelled')
[title, startTime.toUtc().toIso8601String()], LIMIT 1''',
[
title,
startUtc.subtract(tolerance).toIso8601String(),
startUtc.add(tolerance).toIso8601String(),
],
); );
return (result.first['cnt'] as int) > 0; if (result.isEmpty) return null;
return Recording.fromMap(result.first);
} }
/// Close database connection /// Close database connection
+25 -19
View File
@@ -42,31 +42,37 @@ void main(List<String> args) async {
await db.seedAdmin(); await db.seedAdmin();
// Initialize and start Recording Scheduler // Initialize and start Recording Scheduler
// (les Season Passes résolvent la playlist de leur propriétaire à chaque
// scan : plus d'injection figée du premier utilisateur au démarrage)
final recordingScheduler = RecordingScheduler(db); final recordingScheduler = RecordingScheduler(db);
recordingScheduler.start(); recordingScheduler.start();
// Injecter la config playlist dans le scheduler pour les Season Passes
// (on prend la playlist du premier utilisateur disponible)
Future<void> injectPlaylistToScheduler() async {
final users = db.getAllUsers();
if (users.isNotEmpty) {
final playlists = db.getPlaylists(users[0].id);
if (playlists.isNotEmpty) {
final p = playlists.first;
recordingScheduler.playlistDns = p.serverUrl;
recordingScheduler.playlistUsername = p.username;
recordingScheduler.playlistPassword = p.password;
print('[Server] Playlist injectée dans le scheduler: ${p.name}');
}
}
}
// Injecter après 5s pour attendre l'initialisation complète
Future.delayed(const Duration(seconds: 5), injectPlaylistToScheduler);
// Initialize Streaming Subsystem // Initialize Streaming Subsystem
await initStreaming(); await initStreaming();
// Arrêt gracieux unique (docker stop / Ctrl+C) : clôturer d'abord les
// enregistrements (kill FFmpeg, fusion des parties, statut en base), puis
// les sessions de streaming, puis sortir. Sans cela les enregistrements
// restaient au statut « recording » et la reprise d'orphelins devait
// systématiquement rattraper au redémarrage.
var shuttingDown = false;
Future<void> shutdownServer(String signal) async {
if (shuttingDown) return;
shuttingDown = true;
print('[Server] $signal reçu, arrêt en cours…');
try {
await recordingScheduler.shutdown();
} catch (e) {
print('[Server] Erreur à l\'arrêt du scheduler: $e');
}
sessionManager.killAll();
db.close();
exit(0);
}
ProcessSignal.sigterm.watch().listen((_) => shutdownServer('SIGTERM'));
ProcessSignal.sigint.watch().listen((_) => shutdownServer('SIGINT'));
// Helper to get playlist from request // Helper to get playlist from request
Future<PlaylistConfig?> getPlaylist(Request request) async { Future<PlaylistConfig?> getPlaylist(Request request) async {
Playlist? playlist; Playlist? playlist;
+4 -9
View File
@@ -59,15 +59,10 @@ class FfmpegSessionManager {
_reaper = Timer.periodic(const Duration(seconds: 60), (_) => _reap()); _reaper = Timer.periodic(const Duration(seconds: 60), (_) => _reap());
// Clean shutdown for docker stop / Ctrl+C // L'arrêt propre (docker stop / Ctrl+C) est orchestré par server.dart :
ProcessSignal.sigterm.watch().listen((_) { // il clôture d'abord les enregistrements puis appelle killAll(). Un
killAll(); // handler local qui ferait exit(0) immédiatement court-circuiterait
exit(0); // cette clôture.
});
ProcessSignal.sigint.watch().listen((_) {
killAll();
exit(0);
});
} }
FfmpegSession? get(String id) => _sessions[id]; FfmpegSession? get(String id) => _sessions[id];
+223 -47
View File
@@ -127,14 +127,34 @@ class RecordingScheduler {
) ?? ) ??
2; 2;
// Playlist config pour les appels EPG des season passes /// Espace libre minimal (Mo) exigé pour démarrer une capture.
// Rempli depuis server.dart après initialisation final int minFreeDiskMb = int.tryParse(
String? playlistDns; Platform.environment['MIN_FREE_DISK_MB'] ?? '',
String? playlistUsername; ) ??
String? playlistPassword; 500;
/// Quota du dossier d'enregistrements en Go (0 = rotation désactivée).
/// Remplace l'ancienne rotation « max 50 fichiers » qui supprimait
/// aveuglément, y compris des enregistrements en cours d'écriture.
final int recordingsQuotaGb = int.tryParse(
Platform.environment['RECORDINGS_QUOTA_GB'] ?? '',
) ??
0;
/// Nombre maximal d'enregistrements créés par season pass et par scan.
final int seasonPassMaxPerScan = int.tryParse(
Platform.environment['SEASON_PASS_MAX_PER_SCAN'] ?? '',
) ??
10;
RecordingScheduler(this._db); RecordingScheduler(this._db);
/// Un enregistrement est-il activement capturé par un processus FFmpeg ?
bool isCapturing(String id) => _active.containsKey(id);
/// Nombre de relances FFmpeg de l'enregistrement actif [id] (null si inactif).
int? retryCountOf(String id) => _active[id]?.consecutiveFailures;
void start() { void start() {
print( print(
'[RecordingScheduler] Démarrage du planificateur d\'enregistrements TV ' '[RecordingScheduler] Démarrage du planificateur d\'enregistrements TV '
@@ -163,6 +183,31 @@ class RecordingScheduler {
print('[RecordingScheduler] Arrêté'); print('[RecordingScheduler] Arrêté');
} }
/// Arrêt gracieux pour un `docker stop` : arrête les timers, clôt chaque
/// enregistrement actif (kill FFmpeg, fusion des parties, statut en base)
/// et attend la fin des clôtures dans la limite de [timeout].
///
/// Sans cette attente, le conteneur meurt avant la clôture : les
/// enregistrements restent au statut « recording » et la reprise d'orphelins
/// doit systématiquement rattraper au redémarrage.
Future<void> shutdown({Duration timeout = const Duration(seconds: 8)}) async {
_timer?.cancel();
_seasonPassTimer?.cancel();
final closings = <Future<void>>[];
for (final id in _active.keys.toList()) {
final future =
_stopActiveRecording(id, reason: 'Arrêt du serveur');
if (future != null) closings.add(future);
}
if (closings.isNotEmpty) {
print(
'[RecordingScheduler] Clôture de ${closings.length} enregistrement(s)…',
);
await Future.wait(closings).timeout(timeout, onTimeout: () => const []);
}
print('[RecordingScheduler] Arrêté proprement');
}
Future<void> _checkAndRunRecordings() async { Future<void> _checkAndRunRecordings() async {
if (_isRunning) return; if (_isRunning) return;
_isRunning = true; _isRunning = true;
@@ -185,7 +230,7 @@ class RecordingScheduler {
// Lire la base APRÈS les arrêts : un instantané pris avant ferait passer // 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 // l'enregistrement tout juste terminé pour un orphelin (il n'est plus
// dans `_active` alors que l'instantané le dit encore « recording »). // dans `_active` alors que l'instantané le dit encore « recording »).
final recordings = _db.getAllRecordings(); final recordings = _db.getPendingRecordings();
// Rechercher les enregistrements planifiés // Rechercher les enregistrements planifiés
for (final recording in recordings) { for (final recording in recordings) {
@@ -286,21 +331,22 @@ class RecordingScheduler {
} }
} }
/// Normalise un titre pour la correspondance : minuscules, espaces réduits.
static String _normalizeTitle(String title) =>
title.toLowerCase().trim().replaceAll(RegExp(r'\s+'), ' ');
/// Le titre EPG [title] correspond-il au pass selon son [matchMode] ?
static bool _titleMatches(String title, String showTitle, String matchMode) {
final t = _normalizeTitle(title);
final s = _normalizeTitle(showTitle);
if (s.isEmpty) return false;
return matchMode == 'contains' ? t.contains(s) : t == s;
}
Future<void> _checkSeasonPasses() async { Future<void> _checkSeasonPasses() async {
final passes = _db.getAllSeasonPasses(); final passes = _db.getAllSeasonPasses();
if (passes.isEmpty) return; if (passes.isEmpty) return;
final dns = playlistDns;
final username = playlistUsername;
final password = playlistPassword;
if (dns == null || username == null || password == null) {
print(
'[SeasonPass] Config playlist non disponible, vérification annulée',
);
return;
}
print('[SeasonPass] Vérification de ${passes.length} Season Pass(s)...'); print('[SeasonPass] Vérification de ${passes.length} Season Pass(s)...');
for (final pass in passes) { for (final pass in passes) {
@@ -309,9 +355,28 @@ class RecordingScheduler {
final streamUrl = pass['stream_url'] as String; final streamUrl = pass['stream_url'] as String;
final showTitle = pass['show_title'] as String; final showTitle = pass['show_title'] as String;
final userId = pass['user_id'] as String; final userId = pass['user_id'] as String;
final matchMode = pass['match_mode'] as String? ?? 'contains';
// Récupérer l'EPG de la chaîne (48 prochaines heures) // Résoudre la playlist du PROPRIÉTAIRE du pass au moment du scan.
final url = '$dns/player_api.php?username=$username&password=$password' // L'ancienne config injectée au démarrage venait du premier utilisateur
// de la base et n'était jamais réactualisée : une playlist ajoutée
// après coup rendait les passes muets, et en multi-utilisateurs les
// credentials de l'un servaient aux passes d'un autre.
final playlists = _db.getPlaylists(userId);
if (playlists.isEmpty) {
print(
'[SeasonPass] Aucune playlist pour le propriétaire du pass '
'"$showTitle", scan ignoré',
);
continue;
}
final playlist = playlists.first;
// Récupérer l'EPG de la chaîne (les prochains programmes ; `limit`
// est un nombre de programmes, pas des heures)
final url =
'${playlist.serverUrl}/player_api.php?username=${playlist.username}'
'&password=${playlist.password}'
'&action=get_simple_data_table&stream_id=$channelId&type=epg&limit=48'; '&action=get_simple_data_table&stream_id=$channelId&type=epg&limit=48';
final response = final response =
@@ -322,14 +387,14 @@ class RecordingScheduler {
final listings = final listings =
(raw is Map ? raw['epg_listings'] : raw) as List<dynamic>? ?? []; (raw is Map ? raw['epg_listings'] : raw) as List<dynamic>? ?? [];
var createdThisScan = 0;
for (final item in listings) { for (final item in listings) {
String title = item['title'] as String? ?? ''; String title = item['title'] as String? ?? '';
try { try {
title = utf8.decode(base64Decode(title)); title = utf8.decode(base64Decode(title));
} catch (_) {} } catch (_) {}
// Vérifier si le titre correspond au Season Pass (insensible casse, recherche partielle) if (!_titleMatches(title, showTitle, matchMode)) continue;
if (!title.toLowerCase().contains(showTitle.toLowerCase())) continue;
// Parser les heures de début/fin // Parser les heures de début/fin
final startStr = item['start'] as String? ?? ''; final startStr = item['start'] as String? ?? '';
@@ -348,12 +413,31 @@ class RecordingScheduler {
// Ne pas créer pour les programmes déjà terminés // Ne pas créer pour les programmes déjà terminés
if (endTime.isBefore(DateTime.now().toUtc())) continue; if (endTime.isBefore(DateTime.now().toUtc())) continue;
// Déduplication : vérifier si cet épisode est déjà planifié/enregistré // Déduplication : cet épisode est-il déjà planifié/enregistré ?
if (_db.existsRecordingForEpisode(title, startTime)) { final existing = _db.findRecordingForEpisode(title, startTime);
print('[SeasonPass] "$title" déjà enregistré, skip.'); if (existing != null) {
// Programme déplacé dans l'EPG depuis la planification :
// réaligner la fenêtre tant que la capture n'a pas commencé.
if (existing.status == 'scheduled' &&
(existing.startTime.toUtc() != startTime ||
existing.endTime.toUtc() != endTime)) {
_db.updateRecordingWindow(existing.id, startTime, endTime);
print(
'[SeasonPass] "$title" réaligné sur le nouvel horaire '
'${startTime.toLocal()}',
);
}
continue; continue;
} }
if (createdThisScan >= seasonPassMaxPerScan) {
print(
'[SeasonPass] Plafond de $seasonPassMaxPerScan créations atteint '
'pour "$showTitle" sur ce scan',
);
break;
}
// Créer l'enregistrement automatiquement // Créer l'enregistrement automatiquement
_db.createRecording( _db.createRecording(
userId: userId, userId: userId,
@@ -363,12 +447,18 @@ class RecordingScheduler {
startTime: startTime, startTime: startTime,
endTime: endTime, endTime: endTime,
); );
createdThisScan++;
print( print(
'[SeasonPass] ✓ Planifié automatiquement: "$title" le ${startTime.toLocal()}', '[SeasonPass] ✓ Planifié automatiquement: "$title" le ${startTime.toLocal()}',
); );
} }
} catch (e) { } catch (e) {
print('[SeasonPass] Erreur pour le pass "${pass['show_title']}": $e'); // Ne jamais imprimer l'exception brute : une ClientException peut
// contenir l'URL amont avec les credentials Xtream.
print(
'[SeasonPass] Erreur pour le pass "${pass['show_title']}": '
'${LogRedactor.redactUrl('$e')}',
);
} }
} }
} }
@@ -392,6 +482,24 @@ class RecordingScheduler {
// Nettoyer l'espace disque si nécessaire // Nettoyer l'espace disque si nécessaire
await _checkDiskSpaceAndRotate(recordingsDir); await _checkDiskSpaceAndRotate(recordingsDir);
// Refuser de démarrer sur un volume plein : mieux vaut un échec
// explicite immédiat qu'une capture qui meurt à mi-parcours.
final freeBytes = await _freeDiskBytes(recordingsDir.path);
if (freeBytes != null && freeBytes < minFreeDiskMb * 1024 * 1024) {
final freeMb = freeBytes ~/ (1024 * 1024);
print(
'[RecordingScheduler] Espace disque insuffisant ($freeMb Mo libres, '
'minimum $minFreeDiskMb Mo) : "${recording.title}" refusé',
);
_db.updateRecordingStatus(
recording.id,
'failed',
errorReason:
'Espace disque insuffisant ($freeMb Mo libres, minimum $minFreeDiskMb Mo)',
);
return;
}
final filePath = p.join(recordingsDir.path, _fileNameFor(recording)); final filePath = p.join(recordingsDir.path, _fileNameFor(recording));
final logPath = p.setExtension(filePath, '.log'); final logPath = p.setExtension(filePath, '.log');
@@ -673,9 +781,11 @@ class RecordingScheduler {
return false; // Pas d'enregistrement actif avec cet ID return false; // Pas d'enregistrement actif avec cet ID
} }
void _stopActiveRecording(String id, {String? reason}) { /// Arrête l'enregistrement actif [id] et retourne la future de clôture
/// (fusion des parties + statut), ou null si aucun n'est actif.
Future<void>? _stopActiveRecording(String id, {String? reason}) {
final active = _active.remove(id); final active = _active.remove(id);
if (active == null) return; if (active == null) return null;
active.stopping = true; active.stopping = true;
if (reason != null) { if (reason != null) {
print('[RecordingScheduler] Arrêt: $reason (${active.recording.title})'); print('[RecordingScheduler] Arrêt: $reason (${active.recording.title})');
@@ -685,7 +795,9 @@ class RecordingScheduler {
); );
active.process?.kill(ProcessSignal.sigterm); active.process?.kill(ProcessSignal.sigterm);
_db.updateRecordingStatus(id, 'completed'); _db.updateRecordingStatus(id, 'completed');
unawaited(_finalizeStopped(active)); final closing = _finalizeStopped(active);
unawaited(closing);
return closing;
} }
/// Attend la fin effective de FFmpeg puis clôture proprement (fusion + log). /// Attend la fin effective de FFmpeg puis clôture proprement (fusion + log).
@@ -709,7 +821,9 @@ class RecordingScheduler {
} }
String _fileNameFor(Recording recording) { String _fileNameFor(Recording recording) {
// Génération d'un nom de fichier unique et sûr // Génération d'un nom de fichier unique et sûr. Le fragment d'id garantit
// l'unicité : deux utilisateurs enregistrant le même programme sur la même
// chaîne s'écrasaient mutuellement (FFmpeg est lancé avec -y).
final safeTitle = final safeTitle =
recording.title.replaceAll(RegExp(r'[^a-zA-Z0-9_\-]'), '_'); recording.title.replaceAll(RegExp(r'[^a-zA-Z0-9_\-]'), '_');
final dateStr = recording.startTime final dateStr = recording.startTime
@@ -717,7 +831,10 @@ class RecordingScheduler {
.toIso8601String() .toIso8601String()
.replaceAll(':', '') .replaceAll(':', '')
.split('.')[0]; .split('.')[0];
return '${safeTitle}_$dateStr.mkv'; final idFragment = recording.id.length >= 8
? recording.id.substring(0, 8)
: recording.id;
return '${safeTitle}_${dateStr}_$idFragment.mkv';
} }
/// Chemin de la n-ième partie (la partie 1 étant le fichier principal). /// Chemin de la n-ième partie (la partie 1 étant le fichier principal).
@@ -755,28 +872,87 @@ class RecordingScheduler {
} catch (_) {} } catch (_) {}
} }
Future<void> _checkDiskSpaceAndRotate(Directory dir) async { /// Espace libre (octets) sur le volume qui porte [path], ou null si `df`
// Cette fonction pourrait invoquer une commande système `df` ou simplement lister les fichiers /// n'est pas disponible.
// et supprimer les plus anciens si un quota (ex: max 20 Go) est atteint. Future<int?> _freeDiskBytes(String path) async {
// Pour l'implémentation initiale, nous pouvons lister et supprimer si plus de X fichiers
try { try {
const maxFiles = 50; // Nombre max d'enregistrements (exemple simpliste) final result = await Process.run('df', ['-B1', '--output=avail', path]);
final files = dir.listSync().whereType<File>().toList(); if (result.exitCode != 0) return null;
final lines = (result.stdout as String).trim().split('\n');
return int.tryParse(lines.last.trim());
} catch (_) {
return null;
}
}
if (files.length > maxFiles) { /// Fichiers sur disque associés à un enregistrement : fichier principal,
print( /// log, et parties issues des relances.
'[RecordingScheduler] Rotation de l\'espace disque : suppression des anciens enregistrements', List<String> filesFor(String filePath) {
); final paths = <String>[filePath, p.setExtension(filePath, '.log')];
files.sort( for (var attempt = 1;; attempt++) {
(a, b) => a.statSync().modified.compareTo(b.statSync().modified), final part = _partPath(filePath, attempt);
); // Du plus vieux au plus récent if (!File(part).existsSync()) break;
paths.add(part);
}
return paths;
}
// Supprimer les plus anciens pour revenir sous la limite /// Supprime les fichiers d'un enregistrement (appelé par l'API à la
final filesToDelete = files.take(files.length - maxFiles); /// suppression, et par la rotation disque).
for (var file in filesToDelete) { Future<void> deleteRecordingFiles(String filePath) async {
file.deleteSync(); for (final path in filesFor(filePath)) {
await _deleteQuietly(path);
}
}
/// Rotation par quota d'octets (env RECORDINGS_QUOTA_GB, 0 = désactivée).
///
/// Remplace l'ancienne rotation « max 50 fichiers » qui supprimait les plus
/// anciens fichiers du dossier sans distinction : elle pouvait effacer une
/// partie en cours d'écriture par FFmpeg et laissait en base des lignes
/// pointant vers des fichiers disparus. Ici on ne supprime que des
/// enregistrements TERMINÉS connus de la base, du plus ancien au plus
/// récent, fichiers et ligne BDD ensemble, jamais un enregistrement actif.
Future<void> _checkDiskSpaceAndRotate(Directory dir) async {
if (recordingsQuotaGb <= 0) return;
try {
final quotaBytes = recordingsQuotaGb * 1024 * 1024 * 1024;
var totalBytes = 0;
await for (final entity in dir.list()) {
if (entity is File) {
try {
totalBytes += await entity.length();
} catch (_) {}
} }
} }
if (totalBytes <= quotaBytes) return;
print(
'[RecordingScheduler] Quota disque dépassé '
'(${totalBytes ~/ (1024 * 1024)} Mo > $recordingsQuotaGb Go) : '
'rotation des enregistrements terminés les plus anciens',
);
for (final old in _db.getFinishedRecordingsOldestFirst()) {
if (totalBytes <= quotaBytes) break;
final path = old.filePath;
if (path == null) continue;
// Jamais un enregistrement encore capturé (statut périmé en base).
if (_active.containsKey(old.id)) continue;
var freed = 0;
for (final f in filesFor(path)) {
try {
freed += File(f).existsSync() ? File(f).lengthSync() : 0;
} catch (_) {}
}
await deleteRecordingFiles(path);
_db.deleteRecording(old.id);
totalBytes -= freed;
print(
'[RecordingScheduler] Rotation : "${old.title}" supprimé '
'(${freed ~/ (1024 * 1024)} Mo libérés)',
);
}
} catch (e) { } catch (e) {
print( print(
'[RecordingScheduler] Erreur lors de la rotation de l\'espace disque : $e', '[RecordingScheduler] Erreur lors de la rotation de l\'espace disque : $e',
+42
View File
@@ -0,0 +1,42 @@
import 'dart:convert';
import 'package:http/http.dart' as http;
import 'authed_http.dart';
/// Point d'entrée unique pour créer un enregistrement via POST /api/recordings.
///
/// Le backend exige désormais des dates ISO-8601 AVEC fuseau (suffixe 'Z' ou
/// offset) et rejette les dates naïves en 400 : trois conventions d'envoi
/// coexistaient dans l'app (UTC, local naïf, UTC naïf), d'où des
/// enregistrements décalés de 1-2 h selon l'écran utilisé.
///
/// [wallClockIsUtc] : les horaires issus de l'EPG sont des heures UTC
/// « naïves » (parsées sans fuseau par Dart mais comparées à `now.toUtc()`
/// partout dans l'app). true les re-tague en UTC sans décalage ; false (par
/// défaut) convertit depuis l'heure locale réelle (cas d'un DateTime.now()).
Future<http.Response> postRecording({
required String channelId,
required String title,
required DateTime start,
required DateTime end,
String? streamUrl,
bool wallClockIsUtc = false,
}) {
DateTime asUtc(DateTime d) {
if (d.isUtc) return d;
return wallClockIsUtc
? DateTime.utc(d.year, d.month, d.day, d.hour, d.minute, d.second)
: d.toUtc();
}
return AuthedHttp.post(
Uri.parse('/api/recordings'),
headers: {'Content-Type': 'application/json'},
body: json.encode({
'channel_id': channelId,
'stream_url': streamUrl ?? '/api/live/$channelId.ts',
'title': title,
'start_time': asUtc(start).toIso8601String(),
'end_time': asUtc(end).toIso8601String(),
}),
);
}
+9 -15
View File
@@ -1,7 +1,6 @@
import 'dart:convert';
import 'package:flutter/material.dart'; import 'package:flutter/material.dart';
import 'package:google_fonts/google_fonts.dart'; import 'package:google_fonts/google_fonts.dart';
import '../../../core/api/authed_http.dart'; import '../../../core/api/recording_requests.dart';
import '../../../core/models/iptv_models.dart'; import '../../../core/models/iptv_models.dart';
import '../../../core/theme/app_colors.dart'; import '../../../core/theme/app_colors.dart';
import '../../../core/widgets/glass_container.dart'; import '../../../core/widgets/glass_container.dart';
@@ -43,20 +42,15 @@ class _RecordingModalState extends State<RecordingModal> {
setState(() => _isLoading = true); setState(() => _isLoading = true);
final endTime = _startTime.add(Duration(minutes: _durationMinutes)); final endTime = _startTime.add(Duration(minutes: _durationMinutes));
try { try {
// Utilisation d'une URL relative en Web (ou d'une configuration pour autres plateformes) // _startTime est une vraie heure locale (pickers) : postRecording la
final response = await AuthedHttp.post( // convertit en UTC — seule convention acceptée par le backend.
Uri.parse('/api/recordings'), final response = await postRecording(
headers: {'Content-Type': 'application/json'}, channelId: widget.channel.streamId,
body: json.encode({ title: widget.channel.name,
'channel_id': widget.channel.streamId, start: _startTime,
'stream_url': '/api/live/${widget.channel.streamId}.ts', end: endTime,
'title': widget.channel.name,
// Forcer UTC pour éviter le décalage +01:00 (France) vs UTC (serveur Docker)
'start_time': _startTime.toUtc().toIso8601String(),
'end_time': endTime.toUtc().toIso8601String(),
}),
); );
if (response.statusCode == 200) { if (response.statusCode == 200) {
+45 -11
View File
@@ -4,6 +4,7 @@ import 'package:flutter/material.dart';
import 'package:flutter_riverpod/flutter_riverpod.dart'; import 'package:flutter_riverpod/flutter_riverpod.dart';
import 'package:google_fonts/google_fonts.dart'; import 'package:google_fonts/google_fonts.dart';
import '../../../core/api/authed_http.dart'; import '../../../core/api/authed_http.dart';
import '../../../core/api/recording_requests.dart';
import '../../../core/models/iptv_models.dart'; import '../../../core/models/iptv_models.dart';
import '../../../core/models/playlist_config.dart'; import '../../../core/models/playlist_config.dart';
import '../../../core/theme/app_colors.dart'; import '../../../core/theme/app_colors.dart';
@@ -594,16 +595,15 @@ class _ProgrammeCard extends StatelessWidget {
DateTime end, DateTime end,
) async { ) async {
try { try {
final response = await AuthedHttp.post( // Les horaires EPG sont des heures UTC naïves : wallClockIsUtc les
Uri.parse('/api/recordings'), // re-tague sans décalage (l'ancien envoi naïf était interprété dans le
headers: {'Content-Type': 'application/json'}, // fuseau du serveur → enregistrement décalé de 1-2 h).
body: json.encode({ final response = await postRecording(
'channel_id': channel.streamId, channelId: channel.streamId,
'stream_url': '/api/live/${channel.streamId}.ts', title: title,
'title': title, start: start,
'start_time': start.toIso8601String(), end: end,
'end_time': end.toIso8601String(), wallClockIsUtc: true,
}),
); );
if (response.statusCode == 200) notifyRecordingsChanged(); if (response.statusCode == 200) notifyRecordingsChanged();
if (context.mounted) { if (context.mounted) {
@@ -914,6 +914,7 @@ class _RecordingsListViewState extends State<_RecordingsListView> {
'recording' => AppColors.live, 'recording' => AppColors.live,
'completed' => AppColors.success, 'completed' => AppColors.success,
'failed' => AppColors.warning, 'failed' => AppColors.warning,
'cancelled' => AppColors.onSurface38,
_ => AppColors.primaryContainer, _ => AppColors.primaryContainer,
}; };
@@ -922,9 +923,20 @@ class _RecordingsListViewState extends State<_RecordingsListView> {
'recording' => '● En cours', 'recording' => '● En cours',
'completed' => 'Terminé', 'completed' => 'Terminé',
'failed' => 'Échoué', 'failed' => 'Échoué',
'cancelled' => 'Annulé',
_ => status, _ => status,
}; };
String _fmtSize(int bytes) {
if (bytes >= 1024 * 1024 * 1024) {
return '${(bytes / (1024 * 1024 * 1024)).toStringAsFixed(1)} Go';
}
if (bytes >= 1024 * 1024) {
return '${(bytes / (1024 * 1024)).toStringAsFixed(0)} Mo';
}
return '${(bytes / 1024).toStringAsFixed(0)} Ko';
}
String _fmtDate(dynamic raw) { String _fmtDate(dynamic raw) {
if (raw == null) return '?'; if (raw == null) return '?';
try { try {
@@ -1040,12 +1052,34 @@ class _RecordingsListViewState extends State<_RecordingsListView> {
CrossAxisAlignment.start, CrossAxisAlignment.start,
children: [ children: [
Text( Text(
'${_fmtDate(rec['start_time'])} → ${_fmtDate(rec['end_time'])}', '${_fmtDate(rec['start_time'])} → ${_fmtDate(rec['end_time'])}'
'${rec['file_size_bytes'] is int ? ' · ${_fmtSize(rec['file_size_bytes'] as int)}' : ''}',
style: const TextStyle( style: const TextStyle(
color: AppColors.onSurface54, color: AppColors.onSurface54,
fontSize: 12, fontSize: 12,
), ),
), ),
if (status == 'recording' &&
rec['progress_pct'] is int) ...[
const SizedBox(height: 6),
ClipRRect(
borderRadius:
BorderRadius.circular(3),
child: LinearProgressIndicator(
value:
(rec['progress_pct'] as int) /
100,
minHeight: 4,
backgroundColor: AppColors
.onSurface
.withOpacity(0.1),
valueColor:
const AlwaysStoppedAnimation(
AppColors.live,
),
),
),
],
if (rec['error_reason'] != null) if (rec['error_reason'] != null)
Text( Text(
'⚠ ${rec['error_reason']}', '⚠ ${rec['error_reason']}',