diff --git a/bin/api/upstream_probe.dart b/bin/api/upstream_probe.dart new file mode 100644 index 0000000..19e1476 --- /dev/null +++ b/bin/api/upstream_probe.dart @@ -0,0 +1,133 @@ +import 'dart:async'; +import 'dart:convert'; + +import 'package:http/http.dart' as http; + +import '../models/playlist_config.dart'; + +/// Diagnostic admin : chronomètre l'ouverture d'une chaîne chez le +/// fournisseur, en `.ts` et en `.m3u8`. +/// +/// POURQUOI : ~4,3 s des ~4,6 s d'un zap se passent avant le premier octet +/// du panneau (mesuré en prod). Le panneau autorise aussi le HLS ; reste à +/// savoir lequel des deux formats démarre le plus vite, et où part le temps +/// (redirection, en-têtes, premier octet). Les identifiants n'apparaissent +/// jamais dans la réponse : seuls l'hôte et les durées sont rendus. +Future> probeUpstream( + PlaylistConfig playlist, + String streamId, { + http.Client? client, +}) async { + final c = client ?? http.Client(); + try { + final base = '${playlist.dns}/live/${playlist.username}/' + '${playlist.password}/$streamId'; + final ts = await _probeStream(c, Uri.parse('$base.ts')); + // Laisser le panneau libérer la connexion (comptes à une connexion). + await Future.delayed(const Duration(seconds: 2)); + final hls = await _probeHls(c, Uri.parse('$base.m3u8')); + return {'stream': streamId, 'ts': ts, 'm3u8': hls}; + } finally { + if (client == null) c.close(); + } +} + +/// Suit les redirections à la main pour chronométrer chaque saut. +Future<({http.StreamedResponse? response, List> hops, Uri url})> + _open(http.Client c, Uri url, Stopwatch sw) async { + final hops = >[]; + var current = url; + for (var i = 0; i < 5; i++) { + final request = http.Request('GET', current) + ..followRedirects = false + ..headers['User-Agent'] = 'VLC/3.0.18 LibVLC/3.0.18'; + final response = + await c.send(request).timeout(const Duration(seconds: 20)); + hops.add({ + 'host': current.host, + 'status': response.statusCode, + 'headersMs': sw.elapsedMilliseconds, + }); + final location = response.headers['location']; + if (const {301, 302, 303, 307, 308}.contains(response.statusCode) && + location != null) { + await response.stream.drain().catchError((_) {}); + current = current.resolve(location); + continue; + } + return (response: response, hops: hops, url: current); + } + return (response: null, hops: hops, url: current); +} + +Future> _probeStream(http.Client c, Uri url) async { + final sw = Stopwatch()..start(); + try { + final opened = await _open(c, url, sw); + final response = opened.response; + if (response == null || response.statusCode != 200) { + return {'hops': opened.hops, 'error': 'status ${response?.statusCode}'}; + } + var bytes = 0; + int? firstByteMs; + final sub = response.stream.listen(null); + final done = Completer(); + sub.onData((chunk) { + firstByteMs ??= sw.elapsedMilliseconds; + bytes += chunk.length; + // 1 Mo suffit à voir le premier envoi (préchargement du panneau). + if (bytes > 1000000 && !done.isCompleted) done.complete(); + }); + sub.onDone(() { + if (!done.isCompleted) done.complete(); + }); + await done.future.timeout(const Duration(seconds: 15), onTimeout: () {}); + await sub.cancel(); + return { + 'hops': opened.hops, + 'firstByteMs': firstByteMs, + 'firstMegabyteMs': bytes > 1000000 ? sw.elapsedMilliseconds : null, + }; + } catch (e) { + return {'error': e.runtimeType.toString(), 'afterMs': sw.elapsedMilliseconds}; + } +} + +Future> _probeHls(http.Client c, Uri url) async { + final sw = Stopwatch()..start(); + try { + final opened = await _open(c, url, sw); + final response = opened.response; + if (response == null || response.statusCode != 200) { + return {'hops': opened.hops, 'error': 'status ${response?.statusCode}'}; + } + final body = await response.stream.bytesToString(); + final playlistMs = sw.elapsedMilliseconds; + final lines = const LineSplitter().convert(body); + final segments = + lines.where((l) => l.isNotEmpty && !l.startsWith('#')).toList(); + final target = lines + .firstWhere((l) => l.startsWith('#EXT-X-TARGETDURATION'), + orElse: () => '') + .split(':') + .last; + if (segments.isEmpty) { + return {'hops': opened.hops, 'playlistMs': playlistMs, 'segments': 0}; + } + // Un lecteur démarre près du direct : chronométrer l'avant-dernier. + final pick = segments.length >= 2 + ? segments[segments.length - 2] + : segments.last; + final segment = await _probeStream(c, opened.url.resolve(pick)); + return { + 'hops': opened.hops, + 'playlistMs': playlistMs, + 'segments': segments.length, + 'targetDuration': target, + 'segment': segment, + 'segmentStartedAtMs': playlistMs, + }; + } catch (e) { + return {'error': e.runtimeType.toString(), 'afterMs': sw.elapsedMilliseconds}; + } +} diff --git a/bin/server.dart b/bin/server.dart index 1a0112e..951fc17 100644 --- a/bin/server.dart +++ b/bin/server.dart @@ -28,6 +28,7 @@ import 'middleware/security_middleware.dart'; import 'services/cleanup_service.dart'; import 'services/recording_scheduler.dart'; import 'utils/asset_versioning.dart'; +import 'api/upstream_probe.dart'; void main(List args) async { // Parse command line arguments @@ -269,6 +270,29 @@ void main(List args) async { ); }); + // GET /api/admin/upstream-probe?stream= : chronomètre l'ouverture + // d'une chaîne chez le fournisseur en .ts et en .m3u8 (diagnostic du + // délai de zap ; aucun identifiant dans la réponse). + router.get('/upstream-probe', (Request req) async { + final user = req.context['user'] as User?; + if (user == null || !user.isAdmin) { + return Response.forbidden( + jsonEncode({'error': 'Admin access required'}), + ); + } + final streamId = req.url.queryParameters['stream'] ?? ''; + if (!RegExp(r'^[0-9]+$').hasMatch(streamId)) { + return Response.badRequest(body: 'stream= requis'); + } + final playlist = await getPlaylist(req); + if (playlist == null) return Response.forbidden('No playlist'); + final result = await probeUpstream(playlist, streamId); + return Response.ok( + jsonEncode(result), + headers: {'content-type': 'application/json'}, + ); + }); + return router(request); }), ); diff --git a/bin/test/upstream_probe_test.dart b/bin/test/upstream_probe_test.dart new file mode 100644 index 0000000..cef3566 --- /dev/null +++ b/bin/test/upstream_probe_test.dart @@ -0,0 +1,62 @@ +import 'dart:convert'; + +import 'package:http/http.dart' as http; +import 'package:http/testing.dart'; +import 'package:test/test.dart'; + +import '../api/upstream_probe.dart'; +import '../models/playlist_config.dart'; + +void main() { + final playlist = PlaylistConfig( + id: 'p', + name: 'n', + dns: 'http://panel.test', + username: 'secretuser', + password: 'secretpass', + createdAt: DateTime(2026), + ); + + // Panneau simulé : redirige vers un serveur de diffusion avec jeton. + final client = MockClient.streaming((request, _) async { + final path = request.url.path; + if (request.url.host == 'panel.test') { + return http.StreamedResponse( + const Stream.empty(), + 302, + headers: { + 'location': 'http://lb.test$path?token=secrettoken', + }, + ); + } + if (path.endsWith('.m3u8')) { + return http.StreamedResponse( + Stream.value(utf8.encode( + '#EXTM3U\n#EXT-X-TARGETDURATION:6\n#EXTINF:6,\nseg1.ts\n' + '#EXTINF:6,\nseg2.ts\n', + )), + 200, + ); + } + return http.StreamedResponse(Stream.value(List.filled(2000, 0x47)), 200); + }); + + test('chronomètre .ts et .m3u8 en suivant les redirections', () async { + final result = await probeUpstream(playlist, '123', client: client); + final ts = result['ts'] as Map; + expect((ts['hops'] as List).map((h) => (h as Map)['status']), [302, 200]); + expect(ts['firstByteMs'], isNotNull); + + final hls = result['m3u8'] as Map; + expect(hls['segments'], 2); + expect(hls['targetDuration'], '6'); + expect((hls['segment'] as Map)['firstByteMs'], isNotNull); + }); + + test('aucun identifiant ni jeton dans la réponse', () async { + final json = jsonEncode(await probeUpstream(playlist, '123', client: client)); + expect(json, isNot(contains('secretuser'))); + expect(json, isNot(contains('secretpass'))); + expect(json, isNot(contains('secrettoken'))); + }); +}