add: yt-dlp proxy by warp
This commit is contained in:
1 parent
a727be4ef3
commit
164a1da009
5 files changed
+318
-31
No files matched your search
+2
-1
@@ -21,4 +21,5 @@ RUN npm install
|
||||
COPY . .
|
||||
|
||||
EXPOSE 4000
|
||||
CMD ["npm", "run", "start"]
|
||||
# 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"]
|
||||
@@ -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:
|
||||
@@ -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")
|
||||
|
||||
+120
-30
@@ -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,15 +32,13 @@ 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 {
|
||||
@@ -43,24 +47,34 @@ function killCurrentProcess() {
|
||||
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 = [
|
||||
song.url,
|
||||
@@ -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;
|
||||
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();
|
||||
}
|
||||
out.on('close', () => {
|
||||
if (yt.exitCode === null && !yt.killed) yt.kill();
|
||||
});
|
||||
|
||||
yt.stdout.on('error', () => {
|
||||
if (!yt.killed && currentYtProcess === yt) yt.kill();
|
||||
if (!yt.killed) yt.kill();
|
||||
});
|
||||
|
||||
resolve(yt.stdout);
|
||||
} else {
|
||||
reject(new Error("Aucun flux stdout généré."));
|
||||
}
|
||||
yt.stdout.once('data', () => {
|
||||
if (settled) return;
|
||||
settled = true;
|
||||
clearTimeout(timeout);
|
||||
resolve({ stream: out });
|
||||
});
|
||||
|
||||
yt.stdout.pipe(out);
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* @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 };
|
||||
@@ -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 };
|
||||
Reference in new issue
Block a user