diff --git a/Dockerfile b/Dockerfile index 0ce0367..a21e13e 100644 --- a/Dockerfile +++ b/Dockerfile @@ -21,4 +21,5 @@ RUN npm install COPY . . EXPOSE 4000 -CMD ["npm", "run", "start"] \ No newline at end of file +# Mise à jour de yt-dlp à chaque démarrage : une version obsolète provoque des erreurs 403 +CMD ["sh", "-c", "yt-dlp -U || true; npm run start"] \ No newline at end of file diff --git a/docker-compose.yml b/docker-compose.yml index 271caa7..9043093 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -1,6 +1,30 @@ version: "3.9" services: + # Proxy Cloudflare WARP (HTTP + SOCKS5 sur le port 1080) utilisé par yt-dlp + warp: + image: caomingjun/warp + container_name: subsonics-warp + restart: unless-stopped + # Requis pour que WARP puisse créer son interface TUN + device_cgroup_rules: + - "c 10:200 rwm" + cap_add: + - MKNOD + - AUDIT_WRITE + - NET_ADMIN + sysctls: + - net.ipv6.conf.all.disable_ipv6=0 + - net.ipv4.conf.all.src_valid_mark=1 + environment: + - WARP_SLEEP=2 + # - WARP_LICENSE_KEY= # optionnel : clé WARP+ + # Exposé uniquement en local (pour le dev hors Docker), jamais sur l'extérieur : ce serait un proxy ouvert + ports: + - "127.0.0.1:1080:1080" + volumes: + - subsonics-warp-data:/var/lib/cloudflare-warp + subsonics-backend: build: context: . # dossier backend @@ -8,9 +32,16 @@ services: container_name: subsonics-backend ports: - "4000:4000" + environment: + - WARP_PROXY=http://warp:1080 + # Si WARP tombe ou se fait bloquer, on retente en connexion directe + - WARP_FALLBACK_DIRECT=true + depends_on: + - warp volumes: - subsonics-backend-data:/app/data restart: unless-stopped volumes: subsonics-backend-data: + subsonics-warp-data: diff --git a/src/main.js b/src/main.js index 07d42c1..f3df33b 100644 --- a/src/main.js +++ b/src/main.js @@ -17,6 +17,7 @@ metric.publishMetrics("8001", "subsonicsMetricsRaph") setup(); async function setup() { + await require("./utils/WarpProxy").start() const DiscordBot = require("./discord/Bot") await DiscordBot.init() const Server = require("./server/Server") diff --git a/src/player/Method/Youtube.js b/src/player/Method/Youtube.js index 10cf952..a7183fb 100644 --- a/src/player/Method/Youtube.js +++ b/src/player/Method/Youtube.js @@ -2,22 +2,28 @@ const { LogType } = require('loguix'); const clog = new LogType("Youtube-Stream"); const { spawn, exec } = require('child_process'); const fs = require('fs'); +const { PassThrough } = require('stream'); const { __glob } = require('../../utils/GlobalVars'); +const warp = require('../../utils/WarpProxy'); + +// Erreurs liées à l'IP / la session : une autre route a des chances de passer +const RETRYABLE_ERRORS = /HTTP Error 403|403: Forbidden|HTTP Error 429|Sign in to confirm/i; +// Délai max avant de recevoir les premières données audio +const FIRST_DATA_TIMEOUT = 30 * 1000; + // Variable globale pour stocker le processus actif let currentYtProcess = null; +// Incrémenté à chaque appel de getStream, permet de savoir si un appel a été remplacé par un plus récent +let generation = 0; /** - * Tue le processus yt-dlp en cours proprement et attend sa fin réelle. - * Cela garantit qu'aucun flux ne se chevauche. + * Tue un processus yt-dlp et ses enfants (ffmpeg) */ -function killCurrentProcess() { +function killProcess(proc) { return new Promise((resolve) => { - if (!currentYtProcess || currentYtProcess.exitCode !== null) { - currentYtProcess = null; - return resolve(); - } + if (!proc || proc.exitCode !== null) return resolve(); - const pid = currentYtProcess.pid; + const pid = proc.pid; clog.log(`[YT-DLP] Nettoyage violent du processus PID: ${pid}`); // Détection de l'OS pour utiliser la bonne commande de kill @@ -26,50 +32,58 @@ function killCurrentProcess() { if (isWindows) { // Sur Windows, taskkill /T (Tree) /F (Force) est nécessaire pour tuer les enfants (ffmpeg) try { - exec(`taskkill /pid ${pid} /T /F`, (err) => { + exec(`taskkill /pid ${pid} /T /F`, () => { // Peu importe l'erreur (ex: processus déjà mort), on considère que c'est fini - currentYtProcess = null; resolve(); }); } catch (e) { // Fallback si taskkill échoue - try { currentYtProcess.kill('SIGKILL'); } catch (e2) {} - currentYtProcess = null; + try { proc.kill('SIGKILL'); } catch (e2) {} resolve(); } } else { // Sur Linux/Mac, on tente de tuer le groupe de processus try { - process.kill(-pid, 'SIGKILL'); + process.kill(-pid, 'SIGKILL'); } catch (e) { // Fallback si le group kill échoue - try { currentYtProcess.kill('SIGKILL'); } catch (e2) {} + try { proc.kill('SIGKILL'); } catch (e2) {} } - currentYtProcess = null; resolve(); } }); } /** - * @param {Object} song - L'objet contenant l'URL - * @param {number} seekTime - Le temps de démarrage en secondes (par défaut 0) + * Tue le processus yt-dlp en cours proprement et attend sa fin réelle. + * Cela garantit qu'aucun flux ne se chevauche. */ -async function getStream(song, seekTime = 0) { - // ÉTAPE 1 : On s'assure que l'ancien processus est BIEN mort avant de faire quoi que ce soit. - await killCurrentProcess(); +async function killCurrentProcess() { + const proc = currentYtProcess; + currentYtProcess = null; + await killProcess(proc); +} +function routeLabel(proxy) { + return proxy ? `WARP (${proxy})` : "directe"; +} + +/** + * Lance yt-dlp sur une route donnée et attend les premières données. + * Résout { stream } en cas de succès, ou { stream: null, reason, retryable } si la route a échoué avant le début du flux. + */ +function runAttempt(song, seekTime, proxy) { return new Promise((resolve, reject) => { - clog.log(`[YT-DLP] Lancement pour : ${song.url} (Début: ${seekTime}s)`); + clog.log(`[YT-DLP] Lancement pour : ${song.url} (Début: ${seekTime}s, Route: ${routeLabel(proxy)})`); - const ytArgs = [ + const ytArgs = [ song.url, '-o', '-', - + // Sélecteur large : Audio pur OU Vidéo (ffmpeg se débrouillera) - '-f', 'bestaudio/best', + '-f', 'bestaudio/best', '--buffer-size', '16K', - '--js-runtimes', 'node', + '--js-runtimes', 'node', ]; // --- GESTION DU TIMECODE (SEEK) --- @@ -83,13 +97,39 @@ async function getStream(song, seekTime = 0) { ytArgs.push('--no-cache-dir'); } + // --- GESTION DU PROXY WARP --- + if (proxy) { + ytArgs.push('--proxy', proxy); + } + // Lancement du nouveau processus const yt = spawn('yt-dlp', ytArgs); - + // On met à jour la variable globale tout de suite currentYtProcess = yt; + // Flux renvoyé au player : il ne reçoit rien tant qu'on ne sait pas si la route fonctionne + const out = new PassThrough(); let errorLogs = ""; + let settled = false; + + const fail = (reason) => { + if (settled) return; + settled = true; + clearTimeout(timeout); + yt.stdout.unpipe(out); + out.destroy(); + if (currentYtProcess === yt) currentYtProcess = null; + killProcess(yt); + resolve({ + stream: null, + reason, + // Via WARP, tout échec justifie un essai en direct ; en direct, seulement les erreurs d'IP / session + retryable: Boolean(proxy) || RETRYABLE_ERRORS.test(errorLogs) + }); + }; + + const timeout = setTimeout(() => fail(`aucune donnée après ${FIRST_DATA_TIMEOUT / 1000}s`), FIRST_DATA_TIMEOUT); yt.stderr.on('data', (data) => { const msg = data.toString(); @@ -97,13 +137,21 @@ async function getStream(song, seekTime = 0) { if (!msg.includes('[download]') && !msg.includes('[youtube]')) { errorLogs += msg; } + // On n'attend pas la fin du processus pour changer de route + if (!settled && RETRYABLE_ERRORS.test(errorLogs)) { + fail(errorLogs.trim().split('\n').pop()); + } }); yt.on('error', (err) => { clog.error("[YT-DLP] Erreur au lancement.", err); // Si c'est ce processus qui est en cours, on le clean if (currentYtProcess === yt) currentYtProcess = null; - reject(err); + clearTimeout(timeout); + if (!settled) { + settled = true; + reject(err); + } }); yt.on('close', (code) => { @@ -113,26 +161,68 @@ async function getStream(song, seekTime = 0) { if (code !== 0 && code !== null && code !== 143 && code !== 137) { // 137 = SIGKILL clog.warn(`[YT-DLP] Arrêt code ${code}. Logs : ${errorLogs}`); } + fail(`arrêt code ${code} avant toute donnée`); }); - if (yt.stdout) { - // --- SÉCURITÉ ANTI-OVERRIDE SUR LE FLUX --- - // Si le stream se ferme (le bot quitte le vocal ou skip), on tue yt-dlp - yt.stdout.on('close', () => { - if (!yt.killed && currentYtProcess === yt) { - yt.kill(); - } - }); - - yt.stdout.on('error', () => { - if (!yt.killed && currentYtProcess === yt) yt.kill(); - }); + // --- SÉCURITÉ ANTI-OVERRIDE SUR LE FLUX --- + // Si le stream se ferme (le bot quitte le vocal ou skip), on tue yt-dlp + out.on('close', () => { + if (yt.exitCode === null && !yt.killed) yt.kill(); + }); - resolve(yt.stdout); - } else { - reject(new Error("Aucun flux stdout généré.")); - } + yt.stdout.on('error', () => { + if (!yt.killed) yt.kill(); + }); + + yt.stdout.once('data', () => { + if (settled) return; + settled = true; + clearTimeout(timeout); + resolve({ stream: out }); + }); + + yt.stdout.pipe(out); }); } -module.exports = { getStream }; \ No newline at end of file +/** + * @param {Object} song - L'objet contenant l'URL + * @param {number} seekTime - Le temps de démarrage en secondes (par défaut 0) + */ +async function getStream(song, seekTime = 0) { + // Incrémenté avant le kill pour que l'appel précédent sache tout de suite qu'il est remplacé + const myGeneration = ++generation; + + // ÉTAPE 1 : On s'assure que l'ancien processus est BIEN mort avant de faire quoi que ce soit. + await killCurrentProcess(); + + // ÉTAPE 2 : On essaie chaque route (WARP puis directe, selon l'état du proxy) + const routes = warp.getRoutes(); + + for (let i = 0; i < routes.length; i++) { + // Un nouvel appel (skip, seek...) a pris la main : on n'envoie rien au player + if (myGeneration !== generation) return null; + + const proxy = routes[i]; + const result = await runAttempt(song, seekTime, proxy); + + if (result.stream) return result.stream; + if (myGeneration !== generation) return null; + + if (proxy) warp.reportFailure(); + + const next = routes[i + 1]; + if (!result.retryable || next === undefined) { + clog.error(`[YT-DLP] Échec via la route ${routeLabel(proxy)} : ${result.reason}`); + break; + } + clog.warn(`[YT-DLP] Échec via la route ${routeLabel(proxy)} : ${result.reason} - Nouvelle tentative via la route ${routeLabel(next)}`); + } + + // Flux vide : le player passe en Idle et enchaîne sur la musique suivante + const empty = new PassThrough(); + empty.end(); + return empty; +} + +module.exports = { getStream }; diff --git a/src/utils/WarpProxy.js b/src/utils/WarpProxy.js new file mode 100644 index 0000000..ea13f61 --- /dev/null +++ b/src/utils/WarpProxy.js @@ -0,0 +1,164 @@ +const { LogType } = require('loguix'); +const clog = new LogType("WARP-Proxy"); +const net = require('net'); +const fs = require('fs'); +const { __glob } = require('./GlobalVars'); + +const TRACE_URL = "https://www.cloudflare.com/cdn-cgi/trace"; +const CHECK_INTERVAL = 60 * 1000; +const CHECK_TIMEOUT = 8 * 1000; + +const state = { + url: null, + fallbackDirect: true, + healthy: false, + warp: null, + ip: null, + lastCheck: null, + lastError: null +}; + +let started = false; +let pendingCheck = null; + +/** + * Priorité : variable d'environnement WARP_PROXY, puis data/proxy.json. + * Format de proxy.json : { "enabled": true, "url": "http://127.0.0.1:1080", "fallbackDirect": true } + */ +function loadConfig() { + let fileConfig = {}; + if (fs.existsSync(__glob.PROXY)) { + try { + fileConfig = JSON.parse(fs.readFileSync(__glob.PROXY, "utf-8")); + } catch (e) { + clog.error(`Impossible de lire ${__glob.PROXY} : ${e.message}`); + } + } + + const url = process.env.WARP_PROXY || (fileConfig.enabled !== false ? fileConfig.url : null) || null; + const fallbackDirect = process.env.WARP_FALLBACK_DIRECT !== undefined + ? process.env.WARP_FALLBACK_DIRECT !== "false" + : fileConfig.fallbackDirect !== false; + + return { url, fallbackDirect }; +} + +/** + * Récupère la trace Cloudflare à travers le proxy (champ "warp" = on/plus si le trafic sort par WARP) + */ +async function fetchTrace(proxyUrl) { + const { request, ProxyAgent } = require('undici'); + const dispatcher = new ProxyAgent(proxyUrl); + try { + const { statusCode, body } = await request(TRACE_URL, { + dispatcher, + signal: AbortSignal.timeout(CHECK_TIMEOUT) + }); + const text = await body.text(); + if (statusCode !== 200) throw new Error(`Trace Cloudflare HTTP ${statusCode}`); + + const trace = {}; + for (const line of text.trim().split("\n")) { + const [key, ...value] = line.split("="); + trace[key] = value.join("="); + } + return trace; + } finally { + dispatcher.close().catch(() => {}); + } +} + +/** + * Pour les proxys SOCKS (non supportés par undici), on vérifie seulement que le port répond + */ +function tcpProbe(host, port) { + return new Promise((resolve, reject) => { + const socket = net.connect({ host, port: Number(port) }); + socket.setTimeout(CHECK_TIMEOUT); + socket.once('connect', () => { socket.destroy(); resolve(); }); + socket.once('timeout', () => { socket.destroy(); reject(new Error("Délai de connexion dépassé")); }); + socket.once('error', (err) => { socket.destroy(); reject(err); }); + }); +} + +async function runCheck() { + const wasHealthy = state.healthy; + const firstCheck = state.lastCheck === null; + + try { + const parsed = new URL(state.url); + if (parsed.protocol.startsWith("http")) { + const trace = await fetchTrace(state.url); + state.warp = trace.warp; + state.ip = trace.ip; + if (trace.warp !== "on" && trace.warp !== "plus") { + throw new Error(`Proxy joignable mais le trafic ne passe pas par WARP (warp=${trace.warp})`); + } + } else { + await tcpProbe(parsed.hostname, parsed.port); + } + state.healthy = true; + state.lastError = null; + } catch (e) { + state.healthy = false; + state.lastError = e.cause?.message || e.message; + } + state.lastCheck = Date.now(); + + if (state.healthy && (!wasHealthy || firstCheck)) { + clog.log(`Proxy WARP opérationnel (${state.url}) - IP de sortie : ${state.ip ?? "inconnue"}`); + } + if (!state.healthy && (wasHealthy || firstCheck)) { + const fallback = state.fallbackDirect ? "connexion directe utilisée en attendant" : "aucun repli configuré"; + clog.warn(`Proxy WARP indisponible (${state.url}) : ${state.lastError} - ${fallback}`); + } + + return state.healthy; +} + +function checkHealth() { + if (!state.url) return Promise.resolve(false); + if (!pendingCheck) { + pendingCheck = runCheck().finally(() => { pendingCheck = null; }); + } + return pendingCheck; +} + +async function start() { + if (started) return; + started = true; + + const config = loadConfig(); + state.url = config.url; + state.fallbackDirect = config.fallbackDirect; + + if (!state.url) { + clog.log("Aucun proxy WARP configuré (WARP_PROXY ou data/proxy.json), connexion directe"); + return; + } + + await checkHealth(); + setInterval(checkHealth, CHECK_INTERVAL).unref(); +} + +/** + * Ordre des routes à essayer pour yt-dlp. null = connexion directe. + */ +function getRoutes() { + if (!state.url) return [null]; + if (state.healthy) return state.fallbackDirect ? [state.url, null] : [state.url]; + return state.fallbackDirect ? [null] : [state.url]; +} + +/** + * Appelé quand une requête via le proxy échoue : on revérifie tout de suite au lieu d'attendre l'intervalle + */ +function reportFailure() { + checkHealth(); +} + +function getStatus() { + return { ...state }; +} + +module.exports = { start, getRoutes, reportFailure, checkHealth, getStatus };