Compare commits

...
4 Commits
16 changed files with 2346 additions and 415 deletions

No files matched your search

+20 -1
View File
@@ -953,7 +953,26 @@ 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.2:长视频续播与语音发送故障处理
- 修复播放停滞的误判:一帧音频成功发送后,即使缓冲恰好被读空,也会重置断流计数。真实断流仍保留超时检测。
- B站 CDN 按实际域名识别,补齐 `szbdyd.com` 等回退节点的请求头和断点定位,避免续播时从头下载、解码长视频。
- B站提前结束后,获取新播放地址失败会在 1 秒、2 秒后重试,单次恢复最多查询 3 次;期间的切歌、停止、暂停和断开连接均会保留用户意图。诊断日志不记录原始地址或凭据。
- 语音发送异常不再被静默忽略:短暂失败后可继续发送;连续失败 2 秒后暂停,保留曲目与进度。处理网络或连接问题后可手动继续播放,持续故障日志有输出上限。
- 开启“频道无人时自动暂停”后,过期的频道人数查询和旧连接的空闲定时器不会再错误暂停或断开当前播放。
- 修复重连期间旧频道查询和 DNS 结果覆盖新连接的竞态,避免已经连接的机器人被旧请求移回频道或记录错误的服务器地址。
无配置或数据库迁移。此版本修复了已复现的故障路径;[#161](https://github.com/ZHANGTIANYAO1/teamspeak-music-bot/issues/161) 中新分支“67 分钟视频播到 37 分钟停止”的现场原因仍缺少停止时日志,尚未确认。
### 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))**
+83
View File
@@ -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");
});
});
+93
View File
@@ -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() };
}
+407 -3
View File
@@ -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 "";
@@ -30,12 +40,56 @@ describe("buildFfmpegArgs", () => {
expect(headers).toContain("User-Agent: Mozilla/5.0");
});
it.each([
"https://bilivideo.com/audio.m4s",
"https://upos-sz-mirrorcos.bilivideo.com/audio.m4s",
"https://bilivideo.cn/audio.m4s",
"https://cn-example-live-01.bilivideo.cn/audio.m4s",
"https://bilibili.com/audio.m4s",
"https://www.bilibili.com/audio.m4s",
"https://szbdyd.com/audio.m4s",
"https://stream.mcdn.szbdyd.com/audio.m4s",
"https://xy219x131x72x38xy.mcdn.bilivideo.cn.szbdyd.com/audio.m4s",
"HTTPS://UPOS-SZ-MIRRORCOS.BILIVIDEO.COM:443/audio.m4s",
"https://upos-sz-mirrorcos.bilivideo.com./audio.m4s",
])("uses Bilibili headers and input-side seeking for the actual CDN host in %s", (url) => {
const args = buildFfmpegArgs(url, 2220);
const headers = getHeadersArg(args);
expect(headers).toContain("Referer: https://www.bilibili.com");
expect(headers).toContain("User-Agent: Mozilla/5.0");
expect(args[args.indexOf("-ss") + 1]).toBe("2220");
expect(args.indexOf("-ss")).toBeLessThan(args.indexOf("-i"));
expect(args.lastIndexOf("-ss")).toBe(args.indexOf("-ss"));
expect(args).toContain("-reconnect_at_eof");
});
it.each([
"https://bilivideo.com.evil.example/audio.m4s",
"https://evilbilivideo.com/audio.m4s",
"https://bilivideo.cn.evil.example/audio.m4s",
"https://bilibili.com.evil.example/audio.m4s",
"https://evilbilibili.com/audio.m4s",
"https://szbdyd.com.evil.example/audio.m4s",
"https://evil.example/audio.m4s?redirect=https://upos-sz-mirrorcos.bilivideo.com/x",
"https://bilibili.com@evil.example/audio.m4s",
])("does not trust Bilibili text outside an allowed hostname in %s", (url) => {
const args = buildFfmpegArgs(url, 2220);
expect(args).not.toContain("-headers");
expect(args.indexOf("-ss")).toBeGreaterThan(args.indexOf("-i"));
});
it("does not set custom headers for unknown URLs", () => {
const url = "https://example.com/song.mp3";
const args = buildFfmpegArgs(url, 0);
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 +278,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 }
@@ -566,7 +863,7 @@ describe("AudioPlayer stall/EOF end-detection is gated on playing state (R3-4)",
// PCM frame (unknown duration -> isNearEnd forced true). startFrameLoop() runs
// the genuine loop; no real process is spawned (fake ffmpeg has no pid, so the
// end path never touches forceCleanup/process.kill).
function makeStalledPlaying(): AudioPlayer {
function makeStalledPlaying(duration = 0): AudioPlayer {
const player = new AudioPlayer(silentLogger);
const p = player as unknown as {
ffmpeg: unknown;
@@ -578,7 +875,7 @@ describe("AudioPlayer stall/EOF end-detection is gated on playing state (R3-4)",
startFrameLoop(): void;
};
p.ffmpeg = { pid: undefined }; // live ffmpeg, but delivers no PCM
p.currentSongDuration = 0; // unknown duration -> isNearEnd === true
p.currentSongDuration = duration;
p.pcmBuffer = Buffer.alloc(0); // always < one PCM frame
p.emptyFrameAttempts = 0;
p.framesPlayed = 0;
@@ -587,6 +884,113 @@ describe("AudioPlayer stall/EOF end-detection is gated on playing state (R3-4)",
return player;
}
it.each([0, 4020])("keeps emitting one full PCM frame per tick beyond 60 seconds with duration %s", (duration) => {
vi.useFakeTimers(FAKE_TIMER_OPTS);
const player = makeStalledPlaying(duration);
let ended = 0;
let frames = 0;
player.on("trackEnd", () => ended++);
player.on("frame", () => frames++);
try {
const internal = player as unknown as { pcmBuffer: Buffer };
const frame = Buffer.alloc(FRAME_BYTES);
for (let tick = 0; tick < 3100; tick++) {
internal.pcmBuffer = frame;
vi.advanceTimersByTime(20);
}
expect(frames).toBe(3100);
expect(player.getElapsed()).toBe(62);
expect(ended).toBe(0);
expect(player.getState()).toBe("playing");
} finally {
player.stop();
vi.useRealTimers();
}
});
it("resets the true-underrun budget after a successful frame even when no PCM reserve remains", () => {
vi.useFakeTimers(FAKE_TIMER_OPTS);
const player = makeStalledPlaying();
let ended = 0;
player.on("trackEnd", () => ended++);
try {
vi.advanceTimersByTime(20 * 249);
(player as unknown as { pcmBuffer: Buffer }).pcmBuffer = Buffer.alloc(FRAME_BYTES);
vi.advanceTimersByTime(20);
expect(ended).toBe(0);
expect(player.getElapsed()).toBe(0.02);
vi.advanceTimersByTime(20 * 249);
expect(ended).toBe(0);
vi.advanceTimersByTime(20);
expect(ended).toBe(1);
expect(player.getState()).toBe("idle");
} finally {
player.stop();
vi.useRealTimers();
}
});
it("still ends a genuine far-from-end stall after 60 seconds without a frame", () => {
vi.useFakeTimers(FAKE_TIMER_OPTS);
const player = makeStalledPlaying(4020);
let ended = 0;
player.on("trackEnd", () => ended++);
try {
vi.advanceTimersByTime(59980);
expect(ended).toBe(0);
expect(player.getState()).toBe("playing");
vi.advanceTimersByTime(20);
expect(ended).toBe(1);
expect(player.getState()).toBe("idle");
} finally {
player.stop();
vi.useRealTimers();
}
});
it("still emits natural EOF after delivering the final buffered frame", () => {
vi.useFakeTimers(FAKE_TIMER_OPTS);
const player = makeStalledPlaying(4020);
let ended = 0;
let frames = 0;
player.on("trackEnd", () => ended++);
player.on("frame", () => frames++);
try {
Object.assign(player, { ffmpeg: null, pcmBuffer: Buffer.alloc(FRAME_BYTES) });
vi.advanceTimersByTime(20);
expect(frames).toBe(1);
expect(ended).toBe(1);
expect(player.getState()).toBe("idle");
} finally {
player.stop();
vi.useRealTimers();
}
});
it("does not count a consumed frame as successful when the encoder throws", () => {
vi.useFakeTimers(FAKE_TIMER_OPTS);
const player = makeStalledPlaying(4020);
const failure = new Error("offline encoder failure");
const errors: Error[] = [];
let frames = 0;
player.on("error", error => errors.push(error));
player.on("frame", () => frames++);
try {
Object.assign(player, {
pcmBuffer: Buffer.alloc(FRAME_BYTES),
encoder: { encode() { throw failure; } },
});
vi.advanceTimersByTime(20);
expect(errors).toEqual([failure]);
expect(frames).toBe(0);
expect(player.getElapsed()).toBe(0);
expect((player as unknown as { emptyFrameAttempts: number }).emptyFrameAttempts).toBe(1);
} finally {
player.stop();
vi.useRealTimers();
}
});
it("does NOT emit trackEnd (and stays paused) when a stalled unknown-duration stream is paused past the stall threshold", () => {
vi.useFakeTimers(FAKE_TIMER_OPTS);
try {
+72 -13
View File
@@ -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,9 +86,21 @@ export function cleanupTempDir(dir: string): void {
}
export function buildFfmpegArgs(url: string, seekSeconds: number): string[] {
const args: string[] = [];
const isHttp = /^https?:\/\//i.test(url);
const isBilibili = isHttp && (url.includes("bilivideo") || url.includes("bilibili"));
const args: string[] = ["-nostats"];
let httpHostname: string | null = null;
try {
const parsed = new URL(url);
if (parsed.protocol === "http:" || parsed.protocol === "https:") {
httpHostname = parsed.hostname.toLowerCase().replace(/\.$/, "");
}
} catch {
// Local paths and malformed URLs must not inherit CDN-specific options.
}
const isHttp = httpHostname !== null;
const isBilibili = httpHostname !== null &&
["bilivideo.com", "bilivideo.cn", "bilibili.com", "szbdyd.com"].some(
domain => httpHostname === domain || httpHostname.endsWith(`.${domain}`),
);
if (isBilibili) {
args.push(
@@ -167,6 +191,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 +284,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 +309,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 +316,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 +370,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 +379,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 +390,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 +426,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 +448,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 +455,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 +594,7 @@ export class AudioPlayer extends EventEmitter {
// 立即清空缓冲区,确保切歌瞬间静音 (
this.pcmBuffer = Buffer.alloc(0);
this.ffmpegDiagnostics = null;
if (this.ffmpeg) {
const procToKill = this.ffmpeg;
@@ -557,6 +609,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 +633,7 @@ export class AudioPlayer extends EventEmitter {
}
private forceCleanup(proc: ChildProcess, pid: number): void {
this.intentionalCleanup.add(proc);
if (!globalActivePids.has(pid)) return;
try {
@@ -620,8 +674,10 @@ export class AudioPlayer extends EventEmitter {
// 这里的校验能防止旧的定时器回调处理新 Session 的逻辑 (
if (loopSessionId !== this.sessionId || !this.frameLoopRunning) return;
const framesBeforeTick = this.framesPlayed;
if (this.state === "playing") this.sendNextFrame();
else if (this.state === "paused") this.nextFrameTime = performance.now();
const frameSent = this.framesPlayed > framesBeforeTick;
// 检测pcmBuffer不足PCM_FRAME_BYTES导致连续循环卡死:
// 条件1: FFmpeg仍在运行但缓冲区不足一帧,且连续多次无法获取数据
@@ -641,7 +697,9 @@ export class AudioPlayer extends EventEmitter {
// unknown-duration stream would auto-advance ~5s later. Because the if is
// now false while paused, the else resets emptyFrameAttempts to 0, so a
// resumed healthy stream starts fresh and never ends instantly.
if (this.state === "playing" && !this.externalMode && this.ffmpeg !== null && this.pcmBuffer.length < PCM_FRAME_BYTES) {
// A healthy paced source can supply exactly one frame per tick, leaving
// no reserve after sendNextFrame. Count only ticks without emitted audio.
if (this.state === "playing" && !this.externalMode && !frameSent && this.ffmpeg !== null && this.pcmBuffer.length < PCM_FRAME_BYTES) {
this.emptyFrameAttempts++;
// End the track when FFmpeg has gone silent: quickly if we're near the
@@ -664,6 +722,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"
+393
View File
@@ -0,0 +1,393 @@
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
import { EventEmitter } from "node:events";
import { BotInstance } from "./instance.js";
import { AudioPlayer } from "../audio/player.js";
import { PlayQueue } from "../audio/queue.js";
import { ManagedVoiceClientRegistry } from "./managed-voice-clients.js";
function deferred<T>() {
let resolve!: (value: T) => void;
const promise = new Promise<T>(done => { resolve = done; });
return { promise, resolve };
}
async function flush() {
await new Promise<void>(resolve => setImmediate(resolve));
}
function makeBot() {
const queue = new PlayQueue();
queue.add({ id: "BV1example", name: "Long", artist: "A", album: "", coverUrl: "", platform: "bilibili", duration: 10_000, url: "old" });
queue.playAt(0);
const player = new EventEmitter() as any;
player.state = "idle";
player.sessionId = 1;
player.elapsed = 1000;
player.getState = AudioPlayer.prototype.getState;
player.getPlaybackSessionId = AudioPlayer.prototype.getPlaybackSessionId;
player.getElapsed = () => player.elapsed;
player.getVolume = () => 50;
player.pause = AudioPlayer.prototype.pause;
player.resume = AudioPlayer.prototype.resume;
player.play = vi.fn((_url: string, position: number) => { player.sessionId++; player.elapsed = position; player.state = "playing"; });
player.stop = vi.fn(() => { player.sessionId++; player.state = "idle"; });
const provider = { getSongUrl: vi.fn(async (): Promise<{ url: string } | null> => ({ url: "fresh" })) };
const tsClient = new EventEmitter() as any;
tsClient.getClientsInChannel = vi.fn(async () => [{ id: 1 }, { id: 2 }]);
tsClient.getClientId = () => 1;
tsClient.sendVoiceData = vi.fn();
tsClient.connect = vi.fn(async () => { tsClient.emit("connected"); });
tsClient.disconnect = vi.fn(() => { tsClient.emit("disconnected"); });
tsClient.getResolvedVoiceEndpoint = () => null;
const spotifyController = new EventEmitter() as any;
spotifyController.stop = vi.fn();
const bot = Object.assign(Object.create(BotInstance.prototype), {
id: "bot", name: "Bot", queue, player, tsClient, provider, spotifyController,
connected: true, disconnectEmitted: false, effectiveDuration: 10_000,
streamRecovery: null, lifecycleGeneration: 0, occupancyRequest: 0,
config: { autoPauseOnEmpty: true, idleTimeoutMinutes: 1 }, autoPaused: false,
idleTimer: null, snapshotTimer: null, currentSourceIsSpotify: false,
localProvider: {}, voiceDucking: { reset: vi.fn(), removeSpeaker: vi.fn() },
managedVoiceClients: new ManagedVoiceClientRegistry(),
configuredVoiceServerScope: { host: "localhost", voicePort: 9987 },
profileManager: { onConnect: vi.fn(), onChannelMoved: vi.fn(async () => {}) },
unregisterManagedVoiceClient: vi.fn(), registerManagedVoiceClient: vi.fn(),
restoreQueueFromSnapshot: vi.fn(async () => {}), _startJellyfinReportPoller: vi.fn(),
logger: { warn: vi.fn(), info: vi.fn(), error: vi.fn(), debug: vi.fn() },
emit: vi.fn(), getProviderFor: () => provider,
playNext: vi.fn(async () => { const next = queue.next(); if (!next) queue.clear(); return !!next; }),
}) as any;
bot.setupPlayerEvents();
bot.setupTsEvents();
return bot;
}
function cancel(bot: any, action: string) {
if (action === "disconnect") bot.disconnect();
else if (action === "stop") { bot.queue.clear(); bot.player.stop(); }
else if (action === "skip") { if (!bot.queue.next()) bot.queue.clear(); bot.player.stop(); }
else if (action === "replace") {
bot.queue.clear();
bot.queue.add({ id: "other", name: "Other", artist: "A", album: "", coverUrl: "", platform: "bilibili", duration: 10_000 });
bot.queue.playAt(0);
bot.player.play("replacement", 0);
} else { bot.player.sessionId++; bot.player.state = "idle"; }
}
beforeEach(() => vi.useFakeTimers({ toFake: ["setTimeout", "clearTimeout"] }));
afterEach(() => { vi.clearAllTimers(); vi.useRealTimers(); });
describe("Bilibili recovery retries transient URL failures", () => {
it("keeps the same queue and seek through a null lookup, then resumes after one second", async () => {
const bot = makeBot();
const song = bot.queue.current();
bot.provider.getSongUrl.mockResolvedValueOnce(null).mockResolvedValue({ url: "fresh" });
bot.player.emit("trackEnd");
await flush();
expect(bot.queue.current()).toBe(song);
expect(bot.playNext).not.toHaveBeenCalled();
await vi.advanceTimersByTimeAsync(999);
expect(bot.player.play).not.toHaveBeenCalled();
await vi.advanceTimersByTimeAsync(1);
expect(bot.player.play).toHaveBeenCalledWith("fresh", 1000, 10_000);
expect(bot.queue.current()).toBe(song);
expect(bot.provider.getSongUrl.mock.calls.map((args: unknown[]) => args[0])).toEqual(["BV1example", "BV1example"]);
expect(bot.streamRecovery.attempts).toBe(1);
});
it("retries a throw and null at one and two seconds without logging the sensitive exception", async () => {
const bot = makeBot();
const secret = "SYNTHETIC_SIGNED_URL_AND_COOKIE";
bot.provider.getSongUrl.mockRejectedValueOnce(new Error(`https://cdn.invalid/?token=${secret}`))
.mockResolvedValueOnce(null).mockResolvedValue({ url: "fresh" });
bot.player.emit("trackEnd");
await flush();
expect(bot.playNext).not.toHaveBeenCalled();
await vi.advanceTimersByTimeAsync(1000);
expect(bot.player.play).not.toHaveBeenCalled();
await vi.advanceTimersByTimeAsync(1999);
expect(bot.provider.getSongUrl).toHaveBeenCalledTimes(2);
await vi.advanceTimersByTimeAsync(1);
expect(bot.player.play).toHaveBeenCalledWith("fresh", 1000, 10_000);
const logs = JSON.stringify(bot.logger.warn.mock.calls, (_key, value) => value instanceof Error ? { message: value.message, stack: value.stack } : value);
expect(logs).not.toContain(secret);
expect(bot.logger.warn.mock.calls.map((args: any[]) => args[0])).toContainEqual(expect.objectContaining({
platform: "bilibili", sessionId: 1, lookupAttempt: 1, reason: "lookup-error",
}));
expect(bot.playNext).not.toHaveBeenCalled();
});
it.each(["null", "throw"])("advances only after three exhausted %s lookups", async failure => {
const bot = makeBot();
if (failure === "null") bot.provider.getSongUrl.mockResolvedValue(null);
else bot.provider.getSongUrl.mockRejectedValue(new Error("temporary"));
bot.player.emit("trackEnd");
await flush();
expect(bot.playNext).not.toHaveBeenCalled();
await vi.advanceTimersByTimeAsync(2999);
expect(bot.playNext).not.toHaveBeenCalled();
await vi.advanceTimersByTimeAsync(1);
expect(bot.provider.getSongUrl).toHaveBeenCalledTimes(3);
expect(bot.player.play).not.toHaveBeenCalled();
expect(bot.playNext).toHaveBeenCalledTimes(1);
expect(bot.queue.current()).toBeNull();
});
it.each(["skip", "stop", "replace", "restart", "disconnect"])("does not retry or advance after %s during backoff", async action => {
const bot = makeBot();
bot.provider.getSongUrl.mockResolvedValueOnce(null).mockResolvedValue({ url: "fresh" });
bot.player.emit("trackEnd");
await flush();
cancel(bot, action);
bot.player.play.mockClear();
await vi.advanceTimersByTimeAsync(3000);
expect(bot.provider.getSongUrl).toHaveBeenCalledTimes(1);
expect(bot.player.play).not.toHaveBeenCalled();
expect(bot.playNext).not.toHaveBeenCalled();
});
it.each(["skip", "stop", "replace", "restart", "disconnect"])("does not write a late retry result after %s", async action => {
const bot = makeBot();
const retry = deferred<{ url: string }>();
bot.provider.getSongUrl.mockResolvedValueOnce(null).mockReturnValueOnce(retry.promise);
bot.player.emit("trackEnd");
await flush();
await vi.advanceTimersByTimeAsync(1000);
expect(bot.provider.getSongUrl).toHaveBeenCalledTimes(2);
cancel(bot, action);
bot.player.play.mockClear();
retry.resolve({ url: "stale" });
await flush();
expect(bot.player.play).not.toHaveBeenCalled();
expect(bot.playNext).not.toHaveBeenCalled();
});
it("cancels a backoff when a new connection lifecycle begins", async () => {
const bot = makeBot();
bot.provider.getSongUrl.mockResolvedValueOnce(null).mockResolvedValue({ url: "stale" });
bot.player.emit("trackEnd");
await flush();
await bot.connect();
await vi.advanceTimersByTimeAsync(3000);
expect(bot.provider.getSongUrl).toHaveBeenCalledTimes(1);
expect(bot.player.play).not.toHaveBeenCalled();
expect(bot.playNext).not.toHaveBeenCalled();
});
it("shows pause intent in the snapshot during backoff and pauses the recovered stream", async () => {
const bot = makeBot();
bot.provider.getSongUrl.mockResolvedValueOnce(null).mockResolvedValue({ url: "fresh" });
bot.player.emit("trackEnd");
await flush();
bot.cmdPause();
expect(bot.getStatus().paused).toBe(true);
expect(bot.player.getState()).toBe("idle");
await vi.advanceTimersByTimeAsync(1000);
expect(bot.player.getState()).toBe("paused");
expect(bot.playNext).not.toHaveBeenCalled();
bot.cmdResume();
expect(bot.getStatus().paused).toBe(false);
expect(bot.player.getState()).toBe("playing");
expect(bot.player.play).toHaveBeenCalledTimes(1);
});
it("retains a paused exhausted recovery until explicit resume retries", async () => {
const bot = makeBot();
bot.provider.getSongUrl.mockResolvedValue(null);
bot.player.emit("trackEnd");
bot.cmdPause();
await vi.advanceTimersByTimeAsync(3000);
expect(bot.playNext).not.toHaveBeenCalled();
expect(bot.getStatus().paused).toBe(true);
bot.provider.getSongUrl.mockResolvedValue({ url: "fresh" });
bot.cmdResume();
await flush();
expect(bot.player.getState()).toBe("playing");
expect(bot.player.play).toHaveBeenCalledWith("fresh", 1000, 10_000);
expect(bot.playNext).not.toHaveBeenCalled();
});
it("coalesces duplicate ends throughout a lookup and its backoff", async () => {
const bot = makeBot();
const lookup = deferred<{ url: string } | null>();
bot.provider.getSongUrl.mockReturnValueOnce(lookup.promise).mockResolvedValue({ url: "fresh" });
bot.player.emit("trackEnd");
bot.player.emit("trackEnd");
lookup.resolve(null);
await flush();
bot.player.emit("trackEnd");
await vi.advanceTimersByTimeAsync(1000);
expect(bot.provider.getSongUrl).toHaveBeenCalledTimes(2);
expect(bot.player.play).toHaveBeenCalledTimes(1);
expect(bot.streamRecovery.attempts).toBe(1);
expect(bot.playNext).not.toHaveBeenCalled();
});
it("an exhausted end cannot advance over an explicit paused-recovery retry already started", async () => {
const bot = makeBot();
const song = bot.queue.current();
const retry = deferred<{ url: string }>();
bot.provider.getSongUrl.mockResolvedValueOnce(null).mockResolvedValueOnce(null)
.mockResolvedValueOnce(null).mockReturnValueOnce(retry.promise);
const recover = bot.resumeInterruptedStream.bind(bot);
let first = true;
bot.resumeInterruptedStream = () => {
const result = recover();
if (first) { first = false; result.then(() => bot.cmdResume()); }
return result;
};
bot.player.emit("trackEnd");
bot.cmdPause();
await vi.advanceTimersByTimeAsync(3000);
expect(bot.provider.getSongUrl).toHaveBeenCalledTimes(4);
expect(bot.queue.current()).toBe(song);
expect(bot.playNext).not.toHaveBeenCalled();
retry.resolve({ url: "fresh" });
await flush();
expect(bot.player.getState()).toBe("playing");
});
it("lookup retries do not consume the three actual decoder resume attempts", async () => {
const bot = makeBot();
for (let attempt = 0; attempt < 3; attempt++) {
bot.provider.getSongUrl.mockResolvedValueOnce(null).mockResolvedValueOnce({ url: `fresh-${attempt}` });
bot.player.state = "idle";
bot.player.emit("trackEnd");
await vi.advanceTimersByTimeAsync(1000);
expect(bot.queue.current()?.id).toBe("BV1example");
}
bot.player.state = "idle";
bot.player.emit("trackEnd");
await flush();
expect(bot.player.play).toHaveBeenCalledTimes(3);
expect(bot.provider.getSongUrl).toHaveBeenCalledTimes(6);
expect(bot.playNext).toHaveBeenCalledTimes(1);
});
});
describe("occupancy responses belong to the current request and connection", () => {
it("ignores an older alone response after a newer occupied response", async () => {
const bot = makeBot();
bot.player.state = "playing";
const older = deferred<any[]>(), newer = deferred<any[]>();
bot.tsClient.getClientsInChannel.mockReturnValueOnce(older.promise).mockReturnValueOnce(newer.promise);
const first = bot.refreshOccupancy(), second = bot.refreshOccupancy();
newer.resolve([{ id: 1 }, { id: 2 }]);
await second;
older.resolve([{ id: 1 }]);
await first;
expect(bot.player.getState()).toBe("playing");
expect(bot.autoPaused).toBe(false);
expect(bot.idleTimer).toBeNull();
});
it.each(["clientEnter", "clientLeave", "clientMoved"])("a newer %s event invalidates a pending response even if its new query fails", async event => {
const bot = makeBot();
bot.player.state = "playing";
const older = deferred<any[]>();
bot.tsClient.getClientsInChannel.mockReturnValueOnce(older.promise).mockResolvedValueOnce([]);
const first = bot.refreshOccupancy();
bot.tsClient.emit(event, { id: 2, targetChannelID: 3n });
await flush();
older.resolve([{ id: 1 }]);
await first;
expect(bot.player.getState()).toBe("playing");
expect(bot.autoPaused).toBe(false);
expect(bot.idleTimer).toBeNull();
});
it("a listener's return cannot be undone by an earlier alone response", async () => {
const bot = makeBot();
bot.player.state = "paused";
bot.autoPaused = true;
const older = deferred<any[]>(), newer = deferred<any[]>();
bot.tsClient.getClientsInChannel.mockReturnValueOnce(older.promise).mockReturnValueOnce(newer.promise);
const first = bot.refreshOccupancy();
bot.tsClient.emit("clientEnter");
expect(bot.player.getState()).toBe("playing");
older.resolve([{ id: 1 }]);
await first;
expect(bot.player.getState()).toBe("playing");
expect(bot.autoPaused).toBe(false);
newer.resolve([]);
await flush();
});
it.each([false, true])("ignores a response from before disconnect (reconnected=%s)", async reconnect => {
const bot = makeBot();
const older = deferred<any[]>();
bot.tsClient.getClientsInChannel.mockReturnValueOnce(older.promise);
const first = bot.refreshOccupancy();
bot.disconnect();
if (reconnect) { await bot.connect(); bot.player.state = "playing"; }
older.resolve([{ id: 1 }]);
await first;
expect(bot.idleTimer).toBeNull();
expect(bot.autoPaused).toBe(false);
expect(bot.player.getState()).toBe(reconnect ? "playing" : "idle");
});
it("does not let an old lifecycle's poll timer query after reconnect", async () => {
const bot = makeBot();
bot._startIdlePoller();
bot.disconnect();
await bot.connect();
await vi.advanceTimersByTimeAsync(30_000);
expect(bot.tsClient.getClientsInChannel).toHaveBeenCalledTimes(1);
});
it("does not apply or reschedule an old pending poll after reconnect", async () => {
const bot = makeBot();
const oldPoll = deferred<any[]>();
bot.tsClient.getClientsInChannel.mockReturnValueOnce(oldPoll.promise);
bot._startIdlePoller();
await vi.advanceTimersByTimeAsync(30_000);
bot.disconnect();
await bot.connect();
bot.player.state = "playing";
oldPoll.resolve([{ id: 1 }]);
await flush();
expect(bot.player.getState()).toBe("playing");
await vi.advanceTimersByTimeAsync(30_000);
expect(bot.tsClient.getClientsInChannel).toHaveBeenCalledTimes(2);
});
it("does not disconnect a new connection when an old idle deadline expires", async () => {
const bot = makeBot();
bot.player.state = "playing";
bot.tsClient.getClientsInChannel.mockResolvedValueOnce([{ id: 1 }]).mockResolvedValue([]);
await bot.refreshOccupancy();
bot.tsClient.emit("disconnected");
await bot.connect();
bot.player.state = "playing";
await vi.advanceTimersByTimeAsync(60_000);
expect(bot.connected).toBe(true);
expect(bot.tsClient.disconnect).not.toHaveBeenCalled();
});
it("cancels an old idle deadline before awaiting a new handshake", async () => {
const bot = makeBot();
const handshake = deferred<void>();
bot.tsClient.getClientsInChannel.mockResolvedValue([{ id: 1 }]);
await bot.refreshOccupancy();
bot.tsClient.connect.mockReturnValue(handshake.promise);
const connecting = bot.connect();
await vi.advanceTimersByTimeAsync(60_000);
handshake.resolve(undefined);
await expect(connecting).resolves.toBeUndefined();
expect(bot.connected).toBe(true);
expect(bot.tsClient.disconnect).not.toHaveBeenCalled();
});
it("still pauses for a fresh authoritative alone response and ignores unknown occupancy", async () => {
const bot = makeBot();
bot.player.state = "playing";
bot.tsClient.getClientsInChannel.mockResolvedValueOnce([]).mockResolvedValueOnce([{ id: 1 }]);
await bot.refreshOccupancy();
expect(bot.player.getState()).toBe("playing");
await bot.refreshOccupancy();
expect(bot.player.getState()).toBe("paused");
expect(bot.autoPaused).toBe(true);
expect(bot.idleTimer).not.toBeNull();
});
});
+16 -4
View File
@@ -1,4 +1,4 @@
import { describe, it, expect, vi } from "vitest";
import { afterEach, beforeEach, describe, it, expect, vi } from "vitest";
import { EventEmitter } from "node:events";
import { BotInstance, COMMAND_DENIED_MESSAGE, spotifyPortsForBotId } from "./instance.js";
import type { BotInstanceOptions } from "./instance.js";
@@ -134,6 +134,9 @@ describe("BotInstance voice-ducking lifecycle integration", () => {
return {
disconnectEmitted: false,
connected: false,
lifecycleGeneration: 0,
idleTimer: null,
_cancelIdleTimer: (BotInstance.prototype as any)._cancelIdleTimer,
tsClient: {
connect: vi.fn(() => connectPromise),
getResolvedVoiceEndpoint: vi.fn(() => ({ host: "203.0.113.20", port: 12000 })),
@@ -1660,6 +1663,8 @@ describe("cmdPlaylist with a playlist link (#160)", () => {
});
describe("resumeInterruptedStream — long B站 streams dying mid-play (#161)", () => {
beforeEach(() => vi.useFakeTimers({ toFake: ["setTimeout", "clearTimeout"] }));
afterEach(() => { vi.clearAllTimers(); vi.useRealTimers(); });
const resumeInterruptedStream = (BotInstance.prototype as any).resumeInterruptedStream as (
this: unknown,
) => Promise<boolean>;
@@ -1678,6 +1683,7 @@ describe("resumeInterruptedStream — long B站 streams dying mid-play (#161)",
song,
provider,
connected: true,
lifecycleGeneration: 0,
effectiveDuration: song.duration,
streamRecovery: null,
queue: { current: vi.fn(() => song) },
@@ -1739,7 +1745,10 @@ describe("resumeInterruptedStream — long B站 streams dying mid-play (#161)",
it("falls through to advancing when no fresh URL can be fetched", async () => {
const ctx = makeCtx({ url: null });
expect(await resumeInterruptedStream.call(ctx)).toBe(false);
const recovery = resumeInterruptedStream.call(ctx);
await vi.advanceTimersByTimeAsync(3000);
expect(await recovery).toBe(false);
expect(ctx.provider.getSongUrl).toHaveBeenCalledTimes(3);
expect(ctx.player.play).not.toHaveBeenCalled();
});
@@ -1755,6 +1764,8 @@ describe("resumeInterruptedStream — long B站 streams dying mid-play (#161)",
});
describe("BotInstance trackEnd — stale playback sessions", () => {
beforeEach(() => vi.useFakeTimers({ toFake: ["setTimeout", "clearTimeout"] }));
afterEach(() => { vi.clearAllTimers(); vi.useRealTimers(); });
function makeEndedCtx(platform = "bilibili", duration = 10_000) {
const song = {
id: "ended", name: "Ended", artist: "A", album: "", coverUrl: "",
@@ -1774,6 +1785,7 @@ describe("BotInstance trackEnd — stale playback sessions", () => {
const advances: string[] = [];
const ctx: any = {
song, provider, player, connected: true, effectiveDuration: duration,
lifecycleGeneration: 0,
streamRecovery: null, queue: { current: () => current },
spotifyController: new EventEmitter(), tsClient: { sendVoiceData: vi.fn() },
logger: { warn: vi.fn(), debug: vi.fn(), error: vi.fn() }, emit: vi.fn(),
@@ -1893,7 +1905,7 @@ describe("BotInstance trackEnd — stale playback sessions", () => {
expect(ctx.advances).toEqual([]);
});
it("resume after a paused failed lookup retries recovery instead of remaining idle", async () => {
it("resume while a paused failed lookup waits for retry continues after its backoff", async () => {
const ctx = makeEndedCtx();
const lookup = deferred<{ url: string }>();
ctx.provider.getSongUrl.mockReturnValue(lookup.promise);
@@ -1903,7 +1915,7 @@ describe("BotInstance trackEnd — stale playback sessions", () => {
await flushEvents();
ctx.provider.getSongUrl.mockResolvedValue({ url: "recovered" });
ctx.resume();
await flushEvents();
await vi.advanceTimersByTimeAsync(1000);
expect(ctx.player.getState()).toBe("playing");
expect(ctx.player.play).toHaveBeenCalledWith("recovered", 1000, 10_000);
expect(ctx.advances).toEqual([]);
+102 -42
View File
@@ -4,6 +4,7 @@ import {
type TS3ClientOptions,
type TS3TextMessage,
type TS3VoiceActivity,
type TS3VoiceSendFailure,
} from "../ts-protocol/client.js";
import { AudioPlayer } from "../audio/player.js";
import { PlayQueue, PlayMode, type QueuedSong } from "../audio/queue.js";
@@ -182,6 +183,9 @@ export class BotInstance extends EventEmitter {
private logger: Logger;
private avatarStore: AvatarStore;
private connected = false;
/** Fences async work and timers from earlier TeamSpeak connections. */
private lifecycleGeneration = 0;
private occupancyRequest = 0;
private disconnectEmitted = false;
private voteSkipUsers = new Set<string>();
private isAdvancing = false;
@@ -200,7 +204,7 @@ export class BotInstance extends EventEmitter {
/** 当前曲实际播放时长(试听片段秒数或完整 duration);resolveAndPlay 赋值。 */
private effectiveDuration: number | undefined;
/** Resume attempts for the current song's stream (#161); see resumeInterruptedStream. */
private streamRecovery: { song: QueuedSong; attempts: number; position: number; session: number; pauseRequested: boolean; inFlight: boolean } | null = null;
private streamRecovery: { song: QueuedSong; attempts: number; position: number; session: number; generation: number; pauseRequested: boolean; inFlight: boolean } | null = null;
private playGate: Promise<unknown> = Promise.resolve();
/** Per-bot Jellyfin playback-report session (start / ~10s progress / stop).
* null when the wired provider has no reporting capability. */
@@ -326,15 +330,21 @@ export class BotInstance extends EventEmitter {
private setupPlayerEvents(): void {
this.player.on("frame", (opusFrame: Buffer) => {
this.tsClient.sendVoiceData(opusFrame);
const result = this.tsClient.sendVoiceData(opusFrame);
// A terminal fault can outlive a pause, queue change, or URL recovery.
// Retry the actual send first so a repaired transport can clear it.
if (result === "failed" && this.connected && this.player.getState() === "playing") {
this.logger.warn("Voice transport is still failing; playback paused");
this.cmdPause();
}
});
this.player.on("trackEnd", () => {
const endedSong = this.queue.current();
const endedSession = this.player.getPlaybackSessionId();
this.resumeInterruptedStream()
.catch((err) => {
this.logger.warn({ err }, "Stream resume failed");
.catch(() => {
this.logger.warn({ sessionId: endedSession, reason: "recovery-error" }, "Stream resume failed");
return false;
})
.then((resumed) => {
@@ -346,7 +356,8 @@ export class BotInstance extends EventEmitter {
this.queue.current() !== endedSong ||
this.player.getPlaybackSessionId() !== endedSession ||
this.player.getState() !== "idle" ||
(this.streamRecovery?.song === endedSong && this.streamRecovery.pauseRequested)
(this.streamRecovery?.song === endedSong && this.streamRecovery.session === endedSession &&
(this.streamRecovery.pauseRequested || this.streamRecovery.inFlight))
) return;
this.logger.debug("Track ended, advancing queue");
return this.playNext();
@@ -413,6 +424,22 @@ export class BotInstance extends EventEmitter {
}
private setupTsEvents(): void {
this.tsClient.on("voiceSendFailed", (failure: TS3VoiceSendFailure) => {
if (!this.connected) return;
const recovery = this.streamRecovery;
const pendingRecovery = this.player.getState() === "idle" && recovery &&
recovery.song === this.queue.current() && recovery.generation === this.lifecycleGeneration &&
recovery.session === this.player.getPlaybackSessionId();
if (this.player.getState() !== "playing" && !pendingRecovery) return;
this.logger.warn(
{ code: failure.code, consecutiveFailures: failure.consecutiveFailures, durationMs: failure.durationMs },
"Voice transmission failed; playback paused. Restore the connection before resuming",
);
// Preserve the queue and seek position. Use the normal pause path so
// Spotify also pauses and occupancy cannot resume a broken transport.
this.cmdPause();
});
this.tsClient.on("textMessage", (msg: TS3TextMessage) => {
this.handleTextMessage(msg).catch((err) => {
this.logger.error({ err }, "Unhandled error in text message handler");
@@ -424,7 +451,9 @@ export class BotInstance extends EventEmitter {
// completed (hanging handshake → 60s library idle timeout) and
// this.connected was never flipped to true. Previously this handler
// short-circuited on !this.connected, leaving player stuck as "playing".
this.lifecycleGeneration++;
this.connected = false;
this._cancelIdleTimer();
this.unregisterManagedVoiceClient(MANAGED_VOICE_CLIENT_RELEASE_GRACE_MS);
this.voiceDucking.reset(true);
// Cancel any pending live-queue snapshot BEFORE clearing the queue: a
@@ -452,6 +481,8 @@ export class BotInstance extends EventEmitter {
});
this.tsClient.on("connected", () => {
this.lifecycleGeneration++;
this._cancelIdleTimer();
// Fresh connection — clear any stale auto-pause flag from a prior session.
this.autoPaused = false;
this._startIdlePoller();
@@ -562,8 +593,13 @@ export class BotInstance extends EventEmitter {
private async refreshOccupancy(): Promise<void> {
if (!this.connected) return;
const request = ++this.occupancyRequest;
const generation = this.lifecycleGeneration;
const client = this.tsClient;
try {
const clients = await this.tsClient.getClientsInChannel();
const clients = await client.getClientsInChannel();
if (!this.connected || this.lifecycleGeneration !== generation ||
this.occupancyRequest !== request || this.tsClient !== client) return;
// A 0-length result means the clientlist query failed (the bot is always
// in its own channel) — occupancy is unknown, so don't act. Acting on it
// would mis-read it as "empty" and falsely auto-pause / idle-disconnect.
@@ -575,6 +611,8 @@ export class BotInstance extends EventEmitter {
}
async connect(): Promise<void> {
this.lifecycleGeneration++;
this._cancelIdleTimer();
this.disconnectEmitted = false;
await this.tsClient.connect();
const resolvedEndpoint = this.tsClient.getResolvedVoiceEndpoint();
@@ -607,6 +645,7 @@ export class BotInstance extends EventEmitter {
}
disconnect(): void {
this.lifecycleGeneration++;
this._cancelIdleTimer();
this.voiceDucking.reset(true);
// Cancel any pending live-queue snapshot before clearing so it can't fire
@@ -662,16 +701,12 @@ export class BotInstance extends EventEmitter {
}
private _startIdlePoller(): void {
const generation = this.lifecycleGeneration;
// 每 30 秒检查一次频道人数
const poll = async () => {
if (!this.connected) return;
try {
const clients = await this.tsClient.getClientsInChannel();
// null = clientlist query failed (occupancy unknown) → don't act.
const userCount = occupancyFromClientList(clients.length);
if (userCount !== null) this.handleOccupancy(userCount);
} catch { /* ignore */ }
setTimeout(poll, 30_000);
if (!this.connected || this.lifecycleGeneration !== generation) return;
await this.refreshOccupancy();
if (this.connected && this.lifecycleGeneration === generation) setTimeout(poll, 30_000);
};
setTimeout(poll, 30_000);
}
@@ -736,8 +771,9 @@ export class BotInstance extends EventEmitter {
if (this.idleTimer !== null) return; // 已经在倒计时,不重复创建
const minutes = this.config.idleTimeoutMinutes ?? 0;
if (!this.connected || minutes <= 0) return;
const generation = this.lifecycleGeneration;
this.idleTimer = setTimeout(() => {
if (!this.connected) return;
if (!this.connected || this.lifecycleGeneration !== generation) return;
this.logger.info({ idleMinutes: minutes }, "Channel empty, disconnecting due to idle timeout");
this.disconnect();
}, minutes * 60 * 1000);
@@ -1154,6 +1190,7 @@ export class BotInstance extends EventEmitter {
/** A track that ends within this many seconds of its duration ended normally. */
private static readonly STREAM_END_TOLERANCE_S = 30;
private static readonly MAX_STREAM_RESUMES = 3;
private static readonly MAX_STREAM_URL_LOOKUPS = 3;
/**
* Called when the player reports a track end. If a B站 stream ended long
@@ -1176,15 +1213,18 @@ export class BotInstance extends EventEmitter {
if (!(duration > 0) || duration - position <= BotInstance.STREAM_END_TOLERANCE_S) {
return false;
}
if (this.player.getState() !== "idle") return true;
const generation = this.lifecycleGeneration;
const recovery = this.streamRecovery;
if (
!recovery ||
recovery.song !== song ||
recovery.session !== endedSession ||
recovery.generation !== generation ||
position - recovery.position > BotInstance.STREAM_END_TOLERANCE_S
) {
this.streamRecovery = { song, attempts: 0, position, session: endedSession, pauseRequested: false, inFlight: false };
this.streamRecovery = { song, attempts: 0, position, session: endedSession, generation, pauseRequested: false, inFlight: false };
}
const state = this.streamRecovery!;
if (state.inFlight) return true;
@@ -1196,40 +1236,55 @@ export class BotInstance extends EventEmitter {
this.streamRecovery = null;
return false;
}
state.attempts++;
state.position = position;
this.logger.warn(
{ songId: song.id, position, duration, attempt: state.attempts },
{ songId: song.id, position, duration, attempt: state.attempts + 1 },
"Stream ended before the track did — resuming with a fresh URL",
);
state.inFlight = true;
let result: Awaited<ReturnType<MusicProvider["getSongUrl"]>>;
const isCurrent = () => this.connected && this.lifecycleGeneration === generation &&
this.streamRecovery === state && this.queue.current() === song &&
this.player.getPlaybackSessionId() === endedSession && this.player.getState() === "idle";
const cancelled = () => {
if (this.streamRecovery === state) this.streamRecovery = null;
return true;
};
try {
result = await this.getProviderFor(song.platform).getSongUrl(song.id);
for (let lookup = 0; lookup < BotInstance.MAX_STREAM_URL_LOOKUPS; lookup++) {
if (!isCurrent()) return cancelled();
if (lookup > 0) {
await new Promise<void>(resolve => setTimeout(resolve, lookup * 1000));
if (!isCurrent()) return cancelled();
}
let result: Awaited<ReturnType<MusicProvider["getSongUrl"]>> = null;
let reason = "no-url";
try {
result = await this.getProviderFor(song.platform).getSongUrl(song.id);
} catch {
reason = "lookup-error";
}
// Skip, stop, reconnect, or a new playback session supersedes the lookup.
if (!isCurrent()) return cancelled();
if (result?.url) {
// Lookup retries do not consume the actual decoder resume budget.
state.attempts++;
song.url = result.url;
this.player.play(result.url, position, duration);
state.session = this.player.getPlaybackSessionId();
if (state.pauseRequested) this.player.pause();
this.emit("stateChange");
return true;
}
this.logger.warn(
{ platform: "bilibili", sessionId: endedSession, position, duration, lookupAttempt: lookup + 1, reason },
"Fresh stream URL lookup failed",
);
}
return false;
} finally {
state.inFlight = false;
}
// The user may have skipped/stopped while we were resolving; never
// clobber whatever is playing now.
if (
this.queue.current() !== song ||
this.player.getPlaybackSessionId() !== endedSession ||
this.player.getState() !== "idle"
) {
if (this.streamRecovery === state) this.streamRecovery = null;
return true;
}
if (!result?.url || !this.connected) return false;
song.url = result.url;
this.player.play(result.url, position, duration);
state.session = this.player.getPlaybackSessionId();
// During the lookup the ended player is idle, so pause() alone cannot
// remember the user's intent. Pause the recovered stream before it emits frames.
if (state.pauseRequested) this.player.pause();
this.emit("stateChange");
return true;
}
private async syncProfileToSong(song: QueuedSong | null): Promise<void> {
@@ -2135,12 +2190,17 @@ export class BotInstance extends EventEmitter {
}
getStatus(): BotStatus {
const playerState = this.player.getState();
const recovery = this.streamRecovery;
const recoveryPaused = playerState === "idle" && recovery?.pauseRequested &&
recovery.song === this.queue.current() && recovery.generation === this.lifecycleGeneration &&
recovery.session === this.player.getPlaybackSessionId();
return {
id: this.id,
name: this.name,
connected: this.connected,
playing: this.player.getState() === "playing",
paused: this.player.getState() === "paused",
playing: playerState === "playing",
paused: playerState === "paused" || !!recoveryPaused,
currentSong: this.queue.current(),
queueSize: this.queue.size(),
volume: this.player.getVolume(),
+145
View File
@@ -0,0 +1,145 @@
import { afterEach, describe, expect, it, vi } from "vitest";
import { EventEmitter } from "node:events";
import pino from "pino";
import { BotInstance } from "./instance.js";
import { AudioPlayer } from "../audio/player.js";
import { PlayQueue } from "../audio/queue.js";
import { PCM_FRAME_BYTES } from "../audio/encoder.js";
import { TS3Client } from "../ts-protocol/client.js";
function makeHarness(platform: "bilibili" | "spotify" = "bilibili") {
const logger = pino({ level: "silent" });
const player = new AudioPlayer(logger);
const queue = new PlayQueue();
queue.add({ id: "offline-fixture", name: "67-minute fixture", artist: "fixture", album: "", coverUrl: "", platform, duration: 4020 });
queue.play();
Object.assign(player, { state: "playing", framesPlayed: 111000 });
const tsClient = new TS3Client({ host: "localhost", port: 9987, queryPort: 10011, nickname: "OfflineFixture" }, logger);
let sidecarPauses = 0;
const bot = Object.assign(new EventEmitter(), {
logger, player, queue, tsClient, connected: true, autoPaused: true,
streamRecovery: null, lifecycleGeneration: 0,
spotifyController: Object.assign(new EventEmitter(), { pause: async () => { sidecarPauses++; } }),
});
const methods = BotInstance.prototype as unknown as {
setupTsEvents(this: typeof bot): void;
cmdPause(this: typeof bot): string;
};
Object.assign(bot, { cmdPause: methods.cmdPause });
methods.setupTsEvents.call(bot);
return { bot, player, queue, tsClient, sidecarPauses: () => sidecarPauses };
}
afterEach(() => { vi.clearAllTimers(); vi.restoreAllMocks(); vi.useRealTimers(); });
describe("BotInstance persistent voice send failure", () => {
it.each(["bilibili", "spotify"] as const)("pauses %s while retaining the song and playback position", platform => {
const h = makeHarness(platform);
const song = h.queue.current();
let changes = 0;
h.bot.on("stateChange", () => changes++);
try {
h.tsClient.emit("voiceSendFailed", { code: "ERR_SOCKET_DGRAM_NOT_RUNNING", consecutiveFailures: 100, durationMs: 2000 });
expect(h.player.getState()).toBe("paused");
expect(h.player.getElapsed()).toBe(2220);
expect(h.queue.current()).toBe(song);
expect(h.queue.size()).toBe(1);
expect(h.bot.autoPaused).toBe(false);
expect(changes).toBe(1);
expect(h.sidecarPauses()).toBe(platform === "spotify" ? 1 : 0);
h.tsClient.emit("voiceSendFailed", { consecutiveFailures: 200, durationMs: 4000 });
expect(changes).toBe(1);
} finally { h.player.stop(); }
});
it("stops the actual audio scheduler from advancing elapsed after terminal transmission failure", () => {
vi.useFakeTimers({ toFake: ["setTimeout", "clearTimeout", "performance"] });
const h = makeHarness();
const internal = h.player as unknown as { startFrameLoop(): void };
let frames = 0;
h.player.on("frame", () => frames++);
Object.assign(h.player, { framesPlayed: 0, ffmpeg: { pid: undefined }, pcmBuffer: Buffer.alloc(PCM_FRAME_BYTES * 250), currentSongDuration: 4020 });
try {
internal.startFrameLoop();
vi.advanceTimersByTime(1000);
expect(frames).toBe(50);
h.tsClient.emit("voiceSendFailed", { consecutiveFailures: 100, durationMs: 2000 });
vi.advanceTimersByTime(3000);
expect(frames).toBe(50);
expect(h.player.getElapsed()).toBe(1);
expect(h.player.getState()).toBe("paused");
} finally { h.player.stop(); }
});
it("lets transient failure recover without pausing and does not auto-resume a user pause", () => {
const h = makeHarness();
try {
h.tsClient.emit("voiceSendFailure", { consecutiveFailures: 1, durationMs: 0 });
h.tsClient.emit("voiceSendRecovered", { consecutiveFailures: 1, durationMs: 20 });
expect(h.player.getState()).toBe("playing");
h.player.pause();
h.tsClient.emit("voiceSendRecovered", { consecutiveFailures: 1, durationMs: 20 });
expect(h.player.getState()).toBe("paused");
} finally { h.player.stop(); }
});
it("preserves pause intent when terminal failure occurs during a pending fresh URL lookup", async () => {
const h = makeHarness();
Object.assign(h.player, { state: "idle" });
let resolve!: (result: { url: string }) => void;
const lookup = new Promise<{ url: string }>(done => { resolve = done; });
Object.assign(h.bot, { getProviderFor: () => ({ getSongUrl: () => lookup }) });
// Isolate only process creation; the actual stop, seek, state, session,
// recovery method and pause method remain in use.
vi.spyOn(h.player, "play").mockImplementation((_url, seek, duration) => {
h.player.stop();
Object.assign(h.player, { state: "playing", seekOffset: seek, currentSongDuration: duration });
});
const resume = (BotInstance.prototype as unknown as { resumeInterruptedStream(this: typeof h.bot): Promise<boolean> }).resumeInterruptedStream;
const recovery = resume.call(h.bot);
try {
h.tsClient.emit("voiceSendFailed", { code: "EPIPE", consecutiveFailures: 100, durationMs: 2000 });
resolve({ url: "https://offline.invalid/fresh.m4s" });
expect(await recovery).toBe(true);
expect(h.player.getState()).toBe("paused");
expect(h.player.getElapsed()).toBe(2220);
expect(h.queue.current()?.id).toBe("offline-fixture");
} finally { h.player.stop(); }
});
it("pauses again if manual resume encounters the same persistent voice fault", () => {
vi.useFakeTimers({ toFake: ["setTimeout", "clearTimeout", "performance", "Date"] });
const h = makeHarness();
let accepted = 0;
Object.assign(h.tsClient, { client: { sendVoice() {
if (accepted++ > 0) throw Object.assign(new Error("offline fixture"), { code: "EPIPE" });
} } });
const setupPlayerEvents = (BotInstance.prototype as unknown as { setupPlayerEvents(this: typeof h.bot): void }).setupPlayerEvents;
setupPlayerEvents.call(h.bot);
Object.assign(h.player, { framesPlayed: 0, ffmpeg: { pid: undefined }, pcmBuffer: Buffer.alloc(PCM_FRAME_BYTES * 250), currentSongDuration: 4020 });
try {
(h.player as unknown as { startFrameLoop(): void }).startFrameLoop();
vi.advanceTimersByTime(2100);
expect(h.player.getState()).toBe("paused");
const pausedAt = h.player.getElapsed();
h.player.resume();
vi.advanceTimersByTime(40);
expect(h.player.getState()).toBe("paused");
expect(h.player.getElapsed()).toBeLessThanOrEqual(pausedAt + 0.02);
expect(h.queue.current()?.id).toBe("offline-fixture");
} finally { h.player.stop(); }
});
it("ignores terminal notifications when disconnected or already idle", () => {
const h = makeHarness();
try {
h.bot.connected = false;
h.tsClient.emit("voiceSendFailed", { consecutiveFailures: 100, durationMs: 2000 });
expect(h.player.getState()).toBe("playing");
h.bot.connected = true;
h.player.stop();
h.tsClient.emit("voiceSendFailed", { consecutiveFailures: 100, durationMs: 2000 });
expect(h.player.getState()).toBe("idle");
} finally { h.player.stop(); }
});
});
+43
View File
@@ -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);
+86
View File
@@ -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" });
});
});
+92
View File
@@ -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
View File
@@ -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
View File
@@ -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}`,
};
}
+295
View File
@@ -0,0 +1,295 @@
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
import { Client as SDKClient, type ResolvedAddr } from "@honeybbq/teamspeak-client";
import { Resolver } from "@honeybbq/teamspeak-client/discovery";
import { TS3Client, type TS3VoiceSendFailure } from "./client.js";
import type { TrackingVoiceEndpointResolver } from "./voice-endpoint.js";
import type { Logger } from "../logger.js";
type Failure = TS3VoiceSendFailure;
type Harness = { client: SDKClient | null; voiceFramesSent: number };
describe("TS3Client voice send failure lifecycle", () => {
let client: TS3Client;
let sdk: SDKClient;
let logger: Logger;
let records: Array<{ level: string; args: unknown[] }>;
beforeEach(async () => {
vi.useFakeTimers();
// Keep the real SDK object, logger bridge and event dispatch. Replace only
// its network handshake and explicit shutdown so no UDP socket is opened.
vi.spyOn(SDKClient.prototype, "connect").mockImplementation(async function (this: SDKClient) {
this.clid = 42;
this._markConnected();
});
vi.spyOn(SDKClient.prototype, "disconnect").mockImplementation(async function (this: SDKClient) {
this.handler.onClosed?.(null);
});
records = [];
const capture = (level: string) => (...args: unknown[]) => records.push({ level, args });
logger = { info: capture("info"), warn: capture("warn"), error: capture("error"), debug() {} } as unknown as Logger;
client = new TS3Client({ host: "localhost", port: 9987, queryPort: 10011, nickname: "VoiceTest", serverProtocol: "ts3" }, logger);
await client.connect();
sdk = (client as unknown as Harness).client!;
});
afterEach(async () => {
client.disconnect();
await vi.advanceTimersByTimeAsync(0);
vi.restoreAllMocks();
vi.useRealTimers();
});
it.each([false, true])("reports a sanitized first failure after earlier accepted send=%s without throwing", (acceptedFirst) => {
const send = vi.spyOn(sdk, "sendVoice").mockImplementation(() => {});
if (acceptedFirst) client.sendVoiceData(Buffer.from([1]));
const failures: Failure[] = [];
client.on("voiceSendFailure", event => failures.push(event));
send.mockImplementation(() => { throw Object.assign(new Error("https://user:password@host/?token=credential-secret"), { code: "ERR_SOCKET_DGRAM_NOT_RUNNING", spawnargs: ["credential-secret"] }); });
let result: unknown;
expect(() => { result = client.sendVoiceData(Buffer.from([2])); }).not.toThrow();
expect(result).toBe("retrying");
expect(failures).toEqual([{ code: "ERR_SOCKET_DGRAM_NOT_RUNNING", consecutiveFailures: 1, durationMs: 0 }]);
expect((client as unknown as Harness).voiceFramesSent).toBe(acceptedFirst ? 1 : 0);
expect(records.filter(record => record.level === "warn")).toHaveLength(1);
expect(JSON.stringify(records)).not.toContain("password");
expect(JSON.stringify(records)).not.toContain("credential-secret");
});
it.each([{ code: "credential-secret", message: "credential-secret" }, "credential-secret"])("omits untrusted error codes and primitive error text: %s", error => {
const failures: Failure[] = [];
client.on("voiceSendFailure", event => failures.push(event));
vi.spyOn(sdk, "sendVoice").mockImplementation(() => { throw error; });
client.sendVoiceData(Buffer.from([1]));
expect(failures).toEqual([{ consecutiveFailures: 1, durationMs: 0 }]);
expect(JSON.stringify(records)).not.toContain("credential-secret");
});
it("emits one recovery and cancels terminal failure when the next send succeeds", async () => {
const recovered: Failure[] = [];
const terminal: Failure[] = [];
client.on("voiceSendRecovered", event => recovered.push(event));
client.on("voiceSendFailed", event => terminal.push(event));
const send = vi.spyOn(sdk, "sendVoice").mockImplementation(() => { throw { code: "ENOBUFS" }; });
client.sendVoiceData(Buffer.from([1]));
await vi.advanceTimersByTimeAsync(20);
send.mockImplementation(() => {});
expect(client.sendVoiceData(Buffer.from([2]))).toBe("accepted");
client.sendVoiceData(Buffer.from([3]));
await vi.advanceTimersByTimeAsync(3000);
expect(recovered).toEqual([{ code: "ENOBUFS", consecutiveFailures: 1, durationMs: 20 }]);
expect(terminal).toEqual([]);
expect((client as unknown as Harness).voiceFramesSent).toBe(2);
});
it("bounds persistent failure reporting to one first event and one terminal event per burst", async () => {
const first: Failure[] = [];
const terminal: Failure[] = [];
client.on("voiceSendFailure", event => first.push(event));
client.on("voiceSendFailed", event => terminal.push(event));
vi.spyOn(sdk, "sendVoice").mockImplementation(() => { throw { code: "ERR_SOCKET_DGRAM_NOT_RUNNING" }; });
for (let i = 0; i < 100; i++) {
client.sendVoiceData(Buffer.from([1]));
await vi.advanceTimersByTimeAsync(20);
}
expect(first).toHaveLength(1);
expect(terminal).toEqual([{ code: "ERR_SOCKET_DGRAM_NOT_RUNNING", consecutiveFailures: 100, durationMs: 2000 }]);
for (let i = 0; i < 100; i++) expect(client.sendVoiceData(Buffer.from([1]))).toBe("failed");
await vi.advanceTimersByTimeAsync(10000);
expect(terminal).toHaveLength(1);
expect(records.filter(record => ["warn", "error"].includes(record.level)).length).toBeLessThanOrEqual(2);
expect((client as unknown as Harness).voiceFramesSent).toBe(0);
});
it("recovers after the terminal event and allows a fresh failure burst", async () => {
const first: Failure[] = [], recovered: Failure[] = [], terminal: Failure[] = [];
client.on("voiceSendFailure", event => first.push(event));
client.on("voiceSendRecovered", event => recovered.push(event));
client.on("voiceSendFailed", event => terminal.push(event));
const send = vi.spyOn(sdk, "sendVoice").mockImplementation(() => { throw { code: "EPIPE" }; });
client.sendVoiceData(Buffer.from([1]));
await vi.advanceTimersByTimeAsync(2000);
send.mockImplementation(() => {});
expect(client.sendVoiceData(Buffer.from([1]))).toBe("accepted");
send.mockImplementation(() => { throw { code: "EPIPE" }; });
client.sendVoiceData(Buffer.from([1]));
await vi.advanceTimersByTimeAsync(2000);
expect(first).toHaveLength(2);
expect(recovered).toHaveLength(1);
expect(terminal).toHaveLength(2);
});
it("cancels failure timers on explicit disconnect", async () => {
const terminal: Failure[] = [];
client.on("voiceSendFailed", event => terminal.push(event));
const send = vi.spyOn(sdk, "sendVoice").mockImplementation(() => {});
client.sendVoiceData(Buffer.from([1]));
send.mockImplementation(() => { throw { code: "EPIPE" }; });
client.sendVoiceData(Buffer.from([1]));
client.disconnect();
expect(client.sendVoiceData(Buffer.from([1]))).toBe("unavailable");
await vi.advanceTimersByTimeAsync(3000);
expect(terminal).toEqual([]);
expect((client as unknown as Harness).voiceFramesSent).toBe(0);
});
it("cancels failure timers on SDK connection loss", async () => {
const terminal: Failure[] = [];
client.on("voiceSendFailed", event => terminal.push(event));
const send = vi.spyOn(sdk, "sendVoice").mockImplementation(() => {});
client.sendVoiceData(Buffer.from([1]));
send.mockImplementation(() => { throw { code: "EPIPE" }; });
client.sendVoiceData(Buffer.from([1]));
sdk.handler.onClosed?.(new Error("synthetic connection loss"));
await vi.advanceTimersByTimeAsync(3000);
expect(terminal).toEqual([]);
expect((client as unknown as Harness).voiceFramesSent).toBe(0);
send.mockClear();
client.sendVoiceData(Buffer.from([1]));
expect(send).not.toHaveBeenCalled();
});
it("does not revive a connection whose handshake was cancelled", async () => {
let finishHandshake!: () => void;
vi.spyOn(SDKClient.prototype, "connect").mockImplementationOnce(function (this: SDKClient) {
return new Promise<void>(resolve => {
finishHandshake = () => {
this.clid = 99;
this._markConnected();
resolve();
};
});
});
const connected = vi.fn();
client.on("connected", connected);
const pending = client.connect();
await Promise.resolve();
expect(finishHandshake).toBeTypeOf("function");
const connectingSdk = (client as unknown as Harness).client!;
const close = vi.spyOn(connectingSdk.handler, "close");
const send = vi.spyOn(connectingSdk, "sendVoice");
client.disconnect();
finishHandshake();
await pending;
await vi.advanceTimersByTimeAsync(0);
expect(client.getClientId()).toBe(0);
expect(connected).not.toHaveBeenCalled();
expect(close).toHaveBeenCalled();
client.sendVoiceData(Buffer.from([1]));
expect(send).not.toHaveBeenCalled();
});
it("ignores old SDK warnings and disconnect callbacks after replacement", async () => {
const terminal: Failure[] = [];
client.on("voiceSendFailed", event => terminal.push(event));
vi.spyOn(sdk, "sendVoice").mockImplementation(() => { throw { code: "EPIPE" }; });
client.sendVoiceData(Buffer.from([1]));
await client.connect();
// Old disconnect dispatch is queued; old socket callbacks may follow it.
const before = records.length;
sdk.logger.warn("udp send error", Object.assign(new Error("credential-secret"), { code: "ECONNREFUSED" }));
await vi.advanceTimersByTimeAsync(3000);
expect(client.getClientId()).toBe(42);
expect(records.slice(before).filter(record => record.level === "warn")).toEqual([]);
expect(terminal).toEqual([]);
});
it("does not let an older named-channel lookup move the replacement connection", async () => {
client.disconnect();
await vi.advanceTimersByTimeAsync(0);
client = new TS3Client({ host: "localhost", port: 9987, queryPort: 10011, nickname: "VoiceTest", serverProtocol: "ts3", defaultChannel: "Music" }, logger);
let nextClientId = 101;
vi.mocked(SDKClient.prototype.connect).mockImplementation(async function (this: SDKClient) {
this.clid = nextClientId++;
this._markConnected();
});
let finishOlderLookup!: (rows: Record<string, string>[]) => void;
vi.spyOn(SDKClient.prototype, "execCommandWithResponse").mockImplementation(function (this: SDKClient) {
if (this.clid === 101) return new Promise(resolve => { finishOlderLookup = resolve; });
return Promise.resolve([{ cid: "20", channel_name: "Music" }]);
});
const moves: Array<{ clientId: number; command: string }> = [];
vi.spyOn(SDKClient.prototype, "execCommand").mockImplementation(async function (this: SDKClient, command: string) {
moves.push({ clientId: this.clid, command });
});
const first = client.connect();
for (let i = 0; i < 8; i++) await Promise.resolve();
expect(finishOlderLookup).toBeTypeOf("function");
await client.connect();
finishOlderLookup([{ cid: "10", channel_name: "Music" }]);
await first;
expect(moves).toEqual([{ clientId: 102, command: "clientmove clid=102 cid=20" }]);
expect(records.filter(record => record.args[1] === "Joined channel").map(record => record.args[0])).toEqual([{ channelName: "Music", cid: "20" }]);
expect(client.getClientId()).toBe(102);
});
it("keeps the replacement endpoint when an older DNS resolution finishes last", async () => {
let finishOlderResolution!: (rows: ResolvedAddr[]) => void;
const row = (addr: string): ResolvedAddr => ({ addr, source: "test", expiry: new Date(0) });
vi.spyOn(Resolver.prototype, "resolve").mockImplementationOnce(() => new Promise(resolve => { finishOlderResolution = resolve; }))
.mockResolvedValue([row("192.0.2.2:9987")]);
let sequence = 0;
vi.mocked(SDKClient.prototype.connect).mockImplementation(async function (this: SDKClient) {
const clientId = 101 + sequence++;
const resolver = (client as unknown as { voiceEndpointResolver: TrackingVoiceEndpointResolver }).voiceEndpointResolver;
await resolver.resolve("voice.test:9987");
this.clid = clientId;
this._markConnected();
});
const first = client.connect();
for (let i = 0; i < 8; i++) await Promise.resolve();
expect(finishOlderResolution).toBeTypeOf("function");
await client.connect();
expect(client.getResolvedVoiceEndpoint()).toEqual({ host: "192.0.2.2", port: 9987 });
finishOlderResolution([row("192.0.2.1:9987")]);
await first;
expect(client.getClientId()).toBe(102);
expect(client.getResolvedVoiceEndpoint()).toEqual({ host: "192.0.2.2", port: 9987 });
});
it.each(["10", "Music"])("does not report an old %s channel move after connection replacement", async channel => {
vi.spyOn(SDKClient.prototype, "execCommandWithResponse").mockResolvedValue([{ cid: "10", channel_name: "Music" }]);
let finishMove!: () => void;
const move = vi.spyOn(sdk, "execCommand").mockImplementation(() => new Promise<void>(resolve => { finishMove = resolve; }));
const joining = client.joinChannel(channel);
for (let i = 0; i < 8; i++) await Promise.resolve();
expect(move).toHaveBeenCalledWith("clientmove clid=42 cid=10", 10000);
await client.connect();
finishMove();
await joining;
expect(records.filter(record => record.args[1] === "Joined channel")).toEqual([]);
});
it("does not report a stale channel lookup failure against the new connection", async () => {
let failLookup!: (error: Error) => void;
vi.spyOn(sdk, "execCommandWithResponse").mockImplementation(() => new Promise((_resolve, reject) => { failLookup = reject; }));
const joining = client.joinChannel("Music");
await client.connect();
failLookup(new Error("old credential-secret failure"));
await joining;
expect(records.filter(record => record.level === "error")).toEqual([]);
expect(JSON.stringify(records)).not.toContain("credential-secret");
});
it("keeps safe async UDP error detail while throttling without terminalizing it", async () => {
const terminal: Failure[] = [];
client.on("voiceSendFailed", event => terminal.push(event));
for (let i = 0; i < 100; i++) sdk.logger.warn("udp send error", Object.assign(new Error("credential-secret"), { code: "ECONNREFUSED" }));
await vi.advanceTimersByTimeAsync(2000);
const warnings = records.filter(record => record.level === "warn");
expect(warnings).toHaveLength(2);
expect(warnings[0].args[0]).toMatchObject({ code: "ECONNREFUSED", count: 1 });
expect(warnings[1].args[0]).toMatchObject({ code: "ECONNREFUSED", count: 100 });
expect(JSON.stringify(records)).not.toContain("credential-secret");
expect(terminal).toEqual([]);
});
it("discards unsafe async UDP codes and cancels their summary on disconnect", async () => {
sdk.logger.warn("udp send error", { code: "credential-secret", message: "credential-secret" });
client.disconnect();
await vi.advanceTimersByTimeAsync(3000);
expect(records.filter(record => record.level === "warn").filter(record => String(record.args[1]).includes("udp send error"))).toHaveLength(1);
expect(JSON.stringify(records)).not.toContain("credential-secret");
});
});
+169 -62
View File
@@ -80,6 +80,35 @@ export interface TS3VoiceActivity {
clientUid?: string;
}
/** Send acceptance only: an accepted UDP call does not confirm delivery. */
export interface TS3VoiceSendFailure {
code?: string;
consecutiveFailures: number;
durationMs: number;
}
/** Sticky failure status lets playback stay paused when the terminal event
* already fired during another track or an idle URL lookup. */
export type TS3VoiceSendResult = "accepted" | "retrying" | "failed" | "unavailable";
const VOICE_SEND_FAILURE_TIMEOUT_MS = 2_000;
const SAFE_SOCKET_ERROR_CODES = new Set([
"ERR_SOCKET_DGRAM_NOT_RUNNING", "ERR_SOCKET_DGRAM_NOT_CONNECTED",
"EPIPE", "ENOBUFS", "ECONNREFUSED", "ECONNRESET", "EHOSTUNREACH",
"ENETUNREACH", "ENETDOWN", "EACCES", "EPERM", "EINVAL", "EMSGSIZE",
"EAGAIN", "ENOTCONN", "EBADF",
]);
function safeSocketErrorCode(error: unknown): string | undefined {
try {
if (!error || typeof error !== "object") return undefined;
const code = (error as { code?: unknown }).code;
return typeof code === "string" && SAFE_SOCKET_ERROR_CODES.has(code) ? code : undefined;
} catch {
return undefined;
}
}
// Command notifications and UDP voice packets can be reordered in flight.
// Retain a leaving client's UID briefly so its final packet is still
// attributable; a new clientEnter for the same id cancels and overwrites it.
@@ -117,7 +146,10 @@ export class TS3Client extends EventEmitter {
private detectedProtocol: ServerProtocol = "unknown";
private httpQuery: TS6HttpQuery | null = null;
private udpErrorTimer: ReturnType<typeof setTimeout> | null = null;
private readonly voiceEndpointResolver = new TrackingVoiceEndpointResolver();
private connectionGeneration = 0;
private voiceFailureTimer: ReturnType<typeof setTimeout> | null = null;
private voiceFailure: { code?: string; consecutiveFailures: number; startedAt: number } | null = null;
private voiceEndpointResolver = new TrackingVoiceEndpointResolver();
constructor(private options: TS3ClientOptions, logger: Logger) {
super();
@@ -142,18 +174,25 @@ export class TS3Client extends EventEmitter {
}
async connect(): Promise<void> {
this.voiceEndpointResolver.reset();
const generation = ++this.connectionGeneration;
this.resetVoiceSendState();
this.disconnecting = false;
const voiceEndpointResolver = new TrackingVoiceEndpointResolver();
this.voiceEndpointResolver = voiceEndpointResolver;
this.clearVisibleClientUids();
this.httpQuery = null;
// Clean up any existing connection before creating a new one
if (this.client) {
const previousClient = this.client;
this.client = null;
this.clientId = 0;
if (previousClient) {
this.logger.info("Cleaning up previous connection before reconnecting");
try {
await this.client.disconnect();
await previousClient.disconnect();
} catch {
// Ignore errors during cleanup
}
this.client = null;
this.clientId = 0;
if (generation !== this.connectionGeneration) return;
}
const addr = `${this.options.host}:${this.options.port}`;
@@ -173,6 +212,7 @@ export class TS3Client extends EventEmitter {
3000,
{ ts3QueryPort: 10011, ts6HttpPort: 10080 },
);
if (generation !== this.connectionGeneration) return;
this.detectedProtocol = detection.protocol;
if (this.detectedProtocol === "unknown") {
this.logger.warn(
@@ -198,66 +238,65 @@ export class TS3Client extends EventEmitter {
});
}
// Guard against calling connect() while already connected.
// Save detectedProtocol first because disconnect() resets it.
if (this.client) {
this.logger.warn("connect() called while already connected, disconnecting first");
const savedProtocol = this.detectedProtocol;
const savedHttpQuery = this.httpQuery;
this.disconnect();
this.detectedProtocol = savedProtocol;
this.httpQuery = savedHttpQuery;
// Give the old client a moment to tear down
await new Promise((r) => setTimeout(r, 100));
}
this.logger.info(
{ addr, protocol: this.detectedProtocol },
"Connecting to TeamSpeak server (full client protocol)",
);
// Throttle repeated "udp send error" warnings (fires every 20ms during playback if UDP breaks)
let sdkClient: TS3FullClient | null = null;
const isCurrent = () => sdkClient !== null && this.client === sdkClient &&
this.connectionGeneration === generation && !this.disconnecting;
let udpErrorCount = 0;
let udpErrorCode: string | undefined;
const throttledWarn = (msg: string, ...args: unknown[]) => {
if (!isCurrent()) return;
if (typeof msg === "string" && msg.includes("udp send error")) {
udpErrorCount++;
udpErrorCode = args.map(safeSocketErrorCode).find(code => code !== undefined) ?? udpErrorCode;
if (udpErrorCount === 1) {
this.logger.warn(msg);
this.logger.warn({ ...(udpErrorCode ? { code: udpErrorCode } : {}), count: 1 }, "udp send error");
// After 2 seconds, log a summary and reset.
// Clear any previous timer to avoid leaking it.
if (this.udpErrorTimer) clearTimeout(this.udpErrorTimer);
this.udpErrorTimer = setTimeout(() => {
if (!isCurrent()) return;
if (udpErrorCount > 1) {
this.logger.warn(`udp send error (repeated ${udpErrorCount} times, connection may be lost)`);
this.logger.warn({ ...(udpErrorCode ? { code: udpErrorCode } : {}), count: udpErrorCount }, "Repeated udp send error; connection may be lost");
}
udpErrorCount = 0;
udpErrorCode = undefined;
this.udpErrorTimer = null;
}, 2000);
this.udpErrorTimer.unref?.();
}
return;
}
this.logger.warn(msg);
};
this.client = new TS3FullClient(this.identity, addr, this.options.nickname, {
sdkClient = new TS3FullClient(this.identity, addr, this.options.nickname, {
// Forward server password to the protocol library so it can be
// included in clientinit for password-protected servers
serverPassword: this.options.serverPassword,
resolver: this.voiceEndpointResolver,
resolver: voiceEndpointResolver,
logger: {
debug: (msg) => this.logger.debug(msg),
info: (msg) => this.logger.info(msg),
debug: (msg) => { if (isCurrent()) this.logger.debug(msg); },
info: (msg) => { if (isCurrent()) this.logger.info(msg); },
warn: throttledWarn,
error: (msg) => this.logger.error(msg),
error: (msg) => { if (isCurrent()) this.logger.error(msg); },
},
});
this.client = sdkClient;
this.client.on("textMessage", (msg: TextMessage) => {
sdkClient.on("textMessage", (msg: TextMessage) => {
if (!isCurrent()) return;
if (msg.invokerID === this.clientId) return;
this.emit("textMessage", toTS3TextMessage(msg));
});
this.client.on("voiceData", (voice: VoiceData) => {
sdkClient.on("voiceData", (voice: VoiceData) => {
if (!isCurrent()) return;
// The library normally suppresses our own packets; retain the explicit
// guard so a future protocol change cannot make a bot duck itself.
if (voice.clientId === this.clientId) return;
@@ -270,14 +309,20 @@ export class TS3Client extends EventEmitter {
this.emit("voiceActivity", activity);
});
this.client.on("disconnected", (err) => {
this.logger.warn({ err: err?.message }, "Connection closed");
sdkClient.on("disconnected", (err) => {
if (!isCurrent()) return;
const code = safeSocketErrorCode(err);
this.logger.warn(code ? { code } : {}, "Connection closed");
++this.connectionGeneration;
this.resetVoiceSendState();
this.client = null;
this.clientId = 0;
this.clearVisibleClientUids();
this.emit("disconnected");
});
this.client.on("clientEnter", (info: ClientInfo) => {
sdkClient.on("clientEnter", (info: ClientInfo) => {
if (!isCurrent()) return;
this.rememberVisibleClientUid(info.id, info.uid);
this.logger.debug(
{ nickname: info.nickname, id: info.id },
@@ -286,13 +331,15 @@ export class TS3Client extends EventEmitter {
this.emit("clientEnter", info);
});
this.client.on("clientLeave", (ev: ClientLeftViewEvent) => {
sdkClient.on("clientLeave", (ev: ClientLeftViewEvent) => {
if (!isCurrent()) return;
this.releaseVisibleClientUid(ev.id);
this.logger.debug({ id: ev.id }, "Client left");
this.emit("clientLeave", ev);
});
this.client.on("clientMoved", (ev: ClientMovedEvent) => {
sdkClient.on("clientMoved", (ev: ClientMovedEvent) => {
if (!isCurrent()) return;
this.logger.debug(
{ id: ev.id, targetChannelID: ev.targetChannelID.toString() },
"Client moved"
@@ -300,15 +347,22 @@ export class TS3Client extends EventEmitter {
this.emit("clientMoved", ev);
});
await this.client.connect();
await sdkClient.connect();
if (!isCurrent()) {
// DNS/socket setup may finish after disconnect() closed the transport.
// Close it again without waiting on command replies from a stale session.
sdkClient.handler.close();
return;
}
// Note: @honeybbq/teamspeak-client 0.2.x ships a universal clientinit
// (client_version "3.?.? [Build: 5680278000]" + matching signature)
// that works against both TS3 and TS6 servers. The old 3.6.2 monkey-
// patch on handler.sendPacket was removed when we bumped to 0.2.1 — it
// would have replaced the library's new correct version with a stale
// signature and made TS6 handshakes fail.
await this.client.waitConnected();
this.clientId = this.client.clientID();
await sdkClient.waitConnected();
if (!isCurrent()) return;
this.clientId = sdkClient.clientID();
this.voiceFramesSent = 0;
this.logger.info(
{ clientId: this.clientId, protocol: this.detectedProtocol },
@@ -325,25 +379,32 @@ export class TS3Client extends EventEmitter {
);
}
this.emit("connected");
if (isCurrent()) this.emit("connected");
}
async joinChannel(channelName: string, password?: string): Promise<void> {
if (!this.client) return;
const client = this.client;
if (!client || this.disconnecting) return;
const clientId = this.clientId;
const generation = this.connectionGeneration;
const isCurrent = () => this.client === client && this.clientId === clientId &&
this.connectionGeneration === generation && !this.disconnecting;
const isNumeric = /^\d+$/.test(channelName);
if (isNumeric) {
try {
await clientMove(this.client, this.clientId, BigInt(channelName), password);
await clientMove(client, clientId, BigInt(channelName), password);
if (!isCurrent()) return;
this.logger.info({ channelName }, "Joined channel");
} catch (err) {
this.logger.error({ err, channelName }, "Failed to join channel");
if (isCurrent()) this.logger.error({ err, channelName }, "Failed to join channel");
}
return;
}
try {
const channels = await listChannels(this.client);
const channels = await listChannels(client);
if (!isCurrent()) return;
const channel = channels.find((ch) => ch.name === channelName);
if (!channel) {
@@ -351,13 +412,14 @@ export class TS3Client extends EventEmitter {
return;
}
await clientMove(this.client, this.clientId, channel.id, password);
await clientMove(client, clientId, channel.id, password);
if (!isCurrent()) return;
this.logger.info(
{ channelName, cid: channel.id.toString() },
"Joined channel"
);
} catch (err) {
this.logger.error({ err, channelName }, "Failed to join channel");
if (isCurrent()) this.logger.error({ err, channelName }, "Failed to join channel");
}
}
@@ -454,19 +516,68 @@ export class TS3Client extends EventEmitter {
private voiceFramesSent = 0;
sendVoiceData(opusFrame: Buffer): void {
if (!this.client || this.disconnecting) return;
sendVoiceData(opusFrame: Buffer): TS3VoiceSendResult {
const client = this.client;
if (!client || this.disconnecting) return "unavailable";
try {
this.client.sendVoice(opusFrame, 5);
this.voiceFramesSent++;
if (this.voiceFramesSent === 1) {
this.logger.info({ opusBytes: opusFrame.length, clientId: this.clientId }, "First voice packet sent to TeamSpeak");
}
client.sendVoice(opusFrame, 5);
} catch (err) {
if (this.voiceFramesSent === 0) {
this.logger.error({ err }, "Failed to send first voice packet");
if (this.voiceFailure) {
this.voiceFailure.consecutiveFailures++;
this.voiceFailure.code = safeSocketErrorCode(err) ?? this.voiceFailure.code;
return this.voiceFailureTimer ? "retrying" : "failed";
}
const code = safeSocketErrorCode(err);
const failure = this.voiceFailure = {
...(code ? { code } : {}), consecutiveFailures: 1, startedAt: Date.now(),
};
const generation = this.connectionGeneration;
this.voiceFailureTimer = setTimeout(() => {
if (this.client !== client || this.connectionGeneration !== generation ||
this.disconnecting || this.voiceFailure !== failure) return;
this.voiceFailureTimer = null;
const event = this.voiceFailureDetails(failure);
this.logger.error(event, "Voice sends have failed continuously; pausing playback is required");
this.emit("voiceSendFailed", event);
}, VOICE_SEND_FAILURE_TIMEOUT_MS);
this.voiceFailureTimer.unref?.();
const event = this.voiceFailureDetails(failure);
this.logger.warn(event, "Voice send failed");
this.emit("voiceSendFailure", event);
return "retrying";
}
this.voiceFramesSent++;
if (this.voiceFramesSent === 1) {
this.logger.info({ opusBytes: opusFrame.length, clientId: this.clientId }, "First voice packet accepted by TeamSpeak client");
}
if (this.voiceFailure) {
const event = this.voiceFailureDetails(this.voiceFailure);
this.clearVoiceFailure();
this.logger.info(event, "Voice send recovered");
this.emit("voiceSendRecovered", event);
}
return "accepted";
}
private voiceFailureDetails(failure: NonNullable<TS3Client["voiceFailure"]>): TS3VoiceSendFailure {
return {
...(failure.code ? { code: failure.code } : {}),
consecutiveFailures: failure.consecutiveFailures,
durationMs: Math.max(0, Date.now() - failure.startedAt),
};
}
private clearVoiceFailure(): void {
if (this.voiceFailureTimer) clearTimeout(this.voiceFailureTimer);
this.voiceFailureTimer = null;
this.voiceFailure = null;
}
private resetVoiceSendState(): void {
this.clearVoiceFailure();
this.voiceFramesSent = 0;
if (this.udpErrorTimer) clearTimeout(this.udpErrorTimer);
this.udpErrorTimer = null;
}
getIdentityExport(): string {
@@ -524,24 +635,20 @@ export class TS3Client extends EventEmitter {
}
disconnect(): void {
if (this.client && !this.disconnecting) {
const generation = ++this.connectionGeneration;
this.resetVoiceSendState();
const client = this.client;
this.client = null;
if (client && !this.disconnecting) {
this.disconnecting = true;
const client = this.client;
client.disconnect().catch(() => {}).finally(() => {
if (this.client === client) {
this.client = null;
}
this.disconnecting = false;
if (generation === this.connectionGeneration) this.disconnecting = false;
});
}
this.clientId = 0;
this.clearVisibleClientUids();
this.httpQuery = null;
this.detectedProtocol = "unknown";
if (this.udpErrorTimer) {
clearTimeout(this.udpErrorTimer);
this.udpErrorTimer = null;
}
this.logger.info("Disconnected from TeamSpeak server");
}
}