Files
xtremflow/bin/server.dart
T
MichaelandClaude Opus 5.5 4d00409de8 Guide TV : régionales France 3, variantes de nom, chaînes +1, CH et BE
Mesuré en prod après le premier correctif : 69 chaînes sur 149 avec
guide. Les vides restantes relevaient surtout de formes de nom que le
dump écrit autrement, ou de pays non couverts :

- « F3 ALPES » → France 3 Alpes (France.3.-.Alpes.fr) : F2 à F5 développés
- clé souple en dernier recours, dans son propre espace de noms :
  sans code pays ni « channel/tv/hd » (« AB3 » → AB3.Channel.fr)
- chaîne « +1 » absente du dump (TF1 +1) : guide de la chaîne principale
  décalé d'une heure ; une +1 qui a son propre guide le garde
- sources par défaut : FR1 puis CH1 (RTS…) et BE2 (même fournisseur),
  le dump français reste prioritaire

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-10-10 23:17:28 +02:00

494 lines
17 KiB
Dart

import 'dart:io';
import 'dart:async';
import 'dart:convert';
import 'package:shelf/shelf.dart';
import 'package:shelf/shelf_io.dart' as shelf_io;
import 'package:shelf_static/shelf_static.dart';
import 'package:shelf_router/shelf_router.dart';
import 'package:args/args.dart';
import 'database/database.dart';
import 'models/user.dart';
import 'models/playlist.dart';
import 'models/playlist_config.dart';
import 'api/auth_handler.dart';
import 'api/users_handler.dart';
import 'api/playlists_handler.dart';
import 'api/settings_handler.dart';
import 'api/streaming_handler.dart';
import 'api/proxy_handler.dart';
import 'api/recordings_api.dart';
import 'api/epg_api.dart';
import 'api/logo_api.dart';
import 'services/xmltv_epg_service.dart';
import 'services/logo_catalog.dart';
import 'api/season_passes_api.dart';
import 'api/xtream_api_handler.dart';
import 'middleware/auth_middleware.dart';
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<String> args) async {
// Parse command line arguments
final parser = ArgParser()
..addOption('port', abbr: 'p', defaultsTo: '8089')
..addOption('path', defaultsTo: '/app/web');
final result = parser.parse(args);
final port = int.parse(result['port']);
final webPath = result['path'];
// Initialize database
final db = AppDatabase();
await db.init();
await db.seedAdmin();
// 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);
recordingScheduler.start();
// Initialize Streaming Subsystem
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
Future<PlaylistConfig?> getPlaylist(Request request) async {
Playlist? playlist;
final user = request.context['user'] as User?;
// If we have a user from auth middleware, prefer their playlists
if (user != null) {
print('[getPlaylist] User from context: ${user.username}');
final playlists = db.getPlaylists(user.id);
if (playlists.isNotEmpty) playlist = playlists.first;
} else {
// Fallback for proxy: get first user's first playlist
print('[getPlaylist] No user in context, using fallback');
final users = db.getAllUsers();
print('[getPlaylist] Total users in DB: ${users.length}');
if (users.isNotEmpty) {
final playlists = db.getPlaylists(users[0].id);
print(
'[getPlaylist] User ${users[0].username} has ${playlists.length} playlists',
);
if (playlists.isNotEmpty) playlist = playlists.first;
}
}
if (playlist != null) {
print(
'[getPlaylist] Returning playlist: ${playlist.name} (DNS: ${playlist.serverUrl})',
);
return PlaylistConfig(
id: playlist.id,
name: playlist.name,
dns: playlist.serverUrl,
username: playlist.username,
password: playlist.password,
createdAt: playlist.createdAt,
isActive: true,
);
}
print('[getPlaylist] WARNING: No playlist found, returning null');
return null;
}
// Create API handlers
final authHandler = AuthHandler(db);
final playlistsHandler = PlaylistsHandler(db);
final usersHandler = UsersHandler(db);
final settingsHandler = SettingsHandler(db);
final proxyHandler = ProxyHandler(getPlaylist, db);
// Logos de chaînes : URL du panneau d'abord, puis repli par nom sur le
// dépôt public tv-logos quand l'hébergeur de picons du revendeur tombe.
final logoApi = LogoApi(proxyHandler.handler, LogoCatalog());
final recordingsApi = RecordingsApi(db, recordingScheduler);
// Source XMLTV de repli. Vider EPG_XMLTV_URLS désactive tout appel sortant :
// l'EPG se limite alors au panneau de l'abonné.
final xmltvUrls = (Platform.environment['EPG_XMLTV_URLS'] ??
'https://epgshare01.online/epgshare01/epg_ripper_FR1.xml.gz,'
'https://epgshare01.online/epgshare01/epg_ripper_CH1.xml.gz,'
'https://epgshare01.online/epgshare01/epg_ripper_BE2.xml.gz')
.split(',')
.map((u) => u.trim())
.where((u) => u.isNotEmpty)
.toList();
if (xmltvUrls.isEmpty) {
print('[EPG] Repli XMLTV désactivé (EPG_XMLTV_URLS vide)');
} else {
print('[EPG] Repli XMLTV : ${xmltvUrls.length} source(s)');
}
final epgApi = EpgApi(
getPlaylist,
xmltv: xmltvUrls.isEmpty ? null : XmltvEpgService(sourceUrls: xmltvUrls),
);
final seasonPassesApi = SeasonPassesApi(db);
final xtreamApiHandler = XtreamApiHandler(getPlaylist);
// TV Recordings sub-router (wrapped with auth below).
// NOTE: /api/recordings/stream/* falls through this router (404) and is
// handled by streamingRouter further down the Cascade.
final recordingsRouter = Router()
..get('/', recordingsApi.handleGetAll)
..post('/', recordingsApi.handlePost)
..delete('/<id>', recordingsApi.handleDelete)
..post('/stop/<id>', recordingsApi.handleStop)
..get('/logs/<id>', recordingsApi.getLogHandler);
final epgRouter = Router()..get('/<channelId>', epgApi.handleGetEpg);
final seasonPassesRouter = Router()
..get('/', seasonPassesApi.handleGetAll)
..post('/', seasonPassesApi.handlePost)
..delete('/<id>', seasonPassesApi.handleDelete);
// Setup router
final apiRouter = Router()
// Auth endpoints (with per-IP brute-force protection on login)
..mount(
'/api/auth',
const Pipeline()
.addMiddleware(loginRateLimitMiddleware())
.addHandler(authHandler.router.call),
)
// Xtream API gateway: injects credentials server-side so the
// frontend never sees them
..mount(
'/api/xtream-api',
const Pipeline()
.addMiddleware(authMiddleware(db))
.addHandler(xtreamApiHandler.handle),
)
// Playlists endpoints
..mount(
'/api/playlists',
const Pipeline()
.addMiddleware(authMiddleware(db))
.addHandler(playlistsHandler.router.call),
)
// Users endpoints
..mount(
'/api/users',
const Pipeline()
.addMiddleware(authMiddleware(db))
.addHandler(usersHandler.router.call),
)
// Settings endpoints
..mount(
'/api/settings',
const Pipeline()
.addMiddleware(authMiddleware(db))
.addHandler(settingsHandler.router.call),
)
// TV Recordings (auth required)
..mount(
'/api/recordings',
const Pipeline()
.addMiddleware(authMiddleware(db))
.addHandler(recordingsRouter.call),
)
// Logos de chaînes avec repli (auth required ; le cookie de session
// accompagne les requêtes d'image du navigateur)
..get(
'/api/logo',
const Pipeline()
.addMiddleware(authMiddleware(db))
.addHandler(logoApi.handle),
)
// EPG - guide TV (auth required)
..mount(
'/api/epg',
const Pipeline()
.addMiddleware(authMiddleware(db))
.addHandler(epgRouter.call),
)
// Season Passes - enregistrements répétés (auth required)
..mount(
'/api/season-passes',
const Pipeline()
.addMiddleware(authMiddleware(db))
.addHandler(seasonPassesRouter.call),
);
// NOTE: /api/xtream is handled by proxyHandler in the Cascade below
// Do NOT mount here as it would intercept and block the actual proxy
// Initialize Cleanup Service
// Ne JAMAIS cibler Directory.systemTemp en récursif : il contient les
// temporaires de la VM Dart et le dossier des sessions HLS — les fichiers
// de plus de 24 h y étaient supprimés aveuglément.
final cleanupService = CleanupService();
cleanupService.addTarget(Directory('/app/data/logs'));
cleanupService.addTarget(Directory('/app/data/tmp'));
cleanupService.start();
// Admin Routes (protected)
apiRouter.mount(
'/api/admin',
const Pipeline()
.addMiddleware(authMiddleware(db))
.addHandler((Request request) {
final router = Router();
// POST /api/admin/purge
router.post('/purge', (Request req) async {
final user = req.context['user'] as User?;
if (user == null || !user.isAdmin) {
return Response.forbidden(
jsonEncode({'error': 'Admin access required'}),
);
}
final result = await cleanupService.runCleanup();
return Response.ok(
jsonEncode(result),
headers: {'content-type': 'application/json'},
);
});
// GET /api/admin/upstream-probe?stream=<id> : 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=<id numérique> 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);
}),
);
// Les libs de lecture sont servies depuis l'image, plus depuis un CDN : si
// elles manquent (build web incomplet, copie Docker partielle), le lecteur
// échoue côté navigateur avec un simple « Échec du chargement » sans que
// rien ne l'ait signalé au démarrage. On le dit ici, une fois, clairement.
for (final lib in const ['hls.min.js', 'mpegts.min.js']) {
final file = File('$webPath/vendor/$lib');
if (!file.existsSync()) {
print(
'[Server] ATTENTION: web/vendor/$lib absent de $webPath — '
'la lecture échouera dans le navigateur.',
);
}
}
// Create static handler
final baseStaticHandler = createStaticHandler(
webPath,
defaultDocument: 'index.html',
listDirectories: false,
);
// Points d'entrée réécrits avec `?v=<empreinte>` sur leurs JS : voir
// asset_versioning.dart (un proxy qui cache les .js servait l'ancienne
// app et l'ancien lecteur après chaque déploiement).
final versionedEntryPoints = buildVersionedEntryPoints(webPath);
print('[Server] Assets versionnés : ${versionedEntryPoints.keys.join(', ')}');
// Wrap static handler to enforce cache policies
FutureOr<Response> staticHandler(Request request) async {
final path = request.url.path;
final versioned =
versionedEntryPoints[path.isEmpty ? 'index.html' : path];
if (versioned != null && request.method == 'GET') {
return Response.ok(
versioned,
headers: {
'Content-Type': path.endsWith('.js')
? 'text/javascript; charset=utf-8'
: 'text/html; charset=utf-8',
'Cache-Control': 'no-store, no-cache, must-revalidate, max-age=0',
'Pragma': 'no-cache',
'Expires': '0',
},
);
}
final response = await baseStaticHandler(request);
// Disable cache for entry points to ensure updates are seen immediately.
// Vendored player libraries (hls.js/mpegts.js, ~750 KB) are pinned
// versions: keep them cacheable or every player open re-downloads them.
final isVendored = path.startsWith('vendor/');
if (!isVendored &&
(path.isEmpty ||
path.endsWith('.html') ||
path.endsWith('.js') ||
path.endsWith('.json'))) {
return response.change(
headers: {
'Cache-Control': 'no-store, no-cache, must-revalidate, max-age=0',
'Pragma': 'no-cache',
'Expires': '0',
},
);
}
// Allow aggressive caching for hashed assets
return response.change(
headers: {
'Cache-Control': 'public, max-age=86400',
},
);
}
// Create Streaming Router (Mount handlers on correct paths)
final streamingRouter = Router()
..mount(
'/api/live',
createLiveStreamHandler(
getPlaylist,
isGpuEnabled: db.isNvidiaGpuEnabled,
),
)
..mount(
'/api/vod',
createVodStreamHandler(
getPlaylist,
isGpuEnabled: db.isNvidiaGpuEnabled,
),
)
..mount(
'/api/recordings/stream',
createRecordingStreamHandler(
db,
isGpuEnabled: db.isNvidiaGpuEnabled,
),
);
// Proxy Handler (auth is now handled INSIDE the handler, after path check)
// This allows non-/api/xtream requests to fall through to static handler
final proxyPipeline = proxyHandler.handler;
// Streaming routes accept the HttpOnly session cookie (hls.js cannot send
// Authorization headers); local FFmpeg loopback fetches bypass auth.
final protectedStreaming = const Pipeline()
.addMiddleware(streamAuthMiddleware(db))
.addHandler(streamingRouter.call);
// Main handler
final handler = Cascade()
.add(apiRouter.call) /* Standard API endpoints */
.add(proxyPipeline) /* Xtream Proxy (auth inside handler) */
.add(protectedStreaming) /* Streaming endpoints */
.add(staticHandler)
.handler;
// Add middleware
// redactedLogRequests remplace logRequests() : l'URI de /api/xtream/<url>
// contient username/password Xtream en clair.
final pipeline = const Pipeline()
.addMiddleware(redactedLogRequests())
.addMiddleware(securityHeadersMiddleware())
.addMiddleware(honeypotMiddleware())
.addMiddleware(rateLimitMiddleware())
.addMiddleware(_corsMiddleware())
.addHandler(handler);
// Start server
final server = await shelf_io.serve(
pipeline,
InternetAddress.anyIPv4,
port,
);
print('Server started on port ${server.port}');
print('Serving static files from: $webPath');
print('REST API available at: /api/auth/* and /api/playlists/*');
print('Xtream proxy available at: /api/xtream/* (SSRF Protected)');
// Clean expired sessions periodically (every hour)
Timer.periodic(const Duration(hours: 1), (_) {
db.cleanExpiredSessions();
print('Cleaned expired sessions');
});
}
/// CORS middleware.
///
/// The app is served same-origin by this server, so CORS headers are only
/// needed for development (Flutter dev server on another port) or an
/// explicitly configured external origin via the ALLOWED_ORIGIN env var.
/// The request Origin is echoed back only when it matches the allowlist —
/// never a wildcard.
Middleware _corsMiddleware() {
final extraOrigin = Platform.environment['ALLOWED_ORIGIN'];
bool isAllowed(String origin) {
if (extraOrigin != null && extraOrigin.isNotEmpty && origin == extraOrigin) {
return true;
}
// Local development origins (flutter run -d chrome, etc.)
final uri = Uri.tryParse(origin);
return uri != null && (uri.host == 'localhost' || uri.host == '127.0.0.1');
}
return (Handler handler) {
return (Request request) async {
final origin = request.headers['origin'];
final headers = <String, String>{
if (origin != null && isAllowed(origin)) ...{
'Access-Control-Allow-Origin': origin,
'Access-Control-Allow-Methods': 'GET, POST, PUT, DELETE, OPTIONS',
'Access-Control-Allow-Headers':
'Origin, Content-Type, Accept, Authorization, Range',
'Access-Control-Expose-Headers': 'Content-Length, Content-Range',
'Access-Control-Allow-Credentials': 'true',
'Vary': 'Origin',
},
};
// Handle preflight requests
if (request.method == 'OPTIONS') {
return Response.ok('', headers: headers);
}
// Process request and add CORS headers to response
final response = await handler(request);
return headers.isEmpty ? response : response.change(headers: headers);
};
};
}