diff --git a/src/audio/player.ts b/src/audio/player.ts index c3dbf5c..8940552 100644 --- a/src/audio/player.ts +++ b/src/audio/player.ts @@ -36,6 +36,10 @@ export class AudioPlayer extends EventEmitter { private currentUrl = ""; private seekOffset = 0; private framesPlayed = 0; // ground truth: number of 20ms frames sent + private sessionId = 0; + private static readonly BUFFER_HIGH_WATER = 960 * 1024; // ~5s of PCM at 48kHz stereo + private static readonly BUFFER_LOW_WATER = 480 * 1024; // ~2.5s + private ffmpegPaused = false; constructor(logger: Logger) { super(); @@ -45,9 +49,12 @@ export class AudioPlayer extends EventEmitter { play(url: string, seekSeconds = 0): void { this.stop(); + this.sessionId++; + const playSessionId = this.sessionId; this.currentUrl = url; this.seekOffset = seekSeconds; this.framesPlayed = 0; + this.ffmpegPaused = false; this.logger.info({ url: url.slice(0, 80), seek: seekSeconds }, "Starting playback"); @@ -72,19 +79,30 @@ export class AudioPlayer extends EventEmitter { this.logger.debug({ ffmpeg: ffmpegBin }, "Using ffmpeg binary"); this.ffmpeg = spawn(ffmpegBin, args, { stdio: ["ignore", "pipe", "pipe"] }); - this.ffmpeg.stderr!.on("data", () => {}); + this.ffmpeg.stderr!.on("data", (data: Buffer) => { + this.logger.debug({ stderr: data.toString().trimEnd() }, "FFmpeg stderr"); + }); this.ffmpeg.stdout!.on("data", (chunk: Buffer) => { this.pcmBuffer = Buffer.concat([this.pcmBuffer, chunk]); + // Backpressure: pause FFmpeg stdout when buffer is too large + if (this.pcmBuffer.length > AudioPlayer.BUFFER_HIGH_WATER && !this.ffmpegPaused && this.ffmpeg?.stdout) { + this.ffmpeg.stdout.pause(); + this.ffmpegPaused = true; + } }); this.ffmpeg.on("close", () => { - this.ffmpeg = null; // Signal frame loop that no more data is coming + if (this.sessionId === playSessionId) { + this.ffmpeg = null; // Signal frame loop that no more data is coming + } }); this.ffmpeg.on("error", (err) => { this.logger.error({ err }, "FFmpeg error"); - this.emit("error", err); + if (this.sessionId === playSessionId) { + this.emit("error", err); + } }); this.state = "playing"; @@ -101,11 +119,15 @@ export class AudioPlayer extends EventEmitter { private scheduleNextFrame(): void { if (!this.frameLoopRunning) return; + const loopSessionId = this.sessionId; + this.nextFrameTime += FRAME_DURATION_MS; const now = performance.now(); const delay = Math.max(0, this.nextFrameTime - now); setTimeout(() => { + // Discard callback from a stale play session + if (loopSessionId !== this.sessionId) return; if (!this.frameLoopRunning) return; if (this.state === "playing") { @@ -133,6 +155,12 @@ export class AudioPlayer extends EventEmitter { const pcmFrame = this.pcmBuffer.subarray(0, PCM_FRAME_BYTES); this.pcmBuffer = this.pcmBuffer.subarray(PCM_FRAME_BYTES); + // Backpressure: resume FFmpeg stdout when buffer drains below low-water mark + if (this.ffmpegPaused && this.pcmBuffer.length < AudioPlayer.BUFFER_LOW_WATER && this.ffmpeg?.stdout) { + this.ffmpeg.stdout.resume(); + this.ffmpegPaused = false; + } + const adjusted = this.applyVolume(pcmFrame); const opusFrame = this.encoder.encode(adjusted); this.emit("frame", opusFrame); @@ -184,12 +212,14 @@ export class AudioPlayer extends EventEmitter { } stop(): void { + this.sessionId++; this.frameLoopRunning = false; if (this.ffmpeg) { this.ffmpeg.kill("SIGTERM"); this.ffmpeg = null; } this.pcmBuffer = Buffer.alloc(0); + this.ffmpegPaused = false; this.state = "idle"; this.currentUrl = ""; this.seekOffset = 0; diff --git a/src/audio/queue.ts b/src/audio/queue.ts index dbe1ba9..5644b07 100644 --- a/src/audio/queue.ts +++ b/src/audio/queue.ts @@ -103,6 +103,8 @@ export class PlayQueue { if (this.songs.length === 0) return null; const prevIndex = this.currentIndex - 1; if (prevIndex < 0) { + // In Sequential mode, don't wrap around + if (this.mode === PlayMode.Sequential) return null; this.currentIndex = this.songs.length - 1; } else { this.currentIndex = prevIndex; diff --git a/src/bot/instance.ts b/src/bot/instance.ts index 1b902d4..6768432 100644 --- a/src/bot/instance.ts +++ b/src/bot/instance.ts @@ -54,6 +54,7 @@ export class BotInstance extends EventEmitter { private logger: Logger; private connected = false; private voteSkipUsers = new Set(); + private isAdvancing = false; constructor(options: BotInstanceOptions) { super(); @@ -142,9 +143,13 @@ export class BotInstance extends EventEmitter { } } catch (err) { this.logger.error({ err, command: parsed.name }, "Command execution error"); - await this.tsClient.sendTextMessage( - `Error: ${(err as Error).message}` - ); + try { + await this.tsClient.sendTextMessage( + `Error: ${(err as Error).message}` + ); + } catch (sendErr) { + this.logger.error({ err: sendErr }, "Failed to send error message to chat"); + } } } @@ -421,7 +426,9 @@ export class BotInstance extends EventEmitter { if (votes >= needed) { this.voteSkipUsers.clear(); - this.playNext(); + this.playNext().catch((err) => { + this.logger.error({ err }, "playNext failed after vote skip"); + }); return `Vote passed (${votes}/${needed}). Skipping to next song.`; } return `Vote to skip: ${votes}/${needed} (need ${needed - votes} more)`; @@ -472,23 +479,29 @@ export class BotInstance extends EventEmitter { } private async playNext(): Promise { - this.voteSkipUsers.clear(); - const next = this.queue.next(); - if (next) { - const ok = await this.resolveAndPlay(next); - if (!ok) { - // Skip to next if URL resolve fails (up to 3 retries) - for (let i = 0; i < 3; i++) { - const retry = this.queue.next(); - if (!retry) break; - if (await this.resolveAndPlay(retry)) return; + if (this.isAdvancing) return; + this.isAdvancing = true; + try { + this.voteSkipUsers.clear(); + const next = this.queue.next(); + if (next) { + const ok = await this.resolveAndPlay(next); + if (!ok) { + // Skip to next if URL resolve fails (up to 3 retries) + for (let i = 0; i < 3; i++) { + const retry = this.queue.next(); + if (!retry) break; + if (await this.resolveAndPlay(retry)) return; + } + this.player.stop(); } + } else { this.player.stop(); } - } else { - this.player.stop(); + this.emit("stateChange"); + } finally { + this.isAdvancing = false; } - this.emit("stateChange"); } private extractId(input: string): string { diff --git a/src/ts-protocol/client.ts b/src/ts-protocol/client.ts index 503c05b..e922f13 100644 --- a/src/ts-protocol/client.ts +++ b/src/ts-protocol/client.ts @@ -39,6 +39,7 @@ export class TS3Client extends EventEmitter { private identity: Identity; private clientId = 0; private logger: Logger; + private disconnecting = false; constructor(private options: TS3ClientOptions, logger: Logger) { super(); @@ -167,9 +168,15 @@ export class TS3Client extends EventEmitter { } disconnect(): void { - if (this.client) { - this.client.disconnect().catch(() => {}); - this.client = null; + if (this.client && !this.disconnecting) { + this.disconnecting = true; + const client = this.client; + client.disconnect().catch(() => {}).finally(() => { + if (this.client === client) { + this.client = null; + } + this.disconnecting = false; + }); } this.clientId = 0; this.logger.info("Disconnected from TeamSpeak server"); diff --git a/src/web/api/player.ts b/src/web/api/player.ts index 9a7868d..11f164f 100644 --- a/src/web/api/player.ts +++ b/src/web/api/player.ts @@ -211,7 +211,11 @@ export function createPlayerRouter( try { const bot = (req as any).bot; const { songId, platform } = req.body; - const provider = (platform === "qq" ? qqProvider : neteaseProvider)!; + const provider = platform === "qq" ? qqProvider : neteaseProvider; + if (!provider) { + res.status(500).json({ error: "Provider not available" }); + return; + } const song = await provider.getSongDetail(songId); if (!song) { @@ -241,7 +245,11 @@ export function createPlayerRouter( try { const bot = (req as any).bot; const { songId, platform } = req.body; - const provider = (platform === "qq" ? qqProvider : neteaseProvider)!; + const provider = platform === "qq" ? qqProvider : neteaseProvider; + if (!provider) { + res.status(500).json({ error: "Provider not available" }); + return; + } const song = await provider.getSongDetail(songId); if (!song) { diff --git a/src/web/websocket.ts b/src/web/websocket.ts index af8934f..ec00d5e 100644 --- a/src/web/websocket.ts +++ b/src/web/websocket.ts @@ -1,5 +1,6 @@ import { WebSocketServer, WebSocket } from "ws"; import type { BotManager } from "../bot/manager.js"; +import type { BotInstance } from "../bot/instance.js"; import type { Logger } from "../logger.js"; export function setupWebSocket( @@ -9,6 +10,13 @@ export function setupWebSocket( ): () => void { const clients = new Set(); + /** Track which bots already have listeners attached */ + const attachedBots = new Map void; + connected: () => void; + disconnected: () => void; + }>(); + wss.on("connection", (ws) => { clients.add(ws); logger.debug("WebSocket client connected"); @@ -36,39 +44,63 @@ export function setupWebSocket( } }; - const attachBotListeners = () => { + function attachBotListener(bot: BotInstance): void { + if (attachedBots.has(bot.id)) return; + + const onStateChange = () => { + broadcast({ + type: "stateChange", + botId: bot.id, + status: bot.getStatus(), + queue: bot.getQueue(), + }); + }; + + const onConnected = () => { + broadcast({ + type: "botConnected", + botId: bot.id, + status: bot.getStatus(), + }); + }; + + const onDisconnected = () => { + broadcast({ type: "botDisconnected", botId: bot.id }); + }; + + bot.on("stateChange", onStateChange); + bot.on("connected", onConnected); + bot.on("disconnected", onDisconnected); + + attachedBots.set(bot.id, { + stateChange: onStateChange, + connected: onConnected, + disconnected: onDisconnected, + }); + } + + /** Attach listeners for any new bots that don't have them yet */ + function ensureAllBotsAttached(): void { for (const bot of botManager.getAllBots()) { - bot.removeAllListeners("stateChange"); - bot.removeAllListeners("connected"); - bot.removeAllListeners("disconnected"); - - bot.on("stateChange", () => { - broadcast({ - type: "stateChange", - botId: bot.id, - status: bot.getStatus(), - queue: bot.getQueue(), - }); - }); - - bot.on("connected", () => { - broadcast({ - type: "botConnected", - botId: bot.id, - status: bot.getStatus(), - }); - }); - - bot.on("disconnected", () => { - broadcast({ type: "botDisconnected", botId: bot.id }); - }); + attachBotListener(bot); } - }; + } - const intervalId = setInterval(attachBotListeners, 5000); - attachBotListeners(); + // Check for newly added bots periodically + const intervalId = setInterval(ensureAllBotsAttached, 5000); + ensureAllBotsAttached(); return () => { clearInterval(intervalId); + // Clean up named listeners + for (const bot of botManager.getAllBots()) { + const listeners = attachedBots.get(bot.id); + if (listeners) { + bot.removeListener("stateChange", listeners.stateChange); + bot.removeListener("connected", listeners.connected); + bot.removeListener("disconnected", listeners.disconnected); + } + } + attachedBots.clear(); }; } diff --git a/web/src/components/Player.vue b/web/src/components/Player.vue index 8471689..85f556b 100644 --- a/web/src/components/Player.vue +++ b/web/src/components/Player.vue @@ -57,7 +57,7 @@ min="0" max="100" :value="activeBot?.volume ?? 75" - @input="onVolumeChange" + @change="onVolumeChange" class="volume-slider" />