fix: resolve critical race conditions and high-severity bugs

- Add isAdvancing guard to playNext() to prevent concurrent calls from
  skipping songs (trackEnd + user !next race condition)
- Replace removeAllListeners in WebSocket setup with tracked named
  listeners, attach only once per bot instead of every 5 seconds
- Add .catch() to unhandled playNext() in cmdVote
- Add sessionId counter to AudioPlayer to discard stale setTimeout
  callbacks after stop()+play() transitions
- Defer nulling TS3 client until disconnect() promise resolves
- Add null checks for providers in play-by-id and add-by-id endpoints
- Add PCM buffer backpressure (pause FFmpeg at 960KB, resume at 480KB)
- Change volume slider from @input to @change to avoid flooding server
- Wrap error-reporting sendTextMessage in try/catch to prevent
  unhandled rejection in error handler
- Clamp elapsed getter to song duration
- Log FFmpeg stderr at debug level instead of swallowing
- Return null from prev() in Sequential mode instead of wrapping

Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
This commit is contained in:
saopig1andClaude Opus 4.6 committed 2026-03-30 14:31:03 +08:00
1 parent 4c71a0a37f
commit 17991b936a
8 files changed
+150 -57

No files matched your search

+33 -3
View File
@@ -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;
+2
View File
@@ -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;
+30 -17
View File
@@ -54,6 +54,7 @@ export class BotInstance extends EventEmitter {
private logger: Logger;
private connected = false;
private voteSkipUsers = new Set<string>();
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<void> {
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 {
+10 -3
View File
@@ -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");
+10 -2
View File
@@ -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) {
+60 -28
View File
@@ -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<WebSocket>();
/** Track which bots already have listeners attached */
const attachedBots = new Map<string, {
stateChange: () => 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();
};
}
+1 -1
View File
@@ -57,7 +57,7 @@
min="0"
max="100"
:value="activeBot?.volume ?? 75"
@input="onVolumeChange"
@change="onVolumeChange"
class="volume-slider"
/>
<button class="control-btn" :class="{ active: showQueue }" @click="showQueue = !showQueue">
+4 -3
View File
@@ -69,9 +69,10 @@ export const usePlayerStore = defineStore('player', {
/** Interpolated elapsed: serverElapsed + time since last sync (if playing) */
elapsed(): number {
if (!this.activeBot?.currentSong) return 0;
if (!this.wasPlaying || this.serverSyncTime === 0) return this.serverElapsed;
if (this.isPaused) return this.serverElapsed;
return this.serverElapsed + (Date.now() - this.serverSyncTime) / 1000;
const maxDuration = this.activeBot.currentSong.duration || Infinity;
if (!this.wasPlaying || this.serverSyncTime === 0) return Math.min(this.serverElapsed, maxDuration);
if (this.isPaused) return Math.min(this.serverElapsed, maxDuration);
return Math.min(this.serverElapsed + (Date.now() - this.serverSyncTime) / 1000, maxDuration);
},
},