mirror of
https://github.com/ZHANGTIANYAO1/teamspeak-music-bot.git
synced 2026-10-04 05:52:50 +08:00
Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
e5879ea106 | ||
|
|
44a1457baa | ||
|
|
a4c0239a59 |
No files matched your search
@@ -953,7 +953,15 @@ A:本项目内置 `/login` 限流(每 IP 每分钟 5 次),但生产部
|
||||
|
||||
> 完整历史请查看 [git log](https://github.com/ZHANGTIANYAO1/teamspeak-music-bot/commits/main) 或 [Releases](https://github.com/ZHANGTIANYAO1/teamspeak-music-bot/releases)。这里只列出重要变更和面向用户的破坏性改动。
|
||||
|
||||
### 最新版本 — v1.15.0:歌手页面 / REST API / TS6 Profile 与 B站续播修复
|
||||
### 最新版本 — v1.15.1:长视频播放阻塞与音乐 API 日志隐私修复
|
||||
|
||||
- 持续读取 FFmpeg 的 stderr,并关闭周期性进度输出,避免错误输出管道写满后卡住音频解码。直接 URL 播放和 Windows 临时文件播放均已处理。
|
||||
- FFmpeg 异常退出和播放停滞日志增加限长、脱敏的诊断摘要;移除 URL 查询参数、用户凭据和认证头,PowerShell 下载失败日志采用同样的处理。
|
||||
- 内置网易云 / QQ 音乐 API 在独立子进程运行,隔离依赖直接输出的原始请求和响应日志,避免其中的 Cookie 等凭据进入机器人控制台日志。仍保留服务启动和退出的安全诊断;独立部署或已占用端口的外部 API 需自行管理日志。
|
||||
|
||||
无配置或数据库迁移。此补丁修复了已复现的管道阻塞;[#161](https://github.com/ZHANGTIANYAO1/teamspeak-music-bot/issues/161) 中“67 分钟视频播到 37 分钟停止”的现场原因仍缺少停止时日志,尚未确认。
|
||||
|
||||
### v1.15.0:歌手页面 / REST API / TS6 Profile 与 B站续播修复
|
||||
|
||||
**歌手搜索与页面([PR #175](https://github.com/ZHANGTIANYAO1/teamspeak-music-bot/pull/175),感谢 [@zzstar101](https://github.com/zzstar101))**
|
||||
|
||||
|
||||
@@ -0,0 +1,83 @@
|
||||
import { describe, it, expect } from "vitest";
|
||||
import { PassThrough } from "node:stream";
|
||||
import { collectFfmpegDiagnostics } from "./ffmpeg-diagnostics.js";
|
||||
|
||||
async function finish(stream: PassThrough): Promise<void> {
|
||||
const ended = new Promise<void>((resolve) => stream.once("end", resolve));
|
||||
stream.end();
|
||||
await ended;
|
||||
}
|
||||
|
||||
describe("bounded FFmpeg diagnostics", () => {
|
||||
for (const [label, line, expected] of [
|
||||
["an apostrophe in the real FFmpeg URL error format", "Error opening input file http://127.0.0.1:9/audio?filename=artist's-song&api_key=quoted-secret.", "Error opening input file [URL omitted]"],
|
||||
["a space in URL userinfo", 'Error opening input file https://user:space secret@cdn.example/audio?token=space-secret.', "Error opening input file [URL omitted]"],
|
||||
["multiple URLs", 'Error opening inputs https://cdn.example/a?filename=artist\'s-song&key=first-secret and https://cdn.example/b?token=second-secret', "Error opening inputs [URL omitted]"],
|
||||
["quotes and spaces in a request target", "GET /audio?filename=artist's song&api_key=request-secret HTTP/1.1", "GET /audio?[query omitted]"],
|
||||
["an apostrophe in a Bearer value", "Token rejected Bearer prefix'quoted bearer-secret", "Token rejected Bearer [omitted]"],
|
||||
]) {
|
||||
it(`omits the entire sensitive suffix after ${label}`, async () => {
|
||||
const stream = new PassThrough();
|
||||
const diagnostics = collectFfmpegDiagnostics(stream);
|
||||
stream.write(line + "\n");
|
||||
await finish(stream);
|
||||
expect(diagnostics.getTail()).toBe(expected);
|
||||
});
|
||||
}
|
||||
|
||||
for (const splitDelimiter of [false, true]) {
|
||||
it(`omits folded authentication values with ${splitDelimiter ? "chunk-split" : "intact"} CRLF`, async () => {
|
||||
const stream = new PassThrough();
|
||||
const diagnostics = collectFfmpegDiagnostics(stream);
|
||||
stream.write(`Authorization: Basic header-secret\r${splitDelimiter ? "" : "\n"}`);
|
||||
if (splitDelimiter) stream.write("\n");
|
||||
stream.write(" continuation-secret\r\nHTTP error 401\r\n");
|
||||
await finish(stream);
|
||||
expect(diagnostics.getTail()).toContain("HTTP error 401");
|
||||
expect(diagnostics.getTail()).not.toContain("header-secret");
|
||||
expect(diagnostics.getTail()).not.toContain("continuation-secret");
|
||||
});
|
||||
}
|
||||
|
||||
it("redacts URLs and headers split across arbitrary byte and UTF-8 boundaries", async () => {
|
||||
const stream = new PassThrough();
|
||||
const diagnostics = collectFfmpegDiagnostics(stream);
|
||||
const payload = Buffer.from("解码失败 https://user:user-secret@cdn.example/audio?token=query-secret#fragment-secret\rCookie: cookie-secret\nAuthorization: Bearer bearer-secret\nGET /audio?token=request-secret HTTP/1.1\nfinal error");
|
||||
for (const byte of payload) stream.write(Buffer.from([byte]));
|
||||
await finish(stream);
|
||||
expect(diagnostics.getTail()).toContain("解码失败");
|
||||
expect(diagnostics.getTail()).toContain("final error");
|
||||
for (const secret of ["user-secret", "query-secret", "fragment-secret", "cookie-secret", "bearer-secret", "request-secret"]) {
|
||||
expect(diagnostics.getTail()).not.toContain(secret);
|
||||
}
|
||||
});
|
||||
|
||||
it("omits oversized raw lines without retaining an unsafe credential suffix", async () => {
|
||||
const stream = new PassThrough();
|
||||
const diagnostics = collectFfmpegDiagnostics(stream);
|
||||
stream.write("Cookie: " + "x".repeat(100000));
|
||||
stream.write("oversized-secret\r\n folded-oversized-secret\r\nHTTP error 403\n");
|
||||
await new Promise<void>((resolve) => setImmediate(resolve));
|
||||
expect(diagnostics.getTail()).toContain("HTTP error 403");
|
||||
expect(diagnostics.getTail()).not.toContain("oversized-secret");
|
||||
stream.write("decoder warning\n".repeat(10000));
|
||||
stream.write("last useful error\n");
|
||||
await finish(stream);
|
||||
expect(diagnostics.getTail().length).toBeLessThanOrEqual(4096);
|
||||
expect(diagnostics.getTail()).toContain("last useful error");
|
||||
expect(diagnostics.getTail()).not.toContain("oversized-secret");
|
||||
});
|
||||
|
||||
it("withholds an incomplete credential line until it can be safely sanitized", async () => {
|
||||
const stream = new PassThrough();
|
||||
const diagnostics = collectFfmpegDiagnostics(stream);
|
||||
stream.write("decoder warning\nhttps://user:partial-secret@");
|
||||
await new Promise<void>((resolve) => setImmediate(resolve));
|
||||
expect(diagnostics.getTail()).toContain("decoder warning");
|
||||
expect(diagnostics.getTail()).not.toContain("partial-secret");
|
||||
stream.write("cdn.example/audio?token=last-secret");
|
||||
await finish(stream);
|
||||
expect(diagnostics.getTail()).not.toContain("partial-secret");
|
||||
expect(diagnostics.getTail()).not.toContain("last-secret");
|
||||
});
|
||||
});
|
||||
@@ -0,0 +1,93 @@
|
||||
import type { Readable } from "node:stream";
|
||||
import { StringDecoder } from "node:string_decoder";
|
||||
|
||||
const MAX_LINE_CHARS = 2048;
|
||||
const MAX_TAIL_CHARS = 4096;
|
||||
|
||||
export interface FfmpegDiagnostics {
|
||||
getTail(): string;
|
||||
}
|
||||
|
||||
/**
|
||||
* Drain independently of the PCM pipe: unread stderr can block FFmpeg even
|
||||
* when stdout is being consumed. Keep only complete, sanitized lines. Never
|
||||
* retain a suffix of an oversized raw line: it may have lost its URL/header
|
||||
* prefix and would no longer be possible to redact safely.
|
||||
*/
|
||||
export function collectFfmpegDiagnostics(stderr: Readable | null): FfmpegDiagnostics {
|
||||
let tail = "";
|
||||
let pending = "";
|
||||
let discardLine = false;
|
||||
let suppressHeaderContinuation = false;
|
||||
let previousCR = false;
|
||||
const decoder = new StringDecoder("utf8");
|
||||
|
||||
const append = (text: string): void => {
|
||||
if (text) previousCR = false;
|
||||
if (discardLine) return;
|
||||
if (pending.length + text.length > MAX_LINE_CHARS) {
|
||||
pending = "";
|
||||
discardLine = true;
|
||||
return;
|
||||
}
|
||||
pending += text;
|
||||
};
|
||||
|
||||
const finishLine = (): void => {
|
||||
if (discardLine) {
|
||||
tail = (tail + "[oversized diagnostic line omitted]\n").slice(-MAX_TAIL_CHARS);
|
||||
// An omitted line may be an authentication header. Omit folded values.
|
||||
suppressHeaderContinuation = true;
|
||||
} else if (pending) {
|
||||
const line = pending
|
||||
.replace(/\x1b\[[0-?]*[ -/]*[@-~]/g, "")
|
||||
.replace(/[\x00-\x08\x0b\x0c\x0e-\x1f\x7f]/g, "");
|
||||
const authHeader = /\b(?:cookie|set-cookie|authorization|proxy-authorization)\s*[:=]/i.test(line);
|
||||
if (authHeader || (suppressHeaderContinuation && /^\s/.test(line))) {
|
||||
tail = (tail + "[authentication header omitted]\n").slice(-MAX_TAIL_CHARS);
|
||||
suppressHeaderContinuation = true;
|
||||
} else {
|
||||
suppressHeaderContinuation = false;
|
||||
const sanitized = line
|
||||
// URLs and credentials can contain quotes or spaces. Keep the error
|
||||
// prefix only; guessing a closing delimiter could expose a suffix.
|
||||
.replace(/\b[a-z][a-z\d+.-]*:\/\/[\s\S]*/i, "[URL omitted]")
|
||||
// FFmpeg can also print a request target without the scheme/host.
|
||||
.replace(/\?[\s\S]*/, "?[query omitted]")
|
||||
.replace(/\bBearer\s+[\s\S]*/i, "Bearer [omitted]");
|
||||
tail = (tail + sanitized + "\n").slice(-MAX_TAIL_CHARS);
|
||||
}
|
||||
} else {
|
||||
suppressHeaderContinuation = false;
|
||||
}
|
||||
pending = "";
|
||||
discardLine = false;
|
||||
};
|
||||
|
||||
const consume = (text: string): void => {
|
||||
const separators = /[\r\n]/g;
|
||||
let start = 0;
|
||||
for (let match = separators.exec(text); match; match = separators.exec(text)) {
|
||||
append(text.slice(start, match.index));
|
||||
if (match[0] === "\n" && previousCR) {
|
||||
previousCR = false;
|
||||
} else {
|
||||
finishLine();
|
||||
previousCR = match[0] === "\r";
|
||||
}
|
||||
start = match.index + 1;
|
||||
}
|
||||
append(text.slice(start));
|
||||
};
|
||||
|
||||
stderr?.on("data", (chunk: Buffer) => consume(decoder.write(chunk)));
|
||||
stderr?.on("end", () => {
|
||||
consume(decoder.end());
|
||||
finishLine();
|
||||
});
|
||||
// Resume explicitly as attaching a listener does not resume an already
|
||||
// paused Readable. No player pause/backpressure operation touches stderr.
|
||||
stderr?.resume();
|
||||
|
||||
return { getTail: () => tail.trimEnd() };
|
||||
}
|
||||
+260
-1
@@ -2,10 +2,20 @@ import { describe, it, expect, vi } from "vitest";
|
||||
import { mkdtempSync, writeFileSync, existsSync } from "node:fs";
|
||||
import { tmpdir } from "node:os";
|
||||
import { join } from "node:path";
|
||||
import { Readable } from "node:stream";
|
||||
import { Readable, PassThrough } from "node:stream";
|
||||
import { EventEmitter } from "node:events";
|
||||
import { spawn, type ChildProcess } from "node:child_process";
|
||||
import pino from "pino";
|
||||
import { buildFfmpegArgs, shouldUsePowerShellDownload, cleanupTempDir, shouldEndOnStall, volumeToFactor, AudioPlayer } from "./player.js";
|
||||
import type { Logger } from "../logger.js";
|
||||
|
||||
// Only replace process creation in the regression cases below. Their pipes,
|
||||
// write callbacks and backpressure are real OS resources, rather than mocks.
|
||||
vi.mock("node:child_process", async (importOriginal) => {
|
||||
const actual = await importOriginal<typeof import("node:child_process")>();
|
||||
return { ...actual, spawn: vi.fn(actual.spawn) };
|
||||
});
|
||||
|
||||
function getHeadersArg(args: string[]): string {
|
||||
const idx = args.indexOf("-headers");
|
||||
if (idx === -1) return "";
|
||||
@@ -36,6 +46,12 @@ describe("buildFfmpegArgs", () => {
|
||||
expect(args).not.toContain("-headers");
|
||||
});
|
||||
|
||||
it("disables periodic progress stats for both network and file inputs", () => {
|
||||
for (const input of ["https://example.com/song.mp3", "C:/temp/song.audio"]) {
|
||||
expect(buildFfmpegArgs(input, 0)).toContain("-nostats");
|
||||
}
|
||||
});
|
||||
|
||||
it("includes resilient reconnect flags for all URLs", () => {
|
||||
const args = buildFfmpegArgs("https://example.com/song.mp3", 0);
|
||||
expect(args).toContain("-reconnect");
|
||||
@@ -224,6 +240,249 @@ const silentLogger = {
|
||||
},
|
||||
} as unknown as Logger;
|
||||
|
||||
describe("AudioPlayer FFmpeg stderr handling", () => {
|
||||
const producer = `
|
||||
const pressure = 'decoder diagnostic\\n'.repeat(180000);
|
||||
process.stderr.write(pressure, () => {
|
||||
process.stderr.write('HTTP error 403 for https://user:password@cdn.example/audio?token=signed-secret#fragment-secret\\n');
|
||||
process.stderr.write('Cookie: cookie-secret\\nAuthorization: Bearer bearer-secret\\n');
|
||||
process.stderr.write('final decoder failure', () => {
|
||||
process.stdout.write(Buffer.alloc(7680), () => process.exit(1));
|
||||
});
|
||||
});
|
||||
`;
|
||||
|
||||
function recordedLogger() {
|
||||
const records: Array<{ level: string; fields: Record<string, unknown>; message: string }> = [];
|
||||
const capture = (level: string) => (fields: Record<string, unknown>, message: string) => {
|
||||
records.push({ level, fields, message });
|
||||
};
|
||||
return {
|
||||
records,
|
||||
logger: { ...silentLogger, info: capture("info"), warn: capture("warn") } as unknown as Logger,
|
||||
};
|
||||
}
|
||||
|
||||
async function pipeProducer(script: string) {
|
||||
const actual = await vi.importActual<typeof import("node:child_process")>("node:child_process");
|
||||
let child!: ChildProcess;
|
||||
let requestedArgs: readonly string[] = [];
|
||||
let closed!: Promise<number | null>;
|
||||
vi.mocked(spawn).mockImplementationOnce((_command, args, options) => {
|
||||
requestedArgs = args ?? [];
|
||||
child = actual.spawn(process.execPath, ["-e", script], options);
|
||||
closed = new Promise((resolve) => child.once("close", resolve));
|
||||
return child;
|
||||
});
|
||||
return {
|
||||
get child() { return child; },
|
||||
get args() { return requestedArgs; },
|
||||
get closed() { return closed; },
|
||||
};
|
||||
}
|
||||
|
||||
for (const path of ["URL", "temp file"] as const) {
|
||||
it(`drains the ${path} child stderr so a large diagnostic write cannot block PCM output`, async () => {
|
||||
const { records, logger } = recordedLogger();
|
||||
const producerProcess = await pipeProducer(producer);
|
||||
const player = new AudioPlayer(logger);
|
||||
let frameCount = 0;
|
||||
player.on("frame", () => frameCount++);
|
||||
let deadline: ReturnType<typeof setTimeout> | undefined;
|
||||
try {
|
||||
if (path === "URL") {
|
||||
player.play("https://cdn.example/audio?token=input-secret");
|
||||
} else {
|
||||
// The real downloader marks playing before calling this file path.
|
||||
const internal = player as unknown as {
|
||||
state: string;
|
||||
spawnFfmpegFromFile(file: string, seek: number, session: number): void;
|
||||
};
|
||||
internal.state = "playing";
|
||||
internal.spawnFfmpegFromFile("downloaded.audio", 0, player.getPlaybackSessionId());
|
||||
}
|
||||
const outcome = await Promise.race([
|
||||
producerProcess.closed,
|
||||
new Promise<string>((resolve) => {
|
||||
deadline = setTimeout(() => resolve("stderr blocked audio output"), 1500);
|
||||
}),
|
||||
]);
|
||||
expect(outcome).toBe(1);
|
||||
expect(producerProcess.args).toContain("-nostats");
|
||||
await vi.waitFor(() => expect(frameCount).toBeGreaterThan(0));
|
||||
const exit = records.find((record) => record.message === "FFmpeg exited");
|
||||
expect(exit?.fields.stderr).toContain("final decoder failure");
|
||||
expect(exit?.fields.stderr).toContain("HTTP error 403");
|
||||
const logged = JSON.stringify(records);
|
||||
for (const secret of ["password", "signed-secret", "fragment-secret", "cookie-secret", "bearer-secret", "input-secret"]) {
|
||||
expect(logged).not.toContain(secret);
|
||||
}
|
||||
expect(String(exit?.fields.stderr).length).toBeLessThanOrEqual(4096);
|
||||
} finally {
|
||||
if (deadline) clearTimeout(deadline);
|
||||
player.stop();
|
||||
if (producerProcess.child.exitCode === null) producerProcess.child.kill("SIGKILL");
|
||||
await producerProcess.closed;
|
||||
}
|
||||
});
|
||||
|
||||
it(`does not expose the ${path} input credentials through serialized spawn errors`, async () => {
|
||||
const actual = await vi.importActual<typeof import("node:child_process")>("node:child_process");
|
||||
let child!: ChildProcess;
|
||||
let closed!: Promise<void>;
|
||||
vi.mocked(spawn).mockImplementationOnce((_command, args, options) => {
|
||||
child = actual.spawn(join(tmpdir(), "tsbot-ffmpeg-does-not-exist"), args, options);
|
||||
closed = new Promise((resolve) => child.once("close", () => resolve()));
|
||||
return child;
|
||||
});
|
||||
const player = new AudioPlayer(silentLogger);
|
||||
const emitted = new Promise<Error>((resolve) => player.once("error", resolve));
|
||||
try {
|
||||
const input = "https://user:spawn-password@cdn.example/audio?token=spawn-secret";
|
||||
if (path === "URL") {
|
||||
player.play(input);
|
||||
} else {
|
||||
(player as unknown as { spawnFfmpegFromFile(file: string, seek: number, session: number): void })
|
||||
.spawnFfmpegFromFile(input, 0, player.getPlaybackSessionId());
|
||||
}
|
||||
const error = await emitted;
|
||||
expect(error).toBeInstanceOf(Error);
|
||||
expect(error.message).toContain("ENOENT");
|
||||
const serialized = JSON.stringify(pino.stdSerializers.err(error));
|
||||
expect(serialized).not.toContain("spawn-secret");
|
||||
expect(serialized).not.toContain("spawn-password");
|
||||
await closed;
|
||||
} finally {
|
||||
player.stop();
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
it("logs intentional stop signals at info level", async () => {
|
||||
const { records, logger } = recordedLogger();
|
||||
const producerProcess = await pipeProducer("process.stdout.write(Buffer.from([0])); setInterval(() => {}, 1000);");
|
||||
const player = new AudioPlayer(logger);
|
||||
try {
|
||||
player.play("https://cdn.example/audio");
|
||||
await new Promise<void>((resolve) => producerProcess.child.stdout!.once("data", () => resolve()));
|
||||
player.stop();
|
||||
await producerProcess.closed;
|
||||
const exit = records.find((record) => record.message === "FFmpeg exited");
|
||||
expect(exit?.level).toBe("info");
|
||||
} finally {
|
||||
player.stop();
|
||||
if (producerProcess.child.exitCode === null) producerProcess.child.kill("SIGKILL");
|
||||
await producerProcess.closed;
|
||||
}
|
||||
});
|
||||
|
||||
it("reports a sanitized diagnostic tail when a child exits from an unexpected signal", async () => {
|
||||
const { records, logger } = recordedLogger();
|
||||
const child = Object.assign(new EventEmitter(), {
|
||||
pid: undefined,
|
||||
stdout: new PassThrough(),
|
||||
stderr: new PassThrough(),
|
||||
});
|
||||
vi.mocked(spawn).mockReturnValueOnce(child as unknown as ChildProcess);
|
||||
const player = new AudioPlayer(logger);
|
||||
try {
|
||||
player.play("https://cdn.example/audio");
|
||||
child.stderr.write("decoder crashed for https://cdn.example/audio?token=signal-secret");
|
||||
const ended = new Promise<void>((resolve) => child.stderr.once("end", resolve));
|
||||
child.stderr.end();
|
||||
await ended;
|
||||
child.emit("exit", null, "SIGSEGV");
|
||||
child.emit("close", null, "SIGSEGV");
|
||||
const exit = records.find((record) => record.message === "FFmpeg exited");
|
||||
expect(exit?.level).toBe("warn");
|
||||
expect(exit?.fields.signal).toBe("SIGSEGV");
|
||||
expect(exit?.fields.stderr).toContain("decoder crashed");
|
||||
expect(JSON.stringify(records)).not.toContain("signal-secret");
|
||||
} finally {
|
||||
player.stop();
|
||||
child.stdout.destroy();
|
||||
child.stderr.destroy();
|
||||
}
|
||||
});
|
||||
|
||||
it("keeps draining an old child's stderr without mixing its late diagnostics or exit into a new session", async () => {
|
||||
const { records, logger } = recordedLogger();
|
||||
const makeChild = () => Object.assign(new EventEmitter(), {
|
||||
pid: undefined,
|
||||
stdout: new PassThrough(),
|
||||
stderr: new PassThrough(),
|
||||
});
|
||||
const oldChild = makeChild();
|
||||
const newChild = makeChild();
|
||||
vi.mocked(spawn)
|
||||
.mockReturnValueOnce(oldChild as unknown as ChildProcess)
|
||||
.mockReturnValueOnce(newChild as unknown as ChildProcess);
|
||||
const player = new AudioPlayer(logger);
|
||||
try {
|
||||
player.play("https://cdn.example/old");
|
||||
const oldSession = player.getPlaybackSessionId();
|
||||
player.play("https://cdn.example/new");
|
||||
const newSession = player.getPlaybackSessionId();
|
||||
oldChild.stderr.write("old late decoder failure https://cdn.example/old?token=old-secret\n".repeat(1000));
|
||||
newChild.stderr.write("new decoder failure\n");
|
||||
await new Promise<void>((resolve) => setImmediate(resolve));
|
||||
expect(oldChild.stderr.readableLength).toBe(0);
|
||||
oldChild.emit("exit", 1, null);
|
||||
oldChild.stderr.end();
|
||||
oldChild.emit("close", 1, null);
|
||||
expect(player.getState()).toBe("playing");
|
||||
vi.useFakeTimers();
|
||||
const internal = player as unknown as { frameLoopRunning: boolean; startFrameLoop(): void };
|
||||
internal.frameLoopRunning = false;
|
||||
internal.startFrameLoop();
|
||||
vi.advanceTimersByTime(6000);
|
||||
const stall = records.find((record) => record.message === "FFmpeg stopped outputting data, ending track");
|
||||
expect(stall?.fields.sessionId).toBe(newSession);
|
||||
expect(stall?.fields.stderr).toContain("new decoder failure");
|
||||
expect(stall?.fields.stderr).not.toContain("old late decoder failure");
|
||||
const oldExit = records.find((record) => record.message === "FFmpeg exited");
|
||||
expect(oldExit?.fields.sessionId).toBe(oldSession);
|
||||
expect(JSON.stringify(records)).not.toContain("old-secret");
|
||||
} finally {
|
||||
vi.useRealTimers();
|
||||
player.stop();
|
||||
oldChild.stdout.destroy();
|
||||
oldChild.stderr.destroy();
|
||||
newChild.stdout.destroy();
|
||||
newChild.stderr.destroy();
|
||||
}
|
||||
});
|
||||
|
||||
it("includes the current child's sanitized diagnostic tail when the stall watchdog ends playback", async () => {
|
||||
const { records, logger } = recordedLogger();
|
||||
const producerProcess = await pipeProducer(`
|
||||
process.stderr.write('HTTP error 403: https://cdn.example/audio?token=stall-secret\\n');
|
||||
process.stdout.write(Buffer.from([0]));
|
||||
setInterval(() => {}, 1000);
|
||||
`);
|
||||
const player = new AudioPlayer(logger);
|
||||
try {
|
||||
player.play("https://cdn.example/audio");
|
||||
await new Promise<void>((resolve) => producerProcess.child.stdout!.once("data", () => resolve()));
|
||||
vi.useFakeTimers();
|
||||
// Restart scheduling under the test clock, without changing EOF state.
|
||||
const internal = player as unknown as { frameLoopRunning: boolean; startFrameLoop(): void };
|
||||
internal.frameLoopRunning = false;
|
||||
internal.startFrameLoop();
|
||||
vi.advanceTimersByTime(6000);
|
||||
const stall = records.find((record) => record.message === "FFmpeg stopped outputting data, ending track");
|
||||
expect(stall?.fields.stderr).toContain("HTTP error 403");
|
||||
expect(JSON.stringify(records)).not.toContain("stall-secret");
|
||||
expect(player.getState()).toBe("idle");
|
||||
} finally {
|
||||
vi.useRealTimers();
|
||||
player.stop();
|
||||
if (producerProcess.child.exitCode === null) producerProcess.child.kill("SIGKILL");
|
||||
await producerProcess.closed;
|
||||
}
|
||||
});
|
||||
});
|
||||
|
||||
function applyPlayerVolume(player: AudioPlayer, pcm: Buffer): Buffer {
|
||||
return (
|
||||
player as unknown as { applyVolume(input: Buffer): Buffer }
|
||||
|
||||
+53
-10
@@ -7,6 +7,7 @@ import { join } from "node:path";
|
||||
import { createOpusEncoder, PCM_FRAME_BYTES, type Encoder } from "./encoder.js";
|
||||
import type { Readable } from "node:stream";
|
||||
import type { Logger } from "../logger.js";
|
||||
import { collectFfmpegDiagnostics, type FfmpegDiagnostics } from "./ffmpeg-diagnostics.js";
|
||||
|
||||
const require = createRequire(import.meta.url);
|
||||
const ffmpegPath: string | null = require("ffmpeg-static");
|
||||
@@ -52,6 +53,17 @@ export function getFfmpegCommand(): string {
|
||||
return resolvedFfmpeg;
|
||||
}
|
||||
|
||||
function safeFfmpegSpawnError(err: Error): Error {
|
||||
// Node spawn errors include spawnargs; Pino's Error serializer copies them,
|
||||
// including the signed input URL. Preserve a known OS category only.
|
||||
const allowedCodes = new Set(["ENOENT", "EACCES", "EPERM", "ENOEXEC", "EMFILE", "ENFILE", "ENOMEM", "EAGAIN", "EINVAL"]);
|
||||
const code = (err as NodeJS.ErrnoException).code;
|
||||
const safeCode = typeof code === "string" && allowedCodes.has(code) ? code : undefined;
|
||||
const safeError = new Error(`FFmpeg failed to start${safeCode ? ` (${safeCode})` : ""}`);
|
||||
if (safeCode) Object.assign(safeError, { code: safeCode });
|
||||
return safeError;
|
||||
}
|
||||
|
||||
const BROWSER_UA =
|
||||
"Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/120.0.0.0 Safari/537.36";
|
||||
|
||||
@@ -74,7 +86,7 @@ export function cleanupTempDir(dir: string): void {
|
||||
}
|
||||
|
||||
export function buildFfmpegArgs(url: string, seekSeconds: number): string[] {
|
||||
const args: string[] = [];
|
||||
const args: string[] = ["-nostats"];
|
||||
const isHttp = /^https?:\/\//i.test(url);
|
||||
const isBilibili = isHttp && (url.includes("bilivideo") || url.includes("bilibili"));
|
||||
|
||||
@@ -167,6 +179,8 @@ const FRAME_DURATION_MS = 20;
|
||||
|
||||
export class AudioPlayer extends EventEmitter {
|
||||
private ffmpeg: ChildProcess | null = null;
|
||||
private ffmpegDiagnostics: FfmpegDiagnostics | null = null;
|
||||
private readonly intentionalCleanup = new WeakSet<ChildProcess>();
|
||||
private encoder: Encoder;
|
||||
private state: PlayerState = "idle";
|
||||
private volume = 75;
|
||||
@@ -258,6 +272,9 @@ export class AudioPlayer extends EventEmitter {
|
||||
|
||||
const ffmpegBin = getFfmpegCommand();
|
||||
this.ffmpeg = spawn(ffmpegBin, args, { stdio: ["ignore", "pipe", "pipe"] });
|
||||
const child = this.ffmpeg;
|
||||
const diagnostics = collectFfmpegDiagnostics(child.stderr);
|
||||
this.ffmpegDiagnostics = diagnostics;
|
||||
|
||||
const currentPid = this.ffmpeg.pid;
|
||||
if (currentPid) {
|
||||
@@ -280,7 +297,6 @@ export class AudioPlayer extends EventEmitter {
|
||||
|
||||
this.ffmpeg.on("exit", (code, signal) => {
|
||||
if (currentPid) globalActivePids.delete(currentPid);
|
||||
this.logger.info({ pid: currentPid, code, signal }, "FFmpeg exited");
|
||||
|
||||
// 只有当前会话的进程结束才置空变量
|
||||
if (this.sessionId === currentSessionId) {
|
||||
@@ -288,11 +304,21 @@ export class AudioPlayer extends EventEmitter {
|
||||
}
|
||||
});
|
||||
|
||||
// close follows stderr's end, so final unterminated diagnostics are ready.
|
||||
this.ffmpeg.on("close", (code, signal) => {
|
||||
const context = { pid: currentPid, sessionId: currentSessionId, code, signal };
|
||||
if (!this.intentionalCleanup.has(child) && ((typeof code === "number" && code !== 0) || signal !== null)) {
|
||||
this.logger.warn({ ...context, stderr: diagnostics.getTail() }, "FFmpeg exited");
|
||||
} else {
|
||||
this.logger.info(context, "FFmpeg exited");
|
||||
}
|
||||
});
|
||||
|
||||
this.ffmpeg.on("error", (err) => {
|
||||
if (this.sessionId === currentSessionId) {
|
||||
this.spawnFailed = true;
|
||||
this.consecutiveFailures++;
|
||||
this.emit("error", err);
|
||||
this.emit("error", safeFfmpegSpawnError(err));
|
||||
}
|
||||
});
|
||||
|
||||
@@ -332,10 +358,7 @@ export class AudioPlayer extends EventEmitter {
|
||||
);
|
||||
this.downloader = ps;
|
||||
|
||||
let stderrTail = "";
|
||||
ps.stderr!.on("data", (chunk: Buffer) => {
|
||||
stderrTail = (stderrTail + chunk.toString()).slice(-500);
|
||||
});
|
||||
const diagnostics = collectFfmpegDiagnostics(ps.stderr);
|
||||
|
||||
ps.on("exit", (code, signal) => {
|
||||
if (this.sessionId !== sessionId) {
|
||||
@@ -344,7 +367,6 @@ export class AudioPlayer extends EventEmitter {
|
||||
}
|
||||
this.downloader = null;
|
||||
if (code !== 0) {
|
||||
this.logger.warn({ code, signal, stderr: stderrTail }, "PowerShell download failed");
|
||||
this.spawnFailed = true;
|
||||
this.consecutiveFailures++;
|
||||
this.state = "idle";
|
||||
@@ -356,6 +378,12 @@ export class AudioPlayer extends EventEmitter {
|
||||
this.spawnFfmpegFromFile(tempFile, seekSeconds, sessionId);
|
||||
});
|
||||
|
||||
ps.on("close", (code, signal) => {
|
||||
if (!this.intentionalCleanup.has(ps) && ((typeof code === "number" && code !== 0) || signal !== null)) {
|
||||
this.logger.warn({ pid: ps.pid, sessionId, code, signal, stderr: diagnostics.getTail() }, "PowerShell download failed");
|
||||
}
|
||||
});
|
||||
|
||||
ps.on("error", (err) => {
|
||||
if (this.sessionId !== sessionId) return;
|
||||
this.downloader = null;
|
||||
@@ -386,6 +414,9 @@ export class AudioPlayer extends EventEmitter {
|
||||
const args = buildFfmpegArgs(tempFile, seekSeconds);
|
||||
const ffmpegBin = getFfmpegCommand();
|
||||
this.ffmpeg = spawn(ffmpegBin, args, { stdio: ["ignore", "pipe", "pipe"] });
|
||||
const child = this.ffmpeg;
|
||||
const diagnostics = collectFfmpegDiagnostics(child.stderr);
|
||||
this.ffmpegDiagnostics = diagnostics;
|
||||
|
||||
const currentPid = this.ffmpeg.pid;
|
||||
if (currentPid) {
|
||||
@@ -405,7 +436,6 @@ export class AudioPlayer extends EventEmitter {
|
||||
|
||||
this.ffmpeg.on("exit", (code, signal) => {
|
||||
if (currentPid) globalActivePids.delete(currentPid);
|
||||
this.logger.info({ pid: currentPid, code, signal }, "FFmpeg exited");
|
||||
if (this.sessionId === sessionId) {
|
||||
this.ffmpeg = null;
|
||||
if (this.currentTempDir === tempDirToCleanup) this.currentTempDir = null;
|
||||
@@ -413,11 +443,20 @@ export class AudioPlayer extends EventEmitter {
|
||||
if (tempDirToCleanup) cleanupTempDir(tempDirToCleanup);
|
||||
});
|
||||
|
||||
this.ffmpeg.on("close", (code, signal) => {
|
||||
const context = { pid: currentPid, sessionId, code, signal };
|
||||
if (!this.intentionalCleanup.has(child) && ((typeof code === "number" && code !== 0) || signal !== null)) {
|
||||
this.logger.warn({ ...context, stderr: diagnostics.getTail() }, "FFmpeg exited");
|
||||
} else {
|
||||
this.logger.info(context, "FFmpeg exited");
|
||||
}
|
||||
});
|
||||
|
||||
this.ffmpeg.on("error", (err) => {
|
||||
if (this.sessionId === sessionId) {
|
||||
this.spawnFailed = true;
|
||||
this.consecutiveFailures++;
|
||||
this.emit("error", err);
|
||||
this.emit("error", safeFfmpegSpawnError(err));
|
||||
}
|
||||
});
|
||||
|
||||
@@ -543,6 +582,7 @@ export class AudioPlayer extends EventEmitter {
|
||||
|
||||
// 立即清空缓冲区,确保切歌瞬间静音 (
|
||||
this.pcmBuffer = Buffer.alloc(0);
|
||||
this.ffmpegDiagnostics = null;
|
||||
|
||||
if (this.ffmpeg) {
|
||||
const procToKill = this.ffmpeg;
|
||||
@@ -557,6 +597,7 @@ export class AudioPlayer extends EventEmitter {
|
||||
if (this.downloader) {
|
||||
const ps = this.downloader;
|
||||
this.downloader = null;
|
||||
this.intentionalCleanup.add(ps);
|
||||
try { ps.kill("SIGTERM"); } catch { /* already gone */ }
|
||||
}
|
||||
|
||||
@@ -580,6 +621,7 @@ export class AudioPlayer extends EventEmitter {
|
||||
}
|
||||
|
||||
private forceCleanup(proc: ChildProcess, pid: number): void {
|
||||
this.intentionalCleanup.add(proc);
|
||||
if (!globalActivePids.has(pid)) return;
|
||||
|
||||
try {
|
||||
@@ -664,6 +706,7 @@ export class AudioPlayer extends EventEmitter {
|
||||
duration: this.currentSongDuration,
|
||||
remaining: Math.round(this.currentSongDuration - elapsed),
|
||||
nearEnd: isNearEnd,
|
||||
stderr: this.ffmpegDiagnostics?.getTail() ?? "",
|
||||
}, "FFmpeg stopped outputting data, ending track");
|
||||
this.frameLoopRunning = false;
|
||||
// The outer gate guarantees state==="playing" here, so no !=="idle"
|
||||
|
||||
@@ -0,0 +1,43 @@
|
||||
import type { Server } from "node:http";
|
||||
import { closeEmbeddedApi, getSafeApiStartupError, startEmbeddedApi, type ApiChildMessage, type ApiProvider } from "./api-server-runtime.js";
|
||||
|
||||
const providerArg = process.argv[2];
|
||||
const port = Number(process.argv[3]);
|
||||
if ((providerArg !== "netease" && providerArg !== "qq") || !Number.isInteger(port) || port < 1 || port > 65535 || !process.send) process.exit(1);
|
||||
const provider = providerArg as ApiProvider;
|
||||
let server: Server | null = null;
|
||||
let stopping = false;
|
||||
|
||||
function shutdown(exitCode = 0): void {
|
||||
if (stopping) return;
|
||||
stopping = true;
|
||||
// Also covers a legacy auto-start listener and an import still in flight.
|
||||
const deadline = setTimeout(() => process.exit(exitCode), 1000);
|
||||
closeEmbeddedApi(server).finally(() => { clearTimeout(deadline); process.exit(exitCode); });
|
||||
}
|
||||
|
||||
function fail(error: unknown): void {
|
||||
if (stopping) return;
|
||||
const message: ApiChildMessage = { type: "error", provider, port, ...getSafeApiStartupError(error) };
|
||||
if (!process.connected) { shutdown(1); return; }
|
||||
try { process.send!(message, () => shutdown(1)); }
|
||||
catch { shutdown(1); }
|
||||
}
|
||||
|
||||
process.on("message", (message: unknown) => {
|
||||
if (message && typeof message === "object" && (message as { type?: unknown }).type === "stop") shutdown();
|
||||
});
|
||||
process.on("disconnect", () => shutdown());
|
||||
process.on("SIGTERM", () => shutdown());
|
||||
process.on("SIGINT", () => shutdown());
|
||||
process.on("uncaughtException", fail);
|
||||
process.on("unhandledRejection", fail);
|
||||
|
||||
startEmbeddedApi(provider, port).then((runtime) => {
|
||||
server = runtime.server;
|
||||
if (stopping || !process.connected) { shutdown(); return; }
|
||||
server?.on("error", fail);
|
||||
const message: ApiChildMessage = { type: "ready", provider, port };
|
||||
try { process.send!(message, (error) => { if (error) shutdown(1); }); }
|
||||
catch { shutdown(1); }
|
||||
}).catch(fail);
|
||||
@@ -0,0 +1,86 @@
|
||||
import { EventEmitter } from "node:events";
|
||||
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
|
||||
import { closeEmbeddedApi, getSafeApiStartupError, startEmbeddedApi } from "./api-server-runtime.js";
|
||||
|
||||
const state = vi.hoisted(() => ({ serveNcmApi: vi.fn(), qqExport: null as any, portFree: true }));
|
||||
vi.mock("NeteaseCloudMusicApi", () => ({ server: { serveNcmApi: state.serveNcmApi } }));
|
||||
vi.mock("@sansenjian/qq-music-api", () => ({ get default() { return state.qqExport; } }));
|
||||
vi.mock("node:net", async () => {
|
||||
const { EventEmitter } = await import("node:events");
|
||||
return { default: { createServer: () => {
|
||||
const probe = new EventEmitter() as EventEmitter & { close(done: () => void): void; listen(): void };
|
||||
probe.close = (done) => queueMicrotask(done);
|
||||
probe.listen = () => queueMicrotask(() => probe.emit(state.portFree ? "listening" : "error"));
|
||||
return probe;
|
||||
} } };
|
||||
});
|
||||
|
||||
class FakeServer extends EventEmitter {
|
||||
listening = false;
|
||||
close = vi.fn((done: () => void) => { this.listening = false; queueMicrotask(done); return this; });
|
||||
closeAllConnections = vi.fn();
|
||||
}
|
||||
|
||||
describe("embedded API runtime", () => {
|
||||
let server: FakeServer;
|
||||
let listen: ReturnType<typeof vi.fn>;
|
||||
let previousPort: string | undefined;
|
||||
beforeEach(() => {
|
||||
previousPort = process.env.PORT;
|
||||
server = new FakeServer();
|
||||
state.portFree = true;
|
||||
state.serveNcmApi.mockReset();
|
||||
state.serveNcmApi.mockResolvedValue({ server });
|
||||
listen = vi.fn(() => { queueMicrotask(() => { server.listening = true; server.emit("listening"); }); return server; });
|
||||
state.qqExport = { listen };
|
||||
});
|
||||
afterEach(() => { if (previousPort === undefined) delete process.env.PORT; else process.env.PORT = previousPort; });
|
||||
|
||||
it("binds NetEase to loopback/configured port without version checks and waits for listening", async () => {
|
||||
let ready = false;
|
||||
const starting = startEmbeddedApi("netease", 39218).then((result) => { ready = true; return result; });
|
||||
await vi.waitFor(() => expect(server.listenerCount("listening")).toBe(1));
|
||||
expect(state.serveNcmApi).toHaveBeenCalledWith({ port: 39218, host: "127.0.0.1", checkVersion: false });
|
||||
expect(ready).toBe(false);
|
||||
server.listening = true; server.emit("listening");
|
||||
expect((await starting).server).toBe(server);
|
||||
});
|
||||
it("closes the HTTP server returned by NetEase rather than the Express app", async () => {
|
||||
server.listening = true;
|
||||
const runtime = await startEmbeddedApi("netease", 39218);
|
||||
await closeEmbeddedApi(runtime.server);
|
||||
expect(server.close).toHaveBeenCalledTimes(1);
|
||||
expect(server.closeAllConnections).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
it("rejects failed listening and cleans up the startup handle", async () => {
|
||||
const starting = startEmbeddedApi("netease", 39218);
|
||||
const rejected = expect(starting).rejects.toMatchObject({ code: "EADDRINUSE" });
|
||||
await vi.waitFor(() => expect(server.listenerCount("error")).toBe(1));
|
||||
server.emit("error", Object.assign(new Error("bind failure"), { code: "EADDRINUSE" }));
|
||||
await rejected;
|
||||
expect(server.close).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
it("binds QQ to its configured loopback port and restores injected PORT", async () => {
|
||||
process.env.PORT = "39999";
|
||||
expect((await startEmbeddedApi("qq", 39217)).server).toBe(server);
|
||||
expect(listen).toHaveBeenCalledWith(39217, "127.0.0.1");
|
||||
expect(process.env.PORT).toBe("39999");
|
||||
});
|
||||
it("restores an absent PORT and supports the legacy nested export", async () => {
|
||||
delete process.env.PORT; state.qqExport = { default: { listen } };
|
||||
await startEmbeddedApi("qq", 39217);
|
||||
expect(process.env.PORT).toBeUndefined();
|
||||
expect(listen).toHaveBeenCalledWith(39217, "127.0.0.1");
|
||||
});
|
||||
it("reuses a legacy module that auto-started on import without a duplicate listen", async () => {
|
||||
state.portFree = false;
|
||||
expect(await startEmbeddedApi("qq", 39217)).toEqual({ server: null });
|
||||
expect(listen).not.toHaveBeenCalled();
|
||||
});
|
||||
it("never reflects arbitrary startup message, stack or code values", () => {
|
||||
expect(getSafeApiStartupError({ message: "synthetic-credential", stack: "synthetic-credential", code: "synthetic-credential" })).toEqual({ category: "startup" });
|
||||
expect(getSafeApiStartupError({ code: "ERR_REQUIRE_ESM", message: "synthetic-credential" })).toEqual({ category: "esm", code: "ERR_REQUIRE_ESM" });
|
||||
expect(getSafeApiStartupError({ code: "EBADENGINE" })).toEqual({ category: "node-engine", code: "EBADENGINE" });
|
||||
expect(getSafeApiStartupError({ code: "EADDRINUSE" })).toEqual({ category: "port-in-use", code: "EADDRINUSE" });
|
||||
});
|
||||
});
|
||||
@@ -0,0 +1,92 @@
|
||||
import net from "node:net";
|
||||
import type { Server } from "node:http";
|
||||
|
||||
export type ApiProvider = "netease" | "qq";
|
||||
export type ApiStartupCategory = "esm" | "node-engine" | "port-in-use" | "startup" | "timeout" | "cancelled";
|
||||
export interface SafeApiStartupError { category: ApiStartupCategory; code?: string }
|
||||
export type ApiChildMessage =
|
||||
| { type: "ready"; provider: ApiProvider; port: number }
|
||||
| ({ type: "error"; provider: ApiProvider; port: number } & SafeApiStartupError);
|
||||
|
||||
const SAFE_ERROR_CODES = new Set(["ERR_REQUIRE_ESM", "EBADENGINE", "EADDRINUSE", "EACCES", "ENOENT", "MODULE_NOT_FOUND", "ERR_MODULE_NOT_FOUND"]);
|
||||
export function safeApiErrorCode(code: unknown): string | undefined {
|
||||
return typeof code === "string" && SAFE_ERROR_CODES.has(code) ? code : undefined;
|
||||
}
|
||||
|
||||
/** Classification may inspect a message locally, but IPC never contains it. */
|
||||
export function getSafeApiStartupError(err: unknown): SafeApiStartupError {
|
||||
const error = (err ?? {}) as { code?: unknown; message?: unknown };
|
||||
const code = safeApiErrorCode(error.code);
|
||||
const message = typeof error.message === "string" ? error.message : "";
|
||||
let category: ApiStartupCategory = "startup";
|
||||
if (code === "ERR_REQUIRE_ESM" || /ERR_REQUIRE_ESM|require\(\) of ES ?Module/i.test(message)) category = "esm";
|
||||
else if (code === "EBADENGINE" || /Unsupported engine|EBADENGINE|requires Node|Node\.js version/i.test(message)) category = "node-engine";
|
||||
else if (code === "EADDRINUSE") category = "port-in-use";
|
||||
return code ? { category, code } : { category };
|
||||
}
|
||||
|
||||
export function isApiPortFree(port: number): Promise<boolean> {
|
||||
return new Promise((resolve) => {
|
||||
const server = net.createServer();
|
||||
server.once("error", () => server.close(() => resolve(false)));
|
||||
server.once("listening", () => server.close(() => resolve(true)));
|
||||
server.listen(port, "127.0.0.1");
|
||||
});
|
||||
}
|
||||
|
||||
function waitForListening(server: Server): Promise<void> {
|
||||
if (server.listening) return Promise.resolve();
|
||||
return new Promise((resolve, reject) => {
|
||||
const ready = () => { cleanup(); resolve(); };
|
||||
const failed = (error: Error) => { cleanup(); reject(error); };
|
||||
const cleanup = () => { server.off("listening", ready); server.off("error", failed); };
|
||||
server.once("listening", ready);
|
||||
server.once("error", failed);
|
||||
});
|
||||
}
|
||||
|
||||
export async function closeEmbeddedApi(server: Server | null): Promise<void> {
|
||||
if (!server) return;
|
||||
await new Promise<void>((resolve) => {
|
||||
try {
|
||||
server.close(() => resolve());
|
||||
server.closeAllConnections?.();
|
||||
} catch { resolve(); }
|
||||
});
|
||||
}
|
||||
|
||||
/** Only call in the isolated child: these dependencies write raw request URLs
|
||||
* and response cookies directly to console, outside the bot's logger. */
|
||||
export async function startEmbeddedApi(provider: ApiProvider, port: number): Promise<{ server: Server | null }> {
|
||||
let server: Server | null = null;
|
||||
try {
|
||||
if (provider === "netease") {
|
||||
const imported = await import("NeteaseCloudMusicApi") as any;
|
||||
const api = imported.server ?? imported.default?.server;
|
||||
const app = await api.serveNcmApi({ port, host: "127.0.0.1", checkVersion: false });
|
||||
server = app.server;
|
||||
if (!server) throw new Error("NetEase API did not expose its HTTP server");
|
||||
} else {
|
||||
const previousPort = process.env.PORT;
|
||||
process.env.PORT = String(port);
|
||||
let imported: any;
|
||||
try { imported = await import("@sansenjian/qq-music-api"); }
|
||||
finally {
|
||||
if (previousPort === undefined) delete process.env.PORT;
|
||||
else process.env.PORT = previousPort;
|
||||
}
|
||||
const candidate = imported.default ?? imported;
|
||||
const app = typeof candidate.listen === "function" ? candidate : candidate.default;
|
||||
if (!app || typeof app.listen !== "function") throw new Error("QQ API did not expose a Koa app");
|
||||
// Historical packages listened during import. Their listener remains
|
||||
// owned by this child and closes when the child exits.
|
||||
if (!(await isApiPortFree(port))) return { server: null };
|
||||
server = app.listen(port, "127.0.0.1");
|
||||
}
|
||||
await waitForListening(server!);
|
||||
return { server };
|
||||
} catch (error) {
|
||||
await closeEmbeddedApi(server);
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
+170
-112
@@ -1,133 +1,191 @@
|
||||
import { describe, it, expect, vi, beforeEach } from "vitest";
|
||||
import { EventEmitter } from "node:events";
|
||||
import type { ChildProcess } from "node:child_process";
|
||||
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
|
||||
import { createApiServerManager, describeQqApiStartupError } from "./api-server.js";
|
||||
import type { Logger } from "../logger.js";
|
||||
|
||||
// Record every listen() the QQ sidecar makes so we can assert it is always
|
||||
// pinned to the configured port (regression coverage for issue #122).
|
||||
const mockState = vi.hoisted(() => ({
|
||||
listenCalls: [] as Array<{ port: number; host: string }>,
|
||||
}));
|
||||
|
||||
vi.mock("@sansenjian/qq-music-api", () => {
|
||||
const app = {
|
||||
listen(port: number, host: string, cb?: () => void) {
|
||||
mockState.listenCalls.push({ port, host });
|
||||
const server = {
|
||||
address: () => ({ port, address: host, family: "IPv4" as const }),
|
||||
on() {
|
||||
return server;
|
||||
},
|
||||
close(done?: () => void) {
|
||||
done?.();
|
||||
},
|
||||
};
|
||||
// Real net/Koa fire the listening callback on a later tick, after the
|
||||
// caller has captured the returned server handle.
|
||||
if (cb) setImmediate(cb);
|
||||
return server;
|
||||
},
|
||||
};
|
||||
return { default: app };
|
||||
const state = vi.hoisted(() => ({ fork: vi.fn(), probes: [] as EventEmitter[], probeAutomatically: true, portFree: true, directImports: 0 }));
|
||||
vi.mock("node:child_process", () => ({ fork: state.fork }));
|
||||
vi.mock("node:net", async () => {
|
||||
const { EventEmitter } = await import("node:events");
|
||||
return { default: { createServer: () => {
|
||||
const probe = new EventEmitter() as EventEmitter & { close(done: () => void): void; listen(): void };
|
||||
probe.close = (done) => queueMicrotask(done);
|
||||
probe.listen = () => {
|
||||
state.probes.push(probe);
|
||||
if (state.probeAutomatically) queueMicrotask(() => probe.emit(state.portFree ? "listening" : "error"));
|
||||
};
|
||||
return probe;
|
||||
} } };
|
||||
});
|
||||
vi.mock("@sansenjian/qq-music-api", () => {
|
||||
state.directImports++;
|
||||
return { default: { listen: () => { throw new Error("sidecar imported in parent"); } } };
|
||||
});
|
||||
vi.mock("NeteaseCloudMusicApi", () => {
|
||||
state.directImports++;
|
||||
return { server: { serveNcmApi: () => { throw new Error("sidecar imported in parent"); } } };
|
||||
});
|
||||
|
||||
class FakeChild extends EventEmitter {
|
||||
connected = true;
|
||||
exitOnStop = true;
|
||||
send = vi.fn((message: { type: string }) => {
|
||||
if (message.type === "stop" && this.exitOnStop) queueMicrotask(() => this.finish(0, null));
|
||||
return true;
|
||||
});
|
||||
kill = vi.fn((signal: string = "SIGTERM") => { queueMicrotask(() => this.finish(null, signal)); return true; });
|
||||
finish(code: number | null, signal: string | null) { this.connected = false; this.emit("exit", code, signal); }
|
||||
}
|
||||
|
||||
describe("describeQqApiStartupError", () => {
|
||||
it("flags ERR_REQUIRE_ESM by error code with version-pin guidance", () => {
|
||||
const hint = describeQqApiStartupError({ code: "ERR_REQUIRE_ESM", message: "..." });
|
||||
expect(hint).toMatch(/ERR_REQUIRE_ESM/);
|
||||
expect(hint).toMatch(/~2\.4\.0/);
|
||||
expect(hint).toMatch(/~2\.2\.10/);
|
||||
it("retains ESM diagnostics by code and message", () => {
|
||||
expect(describeQqApiStartupError({ code: "ERR_REQUIRE_ESM" })).toMatch(/~2\.4\.0/);
|
||||
expect(describeQqApiStartupError(new Error("require() of ES Module is unsupported"))).toMatch(/ERR_REQUIRE_ESM/);
|
||||
});
|
||||
|
||||
it("flags ERR_REQUIRE_ESM by message when the code is absent", () => {
|
||||
const hint = describeQqApiStartupError(
|
||||
new Error("require() of ES Module .../@sansenjian/qq-music-api/dist/index.js not supported")
|
||||
);
|
||||
expect(hint).toMatch(/incompatible @sansenjian\/qq-music-api/);
|
||||
});
|
||||
|
||||
it("flags a Node engine mismatch with a Node-upgrade hint", () => {
|
||||
const hint = describeQqApiStartupError(new Error("Unsupported engine: requires Node >=20.17"));
|
||||
expect(hint).toMatch(/Node >=20\.17/);
|
||||
expect(hint).toMatch(/~2\.2\.10/);
|
||||
});
|
||||
|
||||
it("returns null for an unrelated startup error (falls back to the generic warning)", () => {
|
||||
expect(describeQqApiStartupError(new Error("EADDRINUSE: port in use"))).toBeNull();
|
||||
expect(describeQqApiStartupError(undefined)).toBeNull();
|
||||
expect(describeQqApiStartupError(null)).toBeNull();
|
||||
it("retains engine diagnostics and ignores unrelated failures", () => {
|
||||
expect(describeQqApiStartupError(new Error("Unsupported engine: requires Node >=20.17"))).toMatch(/Node >=20\.17/);
|
||||
expect(describeQqApiStartupError(new Error("EADDRINUSE"))).toBeNull();
|
||||
});
|
||||
});
|
||||
|
||||
// Regression coverage for issue #122: the QQ Music API sidecar must listen on
|
||||
// the same port the client base URL targets (config.qqMusicApiPort). A stale
|
||||
// build once bound 3300 while the client requested 3200, silently breaking the
|
||||
// QQ login QR / search flow with ECONNREFUSED on 127.0.0.1:3200.
|
||||
describe("createApiServerManager — QQ sidecar port binding", () => {
|
||||
const noopLogger = {
|
||||
info() {},
|
||||
warn() {},
|
||||
error() {},
|
||||
debug() {},
|
||||
trace() {},
|
||||
fatal() {},
|
||||
} as unknown as Logger;
|
||||
|
||||
describe("embedded API child lifecycle", () => {
|
||||
let children: FakeChild[];
|
||||
let logger: Logger;
|
||||
let manager: ReturnType<typeof createApiServerManager>;
|
||||
let automaticReady: boolean;
|
||||
const options = { neteasePort: 39218, qqMusicPort: 39217, neteaseEnabled: true, qqEnabled: true };
|
||||
const flush = async () => { for (let i = 0; i < 12; i++) await Promise.resolve(); };
|
||||
beforeEach(() => {
|
||||
mockState.listenCalls = [];
|
||||
vi.useRealTimers(); children = []; automaticReady = true;
|
||||
state.probes = []; state.portFree = true; state.probeAutomatically = true; state.fork.mockReset();
|
||||
state.fork.mockImplementation((_entry: string, args: string[]) => {
|
||||
const child = new FakeChild(); children.push(child);
|
||||
if (automaticReady) queueMicrotask(() => child.emit("message", { type: "ready", provider: args[0], port: Number(args[1]) }));
|
||||
return child as unknown as ChildProcess;
|
||||
});
|
||||
logger = { info: vi.fn(), warn: vi.fn(), error: vi.fn() } as unknown as Logger;
|
||||
manager = createApiServerManager(options, logger);
|
||||
});
|
||||
afterEach(async () => { manager.stop(); await flush(); vi.useRealTimers(); });
|
||||
|
||||
it("listens on the configured qqMusicPort and exposes a matching base URL", async () => {
|
||||
const port = 39217; // uncommon port to avoid clashing with a real instance
|
||||
const manager = createApiServerManager(
|
||||
{ neteasePort: 39218, qqMusicPort: port, neteaseEnabled: false, qqEnabled: true },
|
||||
noopLogger
|
||||
);
|
||||
it("isolates both APIs with ignored stdio, configured ports and IPC", async () => {
|
||||
await manager.start();
|
||||
manager.stop();
|
||||
|
||||
expect(manager.getQQMusicBaseUrl()).toBe(`http://127.0.0.1:${port}`);
|
||||
expect(mockState.listenCalls).toEqual([{ port, host: "127.0.0.1" }]);
|
||||
expect(state.fork).toHaveBeenCalledTimes(2);
|
||||
expect(state.fork.mock.calls.map((call) => call[1])).toEqual([["netease", "39218"], ["qq", "39217"]]);
|
||||
for (const call of state.fork.mock.calls) {
|
||||
expect(String(call[0])).toMatch(/api-server-child\.ts$/);
|
||||
expect(call[2].stdio).toEqual(["ignore", "ignore", "ignore", "ipc"]);
|
||||
expect(call[2].execArgv).not.toContain("--eval");
|
||||
expect(call[2].execArgv).not.toContain("--input-type=module");
|
||||
}
|
||||
expect(state.directImports).toBe(0);
|
||||
expect(manager.getNeteaseBaseUrl()).toBe("http://127.0.0.1:39218");
|
||||
expect(manager.getQQMusicBaseUrl()).toBe("http://127.0.0.1:39217");
|
||||
});
|
||||
|
||||
it("follows qqMusicPort — not an injected PORT — and restores PORT afterwards", async () => {
|
||||
const port = 39219;
|
||||
const previous = process.env.PORT;
|
||||
// Simulate a hosting platform / compose file injecting a stray PORT that
|
||||
// must NOT leak into the QQ sidecar's chosen port.
|
||||
process.env.PORT = "39999";
|
||||
const manager = createApiServerManager(
|
||||
{ neteasePort: 39220, qqMusicPort: port, neteaseEnabled: false, qqEnabled: true },
|
||||
noopLogger
|
||||
);
|
||||
it("preserves provider gating and externally bound port reuse", async () => {
|
||||
manager = createApiServerManager({ ...options, neteaseEnabled: false, qqEnabled: false }, logger);
|
||||
await manager.start(); expect(state.probes).toHaveLength(0); expect(state.fork).not.toHaveBeenCalled();
|
||||
state.portFree = false;
|
||||
manager = createApiServerManager({ ...options, neteaseEnabled: false }, logger);
|
||||
await manager.start(); expect(state.fork).not.toHaveBeenCalled();
|
||||
expect(logger.info).toHaveBeenCalledWith({ port: 39217 }, expect.stringContaining("reusing"));
|
||||
});
|
||||
it("inherits tsx loader arguments without unrelated parent runner flags", async () => {
|
||||
const previous = process.execArgv;
|
||||
process.execArgv = ["--require", "C:\\app\\node_modules\\tsx\\dist\\preflight.cjs", "--import", "file:///app/node_modules/tsx/dist/loader.mjs", "--eval", "synthetic-evaluation", "--conditions", "vitest", "--input-type=module", "--inspect"];
|
||||
try {
|
||||
await manager.start();
|
||||
// The sidecar follows qqMusicPort, never the injected PORT.
|
||||
expect(mockState.listenCalls).toEqual([{ port, host: "127.0.0.1" }]);
|
||||
// The injected PORT is restored so nothing else in the process is affected.
|
||||
expect(process.env.PORT).toBe("39999");
|
||||
} finally {
|
||||
manager.stop();
|
||||
if (previous === undefined) delete process.env.PORT;
|
||||
else process.env.PORT = previous;
|
||||
}
|
||||
expect(state.fork.mock.calls[0][2].execArgv).toEqual(process.execArgv.slice(0, 4));
|
||||
} finally { process.execArgv = previous; }
|
||||
});
|
||||
|
||||
it("leaves an absent PORT env unset after importing the sidecar", async () => {
|
||||
const port = 39221;
|
||||
const previous = process.env.PORT;
|
||||
delete process.env.PORT;
|
||||
const manager = createApiServerManager(
|
||||
{ neteasePort: 39222, qqMusicPort: port, neteaseEnabled: false, qqEnabled: true },
|
||||
noopLogger
|
||||
);
|
||||
try {
|
||||
await manager.start();
|
||||
// Was unset before importing — must be unset again, no leaked override.
|
||||
expect(process.env.PORT).toBeUndefined();
|
||||
} finally {
|
||||
manager.stop();
|
||||
if (previous === undefined) delete process.env.PORT;
|
||||
else process.env.PORT = previous;
|
||||
it("does not duplicate concurrent or repeated starts", async () => {
|
||||
await Promise.all([manager.start(), manager.start()]); await manager.start();
|
||||
expect(state.fork).toHaveBeenCalledTimes(2);
|
||||
});
|
||||
it("fences a stop during pending port preflight", async () => {
|
||||
state.probeAutomatically = false;
|
||||
const starting = manager.start(); await flush(); manager.stop();
|
||||
state.probes[0].emit("listening"); await starting;
|
||||
expect(state.fork).not.toHaveBeenCalled();
|
||||
});
|
||||
it("cancels a pending handshake and ignores its late ready", async () => {
|
||||
automaticReady = false;
|
||||
const starting = manager.start(); await flush(); expect(children).toHaveLength(1);
|
||||
manager.stop(); children[0].emit("message", { type: "ready", provider: "netease", port: 39218 }); await starting;
|
||||
expect(children[0].send).toHaveBeenCalledWith({ type: "stop" }, expect.any(Function));
|
||||
expect(state.fork).toHaveBeenCalledTimes(1);
|
||||
expect(logger.info).not.toHaveBeenCalledWith({ port: 39218 }, "NetEase Cloud Music API started");
|
||||
});
|
||||
it("waits for a cancelled preflight to release its probe before restart", async () => {
|
||||
state.probeAutomatically = false;
|
||||
const first = manager.start(); await flush(); manager.stop();
|
||||
const restarting = manager.start(); await flush();
|
||||
expect(state.probes).toHaveLength(1);
|
||||
state.probeAutomatically = true; state.probes[0].emit("listening");
|
||||
await Promise.all([first, restarting]);
|
||||
expect(state.fork).toHaveBeenCalledTimes(2);
|
||||
});
|
||||
it("waits for old children to exit before restart", async () => {
|
||||
await manager.start(); children.forEach((child) => { child.exitOnStop = false; }); manager.stop();
|
||||
const restarting = manager.start(); await flush(); expect(state.fork).toHaveBeenCalledTimes(2);
|
||||
children.slice(0, 2).forEach((child) => child.finish(0, null)); await restarting;
|
||||
expect(state.fork).toHaveBeenCalledTimes(4);
|
||||
});
|
||||
it("reports unexpected post-ready exits with safe fields", async () => {
|
||||
await manager.start(); children[1].finish(7, "SIGTERM");
|
||||
expect(logger.error).toHaveBeenCalledWith({ provider: "qq", port: 39217, code: 7, signal: "SIGTERM" }, expect.stringContaining("exited unexpectedly"));
|
||||
});
|
||||
it("retains static QQ diagnostics and discards arbitrary IPC fields", async () => {
|
||||
automaticReady = false; manager = createApiServerManager({ ...options, neteaseEnabled: false }, logger);
|
||||
const starting = manager.start(); await flush();
|
||||
children[0].emit("message", { type: "error", provider: "qq", port: 39217, category: "esm", code: "ERR_REQUIRE_ESM", message: "synthetic-credential", stack: "synthetic-credential" }); await starting;
|
||||
expect(logger.error).toHaveBeenCalledWith({ provider: "qq", port: 39217, category: "esm", code: "ERR_REQUIRE_ESM" }, expect.stringContaining("ERR_REQUIRE_ESM"));
|
||||
expect(JSON.stringify([...(logger.error as ReturnType<typeof vi.fn>).mock.calls, ...(logger.warn as ReturnType<typeof vi.fn>).mock.calls])).not.toContain("synthetic-credential");
|
||||
expect(children[0].send).toHaveBeenCalledWith({ type: "stop" }, expect.any(Function));
|
||||
});
|
||||
it("ignores a ready message for a different provider or port", async () => {
|
||||
automaticReady = false; manager = createApiServerManager({ ...options, neteaseEnabled: false }, logger);
|
||||
const starting = manager.start(); await flush();
|
||||
children[0].emit("message", { type: "ready", provider: "netease", port: 39217 });
|
||||
children[0].emit("message", { type: "ready", provider: "qq", port: 39999 });
|
||||
await flush();
|
||||
expect(logger.info).not.toHaveBeenCalledWith({ port: 39217 }, "QQ Music API started");
|
||||
children[0].emit("message", { type: "ready", provider: "qq", port: 39217 }); await starting;
|
||||
expect(logger.info).toHaveBeenCalledWith({ port: 39217 }, "QQ Music API started");
|
||||
});
|
||||
it("cleans up a failed fork that closes without an exit event", async () => {
|
||||
automaticReady = false; manager = createApiServerManager({ ...options, neteaseEnabled: false }, logger);
|
||||
const starting = manager.start(); await flush();
|
||||
children[0].emit("error", Object.assign(new Error("synthetic-credential"), { code: "ENOENT" }));
|
||||
children[0].emit("close", null, null); await starting;
|
||||
expect(logger.error).toHaveBeenCalledWith({ provider: "qq", port: 39217, category: "startup", code: "ENOENT" }, expect.stringContaining("start"));
|
||||
automaticReady = true; await manager.start();
|
||||
expect(state.fork).toHaveBeenCalledTimes(2);
|
||||
});
|
||||
it("times out and terminates a silent child", async () => {
|
||||
vi.useFakeTimers(); automaticReady = false; manager = createApiServerManager({ ...options, neteaseEnabled: false }, logger);
|
||||
const starting = manager.start(); await flush(); await vi.advanceTimersByTimeAsync(30000); await starting;
|
||||
expect(logger.error).toHaveBeenCalledWith({ provider: "qq", port: 39217, category: "timeout" }, expect.stringContaining("start"));
|
||||
expect(children[0].send).toHaveBeenCalledWith({ type: "stop" }, expect.any(Function));
|
||||
});
|
||||
it("forces shutdown if a child ignores stop", async () => {
|
||||
vi.useFakeTimers(); await manager.start(); children.forEach((child) => { child.exitOnStop = false; }); manager.stop();
|
||||
await vi.advanceTimersByTimeAsync(2000);
|
||||
expect(children.every((child) => child.kill.mock.calls.length > 0)).toBe(true);
|
||||
expect(logger.error).not.toHaveBeenCalled();
|
||||
});
|
||||
it("escalates to SIGKILL if stop and SIGTERM are ignored", async () => {
|
||||
vi.useFakeTimers(); await manager.start();
|
||||
for (const child of children) {
|
||||
child.exitOnStop = false;
|
||||
child.kill.mockImplementation((signal: string = "SIGTERM") => {
|
||||
if (signal === "SIGKILL") queueMicrotask(() => child.finish(null, signal));
|
||||
return true;
|
||||
});
|
||||
}
|
||||
manager.stop(); await vi.advanceTimersByTimeAsync(2000);
|
||||
expect(children.every((child) => child.kill.mock.calls.some(([signal]) => signal === "SIGKILL"))).toBe(true);
|
||||
expect(logger.error).not.toHaveBeenCalled();
|
||||
});
|
||||
});
|
||||
+160
-178
@@ -1,16 +1,13 @@
|
||||
import net from "node:net";
|
||||
import { fork, type ChildProcess } from "node:child_process";
|
||||
import type { Logger } from "../logger.js";
|
||||
import type { Server } from "node:http";
|
||||
import { getSafeApiStartupError, isApiPortFree, safeApiErrorCode, type ApiProvider, type SafeApiStartupError } from "./api-server-runtime.js";
|
||||
|
||||
export interface ApiServerOptions {
|
||||
neteasePort: number;
|
||||
qqMusicPort: number;
|
||||
/** Provider gating (#enabledProviders): when false, the corresponding
|
||||
* embedded sidecar API server is never started and its port never bound. */
|
||||
neteaseEnabled?: boolean;
|
||||
qqEnabled?: boolean;
|
||||
}
|
||||
|
||||
export interface ApiServerManager {
|
||||
start(): Promise<void>;
|
||||
stop(): void;
|
||||
@@ -18,196 +15,181 @@ export interface ApiServerManager {
|
||||
getQQMusicBaseUrl(): string;
|
||||
}
|
||||
|
||||
/**
|
||||
* Classify a QQ Music API (@sansenjian/qq-music-api) startup failure into
|
||||
* actionable operator guidance, or null when it isn't a recognised
|
||||
* dependency/runtime mismatch. Exported for testing.
|
||||
*
|
||||
* Background: the package became ESM in 2.3.x. A loose `^` range could pull an
|
||||
* ESM-only build (2.3.0/2.3.1) that throws ERR_REQUIRE_ESM, or a 2.4.x build
|
||||
* that needs Node >=20.17 — either way the embedded server never binds, so
|
||||
* every QQ request fails downstream with ECONNREFUSED on the API port.
|
||||
*/
|
||||
export function describeQqApiStartupError(err: unknown): string | null {
|
||||
const e = (err ?? {}) as { code?: string; message?: string };
|
||||
const code = String(e.code ?? "");
|
||||
const msg = String(e.message ?? "");
|
||||
if (code === "ERR_REQUIRE_ESM" || /ERR_REQUIRE_ESM|require\(\) of ES ?Module/i.test(msg)) {
|
||||
return (
|
||||
"an incompatible @sansenjian/qq-music-api build is installed (ERR_REQUIRE_ESM). " +
|
||||
"Pin it to ~2.4.0 (needs Node >=20.17) or ~2.2.10 in package.json, then reinstall"
|
||||
);
|
||||
}
|
||||
if (/Unsupported engine|EBADENGINE|requires Node|Node\.js version/i.test(msg)) {
|
||||
return "@sansenjian/qq-music-api 2.4.x requires Node >=20.17 (or >=22.9) — upgrade Node, or pin the package to ~2.2.10";
|
||||
}
|
||||
const { category } = getSafeApiStartupError(err);
|
||||
if (category === "esm") return "an incompatible @sansenjian/qq-music-api build is installed (ERR_REQUIRE_ESM). Pin it to ~2.4.0 (needs Node >=20.17) or ~2.2.10 in package.json, then reinstall";
|
||||
if (category === "node-engine") return "@sansenjian/qq-music-api 2.4.x requires Node >=20.17 (or >=22.9) — upgrade Node, or pin the package to ~2.2.10";
|
||||
return null;
|
||||
}
|
||||
|
||||
function isPortFree(port: number): Promise<boolean> {
|
||||
return new Promise((resolve) => {
|
||||
const server = net.createServer();
|
||||
server.once("error", () => {
|
||||
server.close(() => resolve(false));
|
||||
});
|
||||
server.once("listening", () => {
|
||||
server.close(() => resolve(true));
|
||||
});
|
||||
server.listen(port, "127.0.0.1");
|
||||
});
|
||||
/** Carry only tsx loader arguments into a source child. CLI evaluation,
|
||||
* inspector and test-runner flags have unrelated meanings in a fork. */
|
||||
function childExecArgv(source: boolean): string[] {
|
||||
if (!source) return [];
|
||||
const args: string[] = [];
|
||||
for (let i = 0; i < process.execArgv.length; i++) {
|
||||
const arg = process.execArgv[i];
|
||||
if (arg === "--import" || arg === "--require" || arg === "-r") {
|
||||
const value = process.execArgv[++i];
|
||||
if (value && (value === "tsx" || /[/\\]tsx[/\\]/.test(value))) args.push(arg, value);
|
||||
} else if (arg.startsWith("--import=") && (arg === "--import=tsx" || /[/\\]tsx[/\\]/.test(arg))) args.push(arg);
|
||||
}
|
||||
return args.length ? args : ["--import", "tsx"];
|
||||
}
|
||||
|
||||
export function createApiServerManager(
|
||||
options: ApiServerOptions,
|
||||
logger: Logger
|
||||
): ApiServerManager {
|
||||
let neteaseServer: Server | null = null;
|
||||
let qqMusicServer: Server | null = null;
|
||||
class StartupFailure extends Error {
|
||||
constructor(readonly details: SafeApiStartupError) { super("Embedded music API startup failed"); }
|
||||
}
|
||||
interface ManagedChild {
|
||||
child: ChildProcess;
|
||||
stop(): Promise<void>;
|
||||
}
|
||||
const STARTUP_TIMEOUT_MS = 15000;
|
||||
const ERROR_CATEGORIES = new Set(["esm", "node-engine", "port-in-use", "startup"]);
|
||||
|
||||
const neteaseBaseUrl = `http://127.0.0.1:${options.neteasePort}`;
|
||||
const qqMusicBaseUrl = `http://127.0.0.1:${options.qqMusicPort}`;
|
||||
export function createApiServerManager(options: ApiServerOptions, logger: Logger): ApiServerManager {
|
||||
const children = new Map<ApiProvider, ManagedChild>();
|
||||
let generation = 0;
|
||||
let starting: Promise<void> | null = null;
|
||||
let stopping: Promise<void> = Promise.resolve();
|
||||
|
||||
function launch(provider: ApiProvider, port: number, launchGeneration: number): Promise<void> {
|
||||
const source = import.meta.url.endsWith(".ts");
|
||||
const entry = new URL(source ? "./api-server-child.ts" : "./api-server-child.js", import.meta.url);
|
||||
const child = fork(entry, [provider, String(port)], {
|
||||
stdio: ["ignore", "ignore", "ignore", "ipc"],
|
||||
execArgv: childExecArgv(source),
|
||||
});
|
||||
let ready = false;
|
||||
let expectedExit = false;
|
||||
let settled = false;
|
||||
let hasExited = false;
|
||||
let startupTimer: ReturnType<typeof setTimeout>;
|
||||
let terminateTimer: ReturnType<typeof setTimeout> | undefined;
|
||||
let killTimer: ReturnType<typeof setTimeout> | undefined;
|
||||
let resolveExit!: () => void;
|
||||
const exited = new Promise<void>((resolve) => { resolveExit = resolve; });
|
||||
let resolveStart!: () => void;
|
||||
let rejectStart!: (error: StartupFailure) => void;
|
||||
const started = new Promise<void>((resolve, reject) => { resolveStart = resolve; rejectStart = reject; });
|
||||
const settle = (error?: SafeApiStartupError) => {
|
||||
if (settled) return;
|
||||
settled = true;
|
||||
clearTimeout(startupTimer);
|
||||
if (error) rejectStart(new StartupFailure(error)); else resolveStart();
|
||||
};
|
||||
const record: ManagedChild = {
|
||||
child,
|
||||
stop() {
|
||||
if (expectedExit) return exited;
|
||||
expectedExit = true;
|
||||
settle({ category: "cancelled" });
|
||||
if (children.get(provider) === record) children.delete(provider);
|
||||
stopping = Promise.all([stopping, exited]).then(() => {});
|
||||
if (hasExited) return exited;
|
||||
try {
|
||||
if (child.connected) child.send({ type: "stop" }, (error) => { if (error) child.kill("SIGTERM"); });
|
||||
else child.kill("SIGTERM");
|
||||
} catch { child.kill("SIGTERM"); }
|
||||
terminateTimer = setTimeout(() => child.kill("SIGTERM"), 1000);
|
||||
killTimer = setTimeout(() => child.kill("SIGKILL"), 2000);
|
||||
terminateTimer.unref(); killTimer.unref();
|
||||
return exited;
|
||||
},
|
||||
};
|
||||
children.set(provider, record);
|
||||
child.on("message", (message: unknown) => {
|
||||
if (!message || typeof message !== "object" || expectedExit || launchGeneration !== generation) return;
|
||||
const data = message as Record<string, unknown>;
|
||||
if (data.provider !== provider || data.port !== port) return;
|
||||
if (data.type === "ready") { ready = true; settle(); }
|
||||
else if (data.type === "error" && typeof data.category === "string" && ERROR_CATEGORIES.has(data.category)) {
|
||||
const code = safeApiErrorCode(data.code);
|
||||
const details: SafeApiStartupError = { category: data.category as SafeApiStartupError["category"], ...(code ? { code } : {}) };
|
||||
if (!ready) settle(details);
|
||||
else logger.error({ provider, port, ...details }, "Embedded music API reported a runtime failure");
|
||||
void record.stop();
|
||||
}
|
||||
});
|
||||
child.on("error", (error) => {
|
||||
if (expectedExit) return;
|
||||
const details = getSafeApiStartupError(error);
|
||||
if (!ready) settle(details);
|
||||
else logger.error({ provider, port, ...details }, "Embedded music API child failed");
|
||||
void record.stop();
|
||||
});
|
||||
const onExit = (code: number | null, signal: NodeJS.Signals | null) => {
|
||||
if (hasExited) return;
|
||||
hasExited = true;
|
||||
clearTimeout(startupTimer); clearTimeout(terminateTimer); clearTimeout(killTimer);
|
||||
if (children.get(provider) === record) children.delete(provider);
|
||||
if (!expectedExit) {
|
||||
if (ready) logger.error({ provider, port, code, signal }, "Embedded music API exited unexpectedly");
|
||||
else settle({ category: "startup" });
|
||||
}
|
||||
resolveExit();
|
||||
};
|
||||
child.once("exit", onExit);
|
||||
// A failed fork emits close without exit.
|
||||
child.once("close", onExit);
|
||||
startupTimer = setTimeout(() => { settle({ category: "timeout" }); void record.stop(); }, STARTUP_TIMEOUT_MS);
|
||||
return started;
|
||||
}
|
||||
|
||||
return {
|
||||
async start(): Promise<void> {
|
||||
// Provider gating: with the jellyfin-only default config neither legacy
|
||||
// sidecar starts, so ports 3001/3200 are never opened.
|
||||
if (options.neteaseEnabled === false && options.qqEnabled === false) {
|
||||
logger.info("NetEase/QQ providers disabled — embedded music API servers not started");
|
||||
return;
|
||||
}
|
||||
logger.info("Starting embedded music API servers...");
|
||||
|
||||
// Start NetEase Cloud Music API
|
||||
if (options.neteaseEnabled !== false) {
|
||||
try {
|
||||
const portFree = await isPortFree(options.neteasePort);
|
||||
if (!portFree) {
|
||||
logger.info(
|
||||
{ port: options.neteasePort },
|
||||
"NetEase API port already in use — reusing existing instance"
|
||||
);
|
||||
} else {
|
||||
const ncmModule = await import("NeteaseCloudMusicApi") as any;
|
||||
const serverObj = ncmModule.server ?? ncmModule.default?.server;
|
||||
const app = await serverObj.serveNcmApi({ port: options.neteasePort });
|
||||
neteaseServer = app;
|
||||
logger.info(
|
||||
{ port: options.neteasePort },
|
||||
"NetEase Cloud Music API started"
|
||||
);
|
||||
}
|
||||
} catch (err) {
|
||||
logger.error({ err }, "Failed to start NetEase Cloud Music API");
|
||||
start(): Promise<void> {
|
||||
if (starting) return starting;
|
||||
const startGeneration = generation;
|
||||
const run = async () => {
|
||||
await stopping;
|
||||
if (startGeneration !== generation) return;
|
||||
if (options.neteaseEnabled === false && options.qqEnabled === false) {
|
||||
logger.info("NetEase/QQ providers disabled — embedded music API servers not started"); return;
|
||||
}
|
||||
}
|
||||
|
||||
// Start QQ Music API. Older versions auto-started on import; the
|
||||
// current fork (2.2.11+) only listens when run as `require.main`,
|
||||
// so we explicitly call .listen() on the imported Koa app and keep
|
||||
// the server handle for clean shutdown.
|
||||
if (options.qqEnabled === false) return;
|
||||
try {
|
||||
const portFree = await isPortFree(options.qqMusicPort);
|
||||
if (!portFree) {
|
||||
logger.info(
|
||||
{ port: options.qqMusicPort },
|
||||
"QQ Music API port already in use — reusing existing instance"
|
||||
);
|
||||
} else {
|
||||
// Pin the upstream server to the configured port before importing.
|
||||
// The package derives its default port from process.env.PORT (falling
|
||||
// back to 3200) and, in some historical versions, auto-started that
|
||||
// server as an import side effect. Aligning PORT with qqMusicApiPort
|
||||
// guarantees the sidecar can never bind a different port than the one
|
||||
// the client base URL (getQQMusicBaseUrl) targets — the root cause of
|
||||
// issue #122, where an old build listened on 3300 while the client
|
||||
// requested 3200. Restore the previous value right after import so we
|
||||
// never leak the override into the rest of the process (e.g. the web
|
||||
// server or the NetEase sidecar, which also read PORT as a fallback).
|
||||
const prevPortEnv = process.env.PORT;
|
||||
process.env.PORT = String(options.qqMusicPort);
|
||||
let qqModule: any;
|
||||
logger.info("Starting embedded music API servers...");
|
||||
const providers: Array<{ provider: ApiProvider; port: number; enabled: boolean; name: string }> = [
|
||||
{ provider: "netease", port: options.neteasePort, enabled: options.neteaseEnabled !== false, name: "NetEase Cloud Music" },
|
||||
{ provider: "qq", port: options.qqMusicPort, enabled: options.qqEnabled !== false, name: "QQ Music" },
|
||||
];
|
||||
for (const { provider, port, enabled, name } of providers) {
|
||||
if (startGeneration !== generation) return;
|
||||
if (!enabled || children.has(provider)) continue;
|
||||
try {
|
||||
qqModule = (await import("@sansenjian/qq-music-api")) as any;
|
||||
} finally {
|
||||
if (prevPortEnv === undefined) delete process.env.PORT;
|
||||
else process.env.PORT = prevPortEnv;
|
||||
}
|
||||
// The module's export structure varies between versions:
|
||||
// 2.2.11+: default → Koa app (has .listen)
|
||||
// 2.2.10: default → wrapper object whose .default is the Koa app
|
||||
// older: module itself may be the Koa app
|
||||
const candidate = qqModule.default ?? qqModule;
|
||||
const koaApp = typeof candidate.listen === "function"
|
||||
? candidate
|
||||
: candidate.default ?? null;
|
||||
if (koaApp && typeof koaApp.listen === "function") {
|
||||
// A version that auto-started on import has already bound the
|
||||
// configured port (thanks to the PORT alignment above); reuse it
|
||||
// rather than racing a second listen that would fail EADDRINUSE.
|
||||
const stillFree = await isPortFree(options.qqMusicPort);
|
||||
if (!stillFree) {
|
||||
logger.info(
|
||||
{ port: options.qqMusicPort },
|
||||
"QQ Music API already listening on the configured port (auto-started on import) — reusing embedded instance"
|
||||
);
|
||||
} else {
|
||||
qqMusicServer = await new Promise<Server>((resolve, reject) => {
|
||||
const srv = koaApp.listen(options.qqMusicPort, "127.0.0.1", () =>
|
||||
resolve(srv)
|
||||
);
|
||||
srv.on("error", reject);
|
||||
});
|
||||
// Log the port actually bound (read from the socket) rather than
|
||||
// the requested one, so operators can spot a mismatch in the logs.
|
||||
const addr = qqMusicServer.address();
|
||||
const boundPort =
|
||||
addr && typeof addr === "object" && addr !== null
|
||||
? addr.port
|
||||
: options.qqMusicPort;
|
||||
logger.info(
|
||||
{ port: boundPort },
|
||||
"QQ Music API started"
|
||||
);
|
||||
const free = await isApiPortFree(port);
|
||||
if (startGeneration !== generation) return;
|
||||
if (!free) {
|
||||
logger.info({ port }, `${provider === "netease" ? "NetEase" : "QQ Music"} API port already in use — reusing existing instance`);
|
||||
continue;
|
||||
}
|
||||
} else {
|
||||
logger.warn("QQ Music API module does not expose a Koa app");
|
||||
await launch(provider, port, startGeneration);
|
||||
if (startGeneration !== generation) return;
|
||||
logger.info({ port }, `${name} API started`);
|
||||
} catch (error) {
|
||||
if (startGeneration !== generation) return;
|
||||
const details = error instanceof StartupFailure ? error.details : getSafeApiStartupError(error);
|
||||
if (details.category === "cancelled") return;
|
||||
const hint = provider === "qq" ? describeQqApiStartupError(details.category === "esm" ? { code: "ERR_REQUIRE_ESM" } : details.category === "node-engine" ? { code: "EBADENGINE" } : {}) : null;
|
||||
logger.error({ provider, port, ...details }, hint ? `QQ Music API failed to start — ${hint}. QQ features (search/play/login) will be unavailable until fixed; port ${port} is down.` : `Failed to start ${name} API`);
|
||||
}
|
||||
}
|
||||
} catch (err) {
|
||||
const hint = describeQqApiStartupError(err);
|
||||
if (hint) {
|
||||
logger.error(
|
||||
{ err },
|
||||
`QQ Music API failed to start — ${hint}. QQ features (search/play/login) will be unavailable until fixed; port ${options.qqMusicPort} is down.`
|
||||
);
|
||||
} else {
|
||||
logger.warn(
|
||||
{ err },
|
||||
"QQ Music API not available — QQ Music features may be limited"
|
||||
);
|
||||
}
|
||||
}
|
||||
};
|
||||
const promise = run();
|
||||
starting = promise;
|
||||
void promise.finally(() => { if (starting === promise) starting = null; });
|
||||
return promise;
|
||||
},
|
||||
|
||||
stop(): void {
|
||||
generation++;
|
||||
const pendingStart = starting;
|
||||
starting = null;
|
||||
logger.info("Stopping music API servers");
|
||||
if (neteaseServer && typeof (neteaseServer as any).close === "function") {
|
||||
(neteaseServer as any).close();
|
||||
}
|
||||
neteaseServer = null;
|
||||
if (qqMusicServer && typeof (qqMusicServer as any).close === "function") {
|
||||
(qqMusicServer as any).close();
|
||||
}
|
||||
qqMusicServer = null;
|
||||
},
|
||||
|
||||
getNeteaseBaseUrl(): string {
|
||||
return neteaseBaseUrl;
|
||||
},
|
||||
|
||||
getQQMusicBaseUrl(): string {
|
||||
return qqMusicBaseUrl;
|
||||
const retiring = [...children.values()];
|
||||
children.clear();
|
||||
// A cancelled preflight still owns a temporary listening socket until
|
||||
// its callback closes it. Restart must wait for that work as well.
|
||||
stopping = Promise.all([stopping, pendingStart, ...retiring.map((record) => record.stop())]).then(() => {});
|
||||
},
|
||||
getNeteaseBaseUrl: () => `http://127.0.0.1:${options.neteasePort}`,
|
||||
getQQMusicBaseUrl: () => `http://127.0.0.1:${options.qqMusicPort}`,
|
||||
};
|
||||
}
|
||||
Reference in new issue
Block a user