fix: ffmpeg残留 每次只留一个ffmpeg进程 现已加入闭包校验 防止歌曲帧异常串入

This commit is contained in:
阿梓喵_あずにゃん committed 2026-04-21 03:53:50 +08:00
1 parent 82a23d291e
commit 364112d94f
1 file changed
+105 -228
+105 -228
View File
@@ -5,11 +5,12 @@ import { accessSync, chmodSync, constants } from "node:fs";
import { createOpusEncoder, PCM_FRAME_BYTES, type Encoder } from "./encoder.js";
import type { Logger } from "../logger.js";
// ffmpeg-static is a CJS module that exports the path to the bundled ffmpeg binary.
const require = createRequire(import.meta.url);
const ffmpegPath: string | null = require("ffmpeg-static");
/** Ensure the given binary has execute permission. */
/** 全局 PID 追踪器,防止进程在类实例切换时沦为孤儿进程 ( */
const globalActivePids = new Set<number>();
function isExecutable(binPath: string): boolean {
try {
accessSync(binPath, constants.X_OK);
@@ -25,7 +26,6 @@ function isExecutable(binPath: string): boolean {
}
}
/** Test if an ffmpeg binary actually works by running -version. */
function ffmpegWorks(bin: string): boolean {
try {
execSync(`"${bin}" -version`, { timeout: 5000, stdio: "pipe" });
@@ -35,39 +35,16 @@ function ffmpegWorks(bin: string): boolean {
}
}
/**
* Resolved once at module load.
*
* Priority: system FFmpeg → bundled ffmpeg-static.
*
* System-installed FFmpeg is always compatible with the running OS/container,
* while the pre-compiled binary from ffmpeg-static can SIGSEGV in Docker
* (passes `ffmpeg -version` but crashes during actual audio processing due to
* incompatible glibc or missing shared libraries).
*/
const resolvedFfmpeg: string = (() => {
// 1. Prefer system-installed FFmpeg (always compatible with the runtime)
if (ffmpegWorks("ffmpeg")) {
return "ffmpeg";
}
// 2. Fall back to bundled ffmpeg-static binary
// On Windows, ffmpeg-static may return a path with backslashes; on Linux/macOS
// it may return a Windows .exe path if node_modules was copied cross-platform.
if (ffmpegWorks("ffmpeg")) return "ffmpeg";
const isWinPath = ffmpegPath ? /\\/.test(ffmpegPath) || ffmpegPath.endsWith(".exe") : false;
const onWindows = process.platform === "win32";
if (ffmpegPath && (onWindows === isWinPath)) {
if (isExecutable(ffmpegPath) && ffmpegWorks(ffmpegPath)) {
return ffmpegPath;
if (isExecutable(ffmpegPath) && ffmpegWorks(ffmpegPath)) return ffmpegPath;
}
}
// Last resort: always use "ffmpeg" so spawn error is clear, never use a cross-platform path
return "ffmpeg";
})();
/** Resolve ffmpeg binary: prefer system PATH, fall back to bundled ffmpeg-static. */
function getFfmpegCommand(): string {
return resolvedFfmpeg;
}
@@ -93,12 +70,12 @@ export class AudioPlayer extends EventEmitter {
private nextFrameTime = 0;
private currentUrl = "";
private seekOffset = 0;
private framesPlayed = 0; // ground truth: number of 20ms frames sent
private framesPlayed = 0;
private sessionId = 0;
private static readonly BUFFER_HIGH_WATER = 640 * 1024; // ~5s of PCM at 48kHz stereo
private static readonly BUFFER_LOW_WATER = 256 * 1024; // ~2.5s
private static readonly BUFFER_HIGH_WATER = 640 * 1024;
private static readonly BUFFER_LOW_WATER = 256 * 1024;
private ffmpegPaused = false;
private spawnFailed = false; // true if ffmpeg spawn errored (prevent trackEnd cascade)
private spawnFailed = false;
private consecutiveFailures = 0;
private static readonly MAX_CONSECUTIVE_FAILURES = 3;
@@ -109,112 +86,125 @@ export class AudioPlayer extends EventEmitter {
}
play(url: string, seekSeconds = 0): void {
// 1. 停止当前所有播放,自增 sessionId 屏蔽旧回调 (
this.stop();
this.sessionId++;
const playSessionId = this.sessionId;
const currentSessionId = this.sessionId;
this.currentUrl = url;
this.seekOffset = seekSeconds;
this.framesPlayed = 0;
this.ffmpegPaused = false;
this.spawnFailed = false;
// Prevent rapid-fire spawn attempts when ffmpeg is broken
if (this.consecutiveFailures >= AudioPlayer.MAX_CONSECUTIVE_FAILURES) {
this.logger.error(
{ failures: this.consecutiveFailures, ffmpeg: getFfmpegCommand() },
"Too many consecutive ffmpeg failures — ffmpeg binary may be missing or broken. Stopping playback."
);
this.logger.error({ failures: this.consecutiveFailures }, "FFmpeg failures limit reached");
this.state = "idle";
this.emit("error", new Error("ffmpeg unavailable after repeated failures"));
this.emit("error", new Error("ffmpeg unavailable"));
return;
}
this.logger.info({ url: url.slice(0, 80), seek: seekSeconds }, "Starting playback");
const args: string[] = [];
// BiliBili CDN requires Referer header for audio playback
if (url.includes("bilivideo") || url.includes("bilibili")) {
args.push(
"-headers",
"Referer: https://www.bilibili.com\r\nUser-Agent: Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/120.0.0.0 Safari/537.36\r\n"
);
args.push("-headers", "Referer: https://www.bilibili.com\r\nUser-Agent: Mozilla/5.0...\r\n");
}
args.push(
"-reconnect", "1",
"-reconnect_streamed", "1",
"-reconnect_delay_max", "5",
);
if (seekSeconds > 0) {
args.push("-ss", String(seekSeconds));
}
args.push(
"-i", url,
"-f", "s16le",
"-ar", "48000",
"-ac", "2",
"-acodec", "pcm_s16le",
"-",
);
args.push("-reconnect", "1", "-reconnect_streamed", "1", "-reconnect_delay_max", "5");
if (seekSeconds > 0) args.push("-ss", String(seekSeconds));
args.push("-i", url, "-f", "s16le", "-ar", "48000", "-ac", "2", "-acodec", "pcm_s16le", "-");
const ffmpegBin = getFfmpegCommand();
this.logger.info({ ffmpeg: ffmpegBin }, "Using ffmpeg binary");
this.ffmpeg = spawn(ffmpegBin, args, { stdio: ["ignore", "pipe", "pipe"] });
// Prevent stream errors from crashing the process
this.ffmpeg.stdout!.on("error", (err) => {
this.logger.warn({ err }, "FFmpeg stdout error");
});
this.ffmpeg.stderr!.on("error", (err) => {
this.logger.warn({ err }, "FFmpeg stderr error");
});
let gotFirstData = false;
this.ffmpeg.stdout!.on("data", (chunk: Buffer) => {
if (!gotFirstData) {
gotFirstData = true;
this.logger.info({ bytes: chunk.length }, "FFmpeg: first PCM data received");
const currentPid = this.ffmpeg.pid;
if (currentPid) {
globalActivePids.add(currentPid);
this.logger.debug({ pid: currentPid, sessionId: currentSessionId }, "FFmpeg spawned");
}
this.ffmpeg.stdout!.on("data", (chunk: Buffer) => {
// 2. 严格校验 sessionId,防止老进程的数据混入新播放请求 (
if (this.sessionId !== currentSessionId) {
return;
}
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", (code, signal) => {
this.logger.info({ exitCode: code, signal, gotData: gotFirstData, framesPlayed: this.framesPlayed }, "FFmpeg process closed");
if (this.sessionId === playSessionId) {
this.ffmpeg = null; // Signal frame loop that no more data is coming
this.ffmpeg.on("exit", (code, signal) => {
if (currentPid) globalActivePids.delete(currentPid);
this.logger.info({ pid: currentPid, code, signal }, "FFmpeg exited");
// 只有当前会话的进程结束才置空变量
if (this.sessionId === currentSessionId) {
this.ffmpeg = null;
}
});
this.ffmpeg.on("error", (err) => {
this.logger.error({ err }, "FFmpeg error");
if (this.sessionId === playSessionId) {
if (this.sessionId === currentSessionId) {
this.spawnFailed = true;
this.consecutiveFailures++;
this.emit("error", err);
}
});
// Log FFmpeg stderr at info level for debugging playback issues
this.ffmpeg.stderr!.on("data", (data: Buffer) => {
const msg = data.toString().trimEnd();
// Log important FFmpeg messages at info level
if (msg.includes("Error") || msg.includes("error") || msg.includes("HTTP") || msg.includes("Opening") || msg.includes("Stream")) {
this.logger.info({ ffmpegStderr: msg }, "FFmpeg stderr");
} else {
this.logger.debug({ stderr: msg }, "FFmpeg stderr");
}
});
this.state = "playing";
this.startFrameLoop();
}
stop(): void {
// 3. 递增 ID 是最有效的逻辑“隔离墙”
this.sessionId++;
this.frameLoopRunning = false;
// 立即清空缓冲区,确保切歌瞬间静音 (
this.pcmBuffer = Buffer.alloc(0);
if (this.ffmpeg) {
const procToKill = this.ffmpeg;
const pidToKill = procToKill.pid;
this.ffmpeg = null;
if (pidToKill) {
this.forceCleanup(procToKill, pidToKill);
}
}
this.ffmpegPaused = false;
this.spawnFailed = false;
this.state = "idle";
this.currentUrl = "";
this.seekOffset = 0;
this.framesPlayed = 0;
}
private forceCleanup(proc: ChildProcess, pid: number): void {
if (!globalActivePids.has(pid)) return;
try {
proc.kill("SIGTERM");
} catch (e) { /* ignore */ }
const killTimeout = setTimeout(() => {
try {
process.kill(pid, 0);
process.kill(pid, "SIGKILL");
} catch (e) {
} finally {
globalActivePids.delete(pid);
}
}, 1500);
proc.unref();
proc.once("exit", () => {
clearTimeout(killTimeout);
globalActivePids.delete(pid);
});
}
private startFrameLoop(): void {
if (this.frameLoopRunning) return;
this.frameLoopRunning = true;
@@ -224,50 +214,37 @@ 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);
const delay = Math.max(0, this.nextFrameTime - performance.now());
setTimeout(() => {
// Discard callback from a stale play session
if (loopSessionId !== this.sessionId) return;
if (!this.frameLoopRunning) return;
// 这里的校验能防止旧的定时器回调处理新 Session 的逻辑 (
if (loopSessionId !== this.sessionId || !this.frameLoopRunning) return;
if (this.state === "playing") {
this.sendNextFrame();
} else if (this.state === "paused") {
this.nextFrameTime = performance.now();
}
if (this.state === "playing") this.sendNextFrame();
else if (this.state === "paused") this.nextFrameTime = performance.now();
if (!this.ffmpeg && this.pcmBuffer.length < PCM_FRAME_BYTES) {
this.frameLoopRunning = false;
if (this.state !== "idle") {
this.state = "idle";
// Don't emit trackEnd if ffmpeg spawn failed — prevents infinite retry cascade
if (this.spawnFailed) {
this.logger.warn("Suppressing trackEnd due to ffmpeg spawn failure");
} else {
this.consecutiveFailures = 0; // Reset on successful track completion
if (!this.spawnFailed) {
this.consecutiveFailures = 0;
this.emit("trackEnd");
}
}
return;
}
this.scheduleNextFrame();
}, delay);
}
private sendNextFrame(): void {
if (this.pcmBuffer.length < PCM_FRAME_BYTES) return;
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;
@@ -278,128 +255,28 @@ export class AudioPlayer extends EventEmitter {
const opusFrame = this.encoder.encode(adjusted);
this.emit("frame", opusFrame);
this.framesPlayed++;
if (this.framesPlayed === 1) {
this.logger.info({ opusBytes: opusFrame.length }, "First audio frame encoded and emitted");
}
// Log every ~10 seconds (500 frames * 20ms = 10s)
if (this.framesPlayed % 500 === 0) {
this.logger.debug({ framesPlayed: this.framesPlayed, elapsed: this.getElapsed() }, "Playback progress");
}
} catch (err) {
this.logger.error({ err }, "Error encoding/sending audio frame");
this.emit("error", err as Error);
}
}
private applyVolume(pcm: Buffer): Buffer {
if (this.volume === 100) return Buffer.from(pcm);
// Apply volume + fixed 6dB attenuation (factor ≈ 0.5)
const BASE_ATTENUATION = 0.2; // -6dB (10^(-6/20) ≈ 0.5012)
const factor = (this.volume / 100) * BASE_ATTENUATION;
const factor = (this.volume / 100) * 0.2;
const out = Buffer.alloc(pcm.length);
for (let i = 0; i < pcm.length; i += 2) {
let sample = pcm.readInt16LE(i);
sample = Math.round(sample * factor);
if (sample > 32767) sample = 32767;
else if (sample < -32768) sample = -32768;
out.writeInt16LE(sample, i);
let sample = Math.round(pcm.readInt16LE(i) * factor);
out.writeInt16LE(Math.max(-32768, Math.min(32767, sample)), i);
}
return out;
}
/** Actual elapsed time in seconds (ground truth from frame count) */
getElapsed(): number {
return this.seekOffset + (this.framesPlayed * FRAME_DURATION_MS) / 1000;
}
seek(seconds: number): void {
if (!this.currentUrl) return;
// Reject NaN/Infinity/negative — the HTTP layer validates too, but a
// bad value here would poison seekOffset and leave getElapsed()
// returning NaN until the track ends.
if (!Number.isFinite(seconds) || seconds < 0) {
this.logger.warn({ seek: seconds }, "Ignoring invalid seek position");
return;
}
this.logger.info({ seek: seconds }, "Seeking");
this.play(this.currentUrl, seconds);
}
getSeekOffset(): number {
return this.seekOffset;
}
pause(): void {
if (this.state === "playing") {
this.state = "paused";
this.logger.debug("Playback paused");
}
}
resume(): void {
if (this.state === "paused") {
this.state = "playing";
this.nextFrameTime = performance.now();
this.logger.debug("Playback resumed");
}
}
stop(): void {
this.sessionId++;
this.frameLoopRunning = false;
if (this.ffmpeg) {
const ffmpegPid = this.ffmpeg.pid;
this.logger.debug({ pid: ffmpegPid }, "Stopping ffmpeg process");
// Try SIGTERM first for graceful shutdown
this.ffmpeg.kill("SIGTERM");
// If process doesn't exit after 1 second, force kill with SIGKILL
const killTimeout = setTimeout(() => {
if (this.ffmpeg && this.ffmpeg.pid) {
this.logger.warn({ pid: this.ffmpeg.pid }, "FFmpeg didn't exit from SIGTERM, sending SIGKILL");
try {
this.ffmpeg.kill("SIGKILL");
} catch (err) {
this.logger.debug({ err }, "Error killing ffmpeg process");
}
}
}, 1000);
// Store timeout ID for cleanup on next stop call
const tempProc = this.ffmpeg;
const onClose = () => {
clearTimeout(killTimeout);
tempProc.removeListener("close", onClose);
};
tempProc.once("close", onClose);
this.ffmpeg = null;
}
this.pcmBuffer = Buffer.alloc(0);
this.ffmpegPaused = false;
this.spawnFailed = false;
this.state = "idle";
this.currentUrl = "";
this.seekOffset = 0;
this.framesPlayed = 0;
}
/** Reset the consecutive failure counter (e.g. after user action) */
resetFailures(): void {
this.consecutiveFailures = 0;
}
setVolume(vol: number): void {
this.volume = Math.max(0, Math.min(100, vol));
}
getVolume(): number {
return this.volume;
}
getState(): PlayerState {
return this.state;
}
getElapsed(): number { return this.seekOffset + (this.framesPlayed * FRAME_DURATION_MS) / 1000; }
seek(seconds: number): void { if (this.currentUrl && Number.isFinite(seconds) && seconds >= 0) this.play(this.currentUrl, seconds); }
pause(): void { if (this.state === "playing") this.state = "paused"; }
resume(): void { if (this.state === "paused") { this.state = "playing"; this.nextFrameTime = performance.now(); } }
resetFailures(): void { this.consecutiveFailures = 0; }
setVolume(vol: number): void { this.volume = Math.max(0, Math.min(100, vol)); }
getVolume(): number { return this.volume; }
getState(): PlayerState { return this.state; }
}