mirror of
https://github.com/R0m1k3/xtremflow.git
synced 2026-10-11 17:30:00 +02:00
Lecture : ne couper que les flux orphelins, jamais un flux regardé
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 <noreply@anthropic.com>
This commit is contained in:
1 parent
fdf98f9659
commit
6812a4964d
3 files changed
+81
-10
No files matched your search
@@ -48,17 +48,31 @@ Future<void> _claimUpstream(PlaylistConfig playlist, String id) async {
|
||||
await Future<void>.delayed(const Duration(seconds: 1));
|
||||
}
|
||||
|
||||
/// Délai sans consommation au-delà duquel un flux est tenu pour orphelin.
|
||||
///
|
||||
/// Un lecteur HLS redemande sa playlist toutes les 2 à 4 s ; un flux
|
||||
/// `turbo.ts` reçoit au moins une rafale toutes les ~7 s même sur un panneau
|
||||
/// qui livre par à-coups. 15 s laissent une marge à ces deux rythmes.
|
||||
const _viewerGrace = Duration(seconds: 15);
|
||||
|
||||
/// Inscrit [process] comme connexion amont de [playlist] tant qu'il vit.
|
||||
///
|
||||
/// [lastActivity] : dernier signe de vie du spectateur ; la connexion n'est
|
||||
/// coupée pour en ouvrir une autre que s'il remonte à plus de
|
||||
/// [_viewerGrace].
|
||||
void _registerUpstream(
|
||||
PlaylistConfig playlist,
|
||||
String id,
|
||||
Process process, {
|
||||
required void Function() release,
|
||||
required DateTime Function() lastActivity,
|
||||
}) {
|
||||
final token = _upstreamSlots.register(
|
||||
id: id,
|
||||
account: accountKeyOf(playlist),
|
||||
release: release,
|
||||
isActive: () =>
|
||||
DateTime.now().difference(lastActivity()) < _viewerGrace,
|
||||
);
|
||||
process.exitCode.then((_) => _upstreamSlots.unregister(id, token: token));
|
||||
}
|
||||
@@ -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,
|
||||
|
||||
@@ -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<String, _Slot> _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 = <String>[];
|
||||
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
|
||||
|
||||
@@ -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);
|
||||
|
||||
Reference in new issue
Block a user