Update index.js

This commit is contained in:
kq
2026-08-14 21:00:34 +02:00
parent 56e67eb151
commit 9a05ddb3e9
+168 -24
View File
@@ -30,12 +30,31 @@ const VOLUME = 0.3; // 0.0 - 1.0
const CACHE_DIR = path.join(__dirname, 'cache'); const CACHE_DIR = path.join(__dirname, 'cache');
if (!fs.existsSync(CACHE_DIR)) fs.mkdirSync(CACHE_DIR); if (!fs.existsSync(CACHE_DIR)) fs.mkdirSync(CACHE_DIR);
// On startup, remove any leftover .tmp files from downloads that never // A real encoded opus file is never this small. Anything under this size is
// finished (e.g. the process was killed mid-write). These are never valid // treated as a corrupt/truncated download rather than a usable cache entry —
// cache entries, so there's no reason to keep them around. // used both when finalizing a fresh download and when reading an existing
// cache file back off disk.
const MIN_VALID_CACHE_BYTES = 4096;
// Startup cleanup pass:
// - any leftover .tmp files are from downloads that never finished (e.g.
// the process was killed mid-write) — never valid, always safe to remove.
// - any .opus file smaller than MIN_VALID_CACHE_BYTES is a corrupt/truncated
// cache entry that could have been left behind by an older version of
// this bot (or a crash right at the rename boundary) — remove it so it
// gets cleanly redownloaded next time it's requested instead of being
// served as broken audio.
for (const f of fs.readdirSync(CACHE_DIR)) { for (const f of fs.readdirSync(CACHE_DIR)) {
const full = path.join(CACHE_DIR, f);
if (f.endsWith('.tmp')) { if (f.endsWith('.tmp')) {
try { fs.unlinkSync(path.join(CACHE_DIR, f)); } catch {} try { fs.unlinkSync(full); } catch {}
} else if (f.endsWith('.opus')) {
try {
if (fs.statSync(full).size < MIN_VALID_CACHE_BYTES) {
fs.unlinkSync(full);
console.log(`[Cache] Removed corrupt cache file on startup: ${f}`);
}
} catch {}
} }
} }
@@ -48,7 +67,16 @@ if (fs.existsSync(META_FILE)) {
} }
function saveMetaCache() { function saveMetaCache() {
fs.writeFileSync(META_FILE, JSON.stringify(metaCache, null, 2)); // Write via tmp+rename (same pattern as the audio cache) so a process
// kill mid-write can never leave metadata.json half-written/corrupt
// JSON that fails to parse on next boot.
const tmp = META_FILE + '.tmp';
try {
fs.writeFileSync(tmp, JSON.stringify(metaCache, null, 2));
fs.renameSync(tmp, META_FILE);
} catch (err) {
console.error('[Meta] Failed to save metadata cache:', err.message);
}
} }
function getAudioCachePath(url) { function getAudioCachePath(url) {
@@ -56,6 +84,24 @@ function getAudioCachePath(url) {
return path.join(CACHE_DIR, `${hash}.opus`); return path.join(CACHE_DIR, `${hash}.opus`);
} }
// A cache file only counts as usable if it exists AND is large enough to be
// real audio. Anything smaller gets deleted on sight so a corrupt entry is
// never served twice.
function isCacheFileValid(cachePath) {
try {
return fs.statSync(cachePath).size >= MIN_VALID_CACHE_BYTES;
} catch {
return false;
}
}
function evictIfInvalid(cachePath) {
if (fs.existsSync(cachePath) && !isCacheFileValid(cachePath)) {
console.warn(`[Cache] Discarding corrupt/empty cache file: ${path.basename(cachePath)}`);
try { fs.unlinkSync(cachePath); } catch {}
}
}
// ─── State (per guild) ──────────────────────────────────────────────────────── // ─── State (per guild) ────────────────────────────────────────────────────────
const guilds = new Map(); const guilds = new Map();
@@ -81,7 +127,23 @@ function getState(guildId) {
// cache is populated correctly for next time instead of ending up with // cache is populated correctly for next time instead of ending up with
// a truncated file. // a truncated file.
// - The cache file only ever appears at its final path once fully written // - The cache file only ever appears at its final path once fully written
// (write to .tmp, then atomic rename on success). // AND large enough to be real audio (write to .tmp, size-check, then
// atomic rename on success).
//
// IMPORTANT — multi-guild fan-out:
// The shared ffmpeg output is forwarded to the cache file and to every
// subscribed guild MANUALLY (not via stream.pipe()), and a subscriber's
// write() return value is ignored. This is deliberate: with .pipe(), Node
// ties the production rate of the single shared stream to whichever
// destination is slowest to drain. One guild with a stalled voice
// connection (or a paused player) would backpressure the shared download
// and stall the cache write *and every other guild listening to the same
// song*. Writing manually means a stuck subscriber can only ever hurt
// itself. Subscribers also swallow their own 'error'/'close' quietly —
// previously a subscriber being destroyed (e.g. on skip) while still
// attached to the shared stream could throw an unhandled 'error', which
// crashes the entire Node process and kills playback in every guild, not
// just the one that skipped.
const activeDownloads = new Map(); // url -> DownloadJob const activeDownloads = new Map(); // url -> DownloadJob
class DownloadJob { class DownloadJob {
@@ -90,7 +152,6 @@ class DownloadJob {
this.cachePath = getAudioCachePath(url); this.cachePath = getAudioCachePath(url);
this.tmpPath = this.cachePath + '.tmp'; this.tmpPath = this.cachePath + '.tmp';
this.subscribers = new Set(); this.subscribers = new Set();
this.source = new PassThrough(); // fan-out: ffmpeg output -> all subscribers + cache file
this.failed = false; this.failed = false;
this._start(); this._start();
} }
@@ -119,6 +180,15 @@ class DownloadJob {
ytdlp.stderr.on('data', (d) => console.error('[yt-dlp]', d.toString().trim())); ytdlp.stderr.on('data', (d) => console.error('[yt-dlp]', d.toString().trim()));
ffmpeg.stderr.on('data', (d) => console.error('[ffmpeg]', d.toString().trim())); ffmpeg.stderr.on('data', (d) => console.error('[ffmpeg]', d.toString().trim()));
// Child-process pipe streams can each emit their own 'error'
// (e.g. EPIPE when one side of the pipe dies before the other).
// Without listeners here that's an unhandled exception that
// crashes the whole bot — this is the fix for downloads breaking
// whenever multiple guilds/voice channels were involved.
ytdlp.stdout.on('error', () => {});
ffmpeg.stdin.on('error', () => {});
ffmpeg.stdout.on('error', () => {});
ytdlp.stdout.pipe(ffmpeg.stdin); ytdlp.stdout.pipe(ffmpeg.stdin);
ytdlp.on('error', (err) => this._fail(err)); ytdlp.on('error', (err) => this._fail(err));
@@ -133,9 +203,25 @@ class DownloadJob {
const cacheStream = fs.createWriteStream(this.tmpPath); const cacheStream = fs.createWriteStream(this.tmpPath);
this.cacheStream = cacheStream; this.cacheStream = cacheStream;
cacheStream.on('error', (err) => this._fail(err));
ffmpeg.stdout.pipe(this.source); // Manual fan-out — see class-level comment for why this replaces
this.source.pipe(cacheStream); // ffmpeg.stdout.pipe(...) to the cache file and every subscriber.
ffmpeg.stdout.on('data', (chunk) => {
if (this.failed) return;
if (!cacheStream.destroyed) cacheStream.write(chunk);
for (const sub of this.subscribers) {
if (!sub.destroyed) sub.write(chunk);
}
});
ffmpeg.stdout.on('end', () => {
if (this.failed) return;
if (!cacheStream.destroyed) cacheStream.end();
for (const sub of this.subscribers) {
if (!sub.destroyed) sub.end();
}
});
ffmpeg.on('close', (code) => { ffmpeg.on('close', (code) => {
if (code !== 0 && code !== null) { if (code !== 0 && code !== null) {
@@ -143,17 +229,34 @@ class DownloadJob {
} }
}); });
cacheStream.on('error', (err) => this._fail(err));
cacheStream.on('finish', () => { cacheStream.on('finish', () => {
if (this.failed) return; // already handled via _fail cleanup if (this.failed) return;
this._finalizeCache();
});
}
_finalizeCache() {
fs.stat(this.tmpPath, (err, stats) => {
activeDownloads.delete(this.url);
if (err) {
console.error('[Cache] Stat failed:', err.message);
fs.unlink(this.tmpPath, () => {});
return;
}
// Guard against corrupt/empty cache files — e.g. yt-dlp or
// ffmpeg exited "cleanly" but produced little/no real audio.
// Never let a file this small become a permanent cache entry.
if (stats.size < MIN_VALID_CACHE_BYTES) {
console.error(`[Cache] Discarding suspiciously small file (${stats.size}B): ${this.url}`);
fs.unlink(this.tmpPath, () => {});
return;
}
fs.rename(this.tmpPath, this.cachePath, (err) => { fs.rename(this.tmpPath, this.cachePath, (err) => {
if (err) { if (err) {
console.error('[Cache] Rename failed:', err.message); console.error('[Cache] Rename failed:', err.message);
} else { } else {
console.log(`[Cache] Saved: ${path.basename(this.cachePath)}`); console.log(`[Cache] Saved: ${path.basename(this.cachePath)}`);
} }
activeDownloads.delete(this.url);
}); });
}); });
} }
@@ -166,20 +269,27 @@ class DownloadJob {
try { this.ffmpeg.kill(); } catch {} try { this.ffmpeg.kill(); } catch {}
try { this.cacheStream.destroy(); } catch {} try { this.cacheStream.destroy(); } catch {}
fs.unlink(this.tmpPath, () => {}); // never leave a partial file behind fs.unlink(this.tmpPath, () => {}); // never leave a partial file behind
for (const sub of this.subscribers) sub.destroy(err); for (const sub of this.subscribers) {
if (!sub.destroyed) sub.destroy(err);
}
this.subscribers.clear();
activeDownloads.delete(this.url); activeDownloads.delete(this.url);
} }
// Each consumer (a guild's playback) gets its own tap on the shared // Each consumer (a guild's playback) gets its own independent tap on
// source. Closing this tap (e.g. because the guild skipped the track) // the shared download. Destroying this tap (e.g. because that guild
// only removes that listener — it does not affect the download itself. // skipped/stopped/left) must ONLY remove that one listener — it must
// never affect the shared child processes, the cache write, or any
// other guild's tap.
subscribe() { subscribe() {
const sub = new PassThrough(); const sub = new PassThrough();
this.subscribers.add(sub); this.subscribers.add(sub);
this.source.pipe(sub);
const cleanup = () => this.subscribers.delete(sub); const cleanup = () => this.subscribers.delete(sub);
sub.on('close', cleanup); sub.on('close', cleanup);
sub.on('end', cleanup); sub.on('end', cleanup);
// Swallow subscriber-side errors so a torn-down voice stream can
// never surface as an unhandled 'error' and crash the process.
sub.on('error', cleanup);
return sub; return sub;
} }
} }
@@ -223,6 +333,11 @@ client.on('messageCreate', async (message) => {
if (!query) return message.channel.send('❌ Podaj link lub frazę do wyszukania!'); if (!query) return message.channel.send('❌ Podaj link lub frazę do wyszukania!');
await cmdPlay(message, query); await cmdPlay(message, query);
} else if (command === 'fplay') {
const query = args.join(' ');
if (!query) return message.channel.send('❌ Podaj link lub frazę do wyszukania!');
await cmdPlay(message, query, { force: true });
} else if (command === 'skip') { } else if (command === 'skip') {
await cmdSkip(message); await cmdSkip(message);
@@ -239,6 +354,7 @@ client.on('messageCreate', async (message) => {
await message.channel.send( await message.channel.send(
'🎵 **Komendy:**\n' + '🎵 **Komendy:**\n' +
'`!play <url/szukaj>` – dodaj utwór do kolejki\n' + '`!play <url/szukaj>` – dodaj utwór do kolejki\n' +
'`!fplay <url/szukaj>` – wymuś świeże pobranie (ignoruje cache)\n' +
'`!skip` – pomiń aktualny utwór\n' + '`!skip` – pomiń aktualny utwór\n' +
'`!queue` – pokaż kolejkę\n' + '`!queue` – pokaż kolejkę\n' +
'`!stop` – zatrzymaj i wyczyść kolejkę\n' + '`!stop` – zatrzymaj i wyczyść kolejkę\n' +
@@ -253,7 +369,7 @@ client.on('messageCreate', async (message) => {
// ─── Commands ───────────────────────────────────────────────────────────────── // ─── Commands ─────────────────────────────────────────────────────────────────
async function cmdPlay(message, query) { async function cmdPlay(message, query, { force = false } = {}) {
const voiceChannel = message.member?.voice?.channel; const voiceChannel = message.member?.voice?.channel;
if (!voiceChannel) return message.channel.send('❌ Dołącz do kanału głosowego!'); if (!voiceChannel) return message.channel.send('❌ Dołącz do kanału głosowego!');
@@ -305,9 +421,14 @@ async function cmdPlay(message, query) {
} }
const cacheKey = query.toLowerCase().trim(); const cacheKey = query.toLowerCase().trim();
if (metaCache[cacheKey]) {
// Force play skips the metadata-cache shortcut entirely: always do a
// fresh yt-dlp lookup, and below, always wipe any existing audio cache
// for the resolved URL so playback can't come from a stale/corrupt file.
if (!force && metaCache[cacheKey]) {
const song = metaCache[cacheKey]; const song = metaCache[cacheKey];
const audioCached = fs.existsSync(getAudioCachePath(song.url)); evictIfInvalid(getAudioCachePath(song.url));
const audioCached = isCacheFileValid(getAudioCachePath(song.url));
console.log(`[Meta] HIT: "${query}" → ${song.url} (audio cached: ${audioCached})`); console.log(`[Meta] HIT: "${query}" → ${song.url} (audio cached: ${audioCached})`);
state.queue.push(song); state.queue.push(song);
await message.channel.send( await message.channel.send(
@@ -317,7 +438,7 @@ async function cmdPlay(message, query) {
return; return;
} }
await message.channel.send(`🔍 Szukam: **${query}**…`); await message.channel.send(force ? `🔁 Wymuszam ponowne pobranie: **${query}**…` : `🔍 Szukam: **${query}**…`);
const song = await resolveSong(query); const song = await resolveSong(query);
metaCache[cacheKey] = song; metaCache[cacheKey] = song;
@@ -325,6 +446,22 @@ async function cmdPlay(message, query) {
saveMetaCache(); saveMetaCache();
console.log(`[Meta] Cached: "${query}" → ${song.url}`); console.log(`[Meta] Cached: "${query}" → ${song.url}`);
if (force) {
const cachePath = getAudioCachePath(song.url);
if (fs.existsSync(cachePath)) {
try {
fs.unlinkSync(cachePath);
console.log(`[Cache] Force-cleared: ${path.basename(cachePath)}`);
} catch (err) {
console.error('[Cache] Force-clear failed:', err.message);
}
}
// If a download for this URL happens to already be in flight
// (e.g. another guild requested it a moment ago), that download is
// already fresh, not served from cache — buildAudioResource will
// simply join it like any other cache miss, so nothing more to do.
}
state.queue.push(song); state.queue.push(song);
await message.channel.send(`✅ Dodano do kolejki: **${song.title}** [${song.duration}]`); await message.channel.send(`✅ Dodano do kolejki: **${song.title}** [${song.duration}]`);
@@ -400,14 +537,21 @@ function playNext(guildId) {
function buildAudioResource(url) { function buildAudioResource(url) {
const cachePath = getAudioCachePath(url); const cachePath = getAudioCachePath(url);
// Cache HIT — instant play from disk. Because writes are atomic // Cache HIT — instant play from disk, but only if the file is actually
// (tmp file + rename), a file only exists here if it's complete. // large enough to be real audio. Because writes are atomic (tmp file +
// size-check + rename), a file should only ever exist here if it's
// complete and valid — this check is a last-line-of-defense safety net
// in case an older/corrupt file is still sitting on disk.
if (fs.existsSync(cachePath)) { if (fs.existsSync(cachePath)) {
if (isCacheFileValid(cachePath)) {
console.log(`[Cache] HIT: ${path.basename(cachePath)}`); console.log(`[Cache] HIT: ${path.basename(cachePath)}`);
return createAudioResource(fs.createReadStream(cachePath), { return createAudioResource(fs.createReadStream(cachePath), {
inputType: StreamType.OggOpus, inputType: StreamType.OggOpus,
}); });
} }
console.warn(`[Cache] Corrupt/empty cache file, discarding and redownloading: ${path.basename(cachePath)}`);
try { fs.unlinkSync(cachePath); } catch {}
}
// Cache MISS — join (or start) a shared download for this URL. // Cache MISS — join (or start) a shared download for this URL.
// Skipping/stopping playback later will NOT interrupt this download; // Skipping/stopping playback later will NOT interrupt this download;