From 6812a4964d606667fc1b73f572d1fb857611053b Mon Sep 17 00:00:00 2001 From: Michael Date: Sat, 10 Oct 2026 22:13:58 +0200 Subject: [PATCH] =?UTF-8?q?Lecture=20:=20ne=20couper=20que=20les=20flux=20?= =?UTF-8?q?orphelins,=20jamais=20un=20flux=20regard=C3=A9?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit La règle « le dernier flux gagne » coupait aussi un flux en cours de lecture. Deux lecteurs ouverts sur un compte à une connexion se coupaient alors en boucle (chacun relance aussitôt) : plus aucune image nulle part. Constaté en prod : tant qu'un autre onglet lisait, chaque ouverture échouait après 5 s sans un octet. Une connexion n'est désormais coupée que si son spectateur ne donne plus signe de vie depuis 15 s : aucune requête de playlist/segment (HLS, film) ni octet remis au client (turbo.ts). Les orphelins partent toujours, un second lecteur actif est refusé par le panneau comme avant. Co-Authored-By: Claude Opus 5.5 --- bin/api/streaming_handler.dart | 26 +++++++++++++++++++- bin/services/upstream_slots.dart | 41 ++++++++++++++++++++++++------- bin/test/upstream_slots_test.dart | 24 ++++++++++++++++++ 3 files changed, 81 insertions(+), 10 deletions(-) diff --git a/bin/api/streaming_handler.dart b/bin/api/streaming_handler.dart index cdc5cd7..20b5536 100644 --- a/bin/api/streaming_handler.dart +++ b/bin/api/streaming_handler.dart @@ -48,17 +48,31 @@ Future _claimUpstream(PlaylistConfig playlist, String id) async { await Future.delayed(const Duration(seconds: 1)); } +/// Délai sans consommation au-delà duquel un flux est tenu pour orphelin. +/// +/// Un lecteur HLS redemande sa playlist toutes les 2 à 4 s ; un flux +/// `turbo.ts` reçoit au moins une rafale toutes les ~7 s même sur un panneau +/// qui livre par à-coups. 15 s laissent une marge à ces deux rythmes. +const _viewerGrace = Duration(seconds: 15); + /// Inscrit [process] comme connexion amont de [playlist] tant qu'il vit. +/// +/// [lastActivity] : dernier signe de vie du spectateur ; la connexion n'est +/// coupée pour en ouvrir une autre que s'il remonte à plus de +/// [_viewerGrace]. void _registerUpstream( PlaylistConfig playlist, String id, Process process, { required void Function() release, + required DateTime Function() lastActivity, }) { final token = _upstreamSlots.register( id: id, account: accountKeyOf(playlist), release: release, + isActive: () => + DateTime.now().difference(lastActivity()) < _viewerGrace, ); process.exitCode.then((_) => _upstreamSlots.unregister(id, token: token)); } @@ -274,6 +288,8 @@ Future<_SessionAttempt> _runSession({ id, session.process, release: () => sessionManager.killSession(id), + // Touchée à chaque requête de playlist ou de segment. + lastActivity: () => session.lastAccess, ); } final outcome = @@ -580,11 +596,16 @@ Handler createLiveStreamHandler( return Response(503, body: 'FFmpeg start failed'); } + // Dernier octet remis au client : un lecteur en vie en reçoit à chaque + // rafale ; un client parti (zap mal détecté derrière le proxy) n'en + // reçoit plus, la contre-pression bloquant FFmpeg. + var lastDelivered = DateTime.now(); _registerUpstream( playlist, upstreamId, process, release: () => process.kill(ProcessSignal.sigterm), + lastActivity: () => lastDelivered, ); // Journaliser les erreurs FFmpeg (redactées) sans bloquer le flux. @@ -602,7 +623,10 @@ Handler createLiveStreamHandler( final body = pipeWithBackpressure( process.stdout, onStop: () => process.kill(ProcessSignal.sigterm), - ); + ).map((chunk) { + lastDelivered = DateTime.now(); + return chunk; + }); return Response( 200, diff --git a/bin/services/upstream_slots.dart b/bin/services/upstream_slots.dart index 148d062..4e5ddcc 100644 --- a/bin/services/upstream_slots.dart +++ b/bin/services/upstream_slots.dart @@ -15,10 +15,18 @@ import '../models/playlist_config.dart'; /// le panneau répondait `HTTP 551` ou ne répondait pas (timeout), d'où les /// films « indisponibles » et les zaps qui échouent. /// -/// Règle : le dernier flux demandé gagne. Avant d'ouvrir une nouvelle -/// connexion amont, on coupe les plus anciennes du même compte jusqu'à lui -/// faire de la place. Les enregistrements n'y sont pas inscrits : ils ne -/// sont jamais coupés. +/// Règle : avant d'ouvrir une nouvelle connexion amont, on coupe les +/// connexions ORPHELINES du même compte (plus personne ne les consomme), +/// les plus anciennes d'abord, jusqu'à lui faire de la place. +/// +/// Un flux encore regardé n'est jamais coupé. Une première version coupait +/// sans distinction (« le dernier gagne ») : deux lecteurs ouverts sur un +/// compte à une connexion se coupaient alors en boucle, chacun relançant +/// aussitôt — plus aucune image nulle part (constaté en prod). Le second +/// lecteur est désormais refusé par le panneau, comme avant, et le premier +/// continue. +/// +/// Les enregistrements n'y sont pas inscrits : ils ne sont jamais coupés. class UpstreamSlots { final DateTime Function() _now; final Map _slots = {}; @@ -31,13 +39,17 @@ class UpstreamSlots { /// Renvoie un jeton à repasser à [unregister] : une session relancée sous /// le même identifiant ne doit pas être désinscrite par la fin de /// l'ancien processus. + /// + /// [isActive] : vrai tant qu'un spectateur consomme le flux. Une + /// connexion active n'est jamais coupée par [makeRoom]. int register({ required String id, required String account, required void Function() release, + bool Function()? isActive, }) { final token = ++_nextToken; - _slots[id] = _Slot(token, account, _now(), release); + _slots[id] = _Slot(token, account, _now(), release, isActive); return token; } @@ -60,14 +72,16 @@ class UpstreamSlots { if (max <= 0) return const []; final others = _slots.entries .where((e) => e.value.account == account && e.key != keep) - .toList() - ..sort((a, b) => a.value.since.compareTo(b.value.since)); + .toList(); final excess = others.length - (max - 1); if (excess <= 0) return const []; + final orphans = others.where((e) => !e.value.active).toList() + ..sort((a, b) => a.value.since.compareTo(b.value.since)); + final freed = []; - for (final entry in others.take(excess)) { + for (final entry in orphans.take(excess)) { _slots.remove(entry.key); try { entry.value.release(); @@ -83,7 +97,16 @@ class _Slot { final String account; final DateTime since; final void Function() release; - _Slot(this.token, this.account, this.since, this.release); + final bool Function()? isActive; + _Slot(this.token, this.account, this.since, this.release, this.isActive); + + bool get active { + try { + return isActive?.call() ?? false; + } catch (_) { + return false; + } + } } /// Clé de compte : deux playlists sur les mêmes identifiants partagent le diff --git a/bin/test/upstream_slots_test.dart b/bin/test/upstream_slots_test.dart index 3bc5002..fd8dce4 100644 --- a/bin/test/upstream_slots_test.dart +++ b/bin/test/upstream_slots_test.dart @@ -53,6 +53,30 @@ void main() { expect(slots.countFor('acc'), 1); }); + test('un flux encore regardé n\'est jamais coupé (pas de ping-pong entre deux lecteurs)', () { + slots.register( + id: 'watched', + account: 'acc', + release: () => released.add('watched'), + isActive: () => true, + ); + expect(slots.makeRoom('acc', max: 1, keep: 'new'), isEmpty); + expect(released, isEmpty); + expect(slots.countFor('acc'), 1); + }); + + test('seuls les orphelins partent, même plus récents qu\'un flux regardé', () { + slots.register( + id: 'watched', + account: 'acc', + release: () => released.add('watched'), + isActive: () => true, + ); + add('orphan', 'acc'); + expect(slots.makeRoom('acc', max: 1, keep: 'new'), ['orphan']); + expect(released, ['orphan']); + }); + test('quota inconnu ou nul : rien n\'est coupé', () { add('a', 'acc'); expect(slots.makeRoom('acc', max: 0, keep: 'new'), isEmpty);